Skip to content

[codex] add dag scheduler core - #2465

Merged
Andriy Knysh (aknysh) merged 3 commits into
codex/dag-process-io-foundationfrom
codex/dag-scheduler-core
May 29, 2026
Merged

Andriy Knysh (aknysh) merged 3 commits into
codex/dag-process-io-foundationfrom
codex/dag-scheduler-core

Conversation

@shirkevich

@shirkevich Mikhail Shirkov (shirkevich) commented May 21, 2026 •

Copy link
Copy Markdown
Collaborator

Summary

Adds the PR 2 scheduler foundation as a new generic pkg/scheduler package.

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

  • Added pkg/scheduler with generic scheduling types and Scheduler.Run.
  • Added isolated unit tests for single node, linear chain, diamond, fan-out, fan-in, ready-queue behavior, worker bounds, deterministic aggregate ordering, error propagation, fail-fast, keep-going, validation, cycle rejection, and hooks.
  • Intentionally did not change Terraform routing, add tool adapters, or expose concurrency flags.

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

GOCACHE=/private/tmp/atmos-go-cache go test -count=20 ./pkg/scheduler
GOCACHE=/private/tmp/atmos-go-cache go test ./pkg/scheduler ./pkg/process ./pkg/io ./pkg/dependency
GOCACHE=/private/tmp/atmos-go-cache go test -race ./pkg/scheduler

Summary by CodeRabbit

  • New Features

    • Added a DAG scheduler for deterministic execution of dependency graphs with configurable concurrency.
    • Supports fail-fast and continue-on-error modes; dependent tasks are automatically skipped when prerequisites fail.
    • Execution hooks for task start and completion and deterministic aggregated results with combined error reporting.
  • Tests

    • Comprehensive test suite validating ordering, concurrency bounds, error propagation, hooks, and input validation.

Review Change Stack

@atmos-pro

atmos-pro Bot commented May 21, 2026 •

Copy link
Copy Markdown
Contributor

Tip

Atmos Pro  

No affected stacks workflow was detected for this pull request.
If this is expected, no action is needed.
Learn More. Ask AI.

@github-actions github-actions Bot added the size/m Medium size PR label May 21, 2026
@github-actions

github-actions Bot commented May 21, 2026 •

Copy link
Copy Markdown

Dependency Review

✅ No vulnerabilities or license issues found.

Scanned Files

None

@mergify

mergify Bot commented May 21, 2026

Copy link
Copy Markdown
Contributor

⚠️ The sha of the head commit of this PR conflicts with #2461. Mergify cannot evaluate rules on this PR. Once #2461 is merged or closed, Mergify will resume processing this PR. ⚠️

@codecov

codecov Bot commented May 21, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 87.31884% with 35 lines in your changes missing coverage. Please review.
✅ Project coverage is 78.55%. Comparing base (3ba4c03) to head (9aecbc1).

Files with missing lines Patch % Lines
pkg/scheduler/scheduler.go 87.31% 26 Missing and 9 partials ⚠️
Additional details and impacted files

Impacted file tree graph

@@                         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     
Flag Coverage Δ
unittests 78.55% <87.31%> (-0.03%) ⬇️

Flags with carried forward coverage won't be shown. Click here to find out more.

Files with missing lines Coverage Δ
errors/errors.go 100.00% <ø> (ø)
pkg/scheduler/scheduler.go 87.31% <87.31%> (ø)

... and 29 files with indirect coverage changes

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.

@shirkevich
Mikhail Shirkov (shirkevich) force-pushed the codex/dag-process-io-foundation branch from 7e70c17 to 6a5036e Compare May 25, 2026 07:46
@shirkevich
Mikhail Shirkov (shirkevich) marked this pull request as ready for review May 26, 2026 18:07
@shirkevich
Mikhail Shirkov (shirkevich) requested a review from a team as a code owner May 26, 2026 18:07
@shirkevich

Copy link
Copy Markdown
Collaborator Author

CodeRabbit (@coderabbitai) review

@coderabbitai

coderabbitai Bot commented May 27, 2026

Copy link
Copy Markdown
Contributor
✅ 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.

@coderabbitai

coderabbitai Bot commented May 27, 2026 •

Copy link
Copy Markdown
Contributor

Warning

Review limit reached

@shirkevich, we couldn't start this review because you've reached your PR review rate limit.

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 @coderabbitai review command as a PR comment. Alternatively, push new commits to this PR.

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 configuration

Configuration used: Path: .coderabbit.yaml

Review profile: CHILL

Plan: Pro

Run ID: 68653d10-b9eb-4fe8-9edb-a26e5957a433

📥 Commits

Reviewing files that changed from the base of the PR and between 17ed1dc and 9aecbc1.

📒 Files selected for processing (3)
  • errors/errors.go
  • pkg/scheduler/scheduler.go
  • pkg/scheduler/scheduler_test.go
📝 Walkthrough

Walkthrough

Adds 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.

Changes

DAG Scheduler Implementation

Layer / File(s) Summary
Error and public API contracts
errors/errors.go, pkg/scheduler/scheduler.go
Sentinel scheduler errors; Status constants; Dispatcher/DispatcherFunc; Result/AggregateResult; Scheduler, Option pattern, and New constructor.
Run orchestration and scheduling loop
pkg/scheduler/scheduler.go
(*Scheduler).Run initializes state, validates inputs, computes topological order, starts bounded workers, schedules ready nodes deterministically, enforces fail-fast/keep-going semantics, and returns ordered AggregateResult.
Input validation & topological ordering
pkg/scheduler/scheduler.go
Validates nil inputs and invalid concurrency, rejects cyclic graphs, and computes deterministic topological node IDs (Kahn algorithm).
Worker goroutines & dispatch
pkg/scheduler/scheduler.go
Worker loop that fetches nodes, invokes node start hook, dispatches via Dispatcher, normalizes results, invokes completion hook, and emits completion events.
Result constructors & skip propagation
pkg/scheduler/scheduler.go
Creates failed/skipped results, derives skip-cause errors, marks pending nodes skipped on fail-fast or cancellation, and recursively skips dependents of failed/skipped nodes.
Aggregation, in-degrees, ordering helpers
pkg/scheduler/scheduler.go
Aggregates non-nil errors with errors.Join, orders results in topological order, computes in-degrees and root nodes, and provides stable sorting helpers.
Internal run primitives & ready queue
pkg/scheduler/scheduler.go
Defines runEvent, nodeState, and a deterministic readyQueue used by scheduling and topological ordering.
Comprehensive test suite
pkg/scheduler/scheduler_test.go
Tests DAG execution (linear, diamond, fan-in/out), ready-queue behavior, concurrency bounds, deterministic aggregation ordering, failure modes (keep-going, fail-fast, drain), skipped-status propagation, input validation, cycle rejection, and node start/complete hooks. Includes helpers for graph building, timeouts, and assertions.

Sequence Diagram

sequenceDiagram
  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
Loading

Estimated code review effort

🎯 4 (Complex) | ⏱️ ~50 minutes

Suggested labels

minor

Suggested reviewers

  • aknysh
🚥 Pre-merge checks | ✅ 4 | ❌ 1

❌ Failed checks (1 warning)

Check name Status Explanation Resolution
Docstring Coverage ⚠️ Warning Docstring coverage is 11.63% which is insufficient. The required threshold is 80.00%. Write docstrings for the functions missing them to satisfy the coverage threshold.
✅ Passed checks (4 passed)
Check name Status Explanation
Description Check ✅ Passed Check skipped - CodeRabbit’s high-level summary is enabled.
Title check ✅ Passed The title '[codex] add dag scheduler core' directly and clearly summarizes the main change: adding a DAG scheduler implementation to the codebase.
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 unit tests (beta)
  • Create PR with unit tests
  • Commit unit tests in branch codex/dag-scheduler-core

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.

@shirkevich

Copy link
Copy Markdown
Collaborator Author

CodeRabbit (@coderabbitai) review

@coderabbitai

coderabbitai Bot commented May 27, 2026

Copy link
Copy Markdown
Contributor
✅ 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.

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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

📥 Commits

Reviewing files that changed from the base of the PR and between 3ba4c03 and 1c61093.

📒 Files selected for processing (2)
  • pkg/scheduler/scheduler.go
  • pkg/scheduler/scheduler_test.go

Comment thread pkg/scheduler/scheduler.go Outdated
Comment thread pkg/scheduler/scheduler.go
Comment thread pkg/scheduler/scheduler.go Outdated
Comment thread pkg/scheduler/scheduler.go
@github-actions github-actions Bot added size/l Large size PR and removed size/m Medium size PR labels May 27, 2026
@shirkevich

Copy link
Copy Markdown
Collaborator Author

CodeRabbit (@coderabbitai) review

@coderabbitai

coderabbitai Bot commented May 27, 2026

Copy link
Copy Markdown
Contributor
✅ 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.

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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

📥 Commits

Reviewing files that changed from the base of the PR and between 1c61093 and 17ed1dc.

📒 Files selected for processing (3)
  • errors/errors.go
  • pkg/scheduler/scheduler.go
  • pkg/scheduler/scheduler_test.go

Comment thread pkg/scheduler/scheduler.go
@shirkevich

Copy link
Copy Markdown
Collaborator Author

CodeRabbit (@coderabbitai) review

@coderabbitai

coderabbitai Bot commented May 27, 2026

Copy link
Copy Markdown
Contributor
✅ 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.

@aknysh
Andriy Knysh (aknysh) merged commit 3467269 into codex/dag-process-io-foundation May 29, 2026
52 checks passed
@atmos-pro

atmos-pro Bot commented May 29, 2026 •

Copy link
Copy Markdown
Contributor

Tip

Atmos Pro  

No affected stacks workflow was detected for this pull request.
If this is expected, no action is needed.
Learn More. Ask AI.

@aknysh
Andriy Knysh (aknysh) deleted the codex/dag-scheduler-core branch May 29, 2026 23:33
Andriy Knysh (aknysh) added a commit that referenced this pull request May 30, 2026
* 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>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

size/l Large size PR stacked Stacked

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants