Repository navigation
feat: implement parallel foreach spec - #2317
Conversation
|
No actionable comments were generated in the recent review. 🎉 ℹ️ Recent review info⚙️ Run configurationConfiguration used: Organization UI Review profile: CHILL Plan: Pro Run ID: 📒 Files selected for processing (1)
📝 WalkthroughWalkthroughImplements Spec 018: ChangesSpec 018: Parallel Fan-Out and Foreach Iteration
Sequence Diagram(s)sequenceDiagram
rect rgba(100, 149, 237, 0.5)
Note over foreachExecutor,aggregateWriter: Foreach Executor Runtime
end
participant foreachExecutor as foreach Executor
participant expandItems as expandItems
participant semaphore as semaphore (MaxConcurrent)
participant runItem as runItem[i]
participant runner as runtime Runner
participant collectOutputs as collectOutputs
participant writeAggregate as writeAggregate
foreachExecutor->>expandItems: Foreach.Items / ItemsExpr
expandItems-->>foreachExecutor: []expandedItem (key, value, index)
loop per item (bounded concurrency)
foreachExecutor->>semaphore: acquire slot
foreachExecutor->>runItem: item[i]
runItem->>runner: NewPlanEnv (ctx with Foreach scope) + Run cloned steps
runner-->>runItem: plan result / error
runItem->>collectOutputs: resolve Collect expressions via plan env
collectOutputs-->>runItem: map[string]string outputs
runItem-->>foreachExecutor: itemResult{status, outputs}
foreachExecutor->>semaphore: release slot
end
foreachExecutor->>writeAggregate: results[]
writeAggregate-->>foreachExecutor: JSON aggregate written to stdout
Estimated code review effort🎯 5 (Critical) | ⏱️ ~120 minutes Possibly related PRs
🚥 Pre-merge checks | ✅ 4 | ❌ 1❌ Failed checks (1 warning)
✅ Passed checks (4 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches📝 Generate docstrings
🧪 Generate unit tests (beta)
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
There was a problem hiding this comment.
Actionable comments posted: 3
🧹 Nitpick comments (1)
internal/cmn/schema/dag_schema_spec018_test.go (1)
186-186: 🧹 Nitpick | 🔵 TrivialAlign test schema resolution options with production resolver.
Line 186 resolves with empty options, but production schema resolution enables
ValidateDefaults: true. Update the test helper to match:Required change
resolved, err := schema.Resolve(&jsonschema.ResolveOptions{ ValidateDefaults: true, })This mismatch can let schema tests pass while runtime resolution fails on invalid defaults.
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@internal/cmn/schema/dag_schema_spec018_test.go` at line 186, The schema.Resolve() call in the test is using empty jsonschema.ResolveOptions which does not validate defaults, but the production code enables ValidateDefaults: true. Update the ResolveOptions passed to the Resolve method to include ValidateDefaults: true so that the test behavior matches production behavior and can catch invalid defaults during test execution rather than at runtime.
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@internal/cmn/schema/dag.schema.json`:
- Around line 2080-2083: The foreach.items schema definition at line 2082
currently references stepOutputJSONValue which permits nested objects, arrays,
and null values, but Spec 018 specifies that foreach.items should only accept
scalar values. Replace the $ref from stepOutputJSONValue with a reference to a
more restrictive definition that only allows scalar types (string, number,
boolean, integer) to align the schema validation with the downstream logic
constraints.
In `@internal/core/validator.go`:
- Around line 295-300: The validateForeachConfig function currently validates
body step dependencies but does not check the approval.rewind_to target
semantics for foreach body steps, allowing invalid or out-of-scope rewind
targets to pass validation. Enhance the validateForeachConfig function to
validate that any approval.rewind_to targets in foreach body steps reference
valid step names from the bodyNames set and exist within the foreach scope using
the already-available bodyNames and bodyIDs parameters for validation.
In `@internal/runtime/builtin/foreach/foreach.go`:
- Around line 194-207: The function in foreach.go can return early when context
is cancelled (in the select statement checking ctx.Done()) without waiting for
in-flight worker goroutines to complete, causing the running goroutines to
mutate the results map after the function has already returned. This creates a
race condition with aggregate marshaling. Fix this by ensuring wg.Wait() is
called before any return statement, including the early return on context
cancellation. Restructure the control flow so that the function waits for all
workers to complete their work before returning results, regardless of whether
the context was cancelled.
---
Nitpick comments:
In `@internal/cmn/schema/dag_schema_spec018_test.go`:
- Line 186: The schema.Resolve() call in the test is using empty
jsonschema.ResolveOptions which does not validate defaults, but the production
code enables ValidateDefaults: true. Update the ResolveOptions passed to the
Resolve method to include ValidateDefaults: true so that the test behavior
matches production behavior and can catch invalid defaults during test execution
rather than at runtime.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: CHILL
Plan: Pro
Run ID: 5bae3f60-e54c-4f60-9d1b-5c41cc4d632a
📒 Files selected for processing (50)
conformance/spec018_parallel_foreach/parallel_foreach_test.goconformance/spec018_parallel_foreach/testdata/foreach_duplicate_keys.yamlconformance/spec018_parallel_foreach/testdata/foreach_non_json_string_items.yamlconformance/spec018_parallel_foreach/testdata/foreach_string_items.yamlconformance/spec018_parallel_foreach/testdata/foreach_success_collect.yamlconformance/spec018_parallel_foreach/testdata/foreach_zero_items.yamlconformance/spec018_parallel_foreach/testdata/parallel_dag_run_object_items.yamlconformance/spec018_parallel_foreach/testdata/parallel_flat_mapping_valid.yamlconformance/spec018_parallel_foreach/testdata/parallel_max_concurrent_float.yamlconformance/spec018_parallel_foreach/testdata/parallel_max_concurrent_too_high.yamlconformance/spec018_parallel_foreach/testdata/parallel_nested_array_item.yamlconformance/spec018_parallel_foreach/testdata/parallel_nested_mapping_item.yamlconformance/spec018_parallel_foreach/testdata/parallel_string_items.yamlconformance/spec018_parallel_foreach/testdata/parallel_unknown_object_field.yamlinternal/cmn/schema/dag.schema.jsoninternal/cmn/schema/dag_schema_spec018_test.gointernal/cmn/value/resolver.gointernal/cmn/value/scope.gointernal/cmn/value/template.gointernal/cmn/value/template_test.gointernal/core/foreach.gointernal/core/foreach_validation_spec018_test.gointernal/core/parallel.gointernal/core/parallel_test.gointernal/core/spec/foreach_spec018_test.gointernal/core/spec/parallel_spec018_test.gointernal/core/spec/step.gointernal/core/spec/step_test.gointernal/core/spec/step_types.gointernal/core/step.gointernal/core/validator.gointernal/core/validator_test.gointernal/core/value_fields.gointernal/core/value_fields_test.gointernal/core/value_notices.gointernal/runtime/builtin/builtin.gointernal/runtime/builtin/dag/enqueue.gointernal/runtime/builtin/dag/enqueue_test.gointernal/runtime/builtin/foreach/foreach.gointernal/runtime/env.gointernal/runtime/eval.gointernal/runtime/foreach_spec018_test.gointernal/runtime/runner.gospecs/003-value-resolution.mdspecs/005-value-resolution-params.mdspecs/006-value-resolution-env.mdspecs/007-value-resolution-steps.mdspecs/009-step-reference.mdspecs/018-parallel-and-foreach.mdspecs/README.md
|
@coderabbitai review |
✅ Action performedReview finished.
|
|
@coderabbitai review |
✅ Action performedReview finished.
|
|
@coderabbitai review |
✅ Action performedReview finished.
|
Summary:
Validation:
Summary by cubic
Implements Spec 018 for parallel fan‑out and foreach iteration. Adds a
foreachruntime with item-scoped${foreach.*}bindings and aggregatecollect, tightensparallel.max_concurrentto 1–1000 with item shape checks, improvesdag.enqueueconcurrency, preserves child DAG env scope, and stabilizes the approval inputs test.New Features
foreachexecutor: iterateitems(array or JSON string) with optionalas/key,max_concurrent, andcollectfor aggregate output.${foreach.index},${foreach.key}, and${foreach.item}(or alias) to body steps andcollect; propagateforeachscope through env/resolver; preserve child DAG env scope over inherited step env.parallel.max_concurrentin 1–1000; object item values must be scalars;dag.enqueuehonors sequential/parallel concurrency; validate body scoping (deps by id/name within body, unique ids/names), and suppress out‑of‑scopeforeachnotices.Migration
parallel.max_concurrentmust be an integer from 1 to 1000.foreach.itemsis a string, it must be a valid JSON array.parallel/foreachmust be scalars (string/number/boolean), not nested objects/arrays.foreach; consume the aggregate output instead.keyis set forforeach, item keys must be unique.Written for commit a34adbe. Summary will update on new commits.
Summary by CodeRabbit
foreachstep construct for iterating over item sources and running an inline step body with concurrency control.${foreach.*}template bindings withinforeachbodies (including item fields and per-item key/index).collectto aggregateforeachresults into consolidated outputs.parallel.max_concurrentvalidation (integer-only, range 1–1000) and improved validation/runtime handling for invalidforeachconfigurations and duplicate keys.foreachin the workflow schema and value/runner scoping.