Repository navigation
[codex] add dag scheduler core - #2465
Andriy Knysh (aknysh) merged 3 commits into
Conversation
|
Tip Atmos Pro
No affected stacks workflow was detected for this pull request. |
Dependency Review✅ No vulnerabilities or license issues found.Scanned FilesNone |
Codecov Report❌ Patch coverage is
Additional details and impacted files@@ Coverage Diff @@
## codex/dag-process-io-foundation #2465 +/- ##
===================================================================
- Coverage 78.58% 78.55% -0.03%
===================================================================
Files 1146 1144 -2
Lines 110111 110044 -67
===================================================================
- Hits 86529 86444 -85
- Misses 18778 18798 +20
+ Partials 4804 4802 -2
Flags with carried forward coverage won't be shown. Click here to find out more.
🚀 New features to boost your workflow:
|
7e70c17 to
6a5036e
Compare
cfd2b2a to
e3f69de
Compare
e3f69de to
1c61093
Compare
|
CodeRabbit (@coderabbitai) review |
✅ Actions performedReview triggered.
|
|
Warning Review limit reached
More reviews will be available in 2 minutes and 27 seconds. Learn how PR review limits work. Your organization has run out of usage credits. Purchase more in the billing tab. ⌛ How to resolve this issue?After more reviews become available, a review can be triggered using the We recommend that you space out your commits to avoid hitting the rate limit. 🚦 How do rate limits work?CodeRabbit enforces hourly rate limits for each developer per organization. Our paid plans include higher PR review limits than trial, open-source, and free plans. In all cases, reviews become available again over time. During sustained high-volume PR review activity, CodeRabbit may temporarily slow when the next review becomes available. Please see our Fair Usage Limits Policy for further information. ℹ️ Review info⚙️ Run configurationConfiguration used: Path: .coderabbit.yaml Review profile: CHILL Plan: Pro Run ID: 📒 Files selected for processing (3)
📝 WalkthroughWalkthroughAdds a concurrent DAG scheduler with stable topological ordering, a bounded worker pool, configurable fail-fast behavior and lifecycle hooks, deterministic aggregate results, and a comprehensive test suite exercising execution semantics, concurrency, and error handling. ChangesDAG Scheduler Implementation
Sequence DiagramsequenceDiagram
participant Client
participant Scheduler
participant Validator
participant WorkerPool
participant Dispatcher
participant Aggregator
Client->>Scheduler: Run(ctx)
Scheduler->>Validator: validate graph & topo order
Validator-->>Scheduler: ordered node IDs
Scheduler->>WorkerPool: start workers (maxConcurrency)
loop ready nodes
Scheduler->>WorkerPool: enqueue nodeID
WorkerPool->>Dispatcher: Dispatch(ctx, node)
Dispatcher-->>WorkerPool: Result or error
WorkerPool->>Scheduler: emit completion event (Result)
Scheduler->>Scheduler: recordEvent, enqueue newly-ready dependents
end
Scheduler->>Scheduler: apply fail-fast / skip propagation / cancellation handling
Scheduler->>Aggregator: collect ordered results, join errors
Aggregator-->>Client: AggregateResult
Estimated code review effort🎯 4 (Complex) | ⏱️ ~50 minutes Suggested labels
Suggested reviewers
🚥 Pre-merge checks | ✅ 4 | ❌ 1❌ Failed checks (1 warning)
✅ Passed checks (4 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches🧪 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 |
|
CodeRabbit (@coderabbitai) review |
✅ Actions performedReview triggered.
|
There was a problem hiding this comment.
Actionable comments posted: 4
🤖 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 `@pkg/scheduler/scheduler.go`:
- Around line 13-20: The new sentinel errors (ErrNilGraph, ErrNilDispatcher,
ErrNodeSkipped, ErrNodeNotFound, ErrInvalidGraph, ErrInvalidWorker) in
pkg/scheduler/scheduler.go duplicate the repo's canonical error definitions;
replace these local declarations by importing and reusing the static errors from
errors/errors.go and wrap them at call sites instead of redefining them. Locate
usages of ErrNilGraph, ErrNilDispatcher, ErrNodeSkipped, ErrNodeNotFound,
ErrInvalidGraph, and ErrInvalidWorker in scheduler.go and other scheduler code,
switch to the shared errors package identifiers, and where you need additional
context wrap the shared error with fmt.Errorf("%s: %w", ctx, sharedErr) or
errors.Join as appropriate so callers can still use errors.Is() against the
static errors.
- Around line 275-277: The early return after checking event.result.Status ==
StatusSucceeded ignores StatusSkipped and leaves dependents' in-degree
unchanged; update the dispatcher handling in scheduler.go so that StatusSkipped
is treated as a terminal blocking result: when event.result.Status is
StatusSkipped (in addition to StatusSucceeded) reduce dependents' in-degree and
immediately propagate a skipped Result to downstream nodes (similar to how
failures are handled), ensuring the outer loop's skip-path (which currently
checks StatusFailed) sees and processes skipped dependencies rather than
stalling; adjust logic around event.result.Status checks and any in-degree
decrement/enqueueing paths in the functions handling Result/dispatcher outcomes
(references: event.result.Status, StatusSucceeded, StatusSkipped, Result,
dispatcher, dependents, in-degree).
- Around line 205-209: The fail-fast branch is calling cancel(), which cancels
the shared context and aborts already-running sibling workers; instead remove
the cancel() call and only set the stopping flag and mark pending nodes skipped
(retain the finished += skipPending(states, results, s.graph, fmt.Errorf(...))
line). Update the logic around event.result.Status == StatusFailed (and the use
of s.failFast, stopping) so new scheduling is prevented but existing started
dispatchers are allowed to drain; do not propagate context cancellation from
this handler to running workers.
- Around line 137-138: Add performance tracking to the public Scheduler.Run
method by inserting a deferred perf.Track call at the top of the function: call
defer perf.Track(atmosConfig, "pkg.scheduler.Run")() immediately after the
function signature and add a blank line after that defer. Locate the Run method
on type Scheduler in scheduler.go and ensure the perf.Track import/symbol is
available (or add the import) so the defer compiles.
🪄 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: Path: .coderabbit.yaml
Review profile: CHILL
Plan: Pro
Run ID: 12e17c32-f79a-4ead-88dd-01cfafa9b40e
📒 Files selected for processing (2)
pkg/scheduler/scheduler.gopkg/scheduler/scheduler_test.go
|
CodeRabbit (@coderabbitai) review |
✅ Actions performedReview triggered.
|
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 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 `@pkg/scheduler/scheduler.go`:
- Around line 295-297: When a node returns StatusSkipped or StatusFailed with a
nil Err the aggregate can appear successful; update the logic in Run (and the
similar block around the 365-372 snippet) so that whenever a
Result/AggregateResult has a non-success status
(StatusSkipped/StatusFailed/etc.) but Err == nil you synthesize and attach a
sentinel error (e.g., fmt.Errorf or errors.New that includes the status and node
id) before returning or before calling aggregateError; apply this change for the
branch checking result.Status == StatusSkipped and the analogous branch handling
failed/skipped nodes so AggregateResult.Err always reflects non-success
statuses.
🪄 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: Path: .coderabbit.yaml
Review profile: CHILL
Plan: Pro
Run ID: f77fb7d3-c0d0-416d-8649-071f85e631e2
📒 Files selected for processing (3)
errors/errors.gopkg/scheduler/scheduler.gopkg/scheduler/scheduler_test.go
|
CodeRabbit (@coderabbitai) review |
✅ Actions performedReview triggered.
|
3467269
into
codex/dag-process-io-foundation
|
Tip Atmos Pro
No affected stacks workflow was detected for this pull request. |
* Add process and I/O execution foundation * Address process I/O review feedback * Address CodeRabbit rereview feedback * Rename output pipeline API * [codex] add dag scheduler core (#2465) * add dag scheduler core * Address CodeRabbit scheduler feedback * Ensure scheduler aggregates non-success statuses --------- Co-authored-by: Andriy Knysh <aknysh@users.noreply.github.com>
Summary
Adds the PR 2 scheduler foundation as a new generic
pkg/schedulerpackage.The scheduler consumes the existing dependency graph package and delegates node execution through a generic dispatcher interface. It implements ready-queue DAG scheduling, bounded workers, fail-fast and keep-going behavior, skipped dependent propagation, deterministic aggregate result ordering, and lifecycle hooks for future orchestration layers.
Scope
pkg/schedulerwith generic scheduling types andScheduler.Run.Stacking
This PR is stacked on PR 1 and targets
codex/dag-process-io-foundation.Supersedes the earlier fork-headed draft #2461 now that the stack branches exist in
cloudposse/atmos.Validation
Summary by CodeRabbit
New Features
Tests