diff --git a/apps/web/src/components/ThreadRouteView.tsx b/apps/web/src/components/ThreadRouteView.tsx index 90bff34051dd..f456b8772d6f 100644 --- a/apps/web/src/components/ThreadRouteView.tsx +++ b/apps/web/src/components/ThreadRouteView.tsx @@ -26,6 +26,7 @@ import { useEnvironmentQuery } from "../state/query"; import { environmentShell } from "../state/shell"; import { buildThreadRouteParams, + isThreadRouteSnapshotAuthoritative, resolveThreadRouteRenderState, type ThreadRouteTarget, } from "../threadRoutes"; @@ -82,7 +83,7 @@ export function ThreadRouteView({ target }: { target: ThreadRouteTarget }) { const serverThreadDetail = useThreadDetail(serverThreadRef); const serverThreadStatus = useThreadStatus(serverThreadRef); const environmentThreadRefs = useEnvironmentThreadRefs(serverThreadRef?.environmentId ?? null); - const bootstrapComplete = shell.data?.snapshot._tag === "Some"; + const bootstrapComplete = isThreadRouteSnapshotAuthoritative(shell.data?.status); const draftThread = useComposerDraftStore((store) => serverThreadRef ? store.getDraftThreadByRef(serverThreadRef) : null, ); diff --git a/apps/web/src/threadRoutes.test.ts b/apps/web/src/threadRoutes.test.ts index 3edb2f38dc85..5c2fe1426d53 100644 --- a/apps/web/src/threadRoutes.test.ts +++ b/apps/web/src/threadRoutes.test.ts @@ -6,6 +6,7 @@ import { DraftId } from "./composerDraftStore"; import { buildDraftThreadRouteParams, buildThreadRouteParams, + isThreadRouteSnapshotAuthoritative, resolveActiveThreadRouteRef, resolveThreadRouteRenderState, resolveThreadRouteRef, @@ -13,6 +14,13 @@ import { } from "./threadRoutes"; describe("threadRoutes", () => { + it("does not treat cached thread lists as authoritative for deep links", () => { + expect(isThreadRouteSnapshotAuthoritative("empty")).toBe(false); + expect(isThreadRouteSnapshotAuthoritative("cached")).toBe(false); + expect(isThreadRouteSnapshotAuthoritative("synchronizing")).toBe(false); + expect(isThreadRouteSnapshotAuthoritative("live")).toBe(true); + }); + it("builds canonical thread route params from a scoped ref", () => { const ref = scopeThreadRef("env-1" as never, ThreadId.make("thread-1")); @@ -127,6 +135,26 @@ describe("threadRoutes", () => { ).toBe("ready"); }); + it.each([ + { status: "cached", source: "detail" }, + { status: "synchronizing", source: "detail" }, + { status: "cached", source: "draft" }, + { status: "synchronizing", source: "draft" }, + ] as const)( + "keeps available $source visible while the shell is $status", + ({ status, source }) => { + expect( + resolveThreadRouteRenderState({ + bootstrapComplete: isThreadRouteSnapshotAuthoritative(status), + serverThreadShellExists: false, + serverThreadDetailExists: source === "detail", + serverThreadDetailDeleted: false, + draftThreadExists: source === "draft", + }), + ).toBe("ready"); + }, + ); + it("distinguishes bootstrap loading from a missing thread", () => { expect( resolveThreadRouteRenderState({ diff --git a/apps/web/src/threadRoutes.ts b/apps/web/src/threadRoutes.ts index fd5bc39d836a..0860e8bca454 100644 --- a/apps/web/src/threadRoutes.ts +++ b/apps/web/src/threadRoutes.ts @@ -1,4 +1,5 @@ import { scopeThreadRef } from "@t3tools/client-runtime/environment"; +import type { EnvironmentShellStatus } from "@t3tools/client-runtime/state/shell"; import type { EnvironmentId, ScopedThreadRef, ThreadId } from "@t3tools/contracts"; import type { DraftId } from "./composerDraftStore"; @@ -20,6 +21,12 @@ type DraftThreadRouteState = { export type ThreadRouteRenderState = "loading" | "ready" | "missing"; +export function isThreadRouteSnapshotAuthoritative( + status: EnvironmentShellStatus | undefined, +): boolean { + return status === "live"; +} + export function resolveThreadRouteRenderState(input: { bootstrapComplete: boolean; serverThreadShellExists: boolean; @@ -27,12 +34,12 @@ export function resolveThreadRouteRenderState(input: { serverThreadDetailDeleted: boolean; draftThreadExists: boolean; }): ThreadRouteRenderState { - if (!input.bootstrapComplete) { - return "loading"; - } if (input.serverThreadDetailExists || input.draftThreadExists) { return "ready"; } + if (!input.bootstrapComplete) { + return "loading"; + } if (input.serverThreadDetailDeleted) { return "missing"; } diff --git a/packages/client-runtime/src/state/shell-sync.test.ts b/packages/client-runtime/src/state/shell-sync.test.ts index 9483d86de0ff..bad9b70a1587 100644 --- a/packages/client-runtime/src/state/shell-sync.test.ts +++ b/packages/client-runtime/src/state/shell-sync.test.ts @@ -48,10 +48,15 @@ const LIVE_SHELL_SNAPSHOT: OrchestrationShellSnapshot = { updatedAt: "2026-06-06T00:00:00.000Z", }; -function session(client: WsRpcProtocolClient): RpcSession.RpcSession { +function session( + client: WsRpcProtocolClient, + options?: { readonly completionMarker?: boolean }, +): RpcSession.RpcSession { return { client, - initialConfig: Effect.succeed({ shellResumeCompletionMarker: true } as never), + initialConfig: Effect.succeed( + (options?.completionMarker === false ? {} : { shellResumeCompletionMarker: true }) as never, + ), subscribeServerConfig: (input) => client.subscribeServerConfig(input), ready: Effect.void, probe: Effect.void, @@ -451,4 +456,96 @@ describe("environment shell synchronization", () => { expect(yield* Ref.get(loaderCalls)).toBe(2); }), ); + + it.effect("keeps a legacy HTTP refresh non-authoritative until a socket snapshot arrives", () => + Effect.gen(function* () { + const events = yield* Queue.unbounded(); + const wakeups = yield* Queue.unbounded(); + const subscribeInputs = yield* Queue.unbounded<{ readonly afterSequence?: number }>(); + const loaderCalls = yield* Ref.make(0); + const client = { + [ORCHESTRATION_WS_METHODS.subscribeShell]: (input: { readonly afterSequence?: number }) => + Stream.unwrap( + Queue.offer(subscribeInputs, input).pipe(Effect.as(Stream.fromQueue(events))), + ), + } as unknown as WsRpcProtocolClient; + const supervisorState = yield* SubscriptionRef.make(AVAILABLE_CONNECTION_STATE); + const activeSession = yield* SubscriptionRef.make( + Option.some(session(client, { completionMarker: false })), + ); + const supervisor = EnvironmentSupervisor.EnvironmentSupervisor.of({ + target: TARGET, + state: supervisorState, + session: activeSession, + prepared: yield* SubscriptionRef.make(Option.some(PREPARED)), + connect: Effect.void, + disconnect: Effect.void, + retryNow: Effect.void, + } satisfies EnvironmentSupervisor.EnvironmentSupervisor["Service"]); + const cache = Persistence.EnvironmentCacheStore.of({ + loadShell: () => Effect.succeed(Option.none()), + saveShell: () => Effect.void, + loadThread: () => Effect.succeed(Option.none()), + saveThread: () => Effect.void, + removeThread: () => Effect.void, + loadServerConfig: () => Effect.succeed(Option.none()), + saveServerConfig: () => Effect.void, + loadVcsRefs: () => Effect.succeed(Option.none()), + saveVcsRefs: () => Effect.void, + removeVcsRefs: () => Effect.void, + clearVcsRefs: () => Effect.void, + clear: () => Effect.void, + }); + const snapshotLoader = ShellSnapshotLoader.of({ + load: () => + Ref.updateAndGet(loaderCalls, (count) => count + 1).pipe( + Effect.map((count) => + Option.some({ ...LIVE_SHELL_SNAPSHOT, snapshotSequence: count * 10 }), + ), + ), + }); + const shellState = yield* makeEnvironmentShellState().pipe( + Effect.provideService(EnvironmentSupervisor.EnvironmentSupervisor, supervisor), + Effect.provideService(Persistence.EnvironmentCacheStore, cache), + Effect.provideService(ShellSnapshotLoader, snapshotLoader), + Effect.provideService( + ConnectionWakeups.ConnectionWakeups, + ConnectionWakeups.ConnectionWakeups.of({ changes: Stream.fromQueue(wakeups) }), + ), + ); + + expect(yield* Queue.take(subscribeInputs)).toEqual({}); + const initial = yield* SubscriptionRef.get(shellState); + expect(initial.status).toBe("synchronizing"); + expect(Option.getOrThrow(initial.snapshot).snapshotSequence).toBe(10); + + yield* Queue.offer(events, { + kind: "snapshot", + snapshot: { + ...LIVE_SHELL_SNAPSHOT, + snapshotSequence: 11, + threads: [{ id: "created-after-http-snapshot" } as never], + }, + }); + const live = Option.getOrThrow( + yield* SubscriptionRef.changes(shellState).pipe( + Stream.filter((state) => state.status === "live"), + Stream.runHead, + ), + ); + expect(yield* Queue.size(events)).toBe(0); + expect(Option.getOrThrow(live.snapshot).snapshotSequence).toBe(11); + expect(live.status).toBe("live"); + expect(Option.getOrThrow(live.snapshot).threads).toEqual([ + { id: "created-after-http-snapshot" }, + ]); + expect(yield* Ref.get(loaderCalls)).toBe(1); + + // A same-session resubscribe still requests a full socket snapshot but + // does not pay for another HTTP refresh first. + yield* Queue.offer(wakeups, "application-active"); + expect(yield* Queue.take(subscribeInputs)).toEqual({}); + expect(yield* Ref.get(loaderCalls)).toBe(1); + }), + ); }); diff --git a/packages/client-runtime/src/state/shell.ts b/packages/client-runtime/src/state/shell.ts index b7f39b509462..caa281ed8630 100644 --- a/packages/client-runtime/src/state/shell.ts +++ b/packages/client-runtime/src/state/shell.ts @@ -72,7 +72,7 @@ export const makeEnvironmentShellState = Effect.fn("EnvironmentShellState.make") error: Option.none(), }); const awaitingCompletion = yield* Ref.make(false); - const lastAuthoritativeSession = yield* Ref.make(null); + const lastSnapshotSession = yield* Ref.make(null); const activeSubscriptionSession = yield* Ref.make(null); const persistence = yield* Queue.sliding(1); @@ -177,7 +177,7 @@ export const makeEnvironmentShellState = Effect.fn("EnvironmentShellState.make") if (receivedSnapshot) { const session = yield* Ref.get(activeSubscriptionSession); if (session !== null) { - yield* Ref.set(lastAuthoritativeSession, session); + yield* Ref.set(lastSnapshotSession, session); } } if (next.snapshot !== initial.snapshot && Option.isSome(next.snapshot)) { @@ -205,12 +205,15 @@ export const makeEnvironmentShellState = Effect.fn("EnvironmentShellState.make") yield* setSynchronizing; // Foreground resubscriptions on the same live session can resume from - // the in-memory cursor. A new session reloads the authoritative HTTP - // snapshot so a valid cursor cannot preserve incomplete cached data. - const hasAuthoritativeSnapshot = (yield* Ref.get(lastAuthoritativeSession)) === session; + // the in-memory cursor when the server marks replay completion. Older + // servers need a full socket snapshot because an HTTP snapshot cannot + // account for events committed before the socket subscription starts, + // but a session that already holds one skips the redundant HTTP reload. + const sessionHasSnapshot = (yield* Ref.get(lastSnapshotSession)) === session; + const hasAuthoritativeSnapshot = supportsCompletionMarker && sessionHasSnapshot; let canResume = hasAuthoritativeSnapshot; let current = yield* SubscriptionRef.get(state); - if (!hasAuthoritativeSnapshot || Option.isNone(current.snapshot)) { + if (!sessionHasSnapshot || Option.isNone(current.snapshot)) { const prepared = yield* SubscriptionRef.get(supervisor.prepared).pipe( Effect.flatMap( Option.match({ @@ -227,7 +230,11 @@ export const makeEnvironmentShellState = Effect.fn("EnvironmentShellState.make") ); const httpSnapshot = yield* snapshotLoader.load(prepared); if (Option.isSome(httpSnapshot)) { + yield* Ref.set(awaitingCompletion, true); yield* applyItems([{ kind: "snapshot", snapshot: httpSnapshot.value }]); + if (!supportsCompletionMarker) { + yield* Ref.set(awaitingCompletion, false); + } canResume = true; current = yield* SubscriptionRef.get(state); } @@ -239,17 +246,11 @@ export const makeEnvironmentShellState = Effect.fn("EnvironmentShellState.make") return supportsCompletionMarker ? { requestCompletionMarker: true as const } : {}; } if (!supportsCompletionMarker) { - // Without a completion marker there is no synchronized signal for a - // resumed subscription, so report live immediately, like threads. - yield* SubscriptionRef.update(state, (value) => ({ - ...value, - status: "live" as const, - error: Option.none(), - })); + return {}; } return { afterSequence: current.snapshot.value.snapshotSequence, - ...(supportsCompletionMarker ? { requestCompletionMarker: true as const } : {}), + requestCompletionMarker: true as const, }; }), {