Skip to content

refactor: extract DAG-run intake module - #2178

Merged
yohamta0 merged 2 commits into
mainfrom
refactor/dagrun-intake
May 19, 2026
Merged

yohamta0 merged 2 commits into
mainfrom
refactor/dagrun-intake

Conversation

@yohamta0

@yohamta0 yohamta0 commented May 19, 2026 •

Copy link
Copy Markdown
Member

Summary

  • add internal/dagrun/intake to centralize queued DAG-run and local execution intake
  • migrate CLI, scheduler, sub-DAG executor, and frontend API queue/local setup through the module
  • add focused intake tests for queued status, rollback, priority, and local proc acquisition failure behavior

Testing

  • go test ./internal/dagrun/intake ./internal/cmd ./internal/service/scheduler ./internal/runtime/builtin/dag ./internal/service/frontend/api/v1

Summary by CodeRabbit

  • Refactor
    • Centralized DAG run enqueueing into a single intake flow for more consistent and reliable queue behavior.
    • Centralized local execution preparation and proc acquisition to simplify setup and error handling.
  • Tests
    • Added unit tests covering enqueue and local-prep failure and rollback scenarios to improve resilience and correctness.

Review Change Stack

@coderabbitai

coderabbitai Bot commented May 19, 2026 •

Copy link
Copy Markdown

No actionable comments were generated in the recent review. 🎉

ℹ️ Recent review info
⚙️ Run configuration

Configuration used: Organization UI

Review profile: CHILL

Plan: Pro

Run ID: ff309e94-9449-45dd-b49a-034f8dece3a7

📥 Commits

Reviewing files that changed from the base of the PR and between 96ff2a6 and 6eb9eb9.

📒 Files selected for processing (4)
  • internal/dagrun/intake/local.go
  • internal/dagrun/intake/local_test.go
  • internal/dagrun/intake/queue.go
  • internal/dagrun/intake/queue_test.go
🚧 Files skipped from review as they are similar to previous changes (4)
  • internal/dagrun/intake/local_test.go
  • internal/dagrun/intake/local.go
  • internal/dagrun/intake/queue_test.go
  • internal/dagrun/intake/queue.go

📝 Walkthrough

Walkthrough

Add intake APIs: intake.EnqueueRun centralizes queued-run creation, status persistence, and queue publish; intake.PrepareLocalExecution centralizes local execution setup and proc acquisition. Several callers (cmd, runtime, frontend, scheduler) now delegate enqueue/local-prepare responsibilities to these intake functions.

Changes

DAG Run Intake Abstraction and Adoption

Layer / File(s) Summary
Queue intake foundation
internal/dagrun/intake/queue.go, internal/dagrun/intake/queue_test.go
QueueRequest/QueuedRun defined. EnqueueRun validates inputs, creates a queued attempt, writes queued status, closes attempt (optionally proceeds on close error), enqueues at low priority, and rolls back created run on post-create failure. Tests cover write-order, rollback, and proceed-on-close behavior.
Local execution intake foundation
internal/dagrun/intake/local.go, internal/dagrun/intake/local_test.go
LocalRequest/LocalPreparation defined. PrepareLocalExecution validates inputs, locks proc group, builds attempt via callback, sets DAG context, acquires proc handle, and records failure status if acquisition fails after attempt creation. Tests cover success, acquisition failure, and failure-to-record status.
Queue intake adoption across callers
internal/cmd/enqueue.go, internal/runtime/builtin/dag/enqueue.go, internal/service/frontend/api/v1/dagruns.go, internal/service/scheduler/enqueue.go, internal/service/scheduler/enqueue_webhook.go
Callers refactored to call intake.EnqueueRun with a populated QueueRequest, removing inline log/artifact/status/queue orchestration and updating imports to use intake. Returned QueuedRun.StatusCloseErr is logged/warned when present.
Local execution intake adoption
internal/cmd/local_execution.go
withPreparedLocalExecution now calls intake.PrepareLocalExecution with a LocalRequest containing proc store, DAG/run refs, trigger/schedule values, base log/artifact dirs, and the attempt builder callback.
Cleanup
internal/cmd/start.go
Removed unused package-level error variable previously referenced by the in-file local-exec helper.

Sequence Diagram(s)

sequenceDiagram
  participant Client as Enqueue Caller
  participant ER as intake.EnqueueRun
  participant DRS as DAGRunStore
  participant QS as QueueStore
  Client->>ER: QueueRequest
  ER->>DRS: CreateDAGRun
  DRS-->>ER: DAGRunAttempt
  ER->>ER: write queued status
  ER->>DRS: Close attempt
  ER->>QS: Enqueue reference
  QS-->>ER: success/error
  alt queue fails
    ER->>DRS: RemoveDAGRun
  end
  ER-->>Client: QueuedRun or error
Loading
sequenceDiagram
  participant Client as Local Execution Caller
  participant PE as intake.PrepareLocalExecution
  participant PS as ProcStore
  participant A as Attempt
  Client->>PE: LocalRequest
  PE->>PS: Lock process group
  PE->>A: Build attempt
  PE->>A: Set DAG context
  PE->>PS: Acquire proc handle
  alt acquire succeeds
    PE-->>Client: LocalPreparation
  else acquire fails
    PE->>A: record failed status
    PE->>A: close
    PE->>PS: Unlock
    PE-->>Client: error
  end
Loading

Estimated code review effort

🎯 4 (Complex) | ⏱️ ~45 minutes

Possibly related PRs

  • dagucloud/dagu#2163: Related refactor of internal/runtime/builtin/dag/enqueue.go to use intake.EnqueueRun.
🚥 Pre-merge checks | ✅ 3 | ❌ 2

❌ Failed checks (1 warning, 1 inconclusive)

Check name Status Explanation Resolution
Docstring Coverage ⚠️ Warning Docstring coverage is 22.73% which is insufficient. The required threshold is 80.00%. Write docstrings for the functions missing them to satisfy the coverage threshold.
Description check ❓ Inconclusive The description includes a summary of changes and testing approach, but lacks the structured sections (Changes list, Related Issues, Checklist) specified in the template. Expand the description to follow the template more closely by adding a structured Changes list, clarifying related issues, and completing the checklist items.
✅ Passed checks (3 passed)
Check name Status Explanation
Title check ✅ Passed The title 'refactor: extract DAG-run intake module' clearly and concisely summarizes the main change: extracting intake logic into a centralized module.
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 docstrings
  • Create stacked PR
  • Commit on current branch
🧪 Generate unit tests (beta)
  • Create PR with unit tests
  • Commit unit tests in branch refactor/dagrun-intake

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 (2)
internal/dagrun/intake/queue_test.go (1)

69-91: 💤 Low value

Test could verify StatusCloseErr is propagated.

The test verifies that EnqueueRun succeeds when ProceedOnStatusCloseErr: true and close fails, but it doesn't assert that queued.StatusCloseErr contains the injected error. This would strengthen the test coverage for the close-error propagation behavior.

💡 Suggested assertion
 	require.NoError(t, err)
 	require.NotNil(t, queued)
 	assert.True(t, f.queueStore.enqueued)
 	assert.True(t, f.attempt.closed)
 	assert.False(t, f.runStore.removed)
+	assert.Error(t, queued.StatusCloseErr)
+	assert.Contains(t, queued.StatusCloseErr.Error(), "sync failed")
🤖 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/dagrun/intake/queue_test.go` around lines 69 - 91, Update
TestEnqueueRunCanProceedWhenAttemptCloseFails to assert that the error from the
failing attempt is propagated into the returned queued object's StatusCloseErr:
after calling EnqueueRun (which uses f.attempt.closeErr) add an assertion that
queued.StatusCloseErr reflects f.attempt.closeErr (either by comparing the error
value/string or using an error-contains assertion) so the test verifies
close-error propagation from the attempt into queued.
internal/dagrun/intake/queue.go (1)

191-206: 💤 Low value

Unconventional return signature for writeQueuedStatus.

The (error, error) return type is unusual and can be confusing for callers. The first return value is the close error (which may be non-fatal when proceedOnCloseErr is true), while the second is the fatal error. Consider using a named struct or documenting this more explicitly in a comment.

🤖 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/dagrun/intake/queue.go` around lines 191 - 206, The
writeQueuedStatus function uses an unconventional (error, error) return
signature which is confusing; change it to a clear, single-return type (either a
named struct or a single error) and update callers. For example, define a small
result type (e.g., WriteStatusResult with CloseErr error and Err error) or
return one error where non-fatal close errors are wrapped/annotated (use
exec.DAGRunAttempt, writeQueuedStatus, and proceedOnCloseErr to detect and
populate the appropriate field/message), update the function signature and
propagate the new type or wrapped error to all call sites, and add a short
comment on the semantics of the fields if you choose the struct approach.
🤖 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/dagrun/intake/local.go`:
- Around line 90-91: The code currently discards errors from
recordPreparedAttemptFailure in the error path that returns
ErrProcAcquisitionFailed; change this so the recordPreparedAttemptFailure error
is not ignored — call recordPreparedAttemptFailure(ctx, req, attempt, err),
capture its returned error (e.g., recErr), and surface it by either returning a
wrapped/combined error that includes both ErrProcAcquisitionFailed and the
record error or at minimum logging/returning recErr alongside the original
error; ensure the returned error from this branch references both
ErrProcAcquisitionFailed and the recordPreparedAttemptFailure failure so
terminal-state write failures are observable.

---

Nitpick comments:
In `@internal/dagrun/intake/queue_test.go`:
- Around line 69-91: Update TestEnqueueRunCanProceedWhenAttemptCloseFails to
assert that the error from the failing attempt is propagated into the returned
queued object's StatusCloseErr: after calling EnqueueRun (which uses
f.attempt.closeErr) add an assertion that queued.StatusCloseErr reflects
f.attempt.closeErr (either by comparing the error value/string or using an
error-contains assertion) so the test verifies close-error propagation from the
attempt into queued.

In `@internal/dagrun/intake/queue.go`:
- Around line 191-206: The writeQueuedStatus function uses an unconventional
(error, error) return signature which is confusing; change it to a clear,
single-return type (either a named struct or a single error) and update callers.
For example, define a small result type (e.g., WriteStatusResult with CloseErr
error and Err error) or return one error where non-fatal close errors are
wrapped/annotated (use exec.DAGRunAttempt, writeQueuedStatus, and
proceedOnCloseErr to detect and populate the appropriate field/message), update
the function signature and propagate the new type or wrapped error to all call
sites, and add a short comment on the semantics of the fields if you choose the
struct approach.
🪄 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: f57ac32c-ce61-4e34-8d6a-4931316f39bf

📥 Commits

Reviewing files that changed from the base of the PR and between c170b8d and 96ff2a6.

📒 Files selected for processing (11)
  • internal/cmd/enqueue.go
  • internal/cmd/local_execution.go
  • internal/cmd/start.go
  • internal/dagrun/intake/local.go
  • internal/dagrun/intake/local_test.go
  • internal/dagrun/intake/queue.go
  • internal/dagrun/intake/queue_test.go
  • internal/runtime/builtin/dag/enqueue.go
  • internal/service/frontend/api/v1/dagruns.go
  • internal/service/scheduler/enqueue.go
  • internal/service/scheduler/enqueue_webhook.go
💤 Files with no reviewable changes (1)
  • internal/cmd/start.go

Comment thread internal/dagrun/intake/local.go Outdated

Copy link
Copy Markdown
Member Author

@coderabbitai review

@coderabbitai

coderabbitai Bot commented May 19, 2026

Copy link
Copy Markdown
✅ Actions performed

Review triggered.

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 a955792 into main May 19, 2026
10 checks passed
@yohamta0
yohamta0 deleted the refactor/dagrun-intake branch May 19, 2026 14:48
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