From 25e8bc48c8100db0ae685e159dda64188cd384f3 Mon Sep 17 00:00:00 2001 From: SAPHID Date: Mon, 7 Sep 2026 19:35:21 +0000 Subject: [PATCH 1/5] fix(web): preserve deep links while shell syncs --- apps/web/src/routes/_chat.$environmentId.$threadId.tsx | 8 ++++++-- apps/web/src/threadRoutes.test.ts | 8 ++++++++ apps/web/src/threadRoutes.ts | 6 ++++++ 3 files changed, 20 insertions(+), 2 deletions(-) diff --git a/apps/web/src/routes/_chat.$environmentId.$threadId.tsx b/apps/web/src/routes/_chat.$environmentId.$threadId.tsx index 16c0c4dfcc6d..d1d1687142a9 100644 --- a/apps/web/src/routes/_chat.$environmentId.$threadId.tsx +++ b/apps/web/src/routes/_chat.$environmentId.$threadId.tsx @@ -4,7 +4,11 @@ import { useEffect } from "react"; import ChatView from "../components/ChatView"; import { threadHasStarted } from "../components/ChatView.logic"; import { finalizePromotedDraftThreadByRef, useComposerDraftStore } from "../composerDraftStore"; -import { resolveThreadRouteRef, resolveThreadRouteRenderState } from "../threadRoutes"; +import { + isThreadRouteSnapshotAuthoritative, + resolveThreadRouteRef, + resolveThreadRouteRenderState, +} from "../threadRoutes"; import { resolveThreadSyncPhase } from "../threadSync"; import { useSidebarPendingFileDropStore } from "../sidebarPendingFileDropStore"; import { SidebarInset } from "~/components/ui/sidebar"; @@ -29,7 +33,7 @@ function ChatThreadRouteView() { const serverThreadDetail = useThreadDetail(threadRef); const serverThreadStatus = useThreadStatus(threadRef); const environmentThreadRefs = useEnvironmentThreadRefs(threadRef?.environmentId ?? null); - const bootstrapComplete = shell.data?.snapshot._tag === "Some"; + const bootstrapComplete = isThreadRouteSnapshotAuthoritative(shell.data?.status); const environmentHasServerThreads = environmentThreadRefs.length > 0; const draftThreadExists = useComposerDraftStore((store) => threadRef ? store.getDraftThreadByRef(threadRef) !== null : false, diff --git a/apps/web/src/threadRoutes.test.ts b/apps/web/src/threadRoutes.test.ts index 3edb2f38dc85..327c9df2dce6 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")); diff --git a/apps/web/src/threadRoutes.ts b/apps/web/src/threadRoutes.ts index fd5bc39d836a..dc023a79e6f2 100644 --- a/apps/web/src/threadRoutes.ts +++ b/apps/web/src/threadRoutes.ts @@ -20,6 +20,12 @@ type DraftThreadRouteState = { export type ThreadRouteRenderState = "loading" | "ready" | "missing"; +export function isThreadRouteSnapshotAuthoritative( + status: "empty" | "cached" | "synchronizing" | "live" | undefined, +): boolean { + return status === "live"; +} + export function resolveThreadRouteRenderState(input: { bootstrapComplete: boolean; serverThreadShellExists: boolean; From cd4462de99be605221680cac8c7e83cc79b4055a Mon Sep 17 00:00:00 2001 From: SAPHID Date: Mon, 7 Sep 2026 19:53:46 +0000 Subject: [PATCH 2/5] fix(web): refresh legacy shell resumes --- .../src/state/shell-sync.test.ts | 78 ++++++++++++++++++- packages/client-runtime/src/state/shell.ts | 8 +- 2 files changed, 81 insertions(+), 5 deletions(-) diff --git a/packages/client-runtime/src/state/shell-sync.test.ts b/packages/client-runtime/src/state/shell-sync.test.ts index 0d933c39f8ba..296c21b403fa 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,73 @@ describe("environment shell synchronization", () => { expect(yield* Ref.get(loaderCalls)).toBe(2); }), ); + + it.effect("refreshes the shell before resuming against servers without completion markers", () => + 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({ afterSequence: 10 }); + expect((yield* SubscriptionRef.get(shellState)).status).toBe("live"); + + yield* Queue.offer(wakeups, "application-active"); + expect(yield* Queue.take(subscribeInputs)).toEqual({ afterSequence: 20 }); + const resumed = yield* SubscriptionRef.get(shellState); + expect(resumed.status).toBe("live"); + expect(Option.getOrThrow(resumed.snapshot).snapshotSequence).toBe(20); + expect(yield* Ref.get(loaderCalls)).toBe(2); + }), + ); }); diff --git a/packages/client-runtime/src/state/shell.ts b/packages/client-runtime/src/state/shell.ts index 95d90f9b36f2..fcc079eb4b23 100644 --- a/packages/client-runtime/src/state/shell.ts +++ b/packages/client-runtime/src/state/shell.ts @@ -205,9 +205,11 @@ 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 another authoritative HTTP snapshot because they report + // the resumed subscription as live before its replay can arrive. + const hasAuthoritativeSnapshot = + supportsCompletionMarker && (yield* Ref.get(lastAuthoritativeSession)) === session; let canResume = hasAuthoritativeSnapshot; let current = yield* SubscriptionRef.get(state); if (!hasAuthoritativeSnapshot || Option.isNone(current.snapshot)) { From deaf9577a878141414629581554375f66dd0f5d6 Mon Sep 17 00:00:00 2001 From: SAPHID Date: Mon, 7 Sep 2026 21:40:59 +0000 Subject: [PATCH 3/5] fix(client): wait for legacy shell snapshot --- .../src/state/shell-sync.test.ts | 34 ++++++++++++++----- packages/client-runtime/src/state/shell.ts | 18 +++++----- 2 files changed, 33 insertions(+), 19 deletions(-) diff --git a/packages/client-runtime/src/state/shell-sync.test.ts b/packages/client-runtime/src/state/shell-sync.test.ts index 296c21b403fa..21796be5fb77 100644 --- a/packages/client-runtime/src/state/shell-sync.test.ts +++ b/packages/client-runtime/src/state/shell-sync.test.ts @@ -457,7 +457,7 @@ describe("environment shell synchronization", () => { }), ); - it.effect("refreshes the shell before resuming against servers without completion markers", () => + 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(); @@ -514,15 +514,31 @@ describe("environment shell synchronization", () => { ), ); - expect(yield* Queue.take(subscribeInputs)).toEqual({ afterSequence: 10 }); - expect((yield* SubscriptionRef.get(shellState)).status).toBe("live"); + 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(wakeups, "application-active"); - expect(yield* Queue.take(subscribeInputs)).toEqual({ afterSequence: 20 }); - const resumed = yield* SubscriptionRef.get(shellState); - expect(resumed.status).toBe("live"); - expect(Option.getOrThrow(resumed.snapshot).snapshotSequence).toBe(20); - expect(yield* Ref.get(loaderCalls)).toBe(2); + yield* Queue.offer(events, { + kind: "snapshot", + snapshot: { + ...LIVE_SHELL_SNAPSHOT, + snapshotSequence: 11, + threads: [{ id: "created-after-http-snapshot" } as never], + }, + }); + for (let attempt = 0; attempt < 100; attempt += 1) { + if ((yield* SubscriptionRef.get(shellState)).status === "live") break; + yield* Effect.yieldNow; + } + const live = yield* SubscriptionRef.get(shellState); + 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); }), ); }); diff --git a/packages/client-runtime/src/state/shell.ts b/packages/client-runtime/src/state/shell.ts index fcc079eb4b23..5bc8296c14ba 100644 --- a/packages/client-runtime/src/state/shell.ts +++ b/packages/client-runtime/src/state/shell.ts @@ -206,8 +206,8 @@ export const makeEnvironmentShellState = Effect.fn("EnvironmentShellState.make") // Foreground resubscriptions on the same live session can resume from // the in-memory cursor when the server marks replay completion. Older - // servers need another authoritative HTTP snapshot because they report - // the resumed subscription as live before its replay can arrive. + // servers need a full socket snapshot because an HTTP snapshot cannot + // account for events committed before the socket subscription starts. const hasAuthoritativeSnapshot = supportsCompletionMarker && (yield* Ref.get(lastAuthoritativeSession)) === session; let canResume = hasAuthoritativeSnapshot; @@ -229,7 +229,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); } @@ -241,17 +245,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, }; }), { From c0e437c11ebf8a6b80ac554a3a3c94cf57ee2179 Mon Sep 17 00:00:00 2001 From: Alex Southwell Date: Fri, 11 Sep 2026 17:40:16 +1000 Subject: [PATCH 4/5] test(client): wait on shell state instead of polling Co-Authored-By: Claude Opus 5 (1M context) --- apps/web/src/threadRoutes.ts | 3 ++- packages/client-runtime/src/state/shell-sync.test.ts | 11 ++++++----- 2 files changed, 8 insertions(+), 6 deletions(-) diff --git a/apps/web/src/threadRoutes.ts b/apps/web/src/threadRoutes.ts index dc023a79e6f2..573820f26ec9 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"; @@ -21,7 +22,7 @@ type DraftThreadRouteState = { export type ThreadRouteRenderState = "loading" | "ready" | "missing"; export function isThreadRouteSnapshotAuthoritative( - status: "empty" | "cached" | "synchronizing" | "live" | undefined, + status: EnvironmentShellStatus | undefined, ): boolean { return status === "live"; } diff --git a/packages/client-runtime/src/state/shell-sync.test.ts b/packages/client-runtime/src/state/shell-sync.test.ts index 21796be5fb77..f5ddc8f4f6d7 100644 --- a/packages/client-runtime/src/state/shell-sync.test.ts +++ b/packages/client-runtime/src/state/shell-sync.test.ts @@ -527,11 +527,12 @@ describe("environment shell synchronization", () => { threads: [{ id: "created-after-http-snapshot" } as never], }, }); - for (let attempt = 0; attempt < 100; attempt += 1) { - if ((yield* SubscriptionRef.get(shellState)).status === "live") break; - yield* Effect.yieldNow; - } - const live = yield* SubscriptionRef.get(shellState); + 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"); From 4b9e7bbad258545e0f06e2c57bee9e111f96c127 Mon Sep 17 00:00:00 2001 From: Alex Southwell Date: Sat, 12 Sep 2026 23:07:06 +1000 Subject: [PATCH 5/5] perf(client): skip redundant shell refresh on legacy resubscribe A session that already received a snapshot no longer reloads the shell over HTTP before requesting the full socket snapshot legacy servers need, so foreground wakeups and retries pay one transfer instead of two. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- .../client-runtime/src/state/shell-sync.test.ts | 6 ++++++ packages/client-runtime/src/state/shell.ts | 13 +++++++------ 2 files changed, 13 insertions(+), 6 deletions(-) diff --git a/packages/client-runtime/src/state/shell-sync.test.ts b/packages/client-runtime/src/state/shell-sync.test.ts index f5ddc8f4f6d7..75223780b89a 100644 --- a/packages/client-runtime/src/state/shell-sync.test.ts +++ b/packages/client-runtime/src/state/shell-sync.test.ts @@ -540,6 +540,12 @@ describe("environment shell synchronization", () => { { 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 5bc8296c14ba..b3bb9b23678f 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)) { @@ -207,12 +207,13 @@ export const makeEnvironmentShellState = Effect.fn("EnvironmentShellState.make") // Foreground resubscriptions on the same live session can resume from // 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. - const hasAuthoritativeSnapshot = - supportsCompletionMarker && (yield* Ref.get(lastAuthoritativeSession)) === session; + // 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({