Repository navigation
feat: support variable substitution in worker_selector keys and values - #2417
Conversation
📝 WalkthroughWalkthroughWorker selectors now support variable substitution in keys and values during DAG build and sub-DAG execution. Parallel item context is preserved separately from explicit parameters and used to resolve ChangesWorker selector resolution
Estimated code review effort: 4 (Complex) | ~45 minutes Sequence Diagram(s)sequenceDiagram
participant DAGBuild
participant ParallelExecutor
participant resolveWorkerSelector
participant SubDAGExecutor
DAGBuild->>DAGBuild: Resolve root worker_selector
ParallelExecutor->>resolveWorkerSelector: Resolve selector with ITEM
resolveWorkerSelector-->>ParallelExecutor: Return resolved selector
ParallelExecutor->>SubDAGExecutor: Validate and configure child
SubDAGExecutor-->>ParallelExecutor: Execute sub-DAG
Possibly related PRs
🚥 Pre-merge checks | ✅ 4 | ❌ 1❌ Failed checks (1 warning)
✅ Passed checks (4 passed)
✨ Finishing Touches🧪 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: 1
🧹 Nitpick comments (1)
internal/runtime/builtin/dag/dag.go (1)
65-74: 📐 Maintainability & Code Quality | 🔵 Trivial | ⚡ Quick winExtract shared resolve+validate+cleanup helper for sub-DAG worker selector setup. Both call sites duplicate the identical resolve → cleanup-on-error → validate → cleanup-on-error →
SetWorkerSelectorsequence against*executor.SubDAGExecutor, risking drift between the two paths as this logic evolves.
internal/runtime/builtin/dag/dag.go#L65-L74: extract this sequence into a shared helper (e.g.,resolveAndConfigureWorkerSelector(ctx, child, target, step.WorkerSelector, nil)) reused by both this and the parallel path.internal/runtime/builtin/dag/parallel.go#L564-L574: call the same shared helper withworkerSelectorExtra(runParams)instead of re-implementing the resolve/validate/cleanup steps.♻️ Proposed shared helper
func resolveAndConfigureWorkerSelector( ctx context.Context, child *executor.SubDAGExecutor, target string, rawSelector any, extra map[string]string, ) error { workerSelector, err := resolveWorkerSelector(ctx, rawSelector, extra) if err != nil { _ = child.Cleanup(context.WithoutCancel(ctx)) return err } if err := validateSubDAG(child.DAG, target, workerSelector); err != nil { _ = child.Cleanup(context.WithoutCancel(ctx)) return err } child.SetWorkerSelector(workerSelector) return nil }🤖 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/runtime/builtin/dag/dag.go` around lines 65 - 74, Extract the duplicated worker-selector setup into a shared resolveAndConfigureWorkerSelector helper that resolves, cleans up on resolution or validation failure, validates the target, and calls SetWorkerSelector. In internal/runtime/builtin/dag/dag.go lines 65-74, replace the inline sequence with the helper using the sub-DAG name, step.WorkerSelector, and nil extra; in internal/runtime/builtin/dag/parallel.go lines 564-574, call the same helper with workerSelectorExtra(runParams).
🤖 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/runtime/transform/node.go`:
- Around line 28-40: Persist ParallelItem through every sub-run conversion and
recording path so restart/retry state retains ${ITEM} context: update
internal/runtime/transform/node.go lines 28-40 to copy the serialized value into
runtime.SubDAGRun, lines 89-105 to serialize runtime.SubDAGRun.ParallelItem,
internal/runtime/agent/agent.go lines 1141-1145 to include it when persisting
node sub-runs, and internal/runtime/step_executor.go lines 188-192 to retain it
when recording executor sub-runs.
---
Nitpick comments:
In `@internal/runtime/builtin/dag/dag.go`:
- Around line 65-74: Extract the duplicated worker-selector setup into a shared
resolveAndConfigureWorkerSelector helper that resolves, cleans up on resolution
or validation failure, validates the target, and calls SetWorkerSelector. In
internal/runtime/builtin/dag/dag.go lines 65-74, replace the inline sequence
with the helper using the sub-DAG name, step.WorkerSelector, and nil extra; in
internal/runtime/builtin/dag/parallel.go lines 564-574, call the same helper
with workerSelectorExtra(runParams).
🪄 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 Plus
Run ID: 7e162c3e-1695-404b-847b-eba425789b8c
📒 Files selected for processing (20)
README_SCHEMA.mdinternal/cmn/schema/dag.schema.jsoninternal/core/spec/dag.gointernal/core/spec/dag_test.gointernal/core/spec/loader_test.gointernal/intg/distr/parallel_test.gointernal/runtime/agent/agent.gointernal/runtime/builtin/dag/dag.gointernal/runtime/builtin/dag/enqueue.gointernal/runtime/builtin/dag/enqueue_test.gointernal/runtime/builtin/dag/parallel.gointernal/runtime/builtin/dag/worker_selector.gointernal/runtime/builtin/dag/worker_selector_test.gointernal/runtime/data.gointernal/runtime/executor/executor.gointernal/runtime/node.gointernal/runtime/node_test.gointernal/runtime/step_executor.gointernal/runtime/transform/node.gospecs/003-value-resolution.md
|
@coderabbitai review |
✅ Action performedReview finished.
|
|
@yohamta0 What do you think ? :) |
yohamta0
left a comment
There was a problem hiding this comment.
I found a few issues with when worker_selector is resolved. I left comments inline.
…substitution # Conflicts: # internal/runtime/builtin/dag/dag.go # specs/003-value-resolution.md
|
Resolved DAG-level |
|
Sorry for taking time. I'll look into this PR today. |
|
Found some issues, I'll push to this PR. |
Summary
${VAR}references inworker_selectorkeys and values so routing labels can be defined once (e.g. in workspace base-config env) instead of hardcoded in every DAGenv:(undefined variables stay literal, skipped underBuildFlagNoEval)${ITEM}Changes
resolveWorkerSelectorbuild step for the root selector ininternal/core/spec/dag.go, reusing theevaluatePairsresolver constructionresolveWorkerSelectorhelper ininternal/runtime/builtin/dagand wire it into thedag.run, parallel, and enqueue executors (parallel children resolve with the item context as${ITEM})worker_selector: localliteral; evaluate selector values with the strict workflow field policy, so no command substitution runs in routing metadataspecs/003-value-resolution.md, the JSON schema descriptions, andREADME_SCHEMA.mdImpact
Selectors flow through dispatch and coordinator matching already resolved; no changes to the dispatch decision, proto conversion, or matching logic. DAGs with literal selectors behave exactly as before, and undefined variables stay literal (the label simply never matches, the same failure mode as a typo today). No configuration or API changes.
Related Issues
Closes #2410
Testing
go test ./internal/core/spec/... -count=1(includes newTestWorkerSelectorEvaluation,TestLoad_WorkerSelectorFromBaseConfigEnv)go test ./internal/runtime/builtin/dag/... -count=1(includes newTestResolveWorkerSelector)make test(12k+ tests) andmake lintpass[intraday, batch]withworkload: ${ITEM}lands one child on each workerManual verification
make bin mkdir -p /tmp/wsdemo/{dags,data,logs} cat > /tmp/wsdemo/config.yaml <<'EOF' host: 127.0.0.1 port: 8080 auth: { mode: none } coordinator: { enabled: true, host: 127.0.0.1, advertise: 127.0.0.1, port: 50055, health_port: 0 } queues: { enabled: true } paths: data_dir: /tmp/wsdemo/data dags_dir: /tmp/wsdemo/dags log_dir: /tmp/wsdemo/logs base_config: /tmp/wsdemo/base.yaml EOF # the workspace-level routing label — the only place it is defined printf 'env:\n WORKLOAD: intraday\n' > /tmp/wsdemo/base.yaml cat > /tmp/wsdemo/dags/route-demo.yaml <<'EOF' worker_selector: workload: ${WORKLOAD} steps: - name: report run: echo "routed via workspace env" EOF # run each in its own terminal .local/bin/dagu coordinator -c /tmp/wsdemo/config.yaml .local/bin/dagu server -c /tmp/wsdemo/config.yaml .local/bin/dagu worker -c /tmp/wsdemo/config.yaml --worker.id=w-intraday \ --worker.labels workload=intraday --worker.coordinators=127.0.0.1:50055 --worker.health-port=0 # trigger via the API and check the assigned worker curl -s -X POST localhost:8080/api/v1/dags/route-demo/start -H 'Content-Type: application/json' -d '{}' sleep 3 curl -s 'localhost:8080/api/v1/dag-runs?name=route-demo' | grep -o '"workerId":"[^"]*"' # -> "workerId":"w-intraday"; the worker log shows: worker-selector=map[workload:intraday] # flip /tmp/wsdemo/base.yaml to a label with no worker and re-trigger: # dispatch fails with "no workers match the required selector" showing the RESOLVED valueI've not tested this in real distributed environments yet, will report back though.
Checklist
Summary by cubic
Adds variable substitution for
worker_selectorkeys and values. Resolution happens at DAG build (with base-config) and at sub-DAG runtime (env, params, and${ITEM}), and flows viaexecutor.RunParams.WorkerSelectoracrossdag.run, parallel, and enqueue.New Features
worker_selectorafter base-config composition; honorBuildFlagNoEval; addBuildFlagDeferWorkerSelector; keepworker_selector: localliteral.${ITEM}; if no step override, resolve the child DAG’s selector against runtime params before dispatch/retry and persist on enqueue;dag.run, parallel, and enqueue pass selectors viaexecutor.RunParams.WorkerSelector; docs and schema updated.Bug Fixes
Written for commit 09b28e8. Summary will update on new commits.
Summary by CodeRabbit
${VAR}and${ITEM}.