Skip to content
This repository was archived by the owner on Jun 16, 2026. It is now read-only.

Commit be3ac76

Browse files
authored
Merge branch 'main' into copilot/sub-pr-72
Signed-off-by: Adam Poulemanos <89049923+bashandbone@users.noreply.github.com>
2 parents b729d48 + 074eb3e commit be3ac76

12 files changed

Lines changed: 457 additions & 42 deletions

File tree

.github/workflows/gemini-review.yml

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -49,6 +49,7 @@ jobs:
4949
- name: 'Run Gemini pull request review'
5050
uses: 'google-github-actions/run-gemini-cli@v0' # ratchet:exclude
5151
id: 'gemini_pr_review'
52+
continue-on-error: true
5253
env:
5354
GITHUB_TOKEN: '${{ steps.mint_identity_token.outputs.token || secrets.GITHUB_TOKEN || github.token }}'
5455
ISSUE_TITLE: '${{ github.event.pull_request.title || github.event.issue.title }}'

crates/recoco-core/src/base/value.rs

Lines changed: 47 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1767,4 +1767,51 @@ mod tests {
17671767
let size = value.estimated_byte_size();
17681768
assert_eq!(size, std::mem::size_of::<Value<ScopeValue>>());
17691769
}
1770+
1771+
#[test]
1772+
fn test_range_value_len() {
1773+
let range1 = RangeValue { start: 0, end: 5 };
1774+
assert_eq!(range1.len(), 5);
1775+
1776+
let range2 = RangeValue { start: 5, end: 5 };
1777+
assert_eq!(range2.len(), 0);
1778+
1779+
let range3 = RangeValue { start: 10, end: 25 };
1780+
assert_eq!(range3.len(), 15);
1781+
1782+
#[test]
1783+
fn test_key_part_to_strs() {
1784+
let bytes_part = KeyPart::from(vec![1u8, 2, 3]);
1785+
assert_eq!(bytes_part.to_strs(), vec!["AQID"]);
1786+
1787+
let str_part = KeyPart::from(String::from("hello"));
1788+
assert_eq!(str_part.to_strs(), vec!["hello"]);
1789+
1790+
let bool_part = KeyPart::from(true);
1791+
assert_eq!(bool_part.to_strs(), vec!["true"]);
1792+
1793+
let int64_part = KeyPart::from(42i64);
1794+
assert_eq!(int64_part.to_strs(), vec!["42"]);
1795+
1796+
let range_part = KeyPart::from(RangeValue::new(10, 20));
1797+
assert_eq!(range_part.to_strs(), vec!["10", "20"]);
1798+
1799+
let uuid_val = uuid::Uuid::nil();
1800+
let uuid_part = KeyPart::from(uuid_val);
1801+
assert_eq!(uuid_part.to_strs(), vec![uuid_val.to_string()]);
1802+
1803+
let date_val = chrono::NaiveDate::from_ymd_opt(2023, 10, 15)
1804+
.expect("test date 2023-10-15 should be a valid NaiveDate");
1805+
let date_part = KeyPart::from(date_val);
1806+
assert_eq!(date_part.to_strs(), vec!["2023-10-15"]);
1807+
1808+
let struct_part = KeyPart::from(vec![
1809+
KeyPart::from(String::from("world")),
1810+
KeyPart::from(100i64),
1811+
KeyPart::from(vec![
1812+
KeyPart::from(false),
1813+
]),
1814+
]);
1815+
assert_eq!(struct_part.to_strs(), vec!["world", "100", "false"]);
1816+
}
17701817
}

crates/recoco-core/src/ops/sources/local_file.rs

Lines changed: 53 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,7 @@ pub struct Spec {
3232

3333
struct Executor {
3434
root_path: PathBuf,
35+
canonical_root_path: Option<PathBuf>,
3536
binary: bool,
3637
pattern_matcher: PatternMatcher,
3738
max_file_size: Option<i64>,
@@ -134,14 +135,60 @@ impl SourceExecutor for Executor {
134135
options: &SourceExecutorReadOptions,
135136
) -> Result<PartialSourceRowData> {
136137
let path = key.single_part()?.str_value()?.as_ref();
137-
if !self.pattern_matcher.is_file_included(path) {
138+
let path_obj = Path::new(path);
139+
140+
// Prevent path traversal vulnerabilities by verifying the path
141+
// doesn't contain parent directory or absolute components.
142+
if path_obj.components().any(|c| {
143+
matches!(
144+
c,
145+
std::path::Component::ParentDir
146+
| std::path::Component::RootDir
147+
| std::path::Component::Prefix(_)
148+
)
149+
}) || !self.pattern_matcher.is_file_included(path)
150+
{
138151
return Ok(PartialSourceRowData {
139152
value: Some(SourceValue::NonExistence),
140153
ordinal: Some(Ordinal::unavailable()),
141154
content_version_fp: None,
142155
});
143156
}
157+
144158
let path = self.root_path.join(path);
159+
160+
// Mitigate symlink-based path traversal by canonicalizing and checking boundaries
161+
if let Some(root_canon) = &self.canonical_root_path {
162+
let path_canon = match tokio::fs::canonicalize(&path).await {
163+
Ok(c) => c,
164+
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
165+
// Target file doesn't exist.
166+
return Ok(PartialSourceRowData {
167+
value: Some(SourceValue::NonExistence),
168+
ordinal: Some(Ordinal::unavailable()),
169+
content_version_fp: None,
170+
});
171+
}
172+
Err(e) => Err(e)?,
173+
};
174+
175+
if !path_canon.starts_with(root_canon) {
176+
// Symlink points outside the allowed root directory.
177+
return Ok(PartialSourceRowData {
178+
value: Some(SourceValue::NonExistence),
179+
ordinal: Some(Ordinal::unavailable()),
180+
content_version_fp: None,
181+
});
182+
}
183+
} else {
184+
// Root doesn't exist (failed to canonicalize during setup), so the file cannot exist.
185+
return Ok(PartialSourceRowData {
186+
value: Some(SourceValue::NonExistence),
187+
ordinal: Some(Ordinal::unavailable()),
188+
content_version_fp: None,
189+
});
190+
}
191+
145192
let mut metadata: Option<Metadata> = None;
146193
// Check file size limit
147194
if let Some(max_size) = self.max_file_size
@@ -237,8 +284,12 @@ impl SourceFactoryBase for Factory {
237284
spec: Spec,
238285
_context: Arc<FlowInstanceContext>,
239286
) -> Result<Box<dyn SourceExecutor>> {
287+
let root_path = PathBuf::from(spec.path);
288+
let canonical_root_path = tokio::fs::canonicalize(&root_path).await.ok();
289+
240290
Ok(Box::new(Executor {
241-
root_path: PathBuf::from(spec.path),
291+
root_path,
292+
canonical_root_path,
242293
binary: spec.binary,
243294
pattern_matcher: PatternMatcher::new(spec.included_patterns, spec.excluded_patterns)?,
244295
max_file_size: spec.max_file_size,

crates/recoco-core/src/ops/targets/postgres.rs

Lines changed: 50 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -287,22 +287,59 @@ impl ExportContext {
287287
deletions: &[interface::ExportTargetDeleteEntry],
288288
txn: &mut sqlx::PgTransaction<'_>,
289289
) -> Result<()> {
290-
// TODO: Find a way to batch delete.
291-
for deletion in deletions.iter() {
290+
if deletions.is_empty() {
291+
return Ok(());
292+
}
293+
294+
let num_parameters = self.key_fields_schema.len();
295+
if num_parameters == 0 {
296+
return Ok(());
297+
}
298+
299+
for deletion_chunk in deletions.chunks(BIND_LIMIT / num_parameters) {
292300
let mut query_builder = sqlx::QueryBuilder::new("");
293301
query_builder.push(&self.delete_sql_prefix);
294-
for (i, ((schema, _), value)) in
295-
std::iter::zip(&self.key_fields_schema, &deletion.key).enumerate()
296-
{
302+
303+
if num_parameters > 1 {
304+
query_builder.push("(");
305+
}
306+
for (i, (schema, _)) in self.key_fields_schema.iter().enumerate() {
297307
if i > 0 {
298-
query_builder.push(" AND ");
308+
query_builder.push(", ");
299309
}
300310
query_builder.push("\"");
301311
query_builder.push(schema.name.as_str());
302312
query_builder.push("\"");
303-
query_builder.push("=");
304-
bind_key_field(&mut query_builder, value)?;
305313
}
314+
if num_parameters > 1 {
315+
query_builder.push(")");
316+
}
317+
318+
query_builder.push(" IN (");
319+
320+
for (i, deletion) in deletion_chunk.iter().enumerate() {
321+
if i > 0 {
322+
query_builder.push(", ");
323+
}
324+
if num_parameters > 1 {
325+
query_builder.push("(");
326+
}
327+
for (j, (_schema, _)) in self.key_fields_schema.iter().enumerate() {
328+
if j > 0 {
329+
query_builder.push(", ");
330+
}
331+
if let Some(value) = deletion.key.get(j) {
332+
bind_key_field(&mut query_builder, value)?;
333+
} else {
334+
query_builder.push("NULL");
335+
}
336+
}
337+
if num_parameters > 1 {
338+
query_builder.push(")");
339+
}
340+
}
341+
342+
query_builder.push(")");
306343
query_builder.build().execute(&mut **txn).await?;
307344
}
308345
Ok(())
@@ -993,9 +1030,11 @@ impl AttachmentSetupChange for SqlCommandSetupChange {
9931030
}
9941031

9951032
async fn apply_change(&self) -> Result<()> {
996-
for teardown_sql in self.teardown_sql_to_run.iter() {
997-
sqlx::raw_sql(teardown_sql).execute(&self.db_pool).await?;
998-
}
1033+
let teardown_futs = self
1034+
.teardown_sql_to_run
1035+
.iter()
1036+
.map(|teardown_sql| sqlx::raw_sql(teardown_sql).execute(&self.db_pool));
1037+
futures::future::try_join_all(teardown_futs).await?;
9991038
if let Some(setup_sql) = &self.setup_sql_to_run {
10001039
sqlx::raw_sql(setup_sql).execute(&self.db_pool).await?;
10011040
}

crates/recoco-core/src/ops/targets/qdrant.rs

Lines changed: 40 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,9 @@ use qdrant_client::qdrant::{
2626
const DEFAULT_VECTOR_SIMILARITY_METRIC: spec::VectorSimilarityMetric =
2727
spec::VectorSimilarityMetric::CosineSimilarity;
2828
const DEFAULT_URL: &str = "http://localhost:6334/";
29+
/// Maximum number of setup operations (deletes or creates) to run concurrently.
30+
/// Bounds the number of simultaneous Qdrant client requests during collection setup.
31+
const MAX_CONCURRENT_SETUP_OPS: usize = 16;
2932

3033
////////////////////////////////////////////////////////////
3134
// Public Types
@@ -577,23 +580,45 @@ impl TargetFactoryBase for Factory {
577580
setup_change: Vec<TypedResourceSetupChangeItem<'async_trait, Self>>,
578581
context: Arc<FlowInstanceContext>,
579582
) -> Result<()> {
580-
for setup_change in setup_change.iter() {
581-
let qdrant_client =
582-
self.get_qdrant_client(&setup_change.key.connection, &context.auth_registry)?;
583-
setup_change
584-
.setup_change
585-
.apply_delete(&setup_change.key.collection_name, &qdrant_client)
586-
.await?;
583+
let auth_registry = &context.auth_registry;
584+
// Each item in setup_change has a unique CollectionKey (enforced by the setup
585+
// orchestration layer), so concurrent operations across items are race-free.
586+
587+
// Delete phase: collect futures first, then run with bounded concurrency.
588+
let mut delete_futures = Vec::with_capacity(setup_change.len());
589+
for change in setup_change.iter() {
590+
let client_result = self.get_qdrant_client(&change.key.connection, auth_registry);
591+
delete_futures.push(async move {
592+
let client = client_result?;
593+
change
594+
.setup_change
595+
.apply_delete(&change.key.collection_name, &client)
596+
.await
597+
});
587598
}
588-
for setup_change in setup_change.iter() {
589-
let qdrant_client =
590-
self.get_qdrant_client(&setup_change.key.connection, &context.auth_registry)?;
591-
setup_change
592-
.setup_change
593-
.apply_create(&setup_change.key.collection_name, &qdrant_client)
594-
.await?;
599+
futures::stream::iter(delete_futures)
600+
.buffer_unordered(MAX_CONCURRENT_SETUP_OPS)
601+
.try_collect::<Vec<_>>()
602+
.await
603+
.map(|_| ())?;
604+
605+
// Create phase: collect futures first, then run with bounded concurrency.
606+
let mut create_futures = Vec::with_capacity(setup_change.len());
607+
for change in setup_change.iter() {
608+
let client_result = self.get_qdrant_client(&change.key.connection, auth_registry);
609+
create_futures.push(async move {
610+
let client = client_result?;
611+
change
612+
.setup_change
613+
.apply_create(&change.key.collection_name, &client)
614+
.await
615+
});
595616
}
596-
Ok(())
617+
futures::stream::iter(create_futures)
618+
.buffer_unordered(MAX_CONCURRENT_SETUP_OPS)
619+
.try_collect::<Vec<_>>()
620+
.await
621+
.map(|_| ())
597622
}
598623
}
599624

crates/recoco-core/src/setup/components.rs

Lines changed: 19 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -156,27 +156,35 @@ impl<D: SetupOperator + Send + Sync> ResourceSetupChange for SetupChange<D> {
156156
}
157157
}
158158

159+
/// Maximum number of component operations (deletes or upserts) that may run concurrently.
160+
/// Keeping this bounded prevents overwhelming a database connection pool or
161+
/// network layer when a large number of components change at once.
162+
const COMPONENT_CONCURRENCY_LIMIT: usize = 16;
163+
159164
pub async fn apply_component_changes<D: SetupOperator>(
160165
changes: Vec<&SetupChange<D>>,
161166
context: &D::Context,
162167
) -> Result<()> {
163168
// First delete components that need to be removed
164-
for change in changes.iter() {
165-
for key in &change.keys_to_delete {
166-
change.desc.delete(key, context).await?;
167-
}
168-
}
169+
let delete_futures = changes.iter().flat_map(|change| {
170+
change
171+
.keys_to_delete
172+
.iter()
173+
.map(move |key| change.desc.delete(key, context))
174+
});
175+
futures::future::try_join_all(delete_futures).await?;
169176

170177
// Then upsert components that need to be updated
171-
for change in changes.iter() {
172-
for state in &change.states_to_upsert {
178+
let upsert_futures = changes.iter().flat_map(|change| {
179+
change.states_to_upsert.iter().map(move |state| async move {
173180
if state.already_exists {
174-
change.desc.update(&state.state, context).await?;
181+
change.desc.update(&state.state, context).await
175182
} else {
176-
change.desc.create(&state.state, context).await?;
183+
change.desc.create(&state.state, context).await
177184
}
178-
}
179-
}
185+
})
186+
});
187+
futures::future::try_join_all(upsert_futures).await?;
180188

181189
Ok(())
182190
}

0 commit comments

Comments
 (0)