Skip to content

feat: implement parallel foreach spec - #2317

Merged
yohamta0 merged 8 commits into
mainfrom
implement-spec-018-parallel-foreach
Jun 24, 2026
Merged

yohamta0 merged 8 commits into
mainfrom
implement-spec-018-parallel-foreach

Conversation

@yohamta0

@yohamta0 yohamta0 commented Jun 22, 2026 •

Copy link
Copy Markdown
Member

Summary:

  • Implement Spec 018 parallel and foreach behavior.
  • Add foreach runtime, typed item scope, schema/parser validation, and conformance coverage.
  • Update related value-resolution and step-reference specs.

Validation:

  • go test ./internal/core ./internal/core/spec ./internal/cmn/schema ./internal/cmn/value ./internal/runtime ./internal/runtime/builtin/dag ./internal/runtime/builtin/foreach -count=1
  • make conformance CONFORMANCE_TEST_TARGET=./conformance/spec018_parallel_foreach

Summary by cubic

Implements Spec 018 for parallel fan‑out and foreach iteration. Adds a foreach runtime with item-scoped ${foreach.*} bindings and aggregate collect, tightens parallel.max_concurrent to 1–1000 with item shape checks, improves dag.enqueue concurrency, preserves child DAG env scope, and stabilizes the approval inputs test.

  • New Features

    • Add foreach executor: iterate items (array or JSON string) with optional as/key, max_concurrent, and collect for aggregate output.
    • Expose ${foreach.index}, ${foreach.key}, and ${foreach.item} (or alias) to body steps and collect; propagate foreach scope through env/resolver; preserve child DAG env scope over inherited step env.
    • Enforce bounds and shapes: parallel.max_concurrent in 1–1000; object item values must be scalars; dag.enqueue honors sequential/parallel concurrency; validate body scoping (deps by id/name within body, unique ids/names), and suppress out‑of‑scope foreach notices.
  • Migration

    • parallel.max_concurrent must be an integer from 1 to 1000.
    • When foreach.items is a string, it must be a valid JSON array.
    • Object item values for parallel/foreach must be scalars (string/number/boolean), not nested objects/arrays.
    • Body step outputs are not addressable from outside foreach; consume the aggregate output instead.
    • When key is set for foreach, item keys must be unique.

Written for commit a34adbe. Summary will update on new commits.

Review in cubic

Summary by CodeRabbit

  • New Features
    • Added foreach step construct for iterating over item sources and running an inline step body with concurrency control.
    • Added ${foreach.*} template bindings within foreach bodies (including item fields and per-item key/index).
    • Added collect to aggregate foreach results into consolidated outputs.
  • Improvements
    • Enforced stricter parallel.max_concurrent validation (integer-only, range 1–1000) and improved validation/runtime handling for invalid foreach configurations and duplicate keys.
    • Added support for foreach in the workflow schema and value/runner scoping.
  • Documentation
    • Marked the parallel/foreach spec as implemented in the specs overview.

@coderabbitai

coderabbitai Bot commented Jun 22, 2026 •

Copy link
Copy Markdown

Review Change Stack

No actionable comments were generated in the recent review. 🎉

ℹ️ Recent review info
⚙️ Run configuration

Configuration used: Organization UI

Review profile: CHILL

Plan: Pro

Run ID: 9d5c9a16-70c3-4c0b-b0e3-099cdc58c4c1

📥 Commits

Reviewing files that changed from the base of the PR and between d4f60cc and a34adbe.

📒 Files selected for processing (1)
  • internal/service/frontend/api/v1/dagruns_test.go

📝 Walkthrough

Walkthrough

Implements Spec 018: parallel fan-out and foreach inline iteration for Dagu workflow steps. Adds the ForeachConfig data model, JSON schema validation, YAML parser with mutual-exclusion enforcement, ${foreach.*} value-resolution namespace, body-step isolation validator, a full concurrent foreach executor, bounded-concurrency parallel enqueue, environment propagation for foreach item scope, conformance tests, and specification documents.

Changes

Spec 018: Parallel Fan-Out and Foreach Iteration

Layer / File(s) Summary
Core data types and constants
internal/core/foreach.go, internal/core/step.go, internal/core/parallel.go, internal/core/spec/step_types.go
ForeachConfig struct with Items, ItemsExpr, As, Key, MaxConcurrent, Steps, and Collect fields; Foreach *ForeachConfig field on Step; ExecutorTypeForeach constant; MaxExpansionConcurrency = 1000; and "foreach" registered in builtinStepTypeNames.
JSON schema and YAML spec parsing
internal/cmn/schema/dag.schema.json, internal/cmn/schema/dag_schema_spec018_test.go, internal/core/spec/step.go, internal/core/spec/step_test.go, internal/core/spec/foreach_spec018_test.go, internal/core/spec/parallel_spec018_test.go
dag.schema.json adds foreach block schema, mutual-exclusion if/then rules, scalar-only additionalProperties for parallel items, and foreach in builtinExecutorType; spec/step.go adds Foreach YAML field and buildStepForeach with integer-only max_concurrent range validation (1–1000); schema and spec tests validate both valid and invalid parallel/foreach configurations.
Foreach value resolution namespace
internal/cmn/value/scope.go, internal/cmn/value/resolver.go, internal/cmn/value/template.go, internal/cmn/value/template_test.go
RuntimeScope gains a Foreach Values field; bindingScope() activates the runtime scope when Foreach is non-nil; bindingValue dispatches "foreach" paths to bindingForeachValue (field extraction + JSON serialization for composites); supportedStrictBinding validates 2–3 segment foreach paths; template tests verify both resolved and namespace-unavailable reference behavior.
Validator, value field walker, and value notices
internal/core/validator.go, internal/core/validator_test.go, internal/core/parallel_test.go, internal/core/value_fields.go, internal/core/value_fields_test.go, internal/core/value_notices.go, internal/core/foreach_validation_spec018_test.go
ValidateSteps calls resolveForeachStepDependencies and validateForeachConfig for body-step isolation, collision detection, and dependency containment; parallel.max_concurrent clamped to 1–MaxExpansionConcurrency with updated error text; walkStep delegates to walkForeach for reference-field enumeration; valueReferenceNoticeFieldSink suppresses ${foreach.*} notices for item-scope fields; comprehensive validation and validator tests.
Parallel bounded-concurrency enqueue
internal/runtime/builtin/dag/enqueue.go, internal/runtime/builtin/dag/enqueue_test.go
enqueueExecutor.Run delegates to enqueueAll → enqueueSequential or enqueueParallel; enqueueParallel uses a semaphore channel and cancellable context to bound concurrent enqueues at MaxConcurrent while preserving result order; TestEnqueueExecutorParallelHonorsMaxConcurrent validates the bound with a concurrency-recording queue store wrapper.
Foreach runtime executor and env propagation
internal/runtime/builtin/foreach/foreach.go, internal/runtime/builtin/builtin.go, internal/runtime/env.go, internal/runtime/eval.go, internal/runtime/runner.go
Full foreach executor: item expansion from static array or JSON string, per-item keying, semaphore-bounded concurrent body execution with cloned steps, foreach.collect output resolution, and JSON aggregate output; Env gains a Foreach field propagated via newEnv context inheritance, resolverFromEnv, and inherited StepMap in plan-env constructors; registered via blank import.
Foreach runtime tests
internal/runtime/foreach_spec018_test.go, internal/runtime/env_test.go
TestForeachRuntimeRunsBodyAndPublishesAggregate verifies per-item outputs and aggregate JSON totals with probe executor; TestForeachRuntimeHonorsMaxConcurrent verifies the concurrency bound via probe executor with configurable delay and mutex-tracked active count; TestNewEnvUsesDAGScopeWhenContextHasInheritedEnv validates env inheritance preserving foreach scope.
Specification documents
specs/018-parallel-and-foreach.md, specs/README.md, specs/003-value-resolution.md, specs/005-value-resolution-params.md, specs/006-value-resolution-env.md, specs/007-value-resolution-steps.md, specs/009-step-reference.md
New Spec 018 (796 lines) documents parallel and foreach field shapes, item-scope resolution, body scoping, aggregate output schemas, validation/runtime error conditions, and YAML examples; existing specs extended with foreach namespace, field-matrix rows, body-scope rules, and step-reference scoping constraints.
Conformance tests and testdata
conformance/spec018_parallel_foreach/parallel_foreach_test.go, conformance/spec018_parallel_foreach/testdata/*
TestParallelValidation, TestParallelDAGRunRuntime, TestForeachValidation, and TestForeachRuntime cover CLI validate and start flows; testdata fixtures cover valid parallel flat mapping, invalid nested items, max_concurrent bounds, foreach success collect, string/zero/duplicate-key/non-JSON-string item scenarios.
API approval test update
internal/service/frontend/api/v1/dagruns_test.go
TestApproveDAGRunStepWithInputs extended with an additional approval gate (hold-step) and validates intermediate state persistence without immediate DAG completion; approval inputs and output variable transformation verified.

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
Loading

Estimated code review effort

🎯 5 (Critical) | ⏱️ ~120 minutes

Possibly related PRs

  • dagucloud/dagu#2302: Both PRs modify the shared value-reference notice system in internal/cmn/value and internal/core/value_notices.go; this PR adds ${foreach...} suppression behavior while the related PR extends ValueReferenceNotice with a Reason field and step-output notice reporting.

  • dagucloud/dagu#2306: Both PRs extend the same value-resolution pipeline by modifying shared runtime scoping and template binding dispatch in internal/cmn/value to add new template namespaces (${foreach...} vs ${run...}/${paths...} built-in run context).

🚥 Pre-merge checks | ✅ 4 | ❌ 1

❌ Failed checks (1 warning)

Check name Status Explanation Resolution
Docstring Coverage ⚠️ Warning Docstring coverage is 8.89% which is insufficient. The required threshold is 80.00%. Write docstrings for the functions missing them to satisfy the coverage threshold.
✅ Passed checks (4 passed)
Check name Status Explanation
Title check ✅ Passed The title 'feat: implement parallel foreach spec' clearly summarizes the main change - implementing Spec 018 for parallel and foreach functionality.
Description check ✅ Passed The PR description includes a summary, detailed validation steps, and comprehensive documentation of new features and migration requirements, covering all essential information despite not following the template structure exactly.
Linked Issues check ✅ Passed Check skipped because no linked issues were found for this pull request.
Out of Scope Changes check ✅ Passed Check skipped because no linked issues were found for this pull request.

✏️ Tip: You can configure your own custom pre-merge checks in the settings.

✨ Finishing Touches
📝 Generate docstrings
  • Create stacked PR
  • Commit on current branch
🧪 Generate unit tests (beta)
  • Create PR with unit tests
  • Commit unit tests in branch implement-spec-018-parallel-foreach

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.

❤️ Share

Comment @coderabbitai help to get the list of available commands and usage tips.

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Actionable comments posted: 3

🧹 Nitpick comments (1)
internal/cmn/schema/dag_schema_spec018_test.go (1)

186-186: 🧹 Nitpick | 🔵 Trivial

Align 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

📥 Commits

Reviewing files that changed from the base of the PR and between 2f08f01 and 44547cf.

📒 Files selected for processing (50)
  • conformance/spec018_parallel_foreach/parallel_foreach_test.go
  • conformance/spec018_parallel_foreach/testdata/foreach_duplicate_keys.yaml
  • conformance/spec018_parallel_foreach/testdata/foreach_non_json_string_items.yaml
  • conformance/spec018_parallel_foreach/testdata/foreach_string_items.yaml
  • conformance/spec018_parallel_foreach/testdata/foreach_success_collect.yaml
  • conformance/spec018_parallel_foreach/testdata/foreach_zero_items.yaml
  • conformance/spec018_parallel_foreach/testdata/parallel_dag_run_object_items.yaml
  • conformance/spec018_parallel_foreach/testdata/parallel_flat_mapping_valid.yaml
  • conformance/spec018_parallel_foreach/testdata/parallel_max_concurrent_float.yaml
  • conformance/spec018_parallel_foreach/testdata/parallel_max_concurrent_too_high.yaml
  • conformance/spec018_parallel_foreach/testdata/parallel_nested_array_item.yaml
  • conformance/spec018_parallel_foreach/testdata/parallel_nested_mapping_item.yaml
  • conformance/spec018_parallel_foreach/testdata/parallel_string_items.yaml
  • conformance/spec018_parallel_foreach/testdata/parallel_unknown_object_field.yaml
  • internal/cmn/schema/dag.schema.json
  • internal/cmn/schema/dag_schema_spec018_test.go
  • internal/cmn/value/resolver.go
  • internal/cmn/value/scope.go
  • internal/cmn/value/template.go
  • internal/cmn/value/template_test.go
  • internal/core/foreach.go
  • internal/core/foreach_validation_spec018_test.go
  • internal/core/parallel.go
  • internal/core/parallel_test.go
  • internal/core/spec/foreach_spec018_test.go
  • internal/core/spec/parallel_spec018_test.go
  • internal/core/spec/step.go
  • internal/core/spec/step_test.go
  • internal/core/spec/step_types.go
  • internal/core/step.go
  • internal/core/validator.go
  • internal/core/validator_test.go
  • internal/core/value_fields.go
  • internal/core/value_fields_test.go
  • internal/core/value_notices.go
  • internal/runtime/builtin/builtin.go
  • internal/runtime/builtin/dag/enqueue.go
  • internal/runtime/builtin/dag/enqueue_test.go
  • internal/runtime/builtin/foreach/foreach.go
  • internal/runtime/env.go
  • internal/runtime/eval.go
  • internal/runtime/foreach_spec018_test.go
  • internal/runtime/runner.go
  • specs/003-value-resolution.md
  • specs/005-value-resolution-params.md
  • specs/006-value-resolution-env.md
  • specs/007-value-resolution-steps.md
  • specs/009-step-reference.md
  • specs/018-parallel-and-foreach.md
  • specs/README.md

Comment thread internal/cmn/schema/dag.schema.json
Comment thread internal/core/validator.go
Comment thread internal/runtime/builtin/foreach/foreach.go
@yohamta0

Copy link
Copy Markdown
Member Author

@coderabbitai review

@coderabbitai

coderabbitai Bot commented Jun 22, 2026 •

Copy link
Copy Markdown
✅ Action performed

Review finished.

Note: CodeRabbit is an incremental review system and does not re-review already reviewed commits. This command is applicable only when automatic reviews are paused.

@yohamta0

Copy link
Copy Markdown
Member Author

@coderabbitai review

@coderabbitai

coderabbitai Bot commented Jun 22, 2026 •

Copy link
Copy Markdown
✅ Action performed

Review finished.

Note: CodeRabbit is an incremental review system and does not re-review already reviewed commits. This command is applicable only when automatic reviews are paused.

@yohamta0

Copy link
Copy Markdown
Member Author

@coderabbitai review

@coderabbitai

coderabbitai Bot commented Jun 22, 2026 •

Copy link
Copy Markdown
✅ Action performed

Review finished.

Note: CodeRabbit is an incremental review system and does not re-review already reviewed commits. This command is applicable only when automatic reviews are paused.

@yohamta0
yohamta0 merged commit f3af1e9 into main Jun 24, 2026
11 checks passed
@yohamta0
yohamta0 deleted the implement-spec-018-parallel-foreach branch June 24, 2026 00:21
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant