π€ feat: stream engine becomes the AppFiberScope occupant β dispose() aborts and awaits in-flight streams (Wave 4 PR 1) - #4070
Conversation
This comment has been minimized.
This comment has been minimized.
Dogfooding evidence (PR 1 notes β long form)Setup: Branch, two flowing streams in flight (run 2)Branch, one flowing stream (run 1)(The per-stream Baseline
|
|
@codex review |
|
@codex security review |
|
Codex Review: Didn't find any major issues. More of your lovely PRs please. Reviewed commit: βΉοΈ About Codex in GitHubYour team has set up Codex to review pull requests in this repo. Reviews are triggered when you
If Codex has suggestions, it will comment; otherwise it will react with π. Codex can also answer questions or update the PR. Try commenting "@codex address that feedback". |
This comment has been minimized.
This comment has been minimized.
This comment has been minimized.
This comment has been minimized.
|
CI status note: both Codex loops are clean on |
β¦tics for the engine supervisor
β¦tch cancels; guard completed streams from abort bookkeeping
β¦once the final message is committed
β¦ettles every stream exactly once
β¦efore the explicit teardown
β¦per supervised stream
728df0d to
b1d2a55
Compare
|
Rebased onto @codex review |
|
@codex security review |
This comment has been minimized.
This comment has been minimized.
There was a problem hiding this comment.
π‘ Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: b1d2a55767
βΉοΈ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with π.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
β¦ treat post-abort iterator rejections as cancellation
|
@codex security review |
|
Codex Review: Didn't find any major issues. Delightful! Reviewed commit: βΉοΈ About Codex in GitHubYour team has set up Codex to review pull requests in this repo. Reviews are triggered when you
If Codex has suggestions, it will comment; otherwise it will react with π. Codex can also answer questions or update the PR. Try commenting "@codex address that feedback". |
This comment has been minimized.
This comment has been minimized.
β¦r-step timeouts (Wave 4 PR 3) (coder#4076) ## Summary Wave 4 PR 3 (plan Β§2 D6 / Β§3 "PR 3"): `ServiceContainer.initializeCore()` is now a Promise facade over a **startup effect run on the app runtime** β one `Effect.gen` over an ordered table of the five hard startup steps, each `Effect.tryPromise` (identity catch) bounded by `Effect.timeoutOrElse(STARTUP_STEP_TIMEOUT_MS)` on the runtime's `Clock`. A hung step no longer pins the splash screen / listener bind forever: after 60 s it fails startup with a `StartupStepTimeoutError` (`"<step> exceeded <ms> ms"`, `step`/`timeoutMs` fields) through the roots' **unchanged** failure paths. Step errors keep their identity; later steps do not run; the five step names/order and the `[startup] ServiceContainer.initialize completed { stepDurationsMs }` payload are unchanged. Every process root now runs the bounded `dispose()` before exiting on a rejected startup (D6 abandon-and-quit safety): `cli/server.ts` gains it, and the ACP root's existing catch **did not** dispose after a rejected `initialize()` (`if (initialized)` guard β the plan's claim that it did was wrong), so it does now. ## Background - Plan: `~/.xum/plans/mux/effect-wave4-lifecycle-core-plan.md` (collapsed below); PR 1 = coder#4070; the plan predates coder#4058, which split `initialize()` into `initializeCore()` (hard steps, gate the server listener) + `runStartupHousekeeping()` (best-effort, abort-signal cancellable). This PR targets `initializeCore()` only. - `runStartupHousekeeping()` is deliberately untouched: its steps are already cancellable through the dispose abort signal and non-fatal by policy, so per-step timeouts there would be a policy change the plan forbids (recorded in the `appRuntime.ts` startup contract). ## Implementation - `serviceContainer.ts`: `startupCoreSteps: readonly StartupStep[]` (names asserted unique in the constructor β they key `stepDurationsMs`); `initializeCore()` = `assert(not disposed)` + `runtime.managed.runPromise(startupCoreEffect())`; `timedStartupStep` = `Effect.suspend` β `tryPromise({ try: async () => run(), catch: identity })` β `timeoutOrElse` β `Effect.ensuring` (duration recorded when the *wait* ends β settled, failed, or abandoned β never by the abandoned step later). `timeoutOrElse` rather than `timeout` + `catchTag` because the error channel is `unknown`. No `forkDetach`: the zero-arity thunk gets no AbortSignal, so the step promise keeps running; its late settlement is a no-op on the exited fiber and `tryPromise` keeps a rejection handler attached (pinned by test). - `StartupStepTimeoutError extends Error` with `name` set in the constructor, so the desktop `Startup Failed` dialog and `Failed to initialize server:` print class + step without formatting changes. - `src/constants/terminationTimeouts.ts`: `STARTUP_STEP_TIMEOUT_MS = 60 s` (evidence below) and `SERVICE_TEARDOWN_BUDGET_MS = 5 s` (the outer teardown budget `cli/server.ts` already used as a literal, now shared with the ACP root). - `cli/server.ts`: `main().catch` runs `dispose()` bounded by the same 5 s budget as the SIGTERM cleanup when the container was constructed, then exits 1. `acp/serverConnection.ts`: the catch disposes unconditionally (bounded). `desktop/main.ts` unchanged β the catch at `:1249β1265` quits and the `before-quit` listener (`:1271β1305`) races `dispose()` against 5 s. - `appRuntime.ts`: new "Startup" contract section; "Deliberately not done" no longer lists `initialize()`. - Decisions pinned by tests: re-entry is **not** guarded (parity with the promise chain); `initializeCore()` after `dispose()` fails fast with an assertion instead of the disposed runtime's bare `"ManagedRuntime disposed"` string defect. ## Validation - `serviceContainer.test.ts` (+5): TestClock timeout β rejects with `StartupStepTimeoutError` naming the step **exactly** at `TestClock.adjust(STARTUP_STEP_TIMEOUT_MS)` (still pending at β1 ms), later steps not called, the abandoned step's late rejection is not unhandled and changes nothing on the container; error identity (`toBe`) for a rejecting and for a synchronously throwing step with later steps skipped; happy path records the five keys in order and a second call re-runs the steps; after `dispose()` fails fast without running a step. **Red check** against the old promise chain: the timeout test (hangs β 5 s test timeout) and the after-dispose test fail; the identity/sync/re-run tests are parity pins (green on both, by design). - Gates: `serviceContainer.test.ts` (23), `di/*.test.ts` + `coreServicesRoot.test.ts` + `acp/*.test.ts` (55), `TEST_INTEGRATION=1 bun x jest tests/ipc/{doubleRegister,windowTitle,savedQueries,acp.sessionMethods,acp.disconnectCleanup}` (32), `bun test src/cli/` (185), `make static-check`. - Dogfooding (first comment, with screenshots): 3 + 3 `xum server` cold starts on `main` vs branch β same ten `stepDurationsMs` keys in the same order, slowest core step `taskService.recoverInterruptedTasks` 49β61 ms; a throwaway build with the constant forced to 1 ms on `xum server` and `xum acp` shows `Failed to initialize server: StartupStepTimeoutError: extensionMetadata.initialize exceeded 1 ms`, the full bounded `[shutdown]` dispose sequence (28 ms / 36 ms), exit code 1. - `cli/server.ts` has no unit seam for `main().catch` (`server.test.ts` exercises `createOrpcServer`); the throwaway transcript is the evidence for that path. ## Risks - False timeout on a pathologically slow host turns a slow-but-fine start into a crash β bounded at 60 s, β₯ 1000Γ the measured maximum and above the policy service's own 10 s fetch timeout; the only potentially unbounded step, `taskService.recoverInterruptedTasks`, scales with active agent tasks, not deployment size. - An abandoned step (e.g. mid-`editConfig`) keeps running until the process exits; the bounded `dispose()` applies the same latches as a quit (`beginShutdown`s) and config writes are lock/journal protected β identical exposure to a SIGTERM mid-startup today. - Disposing the runtime concurrently with an in-flight `initializeCore()` does not interrupt the startup fiber (root fibers are not scope children) β same as the promise chain it replaces; documented, not changed. OFF-RAMP did not fire: error identity, `tests/ipc` and the ACP path are preserved. --- <details> <summary>π Implementation Plan (Wave 4 β effect-wave4-lifecycle-core-plan.md)</summary> # Effect migration β Wave 4: finish the concurrency/lifecycle core Bounded wave: **4 PRs (PR 4 optional), explicit STOP criterion, explicit OFF-RAMPs.** Plan only; nothing here is implemented. > **Review status:** Independently reviewed (adversarial Reviewer sub-agent, advisor unavailable): APPROVE WITH REQUIRED EDITS β both edits applied; verified claims: Effect.promise 0-arity thunk allocates no AbortController (internal/effect.js:741β776); closed-scope forkIn+startImmediately runs onInterrupt (2237β2274, 391β409); forkIn observer removes the scope finalizer on exit (2270β2271); Effect.timeoutOrElse exists (Effect.d.ts:7833); 0 line drift at b87f627; 11 settleWorkspaceTurn callers confirmed; 'aborted' β NON_RETRYABLE_STREAM_ERRORS. Line references are to `main` @ `b87f62729`. ## 0. Thesis check (coordinator's judgment vs. evidence) **Thesis:** Effect's payoff in this app is structured concurrency + interruption-safe lifecycles in the orchestration core (still Promise + AbortController). **Verdict: holds for the stream engine; only half-holds for turn handles.** - Stream engine β **holds.** `ServiceContainer.dispose()` never stops or awaits in-flight streams (`serviceContainer.ts:478β537` has no `streamManager` step); an in-flight stream dies with the process and is recovered on next load from `partial.json` (β€ 500 ms stale, `PARTIAL_WRITE_THROTTLE_MS`, `streamManager.ts:776`). `AppFiberScope` exists precisely for this and has no occupant. A supervised per-stream fiber is the right tool. - Turn handles β **half-holds.** The 7Γ "superseded by an uncorrelated workspace stream-end" false-settle is a *correlation-predicate* bug (`interruptWorkspaceTurnFromUncorrelatedStreamEnd`, `workspaceTurnManager.ts:4220β4307`: any uncorrelated stream-end after the prompt index settles the handle `interrupted`), not a Promise-vs-fiber structure bug. Turn handles are **persisted records** (`taskHandleStore.upsertWorkspaceTurn`) spanning **multiple streams** (tool-call continuations are deferred via `hasSameTurnContinuation`, `:4491`) and surviving restarts; a fiber/Deferred can only model the in-process waiter and would not fix correlation. Open PR **coder#3949** fixes the predicate in Promise idiom and is Codex-green. Wave 4's turn-handle PR therefore becomes **"codify the settlement invariant + prove the class is gone"**, not "fiberize handles" (D5 below). Corrected baseline numbers (measured this workspace): 46/470 `src/node` non-test files import `effect` (coordinator said 35); 113 direct `Effect.run*` sites outside `di/` in 16 files; 226 `Effect.gen`; 9 `TaggedError` classes; effect `4.0.0-rc.112`, `@orpc/*` `1.14.11`; **effect v4 is not GA** (rc line still current). ## 1. Verified current state (evidence the design rests on) <details> <summary>Stream engine (streamManager.ts)</summary> - `startStream` (`:4723β4901`): per-workspace mutex β `new AbortController()` + `linkAbortSignal` (`:4771β4772`) β `resourceScope = Scope.makeUnsafe()` (`:4777`) β temp-dir `Effect.acquireRelease` (`:4802β4824`) β `createStreamAtomically` β `streamText` (`:2244`, `abortSignal: abortController.signal` `:2250`) β registered in `workspaceStreams` (`:2463`) β **`streamInfo.processingPromise = this.processStreamWithCleanup(...)` fire-and-forget (`:4876β4882`)** β returns `Ok({ messageId, completion })`. - `processStreamWithCleanup` (`:3331β4089`, plain async): `while(true)` retry loop; `for await (part of fullStream)` (`:3358β3837`) with abort check at loop head (`:3361`); post-loop `if (!signal.aborted)` gate (`:3849`) β completion path (`deletePartial` `:3981`, `updateHistory` `:3989`, `recordSessionUsage` `:4001`, `state = COMPLETED` `:4017`, emit `stream-end` `:4023`, `terminalCompletion` `:4024`); error path β `handleStreamFailure` (`:4094β4112`) β `persistStreamError` writes error partial; `finally` (`:4052β4088`): release MCP lease, `Effect.runFork(Scope.close(resourceScope))` (`:4064β4066`), unlink abort, `workspaceStreams.delete`, `eventSpine.emit("stream.end")`, `completionController.settle`. - Cancellation: `stopStream` (`:5043β5111`) β `cancelStreamSafely` (`:1766β1800`): `if (state === COMPLETED) { await processingPromise; return }` β `state = STOPPING` β `flushPartialWrite` β `abortController.abort()` β `cleanupAbortedStream` (`:1828β1951`): `await processingPromise` β usage β `writePartial` (`:1876β1910`) β `emitStreamAbort` β `settle({status:"aborted"})`. **No completed-guard after the await** (verified `:1838β1951`): a cancel landing between `:3849` and `:4017` re-writes `partial.json` after `deletePartial` and emits `stream-abort` after `stream-end` (pre-existing window; dispose() will widen its exposure). `cancelStreamSafely` is also not idempotent for concurrent callers (only `COMPLETED` is checked). - AIService on `stream-abort` (`aiService.ts:355β377`): `abandonPartial ? deletePartial : commitPartial β deletePartial` (fire-and-forget listener). - Crash recovery: `HistoryService.commitPartial` (`historyService.ts:1963β2061`) β strips error metadata, `hasCommitWorthyParts`, stale-epoch check, update-or-append by `historySequence`, delete partial; invoked from `agentSession.init` (`:5002`), `aiService.streamMessage` (`:886`), stream-abort (`:364`), `duplicateWorkspace`. - `StreamAbortReason = "user" | "startup" | "system"` (`src/common/orpc/schemas/stream.ts:295`). - Pinned seams: chaos test `Reflect.set(streamManager, "tokenTracker" | "createStreamResult")` (`streamManager.chaos.test.ts:130β134, 238β242`); `streamManager.test.ts` `Reflect.set` on `processStreamWithCleanup` (`:2787`), `createStreamAtomically` (`:2783`), `createTempDirForStream`, `cleanupStreamTempDir`, `Reflect.get` on `workspaceStreams`, `schedulePartialWrite`, β¦; `modelOnlyNotifications.test.ts` calls `processStreamWithCleanup` directly (`:93, :187`); `aiService.test.ts` spies `startStream`, `generateStreamToken`, `createTempDirForStream`, `isResponseIdLost`. Constructor: `(historyService, sessionUsageService?, getProvidersConfig?, eventSink = noop, runner = defaultEffectRunner)` (`:801β813`); `effectRunner` used at `:1152, :1154, :1169` only. - Every stream event carries `workspaceId` + `messageId`; `stream-end`/`stream-abort`/`error` carry `metadata.muxMetadata` when the prompt had it. </details> <details> <summary>DI / shutdown / startup</summary> - `AppFiberScopeLive` is in `CoreLive`'s `runtimeSeams` (`di/layers/core.ts:644`), so both roots have it; `StreamManagerLive` (`core.ts:226β238`, stage S2b) already yields `EffectRunnerTag`. CLI cleanup lists include `appFiberScope.close` (`cli/run.ts:1579`, `cli/workflow.ts:286`). - Bounds: `APP_FIBER_SCOPE_CLOSE_TIMEOUT_MS = 2000`, `APP_RUNTIME_DISPOSE_TIMEOUT_MS = 2000`; outer budgets are **5000 ms** on both desktop (`desktop/main.ts:1297` `Promise.race` vs `setTimeout(5000)`) and `xum server` (`cli/server.ts:236β243` force-exit). **The scope bound cannot grow** without changing outer budgets. - rc.112 semantics verified in `node_modules/effect/dist/internal/effect.js:2264`: `forkIn` registers a scope finalizer and **removes it when the fiber completes** (no leak), and **interrupts immediately if the scope is already closed** (streams starting mid-shutdown fail closed). `Effect.promise(evaluate: (signal) => PromiseLike)`, `Effect.onInterrupt`, `Effect.forkIn(_, scope, { startImmediately? })`, `Stream.toAsyncIterableWith(context)`, `Stream.provideContext` all exist. - `ServiceContainer.initialize()` (`serviceContainer.ts:297β362`): six awaited `initialize()`s wrapped in `recordStep` (durations only, no catch, no timeout) + three sync `start()`s + two fire-and-forget sweeps. Failure handling: desktop `Startup Failed` dialog + `app.quit()` (`desktop/main.ts:1249β1265`); `cli/server.ts:136` uncontained; ACP `serverConnection.ts:205β216` dispose + rethrow; `tests/ipc/setup.ts:85` no catch. No outer timeout anywhere. - `streamBridge.subscriptionIterable` (`orpc/streamBridge.ts:176`) β `Stream.toAsyncIterable(...)` on the global runtime, 19 call sites in `routerSubscriptions.ts`; heartbeat via `Effect.sleep` in `forkScoped` (`:145β152`). `streamBridge.test.ts` has 11 real-time waits, but only **3 are clock-bound** (`:207` 1 ms initial delay, `:241` heartbeat 10 ms, `:255` 10 ms laziness); 8 are `waitFor(listenerCountβ¦)` readiness polls that TestClock cannot replace. </details> <details> <summary>Turn handles + open PRs</summary> - Handle record `{ handleId "wst_β¦", ownerWorkspaceId, workspaceId, turnId, messageId, status, attentionPolicy, disposableWorkspace }`; prompt carries `muxMetadata: { type:"workspace-turn-task", taskHandleId, ownerWorkspaceId, turnId }` (`workspaceTurnManager.ts:1420`). `TaskService` forwards `aiService` `stream-end`/`stream-abort`/`error` to `finalizeWorkspaceTurnFromStreamEnd` (`:4442β4544`): correlated branch matches `record.workspaceId && record.turnId` (`:4472`); **uncorrelated branch** (`metadata == null`, not `agentId === "compact"`) β `interruptWorkspaceTurnFromUncorrelatedStreamEnd` β settles `interrupted` whenever `streamEndIndex >= promptIndex` (`:4293β4305`). Producers of such uncorrelated ends: bash-monitor wake continuations, child terminal-attention deliveries, heartbeat, peer messages, parent auto-resume. - Cascade: disposable child β `cleanupDisposableWorkspaceTurn` β `workspaceService.remove(β¦, true)` kills its background processes; persistent child β parent sees `interrupted` β `task_stop` β `backgroundProcessManager.stopMonitor(β¦, "canceled")`. This is the observed "monitors died afterwards". - Settlement chokepoint: `settleWorkspaceTurn(params)` (`:2085`), **11 callers**, guarded by `workspaceTurnSettlementLocks.withLock(handleId)`; waiters in `pendingWorkspaceTurnWaitersByHandleId` with `setTimeout` timeouts (`:2425β2496`). - **coder#3949** "preserve turns across synthetic wake ends" (coadler): rewrites the uncorrelated branch β walks history from the turn anchor to the stream-end and settles *only if a manual child input intervened* (`isManualChildWorkspaceInput`); otherwise ignores the end. Touches `:281β295, :4217β4355` + tests (+286/β31). Codex: "Didn't find any major issues" + clean security on `f9baa2fc9`. `mergeable: MERGEABLE`, but `Test / Unit` and `Codex Comments` red, 19 commits behind main. - **coder#3915** "correlate workspace-turn liveness" (coadler): creation reservations + `getWorkspaceTurnLiveness`/`getWorkspaceTurnRuntimeActivity` (identity-matches the active stream's `muxMetadata` against the record) for staleness/capacity. Touches `:442β486, :1298, :2573, :3761, :3800β4064` (+494/β60). `BLOCKED`, latest Codex review has open findings, `Test / Unit` red, 19 behind. - Together they are the identity-correlated model the coordinator wants: coder#3915 = identity-correlated *liveness*, coder#3949 = identity-gated *settlement*. </details> ## 2. Design decisions **D1 β Fibers WRAP the AbortController; they do not replace it.** The AI SDK is cancelled only via `AbortSignal`; the `for await` loop, soft-interrupt at step boundaries, retry/fallback re-creation of `streamResult`, and ~30 abort touchpoints (coder#4032) all key off the signal. Converting the 750-line loop to `Stream.fromAsyncIterable` + fiber interruption would touch hundreds of `WorkspaceStreamInfo` transitions and break the `processStreamWithCleanup`/`createStreamResult` spy seams. Instead: **the fiber is the ownership/supervision unit; the signal stays the cancellation transport.** The dual-cancellation glue coder#4032 feared is confined to **one** point β the supervisor's `onInterrupt` β which routes through the existing user-stop path (`cancelStreamSafely`), so shutdown β‘ "user pressed stop" semantically (partial flushed with usage, `stream-abort` emitted, `completion` settles `aborted`, AIService commits the partial). **D2 β Supervisor topology: one supervisor fiber per stream in `AppFiberScope`, wrapping the already-started `processingPromise`.** `streamInfo.processingPromise = this.processStreamWithCleanup(...)` stays byte-identical (sync-start preserved; `Reflect.set(processStreamWithCleanup)` seam preserved; `cleanupAbortedStream`'s `await processingPromise` unchanged). Immediately after it: ```ts // startStream, after processingPromise is assigned (unsupervised path unchanged when no scope) this.superviseEngine(typedWorkspaceId, streamInfo); private superviseEngine(workspaceId: WorkspaceId, streamInfo: WorkspaceStreamInfo): void { if (this.engineScope === undefined) return; // direct construction / CLI tests: today's behavior assert(streamInfo.engineFiber === undefined, "engine already supervised"); // Zero-arity thunk on purpose: rc.112 allocates an internal AbortController only // when `evaluate.length !== 0`; the stream's own controller stays the sole signal. const supervisor = Effect.promise(() => streamInfo.processingPromise).pipe( Effect.onInterrupt(() => Effect.uninterruptible( // explicit, per house doctrine (finalizers are already uninterruptible) Effect.promise(async () => this.cancelStreamSafely(workspaceId, streamInfo, "system")) ) ), Effect.catchDefect((d) => Effect.sync(() => log.warn("[stream] engine supervisor defect", { workspaceId, error: d }))) ); streamInfo.engineFiber = this.effectRunner.runSync( Effect.forkIn(supervisor, this.engineScope, { startImmediately: true }) ); } ``` - `Effect.promise` is interruptible while suspended (`internal/effect.js:741β801`, Async op); `onInterrupt` = `onErrorFilter(causeFilterInterruptors, β¦)` (`:1762`); `forkIn` registers `fiberInterrupt(fiber)` as the scope finalizer (`:2264β2275`), `fiberInterrupt` awaits the fiber (`:635β642`), and parallel `scopeClose` awaits all finalizers via `fiberAwaitAll` (`:1590β1601`) β `closeScopeBounded` at dispose step 2 gives "interrupt **and** await" while `historyService`/`sessionUsage`/`eventSink β AIService β bridge servers` are still alive (bridges stop in step 3, so clients receive `stream-abort`). - Normal completion: fiber exits β `forkIn`'s observer removes the scope finalizer (verified) β no per-stream residue. - Stream started after step 2: `forkIn` on a closed scope calls `fiber.interruptUnsafe` synchronously and returns the fiber (`:2272β2274`, `runSync` does not defect). With `startImmediately: true`, `forkUnsafe` runs `child.evaluate` synchronously (`:2233β2247`) up to the `Effect.promise` Async op (`:772β801`), so the fiber is suspended (`_running=false`) when the interrupt lands and `interruptUnsafe` (`:391β409`) unwinds the stack through the `onInterrupt` handler β the stream is aborted as `system` (fail-closed during shutdown). Verified in rc.112 internals (`effect.js:2233β2247, 391β409`); **pin with a test** ("stream started after scope close is aborted") so an RC bump cannot silently change it. - `"system"` is semantically exact: `"user"`/`"startup"` suppress next-startup recovery (`retryEligibility.ts:114β118, 284β287`), `"system"` marks an involuntary backend interruption (as `taskService.ts:8100, 8223` use it). **No in-session retry loop is possible:** the `stream-abort` handler (`agentSession.ts:6010`) routes `{ type: "aborted" }` to `retryManager.handleStreamFailure`, and `"aborted"` is in `NON_RETRYABLE_STREAM_ERRORS` (`retryEligibility.ts:49β59, 106`) β `retryManager.ts:99β104` abandons immediately, never schedules a fiber. Dogfooding still checks the *restart* UX (the recovered partial is shown as interrupted; note whether any next-startup recovery re-sends β same class as today's `system` aborts from `taskService`). - `engineScope` arrives as an **optional 6th constructor parameter** (`engineScope?: Scope.Closeable`), wired from `AppFiberScopeTag` in `StreamManagerLive` (`core.ts:226`). Default `undefined` keeps every direct-construction test and `aiService.ts:174` path identical (I4). `AppFiberScopeLive` already sits beneath S2b in `runtimeSeams`, so no staging change (I6). - Abort reason: reuse **`"system"`** β no wire/schema change; UI copy for `system` already exists. - Pending-start window (`pendingStreamStarts`, before registration) is **not** supervised: nothing is persisted for it yet, and `stopStream` already aborts pending controllers. Documented, not fixed. **D3 β Fix the two adjacent cancel races in the same PR (closely-related bugs, not deferrals).** (a) `cleanupAbortedStream`: after `await processingPromise`, if `streamInfo.terminalCompletion !== undefined` (completed/failed while the cancel was in flight) β return without abort bookkeeping (prevents `partial.json` resurrection after `deletePartial` and a `stream-abort` after `stream-end`). (b) `cancelStreamSafely` (`:1766`): latch a per-stream `cancelPromise` so concurrent cancellers (user stop racing dispose) join one cleanup β exactly one `stream-abort`, one `settle`. **Zero-suspension requirement:** the latch must be checked and assigned **synchronously at function entry, before any `await`** (the current first await is `flushPartialWrite` at `:1789`) β otherwise racing callers can both enter `cleanupAbortedStream`. Shape: ```ts if (streamInfo.cancelPromise) return streamInfo.cancelPromise; streamInfo.cancelPromise = (async () => { /* existing body, unchanged */ })(); return streamInfo.cancelPromise; ``` Both are β€ 10 LoC and get behavioral tests. **D4 β Shutdown bound stays 2 s; the finalizer must be fast or abandoned.** Outer budgets are 5 s; 2 s + 2 s already consume 4 s. A flowing stream aborts within one chunk; a wedged provider (no chunks, ignores abort) hits the existing `boundedTeardown` timeout: warning, continue, process exit β identical to today's outcome. Dogfooding measures the actual `[shutdown] AppFiberScope closed { ms }` with a live stream. **D5 β Turn handles: codify the settlement invariant; do not fiberize.** Invariant: *a workspace-turn handle settles terminally only by (i) a stream terminal event whose `muxMetadata` correlates `{taskHandleId, ownerWorkspaceId, turnId}` to the record; (ii) an explicit interrupt (`task_stop`/`interruptWorkspaceTurn`); (iii) manual supersession β a manual child input after the turn anchor; (iv) stale-liveness reconciliation.* An uncorrelated stream-end is **never** terminal by itself. coder#3949 makes (iii) the only uncorrelated outcome; coder#3915 implements (iv) by identity. Wave 4 adds a `cause` discriminant to `settleWorkspaceTurn` (the single chokepoint) with a runtime assertion, plus the regression harness. Rationale for not converting waiters to `Deferred`/fibers: no behavioral gain, 4.9k-line file, and the coordinator's "settle only on the owning stream's termination" is over-specified β a turn owns *several* streams. **D6 β Startup: `initialize()` stays a Promise facade over a runtime-run startup effect; timeout β same failure path as a thrown step.** Each step is `Effect.tryPromise({ try: async () => step(), catch: identity }).pipe(Effect.timeoutOrElse({ duration: STARTUP_STEP_TIMEOUT_MS, orElse: () => Effect.fail(new StartupStepTimeoutError(name, ms)) }))` (`timeoutOrElse` exists in rc.112, `Effect.d.ts:7833`; chosen over `timeout` + `catchTag` because the step's error channel is `unknown`, which `catchTag` cannot narrow). No `forkDetach` needed: a Promise step keeps running on its own when the waiting fiber times out (not inside an uninterruptible region, so the timeout interrupts the wait directly). `StartupStepTimeoutError extends Error` with `name = "StartupStepTimeoutError"` set in the constructor and message `"<step> exceeded <ms> ms"` (so the desktop dialog's error formatting shows both the class and the step name) β desktop shows it in the existing `Startup Failed` dialog; CLI/ACP/tests paths unchanged. Step **errors keep their identity** (v4 `runPromise` rejects with the raw failure). Downgrading any step to best-effort is a **policy change, out of scope** (audit of the six implementations: extensionMetadata/telemetry/experiments are local fs, <50 ms; policy has its own 10 s fetch timeout; workspaceService bounds its sync internally; only `taskService.initialize` β config scan + `editConfig` + recovery `sendMessage`s β is potentially unbounded). The three `start()`s stay sync (`Effect.sync`), the two fire-and-forget sweeps stay outside the effect. `stepDurationsMs` is preserved. **Abandon-and-quit safety:** an abandoned `taskService.initialize` may be mid-`editConfig` when the root exits. Parity requirement for PR 3: after a rejected `initialize()`, every root runs the bounded `dispose()` before exiting. Verified: desktop already does β `services` is assigned before the await (`main.ts:653β656`), the catch calls `app.quit()`, and the `before-quit` listener (`:1271β1305`, guard `if (isDisposing || !services) return`) races `services.dispose()` against 5 s; ACP does (`serverConnection.ts:205β216`); **`cli/server.ts` does not** (`:133β136` awaited at top level, `main().catch` at `:282` only logs) β PR 3 adds a bounded `dispose()` there (β€ 10 LoC, same 5 s budget). **D7 β streamBridge: thread the runtime context, not a runner.** `subscriptionIterable` gains `context?: Context.Context<never>` β `Stream.toAsyncIterableWith(context)`; `routerSubscriptions` passes the handler's `"effect/context"`. Production behavior identical; heartbeat sleeps on the runtime `Clock`; tests can run the 3 clock-bound waits on `TestClock`. Honest scope: the 8 readiness polls stay. ## 3. PRs (ordered by value Γ· risk; each independently mergeable) ### PR 1 β StreamManager engine core becomes the first `AppFiberScope` occupant **Value:** high (the only remaining shutdown data-integrity gap; the reason `AppFiberScope` exists). **Risk:** medium β low with D1/D2. **Net product LoC β +55** (`superviseEngine` ~25, ctor param/field ~5, `engineFiber` field ~2, D3 guards ~12, `core.ts` wiring ~2, doc updates in `appRuntime.ts`/`appFiberScope.ts` "occupant" text ~10). Files: `src/node/services/streamManager.ts`, `src/node/services/di/layers/core.ts`, `src/node/services/di/appRuntime.ts` + `appFiberScope.ts` (docs), tests below. Pre-work (before writing product code; each yields a note in the PR body): 1. Confirm `processStreamWithCleanup` never rejects (try/catch/finally shape `:3331β4089`); else the supervisor must fold rejections (it already `catchDefect`s). 2. Confirm `Effect.promise` interruption + `onInterrupt` await ordering under `Scope.close` in a 20-line probe test (pattern of `appFiberScope.test.ts:27β47`). 3. Enumerate abort observers that run *after* the finalizer resolves (AIService `stream-abort` listener β `commitPartial`; agentSession completion continuations) and confirm the durable order (`writePartial` β commit β `deletePartial`) makes a mid-flight `process.exit` recoverable on next load (it is: partial survives until commit completes). 4. Measure: `[shutdown] AppFiberScope closed { ms }` with a live stream in the sandbox (D4). 5. Pin (same probe test): a fiber forked with `startImmediately: true` into an already-closed scope still runs its `onInterrupt` finalizer (reviewer-verified in rc.112 internals; the test guards RC bumps). Acceptance (behavioral tests only): - `streamManager.test.ts` (new cases; existing cases untouched): with `engineScope = Scope.makeUnsafe("parallel")` and a fake `createStreamResult` whose `fullStream` yields one `text-delta` then blocks until its `AbortSignal` fires β `closeScopeBounded(engineScope)` resolves; `writePartial` was called with the streamed text; exactly one `stream-abort` (`abortReason: "system"`) and zero `stream-end`; `completion` settles `{status:"aborted"}`; `workspaceStreams` is empty. - Wedged provider (fullStream never yields, ignores abort): `closeScopeBounded` resolves within the bound, never rejects, warns once (assert the returned promise resolves and no throw; do **not** assert log text). - No-scope construction: identical event sequence to today (guards existing suites; no new assertions needed beyond the unchanged suites passing). - D3(a): cancel issued after the loop exits but before `COMPLETED` β history has exactly one final message, `partial.json` absent, event order `stream-end` only. - D3(b): `stopStream` + `closeScopeBounded` racing on one stream β exactly one `stream-abort`, one settle. - Fiber residue: after 50 completed streams, `closeScopeBounded(engineScope)` emits zero `stream-abort` and completes in the same tick class as an empty scope (assert no aborts and `workspaceStreams.size === 0`). - `streamManager.chaos.test.ts` β existing cases byte-identical; **one new fuzz variant** constructs with an engine scope and closes it at a random iteration: every stream settles **exactly once** (count terminal events per `messageId` β€ 1, all `completion` promises settle). - `serviceContainer.test.ts`: "dispose() aborts and awaits an in-flight stream before `desktopBridgeServer.stop()`" (extend the ordering harness at `:295β323`); `coreServicesRoot.test.ts`: `xum run` cleanup list does the same via `appFiberScope.close`. Gate suites: `streamManager.test.ts`, `streamManager.chaos.test.ts`, `streamManager.modelOnlyNotifications.test.ts`, `aiService.test.ts`, `agentSession.disposeRace.test.ts`, `agentSession.sinceReplayContract.test.ts`, `serviceContainer.test.ts`, `coreServicesRoot.test.ts`, `di/*.test.ts`, `taskService.test.ts`, `workspaceService.test.ts`, `turnRequestBuilder.test.ts`; `make static-check`. House pre-review audits: interruption posture (supervisor's only suspension is the promise; finalizer uninterruptible end-to-end incl. `cancelStreamSafely` β `cleanupAbortedStream`); no defect escapes (`catchDefect` on the supervisor; `Effect.promise` thunks `async`); spy-seam check (`processStreamWithCleanup`, `createStreamResult`, `createStreamAtomically`, `startStream` signatures unchanged; constructor arity unchanged, trailing optional); sync-start (`processingPromise` assigned before fork; `runSync(forkIn)` completes synchronously); no constructor side-effects added; zero-suspension check on D3(b) latch (`cancelPromise` checked-and-assigned synchronously at `cancelStreamSafely` entry, before the first `await` at `:1789`; review the diff for any inserted `await`/lookup ahead of the assignment). Rollback: revert the `core.ts` wiring line β `engineScope` undefined β today's behavior; D3 guards can stay (independent bug fixes). ### PR 2 β Turn-settlement invariant + false-settle regression harness (gated on coder#3949) **Value:** high (7Γ production race). **Risk:** low. **Net product LoC β +40** (`WorkspaceTurnSettlementCause` union + `cause` on `settleWorkspaceTurn` params + assert ~10; 11 call sites Γ 1β3 lines). Relationship to open PRs β explicit: - **coder#3949 is the fix and a hard prerequisite.** PR 2 rebases on it, changes none of its logic (`interruptWorkspaceTurnFromUncorrelatedStreamEnd`, `isWorkspaceTurnAnchorForRecord`, `isManualChildWorkspaceInput`), and adds the invariant + proof on top. If coder#3949 has not merged when PRs 1/3 are done: **do not fork a competing fix**; report to the coordinator, offer the regression test file to coder#3949's author as a review artifact, and hold PR 2 (it is not on any other PR's critical path). - **coder#3915 is a soft prerequisite.** Its diff (`:3800β4064` incl. `settleStaleWorkspaceTurn`, a `settleWorkspaceTurn` caller) overlaps PR 2's one-line-per-caller change. Prefer landing after it; if PR 2 must go first, the conflict is a one-line `cause:` addition per caller. PR 2 never edits liveness/reservation code. - Wave 4 does not otherwise touch `workspaceTurnManager.ts`. Design: `type WorkspaceTurnSettlementCause` enumerated from the 11 callers (audited at `main` @ `b87f62729`): `:1373` creation validation failure; `:1476` pre-stream interrupt during launch; `:1500`/`:1518` pre-stream send failure; `:3834`/`:3863` stale-liveness recovery / restart timeout (`settleStaleWorkspaceTurn` β coder#3915's region); `:4301` **uncorrelated-stream-end manual supersession** (the only uncorrelated settle in the codebase; the path coder#3949 rewrites); `:4529` correlated terminal; `:4571` stream-abort; `:4675` deferred stream error; `:4736` terminal stream error. `settleWorkspaceTurn` asserts `params.cause` is a member and, for `manual-supersession`, that the superseding input's `messageId` is supplied β turning D5 into an exhaustive `Record<Cause, β¦>` check rather than prose, so a future "settle on uncorrelated end" cannot be added without naming (and justifying) a cause. Acceptance: - New `workspaceTurnManager.uncorrelatedStreamEnd.test.ts` (real `WorkspaceTurnManager` + `TaskHandleStore` + fake `aiService` emitter, following the existing suite's harness): (1) create turn β correlated `stream-start` β **synthetic wake stream** on the same child ends uncorrelated after the anchor β handle stays `running`, waiter unresolved, no disposable cleanup, no terminal attention β correlated `stream-end` β `completed`. (2) same with `finishReason:"tool-calls"` continuation in between. (3) manual child input between anchor and end β `interrupted` with `cause: manual-supersession`. (4) explicit `interruptWorkspaceTurn` β `interrupted`, `cause: explicit-interrupt`. Case (1) is the **scripted reproduction**: it must **fail on the pre-coder#3949 merge-base** (run the file from a sibling worktree at `git merge-base origin/main <coder#3949 head>`; record the failing assertion in the PR body) and pass after. - Existing 113 `workspaceTurnManager.test.ts` cases and `taskService.test.ts` turn cases unchanged. Gate suites: `workspaceTurnManager.test.ts`, `taskService.test.ts`, `taskHandleStore.test.ts`, `tools/task*.test.ts`; `make static-check`. Audits: spy-seam (`getWorkspaceTurn`, `listAllWorkspaceTurns`, `enqueueTerminalAttention`, `deliverPersistentChildWorkspaceTurnResult` untouched); settlement lock held across the assert; no new suspension inside `withLock`. Rollback: revert; the test file stays valid against coder#3949 alone (drop the `cause` assertions). ### PR 3 β `ServiceContainer.initialize()` as a runtime-run startup effect with per-step timeouts **Value:** medium (a hung `taskService.initialize()` currently pins the splash screen forever; deterministic TestClock tests of startup). **Risk:** lowβmedium. **Net product LoC β +80** (step table ~20, timed-step helper ~15, `StartupStepTimeoutError` ~8, constant ~3, facade ~10, root dispose-on-failure parity β€ 10, doc update ~10). Files: `serviceContainer.ts`, `src/constants/terminationTimeouts.ts` (keep with the termination constants so the budget doc stays in one place), `cli/server.ts` (dispose in the startup catch if missing), `di/appRuntime.ts` doc ("Deliberately not done" β remove the initialize() line; add startup contract). Design (D6): `initialize(): Promise<void>` β `this.runtime.managed.runPromise(this.startupEffect())`. `startupEffect = Effect.gen` over an ordered `readonly steps: ReadonlyArray<{ name, run: () => Promise<void> }>` (assert names unique); each step: `recordStep` timing kept, `Effect.tryPromise({ try: async () => run(), catch: identity }).pipe(Effect.timeoutOrElse({ duration: STARTUP_STEP_TIMEOUT_MS, orElse: () => Effect.fail(new StartupStepTimeoutError(name, STARTUP_STEP_TIMEOUT_MS)) }))`. Then `Effect.sync` for the three `start()`s; the sweeps remain after `runPromise`. Constant `STARTUP_STEP_TIMEOUT_MS` β pre-work measures `[startup] <step> { ms }` across sandbox cold starts and picks β₯ 10Γ the slowest observed (propose 60 s; must be generous β a false timeout turns a slow-but-fine start into a crash). Roots dispose after a rejected `initialize()` (D6 abandon-and-quit safety). Acceptance (all in `serviceContainer.test.ts`, TestClock via the existing `AppLive` spy at `:355β395`): - A step that never resolves β `initialize()` rejects with `StartupStepTimeoutError` naming the step after exactly `TestClock.adjust(STARTUP_STEP_TIMEOUT_MS)`; later steps did not run. - A rejecting step β `initialize()` rejects with **the same error object** (identity), later steps did not run (parity with today). - Happy path β `stepDurationsMs` has all six keys; `start()`s called once each; second `initialize()` call behavior unchanged from today (verify whether re-entry is guarded today; preserve). - `tests/ipc` harness and ACP entry still pass unchanged. Gate suites: `serviceContainer.test.ts`, `coreServicesRoot.test.ts`, `di/*.test.ts`, `src/node/acp/*.test.ts`, `TEST_INTEGRATION=1 bun x jest tests/ipc` (smoke subset); `make static-check`. Audits: I1 untouched (no layer body changes); I2 (only the composition root touches the runtime); error identity preserved (no wrapping); abandoned-step safety β every root runs bounded `dispose()` after a rejected `initialize()` (D6; `cli/server.ts` gains it in this PR); no sweep moved into the effect; the six-step order and `[startup] <step>` names unchanged. Rollback: revert; constant removal. ### PR 4 (optional β cut if budget is exhausted) β `streamBridge` on the runtime context **Value:** low (closes the last documented "global runtime" exception; enables TestClock for the heartbeat). **Risk:** low. **Net product LoC β +30** (`context?` option + `toAsyncIterableWith` ~8; 19 call sites Γ 1 line via one shared helper in `routerSubscriptions.ts` ~3). Acceptance: the 3 clock-bound waits (`:207, :241, :255`) run on `TestClock`; the heartbeat test asserts N heartbeats after `TestClock.adjust(N Γ interval)` with zero real time; existing behavioral assertions unchanged; `tests/ipc` subscription tests pass. Do not rewrite the 8 readiness polls. Gate suites: `streamBridge.test.ts`, `routerSubscriptions*.test.ts`, `orpc/*.test.ts`, `TEST_INTEGRATION=1 bun x jest tests/ipc` (subscription subset); `make static-check`. Audits: `Stream.toAsyncIterableWith` preserves double-close safety (pin with the existing test); `Cause.Done` typing unchanged; no `Scope`/`MemoMap`/`Scheduler` captured (pass the oRPC `effect/context`, which the DI layer already strips per `EffectRunnerLive`); `context` stays optional so direct callers/tests without a runtime keep today's global-runtime path. Rollback: revert; the optional `context` default (`Context.empty()`) is exactly today's `toAsyncIterable`, so a partial revert of call sites is also safe. ### Execution order and size Net product LoC for the wave β **+205** (PR 1 β +55, PR 2 β +40, PR 3 β +80, PR 4 β +30); tests β +600β800. PR 1 starts immediately. PR 2 starts the moment coder#3949 merges (parallel with PR 1/3 β disjoint files). PR 3 after PR 1 merges (both touch `appRuntime.ts` docs; PR 3 also touches `serviceContainer.ts`). PR 4 last, only if PRs 1β3 landed and no OFF-RAMP fired. Each PR: Codex dual review, `Codex Comments` minimization, merge queue; commit WIP early (`/tmp` wipes). ## 4. STOP criterion (measurable) and OFF-RAMPs Wave 4 is **done** β and the Effect migration line **stops** without a new RFC β when all hold: 1. **dispose() awaits in-flight streams:** `serviceContainer.test.ts` ordering test + `coreServicesRoot.test.ts` pass on main; a sandbox `script -f` transcript of `xum server` receiving SIGTERM mid-stream shows `stream-abort` β `[shutdown] AppFiberScope closed { ms }` β `[shutdown] desktopBridgeServer.stop`, and immediately after exit `partial.json` is absent while `chat.jsonl` contains the interrupted assistant message (baseline on `main`: `partial.json` present, message absent until next load). `{ ms }` < 2000 in the flowing-stream case. 2. **False-settle class eliminated:** the scripted reproduction fails on the pre-coder#3949 merge-base and passes on main after PR 2; `settleWorkspaceTurn` rejects any settlement without an enumerated cause; the coordinator's own Mux sessions show zero "superseded by an uncorrelated workspace stream-end" in the two weeks after PR 2 (soft signal, logged in the wave summary). 3. **Startup:** timeout and error-identity tests pass under TestClock; `[startup]` per-step lines unchanged in the sandbox transcript; a throwaway build with the constant set to 1 ms shows `Startup failed: StartupStepTimeoutError: <step> exceeded 1 ms` and a clean exit. 4. **No new lifecycle flakes:** 0 failures attributable to the touched suites across **N = 20** consecutive *completed* `Test / Unit` runs on `main` after the last Wave 4 merge β query `gh run list --workflow pr.yml --branch main --limit 60 --json databaseId,status,conclusion,event,headSha` (note: `gh run list --json` serializes these fields in **lowercase**, e.g. `{"status":"completed","conclusion":"success"}`, unlike `statusCheckRollup`), keep `status === "completed"` (pending runs have an empty `conclusion`, not null), take the newest 20, and for any run with `conclusion !== "success"` (case-insensitive normalization acceptable) inspect the failing job's log for the touched suite names (job `timeout`/`cancelled` from the 15-min budget is not a flake); plus green merge-queue runs for each PR. Any attributable flake β fix or revert before declaring done. **OFF-RAMP (PR 1):** fires if pre-work 1β3 shows (a) routing shutdown through `cancelStreamSafely` cannot preserve crash-recovery semantics without changing `cleanupAbortedStream`'s contract beyond D3, (b) the chaos variant exposes a double-settle not closable by D3(b), or (c) the finalizer cannot fit the 2 s bound for flowing streams. Then: stop PR 1, keep `AppFiberScope` unoccupied, update `appRuntime.ts` "Deliberately not done" with the concrete blocker and the measured evidence, land D3 alone as a bug-fix PR. PRs 2β4 are independent and proceed. **OFF-RAMP (PR 3):** if error identity or the `tests/ipc`/ACP paths cannot be preserved, keep `initialize()` as is and record why. **OFF-RAMP (PR 2):** coder#3949 not merged β hold (see PR 2). ## 5. Risk register | Risk | Likelihood | Mitigation | |---|---|---| | Crash-recovery regression: double commit / partial resurrection when shutdown-abort races completion | medium (window widens with dispose()) | D3(a) guard + test; `commitPartial`'s `historySequence` update-or-append is idempotent (`historyService.ts:2036β2041`) | | Provider abort emits an `error` chunk β error path instead of abort path | low | identical to today's user-stop path (parity); chaos variant covers hostile streams | | Wedged provider pins the 2 s bound β warning every shutdown | low | `boundedTeardown` already bounds; transcript measures; no budget change possible (5 s outer) | | AIService `stream-abort` listener (`commitPartial`) still in flight when `process.exit` runs | low | durable order writePartial β commit β deletePartial; next-load recovery; PR 1 transcript checks `partial.json` is already gone when `cli/server.ts` logs its final cleanup line before `process.exit(0)` (`:252β266`) | | Chaos-test seams (`createStreamResult`, `tokenTracker`) | none if scope-less construction stays default | new variant added, old cases untouched | | Collision with coder#3915/coder#3949 | medium | PR 2 gated; no edits to their regions; one-line `cause:` conflicts only | | RC churn (rc.113+ renames `forkIn`/`onInterrupt`/`toAsyncIterableWith`) | low | all Effect imports already in `streamManager.ts`/`streamBridge.ts`; pins fixed; GA upgrade is a separate lockstep PR (Β§6) | | Startup false timeout on slow hosts | medium if constant too small | measure first; β₯ 10Γ slowest observed; generous default (60 s) | | Sync-start assumptions in tests that `Reflect.set(processStreamWithCleanup)` | low | promise assigned before fork; forkIn `runSync` synchronous | | `shutdown()` (desktop second `before-quit` listener) still does not await streams | accepted | contract says `shutdown()` never touches the runtime; desktop's dispose race is the covered path | | Streams starting during shutdown | covered | `forkIn` on closed scope interrupts immediately β `system` abort (`startImmediately` semantics verified in rc.112; pinned by test) | | `system` abort triggers an in-session RetryManager retry during shutdown | **unreachable** (verified) | `"aborted"` β `NON_RETRYABLE_STREAM_ERRORS` β `retryManager.ts:99β104` abandons; no fiber scheduled | | PR 3: abandoned `taskService.initialize` mid-`editConfig` when the root exits after a timeout | low | desktop/ACP already dispose on startup failure; PR 3 adds the missing `cli/server.ts` dispose; config writes are lock/journal-protected | | `startImmediately`/`onInterrupt` semantics differ in a later RC | low | pinned by the probe test in `appFiberScope.test.ts` style; RC bumps are a separate lockstep PR | ## 6. Standing item β effect v4 GA + `@orpc/experimental-effect` lockstep (analysis only) v4 is **not GA** (rc.112 is current; v3 `3.x` remains the stable line). No PR this wave. When GA ships: one lockstep PR bumping `effect` + all `@orpc/*` (`1.14.11` today; check the GA-compatible `@orpc/experimental-effect`), canary gates = `di/*.test.ts`, `streamBridge.test.ts`, `streamManager.test.ts`, `serviceContainer.test.ts`, `TEST_INTEGRATION=1 bun x jest tests/ipc`, `make static-check`. The `Context β ServiceMap` rename risk is firewalled: `Context.Service` tags, `Context.omit/get`, `Layer`, `ManagedRuntime`, `TestClock` live only under `di/` + `orpc/effectContext.ts`; `streamManager.ts`/`streamBridge.ts` use `Effect`/`Scope`/`Fiber`/`Exit`/`Stream`/`Queue`/`Cause` only. PR 4 adds one `Context.Context<never>` type reference to `streamBridge.ts` β keep it as a type-only import so a rename is a one-line fix. ## 7. Dogfooding (per PR; evidence attached to the PR with `gh β¦ --attach`) Common setup: `make dev-server-sandbox DEV_SERVER_SANDBOX_ARGS="--clean-projects"` (or `xum server` on a temp `XUM_ROOT`) with `XUM_LOG_LEVEL=debug`, run under `script -f ~/wave4-scratch/<pr>-<scenario>.log`; scratch under `$HOME/wave4-scratch/` (never `/tmp`). Drive the UI with `agent-browser` (`open` β `snapshot -i` β click the explicit "Send message" ref; re-snapshot after typing). Screenshots are primary evidence; record WebM and finalize with `ffmpeg -c copy`. - **PR 1:** (1) start a long stream (prompt that streams ~30 s), wait 3β5 s, `kill -TERM <server pid>`; transcript must show the order in STOP coder#1 and `[shutdown] AppFiberScope closed { ms }`; (2) `ls <XUM_ROOT>/sessions/<ws>/partial.json` (absent) + `tail -n 1 chat.jsonl` (interrupted assistant message); (3) restart, open the workspace in agent-browser, screenshot the persisted interrupted message; (4) same scenario on `main` for the baseline diff; (5) `xum run` Ctrl-C mid-stream transcript (CLI root parity); (6) quality gate between phases: gate suites green before the sandbox run, sandbox evidence before requesting review. - **PR 2:** primary evidence is the regression file run on both the merge-base worktree (failing output) and the branch (passing). Secondary: sandbox parent workspace delegates via `task` kind=workspace to a child that arms a background bash monitor firing within ~10 s and keeps working ~60 s; screenshot the parent's task result (baseline `main`: `interrupted β¦ uncorrelated workspace stream-end`; after: `completed`) and the child's log lines. - **PR 3:** `xum server` cold start transcript with `[startup] <step> { ms }` for the six steps (parity); throwaway worktree build with `STARTUP_STEP_TIMEOUT_MS = 1` β transcript of `Startup failed: StartupStepTimeoutError β¦` and exit code (do not ship); desktop dialog cannot be shown headless β cite the unchanged `desktop/main.ts:1249β1265` catch. - **PR 4:** agent-browser session left idle 60 s with the connection indicator visible (heartbeats keep it green) + screenshot; `streamBridge.test.ts` on TestClock. ## 8. Non-goals (restated; out of this wave) Typed-error propagation sweep / removing the ~113 facades; converting services to `yield*`-based Effect services; PubSub for the internal EventEmitter bus; Schema at persistence boundaries; Effect observability; converting sync read paths, AI-SDK per-request callbacks, cross-process lock interiors, or deterministic try-lock funnels; replacing `AbortController` as the SDK cancellation transport; converting the `fullStream` loop to `Stream`; fiberizing turn-handle waiters; downgrading startup steps to best-effort; changing outer quit budgets; supervising the pre-registration stream-start window; `shutdown()` semantics; the effect GA bump (standing analysis only). </details> --- _Generated with `xum` β’ Model: `anthropic:claude-fable-5-1` β’ Thinking: `xhigh` β’ Cost: `$13.39`_ <!-- mux-attribution: model=anthropic:claude-fable-5-1 thinking=xhigh costs=13.39 -->
Summary
StreamManagerbecomes the first (and intended) occupant ofAppFiberScope: every started stream is wrapped in one supervisor fiber, soServiceContainer.dispose()(desktop,xum server, ACP) and the CLI cleanup lists now abort and await in-flight streams β partial flushed with usage,stream-abort(abortReason: "system") delivered, the interrupted assistant message committed intochat.jsonl,partial.jsonremoved β before the bridges and sessions are torn down. The same PR closes the two adjacent cancel races the shutdown path would have widened (Wave 4 plan D3): a cancel landing after the stream loop finished no longer resurrectspartial.jsonor emits a second terminal event, and concurrent cancellers (user stop racing dispose) join one cleanup.Background
Until now
dispose()had nostreamManagerstep: an in-flight stream died with the process and was reconciled frompartial.json(β€ 500 ms stale) on the next load, leaving an empty assistant placeholder row inchat.jsonluntil then.AppFiberScope(Phase 11) was built for exactly this and had no occupant. This is PR 1 of the Effect migration Wave 4 plan (appended below): D1 β the fiber is the ownership/supervision unit, theAbortSignalstays the cancellation transport; D2 β supervisor topology; D3 β adjacent cancel races; D4 β the 2 s close bound stays (outer budgets are 5 s).Implementation
StreamManager.superviseEngine(D2): afterstreamInfo.processingPromiseis assigned (byte-identical) it forksEffect.promise(() => processingPromise)(zero-arity thunk β rc.112 allocates no internal AbortController) intoengineScopewith{ startImmediately: true }.onInterruptβEffect.uninterruptible(Effect.promise(async () => { await cancelStreamSafely(ws, info, "system"); await info.completionController.promise; }));catchDefectβlog.warn. A stream that finishes on its own exits the fiber, which removes its scope finalizer (no residue). A stream started after the scope closed is interrupted synchronously byforkInand aborted right away (fail-closed during shutdown, pinned inappFiberScope.test.ts). The pre-registration window (pendingStreamStarts) stays unsupervised (documented).engineScope?: Scope.Closeableis a trailing optional 6th constructor parameter, wired fromAppFiberScopeTaginStreamManagerLive(di/layers/core.ts;AppFiberScopeTagadded toCoreInputTagsβ both roots already provide it viaruntimeSeams). Defaultundefined= today's behavior for every direct construction (tests,aiService.tscompat path).cleanupAbortedStream: afterawait processingPromise, return ifterminalCompletion !== undefined(completed/failed while the cancel was in flight). Plus apartialRetiredmarker set in the completion path right beforedeletePartial:flushPartialWritebecomes a no-op afterwards, because the pre-abort flush incancelStreamSafely(not only the abort bookkeeping) re-createdpartial.jsonwhen the cancel landed betweendeletePartialandCOMPLETED.cancelStreamSafely: per-streamcancelPromiselatch, checked and assigned synchronously at entry (after the existingCOMPLETEDearly return, before the firstawait).checkSoftCancelStreamshares the latch (??=) so a soft interrupt racing a hard cancel/dispose also yields exactly onestream-abort.appRuntime.tscontract ("Deliberately not done" updated; step 2 now says what it aborts),appFiberScope.ts,serviceContainer.tsstep comment. New[shutdown] streamManager.abortStream { workspaceId, messageId, ms }debug line per supervised stream (shutdownStep style).Net product change: 5 files, +196/β42 (β 70 non-comment lines added).
PR 1 notes
Pre-work findings
processStreamWithCleanupnever rejects β body istry { β¦ } catch { await handleStreamFailure } finally { β¦ }; a throw fromhandleStreamFailure/finallycould still reject, butstartStreamassignsprocessingPromise = processStreamWithCleanup(...).catch(log.error), so the promise the supervisor wraps never rejects.catchDefectstays as belt-and-braces.appFiberScope.test.ts, +2 cases):closeScopeBoundedinterrupts a fiber suspended onEffect.promisewithout settling the wrapped promise, runs the asynconInterruptfinalizer to completion, and resolves only afterwards (["finalizer-start", "finalizer-end"]).completionController.promise, i.e. aftercleanupAbortedStream(writePartialwith usage) βemitStreamAbortβ AIService sink (readPartialβcommitPartialβdeletePartialβemit("stream-abort")to agentSession/taskService/analytics listeners) β.finally(settle). The only work left after it is the EventEmitter listeners' own async continuations (turn-phase transitions, task-handle settlement) β the same as today'ssystemaborts fromtaskService. Durable orderwritePartial β commitPartial β deletePartialholds; a force-exit between any two steps leaves eitherpartial.json(recovered byagentSession.init'scommitPartial) or an already-committed row whosehistorySequenceupdate-or-append is idempotent. No OFF-RAMP condition fired.[shutdown] AppFiberScope closed { ms }(xum server, SIGTERM,XUM_LOG_LEVEL=debug,script -f): idle 1 ms; one flowing Sonnet 5 stream 56 ms; two flowing streams 95 ms (per-streamabortStream91 / 94 ms, parallel); wedged provider (scratch env-gatedfullStreamthat never yields and ignores abort, not committed) 2050 ms with theteardown timed outwarning,AppRuntime disposed5 ms (idempotent re-close), process exit 2.10 s after SIGTERM β inside the 5 s force-exit. Baselinemain: 1 ms (nothing supervised).forkInpinned (appFiberScope.test.ts):startImmediately: trueinto an already-closed scope runs the body to its first async boundary and then theonInterruptfinalizer synchronously (["body-started", "interrupted"], fiber exit is a failure).streamManager.test.tspins the product consequence (late start β{ status: "aborted", abortReason: "system" }, onestream-abort, registry empty).Deviations from the plan (all additive)
completionController.promise(plan:cancelStreamSafelyonly).cancelStreamSafelyreturns before abort delivery β thestream-abortsink is where AIService commits the partial β so without this the close would resolve withpartial.jsonstill present and the commit racingdesktopBridgeServer.stop(). This is what makes STOP criterion Better authentication UXΒ #1 ("partial.jsonabsent immediately after exit") true.partialRetired(see D3(a) above): the plan's guard alone left the pre-abort flush as a resurrection path; the D3(a) test was red against the guard-only variant on exactly that assertion.checkSoftCancelStreamjoins the latch: dispose during a pending soft interrupt (queued-message preemption at tool-end) is the same double-cleanup shape D3(b) fixes for hard cancels.[shutdown] streamManager.abortStreamdebug line β there was no log line for the abort at all, so the transcript could not show the plan'sstream-abort β AppFiberScope closed β desktopBridgeServer.stoporder; it now does.xum runCtrl-C transcript:xum runinstalls no SIGINT handler (onlycli/server.tsdoes), so Ctrl-C is Node's default immediate termination (exit 130) and the cleanup list β includingappFiberScope.closeβ runs only on normal completion. The CLI-root path is pinned by the newcoreServicesRoot.test.tscase instead; adding a signal handler toxum runis out of scope.Dogfooding evidence (transcripts + before/after in the first comment)
xum server, SIGTERM 3 s after first text)AppFiberScope closedpartial.jsonright after exitchat.jsonlrowmain@ 3c06630 (baseline)partial: true, 1172 charsRestart UX (screenshot below, branch): the persisted interrupted text is shown above the
INTERRUPTEDbarrier; the workspace then auto-resumed the turn on startup (agentSession startup auto-retry β"system"aborts do not suppress it). Verified the baseline does the same after it commits the partial on load ([STREAM MESSAGE]right after startup onmain), so the only behavioral difference is when the partial is committed (at shutdown vs. next load). Not exercised headless: Electronbefore-quit(samedispose(); CI e2e).Lessons for PR 3 (
initialize()as a startup effect)runPromise/facades reject with the raw error, so error identity is free β but anyEffect.promisethunk wrapping anawaited step must beasyncand the whole pipeline needscatchDefectif the facade must never reject.dispose()are cheap (disposeAppRuntimeafter a timed-out scope close took 5 ms becauseScope.closeis idempotent); a startup timeout should equally not try to unwind the abandoned step β parity with the existing "throwing step" path is the whole design.[startup] <step> { ms }lines are already the right measurement surface; this run's cold start showedtaskService.initialize114 ms andworkspaceService.initialize83 ms, so a 60 s per-step bound is β₯ 500Γ the slowest observed step.xum runparity caveat in mind: the CLI roots have no signal handling, so any "dispose on failed initialize()" work only applies where a dispose path exists (cli/server.tsis the one PR 3 adds).Validation
streamManager.test.ts(6: flowing close, wedged bound, late start, D3(a), D3(b) race, 50-stream residue),streamManager.chaos.test.ts(1 new fuzz variant, existing cases untouched; also ran 7 extra seeds locally),serviceContainer.test.ts(dispose aborts + awaits a real stream beforedesktopBridgeServer.stop(), checkspartial.json/chat.jsonl),coreServicesRoot.test.ts(same viacloseScopeBounded(appFiberScope)),appFiberScope.test.ts(2 probes).stream-aborts).streamManager.test.ts,streamManager.chaos.test.ts,streamManager.modelOnlyNotifications.test.ts,aiService.test.ts,agentSession.disposeRace.test.ts,agentSession.sinceReplayContract.test.ts,serviceContainer.test.ts,coreServicesRoot.test.ts,di/*.test.ts,taskService.test.ts,workspaceService.test.ts(one unrelatedFakeAIServicemetadata test flaked once in a combined run, green isolated and on rerun),turnRequestBuilder.test.ts;make static-check.catchDefect; async thunks), spy seams unchanged (processStreamWithCleanup,createStreamResult,createStreamAtomically,startStream; constructor arity trailing-optional;Reflect.set/gettargets valid), sync-start (processingPromiseassigned before the fork;runSync(forkIn)synchronous), no constructor side effects, zero-suspension latch.Risks
xum server/ACP): a flowing stream now costs one chunk (~50β100 ms) at dispose; a wedged provider costs the existing 2 s bound and a warning β identical outcome to today's process exit. Rollback: revert the onecore.tswiring line (engineScopeundefined β today's behavior); the D3 guards are independent bug fixes.abortReason/abandonPartialwin; a discard racing a system cancel by a few ms keeps the partial instead of dropping it (previously both cleanups ran and emitted two aborts)."system"(fail-closed): the caller sees a normalOk(handle)whose completion settlesaborted, and the turn is recovered on next load exactly like any othersystemabort.π Implementation Plan
Effect migration β Wave 4: finish the concurrency/lifecycle core
Bounded wave: 4 PRs (PR 4 optional), explicit STOP criterion, explicit OFF-RAMPs. Plan only; nothing here is implemented.
0. Thesis check (coordinator's judgment vs. evidence)
Thesis: Effect's payoff in this app is structured concurrency + interruption-safe lifecycles in the orchestration core (still Promise + AbortController).
Verdict: holds for the stream engine; only half-holds for turn handles.
ServiceContainer.dispose()never stops or awaits in-flight streams (serviceContainer.ts:478β537has nostreamManagerstep); an in-flight stream dies with the process and is recovered on next load frompartial.json(β€ 500 ms stale,PARTIAL_WRITE_THROTTLE_MS,streamManager.ts:776).AppFiberScopeexists precisely for this and has no occupant. A supervised per-stream fiber is the right tool.interruptWorkspaceTurnFromUncorrelatedStreamEnd,workspaceTurnManager.ts:4220β4307: any uncorrelated stream-end after the prompt index settles the handleinterrupted), not a Promise-vs-fiber structure bug. Turn handles are persisted records (taskHandleStore.upsertWorkspaceTurn) spanning multiple streams (tool-call continuations are deferred viahasSameTurnContinuation,:4491) and surviving restarts; a fiber/Deferred can only model the in-process waiter and would not fix correlation. Open PR [task-service] π€ fix: preserve turns across synthetic wake endsΒ #3949 fixes the predicate in Promise idiom and is Codex-green. Wave 4's turn-handle PR therefore becomes "codify the settlement invariant + prove the class is gone", not "fiberize handles" (D5 below).Corrected baseline numbers (measured this workspace): 46/470
src/nodenon-test files importeffect(coordinator said 35); 113 directEffect.run*sites outsidedi/in 16 files; 226Effect.gen; 9TaggedErrorclasses; effect4.0.0-rc.112,@orpc/*1.14.11; effect v4 is not GA (rc line still current).1. Verified current state (evidence the design rests on)
Stream engine (streamManager.ts)
startStream(:4723β4901): per-workspace mutex βnew AbortController()+linkAbortSignal(:4771β4772) βresourceScope = Scope.makeUnsafe()(:4777) β temp-dirEffect.acquireRelease(:4802β4824) βcreateStreamAtomicallyβstreamText(:2244,abortSignal: abortController.signal:2250) β registered inworkspaceStreams(:2463) βstreamInfo.processingPromise = this.processStreamWithCleanup(...)fire-and-forget (:4876β4882) β returnsOk({ messageId, completion }).processStreamWithCleanup(:3331β4089, plain async):while(true)retry loop;for await (part of fullStream)(:3358β3837) with abort check at loop head (:3361); post-loopif (!signal.aborted)gate (:3849) β completion path (deletePartial:3981,updateHistory:3989,recordSessionUsage:4001,state = COMPLETED:4017, emitstream-end:4023,terminalCompletion:4024); error path βhandleStreamFailure(:4094β4112) βpersistStreamErrorwrites error partial;finally(:4052β4088): release MCP lease,Effect.runFork(Scope.close(resourceScope))(:4064β4066), unlink abort,workspaceStreams.delete,eventSpine.emit("stream.end"),completionController.settle.stopStream(:5043β5111) βcancelStreamSafely(:1766β1800):if (state === COMPLETED) { await processingPromise; return }βstate = STOPPINGβflushPartialWriteβabortController.abort()βcleanupAbortedStream(:1828β1951):await processingPromiseβ usage βwritePartial(:1876β1910) βemitStreamAbortβsettle({status:"aborted"}). No completed-guard after the await (verified:1838β1951): a cancel landing between:3849and:4017re-writespartial.jsonafterdeletePartialand emitsstream-abortafterstream-end(pre-existing window; dispose() will widen its exposure).cancelStreamSafelyis also not idempotent for concurrent callers (onlyCOMPLETEDis checked).stream-abort(aiService.ts:355β377):abandonPartial ? deletePartial : commitPartial β deletePartial(fire-and-forget listener).HistoryService.commitPartial(historyService.ts:1963β2061) β strips error metadata,hasCommitWorthyParts, stale-epoch check, update-or-append byhistorySequence, delete partial; invoked fromagentSession.init(:5002),aiService.streamMessage(:886), stream-abort (:364),duplicateWorkspace.StreamAbortReason = "user" | "startup" | "system"(src/common/orpc/schemas/stream.ts:295).Reflect.set(streamManager, "tokenTracker" | "createStreamResult")(streamManager.chaos.test.ts:130β134, 238β242);streamManager.test.tsReflect.setonprocessStreamWithCleanup(:2787),createStreamAtomically(:2783),createTempDirForStream,cleanupStreamTempDir,Reflect.getonworkspaceStreams,schedulePartialWrite, β¦;modelOnlyNotifications.test.tscallsprocessStreamWithCleanupdirectly (:93, :187);aiService.test.tsspiesstartStream,generateStreamToken,createTempDirForStream,isResponseIdLost. Constructor:(historyService, sessionUsageService?, getProvidersConfig?, eventSink = noop, runner = defaultEffectRunner)(:801β813);effectRunnerused at:1152, :1154, :1169only.workspaceId+messageId;stream-end/stream-abort/errorcarrymetadata.muxMetadatawhen the prompt had it.DI / shutdown / startup
AppFiberScopeLiveis inCoreLive'sruntimeSeams(di/layers/core.ts:644), so both roots have it;StreamManagerLive(core.ts:226β238, stage S2b) already yieldsEffectRunnerTag. CLI cleanup lists includeappFiberScope.close(cli/run.ts:1579,cli/workflow.ts:286).APP_FIBER_SCOPE_CLOSE_TIMEOUT_MS = 2000,APP_RUNTIME_DISPOSE_TIMEOUT_MS = 2000; outer budgets are 5000 ms on both desktop (desktop/main.ts:1297Promise.racevssetTimeout(5000)) andxum server(cli/server.ts:236β243force-exit). The scope bound cannot grow without changing outer budgets.node_modules/effect/dist/internal/effect.js:2264:forkInregisters a scope finalizer and removes it when the fiber completes (no leak), and interrupts immediately if the scope is already closed (streams starting mid-shutdown fail closed).Effect.promise(evaluate: (signal) => PromiseLike),Effect.onInterrupt,Effect.forkIn(_, scope, { startImmediately? }),Stream.toAsyncIterableWith(context),Stream.provideContextall exist.ServiceContainer.initialize()(serviceContainer.ts:297β362): six awaitedinitialize()s wrapped inrecordStep(durations only, no catch, no timeout) + three syncstart()s + two fire-and-forget sweeps. Failure handling: desktopStartup Faileddialog +app.quit()(desktop/main.ts:1249β1265);cli/server.ts:136uncontained; ACPserverConnection.ts:205β216dispose + rethrow;tests/ipc/setup.ts:85no catch. No outer timeout anywhere.streamBridge.subscriptionIterable(orpc/streamBridge.ts:176) βStream.toAsyncIterable(...)on the global runtime, 19 call sites inrouterSubscriptions.ts; heartbeat viaEffect.sleepinforkScoped(:145β152).streamBridge.test.tshas 11 real-time waits, but only 3 are clock-bound (:2071 ms initial delay,:241heartbeat 10 ms,:25510 ms laziness); 8 arewaitFor(listenerCountβ¦)readiness polls that TestClock cannot replace.Turn handles + open PRs
{ handleId "wst_β¦", ownerWorkspaceId, workspaceId, turnId, messageId, status, attentionPolicy, disposableWorkspace }; prompt carriesmuxMetadata: { type:"workspace-turn-task", taskHandleId, ownerWorkspaceId, turnId }(workspaceTurnManager.ts:1420).TaskServiceforwardsaiServicestream-end/stream-abort/errortofinalizeWorkspaceTurnFromStreamEnd(:4442β4544): correlated branch matchesrecord.workspaceId && record.turnId(:4472); uncorrelated branch (metadata == null, notagentId === "compact") βinterruptWorkspaceTurnFromUncorrelatedStreamEndβ settlesinterruptedwheneverstreamEndIndex >= promptIndex(:4293β4305). Producers of such uncorrelated ends: bash-monitor wake continuations, child terminal-attention deliveries, heartbeat, peer messages, parent auto-resume.cleanupDisposableWorkspaceTurnβworkspaceService.remove(β¦, true)kills its background processes; persistent child β parent seesinterruptedβtask_stopβbackgroundProcessManager.stopMonitor(β¦, "canceled"). This is the observed "monitors died afterwards".settleWorkspaceTurn(params)(:2085), 11 callers, guarded byworkspaceTurnSettlementLocks.withLock(handleId); waiters inpendingWorkspaceTurnWaitersByHandleIdwithsetTimeouttimeouts (:2425β2496).isManualChildWorkspaceInput); otherwise ignores the end. Touches:281β295, :4217β4355+ tests (+286/β31). Codex: "Didn't find any major issues" + clean security onf9baa2fc9.mergeable: MERGEABLE, butTest / UnitandCodex Commentsred, 19 commits behind main.getWorkspaceTurnLiveness/getWorkspaceTurnRuntimeActivity(identity-matches the active stream'smuxMetadataagainst the record) for staleness/capacity. Touches:442β486, :1298, :2573, :3761, :3800β4064(+494/β60).BLOCKED, latest Codex review has open findings,Test / Unitred, 19 behind.2. Design decisions
D1 β Fibers WRAP the AbortController; they do not replace it.
The AI SDK is cancelled only via
AbortSignal; thefor awaitloop, soft-interrupt at step boundaries, retry/fallback re-creation ofstreamResult, and ~30 abort touchpoints (#4032) all key off the signal. Converting the 750-line loop toStream.fromAsyncIterable+ fiber interruption would touch hundreds ofWorkspaceStreamInfotransitions and break theprocessStreamWithCleanup/createStreamResultspy seams. Instead: the fiber is the ownership/supervision unit; the signal stays the cancellation transport. The dual-cancellation glue #4032 feared is confined to one point β the supervisor'sonInterruptβ which routes through the existing user-stop path (cancelStreamSafely), so shutdown β‘ "user pressed stop" semantically (partial flushed with usage,stream-abortemitted,completionsettlesaborted, AIService commits the partial).D2 β Supervisor topology: one supervisor fiber per stream in
AppFiberScope, wrapping the already-startedprocessingPromise.streamInfo.processingPromise = this.processStreamWithCleanup(...)stays byte-identical (sync-start preserved;Reflect.set(processStreamWithCleanup)seam preserved;cleanupAbortedStream'sawait processingPromiseunchanged). Immediately after it:Effect.promiseis interruptible while suspended (internal/effect.js:741β801, Async op);onInterrupt=onErrorFilter(causeFilterInterruptors, β¦)(:1762);forkInregistersfiberInterrupt(fiber)as the scope finalizer (:2264β2275),fiberInterruptawaits the fiber (:635β642), and parallelscopeCloseawaits all finalizers viafiberAwaitAll(:1590β1601) βcloseScopeBoundedat dispose step 2 gives "interrupt and await" whilehistoryService/sessionUsage/eventSink β AIService β bridge serversare still alive (bridges stop in step 3, so clients receivestream-abort).forkIn's observer removes the scope finalizer (verified) β no per-stream residue.forkInon a closed scope callsfiber.interruptUnsafesynchronously and returns the fiber (:2272β2274,runSyncdoes not defect). WithstartImmediately: true,forkUnsaferunschild.evaluatesynchronously (:2233β2247) up to theEffect.promiseAsync op (:772β801), so the fiber is suspended (_running=false) when the interrupt lands andinterruptUnsafe(:391β409) unwinds the stack through theonInterrupthandler β the stream is aborted assystem(fail-closed during shutdown). Verified in rc.112 internals (effect.js:2233β2247, 391β409); pin with a test ("stream started after scope close is aborted") so an RC bump cannot silently change it."system"is semantically exact:"user"/"startup"suppress next-startup recovery (retryEligibility.ts:114β118, 284β287),"system"marks an involuntary backend interruption (astaskService.ts:8100, 8223use it). No in-session retry loop is possible: thestream-aborthandler (agentSession.ts:6010) routes{ type: "aborted" }toretryManager.handleStreamFailure, and"aborted"is inNON_RETRYABLE_STREAM_ERRORS(retryEligibility.ts:49β59, 106) βretryManager.ts:99β104abandons immediately, never schedules a fiber. Dogfooding still checks the restart UX (the recovered partial is shown as interrupted; note whether any next-startup recovery re-sends β same class as today'ssystemaborts fromtaskService).engineScopearrives as an optional 6th constructor parameter (engineScope?: Scope.Closeable), wired fromAppFiberScopeTaginStreamManagerLive(core.ts:226). Defaultundefinedkeeps every direct-construction test andaiService.ts:174path identical (I4).AppFiberScopeLivealready sits beneath S2b inruntimeSeams, so no staging change (I6)."system"β no wire/schema change; UI copy forsystemalready exists.pendingStreamStarts, before registration) is not supervised: nothing is persisted for it yet, andstopStreamalready aborts pending controllers. Documented, not fixed.D3 β Fix the two adjacent cancel races in the same PR (closely-related bugs, not deferrals).
(a)
cleanupAbortedStream: afterawait processingPromise, ifstreamInfo.terminalCompletion !== undefined(completed/failed while the cancel was in flight) β return without abort bookkeeping (preventspartial.jsonresurrection afterdeletePartialand astream-abortafterstream-end). (b)cancelStreamSafely(:1766): latch a per-streamcancelPromiseso concurrent cancellers (user stop racing dispose) join one cleanup β exactly onestream-abort, onesettle. Zero-suspension requirement: the latch must be checked and assigned synchronously at function entry, before anyawait(the current first await isflushPartialWriteat:1789) β otherwise racing callers can both entercleanupAbortedStream. Shape:Both are β€ 10 LoC and get behavioral tests.
D4 β Shutdown bound stays 2 s; the finalizer must be fast or abandoned.
Outer budgets are 5 s; 2 s + 2 s already consume 4 s. A flowing stream aborts within one chunk; a wedged provider (no chunks, ignores abort) hits the existing
boundedTeardowntimeout: warning, continue, process exit β identical to today's outcome. Dogfooding measures the actual[shutdown] AppFiberScope closed { ms }with a live stream.D5 β Turn handles: codify the settlement invariant; do not fiberize.
Invariant: a workspace-turn handle settles terminally only by (i) a stream terminal event whose
muxMetadatacorrelates{taskHandleId, ownerWorkspaceId, turnId}to the record; (ii) an explicit interrupt (task_stop/interruptWorkspaceTurn); (iii) manual supersession β a manual child input after the turn anchor; (iv) stale-liveness reconciliation. An uncorrelated stream-end is never terminal by itself. #3949 makes (iii) the only uncorrelated outcome; #3915 implements (iv) by identity. Wave 4 adds acausediscriminant tosettleWorkspaceTurn(the single chokepoint) with a runtime assertion, plus the regression harness. Rationale for not converting waiters toDeferred/fibers: no behavioral gain, 4.9k-line file, and the coordinator's "settle only on the owning stream's termination" is over-specified β a turn owns several streams.D6 β Startup:
initialize()stays a Promise facade over a runtime-run startup effect; timeout β same failure path as a thrown step.Each step is
Effect.tryPromise({ try: async () => step(), catch: identity }).pipe(Effect.timeoutOrElse({ duration: STARTUP_STEP_TIMEOUT_MS, orElse: () => Effect.fail(new StartupStepTimeoutError(name, ms)) }))(timeoutOrElseexists in rc.112,Effect.d.ts:7833; chosen overtimeout+catchTagbecause the step's error channel isunknown, whichcatchTagcannot narrow). NoforkDetachneeded: a Promise step keeps running on its own when the waiting fiber times out (not inside an uninterruptible region, so the timeout interrupts the wait directly).StartupStepTimeoutError extends Errorwithname = "StartupStepTimeoutError"set in the constructor and message"<step> exceeded <ms> ms"(so the desktop dialog's error formatting shows both the class and the step name) β desktop shows it in the existingStartup Faileddialog; CLI/ACP/tests paths unchanged. Step errors keep their identity (v4runPromiserejects with the raw failure). Downgrading any step to best-effort is a policy change, out of scope (audit of the six implementations: extensionMetadata/telemetry/experiments are local fs, <50 ms; policy has its own 10 s fetch timeout; workspaceService bounds its sync internally; onlytaskService.initializeβ config scan +editConfig+ recoverysendMessages β is potentially unbounded). The threestart()s stay sync (Effect.sync), the two fire-and-forget sweeps stay outside the effect.stepDurationsMsis preserved.Abandon-and-quit safety: an abandoned
taskService.initializemay be mid-editConfigwhen the root exits. Parity requirement for PR 3: after a rejectedinitialize(), every root runs the boundeddispose()before exiting. Verified: desktop already does βservicesis assigned before the await (main.ts:653β656), the catch callsapp.quit(), and thebefore-quitlistener (:1271β1305, guardif (isDisposing || !services) return) racesservices.dispose()against 5 s; ACP does (serverConnection.ts:205β216);cli/server.tsdoes not (:133β136awaited at top level,main().catchat:282only logs) β PR 3 adds a boundeddispose()there (β€ 10 LoC, same 5 s budget).D7 β streamBridge: thread the runtime context, not a runner.
subscriptionIterablegainscontext?: Context.Context<never>βStream.toAsyncIterableWith(context);routerSubscriptionspasses the handler's"effect/context". Production behavior identical; heartbeat sleeps on the runtimeClock; tests can run the 3 clock-bound waits onTestClock. Honest scope: the 8 readiness polls stay.3. PRs (ordered by value Γ· risk; each independently mergeable)
PR 1 β StreamManager engine core becomes the first
AppFiberScopeoccupantValue: high (the only remaining shutdown data-integrity gap; the reason
AppFiberScopeexists). Risk: medium β low with D1/D2. Net product LoC β +55 (superviseEngine~25, ctor param/field ~5,engineFiberfield ~2, D3 guards ~12,core.tswiring ~2, doc updates inappRuntime.ts/appFiberScope.ts"occupant" text ~10).Files:
src/node/services/streamManager.ts,src/node/services/di/layers/core.ts,src/node/services/di/appRuntime.ts+appFiberScope.ts(docs), tests below.Pre-work (before writing product code; each yields a note in the PR body):
processStreamWithCleanupnever rejects (try/catch/finally shape:3331β4089); else the supervisor must fold rejections (it alreadycatchDefects).Effect.promiseinterruption +onInterruptawait ordering underScope.closein a 20-line probe test (pattern ofappFiberScope.test.ts:27β47).stream-abortlistener βcommitPartial; agentSession completion continuations) and confirm the durable order (writePartialβ commit βdeletePartial) makes a mid-flightprocess.exitrecoverable on next load (it is: partial survives until commit completes).[shutdown] AppFiberScope closed { ms }with a live stream in the sandbox (D4).startImmediately: trueinto an already-closed scope still runs itsonInterruptfinalizer (reviewer-verified in rc.112 internals; the test guards RC bumps).Acceptance (behavioral tests only):
streamManager.test.ts(new cases; existing cases untouched): withengineScope = Scope.makeUnsafe("parallel")and a fakecreateStreamResultwhosefullStreamyields onetext-deltathen blocks until itsAbortSignalfires βcloseScopeBounded(engineScope)resolves;writePartialwas called with the streamed text; exactly onestream-abort(abortReason: "system") and zerostream-end;completionsettles{status:"aborted"};workspaceStreamsis empty.closeScopeBoundedresolves within the bound, never rejects, warns once (assert the returned promise resolves and no throw; do not assert log text).COMPLETEDβ history has exactly one final message,partial.jsonabsent, event orderstream-endonly.stopStream+closeScopeBoundedracing on one stream β exactly onestream-abort, one settle.closeScopeBounded(engineScope)emits zerostream-abortand completes in the same tick class as an empty scope (assert no aborts andworkspaceStreams.size === 0).streamManager.chaos.test.tsβ existing cases byte-identical; one new fuzz variant constructs with an engine scope and closes it at a random iteration: every stream settles exactly once (count terminal events permessageIdβ€ 1, allcompletionpromises settle).serviceContainer.test.ts: "dispose() aborts and awaits an in-flight stream beforedesktopBridgeServer.stop()" (extend the ordering harness at:295β323);coreServicesRoot.test.ts:xum runcleanup list does the same viaappFiberScope.close.Gate suites:
streamManager.test.ts,streamManager.chaos.test.ts,streamManager.modelOnlyNotifications.test.ts,aiService.test.ts,agentSession.disposeRace.test.ts,agentSession.sinceReplayContract.test.ts,serviceContainer.test.ts,coreServicesRoot.test.ts,di/*.test.ts,taskService.test.ts,workspaceService.test.ts,turnRequestBuilder.test.ts;make static-check.House pre-review audits: interruption posture (supervisor's only suspension is the promise; finalizer uninterruptible end-to-end incl.
cancelStreamSafelyβcleanupAbortedStream); no defect escapes (catchDefecton the supervisor;Effect.promisethunksasync); spy-seam check (processStreamWithCleanup,createStreamResult,createStreamAtomically,startStreamsignatures unchanged; constructor arity unchanged, trailing optional); sync-start (processingPromiseassigned before fork;runSync(forkIn)completes synchronously); no constructor side-effects added; zero-suspension check on D3(b) latch (cancelPromisechecked-and-assigned synchronously atcancelStreamSafelyentry, before the firstawaitat:1789; review the diff for any insertedawait/lookup ahead of the assignment).Rollback: revert the
core.tswiring line βengineScopeundefined β today's behavior; D3 guards can stay (independent bug fixes).PR 2 β Turn-settlement invariant + false-settle regression harness (gated on #3949)
Value: high (7Γ production race). Risk: low. Net product LoC β +40 (
WorkspaceTurnSettlementCauseunion +causeonsettleWorkspaceTurnparams + assert ~10; 11 call sites Γ 1β3 lines).Relationship to open PRs β explicit:
interruptWorkspaceTurnFromUncorrelatedStreamEnd,isWorkspaceTurnAnchorForRecord,isManualChildWorkspaceInput), and adds the invariant + proof on top. If [task-service] π€ fix: preserve turns across synthetic wake endsΒ #3949 has not merged when PRs 1/3 are done: do not fork a competing fix; report to the coordinator, offer the regression test file to [task-service] π€ fix: preserve turns across synthetic wake endsΒ #3949's author as a review artifact, and hold PR 2 (it is not on any other PR's critical path).:3800β4064incl.settleStaleWorkspaceTurn, asettleWorkspaceTurncaller) overlaps PR 2's one-line-per-caller change. Prefer landing after it; if PR 2 must go first, the conflict is a one-linecause:addition per caller. PR 2 never edits liveness/reservation code.workspaceTurnManager.ts.Design:
type WorkspaceTurnSettlementCauseenumerated from the 11 callers (audited atmain@b87f62729)::1373creation validation failure;:1476pre-stream interrupt during launch;:1500/:1518pre-stream send failure;:3834/:3863stale-liveness recovery / restart timeout (settleStaleWorkspaceTurnβ #3915's region);:4301uncorrelated-stream-end manual supersession (the only uncorrelated settle in the codebase; the path #3949 rewrites);:4529correlated terminal;:4571stream-abort;:4675deferred stream error;:4736terminal stream error.settleWorkspaceTurnassertsparams.causeis a member and, formanual-supersession, that the superseding input'smessageIdis supplied β turning D5 into an exhaustiveRecord<Cause, β¦>check rather than prose, so a future "settle on uncorrelated end" cannot be added without naming (and justifying) a cause.Acceptance:
workspaceTurnManager.uncorrelatedStreamEnd.test.ts(realWorkspaceTurnManager+TaskHandleStore+ fakeaiServiceemitter, following the existing suite's harness): (1) create turn β correlatedstream-startβ synthetic wake stream on the same child ends uncorrelated after the anchor β handle staysrunning, waiter unresolved, no disposable cleanup, no terminal attention β correlatedstream-endβcompleted. (2) same withfinishReason:"tool-calls"continuation in between. (3) manual child input between anchor and end βinterruptedwithcause: manual-supersession. (4) explicitinterruptWorkspaceTurnβinterrupted,cause: explicit-interrupt. Case (1) is the scripted reproduction: it must fail on the pre-[task-service] π€ fix: preserve turns across synthetic wake endsΒ #3949 merge-base (run the file from a sibling worktree atgit merge-base origin/main <#3949 head>; record the failing assertion in the PR body) and pass after.workspaceTurnManager.test.tscases andtaskService.test.tsturn cases unchanged.Gate suites:
workspaceTurnManager.test.ts,taskService.test.ts,taskHandleStore.test.ts,tools/task*.test.ts;make static-check.Audits: spy-seam (
getWorkspaceTurn,listAllWorkspaceTurns,enqueueTerminalAttention,deliverPersistentChildWorkspaceTurnResultuntouched); settlement lock held across the assert; no new suspension insidewithLock.Rollback: revert; the test file stays valid against #3949 alone (drop the
causeassertions).PR 3 β
ServiceContainer.initialize()as a runtime-run startup effect with per-step timeoutsValue: medium (a hung
taskService.initialize()currently pins the splash screen forever; deterministic TestClock tests of startup). Risk: lowβmedium. Net product LoC β +80 (step table ~20, timed-step helper ~15,StartupStepTimeoutError~8, constant ~3, facade ~10, root dispose-on-failure parity β€ 10, doc update ~10).Files:
serviceContainer.ts,src/constants/terminationTimeouts.ts(keep with the termination constants so the budget doc stays in one place),cli/server.ts(dispose in the startup catch if missing),di/appRuntime.tsdoc ("Deliberately not done" β remove the initialize() line; add startup contract).Design (D6):
initialize(): Promise<void>βthis.runtime.managed.runPromise(this.startupEffect()).startupEffect = Effect.genover an orderedreadonly steps: ReadonlyArray<{ name, run: () => Promise<void> }>(assert names unique); each step:recordSteptiming kept,Effect.tryPromise({ try: async () => run(), catch: identity }).pipe(Effect.timeoutOrElse({ duration: STARTUP_STEP_TIMEOUT_MS, orElse: () => Effect.fail(new StartupStepTimeoutError(name, STARTUP_STEP_TIMEOUT_MS)) })). ThenEffect.syncfor the threestart()s; the sweeps remain afterrunPromise. ConstantSTARTUP_STEP_TIMEOUT_MSβ pre-work measures[startup] <step> { ms }across sandbox cold starts and picks β₯ 10Γ the slowest observed (propose 60 s; must be generous β a false timeout turns a slow-but-fine start into a crash). Roots dispose after a rejectedinitialize()(D6 abandon-and-quit safety).Acceptance (all in
serviceContainer.test.ts, TestClock via the existingAppLivespy at:355β395):initialize()rejects withStartupStepTimeoutErrornaming the step after exactlyTestClock.adjust(STARTUP_STEP_TIMEOUT_MS); later steps did not run.initialize()rejects with the same error object (identity), later steps did not run (parity with today).stepDurationsMshas all six keys;start()s called once each; secondinitialize()call behavior unchanged from today (verify whether re-entry is guarded today; preserve).tests/ipcharness and ACP entry still pass unchanged.Gate suites:
serviceContainer.test.ts,coreServicesRoot.test.ts,di/*.test.ts,src/node/acp/*.test.ts,TEST_INTEGRATION=1 bun x jest tests/ipc(smoke subset);make static-check.Audits: I1 untouched (no layer body changes); I2 (only the composition root touches the runtime); error identity preserved (no wrapping); abandoned-step safety β every root runs bounded
dispose()after a rejectedinitialize()(D6;cli/server.tsgains it in this PR); no sweep moved into the effect; the six-step order and[startup] <step>names unchanged.Rollback: revert; constant removal.
PR 4 (optional β cut if budget is exhausted) β
streamBridgeon the runtime contextValue: low (closes the last documented "global runtime" exception; enables TestClock for the heartbeat). Risk: low. Net product LoC β +30 (
context?option +toAsyncIterableWith~8; 19 call sites Γ 1 line via one shared helper inrouterSubscriptions.ts~3).Acceptance: the 3 clock-bound waits (
:207, :241, :255) run onTestClock; the heartbeat test asserts N heartbeats afterTestClock.adjust(N Γ interval)with zero real time; existing behavioral assertions unchanged;tests/ipcsubscription tests pass. Do not rewrite the 8 readiness polls.Gate suites:
streamBridge.test.ts,routerSubscriptions*.test.ts,orpc/*.test.ts,TEST_INTEGRATION=1 bun x jest tests/ipc(subscription subset);make static-check.Audits:
Stream.toAsyncIterableWithpreserves double-close safety (pin with the existing test);Cause.Donetyping unchanged; noScope/MemoMap/Schedulercaptured (pass the oRPCeffect/context, which the DI layer already strips perEffectRunnerLive);contextstays optional so direct callers/tests without a runtime keep today's global-runtime path.Rollback: revert; the optional
contextdefault (Context.empty()) is exactly today'stoAsyncIterable, so a partial revert of call sites is also safe.Execution order and size
Net product LoC for the wave β +205 (PR 1 β +55, PR 2 β +40, PR 3 β +80, PR 4 β +30); tests β +600β800. PR 1 starts immediately. PR 2 starts the moment #3949 merges (parallel with PR 1/3 β disjoint files). PR 3 after PR 1 merges (both touch
appRuntime.tsdocs; PR 3 also touchesserviceContainer.ts). PR 4 last, only if PRs 1β3 landed and no OFF-RAMP fired. Each PR: Codex dual review,Codex Commentsminimization, merge queue; commit WIP early (/tmpwipes).4. STOP criterion (measurable) and OFF-RAMPs
Wave 4 is done β and the Effect migration line stops without a new RFC β when all hold:
serviceContainer.test.tsordering test +coreServicesRoot.test.tspass on main; a sandboxscript -ftranscript ofxum serverreceiving SIGTERM mid-stream showsstream-abortβ[shutdown] AppFiberScope closed { ms }β[shutdown] desktopBridgeServer.stop, and immediately after exitpartial.jsonis absent whilechat.jsonlcontains the interrupted assistant message (baseline onmain:partial.jsonpresent, message absent until next load).{ ms }< 2000 in the flowing-stream case.settleWorkspaceTurnrejects any settlement without an enumerated cause; the coordinator's own Mux sessions show zero "superseded by an uncorrelated workspace stream-end" in the two weeks after PR 2 (soft signal, logged in the wave summary).[startup]per-step lines unchanged in the sandbox transcript; a throwaway build with the constant set to 1 ms showsStartup failed: StartupStepTimeoutError: <step> exceeded 1 msand a clean exit.Test / Unitruns onmainafter the last Wave 4 merge β querygh run list --workflow pr.yml --branch main --limit 60 --json databaseId,status,conclusion,event,headSha(note:gh run list --jsonserializes these fields in lowercase, e.g.{"status":"completed","conclusion":"success"}, unlikestatusCheckRollup), keepstatus === "completed"(pending runs have an emptyconclusion, not null), take the newest 20, and for any run withconclusion !== "success"(case-insensitive normalization acceptable) inspect the failing job's log for the touched suite names (jobtimeout/cancelledfrom the 15-min budget is not a flake); plus green merge-queue runs for each PR. Any attributable flake β fix or revert before declaring done.OFF-RAMP (PR 1): fires if pre-work 1β3 shows (a) routing shutdown through
cancelStreamSafelycannot preserve crash-recovery semantics without changingcleanupAbortedStream's contract beyond D3, (b) the chaos variant exposes a double-settle not closable by D3(b), or (c) the finalizer cannot fit the 2 s bound for flowing streams. Then: stop PR 1, keepAppFiberScopeunoccupied, updateappRuntime.ts"Deliberately not done" with the concrete blocker and the measured evidence, land D3 alone as a bug-fix PR. PRs 2β4 are independent and proceed.OFF-RAMP (PR 3): if error identity or the
tests/ipc/ACP paths cannot be preserved, keepinitialize()as is and record why.OFF-RAMP (PR 2): #3949 not merged β hold (see PR 2).
5. Risk register
commitPartial'shistorySequenceupdate-or-append is idempotent (historyService.ts:2036β2041)errorchunk β error path instead of abort pathboundedTeardownalready bounds; transcript measures; no budget change possible (5 s outer)stream-abortlistener (commitPartial) still in flight whenprocess.exitrunspartial.jsonis already gone whencli/server.tslogs its final cleanup line beforeprocess.exit(0)(:252β266)createStreamResult,tokenTracker)cause:conflicts onlyforkIn/onInterrupt/toAsyncIterableWith)streamManager.ts/streamBridge.ts; pins fixed; GA upgrade is a separate lockstep PR (Β§6)Reflect.set(processStreamWithCleanup)runSyncsynchronousshutdown()(desktop secondbefore-quitlistener) still does not await streamsshutdown()never touches the runtime; desktop's dispose race is the covered pathforkInon closed scope interrupts immediately βsystemabort (startImmediatelysemantics verified in rc.112; pinned by test)systemabort triggers an in-session RetryManager retry during shutdown"aborted"βNON_RETRYABLE_STREAM_ERRORSβretryManager.ts:99β104abandons; no fiber scheduledtaskService.initializemid-editConfigwhen the root exits after a timeoutcli/server.tsdispose; config writes are lock/journal-protectedstartImmediately/onInterruptsemantics differ in a later RCappFiberScope.test.tsstyle; RC bumps are a separate lockstep PR6. Standing item β effect v4 GA +
@orpc/experimental-effectlockstep (analysis only)v4 is not GA (rc.112 is current; v3
3.xremains the stable line). No PR this wave. When GA ships: one lockstep PR bumpingeffect+ all@orpc/*(1.14.11today; check the GA-compatible@orpc/experimental-effect), canary gates =di/*.test.ts,streamBridge.test.ts,streamManager.test.ts,serviceContainer.test.ts,TEST_INTEGRATION=1 bun x jest tests/ipc,make static-check. TheContext β ServiceMaprename risk is firewalled:Context.Servicetags,Context.omit/get,Layer,ManagedRuntime,TestClocklive only underdi/+orpc/effectContext.ts;streamManager.ts/streamBridge.tsuseEffect/Scope/Fiber/Exit/Stream/Queue/Causeonly. PR 4 adds oneContext.Context<never>type reference tostreamBridge.tsβ keep it as a type-only import so a rename is a one-line fix.7. Dogfooding (per PR; evidence attached to the PR with
gh β¦ --attach)Common setup:
make dev-server-sandbox DEV_SERVER_SANDBOX_ARGS="--clean-projects"(orxum serveron a tempXUM_ROOT) withXUM_LOG_LEVEL=debug, run underscript -f ~/wave4-scratch/<pr>-<scenario>.log; scratch under$HOME/wave4-scratch/(never/tmp). Drive the UI withagent-browser(openβsnapshot -iβ click the explicit "Send message" ref; re-snapshot after typing). Screenshots are primary evidence; record WebM and finalize withffmpeg -c copy.kill -TERM <server pid>; transcript must show the order in STOP Better authentication UXΒ #1 and[shutdown] AppFiberScope closed { ms }; (2)ls <XUM_ROOT>/sessions/<ws>/partial.json(absent) +tail -n 1 chat.jsonl(interrupted assistant message); (3) restart, open the workspace in agent-browser, screenshot the persisted interrupted message; (4) same scenario onmainfor the baseline diff; (5)xum runCtrl-C mid-stream transcript (CLI root parity); (6) quality gate between phases: gate suites green before the sandbox run, sandbox evidence before requesting review.taskkind=workspace to a child that arms a background bash monitor firing within ~10 s and keeps working ~60 s; screenshot the parent's task result (baselinemain:interrupted β¦ uncorrelated workspace stream-end; after:completed) and the child's log lines.xum servercold start transcript with[startup] <step> { ms }for the six steps (parity); throwaway worktree build withSTARTUP_STEP_TIMEOUT_MS = 1β transcript ofStartup failed: StartupStepTimeoutError β¦and exit code (do not ship); desktop dialog cannot be shown headless β cite the unchangeddesktop/main.ts:1249β1265catch.streamBridge.test.tson TestClock.8. Non-goals (restated; out of this wave)
Typed-error propagation sweep / removing the ~113 facades; converting services to
yield*-based Effect services; PubSub for the internal EventEmitter bus; Schema at persistence boundaries; Effect observability; converting sync read paths, AI-SDK per-request callbacks, cross-process lock interiors, or deterministic try-lock funnels; replacingAbortControlleras the SDK cancellation transport; converting thefullStreamloop toStream; fiberizing turn-handle waiters; downgrading startup steps to best-effort; changing outer quit budgets; supervising the pre-registration stream-start window;shutdown()semantics; the effect GA bump (standing analysis only).Generated with
xumβ’ Model:anthropic:claude-fable-5-1β’ Thinking:xhighβ’ Cost:$34.90