Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion apps/web/src/components/ThreadRouteView.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ import { useEnvironmentQuery } from "../state/query";
import { environmentShell } from "../state/shell";
import {
buildThreadRouteParams,
isThreadRouteSnapshotAuthoritative,
resolveThreadRouteRenderState,
type ThreadRouteTarget,
} from "../threadRoutes";
Expand Down Expand Up @@ -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,
);
Expand Down
28 changes: 28 additions & 0 deletions apps/web/src/threadRoutes.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,13 +6,21 @@ import { DraftId } from "./composerDraftStore";
import {
buildDraftThreadRouteParams,
buildThreadRouteParams,
isThreadRouteSnapshotAuthoritative,
resolveActiveThreadRouteRef,
resolveThreadRouteRenderState,
resolveThreadRouteRef,
resolveThreadRouteTarget,
} 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"));

Expand Down Expand Up @@ -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({
Expand Down
13 changes: 10 additions & 3 deletions apps/web/src/threadRoutes.ts
Original file line number Diff line number Diff line change
@@ -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";

Expand All @@ -20,19 +21,25 @@ 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;
serverThreadDetailExists: boolean;
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";
}
Expand Down
101 changes: 99 additions & 2 deletions packages/client-runtime/src/state/shell-sync.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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<OrchestrationShellStreamItem>();
const wakeups = yield* Queue.unbounded<ConnectionWakeups.ConnectionWakeup>();
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);
}),
);
});
29 changes: 15 additions & 14 deletions packages/client-runtime/src/state/shell.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<RpcSession | null>(null);
const lastSnapshotSession = yield* Ref.make<RpcSession | null>(null);
const activeSubscriptionSession = yield* Ref.make<RpcSession | null>(null);
const persistence = yield* Queue.sliding<OrchestrationShellSnapshot>(1);

Expand Down Expand Up @@ -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)) {
Expand Down Expand Up @@ -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({
Expand All @@ -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);
}
Expand All @@ -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,
};
}),
{
Expand Down
Loading