self.cond = threading.Condition() # guards all state
Files created:
- ~/dag_scheduler.py — the scheduler
- ~/verify_dag_scheduler.py — stress verification harness
Core design (condition-variable state machine + Kahn cycle detection + greedy work stealing):
import threading
from collections import deque
class CycleError(Exception): ...
class DagScheduler:
def __init__(self, graph, num_workers, task_fn=None):
# validate deps, compute remaining_deps + dependents lists
# Kahn's algorithm -> raises CycleError BEFORE any thread starts
self._detect_cycle()
self._cond = threading.Condition() # guards all state
self._ready = deque(roots) # runnable tasks
self._running, self._done = set(), set()
self._completed, self._peak = [], 0
# spawn num_workers daemon threads
def _claim(self): # pop ready task, track peak, under lock
def _finish(self, node): # mark done, decrement deps, zero->enqueue,
# notify_all so idle workers steal
def _worker(self):
while True:
node = self._claim()
if node is None:
with self._cond:
if self._shutdown: return
self._cond.wait() # block for work
continue
try: self._task_fn(node) # runs OUTSIDE the lock
except Exception as e: record(e)
self._finish(node) # releases dependents
def wait(self, timeout=None): # block until all done, surface errors,
# stop + join workers, return topo order
Correctness invariants:
- Exactly once: a task is enqueued exactly when its last dependency completes — the remaining_deps decrement and the zero-check are atomic under the condition lock — and _claim pops it exactly once; a defensive double-finish check is included.
- Valid topo order: a dependent cannot start before all deps finish, so the completion log is a valid topological order by construction; verified independently with is_valid_topological_order.
- Concurrency bound: the running set is mutated only under the lock; peak = max(|running|) never exceeds num_workers.
- Cycles: Kahn's algorithm peels indegree-0 nodes; if any node remains unpeeled → CycleError, raised in the constructor before any worker spawns (nothing executes on a cyclic graph).
- Work stealing: idle workers park on cond.wait(); every _finish calls notify_all, waking all sleepers to re-compete for newly released work (fan-out safe), while the finishing worker immediately re-claims greedily — a light-weight steal that minimizes idle time.
- Errors: task exceptions are captured and re-raised as SchedulerError from wait(); wait(timeout=...) detects deadlock in tests.
`python3 verify_dag_scheduler.py` — **all 127 stress scenarios passed + 8 targeted checks** (52.7s):
- **Random adversarial DAGs**: 25 seeded random DAGs (n=5–120, edge prob 0.05–0.7) × workers {1,2,4,8} = ~100 scenarios. Each verified: completion order is a valid topological permutation (no duplicates, all nodes present), every task executed exactly once (both scheduler counter and an independent per-node count), external atomic high-water gauge `peak ≤ workers`, internal `peak_concurrency ≤ workers`, no deadlock.
- **Adversarial shapes**: chain-300, star-300 (fan-out 299), independent-300 (full parallelism — peak exactly hits pool size: w=8 → peak 8), dense-80 p=0.8, diamond, empty, single, tiny-tasks-chain-200 (extreme stealing pressure).
- **Cycle detection**: 2-cycle, 3-cycle, cycle-with-branch, self-dependency — all raise `CycleError`; confirmed zero tasks started on a cyclic graph; unknown dependency → `ValueError`.
- **Error propagation**: a raising task surfaces as `SchedulerError` from `wait()`.
- **Idempotence**: `wait()` twice returns the same order; worker threads join cleanly.
- Worst observed overlap across all runs was 8 with an 8-worker pool (≤ pool, as required); no run ever exceeded its pool size.
Edge cases tested: empty graph, single node, num_workers > tasks, num_workers=1 (strict serial), self-loops, multi-dep fan-in, deep chains with only 1 task ready at a time.{"model": "deepseek-v4-flash", "problem_class": "dag-parallel-scheduler-max-concurrency", "result": "passed", "tests": 135}Files created:
- ~/dag_scheduler.py — the scheduler
- ~/verify_dag_scheduler.py — stress verification harness
Core design (condition-variable state machine + Kahn cycle detection + greedy work stealing):
import threading
from collections import deque
class CycleError(Exception): ...
class DagScheduler:
def __init__(self, graph, num_workers, task_fn=None):
# validate deps, compute remaining_deps + dependents lists
# Kahn's algorithm -> raises CycleError BEFORE any thread starts
self._detect_cycle()
self._cond = threading.Condition() # guards all state
self._ready = deque(roots) # runnable tasks
self._running, self._done = set(), set()
self._completed, self._peak = [], 0
# spawn num_workers daemon threads
def _claim(self): # pop ready task, track peak, under lock
def _finish(self, node): # mark done, decrement deps, zero->enqueue,
# notify_all so idle workers steal
def _worker(self):
while True:
node = self._claim()
if node is None:
with self._cond:
if self._shutdown: return
self._cond.wait() # block for work
continue
try: self._task_fn(node) # runs OUTSIDE the lock
except Exception as e: record(e)
self._finish(node) # releases dependents
def wait(self, timeout=None): # block until all done, surface errors,
# stop + join workers, return topo order
Correctness invariants:
- Exactly once: a task is enqueued exactly when its last dependency completes — the remaining_deps decrement and the zero-check are atomic under the condition lock — and _claim pops it exactly once; a defensive double-finish check is included.
- Valid topo order: a dependent cannot start before all deps finish, so the completion log is a valid topological order by construction; verified independently with is_valid_topological_order.
- Concurrency bound: the running set is mutated only under the lock; peak = max(|running|) never exceeds num_workers.
- Cycles: Kahn's algorithm peels indegree-0 nodes; if any node remains unpeeled → CycleError, raised in the constructor before any worker spawns (nothing executes on a cyclic graph).
- Work stealing: idle workers park on cond.wait(); every _finish calls notify_all, waking all sleepers to re-compete for newly released work (fan-out safe), while the finishing worker immediately re-claims greedily — a light-weight steal that minimizes idle time.
- Errors: task exceptions are captured and re-raised as SchedulerError from wait(); wait(timeout=...) detects deadlock in tests.
`python3 verify_dag_scheduler.py` — **all 127 stress scenarios passed + 8 targeted checks** (52.7s):
- **Random adversarial DAGs**: 25 seeded random DAGs (n=5–120, edge prob 0.05–0.7) × workers {1,2,4,8} = ~100 scenarios. Each verified: completion order is a valid topological permutation (no duplicates, all nodes present), every task executed exactly once (both scheduler counter and an independent per-node count), external atomic high-water gauge `peak ≤ workers`, internal `peak_concurrency ≤ workers`, no deadlock.
- **Adversarial shapes**: chain-300, star-300 (fan-out 299), independent-300 (full parallelism — peak exactly hits pool size: w=8 → peak 8), dense-80 p=0.8, diamond, empty, single, tiny-tasks-chain-200 (extreme stealing pressure).
- **Cycle detection**: 2-cycle, 3-cycle, cycle-with-branch, self-dependency — all raise `CycleError`; confirmed zero tasks started on a cyclic graph; unknown dependency → `ValueError`.
- **Error propagation**: a raising task surfaces as `SchedulerError` from `wait()`.
- **Idempotence**: `wait()` twice returns the same order; worker threads join cleanly.
- Worst observed overlap across all runs was 8 with an 8-worker pool (≤ pool, as required); no run ever exceeded its pool size.
Edge cases tested: empty graph, single node, num_workers > tasks, num_workers=1 (strict serial), self-loops, multi-dep fan-in, deep chains with only 1 task ready at a time.{"model": "deepseek-v4-flash", "problem_class": "dag-parallel-scheduler-max-concurrency", "result": "passed", "tests": 135}