Repository navigation
feat: add human task workflow steps - #2396
Conversation
Cover local and distributed waiting, completion, persistence, validation, and retry behavior. Harden atomic resume claims and release distributed admission at waiting checkpoints.
📝 WalkthroughWalkthroughAdds the ChangesHuman Task Feature
Estimated code review effort: 5 (Critical) | ~120 minutes Possibly related PRs
Sequence Diagram(s)sequenceDiagram
participant Operator
participant Runner
participant DAGRunStore
participant Worker
Operator->>Runner: Start root DAG
Runner->>DAGRunStore: Persist human task as Waiting
Operator->>DAGRunStore: Complete task with input
DAGRunStore->>Worker: Queue distributed retry when applicable
Worker->>DAGRunStore: Resume DAG run
DAGRunStore-->>Operator: Report Succeeded status
🚥 Pre-merge checks | ✅ 3 | ❌ 2❌ Failed checks (2 warnings)
✅ Passed checks (3 passed)
✨ Finishing Touches📝 Generate docstrings
🧪 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: 4
🧹 Nitpick comments (2)
internal/runtime/builtin/dag/parallel.go (1)
564-566: 🩺 Stability & Availability | 🔵 Trivial | 💤 Low valueUse
context.WithoutCancelfor cleanup.If
validateSubDAGfails due to a context cancellation that occurred elsewhere,child.Cleanup(ctx)will immediately fail to execute its internal file-removal or logging operations because its context is already cancelled. Wrapping the context ensures the temporary file is still removed, matching the safety pattern used inenqueue.go.♻️ Proposed fix
- if err := validateSubDAG(child.DAG, target, e.step.WorkerSelector); err != nil { - _ = child.Cleanup(ctx) - return nil, err - } + if err := validateSubDAG(child.DAG, target, e.step.WorkerSelector); err != nil { + _ = child.Cleanup(context.WithoutCancel(ctx)) + return nil, err + }🤖 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/parallel.go` around lines 564 - 566, Update the cleanup call in the validateSubDAG failure path to pass a context derived with context.WithoutCancel(ctx) to child.Cleanup, ensuring cleanup operations run even when the original context is cancelled. Preserve the existing error return and validation flow.conformance/spec031_human_task/validation_test.go (1)
60-74: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low valueConsider enabling parallel execution for this top-level test.
Adding
t.Parallel()at the top-level of this test allows it to execute concurrently alongside other tests likeTestHumanTaskShapeValidation. Note: It's best omitted from the innert.Runloop here because the subtests currently share the same--run-id.♻️ Proposed refactor
func TestHumanTaskChildDAGBoundary(t *testing.T) { + t.Parallel() root := harness.NewRunner(t) rootResult := root.Run("start", "--run-id=spec031-child-root", "child_human.yaml")🤖 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 `@conformance/spec031_human_task/validation_test.go` around lines 60 - 74, Enable parallel execution in TestHumanTaskChildDAG by adding t.Parallel() at the start of the top-level test. Do not add parallel execution to the inner t.Run subtests because they share the same --run-id.
🤖 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/cmd/human_task.go`:
- Around line 641-680: Update humanTaskAttemptIsFinalizing and its caller
waitForHumanTaskAttemptFinalization to propagate findHumanTaskNodeByID errors
immediately instead of treating them as a finalizing state. Preserve polling
only when the human-task node is found and remains incomplete, while returning
the lookup error for missing or ambiguous steps.
- Around line 363-381: Update rollbackHumanTaskResumeClaim to capture the error
returned by CompareAndSwapLatestAttemptStatus and log a clear rollback-failure
message when the CAS fails, including the relevant human-task or run context.
Preserve the existing status restoration behavior and do not silently discard
the CAS error.
In `@internal/runtime/builtin/dag/dag.go`:
- Around line 65-66: Update the child-executor failure paths around
validateSubDAG and the subsequent working-directory existence check to call
child.Cleanup() before returning an error. Ensure cleanup occurs for both
validation failures and missing-directory failures while preserving the existing
error returns.
In `@internal/runtime/runner.go`:
- Around line 653-658: In the runHumanTask precondition error branch around
meetsPreconditions, remove the r.Cancel(plan) call while retaining
r.setLastError(err) and the early return. This should mark the human-task node
failure without immediately cancelling the entire DAG, matching the prompt
evaluation and standard prepareNode behavior.
---
Nitpick comments:
In `@conformance/spec031_human_task/validation_test.go`:
- Around line 60-74: Enable parallel execution in TestHumanTaskChildDAG by
adding t.Parallel() at the start of the top-level test. Do not add parallel
execution to the inner t.Run subtests because they share the same --run-id.
In `@internal/runtime/builtin/dag/parallel.go`:
- Around line 564-566: Update the cleanup call in the validateSubDAG failure
path to pass a context derived with context.WithoutCancel(ctx) to child.Cleanup,
ensuring cleanup operations run even when the original context is cancelled.
Preserve the existing error return and validation flow.
🪄 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: 78759d6d-3c07-4bf7-924c-6ae43a809a05
📒 Files selected for processing (98)
README.mdREADME_SCHEMA.mdcmd/main.gocmd/main_test.goconformance/harness/runner.goconformance/spec031_human_task/completion_test.goconformance/spec031_human_task/distributed_test.goconformance/spec031_human_task/helpers_test.goconformance/spec031_human_task/lifecycle_test.goconformance/spec031_human_task/testdata/acknowledgement.yamlconformance/spec031_human_task/testdata/additional_properties.yamlconformance/spec031_human_task/testdata/child_human.yamlconformance/spec031_human_task/testdata/concurrent.yamlconformance/spec031_human_task/testdata/distributed_probe.yamlconformance/spec031_human_task/testdata/distributed_root.yamlconformance/spec031_human_task/testdata/distributed_worker_one_probe.yamlconformance/spec031_human_task/testdata/distributed_worker_two_probe.yamlconformance/spec031_human_task/testdata/independent_branch.yamlconformance/spec031_human_task/testdata/invalid_execution_field.yamlconformance/spec031_human_task/testdata/invalid_foreach_human_task.yamlconformance/spec031_human_task/testdata/invalid_form_additional_properties.yamlconformance/spec031_human_task/testdata/invalid_form_null.yamlconformance/spec031_human_task/testdata/invalid_form_property.yamlconformance/spec031_human_task/testdata/invalid_form_required.yamlconformance/spec031_human_task/testdata/invalid_form_root.yamlconformance/spec031_human_task/testdata/invalid_handler_human_task.yamlconformance/spec031_human_task/testdata/invalid_lifecycle_field.yamlconformance/spec031_human_task/testdata/invalid_missing_id.yamlconformance/spec031_human_task/testdata/invalid_output_field.yamlconformance/spec031_human_task/testdata/invalid_prompt.yamlconformance/spec031_human_task/testdata/invalid_with_field.yamlconformance/spec031_human_task/testdata/max_output_size.yamlconformance/spec031_human_task/testdata/multi_document_human.yamlconformance/spec031_human_task/testdata/multiple_waiting.yamlconformance/spec031_human_task/testdata/parent_dag_enqueue.yamlconformance/spec031_human_task/testdata/parent_dag_run.yamlconformance/spec031_human_task/testdata/parent_parallel.yamlconformance/spec031_human_task/testdata/precondition_skipped.yamlconformance/spec031_human_task/testdata/prompt_snapshot.yamlconformance/spec031_human_task/testdata/prompt_snapshot_changed.yamlconformance/spec031_human_task/testdata/sequential_waiting.yamlconformance/spec031_human_task/testdata/typed_form.yamlconformance/spec031_human_task/testdata/valid_acknowledgement_shape.yamlconformance/spec031_human_task/testdata/valid_form_shape.yamlconformance/spec031_human_task/validation_test.gointernal/cmd/context.gointernal/cmd/human_task.gointernal/cmd/human_task_test.gointernal/cmn/schema/dag.schema.jsoninternal/cmn/schema/dag_schema_test.gointernal/core/dag.gointernal/core/dag_test.gointernal/core/exec/env.gointernal/core/exec/node.gointernal/core/foreach_validation_spec018_test.gointernal/core/spec/dag.gointernal/core/spec/defaults.gointernal/core/spec/human_task.gointernal/core/spec/human_task_form.gointernal/core/spec/human_task_test.gointernal/core/spec/loader.gointernal/core/spec/step.gointernal/core/spec/step_v2.gointernal/core/status.gointernal/core/step.gointernal/core/validator.gointernal/core/validator_test.gointernal/core/value_fields.gointernal/intg/distr/human_task_test.gointernal/output/tree.gointernal/output/tree_test.gointernal/runtime/agent/agent.gointernal/runtime/agent/agent_test.gointernal/runtime/agent/status_masking.gointernal/runtime/agent/status_masking_test.gointernal/runtime/builtin/dag/dag.gointernal/runtime/builtin/dag/enqueue.gointernal/runtime/builtin/dag/enqueue_test.gointernal/runtime/builtin/dag/parallel.gointernal/runtime/data.gointernal/runtime/plan.gointernal/runtime/runner.gointernal/runtime/runner_helper_test.gointernal/runtime/runner_test.gointernal/runtime/transform/node.gointernal/runtime/transform/node_test.gointernal/service/chatbridge/notifications.gointernal/service/coordinator/handler.gointernal/service/frontend/api/v1/dags.gointernal/service/mcp/reference.gointernal/service/mcp/server_test.gollms.txtskills/dagu/SKILL.mdskills/dagu/references/cli.mdskills/dagu/references/context.mdskills/dagu/references/steptypes.mdspecs/031-human-task.mdspecs/README.md
💤 Files with no reviewable changes (1)
- internal/core/spec/loader.go
|
@coderabbitai review |
✅ Action performedReview finished.
|
|
@coderabbitai review |
✅ Action performedReview finished.
|
There was a problem hiding this comment.
Actionable comments posted: 2
🤖 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/foreach_spec018_test.go`:
- Line 135: Update the foreachProbeState initialization in the relevant test to
create activeChanged as an unbuffered channel, matching the broadcast
synchronization pattern used by start() and waitForActive().
- Around line 164-196: Update foreachProbeState’s active-change notification so
each state change broadcasts to all waiters: under s.mu, capture the current
s.activeChanged channel, close it, and replace it with a new channel before
releasing the lock. In waitForActive, read the current notification channel
while holding s.mu and select on that captured channel alongside ctx.Done() and
the timer, preserving the maxActive check and existing timeout behavior.
🪄 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: dbcaec9c-356b-4ee0-8eb0-758db33503ac
📒 Files selected for processing (1)
internal/runtime/foreach_spec018_test.go
|
@coderabbitai review |
✅ Action performedReview finished.
|
Summary
human.taskworkflow stepsWhy
Workflows sometimes need durable human input between automated steps. A human task should end the active attempt, keep the run in
Waiting, and later continue from persisted run data without holding a process or distributed worker slot.Behavior
Complete a waiting task locally:
Human tasks are restricted to root DAG runs. REST API, control-plane, and Web UI changes are intentionally outside this PR.
Validation
DAGU_BIN=.local/bin/dagu go test -race ./conformance/spec031_human_task -count=1go test ./internal/cmd ./internal/core/spec ./internal/service/coordinator -count=1make lintgit diff --checkSummary by cubic
Adds processless
human.tasksteps that pause root DAG runs for operator input, persist a waiting checkpoint, and resume locally or on another distributed worker. Also hardens waiting admission handling (including closing cached attempts) and status compaction, and improves output rendering with secret‑masked prompts.New Features
action: human.taskwith explicitidandwith.prompt; optional flatwith.form(JSON Schema). Required/defaulted properties become${steps.<id>.outputs.*}; no executor/commands/outputs on the step.dagu human-task complete --run-id <id> --step <id> [--input k=v ... | --inputs-json '{...}'] <dag>; inputs are validated/coerced by the form schema; clearer local-context errors. Works for local and distributed root runs.foreach.steps. Defaultretry,timeout, andmail_on_errorare not applied. Status-file compaction tolerates active readers.Migration
id,action: human.task, andwith.prompt; add a flatwith.formwhen you need typed inputs.outputs:on ahuman.task; read values as${steps.<id>.outputs.<name>}downstream.dagu human-task completeto finish waiting tasks; pass inputs via--inputor--inputs-json.foreach.steps; remove executors/commands.Written for commit ffed36f. Summary will update on new commits.
Summary by CodeRabbit
New Features
human.taskoperator-input steps that pause root DAG runs (acknowledgement-only or with a typed/validated form) and resume downstream with submitted/defaulted values.dagu human-task completeto finish waiting tasks using--inputor--inputs-json, including distributed-run resumption.Documentation
human.taskauthoring rules, form/output semantics, and the completion CLI; marked Human Tasks spec implemented.Validation
Tests