Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -512,6 +512,7 @@ Dagu includes built-in actions that run within the Dagu process (or worker). Loc
| [`template.render`](https://docs.dagu.sh/step-types/template) | Text generation with template rendering |
| [`router.route`](https://docs.dagu.sh/step-types/router) | Conditional step routing based on values and patterns |
| [`dag.run`](https://docs.dagu.sh/writing-workflows/control-flow) | Invoke another DAG as a sub-workflow with params and dependencies |
| [`dag.enqueue`](https://docs.dagu.sh/writing-workflows/control-flow) | Queue another DAG asynchronously and continue after enqueue |
| [`harness.run`](https://docs.dagu.sh/step-types/harness) | Run coding agent CLIs such as Claude Code, Codex, Copilot, OpenCode, and Pi |
| [`agent.run`](https://docs.dagu.sh/features/agent/step) | Built-in agent action with tool use |

Expand Down
22 changes: 21 additions & 1 deletion README_SCHEMA.md
Original file line number Diff line number Diff line change
Expand Up @@ -191,6 +191,7 @@ Current builtin actions:
| `redis.<operation>` | Redis operations | Redis config; operation comes from the action suffix |
| `jq.filter` | jq transforms | `filter`, plus `data` or `input` |
| `dag.run` | Child DAG execution | `dag`, optional `params` |
| `dag.enqueue` | Asynchronous child DAG enqueue | `dag`, optional `params`, optional `queue` |
| `router.route` | Conditional routing | `value`, `routes` |
| `chat.completion` | LLM chat completion | `prompt` or `messages`, model config |
| `agent.run` | Agent step execution | `task`, `prompt`, or `messages`, agent config |
Expand Down Expand Up @@ -229,6 +230,8 @@ for an in-memory DuckDB database.

### Child DAG

Use `dag.run` when the parent workflow must wait for the child DAG result:

```yaml
steps:
- id: process_account
Expand All @@ -240,7 +243,24 @@ steps:
REGION: us-east-1
```

`parallel:` currently requires `action: dag.run`:
Use `dag.enqueue` when the parent only needs to create a queued child DAG run
and continue:

```yaml
steps:
- id: queue_account_report
action: dag.enqueue
with:
dag: workflows/account-report
params:
ACCOUNT_ID: acct_123
queue: background
```

`dag.enqueue` accepts the same `with.dag` and `with.params` inputs as
`dag.run`, plus `with.queue` to override the queued child run's queue.

`parallel:` currently requires `action: dag.run` or `action: dag.enqueue`:

```yaml
steps:
Expand Down
3 changes: 3 additions & 0 deletions internal/cmd/dry.go
Original file line number Diff line number Diff line change
Expand Up @@ -84,6 +84,7 @@ func runDry(ctx *Context, args []string) error {
agent.Options{
Dry: true,
DAGRunStore: ctx.DAGRunStore,
QueueStore: ctx.QueueStore,
SecretStore: as.SecretStore,
ServiceRegistry: ctx.ServiceRegistry,
RootDAGRun: exec.NewDAGRunRef(dag.Name, dagRunID),
Expand All @@ -95,6 +96,8 @@ func runDry(ctx *Context, args []string) error {
AgentSoulStore: as.SoulStore,
AgentOAuthManager: as.OAuthManager,
AgentRemoteContextResolver: as.ContextResolver,
DAGRunLogDir: ctx.Config.Paths.LogDir,
DAGRunArtifactDir: ctx.Config.Paths.ArtifactDir,
},
)

Expand Down
3 changes: 3 additions & 0 deletions internal/cmd/restart.go
Original file line number Diff line number Diff line change
Expand Up @@ -175,6 +175,7 @@ func executeDAGWithRunID(ctx *Context, cli runtime.Manager, dag *core.DAG, dagRu
ExtraEnvs: extraEnvs,
PreparedAttempt: preparedAttempt,
DAGRunStore: ctx.DAGRunStore,
QueueStore: ctx.QueueStore,
SecretStore: as.SecretStore,
ServiceRegistry: ctx.ServiceRegistry,
RootDAGRun: exec.NewDAGRunRef(dag.Name, dagRunID),
Expand All @@ -188,6 +189,8 @@ func executeDAGWithRunID(ctx *Context, cli runtime.Manager, dag *core.DAG, dagRu
AgentRemoteContextResolver: as.ContextResolver,
ScheduleTime: scheduleTime,
ArtifactDir: artifactDir,
DAGRunLogDir: ctx.Config.Paths.LogDir,
DAGRunArtifactDir: ctx.Config.Paths.ArtifactDir,
})

listenSignals(ctx, agentInstance)
Expand Down
3 changes: 3 additions & 0 deletions internal/cmd/retry.go
Original file line number Diff line number Diff line change
Expand Up @@ -435,6 +435,7 @@ func executeRetry(ctx *Context, dag *core.DAG, status *exec.DAGRunStatus, rootRu
AttemptID: attemptID,
PreparedAttempt: preparedAttempt,
DAGRunStore: ctx.DAGRunStore,
QueueStore: ctx.QueueStore,
SecretStore: as.SecretStore,
ServiceRegistry: ctx.ServiceRegistry,
RootDAGRun: rootRun,
Expand All @@ -448,6 +449,8 @@ func executeRetry(ctx *Context, dag *core.DAG, status *exec.DAGRunStatus, rootRu
AgentOAuthManager: as.OAuthManager,
AgentRemoteContextResolver: as.ContextResolver,
ArtifactDir: artifactDir,
DAGRunLogDir: ctx.Config.Paths.LogDir,
DAGRunArtifactDir: ctx.Config.Paths.ArtifactDir,
},
)

Expand Down
3 changes: 3 additions & 0 deletions internal/cmd/start.go
Original file line number Diff line number Diff line change
Expand Up @@ -527,6 +527,7 @@ func executeDAGRun(ctx *Context, d *core.DAG, parent exec.DAGRunRef, dagRunID st
QueuedRun: queuedRun,
PreparedAttempt: preparedAttempt,
DAGRunStore: ctx.DAGRunStore,
QueueStore: ctx.QueueStore,
SecretStore: as.SecretStore,
ServiceRegistry: ctx.ServiceRegistry,
RootDAGRun: root,
Expand All @@ -541,6 +542,8 @@ func executeDAGRun(ctx *Context, d *core.DAG, parent exec.DAGRunRef, dagRunID st
AgentRemoteContextResolver: as.ContextResolver,
ScheduleTime: scheduleTime,
ArtifactDir: artifactDir,
DAGRunLogDir: ctx.Config.Paths.LogDir,
DAGRunArtifactDir: ctx.Config.Paths.ArtifactDir,
},
)

Expand Down
48 changes: 47 additions & 1 deletion internal/cmn/schema/dag.schema.json
Original file line number Diff line number Diff line change
Expand Up @@ -1941,7 +1941,7 @@
"description": "Parallel execution configuration with concurrency control"
}
],
"description": "Configuration for parallel execution of child DAGs. Currently only valid with action: dag.run. Allows processing multiple items concurrently using the same workflow definition."
"description": "Configuration for parallel execution of child DAGs. Currently only valid with action: dag.run or dag.enqueue. Allows processing multiple items concurrently using the same workflow definition."
},
"worker_selector": {
"oneOf": [
Expand Down Expand Up @@ -2122,6 +2122,18 @@
}
}
},
{
"if": {
"properties": { "action": { "const": "dag.enqueue" } },
"required": ["action"]
},
"then": {
"required": ["with"],
"properties": {
"with": { "$ref": "#/definitions/dagEnqueueActionConfig" }
}
}
},
{
"if": {
"properties": { "action": { "const": "http.request" } },
Expand Down Expand Up @@ -5928,6 +5940,7 @@
"shell",
"docker",
"container",
"dag_enqueue",
"file",
"git",
"kubernetes",
Expand Down Expand Up @@ -5982,6 +5995,7 @@
"artifact.write",
"chat.completion",
"container.run",
"dag.enqueue",
"dag.run",
"data.convert",
"data.pick",
Expand Down Expand Up @@ -6114,6 +6128,38 @@
},
"description": "Configuration for action: dag.run."
},
"dagEnqueueActionConfig": {
"type": "object",
"additionalProperties": false,
"required": ["dag"],
"properties": {
"dag": {
"type": "string",
"description": "Child DAG name or path."
},
"params": {
"oneOf": [
{ "type": "string" },
{
"type": "array",
"items": {
"oneOf": [
{ "type": "string" },
{ "type": "object", "additionalProperties": true }
]
}
},
{ "type": "object", "additionalProperties": true }
],
"description": "Parameters passed to the queued child DAG."
},
"queue": {
"type": "string",
"description": "Optional queue name override for the queued child DAG run."
}
},
"description": "Configuration for action: dag.enqueue."
},
"httpRequestActionConfig": {
"allOf": [
{ "$ref": "#/definitions/httpExecutorConfig" },
Expand Down
40 changes: 40 additions & 0 deletions internal/core/exec/context.go
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,10 @@ type Context struct {
BaseEnv *config.BaseEnv
EnvScope *eval.EnvScope // Unified environment scope - THE single source for all env vars
CoordinatorCli Dispatcher
DAGRunStore DAGRunStore
QueueStore QueueStore
DAGRunLogDir string
DAGRunArtifactDir string
Shell string // Default shell for this DAG (from DAG.Shell)
LogEncodingCharset string // Character encoding for log files (e.g., "utf-8", "shift_jis", "euc-jp")
LogWriterFactory LogWriterFactory // For remote log streaming (nil = use local files)
Expand Down Expand Up @@ -202,6 +206,10 @@ type contextOptions struct {
logEncodingCharset string
logWriterFactory LogWriterFactory
defaultExecMode config.ExecutionMode
dagRunStore DAGRunStore
queueStore QueueStore
dagRunLogDir string
dagRunArtifactDir string
workDir string
artifactDir string
}
Expand Down Expand Up @@ -273,6 +281,34 @@ func WithDefaultExecMode(mode config.ExecutionMode) ContextOption {
}
}

// WithDAGRunStore sets the dag-run store for executors that persist DAG runs.
func WithDAGRunStore(store DAGRunStore) ContextOption {
return func(o *contextOptions) {
o.dagRunStore = store
}
}

// WithQueueStore sets the queue store for executors that enqueue DAG runs.
func WithQueueStore(store QueueStore) ContextOption {
return func(o *contextOptions) {
o.queueStore = store
}
}

// WithDAGRunLogDir sets the base log directory for newly persisted DAG runs.
func WithDAGRunLogDir(dir string) ContextOption {
return func(o *contextOptions) {
o.dagRunLogDir = dir
}
}

// WithDAGRunArtifactDir sets the base artifact directory for newly persisted DAG runs.
func WithDAGRunArtifactDir(dir string) ContextOption {
return func(o *contextOptions) {
o.dagRunArtifactDir = dir
}
}

// WithWorkDir sets the per-DAG-run working directory path.
func WithWorkDir(dir string) ContextOption {
return func(o *contextOptions) {
Expand Down Expand Up @@ -340,6 +376,10 @@ func NewContext(
DAGRunID: dagRunID,
BaseEnv: config.GetBaseEnv(ctx),
CoordinatorCli: options.coordinator,
DAGRunStore: options.dagRunStore,
QueueStore: options.queueStore,
DAGRunLogDir: options.dagRunLogDir,
DAGRunArtifactDir: options.dagRunArtifactDir,
Shell: dag.Shell,
LogEncodingCharset: options.logEncodingCharset,
LogWriterFactory: options.logWriterFactory,
Expand Down
2 changes: 1 addition & 1 deletion internal/core/parallel_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -133,7 +133,7 @@ steps:
parallel: ${ITEMS}
`,
wantErr: true,
wantErrMsg: "parallel currently requires action: dag.run",
wantErrMsg: "parallel currently requires action: dag.run or dag.enqueue",
},
{
name: "ErrorParallelWithoutCommandOrRun",
Expand Down
4 changes: 2 additions & 2 deletions internal/core/spec/step_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -61,8 +61,8 @@ func TestMain(m *testing.M) {
core.RegisterExecutorCapabilities("wait", core.ExecutorCapabilities{Command: true})
// git: supports command only
core.RegisterExecutorCapabilities("git", core.ExecutorCapabilities{Command: true})
// dag/subworkflow/parallel: support SubDAG and WorkerSelector
for _, t := range []string{"dag", "subworkflow", "parallel"} {
// dag/subworkflow/parallel/dag_enqueue: support SubDAG and WorkerSelector
for _, t := range []string{"dag", "subworkflow", "parallel", core.ExecutorTypeDAGEnqueue} {
core.RegisterExecutorCapabilities(t, core.ExecutorCapabilities{
SubDAG: true, WorkerSelector: true,
})
Expand Down
1 change: 1 addition & 0 deletions internal/core/spec/step_types.go
Original file line number Diff line number Diff line change
Expand Up @@ -77,6 +77,7 @@ var builtinStepTypeNames = map[string]struct{}{
"container": {},
"dag": {},
"data": {},
"dag_enqueue": {},
"docker": {},
"duckdb": {},
"file": {},
Expand Down
29 changes: 29 additions & 0 deletions internal/core/spec/step_v2.go
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,7 @@ var builtinActionNormalizers = map[string]actionNormalizer{
"archive.list": operationAction("archive", "list"),
"chat.completion": normalizeChatAction,
"container.run": optionalCommandAction("container", "command"),
"dag.enqueue": normalizeDagEnqueueAction,
"dag.run": normalizeDagRunAction,
"data.convert": operationAction("data", "convert"),
"data.pick": operationAction("data", "pick"),
Expand Down Expand Up @@ -295,6 +296,34 @@ func normalizeDagRunAction(normalized map[string]any, with map[string]any) error
return nil
}

func normalizeDagEnqueueAction(normalized map[string]any, with map[string]any) error {
dagName, err := requireActionStringField(with, "dag")
if err != nil {
return err
}
for key := range with {
if key != "dag" && key != "params" && key != "queue" {
return core.NewValidationError("with", with, fmt.Errorf("dag.enqueue does not support with.%s", key))
}
}
normalized["type"] = core.ExecutorTypeDAGEnqueue
normalized["call"] = dagName
if params, ok := with["params"]; ok {
normalized["params"] = cloneAny(params)
}
config := make(map[string]any)
if queue, ok := with["queue"]; ok {
config["queue"] = cloneAny(queue)
}
if len(config) > 0 {
normalized["with"] = config
} else {
delete(normalized, "with")
}
delete(normalized, "action")
return nil
}

func normalizeExecAction(normalized map[string]any, with map[string]any) error {
command, err := requireActionStringField(with, "command")
if err != nil {
Expand Down
Loading
Loading