Skip to content

feat: add dag enqueue action - #2163

Merged
yohamta0 merged 5 commits into
mainfrom
feat/dag-enqueue-action
May 16, 2026
Merged

yohamta0 merged 5 commits into
mainfrom
feat/dag-enqueue-action

Conversation

@yohamta0

@yohamta0 yohamta0 commented May 16, 2026 •

Copy link
Copy Markdown
Member

Summary

Add asynchronous child DAG enqueue support with action: dag.enqueue.

Changes

  • Add dag.enqueue action normalization and dag_enqueue executor support.
  • Persist queued child DAG runs with sub-DAG trigger metadata, resolved params, logs, artifacts, and queue selection.
  • Wire queue and DAG-run stores into runtime execution context for executors that create queued runs.
  • Update sub-DAG API lookup paths so asynchronously queued child DAG runs remain visible from parent DAG run details.
  • Add parser, schema, runtime, queue-processing, and API coverage for enqueue behavior.

Related Issues

N/A

Checklist

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

Validation

  • go test -count=1 ./internal/core ./internal/core/spec ./internal/service/frontend/api/v1 ./internal/cmd ./internal/service/scheduler ./internal/runtime/agent ./internal/runtime/builtin/dag
  • make bin
  • git diff --check

Summary by CodeRabbit

  • New Features

    • Added dag.enqueue action to enqueue child DAG runs asynchronously with optional queue configuration and parameter overrides.
    • Enhanced support for parallel execution of enqueued child DAG runs.
  • Improvements

    • Improved sub-DAG run resolution and tracking in API endpoints for better visibility across parent and child DAG runs.

Review Change Stack

@coderabbitai

coderabbitai Bot commented May 16, 2026 •

Copy link
Copy Markdown

Important

Review skipped

Auto incremental reviews are disabled on this repository.

Please check the settings in the CodeRabbit UI or the .coderabbit.yaml file in this repository. To trigger a single review, invoke the @coderabbitai review command.

⚙️ Run configuration

Configuration used: Organization UI

Review profile: CHILL

Plan: Pro

Run ID: 9ce61d1b-282d-42a4-a4d6-93c22283e4d2

You can disable this status message by setting the reviews.review_status to false in the CodeRabbit configuration file.

Use the checkbox below for a quick retry:

  • 🔍 Trigger review
📝 Walkthrough

Walkthrough

This PR introduces asynchronous sub-DAG queueing via action: dag.enqueue. It extends the schema to validate the new action, implements the enqueue executor to persist queued runs into a queue store, wires store and directory configuration throughout the agent runtime, updates the API to resolve queued children, and adds comprehensive tests.

Changes

DAG Enqueue Feature

Layer / File(s) Summary
Schema and executor type definitions
internal/cmn/schema/dag.schema.json, internal/core/step.go, internal/core/spec/step_types.go
JSON schema now validates action: dag.enqueue with required with.dag, optional params and queue; executor type constant ExecutorTypeDAGEnqueue and builtin registry updated.
Action normalization and validation
internal/core/spec/step_v2.go, internal/core/spec/step_v2_test.go, internal/core/spec/step_test.go, internal/core/validator.go, internal/core/validator_test.go, internal/core/parallel_test.go
normalizeDagEnqueueAction converts YAML into step representation; validation error messages updated to include dag.enqueue alongside dag.run; tests added for enqueue action parsing and parallel fanout.
Execution context and agent configuration
internal/core/exec/context.go, internal/runtime/agent/agent.go, internal/runtime/context.go
Context extended with DAG run store, queue store, log/artifact directories; ContextOption helpers provided; Agent.Run and Agent.dryRun conditionally wire these into runtime context.
Command layer integration
internal/cmd/dry.go, internal/cmd/start.go, internal/cmd/restart.go, internal/cmd/retry.go, internal/test/helper.go
All command executors and test helper now pass queue store and log/artifact directories to agent options during initialization.
DAG enqueue executor implementation
internal/runtime/builtin/dag/enqueue.go
New executor validates prerequisites, resolves sub-DAG targets, creates queued attempts with persisted status and log/artifact paths, enqueues items, and outputs summary JSON; supports parameter resolution, worker selector filtering, and queue override.
API sub-DAG run resolution
internal/service/frontend/api/v1/dagruns.go, internal/service/frontend/api/v1/dagruns_test.go
New helpers getReferencedDAGRunStatus and getReferencedDAGRunAttempt scan parent status nodes to locate queued children when direct parent-referenced lookups fail; API endpoints updated to use these resolvers; test validates /sub-dag-runs endpoint includes enqueued children.
Agent queue operation tests
internal/runtime/agent/agent_test.go
Tests verify dag.enqueue queues sub-DAG with correct params/status, and that queue processing executes queued children to completion.

Estimated code review effort

🎯 3 (Moderate) | ⏱️ ~25 minutes

🚥 Pre-merge checks | ✅ 4 | ❌ 1

❌ Failed checks (1 warning)

Check name Status Explanation Resolution
Docstring Coverage ⚠️ Warning Docstring coverage is 34.29% 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: add dag enqueue action' clearly and concisely summarizes the main change: introducing a new dag.enqueue action feature.
Description check ✅ Passed The PR description follows the template structure with all required sections: Summary, Changes (bulleted list), Related Issues, and Checklist. All critical information is present and complete.
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 unit tests (beta)
  • Create PR with unit tests
  • Commit unit tests in branch feat/dag-enqueue-action

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

🧹 Nitpick comments (1)
internal/core/spec/step_test.go (1)

60-62: ⚡ Quick win

Use the executor type constant instead of the raw "dag_enqueue" literal.

Using core.ExecutorTypeDAGEnqueue here avoids drift if the canonical type name changes and keeps tests aligned with production constants.

Suggested patch
-	// dag/subworkflow/parallel/dag_enqueue: support SubDAG and WorkerSelector
-	for _, t := range []string{"dag", "subworkflow", "parallel", "dag_enqueue"} {
+	// dag/subworkflow/parallel/dag_enqueue: support SubDAG and WorkerSelector
+	for _, t := range []string{"dag", "subworkflow", "parallel", core.ExecutorTypeDAGEnqueue} {
🤖 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/core/spec/step_test.go` around lines 60 - 62, Replace the hard-coded
executor type string "dag_enqueue" with the canonical constant
core.ExecutorTypeDAGEnqueue in the RegisterExecutorCapabilities loop; update the
slice of executor types (used where RegisterExecutorCapabilities is called) to
use core.ExecutorTypeDAGEnqueue so tests reference the production constant and
avoid drift with core.RegisterExecutorCapabilities and
core.ExecutorCapabilities.
🤖 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/service/frontend/api/v1/dagruns.go`:
- Around line 3292-3336: Mutation handlers that currently call
dagRunMgr.FindSubDAGRunStatus or dagRunStore.FindSubAttempt/FindAttempt directly
should use the same fallback resolution as reads: replace those direct calls
with the helper methods getReferencedDAGRunStatus(...) and
getReferencedDAGRunAttempt(...) (or replicate their logic) so queued child runs
without parent linkage resolve via findReferencedDAGName before failing; update
all mutation flows that perform approve/reject/push-back/update/resume to call
these helpers (or perform the same trim+findReferencedDAGName then fallback
lookup) and return the original subErr when neither lookup succeeds.

---

Nitpick comments:
In `@internal/core/spec/step_test.go`:
- Around line 60-62: Replace the hard-coded executor type string "dag_enqueue"
with the canonical constant core.ExecutorTypeDAGEnqueue in the
RegisterExecutorCapabilities loop; update the slice of executor types (used
where RegisterExecutorCapabilities is called) to use core.ExecutorTypeDAGEnqueue
so tests reference the production constant and avoid drift with
core.RegisterExecutorCapabilities and core.ExecutorCapabilities.
🪄 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: 3daad6e4-4d86-49e4-a466-82db6810e12c

📥 Commits

Reviewing files that changed from the base of the PR and between 1592e1b and 22c47d2.

📒 Files selected for processing (21)
  • internal/cmd/dry.go
  • internal/cmd/restart.go
  • internal/cmd/retry.go
  • internal/cmd/start.go
  • internal/cmn/schema/dag.schema.json
  • internal/core/exec/context.go
  • internal/core/parallel_test.go
  • internal/core/spec/step_test.go
  • internal/core/spec/step_types.go
  • internal/core/spec/step_v2.go
  • internal/core/spec/step_v2_test.go
  • internal/core/step.go
  • internal/core/validator.go
  • internal/core/validator_test.go
  • internal/runtime/agent/agent.go
  • internal/runtime/agent/agent_test.go
  • internal/runtime/builtin/dag/enqueue.go
  • internal/runtime/context.go
  • internal/service/frontend/api/v1/dagruns.go
  • internal/service/frontend/api/v1/dagruns_test.go
  • internal/test/helper.go

Comment thread internal/service/frontend/api/v1/dagruns.go
@yohamta0
yohamta0 force-pushed the feat/dag-enqueue-action branch from 41fb6b1 to 6a67638 Compare May 16, 2026 15:19
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