Repository navigation
refactor: extract DAG-run intake module - #2178
Conversation
|
No actionable comments were generated in the recent review. 🎉 ℹ️ Recent review info⚙️ Run configurationConfiguration used: Organization UI Review profile: CHILL Plan: Pro Run ID: 📒 Files selected for processing (4)
🚧 Files skipped from review as they are similar to previous changes (4)
📝 WalkthroughWalkthroughAdd 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. ChangesDAG Run Intake Abstraction and Adoption
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
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
Estimated code review effort🎯 4 (Complex) | ⏱️ ~45 minutes Possibly related PRs
🚥 Pre-merge checks | ✅ 3 | ❌ 2❌ Failed checks (1 warning, 1 inconclusive)
✅ Passed checks (3 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ 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: 1
🧹 Nitpick comments (2)
internal/dagrun/intake/queue_test.go (1)
69-91: 💤 Low valueTest could verify
StatusCloseErris propagated.The test verifies that
EnqueueRunsucceeds whenProceedOnStatusCloseErr: trueand close fails, but it doesn't assert thatqueued.StatusCloseErrcontains 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 valueUnconventional 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 whenproceedOnCloseErris 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
📒 Files selected for processing (11)
internal/cmd/enqueue.gointernal/cmd/local_execution.gointernal/cmd/start.gointernal/dagrun/intake/local.gointernal/dagrun/intake/local_test.gointernal/dagrun/intake/queue.gointernal/dagrun/intake/queue_test.gointernal/runtime/builtin/dag/enqueue.gointernal/service/frontend/api/v1/dagruns.gointernal/service/scheduler/enqueue.gointernal/service/scheduler/enqueue_webhook.go
💤 Files with no reviewable changes (1)
- internal/cmd/start.go
|
@coderabbitai review |
✅ Actions performedReview triggered.
|
Summary
Testing
Summary by CodeRabbit