Skip to content

feat: support variable substitution in worker_selector keys and values - #2417

Merged
yohamta0 merged 18 commits into
dagucloud:mainfrom
Giebisch:feat/worker-selector-substitution
Jul 30, 2026
Merged

yohamta0 merged 18 commits into
dagucloud:mainfrom
Giebisch:feat/worker-selector-substitution

Conversation

@Giebisch

@Giebisch Giebisch commented Jul 24, 2026 •

Copy link
Copy Markdown
Contributor

Summary

  • resolve ${VAR} references in worker_selector keys and values so routing labels can be defined once (e.g. in workspace base-config env) instead of hardcoded in every DAG
  • DAG-level selectors resolve at build time after base-config composition, with the same semantics as env: (undefined variables stay literal, skipped under BuildFlagNoEval)
  • step-level selectors resolve at runtime before the sub-DAG run/enqueue request is created; parallel children can route per item via ${ITEM}
  • add spec, loader (base-config inheritance), and runtime regression coverage

Changes

  • add a post-compose resolveWorkerSelector build step for the root selector in internal/core/spec/dag.go, reusing the evaluatePairs resolver construction
  • add a runtime resolveWorkerSelector helper in internal/runtime/builtin/dag and wire it into the dag.run, parallel, and enqueue executors (parallel children resolve with the item context as ${ITEM})
  • keep the string form worker_selector: local literal; evaluate selector values with the strict workflow field policy, so no command substitution runs in routing metadata
  • update the value-resolution matrix in specs/003-value-resolution.md, the JSON schema descriptions, and README_SCHEMA.md

Impact

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 new TestWorkerSelectorEvaluation, TestLoad_WorkerSelectorFromBaseConfigEnv)
  • go test ./internal/runtime/builtin/dag/... -count=1 (includes new TestResolveWorkerSelector)
  • full make test (12k+ tests) and make lint pass
  • verified live on a local distributed stack (coordinator + two labeled workers): DAG-level selector follows base-config env, step-level selector dispatches the child correctly, and a parallel step over [intraday, batch] with workload: ${ITEM} lands one child on each worker

Manual 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 value

I've not tested this in real distributed environments yet, will report back though.

Checklist

  • Code follows the project style guidelines
  • Self-review of the code has been performed
  • Tests have been added or updated as needed
  • Documentation has been updated as needed
  • Changes have been tested locally

Summary by cubic

Adds variable substitution for worker_selector keys and values. Resolution happens at DAG build (with base-config) and at sub-DAG runtime (env, params, and ${ITEM}), and flows via executor.RunParams.WorkerSelector across dag.run, parallel, and enqueue.

  • New Features

    • Resolve root worker_selector after base-config composition; honor BuildFlagNoEval; add BuildFlagDeferWorkerSelector; keep worker_selector: local literal.
    • Resolve step-level selectors during sub-DAG expansion with ${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 via executor.RunParams.WorkerSelector; docs and schema updated.
  • Bug Fixes

    • Validate empty or duplicate resolved keys; error when deduplicated parallel items produce the same sub-run ID with different selectors.
    • Fail fast if a child selector requires missing params; surface the error on the parent before dispatch.
    • Keep the parallel item as runtime-only routing context and do not persist it.

Written for commit 09b28e8. Summary will update on new commits.

Review in cubic

Summary by CodeRabbit

  • New Features
    • Added variable substitution support for worker selector keys and values.
    • Worker selectors can now reference environment variables, parameters, and parallel items using expressions such as ${VAR} and ${ITEM}.
    • Parallel sub-DAG runs preserve their original item context while honoring explicitly provided parameters.
  • Bug Fixes
    • Added validation for unresolved empty or duplicate selector keys.
  • Documentation
    • Updated schema, guide, and value-resolution documentation to describe worker selector substitution behavior.

@coderabbitai

coderabbitai Bot commented Jul 24, 2026 •

Copy link
Copy Markdown

Review Change Stack

📝 Walkthrough

Walkthrough

Worker 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 ${ITEM} for enqueue and parallel child execution.

Changes

Worker selector resolution

Layer / File(s) Summary
Build-time selector resolution
README_SCHEMA.md, internal/cmn/schema/dag.schema.json, internal/core/spec/dag.go, internal/core/spec/*_test.go, specs/003-value-resolution.md
DAG-level and step-level worker selector keys and values are evaluated from environment, parameters, and constants, with validation for empty or duplicate resolved keys and support for WithoutEval().
Parallel item propagation
internal/runtime/data.go, internal/runtime/executor/executor.go, internal/core/exec/node.go, internal/runtime/node.go, internal/runtime/builtin/action/action.go, internal/runtime/transform/node_test.go, internal/runtime/*_test.go
Parallel sub-run items are carried through SubDAGRun, RunParams, executor setup, action records, and persistence independently from explicit sub-DAG parameters.
Runtime selector resolution
internal/runtime/builtin/dag/worker_selector.go, internal/runtime/builtin/dag/{dag,enqueue,parallel}.go, internal/runtime/builtin/dag/*_test.go, internal/intg/distr/parallel_test.go
Runtime selectors are resolved with environment and ${ITEM} context before validation and child executor configuration, while enqueue records the resolved selector and parallel item.

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
Loading

Possibly related PRs

🚥 Pre-merge checks | ✅ 4 | ❌ 1

❌ Failed checks (1 warning)

Check name Status Explanation Resolution
Docstring Coverage ⚠️ Warning Docstring coverage is 56.52% 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
Linked Issues check ✅ Passed The PR implements the requested worker_selector variable/parameter substitution and supports workspace and parallel-item routing as required by #2410.
Out of Scope Changes check ✅ Passed The additional schema, docs, tests, and ParallelItem plumbing all support the worker_selector substitution work and appear in scope.
Title check ✅ Passed The title is concise and accurately states the main change: variable substitution in worker_selector keys and values.
Description check ✅ Passed All required template sections are present and filled with concrete details, including summary, changes, related issue, and checklist.
✨ Finishing Touches
🧪 Generate unit tests (beta)
  • Create PR with unit tests

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.

@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: 1

🧹 Nitpick comments (1)
internal/runtime/builtin/dag/dag.go (1)

65-74: 📐 Maintainability & Code Quality | 🔵 Trivial | ⚡ Quick win

Extract 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 → SetWorkerSelector sequence 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 with workerSelectorExtra(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

📥 Commits

Reviewing files that changed from the base of the PR and between 94847b1 and 7f54801.

📒 Files selected for processing (20)
  • README_SCHEMA.md
  • internal/cmn/schema/dag.schema.json
  • internal/core/spec/dag.go
  • internal/core/spec/dag_test.go
  • internal/core/spec/loader_test.go
  • internal/intg/distr/parallel_test.go
  • internal/runtime/agent/agent.go
  • internal/runtime/builtin/dag/dag.go
  • internal/runtime/builtin/dag/enqueue.go
  • internal/runtime/builtin/dag/enqueue_test.go
  • internal/runtime/builtin/dag/parallel.go
  • internal/runtime/builtin/dag/worker_selector.go
  • internal/runtime/builtin/dag/worker_selector_test.go
  • internal/runtime/data.go
  • internal/runtime/executor/executor.go
  • internal/runtime/node.go
  • internal/runtime/node_test.go
  • internal/runtime/step_executor.go
  • internal/runtime/transform/node.go
  • specs/003-value-resolution.md

Comment thread internal/runtime/transform/node.go Outdated
@Giebisch

Copy link
Copy Markdown
Contributor Author

@coderabbitai review

@coderabbitai

coderabbitai Bot commented Jul 24, 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.

@Giebisch

Copy link
Copy Markdown
Contributor Author

@yohamta0 What do you think ? :)

@yohamta0 yohamta0 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

I found a few issues with when worker_selector is resolved. I left comments inline.

Comment thread internal/core/exec/node.go Outdated
Comment thread internal/runtime/builtin/dag/dag.go Outdated
Comment thread internal/core/spec/dag.go Outdated
Comment thread internal/core/spec/dag.go
@Giebisch

Copy link
Copy Markdown
Contributor Author

Resolved DAG-level worker_selector against the composed base config (so base-config env/params and ${params.*} work) and made the empty step-selector fall back to the child's selector re-resolved with with.params overrides. Also made the parallel ${ITEM} runtime-only instead of persisting it in the sub-DAG run status, since it's re-derived on every run.

@Giebisch
Giebisch requested a review from yohamta0 July 27, 2026 14:35
@yohamta0

Copy link
Copy Markdown
Member

Sorry for taking time. I'll look into this PR today.

@yohamta0

Copy link
Copy Markdown
Member

Found some issues, I'll push to this PR.

@yohamta0 yohamta0 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

LGTM, thank you very much!

@yohamta0
yohamta0 merged commit cba1426 into dagucloud:main Jul 30, 2026
12 of 13 checks passed
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.

feat: Support variable / parameter substitution in worker_selector values

2 participants