From a506661033da19f3abc693891475ed3674ce5e74 Mon Sep 17 00:00:00 2001 From: Lars Nieuwenhuis <35393046+lnieuwenhuis@users.noreply.github.com> Date: Sat, 5 Sep 2026 00:44:50 +0200 Subject: [PATCH 1/7] fix(client-runtime): stop resubscribing threads the server reports missing Subscribe failures carrying threadDisposition not-found now end the subscription terminally, tombstone the thread so foreground/probe wakeups never resubscribe, and drain queued persistence before cache removal so a debounced write cannot resurrect the deleted thread. --- apps/server/src/ws.ts | 1 + .../src/errors/orchestration.test.ts | 60 ++++++++- .../src/errors/orchestration.ts | 12 +- .../client-runtime/src/rpc/client.test.ts | 98 ++++++++++++++ packages/client-runtime/src/rpc/client.ts | 27 +++- .../src/state/threads-sync.test.ts | 123 ++++++++++++++++++ packages/client-runtime/src/state/threads.ts | 44 ++++++- packages/contracts/src/orchestration.ts | 1 + 8 files changed, 352 insertions(+), 14 deletions(-) diff --git a/apps/server/src/ws.ts b/apps/server/src/ws.ts index e1aff4a678c1..86006ed0475c 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -1702,6 +1702,7 @@ const makeWsRpcLayer = ( return yield* new OrchestrationGetSnapshotError({ message: `Thread ${input.threadId} was not found`, cause: input.threadId, + threadDisposition: "not-found", }); } diff --git a/packages/client-runtime/src/errors/orchestration.test.ts b/packages/client-runtime/src/errors/orchestration.test.ts index 0e1a723d5471..539a0d70d0d5 100644 --- a/packages/client-runtime/src/errors/orchestration.test.ts +++ b/packages/client-runtime/src/errors/orchestration.test.ts @@ -1,7 +1,55 @@ -import { OrchestrationDispatchCommandError } from "@t3tools/contracts"; -import { describe, expect, it } from "vite-plus/test"; +import { + OrchestrationDispatchCommandError, + OrchestrationGetSnapshotError, +} from "@t3tools/contracts"; +import { describe, expect, it } from "@effect/vitest"; -import { wasBootstrapThreadDeleted } from "./orchestration.ts"; +import { wasBootstrapThreadDeleted, wasSubscribeThreadNotFound } from "./orchestration.ts"; + +describe("wasSubscribeThreadNotFound", () => { + it("matches the typed not-found error", () => { + expect( + wasSubscribeThreadNotFound( + new OrchestrationGetSnapshotError({ + message: "Thread thread-1 was not found", + cause: "thread-1", + threadDisposition: "not-found", + }), + ), + ).toBe(true); + }); + + it("rejects a missing disposition", () => { + expect( + wasSubscribeThreadNotFound( + new OrchestrationGetSnapshotError({ + message: "Thread thread-1 was not found", + cause: "thread-1", + }), + ), + ).toBe(false); + }); + + it("rejects plain errors with a matching message", () => { + expect(wasSubscribeThreadNotFound(new Error("Thread thread-1 was not found"))).toBe(false); + }); + + it("rejects other snapshot errors", () => { + expect( + wasSubscribeThreadNotFound( + new OrchestrationGetSnapshotError({ + message: "Failed to load thread thread-1", + cause: "thread-1", + }), + ), + ).toBe(false); + }); + + it("rejects unrelated errors", () => { + expect(wasSubscribeThreadNotFound(new Error("boom"))).toBe(false); + expect(wasSubscribeThreadNotFound(null)).toBe(false); + }); +}); describe("wasBootstrapThreadDeleted", () => { it("accepts only a confirmed deleted bootstrap thread", () => { @@ -13,11 +61,17 @@ describe("wasBootstrapThreadDeleted", () => { }), ), ).toBe(true); + }); + + it("rejects a missing disposition", () => { expect( wasBootstrapThreadDeleted( new OrchestrationDispatchCommandError({ message: "Failed to create worktree." }), ), ).toBe(false); + }); + + it("rejects unrelated errors", () => { expect(wasBootstrapThreadDeleted(new Error("connection lost"))).toBe(false); }); }); diff --git a/packages/client-runtime/src/errors/orchestration.ts b/packages/client-runtime/src/errors/orchestration.ts index 39ef26e46576..521d046beebf 100644 --- a/packages/client-runtime/src/errors/orchestration.ts +++ b/packages/client-runtime/src/errors/orchestration.ts @@ -1,4 +1,7 @@ -import { OrchestrationDispatchCommandError } from "@t3tools/contracts"; +import { + OrchestrationDispatchCommandError, + OrchestrationGetSnapshotError, +} from "@t3tools/contracts"; import * as Schema from "effect/Schema"; const isOrchestrationDispatchCommandError = Schema.is(OrchestrationDispatchCommandError); @@ -8,3 +11,10 @@ export function wasBootstrapThreadDeleted(error: unknown): boolean { isOrchestrationDispatchCommandError(error) && error.bootstrapThreadDisposition === "deleted" ); } + +const isOrchestrationGetSnapshotError = Schema.is(OrchestrationGetSnapshotError); + +/** Set by the server when a subscribeThread miss has no snapshot (`apps/server/src/ws.ts`). */ +export function wasSubscribeThreadNotFound(error: unknown): boolean { + return isOrchestrationGetSnapshotError(error) && error.threadDisposition === "not-found"; +} diff --git a/packages/client-runtime/src/rpc/client.test.ts b/packages/client-runtime/src/rpc/client.test.ts index 4e6baba8bef4..917a3beeddd7 100644 --- a/packages/client-runtime/src/rpc/client.test.ts +++ b/packages/client-runtime/src/rpc/client.test.ts @@ -391,6 +391,104 @@ describe("environment RPC", () => { }), ); + it.effect("ends the subscription permanently when the failure is terminal", () => + Effect.gen(function* () { + const notFound = new Error("thread was not found"); + const subscriptionCount = yield* Ref.make(0); + const handled = yield* Ref.make | null>(null); + const client = { + [WS_METHODS.subscribeTerminalEvents]: () => + Stream.unwrap( + Ref.updateAndGet(subscriptionCount, (count) => count + 1).pipe( + Effect.map(() => Stream.fail(notFound)), + ), + ), + } as unknown as WsRpcProtocolClient; + const { activeSession, supervisor } = yield* makeHarness(); + + yield* SubscriptionRef.set(activeSession, Option.some(session(client))); + const subscriptionFiber = yield* subscribe( + WS_METHODS.subscribeTerminalEvents, + {}, + { + terminalFailure: { + matches: (error) => error === notFound, + handle: (cause) => Ref.set(handled, cause), + }, + }, + ).pipe( + Stream.runDrain, + Effect.provideService(EnvironmentSupervisor.EnvironmentSupervisor, supervisor), + Effect.forkChild, + ); + for (let attempt = 0; attempt < 100 && (yield* Ref.get(handled)) === null; attempt += 1) { + yield* Effect.yieldNow; + } + + expect(yield* Ref.get(subscriptionCount)).toBe(1); + expect(yield* Ref.get(handled)).not.toBeNull(); + + // Far past any retry delay: a terminal failure must not re-attempt. + yield* TestClock.adjust("30 seconds"); + yield* Effect.yieldNow; + yield* Fiber.interrupt(subscriptionFiber); + + expect(yield* Ref.get(subscriptionCount)).toBe(1); + }), + ); + + it.effect("does not treat failures rejected by the terminal classifier as terminal", () => + Effect.gen(function* () { + const transient = new Error("transient snapshot failure"); + const subscriptionCount = yield* Ref.make(0); + const client = { + [WS_METHODS.subscribeTerminalEvents]: () => + Stream.unwrap( + Ref.getAndUpdate(subscriptionCount, (count) => count + 1).pipe( + Effect.map((count) => (count === 0 ? Stream.fail(transient) : Stream.never)), + ), + ), + } as unknown as WsRpcProtocolClient; + const { activeSession, supervisor } = yield* makeHarness(); + + yield* SubscriptionRef.set(activeSession, Option.some(session(client))); + const subscriptionFiber = yield* subscribe( + WS_METHODS.subscribeTerminalEvents, + {}, + { + onExpectedFailure: () => Effect.void, + retryExpectedFailureAfter: "100 millis", + terminalFailure: { + matches: () => false, + handle: () => Effect.void, + }, + }, + ).pipe( + Stream.runDrain, + Effect.provideService(EnvironmentSupervisor.EnvironmentSupervisor, supervisor), + Effect.forkChild, + ); + + // The retry sleep must be scheduled before virtual time advances. + for ( + let attempt = 0; + attempt < 100 && (yield* Ref.get(subscriptionCount)) < 1; + attempt += 1 + ) { + yield* Effect.yieldNow; + } + yield* TestClock.adjust("100 millis"); + for (let attempt = 0; attempt < 100; attempt += 1) { + if ((yield* Ref.get(subscriptionCount)) >= 2) break; + yield* Effect.yieldNow; + } + yield* Fiber.interrupt(subscriptionFiber); + + // Classified as a regular expected failure: retried once, handler untouched. + expect(yield* Ref.get(subscriptionCount)).toBe(2); + }), + ); + it.effect("does not classify subscription defects as expected failures", () => Effect.gen(function* () { const defect = new Error("subscription invariant failed"); diff --git a/packages/client-runtime/src/rpc/client.ts b/packages/client-runtime/src/rpc/client.ts index 0d68d2b2d531..ddf27ff3aaa8 100644 --- a/packages/client-runtime/src/rpc/client.ts +++ b/packages/client-runtime/src/rpc/client.ts @@ -176,6 +176,15 @@ interface SubscriptionOptions { ) => Effect.Effect; readonly retryExpectedFailureAfter?: Duration.Input; readonly resubscribe?: Stream.Stream; + /** + * Classifies an all-Fail cause as terminal: the attempt ends for this + * session with no retry, after `handle` runs. Checked after transport + * failures and before expected-failure retry. + */ + readonly terminalFailure?: { + readonly matches: (error: EnvironmentRpcStreamFailure) => boolean; + readonly handle: (cause: Cause.Cause>) => Effect.Effect; + }; } export function subscribeDynamic( @@ -232,14 +241,19 @@ export function subscribeDynamic( return method(input).pipe( Stream.ensuring(completeObservation), Stream.catchCause((cause) => { + const failErrors = cause.reasons.flatMap((reason) => + reason._tag === "Fail" ? [reason.error] : [], + ); const hasOnlyExpectedFailures = - cause.reasons.length > 0 && - cause.reasons.every((reason) => reason._tag === "Fail"); + cause.reasons.length > 0 && failErrors.length === cause.reasons.length; const isTransportFailure = hasOnlyExpectedFailures && - cause.reasons.every( - (reason) => reason._tag === "Fail" && isRpcClientError(reason.error), - ); + failErrors.every((error) => isRpcClientError(error)); + const terminal = options?.terminalFailure; + const isTerminal = + hasOnlyExpectedFailures && + terminal !== undefined && + failErrors.every((error) => terminal.matches(error)); if (isTransportFailure) { return Stream.fromEffect( Effect.logWarning( @@ -252,6 +266,9 @@ export function subscribeDynamic( ), ).pipe(Stream.drain); } + if (isTerminal && terminal !== undefined) { + return Stream.fromEffect(terminal.handle(cause)).pipe(Stream.drain); + } if (hasOnlyExpectedFailures && options?.onExpectedFailure !== undefined) { const handled = Stream.fromEffect( options.onExpectedFailure(cause), diff --git a/packages/client-runtime/src/state/threads-sync.test.ts b/packages/client-runtime/src/state/threads-sync.test.ts index 619b5f43425d..238c97593fcb 100644 --- a/packages/client-runtime/src/state/threads-sync.test.ts +++ b/packages/client-runtime/src/state/threads-sync.test.ts @@ -1,6 +1,7 @@ import { EnvironmentId, EventId, + OrchestrationGetSnapshotError, ORCHESTRATION_WS_METHODS, ProjectId, ProviderInstanceId, @@ -988,4 +989,126 @@ describe("EnvironmentThreads", () => { expect(yield* Ref.get(harness.subscriptionCount)).toBe(3); }), ); + + it.effect("marks the thread deleted and stops subscribing when the thread is missing", () => + Effect.gen(function* () { + const harness = yield* makeHarness(); + yield* Queue.offer( + harness.inputs, + new OrchestrationGetSnapshotError({ + message: `Thread ${THREAD_ID} was not found`, + cause: THREAD_ID, + threadDisposition: "not-found", + }), + ); + + const state = yield* awaitThreadState( + harness.observed, + (value) => value.status === "deleted", + ); + expect(Option.isNone(state.data)).toBe(true); + expect(Option.isNone(state.error)).toBe(true); + // Deletion parity with the normal thread.deleted event path. + expect(yield* Ref.get(harness.removedThreads)).toEqual([THREAD_ID]); + + // No retry storm: far past any retry delay there is exactly one + // subscribe attempt for a deleted thread. + yield* TestClock.adjust("30 seconds"); + yield* Effect.yieldNow; + expect(yield* Ref.get(harness.subscriptionCount)).toBe(1); + }), + ); + + it.effect("does not resubscribe a missing thread on foreground wakeups", () => + Effect.gen(function* () { + const harness = yield* makeHarness({ cached: BASE_THREAD }); + yield* Queue.offer( + harness.inputs, + new OrchestrationGetSnapshotError({ + message: `Thread ${THREAD_ID} was not found`, + cause: THREAD_ID, + threadDisposition: "not-found", + }), + ); + yield* awaitThreadState(harness.observed, (value) => value.status === "deleted"); + + // Outer resubscribe triggers (foreground, probe) must stay gated by the + // tombstone: the terminal attempt remains the only subscribe attempt. + yield* Queue.offer(harness.wakeups, "application-active"); + yield* Queue.offer(harness.wakeups, "application-active-probe"); + yield* TestClock.adjust("30 seconds"); + for (let attempt = 0; attempt < 100; attempt += 1) { + yield* Effect.yieldNow; + } + + expect(yield* Ref.get(harness.subscriptionCount)).toBe(1); + const latest = yield* Ref.get(harness.latest); + expect(latest.status).toBe("deleted"); + expect(Option.isNone(latest.data)).toBe(true); + expect(yield* Ref.get(harness.loaderCalls)).toBe(0); + }), + ); + + it.effect("does not resurrect a missing thread via queued persistence", () => + Effect.gen(function* () { + const harness = yield* makeHarness({ cached: BASE_THREAD }); + // Queue a debounced cache write, then delete before it flushes. + yield* Queue.offer(harness.inputs, snapshot(BASE_THREAD)); + yield* Queue.offer(harness.inputs, titleUpdated("Live title", 2)); + yield* Queue.offer( + harness.inputs, + new OrchestrationGetSnapshotError({ + message: `Thread ${THREAD_ID} was not found`, + cause: THREAD_ID, + threadDisposition: "not-found", + }), + ); + + const state = yield* awaitThreadState( + harness.observed, + (value) => value.status === "deleted", + ); + expect(Option.isNone(state.data)).toBe(true); + + // Flush far past the persistence debounce: the stale write must not + // re-save the thread after its cache entry was removed. + yield* TestClock.adjust("30 seconds"); + for (let attempt = 0; attempt < 100; attempt += 1) { + yield* Effect.yieldNow; + } + + expect(yield* Ref.get(harness.removedThreads)).toEqual([THREAD_ID]); + expect(yield* Ref.get(harness.savedThreads)).toEqual([]); + expect(yield* Ref.get(harness.subscriptionCount)).toBe(1); + }), + ); + + it.effect("retries a snapshot error without the not-found disposition", () => + Effect.gen(function* () { + const harness = yield* makeHarness(); + yield* Queue.offer( + harness.inputs, + new OrchestrationGetSnapshotError({ + message: `Failed to load thread ${THREAD_ID}`, + cause: THREAD_ID, + }), + ); + + const failed = yield* awaitThreadState(harness.observed, (value) => + Option.isSome(value.error), + ); + expect(failed.status).not.toBe("deleted"); + expect(Option.getOrThrow(failed.error)).toBe(`Failed to load thread ${THREAD_ID}`); + expect(yield* Ref.get(harness.subscriptionCount)).toBe(1); + + yield* TestClock.adjust("250 millis"); + for (let attempt = 0; attempt < 100; attempt += 1) { + if ((yield* Ref.get(harness.subscriptionCount)) >= 2) break; + yield* Effect.yieldNow; + } + + expect(yield* Ref.get(harness.subscriptionCount)).toBe(2); + expect(yield* Ref.get(harness.removedThreads)).toEqual([]); + }), + ); }); diff --git a/packages/client-runtime/src/state/threads.ts b/packages/client-runtime/src/state/threads.ts index a0055b6cab3c..752f4bf06d4a 100644 --- a/packages/client-runtime/src/state/threads.ts +++ b/packages/client-runtime/src/state/threads.ts @@ -22,6 +22,7 @@ import { EnvironmentRegistry } from "../connection/registry.ts"; import { connectionProjectionPhase } from "../connection/model.ts"; import { EnvironmentSupervisor } from "../connection/supervisor.ts"; import * as ConnectionWakeups from "../connection/wakeups.ts"; +import { wasSubscribeThreadNotFound } from "../errors/orchestration.ts"; import { EnvironmentCacheStore } from "../platform/persistence.ts"; import { subscribeDynamic } from "../rpc/client.ts"; import { ThreadSnapshotLoader, type ThreadSnapshotWindow } from "./threadSnapshotHttp.ts"; @@ -258,10 +259,19 @@ export const makeEnvironmentThreadState = Effect.fn("EnvironmentThreadState.make readonly epoch: number; } | null>(null); const persistence = yield* Queue.sliding(1); + // Latch set when the server reports the thread missing (subscribe fails + // with threadDisposition "not-found"). Gates the outer foreground/probe + // resubscribe path below so a deleted thread stops after its single + // terminal attempt instead of retrying on every app wakeup. + const terminalNotFound = yield* Ref.make(false); const persist = Effect.fn("EnvironmentThreadState.persist")(function* ( snapshot: OrchestrationThreadDetailSnapshot, ) { + // A debounced write racing a deletion must not resurrect the cache after + // setDeleted removes it; the queue drain in setDeleted drops pending + // offers, this guard drops the offer already past the debounce. + if ((yield* SubscriptionRef.get(state)).status === "deleted") return; if (resumeCache !== undefined && resumeCache.owner !== owner) return; if ( committed.persisted && @@ -394,6 +404,10 @@ export const makeEnvironmentThreadState = Effect.fn("EnvironmentThreadState.make page: Option.none(), }); yield* remember; + // Drop the queued cache write before removing the entry so a debounced + // snapshot cannot re-save the thread after deletion. The queue is + // sliding(1), so a single poll discards the whole backlog. + yield* Queue.poll(persistence); if (resumeCache !== undefined && resumeCache.owner !== owner) return; yield* cache.removeThread(environmentId, threadId).pipe( Effect.catch((error) => @@ -644,7 +658,12 @@ export const makeEnvironmentThreadState = Effect.fn("EnvironmentThreadState.make const foregroundResubscriptions = Option.match(wakeups, { onNone: () => Stream.never, onSome: (service) => - service.changes.pipe(Stream.filter(ConnectionWakeups.shouldResubscribeAfterWakeup)), + service.changes.pipe( + Stream.filter(ConnectionWakeups.shouldResubscribeAfterWakeup), + Stream.filterEffect(() => + Ref.get(terminalNotFound).pipe(Effect.map((missing) => !missing)), + ), + ), }); yield* setSynchronizing; @@ -742,6 +761,19 @@ export const makeEnvironmentThreadState = Effect.fn("EnvironmentThreadState.make onExpectedFailure: setStreamError, retryExpectedFailureAfter: "250 millis", resubscribe: foregroundResubscriptions, + terminalFailure: { + matches: wasSubscribeThreadNotFound, + handle: (cause) => + Ref.set(terminalNotFound, true).pipe( + Effect.andThen( + Effect.logWarning("Subscribed thread is gone; stopping.", { + threadId, + cause: Cause.pretty(cause), + }), + ), + Effect.andThen(setDeleted()), + ), + }, }, ).pipe(Stream.runForEach(applyItem)), ); @@ -816,10 +848,12 @@ export function createEnvironmentThreadStateAtoms( // Cache definitions must outlive collectible live-atom definitions. The // registry retains these nodes without retaining environment or RPC scopes. const resumeFamily = Atom.family((key: string) => - Atom.make((): ThreadResumeCache => ({ - snapshot: undefined, - owner: undefined, - })).pipe( + Atom.make( + (): ThreadResumeCache => ({ + snapshot: undefined, + owner: undefined, + }), + ).pipe( Atom.setIdleTTL(THREAD_SNAPSHOT_IDLE_TTL_MS), Atom.withLabel(`environment-thread-resume:${key}`), ), diff --git a/packages/contracts/src/orchestration.ts b/packages/contracts/src/orchestration.ts index 17cadc6d1d7f..5727b9987ff5 100644 --- a/packages/contracts/src/orchestration.ts +++ b/packages/contracts/src/orchestration.ts @@ -1852,6 +1852,7 @@ export class OrchestrationGetSnapshotError extends Schema.TaggedErrorClass Date: Sat, 5 Sep 2026 01:00:29 +0200 Subject: [PATCH 2/7] fix(client-runtime): terminate durable subscription on terminal thread miss Session replacements re-issued subscribeThread after a not-found tombstone because the terminal latch only filtered foreground wakeups while the outer session stream in subscribeDynamic stayed alive. Signal a halt Deferred from the terminalFailure handler and interrupt the outer session stream so no new subscribe issues; non-matching failures keep session-driven resubscription. --- packages/client-runtime/src/rpc/client.ts | 16 +++- .../src/state/threads-sync.test.ts | 76 +++++++++++++++++++ 2 files changed, 91 insertions(+), 1 deletion(-) diff --git a/packages/client-runtime/src/rpc/client.ts b/packages/client-runtime/src/rpc/client.ts index ddf27ff3aaa8..d465a24431b9 100644 --- a/packages/client-runtime/src/rpc/client.ts +++ b/packages/client-runtime/src/rpc/client.ts @@ -1,6 +1,7 @@ import { ORCHESTRATION_WS_METHODS, WS_METHODS } from "@t3tools/contracts"; import * as Cause from "effect/Cause"; import * as Context from "effect/Context"; +import * as Deferred from "effect/Deferred"; import type * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; import * as Option from "effect/Option"; @@ -200,6 +201,10 @@ export function subscribeDynamic( Effect.gen(function* () { const supervisor = yield* EnvironmentSupervisor; const observer = yield* EnvironmentRpcSubscriptionObserver; + // Signaled once a terminalFailure handler runs. Interrupts the outer + // session stream below so a later supervisor.session replacement cannot + // re-issue the subscription after a terminal tombstone. + const terminalHalt = yield* Deferred.make(); const sessionChanges = SubscriptionRef.changes(supervisor.session); const sessions = options?.resubscribe === undefined @@ -211,6 +216,7 @@ export function subscribeDynamic( ), ); return sessions.pipe( + Stream.interruptWhen(Deferred.await(terminalHalt)), Stream.switchMap( Option.match({ onNone: () => Stream.empty, @@ -267,7 +273,15 @@ export function subscribeDynamic( ).pipe(Stream.drain); } if (isTerminal && terminal !== undefined) { - return Stream.fromEffect(terminal.handle(cause)).pipe(Stream.drain); + return Stream.fromEffect( + terminal + .handle(cause) + .pipe( + Effect.ensuring( + Deferred.succeed(terminalHalt, undefined).pipe(Effect.asVoid), + ), + ), + ).pipe(Stream.drain); } if (hasOnlyExpectedFailures && options?.onExpectedFailure !== undefined) { const handled = Stream.fromEffect( diff --git a/packages/client-runtime/src/state/threads-sync.test.ts b/packages/client-runtime/src/state/threads-sync.test.ts index 238c97593fcb..40fbbea7d116 100644 --- a/packages/client-runtime/src/state/threads-sync.test.ts +++ b/packages/client-runtime/src/state/threads-sync.test.ts @@ -1049,6 +1049,82 @@ describe("EnvironmentThreads", () => { }), ); + it.effect("does not resubscribe a missing thread on session replacement", () => + Effect.gen(function* () { + const harness = yield* makeHarness({ cached: BASE_THREAD }); + yield* Queue.offer( + harness.inputs, + new OrchestrationGetSnapshotError({ + message: `Thread ${THREAD_ID} was not found`, + cause: THREAD_ID, + threadDisposition: "not-found", + }), + ); + yield* awaitThreadState(harness.observed, (value) => value.status === "deleted"); + expect(yield* Ref.get(harness.subscriptionCount)).toBe(1); + + // A supervisor.session replacement drives the outer session stream in + // subscribeDynamic. The terminal tombstone must terminate that path too, + // not just foreground wakeups. + yield* harness.replaceSession; + yield* TestClock.adjust("30 seconds"); + for (let attempt = 0; attempt < 100; attempt += 1) { + yield* Effect.yieldNow; + } + + expect(yield* Ref.get(harness.subscriptionCount)).toBe(1); + const latest = yield* Ref.get(harness.latest); + expect(latest.status).toBe("deleted"); + expect(Option.isNone(latest.data)).toBe(true); + expect(yield* Ref.get(harness.loaderCalls)).toBe(0); + }), + ); + + it.effect("resubscribes on session replacement after a non-matching failure", () => + Effect.gen(function* () { + const harness = yield* makeHarness(); + yield* Queue.offer( + harness.inputs, + new OrchestrationGetSnapshotError({ + message: `Failed to load thread ${THREAD_ID}`, + cause: THREAD_ID, + }), + ); + + const failed = yield* awaitThreadState(harness.observed, (value) => + Option.isSome(value.error), + ); + expect(failed.status).not.toBe("deleted"); + expect(yield* Ref.get(harness.subscriptionCount)).toBe(1); + + // Non-terminal failures keep the session-driven path alive: a + // replacement session must re-issue subscribeThread. + yield* harness.replaceSession; + for (let attempt = 0; attempt < 100; attempt += 1) { + if ((yield* Ref.get(harness.subscriptionCount)) >= 2) break; + yield* Effect.yieldNow; + } + expect(yield* Ref.get(harness.subscriptionCount)).toBe(2); + + yield* Queue.offer( + harness.inputs, + snapshot({ + ...BASE_THREAD, + title: "Recovered thread", + }), + ); + const recovered = yield* awaitThreadState( + harness.observed, + (value) => + value.status === "live" && + Option.isSome(value.data) && + value.data.value.title === "Recovered thread", + ); + expect(Option.isNone(recovered.error)).toBe(true); + expect(yield* Ref.get(harness.removedThreads)).toEqual([]); + }), + ); + it.effect("does not resurrect a missing thread via queued persistence", () => Effect.gen(function* () { const harness = yield* makeHarness({ cached: BASE_THREAD }); From a85cd6d75cd291876c28088b437df9a578e0e086 Mon Sep 17 00:00:00 2001 From: Lars Nieuwenhuis <35393046+lnieuwenhuis@users.noreply.github.com> Date: Sat, 5 Sep 2026 01:07:41 +0200 Subject: [PATCH 3/7] fix(client-runtime): use named resume-cache factory for formatter parity --- packages/client-runtime/src/state/threads.ts | 11 +++++------ 1 file changed, 5 insertions(+), 6 deletions(-) diff --git a/packages/client-runtime/src/state/threads.ts b/packages/client-runtime/src/state/threads.ts index 752f4bf06d4a..9ebdd554f1cf 100644 --- a/packages/client-runtime/src/state/threads.ts +++ b/packages/client-runtime/src/state/threads.ts @@ -847,13 +847,12 @@ export function createEnvironmentThreadStateAtoms( ) { // Cache definitions must outlive collectible live-atom definitions. The // registry retains these nodes without retaining environment or RPC scopes. + const makeThreadResumeCache = (): ThreadResumeCache => ({ + snapshot: undefined, + owner: undefined, + }); const resumeFamily = Atom.family((key: string) => - Atom.make( - (): ThreadResumeCache => ({ - snapshot: undefined, - owner: undefined, - }), - ).pipe( + Atom.make(makeThreadResumeCache).pipe( Atom.setIdleTTL(THREAD_SNAPSHOT_IDLE_TTL_MS), Atom.withLabel(`environment-thread-resume:${key}`), ), From 07a8db4b9fe78b359bb8183eb18f94e6cb2e765e Mon Sep 17 00:00:00 2001 From: Lars Nieuwenhuis <35393046+lnieuwenhuis@users.noreply.github.com> Date: Sat, 5 Sep 2026 01:15:00 +0200 Subject: [PATCH 4/7] fix(client-runtime): signal terminal halt before running the handler A session replacement landing during the terminal handler's cache I/O started a new inner subscribe before the post-handle halt landed. Signal terminalHalt first so the outer session stream is already dead; the handler still drains as the running inner. Muse Spark (opencode) --- packages/client-runtime/src/rpc/client.ts | 24 +++++---- .../src/state/threads-sync.test.ts | 49 ++++++++++++++++++- 2 files changed, 62 insertions(+), 11 deletions(-) diff --git a/packages/client-runtime/src/rpc/client.ts b/packages/client-runtime/src/rpc/client.ts index d465a24431b9..83722f10a23e 100644 --- a/packages/client-runtime/src/rpc/client.ts +++ b/packages/client-runtime/src/rpc/client.ts @@ -201,9 +201,11 @@ export function subscribeDynamic( Effect.gen(function* () { const supervisor = yield* EnvironmentSupervisor; const observer = yield* EnvironmentRpcSubscriptionObserver; - // Signaled once a terminalFailure handler runs. Interrupts the outer - // session stream below so a later supervisor.session replacement cannot - // re-issue the subscription after a terminal tombstone. + // Signaled before a terminalFailure handler runs. Interrupts the + // outer session stream below so a supervisor.session replacement + // landing during the handler's cache I/O cannot re-issue the + // subscription after a terminal tombstone; the handler itself still + // runs to completion as the draining inner. const terminalHalt = yield* Deferred.make(); const sessionChanges = SubscriptionRef.changes(supervisor.session); const sessions = @@ -273,14 +275,16 @@ export function subscribeDynamic( ).pipe(Stream.drain); } if (isTerminal && terminal !== undefined) { + // Halt first: handle performs cache I/O, and a + // session replacement in that window must not + // start a new inner subscribe before the interrupt + // lands. The handler still runs to completion as + // the draining inner. return Stream.fromEffect( - terminal - .handle(cause) - .pipe( - Effect.ensuring( - Deferred.succeed(terminalHalt, undefined).pipe(Effect.asVoid), - ), - ), + Deferred.succeed(terminalHalt, undefined).pipe( + Effect.asVoid, + Effect.andThen(terminal.handle(cause)), + ), ).pipe(Stream.drain); } if (hasOnlyExpectedFailures && options?.onExpectedFailure !== undefined) { diff --git a/packages/client-runtime/src/state/threads-sync.test.ts b/packages/client-runtime/src/state/threads-sync.test.ts index 40fbbea7d116..de59aef6ae32 100644 --- a/packages/client-runtime/src/state/threads-sync.test.ts +++ b/packages/client-runtime/src/state/threads-sync.test.ts @@ -142,6 +142,7 @@ const makeHarness = Effect.fn("TestEnvironmentThreads.makeHarness")(function* (o readonly resumeCache?: NonNullable[1]>; readonly loadCached?: Effect.Effect>; readonly saveThread?: Persistence.EnvironmentCacheStore["Service"]["saveThread"]; + readonly onRemoveThread?: Effect.Effect; }) { const inputs = yield* Queue.unbounded(); const observed = yield* Queue.unbounded(); @@ -224,7 +225,9 @@ const makeHarness = Effect.fn("TestEnvironmentThreads.makeHarness")(function* (o Effect.andThen(options?.saveThread?.(environmentId, thread) ?? Effect.void), ), removeThread: (_environmentId, threadId) => - Ref.update(removedThreads, (current) => [...current, threadId]), + (options?.onRemoveThread ?? Effect.void).pipe( + Effect.andThen(Ref.update(removedThreads, (current) => [...current, threadId])), + ), loadServerConfig: () => Effect.succeed(Option.none()), saveServerConfig: () => Effect.void, loadVcsRefs: () => Effect.succeed(Option.none()), @@ -1080,6 +1083,50 @@ describe("EnvironmentThreads", () => { }), ); + it.effect("does not resubscribe when the session is replaced mid-handle", () => + Effect.gen(function* () { + const entered = yield* Deferred.make(); + const release = yield* Deferred.make(); + const harness = yield* makeHarness({ + cached: BASE_THREAD, + onRemoveThread: Deferred.succeed(entered, undefined).pipe( + Effect.asVoid, + Effect.andThen(Deferred.await(release)), + ), + }); + yield* Queue.offer( + harness.inputs, + new OrchestrationGetSnapshotError({ + message: `Thread ${THREAD_ID} was not found`, + cause: THREAD_ID, + threadDisposition: "not-found", + }), + ); + + // Block the terminal handler inside its cache I/O, then replace the + // session mid-handle. The halt must already have fired, so the dead + // outer session stream cannot start a second subscribe. Deterministic: + // Deferred gates only, no sleeps. + yield* Deferred.await(entered); + yield* harness.replaceSession; + for (let attempt = 0; attempt < 100; attempt += 1) { + if ((yield* Ref.get(harness.subscriptionCount)) > 1) break; + yield* Effect.yieldNow; + } + expect(yield* Ref.get(harness.subscriptionCount)).toBe(1); + + yield* Deferred.succeed(release, undefined); + yield* awaitThreadState(harness.observed, (value) => value.status === "deleted"); + + expect(yield* Ref.get(harness.subscriptionCount)).toBe(1); + const latest = yield* Ref.get(harness.latest); + expect(latest.status).toBe("deleted"); + expect(Option.isNone(latest.data)).toBe(true); + expect(yield* Ref.get(harness.removedThreads)).toEqual([THREAD_ID]); + expect(yield* Ref.get(harness.loaderCalls)).toBe(0); + }), + ); + it.effect("resubscribes on session replacement after a non-matching failure", () => Effect.gen(function* () { const harness = yield* makeHarness(); From 8d74e120d3fea488046b1a932cf4a05ada97e9be Mon Sep 17 00:00:00 2001 From: Lars Nieuwenhuis <35393046+lnieuwenhuis@users.noreply.github.com> Date: Sat, 5 Sep 2026 08:08:55 +0200 Subject: [PATCH 5/7] fix(contracts): distinguish missing thread subscription errors --- apps/server/src/server.test.ts | 3 +- apps/server/src/ws.ts | 7 ++-- .../src/errors/orchestration.test.ts | 28 ++++++++-------- .../src/errors/orchestration.ts | 9 ++--- .../src/state/threads-sync.test.ts | 33 ++++--------------- packages/client-runtime/src/state/threads.ts | 6 ++-- packages/contracts/src/orchestration.ts | 10 +++++- packages/contracts/src/rpc.test.ts | 17 +++++++++- packages/contracts/src/rpc.ts | 7 +++- 9 files changed, 61 insertions(+), 59 deletions(-) diff --git a/apps/server/src/server.test.ts b/apps/server/src/server.test.ts index 92def1774e74..13d95cfbaa71 100644 --- a/apps/server/src/server.test.ts +++ b/apps/server/src/server.test.ts @@ -8610,7 +8610,8 @@ it.layer(NodeServices.layer)("server router seam", (it) => { ); if (oversized) { assertTrue(threadResult._tag === "Failure"); - assert.equal(threadResult.failure._tag, "OrchestrationGetSnapshotError"); + assertTrue(threadResult.failure._tag === "OrchestrationThreadNotFoundError"); + assert.equal(threadResult.failure.threadId, defaultThreadId); assert.equal( threadResult.failure.message, `Thread ${defaultThreadId} was not found`, diff --git a/apps/server/src/ws.ts b/apps/server/src/ws.ts index 86006ed0475c..5590e5c86bbe 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -35,6 +35,7 @@ import { type OrchestrationShellStreamItem, OrchestrationGetFullThreadDiffError, OrchestrationGetSnapshotError, + OrchestrationThreadNotFoundError, OrchestrationSearchThreadsError, OrchestrationGetTurnDiffError, ORCHESTRATION_WS_METHODS, @@ -1699,11 +1700,7 @@ const makeWsRpcLayer = ( if (replayOnMissingSnapshot !== undefined) { return replayOnMissingSnapshot; } - return yield* new OrchestrationGetSnapshotError({ - message: `Thread ${input.threadId} was not found`, - cause: input.threadId, - threadDisposition: "not-found", - }); + return yield* new OrchestrationThreadNotFoundError({ threadId: input.threadId }); } const afterSnapshot = diff --git a/packages/client-runtime/src/errors/orchestration.test.ts b/packages/client-runtime/src/errors/orchestration.test.ts index 539a0d70d0d5..476780374fe1 100644 --- a/packages/client-runtime/src/errors/orchestration.test.ts +++ b/packages/client-runtime/src/errors/orchestration.test.ts @@ -1,27 +1,25 @@ import { OrchestrationDispatchCommandError, OrchestrationGetSnapshotError, + OrchestrationThreadNotFoundError, + ThreadId, } from "@t3tools/contracts"; import { describe, expect, it } from "@effect/vitest"; -import { wasBootstrapThreadDeleted, wasSubscribeThreadNotFound } from "./orchestration.ts"; +import { wasBootstrapThreadDeleted, isOrchestrationThreadNotFoundError } from "./orchestration.ts"; -describe("wasSubscribeThreadNotFound", () => { +describe("isOrchestrationThreadNotFoundError", () => { it("matches the typed not-found error", () => { expect( - wasSubscribeThreadNotFound( - new OrchestrationGetSnapshotError({ - message: "Thread thread-1 was not found", - cause: "thread-1", - threadDisposition: "not-found", - }), + isOrchestrationThreadNotFoundError( + new OrchestrationThreadNotFoundError({ threadId: ThreadId.make("thread-1") }), ), ).toBe(true); }); - it("rejects a missing disposition", () => { + it("rejects a generic snapshot error with a matching message", () => { expect( - wasSubscribeThreadNotFound( + isOrchestrationThreadNotFoundError( new OrchestrationGetSnapshotError({ message: "Thread thread-1 was not found", cause: "thread-1", @@ -31,12 +29,14 @@ describe("wasSubscribeThreadNotFound", () => { }); it("rejects plain errors with a matching message", () => { - expect(wasSubscribeThreadNotFound(new Error("Thread thread-1 was not found"))).toBe(false); + expect(isOrchestrationThreadNotFoundError(new Error("Thread thread-1 was not found"))).toBe( + false, + ); }); it("rejects other snapshot errors", () => { expect( - wasSubscribeThreadNotFound( + isOrchestrationThreadNotFoundError( new OrchestrationGetSnapshotError({ message: "Failed to load thread thread-1", cause: "thread-1", @@ -46,8 +46,8 @@ describe("wasSubscribeThreadNotFound", () => { }); it("rejects unrelated errors", () => { - expect(wasSubscribeThreadNotFound(new Error("boom"))).toBe(false); - expect(wasSubscribeThreadNotFound(null)).toBe(false); + expect(isOrchestrationThreadNotFoundError(new Error("boom"))).toBe(false); + expect(isOrchestrationThreadNotFoundError(null)).toBe(false); }); }); diff --git a/packages/client-runtime/src/errors/orchestration.ts b/packages/client-runtime/src/errors/orchestration.ts index 521d046beebf..565e6bd041c6 100644 --- a/packages/client-runtime/src/errors/orchestration.ts +++ b/packages/client-runtime/src/errors/orchestration.ts @@ -1,6 +1,6 @@ import { OrchestrationDispatchCommandError, - OrchestrationGetSnapshotError, + OrchestrationThreadNotFoundError, } from "@t3tools/contracts"; import * as Schema from "effect/Schema"; @@ -12,9 +12,4 @@ export function wasBootstrapThreadDeleted(error: unknown): boolean { ); } -const isOrchestrationGetSnapshotError = Schema.is(OrchestrationGetSnapshotError); - -/** Set by the server when a subscribeThread miss has no snapshot (`apps/server/src/ws.ts`). */ -export function wasSubscribeThreadNotFound(error: unknown): boolean { - return isOrchestrationGetSnapshotError(error) && error.threadDisposition === "not-found"; -} +export const isOrchestrationThreadNotFoundError = Schema.is(OrchestrationThreadNotFoundError); diff --git a/packages/client-runtime/src/state/threads-sync.test.ts b/packages/client-runtime/src/state/threads-sync.test.ts index de59aef6ae32..988d736e6f6a 100644 --- a/packages/client-runtime/src/state/threads-sync.test.ts +++ b/packages/client-runtime/src/state/threads-sync.test.ts @@ -2,6 +2,7 @@ import { EnvironmentId, EventId, OrchestrationGetSnapshotError, + OrchestrationThreadNotFoundError, ORCHESTRATION_WS_METHODS, ProjectId, ProviderInstanceId, @@ -998,11 +999,7 @@ describe("EnvironmentThreads", () => { const harness = yield* makeHarness(); yield* Queue.offer( harness.inputs, - new OrchestrationGetSnapshotError({ - message: `Thread ${THREAD_ID} was not found`, - cause: THREAD_ID, - threadDisposition: "not-found", - }), + new OrchestrationThreadNotFoundError({ threadId: THREAD_ID }), ); const state = yield* awaitThreadState( @@ -1027,11 +1024,7 @@ describe("EnvironmentThreads", () => { const harness = yield* makeHarness({ cached: BASE_THREAD }); yield* Queue.offer( harness.inputs, - new OrchestrationGetSnapshotError({ - message: `Thread ${THREAD_ID} was not found`, - cause: THREAD_ID, - threadDisposition: "not-found", - }), + new OrchestrationThreadNotFoundError({ threadId: THREAD_ID }), ); yield* awaitThreadState(harness.observed, (value) => value.status === "deleted"); @@ -1057,11 +1050,7 @@ describe("EnvironmentThreads", () => { const harness = yield* makeHarness({ cached: BASE_THREAD }); yield* Queue.offer( harness.inputs, - new OrchestrationGetSnapshotError({ - message: `Thread ${THREAD_ID} was not found`, - cause: THREAD_ID, - threadDisposition: "not-found", - }), + new OrchestrationThreadNotFoundError({ threadId: THREAD_ID }), ); yield* awaitThreadState(harness.observed, (value) => value.status === "deleted"); expect(yield* Ref.get(harness.subscriptionCount)).toBe(1); @@ -1096,11 +1085,7 @@ describe("EnvironmentThreads", () => { }); yield* Queue.offer( harness.inputs, - new OrchestrationGetSnapshotError({ - message: `Thread ${THREAD_ID} was not found`, - cause: THREAD_ID, - threadDisposition: "not-found", - }), + new OrchestrationThreadNotFoundError({ threadId: THREAD_ID }), ); // Block the terminal handler inside its cache I/O, then replace the @@ -1180,11 +1165,7 @@ describe("EnvironmentThreads", () => { yield* Queue.offer(harness.inputs, titleUpdated("Live title", 2)); yield* Queue.offer( harness.inputs, - new OrchestrationGetSnapshotError({ - message: `Thread ${THREAD_ID} was not found`, - cause: THREAD_ID, - threadDisposition: "not-found", - }), + new OrchestrationThreadNotFoundError({ threadId: THREAD_ID }), ); const state = yield* awaitThreadState( @@ -1206,7 +1187,7 @@ describe("EnvironmentThreads", () => { }), ); - it.effect("retries a snapshot error without the not-found disposition", () => + it.effect("retries a generic snapshot error", () => Effect.gen(function* () { const harness = yield* makeHarness(); yield* Queue.offer( diff --git a/packages/client-runtime/src/state/threads.ts b/packages/client-runtime/src/state/threads.ts index 9ebdd554f1cf..faa6fc705777 100644 --- a/packages/client-runtime/src/state/threads.ts +++ b/packages/client-runtime/src/state/threads.ts @@ -22,7 +22,7 @@ import { EnvironmentRegistry } from "../connection/registry.ts"; import { connectionProjectionPhase } from "../connection/model.ts"; import { EnvironmentSupervisor } from "../connection/supervisor.ts"; import * as ConnectionWakeups from "../connection/wakeups.ts"; -import { wasSubscribeThreadNotFound } from "../errors/orchestration.ts"; +import { isOrchestrationThreadNotFoundError } from "../errors/orchestration.ts"; import { EnvironmentCacheStore } from "../platform/persistence.ts"; import { subscribeDynamic } from "../rpc/client.ts"; import { ThreadSnapshotLoader, type ThreadSnapshotWindow } from "./threadSnapshotHttp.ts"; @@ -260,7 +260,7 @@ export const makeEnvironmentThreadState = Effect.fn("EnvironmentThreadState.make } | null>(null); const persistence = yield* Queue.sliding(1); // Latch set when the server reports the thread missing (subscribe fails - // with threadDisposition "not-found"). Gates the outer foreground/probe + // with OrchestrationThreadNotFoundError). Gates the outer foreground/probe // resubscribe path below so a deleted thread stops after its single // terminal attempt instead of retrying on every app wakeup. const terminalNotFound = yield* Ref.make(false); @@ -762,7 +762,7 @@ export const makeEnvironmentThreadState = Effect.fn("EnvironmentThreadState.make retryExpectedFailureAfter: "250 millis", resubscribe: foregroundResubscriptions, terminalFailure: { - matches: wasSubscribeThreadNotFound, + matches: isOrchestrationThreadNotFoundError, handle: (cause) => Ref.set(terminalNotFound, true).pipe( Effect.andThen( diff --git a/packages/contracts/src/orchestration.ts b/packages/contracts/src/orchestration.ts index 5727b9987ff5..7e906aac5a77 100644 --- a/packages/contracts/src/orchestration.ts +++ b/packages/contracts/src/orchestration.ts @@ -1852,10 +1852,18 @@ export class OrchestrationGetSnapshotError extends Schema.TaggedErrorClass()( + "OrchestrationThreadNotFoundError", + { threadId: ThreadId }, +) { + override get message(): string { + return `Thread ${this.threadId} was not found`; + } +} + export class OrchestrationDispatchCommandError extends Schema.TaggedErrorClass()( "OrchestrationDispatchCommandError", { diff --git a/packages/contracts/src/rpc.test.ts b/packages/contracts/src/rpc.test.ts index 6a9adebd85d6..a952ec6cf509 100644 --- a/packages/contracts/src/rpc.test.ts +++ b/packages/contracts/src/rpc.test.ts @@ -2,7 +2,7 @@ import { describe, expect, it } from "vite-plus/test"; import * as Exit from "effect/Exit"; import * as Schema from "effect/Schema"; -import { WsSubscribeServerConfigRpc } from "./rpc.ts"; +import { WsOrchestrationSubscribeThreadRpc, WsSubscribeServerConfigRpc } from "./rpc.ts"; /** * The client always sends `environmentThemes`, including to servers built @@ -29,3 +29,18 @@ describe("subscribeServerConfig payload compatibility", () => { expect(decoded).toEqual({}); }); }); + +const decodeSubscribeThreadError = Schema.decodeUnknownSync( + WsOrchestrationSubscribeThreadRpc.successSchema.error, +); + +describe("subscribeThread errors", () => { + it("decodes a missing thread as a distinct terminal error", () => { + const error = decodeSubscribeThreadError({ + _tag: "OrchestrationThreadNotFoundError", + threadId: "thread-1", + }); + expect(error._tag).toBe("OrchestrationThreadNotFoundError"); + expect(error.message).toBe("Thread thread-1 was not found"); + }); +}); diff --git a/packages/contracts/src/rpc.ts b/packages/contracts/src/rpc.ts index f7f2c2b6faa7..18144ed109b7 100644 --- a/packages/contracts/src/rpc.ts +++ b/packages/contracts/src/rpc.ts @@ -77,6 +77,7 @@ import { OrchestrationGetFullThreadDiffError, OrchestrationGetFullThreadDiffInput, OrchestrationGetSnapshotError, + OrchestrationThreadNotFoundError, OrchestrationSearchThreadsError, OrchestrationSearchThreadsInput, OrchestrationGetTurnDiffError, @@ -1109,7 +1110,11 @@ export const WsOrchestrationSubscribeThreadRpc = Rpc.make( { payload: OrchestrationRpcSchemas.subscribeThread.input, success: OrchestrationRpcSchemas.subscribeThread.output, - error: Schema.Union([OrchestrationGetSnapshotError, EnvironmentAuthorizationError]), + error: Schema.Union([ + OrchestrationGetSnapshotError, + OrchestrationThreadNotFoundError, + EnvironmentAuthorizationError, + ]), stream: true, }, ); From 74ce349339069ae4deff3240e570d5ecc9bb4a95 Mon Sep 17 00:00:00 2001 From: Lars Nieuwenhuis <35393046+lnieuwenhuis@users.noreply.github.com> Date: Sat, 5 Sep 2026 08:19:15 +0200 Subject: [PATCH 6/7] fix(client-runtime): serialize thread cache saves and deletion --- .../src/state/threads-sync.test.ts | 40 +++++++++++++++++++ packages/client-runtime/src/state/threads.ts | 36 ++++++++++------- 2 files changed, 61 insertions(+), 15 deletions(-) diff --git a/packages/client-runtime/src/state/threads-sync.test.ts b/packages/client-runtime/src/state/threads-sync.test.ts index 988d736e6f6a..2ad18a9b5cba 100644 --- a/packages/client-runtime/src/state/threads-sync.test.ts +++ b/packages/client-runtime/src/state/threads-sync.test.ts @@ -1157,6 +1157,46 @@ describe("EnvironmentThreads", () => { }), ); + for (const deletion of ["event", "not-found"] as const) { + it.effect(`waits for an in-flight cache save before ${deletion} removal`, () => + Effect.gen(function* () { + const writing = yield* Deferred.make(); + const release = yield* Deferred.make(); + const written = yield* Deferred.make(); + const removed = yield* Deferred.make(); + const cached = yield* Ref.make(false); + const harness = yield* makeHarness({ + saveThread: () => + Effect.gen(function* () { + yield* Deferred.succeed(writing, undefined); + yield* Deferred.await(release); + yield* Ref.set(cached, true); + yield* Deferred.succeed(written, undefined); + }), + onRemoveThread: Ref.set(cached, false).pipe( + Effect.andThen(Deferred.succeed(removed, undefined)), + Effect.asVoid, + ), + }); + yield* Queue.offer(harness.inputs, snapshot(BASE_THREAD)); + yield* awaitThreadState(harness.observed, (value) => value.status === "live"); + yield* TestClock.adjust("500 millis"); + yield* Deferred.await(writing); + yield* Queue.offer( + harness.inputs, + deletion === "event" + ? deleted() + : new OrchestrationThreadNotFoundError({ threadId: THREAD_ID }), + ); + yield* awaitThreadState(harness.observed, (value) => value.status === "deleted"); + yield* Deferred.succeed(release, undefined); + yield* Deferred.await(written); + yield* Deferred.await(removed); + expect(yield* Ref.get(cached)).toBe(false); + }), + ); + } + it.effect("does not resurrect a missing thread via queued persistence", () => Effect.gen(function* () { const harness = yield* makeHarness({ cached: BASE_THREAD }); diff --git a/packages/client-runtime/src/state/threads.ts b/packages/client-runtime/src/state/threads.ts index faa6fc705777..c8114dee7038 100644 --- a/packages/client-runtime/src/state/threads.ts +++ b/packages/client-runtime/src/state/threads.ts @@ -259,6 +259,9 @@ export const makeEnvironmentThreadState = Effect.fn("EnvironmentThreadState.make readonly epoch: number; } | null>(null); const persistence = yield* Queue.sliding(1); + // Cache removal must finish after any save that already passed its state + // check. Keep storage I/O separate from the stream/history application lock. + const persistenceLock = yield* Semaphore.make(1); // Latch set when the server reports the thread missing (subscribe fails // with OrchestrationThreadNotFoundError). Gates the outer foreground/probe // resubscribe path below so a deleted thread stops after its single @@ -268,9 +271,8 @@ export const makeEnvironmentThreadState = Effect.fn("EnvironmentThreadState.make const persist = Effect.fn("EnvironmentThreadState.persist")(function* ( snapshot: OrchestrationThreadDetailSnapshot, ) { - // A debounced write racing a deletion must not resurrect the cache after - // setDeleted removes it; the queue drain in setDeleted drops pending - // offers, this guard drops the offer already past the debounce. + // Under persistenceLock, delayed writes skip deleted snapshots and + // deletion waits for saves that have already passed this check. if ((yield* SubscriptionRef.get(state)).status === "deleted") return; if (resumeCache !== undefined && resumeCache.owner !== owner) return; if ( @@ -304,7 +306,7 @@ export const makeEnvironmentThreadState = Effect.fn("EnvironmentThreadState.make ), ), ); - }); + }, persistenceLock.withPermits(1)); yield* Stream.fromQueue(persistence).pipe( Stream.debounce("500 millis"), @@ -408,17 +410,21 @@ export const makeEnvironmentThreadState = Effect.fn("EnvironmentThreadState.make // snapshot cannot re-save the thread after deletion. The queue is // sliding(1), so a single poll discards the whole backlog. yield* Queue.poll(persistence); - if (resumeCache !== undefined && resumeCache.owner !== owner) return; - yield* cache.removeThread(environmentId, threadId).pipe( - Effect.catch((error) => - Effect.logWarning("Could not remove the cached thread.").pipe( - Effect.annotateLogs({ - environmentId, - threadId, - error: error.message, - }), - ), - ), + yield* persistenceLock.withPermits(1)( + Effect.gen(function* () { + if (resumeCache !== undefined && resumeCache.owner !== owner) return; + yield* cache.removeThread(environmentId, threadId).pipe( + Effect.catch((error) => + Effect.logWarning("Could not remove the cached thread.").pipe( + Effect.annotateLogs({ + environmentId, + threadId, + error: error.message, + }), + ), + ), + ); + }), ); }); From 5422217883f04ce66efd14f1d6fac1fa54e5edc5 Mon Sep 17 00:00:00 2001 From: Lars Nieuwenhuis <35393046+lnieuwenhuis@users.noreply.github.com> Date: Wed, 9 Sep 2026 11:22:49 +0200 Subject: [PATCH 7/7] test(client-runtime): await retry handler before advancing time --- .../client-runtime/src/rpc/client.test.ts | 29 ++++++++++--------- 1 file changed, 15 insertions(+), 14 deletions(-) diff --git a/packages/client-runtime/src/rpc/client.test.ts b/packages/client-runtime/src/rpc/client.test.ts index 7ed31df1c76c..e5d920d8816b 100644 --- a/packages/client-runtime/src/rpc/client.test.ts +++ b/packages/client-runtime/src/rpc/client.test.ts @@ -515,11 +515,20 @@ describe("environment RPC", () => { Effect.gen(function* () { const transient = new Error("transient snapshot failure"); const subscriptionCount = yield* Ref.make(0); + const retryReady = yield* Deferred.make(); + const retryStarted = yield* Deferred.make(); const client = { [WS_METHODS.subscribeTerminalEvents]: () => Stream.unwrap( Ref.getAndUpdate(subscriptionCount, (count) => count + 1).pipe( - Effect.map((count) => (count === 0 ? Stream.fail(transient) : Stream.never)), + Effect.map((count) => + count === 0 + ? Stream.fail(transient) + : Stream.fromEffect(Deferred.succeed(retryStarted, undefined)).pipe( + Stream.drain, + Stream.concat(Stream.never), + ), + ), ), ), } as unknown as WsRpcProtocolClient; @@ -530,7 +539,7 @@ describe("environment RPC", () => { WS_METHODS.subscribeTerminalEvents, {}, { - onExpectedFailure: () => Effect.void, + onExpectedFailure: () => Deferred.succeed(retryReady, undefined).pipe(Effect.asVoid), retryExpectedFailureAfter: "100 millis", terminalFailure: { matches: () => false, @@ -543,19 +552,11 @@ describe("environment RPC", () => { Effect.forkChild, ); - // The retry sleep must be scheduled before virtual time advances. - for ( - let attempt = 0; - attempt < 100 && (yield* Ref.get(subscriptionCount)) < 1; - attempt += 1 - ) { - yield* Effect.yieldNow; - } + // The handler runs immediately before the retry delay; let that fiber schedule its sleep. + yield* Deferred.await(retryReady); + yield* Effect.yieldNow; yield* TestClock.adjust("100 millis"); - for (let attempt = 0; attempt < 100; attempt += 1) { - if ((yield* Ref.get(subscriptionCount)) >= 2) break; - yield* Effect.yieldNow; - } + yield* Deferred.await(retryStarted); yield* Fiber.interrupt(subscriptionFiber); // Classified as a regular expected failure: retried once, handler untouched.