Skip to content

feat: add human task workflow steps - #2396

Merged
yohamta0 merged 16 commits into
mainfrom
feat/human-task-backend
Jul 21, 2026
Merged

yohamta0 merged 16 commits into
mainfrom
feat/human-task-backend

Conversation

@yohamta0

@yohamta0 yohamta0 commented Jul 20, 2026 •

Copy link
Copy Markdown
Member

Summary

  • add backend support for processless human.task workflow steps
  • persist waiting checkpoints and resume the same root DAG run through the local CLI
  • support typed flat forms whose declared values become step outputs
  • allow distributed root runs to release worker capacity while waiting and resume on another worker
  • document the feature across the schema, CLI reference, skill, MCP reference, and user documentation
  • add local and distributed Spec 031 conformance coverage

Why

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

steps:
  - id: release_review
    action: human.task
    with:
      prompt: Choose the release target
      form:
        type: object
        properties:
          environment:
            type: string
            enum: [staging, production]
        required: [environment]

  - id: deploy
    depends: release_review
    run: deploy "${steps.release_review.outputs.environment}"

Complete a waiting task locally:

dagu human-task complete \
  --run-id <run-id> \
  --step release_review \
  --input environment=production \
  <dag-name>

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=1
  • go test ./internal/cmd ./internal/core/spec ./internal/service/coordinator -count=1
  • make lint
  • git diff --check

Summary by cubic

Adds processless human.task steps 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

    • Step: action: human.task with explicit id and with.prompt; optional flat with.form (JSON Schema). Required/defaulted properties become ${steps.<id>.outputs.*}; no executor/commands/outputs on the step.
    • Runtime: snapshots the prompt, opens Waiting without a process, and makes completion atomic/idempotent. Distributed runs release the admission token and close the cached attempt so another worker can resume. Tree output shows prompt and form; prompts are secret‑masked.
    • CLI: 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.
    • Validation: human tasks are root‑only; forbidden in sub‑DAGs, lifecycle handlers, and foreach.steps. Default retry, timeout, and mail_on_error are not applied. Status-file compaction tolerates active readers.
  • Migration

    • Author steps with id, action: human.task, and with.prompt; add a flat with.form when you need typed inputs.
    • Do not declare outputs: on a human.task; read values as ${steps.<id>.outputs.<name>} downstream.
    • Use dagu human-task complete to finish waiting tasks; pass inputs via --input or --inputs-json.
    • Move any human tasks out of sub‑DAGs, lifecycle handlers, and foreach.steps; remove executors/commands.

Written for commit ffed36f. Summary will update on new commits.

Review in cubic

Summary by CodeRabbit

  • New Features

    • Added human.task operator-input steps that pause root DAG runs (acknowledgement-only or with a typed/validated form) and resume downstream with submitted/defaulted values.
    • Added dagu human-task complete to finish waiting tasks using --input or --inputs-json, including distributed-run resumption.
    • Improved DAG status/output tree rendering for waiting human tasks.
  • Documentation

    • Documented human.task authoring rules, form/output semantics, and the completion CLI; marked Human Tasks spec implemented.
  • Validation

    • Enforced schema and lifecycle restrictions (no sub-DAG/handlers/foreach placement; constrained step fields; rejects invalid form shapes).
  • Tests

    • Added conformance, unit, integration, and CLI coverage for completion behavior, retries, concurrency, and distributed scenarios.

@coderabbitai

coderabbitai Bot commented Jul 20, 2026 •

Copy link
Copy Markdown

Review Change Stack

📝 Walkthrough

Walkthrough

Adds the human.task action for root DAGs, supporting acknowledgement or validated form input, waiting-state persistence, CLI completion, local/distributed resumption, output propagation, rendering, validation, tests, and documentation.

Changes

Human Task Feature

Layer / File(s) Summary
Human task contracts and validation
internal/core/*, internal/cmn/schema/*
Adds human-task types, form validation, derived outputs, action normalization, prohibited-field checks, and root-DAG, handler, and foreach boundaries.
Runtime lifecycle and persistence
internal/runtime/*, internal/service/coordinator/*
Opens tasks in Waiting, persists input, preserves prompts and state, releases distributed attempts, and rejects sub-DAG execution.
Completion CLI
internal/cmd/human_task.go, cmd/*
Registers dagu human-task complete, parses inputs, applies atomic completion, handles conflicts, and resumes locally or through the worker queue.
Rendering and verification
internal/output/*, conformance/spec031_human_task/*, internal/*_test.go
Displays waiting task details and adds unit, integration, lifecycle, distributed, concurrency, retry, and boundary coverage.
Documentation and specification
specs/*, README*, llms.txt, skills/dagu/*, internal/service/mcp/*
Documents authoring, forms, outputs, restrictions, completion commands, and the Spec031 contract.

Estimated code review effort: 5 (Critical) | ~120 minutes

Possibly related PRs

  • dagucloud/dagu#2098 — Both changes touch step validation and validation error behavior.
  • dagucloud/dagu#2122 — The schema and action-normalization infrastructure is extended here with human.task.

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
Loading
🚥 Pre-merge checks | ✅ 3 | ❌ 2

❌ Failed checks (2 warnings)

Check name Status Explanation Resolution
Docstring Coverage ⚠️ Warning Docstring coverage is 3.77% which is insufficient. The required threshold is 80.00%. Write docstrings for the functions missing them to satisfy the coverage threshold.
Description check ⚠️ Warning The description is detailed, but it omits the template's Changes, Related Issues, and Checklist sections. Add the missing Changes, Related Issues, and Checklist sections, and confirm the checklist items or mark N/A where appropriate.
✅ Passed checks (3 passed)
Check name Status Explanation
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.
Title check ✅ Passed The title is concise and accurately summarizes the main change: adding human task workflow steps.
✨ Finishing Touches
📝 Generate docstrings
  • Create stacked PR
  • Commit on current branch
🧪 Generate unit tests (beta)
  • Create PR with unit tests
  • Commit unit tests in branch feat/human-task-backend

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

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

564-566: 🩺 Stability & Availability | 🔵 Trivial | 💤 Low value

Use context.WithoutCancel for cleanup.

If validateSubDAG fails 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 in enqueue.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 value

Consider 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 like TestHumanTaskShapeValidation. Note: It's best omitted from the inner t.Run loop 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

📥 Commits

Reviewing files that changed from the base of the PR and between fce1629 and 1c1ed33.

📒 Files selected for processing (98)
  • README.md
  • README_SCHEMA.md
  • cmd/main.go
  • cmd/main_test.go
  • conformance/harness/runner.go
  • conformance/spec031_human_task/completion_test.go
  • conformance/spec031_human_task/distributed_test.go
  • conformance/spec031_human_task/helpers_test.go
  • conformance/spec031_human_task/lifecycle_test.go
  • conformance/spec031_human_task/testdata/acknowledgement.yaml
  • conformance/spec031_human_task/testdata/additional_properties.yaml
  • conformance/spec031_human_task/testdata/child_human.yaml
  • conformance/spec031_human_task/testdata/concurrent.yaml
  • conformance/spec031_human_task/testdata/distributed_probe.yaml
  • conformance/spec031_human_task/testdata/distributed_root.yaml
  • conformance/spec031_human_task/testdata/distributed_worker_one_probe.yaml
  • conformance/spec031_human_task/testdata/distributed_worker_two_probe.yaml
  • conformance/spec031_human_task/testdata/independent_branch.yaml
  • conformance/spec031_human_task/testdata/invalid_execution_field.yaml
  • conformance/spec031_human_task/testdata/invalid_foreach_human_task.yaml
  • conformance/spec031_human_task/testdata/invalid_form_additional_properties.yaml
  • conformance/spec031_human_task/testdata/invalid_form_null.yaml
  • conformance/spec031_human_task/testdata/invalid_form_property.yaml
  • conformance/spec031_human_task/testdata/invalid_form_required.yaml
  • conformance/spec031_human_task/testdata/invalid_form_root.yaml
  • conformance/spec031_human_task/testdata/invalid_handler_human_task.yaml
  • conformance/spec031_human_task/testdata/invalid_lifecycle_field.yaml
  • conformance/spec031_human_task/testdata/invalid_missing_id.yaml
  • conformance/spec031_human_task/testdata/invalid_output_field.yaml
  • conformance/spec031_human_task/testdata/invalid_prompt.yaml
  • conformance/spec031_human_task/testdata/invalid_with_field.yaml
  • conformance/spec031_human_task/testdata/max_output_size.yaml
  • conformance/spec031_human_task/testdata/multi_document_human.yaml
  • conformance/spec031_human_task/testdata/multiple_waiting.yaml
  • conformance/spec031_human_task/testdata/parent_dag_enqueue.yaml
  • conformance/spec031_human_task/testdata/parent_dag_run.yaml
  • conformance/spec031_human_task/testdata/parent_parallel.yaml
  • conformance/spec031_human_task/testdata/precondition_skipped.yaml
  • conformance/spec031_human_task/testdata/prompt_snapshot.yaml
  • conformance/spec031_human_task/testdata/prompt_snapshot_changed.yaml
  • conformance/spec031_human_task/testdata/sequential_waiting.yaml
  • conformance/spec031_human_task/testdata/typed_form.yaml
  • conformance/spec031_human_task/testdata/valid_acknowledgement_shape.yaml
  • conformance/spec031_human_task/testdata/valid_form_shape.yaml
  • conformance/spec031_human_task/validation_test.go
  • internal/cmd/context.go
  • internal/cmd/human_task.go
  • internal/cmd/human_task_test.go
  • internal/cmn/schema/dag.schema.json
  • internal/cmn/schema/dag_schema_test.go
  • internal/core/dag.go
  • internal/core/dag_test.go
  • internal/core/exec/env.go
  • internal/core/exec/node.go
  • internal/core/foreach_validation_spec018_test.go
  • internal/core/spec/dag.go
  • internal/core/spec/defaults.go
  • internal/core/spec/human_task.go
  • internal/core/spec/human_task_form.go
  • internal/core/spec/human_task_test.go
  • internal/core/spec/loader.go
  • internal/core/spec/step.go
  • internal/core/spec/step_v2.go
  • internal/core/status.go
  • internal/core/step.go
  • internal/core/validator.go
  • internal/core/validator_test.go
  • internal/core/value_fields.go
  • internal/intg/distr/human_task_test.go
  • internal/output/tree.go
  • internal/output/tree_test.go
  • internal/runtime/agent/agent.go
  • internal/runtime/agent/agent_test.go
  • internal/runtime/agent/status_masking.go
  • internal/runtime/agent/status_masking_test.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/data.go
  • internal/runtime/plan.go
  • internal/runtime/runner.go
  • internal/runtime/runner_helper_test.go
  • internal/runtime/runner_test.go
  • internal/runtime/transform/node.go
  • internal/runtime/transform/node_test.go
  • internal/service/chatbridge/notifications.go
  • internal/service/coordinator/handler.go
  • internal/service/frontend/api/v1/dags.go
  • internal/service/mcp/reference.go
  • internal/service/mcp/server_test.go
  • llms.txt
  • skills/dagu/SKILL.md
  • skills/dagu/references/cli.md
  • skills/dagu/references/context.md
  • skills/dagu/references/steptypes.md
  • specs/031-human-task.md
  • specs/README.md
💤 Files with no reviewable changes (1)
  • internal/core/spec/loader.go

Comment thread internal/cmd/human_task.go
Comment thread internal/cmd/human_task.go
Comment thread internal/runtime/builtin/dag/dag.go
Comment thread internal/runtime/runner.go
@yohamta0

Copy link
Copy Markdown
Member Author

@coderabbitai review

@coderabbitai

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

@yohamta0

Copy link
Copy Markdown
Member Author

@coderabbitai review

@coderabbitai

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

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

📥 Commits

Reviewing files that changed from the base of the PR and between 70bc8f5 and 786ea06.

📒 Files selected for processing (1)
  • internal/runtime/foreach_spec018_test.go

Comment thread internal/runtime/foreach_spec018_test.go Outdated
Comment thread internal/runtime/foreach_spec018_test.go
@yohamta0

Copy link
Copy Markdown
Member Author

@coderabbitai review

@coderabbitai

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

@yohamta0
yohamta0 merged commit 35d6178 into main Jul 21, 2026
11 checks passed
@yohamta0
yohamta0 deleted the feat/human-task-backend branch July 21, 2026 04:21
@yohamta0 yohamta0 mentioned this pull request Jul 21, 2026
6 tasks done
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