11 KiB
Explore-note — verify before trusting. Distilled map of a subsystem, not authoritative. Last verified against commit
cc90600(2026-08-06). Drift check:git log --oneline cc90600..HEAD -- src/ClaudeDo.WorkerStable structure only (no line numbers). See docs/explore-notes/README.md.
Worker: Task Execution Pipeline
How a task moves Queued → Running → terminal, across src/ClaudeDo.Worker
(Queue, Runner, Lifecycle, State, Agents, Worktrees, Hub).
End-to-End Flow (Queued → Terminal)
-
Enqueue —
ITaskStateService.EnqueueAsync()(State/TaskStateService.cs)- Idle → Queued, then wakes the dispatcher via
IQueueWaker.Wake().
- Idle → Queued, then wakes the dispatcher via
-
Dispatch —
QueueServiceloop (Queue/QueueService.cs)BackgroundService; waits for a wake signal or a backstop timer.- Reads the max-parallel limit from settings; claims a free slot if under limit.
-
Atomic Claim —
IQueuePicker.ClaimNextAsync()(Queue/QueuePicker.cs)- Raw SQL
UPDATE ... RETURNINGin one transaction: picks an eligible Queued task (unblocked, due or unscheduled; sorted by sort_order/created_at), sets status→Running- started_at, returns the row. Prevents two workers claiming the same task (TOCTOU).
- Raw SQL
-
Slot Execution —
QueueService.RunInSlotAsync()(Queue/QueueService.cs)- For review feedback: resume the prior session if one exists, else fold feedback into
the prompt. Calls
TaskRunner.RunAsync()/ContinueAsync()withalreadyClaimed=true.
- For review feedback: resume the prior session if one exists, else fold feedback into
the prompt. Calls
-
Run Preparation —
TaskRunner.RunAsync()(Runner/TaskRunner.cs)- Loads task, list config, subtasks, attachments from the DB.
StartRunningAsync()(only if not pre-claimed): atomic claim to Running, before any resource is created. A rejected claim (task already Running) bails out immediately — no worktree, no MCP token file. Broadcasts TaskStarted.PrepareRunDirectoryAsync(): worktree (via WorktreeManager) if the list has a WorkingDir, else sandbox. Generates a per-run MCP token, writes MCP config to disk.
-
Claude Execution —
TaskRunner.RunOnceAsync()(Runner/TaskRunner.cs)- Creates a TaskRunEntity, points the task at the run's log path.
- Builds claude CLI args (ClaudeArgsBuilder), spawns the process via
IClaudeProcess.RunAsync()with prompt + working dir + streaming callback. - Stream lines → NDJSON log + broadcast via TaskMessage. MCP tools (AskUser, SuggestImprovement) are scoped by the per-run token.
-
Result Handling —
TaskRunner.HandleSuccess()/MarkFailed()(Runner/TaskRunner.cs)- Success (exit 0 + result markdown): if worktree, commit + broadcast WorktreeUpdated; then transition to Done / WaitingForReview / WaitingForChildren (CompleteAsync / SubmitForReviewAsync / SubmitForChildrenAsync).
- Failure: if a session exists, auto-retry once via ContinueAsync; else MarkFailed → FailAsync.
- All terminal writes use
CancellationToken.Noneso a task is never left Running.
-
Terminal States —
ITaskStateServicetransitions (State/TaskStateService.cs)- Done CompleteAsync (Running → Done) — top-level success.
- WaitingForReview SubmitForReviewAsync (Running → WaitingForReview) — review gate.
- WaitingForChildren SubmitForChildrenAsync (Running → WaitingForChildren) — blocks on children.
Advances to WaitingForReview via
TryAdvanceParentAsynconce every remaining child is terminal (Done/Failed/Cancelled) — including zero children left, e.g. after the last child is deleted (WorkerHub.DeleteTask/ExternalMcpService.DeleteTaskboth call it). - Failed FailAsync (Running/Queued → Failed).
- Cancelled CancelAsync (Running/Queued/WaitingForReview/WaitingForChildren → Cancelled).
Model, effort & max-turns resolution
(section added at commit f6cb825, 2026-08-05; resolver extraction added same day)
The resolution below lives in Runner/EffectiveRunConfigResolver.Resolve (not inlined in
TaskRunner anymore) so TaskRunner.ResolveConfigAsync and the read-only
get_effective_run_config MCP tool (External/ConfigMcpTools.cs) share one codepath and can't
report different numbers for the same task. The tool additionally surfaces, per field, whether
it came from the task/list/preset/global layer, and — for max turns — the raw requested value
plus whether it was clamped.
Step 6 builds the CLI args. Model and turn budget resolve like this:
- Effective model — task override → list config →
AppSettings.DefaultModel. - Preset row —
ModelPresets.For(global.ModelPresets, model, global.DefaultMaxTurns). The model string is resolved throughModelRegistry.TryNormalizeAliasfirst, so a full CLI model id (e.g.claude-sonnet-4-6, not just the baresonnet/opus/haiku/fablealiases) still hits its alias's preset row instead of missing every lookup. Only a model that normalizes to nothing recognized falls back to a synthesized row usingAppSettings.DefaultMaxTurns— never a hardcoded number, and it never throws: an unknown model must not block a run. - The preset supplies
--effortand the global max-turns default. Task/listMaxTurnsoverrides still win over it. - Ceiling clamp —
TaskRunner.ResolveMaxTurnshard-clamps the resolved value toAppSettings.MaxTurnsCeiling(default 80). An override above the ceiling still starts, just capped, and a Warn logs the task id + requested + effective value.
⚠️ Trap: if app_settings.model_presets is somehow null, the fallback path decides the turn
budget — which is why AppSettingsRepository.GetAsync backfills shipping defaults on the first
read after null. Ship preset turns are low (haiku 20, sonnet 30, opus 40, fable 25), so a task
that genuinely needs a long run must set its own MaxTurns.
Prompt composition: TaskPromptComposer.Compose injects attachment absolute paths as a
read-only "## Reference files" section.
Component Responsibilities
Queue/
QueueService— main dispatch loop; slot limit; decides when to start tasks.QueuePicker— atomic Queued→Running claim via raw SQL.QueueWaker— semaphore for non-blocking, idempotent wake signals.OverrideSlotService— owns the RunNow / ContinueTask slot (bypasses the queue).
Runner/
TaskRunner— orchestrates the run (prepare, execute, handle result).WorktreeManager— creates/manages git worktrees; self-heals stale branches.ClaudeProcess— spawns the claude CLI subprocess; manages streams/logs.TaskRunMcpService— runtime MCP tools (AskUser, SuggestImprovement).TaskRunTokenRegistry— per-run MCP identity for tool-access control.InteractiveLaunchSpecService— config for the task's claude run.
State/
TaskStateService— all task status transitions; guards preconditions; signals queue/hub.
Lifecycle/ (startup recovery)
StaleTaskRecovery— tasks stuck Running after a crash/restart → Failed. The underlyingTaskStateService.RecoverStaleRunningAsyncbulk-flips Running→Failed, then re-runs the same chain/parent-advance side effects as a normalFailAsync(per recovered id, best-effort) so a crash mid-chain-child or mid-improvement-child doesn't leave a successor blocked forever or aWaitingForChildrenparent wedged.OrphanRecovery— dequeues children whose parent is no longer planning (stays attached).AttachmentOrphanRecovery— cleans orphaned attachment files.TaskResetService— manual reset to Idle.TaskMergeService— conflict resolution for worktree merges.
Hub/
HubBroadcaster— single SignalR broadcast point (TaskStarted/TaskUpdated/TaskMessage/WorktreeUpdated…).WorkerHub— SignalR hub + client methods.
Agents/
AgentFileService— file I/O for custom agents.DefaultAgentSeeder— seeds built-in agents on startup.
Worktrees/
WorktreeMaintenanceService— cleanup, state tracking, overview reporting.
Entry Points & Call Chain
Program.cs (DI setup)
├─ QueueService (BackgroundService) → ExecuteAsync loop
│ ├─ waits: IQueueWaker.WaitAsync() or timer
│ ├─ claims: IQueuePicker.ClaimNextAsync()
│ └─ runs: TaskRunner.RunAsync() / ContinueAsync()
├─ Hub clients → WorkerHub methods
│ ├─ Enqueue → ITaskStateService.EnqueueAsync() → Wake()
│ ├─ RunNow → OverrideSlotService.RunNow() → TaskRunner.RunAsync()
│ ├─ ContinueTask→ OverrideSlotService.ContinueTask()→ TaskRunner.ContinueAsync()
│ └─ CancelTask → QueueService.CancelTask()
├─ Lifecycle recovery (startup): StaleTaskRecovery / OrphanRecovery / AttachmentOrphanRecovery
└─ State transitions → HubBroadcaster.TaskUpdated()
Invariants & Conventions
- Atomic claiming — QueuePicker's
UPDATE ... RETURNINGmakes Queued→Running atomic. - Slot limit — respects MaxParallelExecutions; a backstop timer wakes even if a Wake() is missed.
- Pre-claimed tasks — the dispatcher pre-claims via the picker; the override slot (RunNow/ContinueTask) must call StartRunningAsync if a task is not pre-claimed.
- Claim before create —
TaskRunner.RunAsync's unclaimed path callsStartRunningAsyncbeforePrepareRunDirectoryAsync. RunNow racing the picker for the same Queued row used to create the worktree first and only claim afterwards, so the losing dispatch could hit WorktreeManager's branch-collision self-heal and force-remove the winner's live worktree mid-run.OverrideSlotService.RunNowalso fast-rejects a task already Running in the DB (defense in depth; the picker's atomic SQL claim is the real arbiter either way).RunCancellationRegistry.Registerrefuses (and logs) a second registration for the same task id instead of silently overwriting the first, so a losing dispatch's cleanup can't unregister the winner's CTS out from under it. - Terminal writes — use
CancellationToken.None; a task is never left Running after crash/cancel. - Per-run MCP tokens — each run gets a unique token scoping tool access; unregistered on end.
- Auto-retry — one automatic retry if a session exists and the first run failed.
- Worktree self-heal — on branch collision, remove phantom worktrees, prune, delete branch, retry add.
- Review feedback — stored on the task; consumed once a run reaches a terminal state; a re-queued task resumes the session or folds feedback into the prompt.
- Child tasks — planning creates draft children; finalization requires no Queued children remain; OrphanRecovery dequeues children if the parent is not planning.
- Lifecycle recovery runs at startup: stale-Running → Failed; orphaned children → dequeued but attached.