diff --git a/apps/server/src/server.test.ts b/apps/server/src/server.test.ts index 2f32b6524d7b..134387bc22c9 100644 --- a/apps/server/src/server.test.ts +++ b/apps/server/src/server.test.ts @@ -9251,7 +9251,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 740be1330817..ecc38b52e6af 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -40,6 +40,7 @@ import { type OrchestrationShellStreamItem, OrchestrationGetFullThreadDiffError, OrchestrationGetSnapshotError, + OrchestrationThreadNotFoundError, OrchestrationSearchThreadsError, OrchestrationGetTurnDiffError, ORCHESTRATION_WS_METHODS, @@ -1736,10 +1737,7 @@ const makeWsRpcLayer = ( if (replayOnMissingSnapshot !== undefined) { return replayOnMissingSnapshot; } - return yield* new OrchestrationGetSnapshotError({ - message: `Thread ${input.threadId} was not found`, - cause: input.threadId, - }); + 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 0e1a723d5471..476780374fe1 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, + OrchestrationThreadNotFoundError, + ThreadId, +} from "@t3tools/contracts"; +import { describe, expect, it } from "@effect/vitest"; -import { wasBootstrapThreadDeleted } from "./orchestration.ts"; +import { wasBootstrapThreadDeleted, isOrchestrationThreadNotFoundError } from "./orchestration.ts"; + +describe("isOrchestrationThreadNotFoundError", () => { + it("matches the typed not-found error", () => { + expect( + isOrchestrationThreadNotFoundError( + new OrchestrationThreadNotFoundError({ threadId: ThreadId.make("thread-1") }), + ), + ).toBe(true); + }); + + it("rejects a generic snapshot error with a matching message", () => { + expect( + isOrchestrationThreadNotFoundError( + new OrchestrationGetSnapshotError({ + message: "Thread thread-1 was not found", + cause: "thread-1", + }), + ), + ).toBe(false); + }); + + it("rejects plain errors with a matching message", () => { + expect(isOrchestrationThreadNotFoundError(new Error("Thread thread-1 was not found"))).toBe( + false, + ); + }); + + it("rejects other snapshot errors", () => { + expect( + isOrchestrationThreadNotFoundError( + new OrchestrationGetSnapshotError({ + message: "Failed to load thread thread-1", + cause: "thread-1", + }), + ), + ).toBe(false); + }); + + it("rejects unrelated errors", () => { + expect(isOrchestrationThreadNotFoundError(new Error("boom"))).toBe(false); + expect(isOrchestrationThreadNotFoundError(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..565e6bd041c6 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, + OrchestrationThreadNotFoundError, +} from "@t3tools/contracts"; import * as Schema from "effect/Schema"; const isOrchestrationDispatchCommandError = Schema.is(OrchestrationDispatchCommandError); @@ -8,3 +11,5 @@ export function wasBootstrapThreadDeleted(error: unknown): boolean { isOrchestrationDispatchCommandError(error) && error.bootstrapThreadDisposition === "deleted" ); } + +export const isOrchestrationThreadNotFoundError = Schema.is(OrchestrationThreadNotFoundError); diff --git a/packages/client-runtime/src/rpc/client.test.ts b/packages/client-runtime/src/rpc/client.test.ts index 9e4e8a600d55..e5d920d8816b 100644 --- a/packages/client-runtime/src/rpc/client.test.ts +++ b/packages/client-runtime/src/rpc/client.test.ts @@ -465,6 +465,105 @@ 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 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.fromEffect(Deferred.succeed(retryStarted, undefined)).pipe( + Stream.drain, + Stream.concat(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: () => Deferred.succeed(retryReady, undefined).pipe(Effect.asVoid), + retryExpectedFailureAfter: "100 millis", + terminalFailure: { + matches: () => false, + handle: () => Effect.void, + }, + }, + ).pipe( + Stream.runDrain, + Effect.provideService(EnvironmentSupervisor.EnvironmentSupervisor, supervisor), + Effect.forkChild, + ); + + // 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"); + yield* Deferred.await(retryStarted); + yield* Fiber.interrupt(subscriptionFiber); + + // Classified as a regular expected failure: retried once, handler untouched. + expect(yield* Ref.get(subscriptionCount)).toBe(2); + }), + ); + it.effect.each(["input", "stream"] as const)( "does not classify %s subscription defects as expected failures", (where) => diff --git a/packages/client-runtime/src/rpc/client.ts b/packages/client-runtime/src/rpc/client.ts index 4ddfb9c4160e..9b51510f9541 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"; @@ -180,6 +181,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; + }; } function subscribeDynamicMapped( @@ -195,6 +205,12 @@ function subscribeDynamicMapped( Effect.gen(function* () { const supervisor = yield* EnvironmentSupervisor; const observer = yield* EnvironmentRpcSubscriptionObserver; + // 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 = options?.resubscribe === undefined @@ -206,6 +222,7 @@ function subscribeDynamicMapped( ), ); return sessions.pipe( + Stream.interruptWhen(Deferred.await(terminalHalt)), Stream.switchMap( Option.match({ onNone: () => Stream.empty, @@ -268,6 +285,21 @@ function subscribeDynamicMapped( ), ).pipe(Stream.drain); } + const terminal = options?.terminalFailure; + if ( + hasOnlyExpectedFailures && + terminal !== undefined && + cause.reasons.every( + (reason) => reason._tag === "Fail" && terminal.matches(reason.error), + ) + ) { + return Stream.fromEffect( + Deferred.succeed(terminalHalt, undefined).pipe( + Effect.asVoid, + Effect.andThen(terminal.handle(cause)), + ), + ).pipe(Stream.drain); + } if (hasOnlyExpectedFailures && options?.onExpectedFailure !== undefined) { const handled = Stream.fromEffect(options.onExpectedFailure(cause)).pipe( Stream.drain, diff --git a/packages/client-runtime/src/state/threads-sync.test.ts b/packages/client-runtime/src/state/threads-sync.test.ts index 619b5f43425d..81bb03f8b1e0 100644 --- a/packages/client-runtime/src/state/threads-sync.test.ts +++ b/packages/client-runtime/src/state/threads-sync.test.ts @@ -1,6 +1,8 @@ import { EnvironmentId, EventId, + OrchestrationGetSnapshotError, + OrchestrationThreadNotFoundError, ORCHESTRATION_WS_METHODS, ProjectId, ProviderInstanceId, @@ -141,6 +143,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(); @@ -223,7 +226,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()), @@ -988,4 +993,281 @@ 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 OrchestrationThreadNotFoundError({ threadId: THREAD_ID }), + ); + + 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 republish buffered snapshots after a terminal tombstone", () => + Effect.gen(function* () { + const harness = yield* makeHarness({ cached: BASE_THREAD }); + yield* Queue.offerAll(harness.inputs, [ + ...Array.from({ length: 400 }, () => snapshot(BASE_THREAD)), + new OrchestrationThreadNotFoundError({ threadId: THREAD_ID }), + ]); + yield* awaitThreadState(harness.observed, (value) => value.status === "deleted"); + yield* TestClock.adjust("30 seconds"); + expect((yield* Ref.get(harness.latest)).status).toBe("deleted"); + expect(yield* Ref.get(harness.savedThreads)).toEqual([]); + 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 OrchestrationThreadNotFoundError({ threadId: THREAD_ID }), + ); + 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 resubscribe a missing thread on session replacement", () => + Effect.gen(function* () { + const harness = yield* makeHarness({ cached: BASE_THREAD }); + yield* Queue.offer( + harness.inputs, + new OrchestrationThreadNotFoundError({ threadId: THREAD_ID }), + ); + 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("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 OrchestrationThreadNotFoundError({ threadId: THREAD_ID }), + ); + + // 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(); + 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([]); + }), + ); + + 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 }); + // 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 OrchestrationThreadNotFoundError({ threadId: THREAD_ID }), + ); + + 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 generic snapshot error", () => + 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 83b85bf02f09..f26fa35f4ec6 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 { isOrchestrationThreadNotFoundError } from "../errors/orchestration.ts"; import { EnvironmentCacheStore } from "../platform/persistence.ts"; import { subscribeDynamic } from "../rpc/client.ts"; import { ThreadSnapshotLoader, type ThreadSnapshotWindow } from "./threadSnapshotHttp.ts"; @@ -266,10 +267,21 @@ 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 + // terminal attempt instead of retrying on every app wakeup. + const terminalNotFound = yield* Ref.make(false); const persist = Effect.fn("EnvironmentThreadState.persist")(function* ( snapshot: OrchestrationThreadDetailSnapshot, ) { + // 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 ( committed.persisted && @@ -302,7 +314,7 @@ export const makeEnvironmentThreadState = Effect.fn("EnvironmentThreadState.make ), ), ); - }); + }, persistenceLock.withPermits(1)); yield* Stream.fromQueue(persistence).pipe( Stream.debounce("500 millis"), @@ -407,17 +419,25 @@ export const makeEnvironmentThreadState = Effect.fn("EnvironmentThreadState.make page: Option.none(), }); yield* remember; - 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, - }), - ), - ), + // 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); + 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, + }), + ), + ), + ); + }), ); }); @@ -510,7 +530,13 @@ export const makeEnvironmentThreadState = Effect.fn("EnvironmentThreadState.make const applyItem = Effect.fn("EnvironmentThreadState.applyItem")(function* ( item: OrchestrationThreadStreamItem, ) { - yield* applyLock.withPermits(1)(applyItemLocked(item).pipe(Effect.andThen(remember))); + yield* applyLock.withPermits(1)( + Effect.gen(function* () { + if (yield* Ref.get(terminalNotFound)) return; + yield* applyItemLocked(item); + yield* remember; + }), + ); }); // Merges an older disjoint page below the currently loaded window. All four @@ -657,7 +683,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)), + ), + ), }); // Only the first subscription after a warm live resume keeps the retained @@ -772,6 +803,19 @@ export const makeEnvironmentThreadState = Effect.fn("EnvironmentThreadState.make onExpectedFailure: (cause) => setStreamError(formatThreadError(cause)), retryExpectedFailureAfter: "250 millis", resubscribe: foregroundResubscriptions, + terminalFailure: { + matches: isOrchestrationThreadNotFoundError, + handle: (cause) => + Ref.set(terminalNotFound, true).pipe( + Effect.andThen( + Effect.logWarning("Subscribed thread is gone; stopping.", { + threadId, + cause: Cause.pretty(cause), + }), + ), + Effect.andThen(applyLock.withPermits(1)(setDeleted())), + ), + }, }, ).pipe(Stream.runForEach(applyItem)), ); @@ -845,11 +889,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}`), ), diff --git a/packages/contracts/src/orchestration.ts b/packages/contracts/src/orchestration.ts index 4d2f80a1101a..ba9e2f3edb1e 100644 --- a/packages/contracts/src/orchestration.ts +++ b/packages/contracts/src/orchestration.ts @@ -2049,6 +2049,15 @@ export class OrchestrationGetSnapshotError extends Schema.TaggedError()( + "OrchestrationThreadNotFoundError", + { threadId: ThreadId }, +) { + override get message(): string { + return `Thread ${this.threadId} was not found`; + } +} + export class OrchestrationDispatchCommandError extends Schema.TaggedError()( "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 9dbcaa9f4164..af8a512d1fdc 100644 --- a/packages/contracts/src/rpc.ts +++ b/packages/contracts/src/rpc.ts @@ -86,6 +86,7 @@ import { OrchestrationGetFullThreadDiffError, OrchestrationGetFullThreadDiffInput, OrchestrationGetSnapshotError, + OrchestrationThreadNotFoundError, OrchestrationSearchThreadsError, OrchestrationSearchThreadsInput, OrchestrationGetTurnDiffError, @@ -1108,12 +1109,19 @@ const WsOrchestrationSubscribeShellRpc = Rpc.make(ORCHESTRATION_WS_METHODS.subsc stream: true, }); -const WsOrchestrationSubscribeThreadRpc = Rpc.make(ORCHESTRATION_WS_METHODS.subscribeThread, { - payload: OrchestrationRpcSchemas.subscribeThread.input, - success: OrchestrationRpcSchemas.subscribeThread.output, - error: Schema.Union([OrchestrationGetSnapshotError, EnvironmentAuthorizationError]), - stream: true, -}); +export const WsOrchestrationSubscribeThreadRpc = Rpc.make( + ORCHESTRATION_WS_METHODS.subscribeThread, + { + payload: OrchestrationRpcSchemas.subscribeThread.input, + success: OrchestrationRpcSchemas.subscribeThread.output, + error: Schema.Union([ + OrchestrationGetSnapshotError, + OrchestrationThreadNotFoundError, + EnvironmentAuthorizationError, + ]), + stream: true, + }, +); const WsSubscribeTerminalEventsRpc = Rpc.make(WS_METHODS.subscribeTerminalEvents, { payload: Schema.Struct({}),