diff --git a/src/client/src/api/clients.test.ts b/src/client/src/api/clients.test.ts index eda5907..26f990a 100644 --- a/src/client/src/api/clients.test.ts +++ b/src/client/src/api/clients.test.ts @@ -253,6 +253,22 @@ describe("session API compatibility", () => { expect(JSON.parse(requestBody(fetchCall(fetchMock, 1)[1]))).toEqual({ sessions: [{ id: "s 1", cwd: "/repo" }] }); }); + it("carries a create's correlation token in the start request body when one is supplied", async () => { + const fetchMock = stubSequenceFetch([ + jsonResponse(sessionInfoResponse("s 1")), + jsonResponse(sessionInfoResponse("s 2")), + ]); + + await sessionsApi.startSession("/repo", "remote a", "pending-session-3-k2x9"); + await sessionsApi.startSession("/repo", "remote a"); + + expect(fetchCall(fetchMock, 0)[0]).toBe("https://pi.example.test/api/machines/remote%20a/sessions"); + expect(JSON.parse(requestBody(fetchCall(fetchMock, 0)[1]))).toEqual({ cwd: "/repo", startupToken: "pending-session-3-k2x9" }); + // The token is optional, so a caller with no row to label sends none rather + // than an empty one. + expect(JSON.parse(requestBody(fetchCall(fetchMock, 1)[1]))).toEqual({ cwd: "/repo" }); + }); + it("keeps legacy session-id calls free of cwd context", async () => { const fetchMock = stubJsonFetch({ accepted: true }); @@ -572,6 +588,10 @@ function requestBody(init: RequestInit | undefined): string { return init.body; } +function sessionInfoResponse(id: string) { + return { id, path: `/tmp/${id}.jsonl`, cwd: "/repo", created: "now", modified: "now", messageCount: 0, firstMessage: "" }; +} + function piWebConfigResponse(config: PiWebConfigValues) { return { path: "/tmp/pi-web/config.json", diff --git a/src/client/src/api/clients.ts b/src/client/src/api/clients.ts index d5ca7d0..e58012e 100644 --- a/src/client/src/api/clients.ts +++ b/src/client/src/api/clients.ts @@ -214,7 +214,7 @@ export const sessionsApi = { notificationInbox: (session: SessionLookup, machineId = "local") => request(sessionQueryPath(session, "notifications", machineId), parseSessionNotificationInboxSnapshot), dismissNotification: (session: SessionLookup, daemonInstanceId: string, notificationId: string, machineId = "local") => request(sessionPath(session, "notifications/dismiss", machineId), parseSessionNotificationInboxSnapshot, { method: "POST", body: sessionBody(session, { daemonInstanceId, notificationId }) }), dismissAllNotifications: (session: SessionLookup, daemonInstanceId: string, through: SessionNotificationDismissThrough, machineId = "local") => request(sessionPath(session, "notifications/dismiss-all", machineId), parseSessionNotificationInboxSnapshot, { method: "POST", body: sessionBody(session, { daemonInstanceId, throughOrder: through.order, throughOverflowWatermark: through.overflowWatermark }) }), - startSession: (cwd: string, machineId = "local") => request(`${machinePrefix(machineId)}/sessions`, parseSessionInfo, { method: "POST", body: JSON.stringify({ cwd }) }), + startSession: (cwd: string, machineId = "local", startupToken?: string) => request(`${machinePrefix(machineId)}/sessions`, parseSessionInfo, { method: "POST", body: JSON.stringify(startupToken === undefined ? { cwd } : { cwd, startupToken }) }), cleanupPreview: (input: SessionCleanupRequest, machineId = "local") => request(`${machinePrefix(machineId)}/sessions/cleanup/preview`, parseSessionCleanupPreviewResponse, { method: "POST", body: JSON.stringify(input) }), cleanup: (input: SessionCleanupRequest, machineId = "local") => request(`${machinePrefix(machineId)}/sessions/cleanup`, parseSessionCleanupExecuteResponse, { method: "POST", body: JSON.stringify(input) }), archiveMany: (sessions: readonly SessionLookup[], machineId = "local") => request(`${machinePrefix(machineId)}/sessions/bulk/archive`, parseSessionBulkArchiveResponse, { method: "POST", body: sessionBulkMutationBody(sessions) }), diff --git a/src/client/src/api/parsers.test.ts b/src/client/src/api/parsers.test.ts index d356dc2..a05ec35 100644 --- a/src/client/src/api/parsers.test.ts +++ b/src/client/src/api/parsers.test.ts @@ -289,18 +289,18 @@ describe("API parsers", () => { })).toThrow("positive safe integer"); }); - it("parses session startup progress with and without a wait detail", () => { + it("parses session startup progress with and without a correlation token", () => { const activity = { sessionId: "session-1", phase: "active", label: "Creating session", detail: "Starting the Pi session", at: "2026-07-20T00:00:01.000Z" }; - expect(parseSessionStartupProgressEvent({ type: "session.startup", cwd: "/repo", activity })).toEqual({ + expect(parseSessionStartupProgressEvent({ type: "session.startup", startupToken: "pending-session-1-abc", activity })).toEqual({ type: "session.startup", - cwd: "/repo", + startupToken: "pending-session-1-abc", activity, }); + // An open carries no token: the activity's own session id is the only route. const idle = { sessionId: "session-1", phase: "idle", label: "idle", at: "2026-07-20T00:00:02.000Z" }; - expect(parseSessionStartupProgressEvent({ type: "session.startup", cwd: "/repo", activity: idle })).toEqual({ + expect(parseSessionStartupProgressEvent({ type: "session.startup", activity: idle })).toEqual({ type: "session.startup", - cwd: "/repo", activity: idle, }); }); @@ -308,15 +308,17 @@ describe("API parsers", () => { it("rejects session startup progress that cannot be routed or rendered honestly", () => { const activity = { sessionId: "session-1", phase: "active", label: "Creating session", at: "2026-07-20T00:00:01.000Z" }; - expect(() => parseSessionStartupProgressEvent({ type: "activity.update", cwd: "/repo", activity })).toThrow("Invalid session startup event type"); - expect(() => parseSessionStartupProgressEvent({ type: "session.startup", activity })).toThrow("Expected string field: cwd"); - expect(() => parseSessionStartupProgressEvent({ type: "session.startup", cwd: "", activity })).toThrow("Expected non-empty string field: cwd"); - expect(() => parseSessionStartupProgressEvent({ type: "session.startup", cwd: "/repo" })).toThrow("Expected object response"); - expect(() => parseSessionStartupProgressEvent({ type: "session.startup", cwd: "/repo", activity: { ...activity, phase: "waiting" } })).toThrow("Expected session activity phase field: phase"); - expect(() => parseSessionStartupProgressEvent({ type: "session.startup", cwd: "/repo", activity: { ...activity, label: 7 } })).toThrow("Expected string field: label"); - expect(() => parseSessionStartupProgressEvent({ type: "session.startup", cwd: "/repo", activity: { ...activity, label: "" } })).toThrow("Expected non-empty string field: label"); - expect(() => parseSessionStartupProgressEvent({ type: "session.startup", cwd: "/repo", activity: { ...activity, detail: 7 } })).toThrow("Expected optional string field: detail"); - expect(() => parseSessionStartupProgressEvent({ type: "session.startup", cwd: "/repo", activity: { ...activity, sessionId: "" } })).toThrow("Expected non-empty string field: sessionId"); + expect(() => parseSessionStartupProgressEvent({ type: "activity.update", activity })).toThrow("Invalid session startup event type"); + expect(() => parseSessionStartupProgressEvent({ type: "session.startup" })).toThrow("Expected object response"); + expect(() => parseSessionStartupProgressEvent({ type: "session.startup", startupToken: 7, activity })).toThrow("Expected optional string field: startupToken"); + // An empty token would match nothing but must still be rejected rather than + // silently carried, so a malformed frame never reaches the routing at all. + expect(() => parseSessionStartupProgressEvent({ type: "session.startup", startupToken: "", activity })).toThrow("Expected non-empty string field: startupToken"); + expect(() => parseSessionStartupProgressEvent({ type: "session.startup", activity: { ...activity, phase: "waiting" } })).toThrow("Expected session activity phase field: phase"); + expect(() => parseSessionStartupProgressEvent({ type: "session.startup", activity: { ...activity, label: 7 } })).toThrow("Expected string field: label"); + expect(() => parseSessionStartupProgressEvent({ type: "session.startup", activity: { ...activity, label: "" } })).toThrow("Expected non-empty string field: label"); + expect(() => parseSessionStartupProgressEvent({ type: "session.startup", activity: { ...activity, detail: 7 } })).toThrow("Expected optional string field: detail"); + expect(() => parseSessionStartupProgressEvent({ type: "session.startup", activity: { ...activity, sessionId: "" } })).toThrow("Expected non-empty string field: sessionId"); }); it("parses session cleanup preview and execute responses", () => { diff --git a/src/client/src/api/parsers.ts b/src/client/src/api/parsers.ts index 9941aa0..4621037 100644 --- a/src/client/src/api/parsers.ts +++ b/src/client/src/api/parsers.ts @@ -275,15 +275,19 @@ export function parseSessionUnreadEvent(value: unknown): SessionUnreadEvent { /** * Validate a startup progress frame. The browser substitutes its own wording * from this event, so a malformed frame must be dropped rather than rendered: - * `cwd` is the routing key, and an activity missing its phase or label could - * otherwise blank out or freeze the text a user is reading while they wait. + * `startupToken` is the routing key when present, and an activity missing its + * phase or label could otherwise blank out or freeze the text a user is reading + * while they wait. An absent token is valid — an open routes by session id — but + * a present empty one is not, since it could match no row honestly. */ export function parseSessionStartupProgressEvent(value: unknown): SessionStartupProgressEvent { const record = requireRecord(value); if (record["type"] !== "session.startup") throw new Error("Invalid session startup event type"); + const startupToken = optionalString(record, "startupToken"); + if (startupToken === "") throw new Error("Expected non-empty string field: startupToken"); return { type: "session.startup", - cwd: requireNonEmptyString(record, "cwd"), + ...optionalField("startupToken", startupToken), activity: parseSessionActivity(record["activity"]), }; } diff --git a/src/client/src/controllers/sessionController.startupProgress.test.ts b/src/client/src/controllers/sessionController.startupProgress.test.ts index 0804e18..8b59c8d 100644 --- a/src/client/src/controllers/sessionController.startupProgress.test.ts +++ b/src/client/src/controllers/sessionController.startupProgress.test.ts @@ -3,8 +3,6 @@ import { initialAppState } from "../appState"; import { SessionController } from "./sessionController"; import { defaultApi, deferred, emptyPage, FakeSocket, oldSession, runPendingAnimationFrames, sessionLookupId, status, workspace, type AppState, type SessionActivity, type SessionInfo } from "./sessionController.testSupport"; -const REMOTE_MACHINE = { id: "remote", name: "Remote", kind: "remote" as const, createdAt: "now", updatedAt: "now" }; - function startupActivity(patch: Partial = {}): SessionActivity { return { sessionId: "backend-session", @@ -21,8 +19,15 @@ function idleStartupActivity(): SessionActivity { return { sessionId: "backend-session", phase: "idle", label: "idle", at: "2026-07-20T00:00:02.000Z" }; } +interface StartCall { + cwd: string; + machineId: string | undefined; + startupToken: string | undefined; +} + function pendingStartController(state: { current: AppState }, api: Partial = {}) { const startRequest = deferred(); + const startCalls: StartCall[] = []; const controller = new SessionController( () => state.current, (patch) => { state.current = { ...state.current, ...patch }; }, @@ -31,7 +36,10 @@ function pendingStartController(state: { current: AppState }, api: Partial startRequest.promise, + startSession: (cwd: string, machineId?: string, startupToken?: string) => { + startCalls.push({ cwd, machineId, startupToken }); + return startRequest.promise; + }, messages: () => Promise.resolve(emptyPage), status: (session) => Promise.resolve(status(sessionLookupId(session))), ...api, @@ -39,7 +47,7 @@ function pendingStartController(state: { current: AppState }, api: Partial { @@ -51,7 +59,7 @@ describe("SessionController session startup progress", () => { const temporaryId = state.current.selectedSession?.id; if (temporaryId === undefined) throw new Error("Expected temporary session id"); - controller.applyGlobalEvent({ type: "session.startup", cwd: workspace.path, activity: startupActivity() }); + controller.applyGlobalEvent({ type: "session.startup", startupToken: temporaryId, activity: startupActivity() }); runPendingAnimationFrames(); // The label changes while the user is waiting, before the start resolves, @@ -59,7 +67,7 @@ describe("SessionController session startup progress", () => { expect(state.current.activity).toMatchObject({ sessionId: temporaryId, phase: "active", label: "Creating session", detail: "Starting the Pi session" }); expect(state.current.sessionActivities[temporaryId]).toMatchObject({ detail: "Starting the Pi session" }); - controller.applyGlobalEvent({ type: "session.startup", cwd: workspace.path, activity: startupActivity({ detail: "Loading session extensions" }) }); + controller.applyGlobalEvent({ type: "session.startup", startupToken: temporaryId, activity: startupActivity({ detail: "Loading session extensions" }) }); runPendingAnimationFrames(); expect(state.current.activity?.detail).toBe("Loading session extensions"); @@ -75,10 +83,10 @@ describe("SessionController session startup progress", () => { const start = controller.startSession(); const temporaryId = state.current.selectedSession?.id; if (temporaryId === undefined) throw new Error("Expected temporary session id"); - controller.applyGlobalEvent({ type: "session.startup", cwd: workspace.path, activity: startupActivity() }); + controller.applyGlobalEvent({ type: "session.startup", startupToken: temporaryId, activity: startupActivity() }); runPendingAnimationFrames(); - controller.applyGlobalEvent({ type: "session.startup", cwd: workspace.path, activity: idleStartupActivity() }); + controller.applyGlobalEvent({ type: "session.startup", startupToken: temporaryId, activity: idleStartupActivity() }); runPendingAnimationFrames(); expect(state.current.activity).toMatchObject({ @@ -98,7 +106,9 @@ describe("SessionController session startup progress", () => { const start = controller.startSession(); await controller.send("queued while starting"); - controller.applyGlobalEvent({ type: "session.startup", cwd: workspace.path, activity: idleStartupActivity() }); + const temporaryId = state.current.selectedSession?.id; + if (temporaryId === undefined) throw new Error("Expected temporary session id"); + controller.applyGlobalEvent({ type: "session.startup", startupToken: temporaryId, activity: idleStartupActivity() }); runPendingAnimationFrames(); expect(state.current.activity?.detail).toBe("1 queued message will send when the backend session is ready"); @@ -119,7 +129,6 @@ describe("SessionController session startup progress", () => { controller.applyGlobalEvent({ type: "session.startup", - cwd: oldSession.cwd, activity: startupActivity({ sessionId: oldSession.id, label: "Opening session" }), }); runPendingAnimationFrames(); @@ -136,12 +145,11 @@ describe("SessionController session startup progress", () => { const temporaryId = state.current.selectedSession?.id; if (temporaryId === undefined) throw new Error("Expected temporary session id"); - // Opening an existing session in the same workspace publishes the same cwd as - // the pending create. The known id is the proof of which row it belongs to, so - // the pending row must keep its own wording instead of the other row's phase. + // Opening an existing session in the same workspace carries no create token, + // so the known id is the only proof of which row it belongs to and the pending + // row must keep its own wording instead of the other row's phase. controller.applyGlobalEvent({ type: "session.startup", - cwd: workspace.path, activity: startupActivity({ sessionId: existing.id, label: "Opening session" }), }); runPendingAnimationFrames(); @@ -153,7 +161,7 @@ describe("SessionController session startup progress", () => { await start; }); - it("keeps the generic wording when the startup progress cannot be attributed to one row", async () => { + it("keeps the generic wording when no pending row's token matches the startup progress", async () => { const state = { current: { ...initialAppState(), selectedWorkspace: workspace, sessions: [] } }; const { controller, startRequest } = pendingStartController(state); @@ -161,29 +169,54 @@ describe("SessionController session startup progress", () => { const temporaryId = state.current.selectedSession?.id; if (temporaryId === undefined) throw new Error("Expected temporary session id"); - // Another workspace's startup. - controller.applyGlobalEvent({ type: "session.startup", cwd: "/elsewhere", activity: startupActivity() }); - // The selected machine's socket is the only feed for these events, so a cwd - // that matches while another machine is selected belongs to a different row. - state.current = { ...state.current, selectedMachine: REMOTE_MACHINE }; - controller.applyGlobalEvent({ type: "session.startup", cwd: workspace.path, activity: startupActivity() }); - state.current = { ...state.current, selectedMachine: undefined }; + // Another browser tab's create, or another workspace's: its token is one this + // browser never minted, so there is no row here that it belongs to. + controller.applyGlobalEvent({ type: "session.startup", startupToken: "pending-session-9-other-tab", activity: startupActivity() }); + // A session this browser has not been told about — an agent's spawned + // subsession, say, whose `session.created` a pending create suppresses — is + // opened rather than created, so it carries no token at all. + controller.applyGlobalEvent({ type: "session.startup", activity: startupActivity({ sessionId: "foreign-session", label: "Opening session", detail: "Loading session extensions" }) }); runPendingAnimationFrames(); expect(state.current.activity?.detail).toBe("Waiting for the backend session to be ready"); + expect(state.current.activity?.label).toBe("Creating session"); - // A second concurrent start in the same workspace makes the target ambiguous, - // so neither row is given a phase that might belong to the other. + startRequest.resolve({ ...oldSession, id: "backend-session", path: "/tmp/backend-session.jsonl" }); + await start; + }); + + it("gives each of two concurrent creates only the progress its own token carries", async () => { + const state = { current: { ...initialAppState(), selectedWorkspace: workspace, sessions: [] } }; + const { controller, startRequest } = pendingStartController(state); + + const start = controller.startSession(); + const firstId = state.current.selectedSession?.id; const secondStart = controller.startSession(); - controller.applyGlobalEvent({ type: "session.startup", cwd: workspace.path, activity: startupActivity() }); + const secondId = state.current.selectedSession?.id; + if (firstId === undefined || secondId === undefined || firstId === secondId) throw new Error("Expected two distinct temporary session ids"); + + // Two creates in the same workspace are indistinguishable by workspace path; + // the token each request carried is what tells them apart. + controller.applyGlobalEvent({ type: "session.startup", startupToken: secondId, activity: startupActivity({ detail: "Loading session extensions" }) }); runPendingAnimationFrames(); - const secondTemporaryId = state.current.selectedSession?.id; - expect(secondTemporaryId).not.toBe(temporaryId); - expect(state.current.sessionActivities[temporaryId]?.detail).toBe("Waiting for the backend session to be ready"); - expect(state.current.sessionActivities[secondTemporaryId ?? ""]?.detail).toBe("Waiting for the backend session to be ready"); + expect(state.current.sessionActivities[secondId]).toMatchObject({ sessionId: secondId, detail: "Loading session extensions" }); + expect(state.current.sessionActivities[firstId]?.detail).toBe("Waiting for the backend session to be ready"); startRequest.resolve({ ...oldSession, id: "backend-session", path: "/tmp/backend-session.jsonl" }); await Promise.all([start, secondStart]); }); + + it("sends the pending row's own id as the create request's correlation token", async () => { + const state = { current: { ...initialAppState(), selectedWorkspace: workspace, sessions: [] } }; + const { controller, startRequest, startCalls } = pendingStartController(state); + + const start = controller.startSession(); + const temporaryId = state.current.selectedSession?.id; + + expect(startCalls).toEqual([{ cwd: workspace.path, machineId: "local", startupToken: temporaryId }]); + + startRequest.resolve({ ...oldSession, id: "backend-session", path: "/tmp/backend-session.jsonl" }); + await start; + }); }); diff --git a/src/client/src/controllers/sessionController.ts b/src/client/src/controllers/sessionController.ts index 2d7f65e..46fbaa0 100644 --- a/src/client/src/controllers/sessionController.ts +++ b/src/client/src/controllers/sessionController.ts @@ -185,7 +185,7 @@ export class SessionController { this.pendingSessionStarts.set(pending.tempId, pending); this.insertAndSelectPendingSession(pending.session); try { - const session = await this.api.startSession(workspace.path, machineId); + const session = await this.api.startSession(workspace.path, machineId, pending.tempId); await this.resolvePendingSessionStart(pending.tempId, session); } catch (error) { this.failPendingSessionStart(pending.tempId, error); @@ -1309,25 +1309,20 @@ export class SessionController { } // Session startup progress arrives while the daemon is still constructing the - // session, so the target row is resolved by session id when the browser knows - // it and by workspace path when it does not: a pending start knows its cwd but - // not the session id the daemon is creating. Once the row is resolved the - // progress goes through the normal activity buffer, so it renders exactly like - // any other activity and stays batched per frame. + // session, so the target row is resolved by exact identity only: a session id + // the browser already knows (an open), else the correlation token this browser + // minted for its own create and the daemon echoed back. Matching neither means + // the row is not one this browser shows — an agent's or another tab's session is + // *deliberately* absent while a create is pending — so it is ignored rather than + // guessed at. A resolved row goes through the normal activity buffer, rendering + // like any other activity and staying batched per frame. private queueStartupProgress(event: SessionStartupProgressEvent): void { - // A known session id is the strongest possible proof of the target, so it is - // checked first: while a create is pending in a workspace, an *existing* - // session in that same workspace can be opened too (another row selected, - // another tab, a subsession), and that open publishes the same cwd. Matching - // on cwd first would paint the pending row with another session's phase. if (this.getState().sessions.some((session) => session.id === event.activity.sessionId)) { this.queueActivityUpdate(event.activity); return; } - // The id is unknown, so this can only be a create whose id the browser has - // not been told yet. Route it by workspace path, the one key both sides share. - const pending = this.startupProgressPendingStart(event.cwd); - if (pending === undefined) return; + const pending = event.startupToken === undefined ? undefined : this.pendingSessionStarts.get(event.startupToken); + if (pending === undefined || pending.discarded) return; // An idle startup phase means the daemon has nothing left to attribute, so // restore this row's own generic wording rather than clearing the text of a // creation request that has not returned yet. @@ -1336,13 +1331,6 @@ export class SessionController { : { ...event.activity, sessionId: pending.tempId }); } - private startupProgressPendingStart(cwd: string): PendingSessionStart | undefined { - const machineId = selectedMachineId(this.getState()); - const matches = Array.from(this.pendingSessionStarts.values()) - .filter((pending) => pending.cwd === cwd && pending.machineId === machineId && !pending.discarded); - return matches.length === 1 ? matches[0] : undefined; - } - private schedulePendingFlush(): void { if (this.pendingFrame !== undefined) return; this.pendingFrame = requestAnimationFrame(() => { diff --git a/src/client/src/sessionSocket.test.ts b/src/client/src/sessionSocket.test.ts index efb2163..0337316 100644 --- a/src/client/src/sessionSocket.test.ts +++ b/src/client/src/sessionSocket.test.ts @@ -93,14 +93,15 @@ describe("notification socket guards", () => { it("accepts validated session startup progress and drops malformed frames", () => { const activity = { sessionId: "session-1", phase: "active", label: "Creating session", detail: "Starting the Pi session", at: "2026-07-20T00:00:01.000Z" }; - expect(parseRealtimeSocketEvent({ type: "session.startup", cwd: "/repo", activity })) - .toMatchObject({ type: "session.startup", cwd: "/repo", activity }); - expect(parseRealtimeSocketEvent({ type: "session.startup", cwd: "", activity })).toBeUndefined(); - expect(parseRealtimeSocketEvent({ type: "session.startup", cwd: "/repo" })).toBeUndefined(); - expect(parseRealtimeSocketEvent({ type: "session.startup", cwd: "/repo", activity: { ...activity, phase: "waiting" } })).toBeUndefined(); + expect(parseRealtimeSocketEvent({ type: "session.startup", startupToken: "pending-session-1-abc", activity })) + .toMatchObject({ type: "session.startup", startupToken: "pending-session-1-abc", activity }); + expect(parseRealtimeSocketEvent({ type: "session.startup", activity })).toMatchObject({ type: "session.startup", activity }); + expect(parseRealtimeSocketEvent({ type: "session.startup", startupToken: "", activity })).toBeUndefined(); + expect(parseRealtimeSocketEvent({ type: "session.startup" })).toBeUndefined(); + expect(parseRealtimeSocketEvent({ type: "session.startup", activity: { ...activity, phase: "waiting" } })).toBeUndefined(); // Startup progress is global-only, so it must not be accepted as a // per-session frame even when it is well formed. - expect(parseSessionSocketEvent({ type: "session.startup", cwd: "/repo", activity })).toBeUndefined(); + expect(parseSessionSocketEvent({ type: "session.startup", activity })).toBeUndefined(); }); it("preserves existing event acceptance without treating unknown types as realtime events", () => { diff --git a/src/server/sessions/piSessionService.startupProgress.test.ts b/src/server/sessions/piSessionService.startupProgress.test.ts index 57bf67f..aa55a8d 100644 --- a/src/server/sessions/piSessionService.startupProgress.test.ts +++ b/src/server/sessions/piSessionService.startupProgress.test.ts @@ -77,7 +77,7 @@ describe("PiSessionService session startup progress", () => { // The proof that matters: the user is told what is being waited on before // the wait ends, not after it. expect(startupText(hub)).toEqual(["Creating session: Starting the Pi session"]); - expect(startupEvents(hub).at(0)).toMatchObject({ cwd: "/workspace", activity: { sessionId: "session-1", phase: "active" } }); + expect(startupEvents(hub).at(0)).toMatchObject({ activity: { sessionId: "session-1", phase: "active" } }); runtimeResult.resolve(fake.runtime); await started; @@ -157,11 +157,42 @@ describe("PiSessionService session startup progress", () => { await service.start("/workspace"); - expect(startupEvents(hub).at(-1)).toMatchObject({ cwd: "/workspace", activity: { sessionId: "session-1", phase: "idle", label: "idle" } }); + expect(startupEvents(hub).at(-1)).toMatchObject({ activity: { sessionId: "session-1", phase: "idle", label: "idle" } }); expect(startupEvents(hub).at(-1)?.activity.detail).toBeUndefined(); await service.dispose(); }); + it("echoes a create's correlation token on every startup report of that construction", async () => { + const { hub, service } = startupService(); + + await service.start("/workspace", { startupToken: "pending-session-3-k2x9" }); + + // The token labels the browser row that is waiting, so it must ride every + // report of this construction, the closing idle one included. + expect(startupEvents(hub).map((event) => event.startupToken)).toEqual([ + "pending-session-3-k2x9", + "pending-session-3-k2x9", + "pending-session-3-k2x9", + ]); + // The token is an opaque throwaway label, never the session's identity. + expect(startupEvents(hub).map((event) => event.activity.sessionId)).toEqual(["session-1", "session-1", "session-1"]); + await service.dispose(); + }); + + it("publishes no correlation token when a create supplies none, and none for an open", async () => { + const created = startupService(); + await created.service.start("/workspace"); + const opened = startupService({ sessionRecords: [sessionRecord("session-1")] }); + await opened.service.status(sessionRef("session-1")); + + for (const hub of [created.hub, opened.hub]) { + expect(startupEvents(hub).length).toBeGreaterThan(0); + expect(startupEvents(hub).every((event) => event.startupToken === undefined)).toBe(true); + } + await created.service.dispose(); + await opened.service.dispose(); + }); + it("ends the startup window when the runtime construction itself fails", async () => { const failure = new Error("runtime unavailable"); const { hub, service } = startupService({ createAgentRuntime: () => Promise.reject(failure) }); diff --git a/src/server/sessions/piSessionService.ts b/src/server/sessions/piSessionService.ts index da3522b..8ccc346 100644 --- a/src/server/sessions/piSessionService.ts +++ b/src/server/sessions/piSessionService.ts @@ -199,6 +199,11 @@ type SessionCreationProvenance = "tracked-subsession"; interface StartSessionOptions { parentSession?: string; initialModel?: AgentModel; + /** + * Opaque label, echoed on this construction's startup progress so a browser + * row with no session id yet can recognise its own. + */ + startupToken?: string; } interface InternalStartSessionOptions extends StartSessionOptions { @@ -390,7 +395,7 @@ interface PendingSessionOpen { promise: Promise>; } -interface CreateSessionRuntimeOptions extends Pick { +interface CreateSessionRuntimeOptions extends Pick { notificationGeneration?: SessionNotificationGeneration; notifications?: "enabled" | "disabled"; /** @@ -992,6 +997,7 @@ export class PiSessionService implements SessionRouteService { cwd, { startupIntent: "create", + ...(options.startupToken === undefined ? {} : { startupToken: options.startupToken }), ...(options.initialModel === undefined ? {} : { initialModel: options.initialModel }), ...(options.creationProvenance === undefined ? {} : { creationProvenance: options.creationProvenance }), }, @@ -2353,7 +2359,7 @@ export class PiSessionService implements SessionRouteService { cwd: string, options: CreateSessionRuntimeOptions = {}, ): Promise> { - const startup = this.startupProgress(sessionManager, cwd, options.startupIntent ?? "open"); + const startup = this.startupProgress(sessionManager, options.startupIntent ?? "open", options.startupToken); try { return await this.createSessionRuntime(sessionManager, cwd, options, startup); } finally { @@ -2999,23 +3005,23 @@ export class PiSessionService implements SessionRouteService { /** * Build the reporter for one session construction. * - * The session id and cwd are both known before any await — a `SessionManager` - * has its id from construction — so the daemon can name what it is starting - * even though the `PiAgentSession` that {@link publishActivity} needs does not - * exist yet. When either is missing there is nothing honest to route on, so - * the reporter stays silent and the browser keeps its own generic wording. + * The session id is known before any await — a `SessionManager` has its id + * from construction — so the daemon can name what it is starting even though + * the `PiAgentSession` that {@link publishActivity} needs does not exist yet. + * Without an id there is nothing to report against, so the reporter stays + * silent and the browser keeps its own generic wording. */ - private startupProgress(sessionManager: PiSessionManager, cwd: string, intent: "create" | "open"): SessionStartupProgressReporter { + private startupProgress(sessionManager: PiSessionManager, intent: "create" | "open", startupToken: string | undefined): SessionStartupProgressReporter { const sessionId = sessionManager.getSessionId(); - if (sessionId === "" || cwd === "") return { report: noop, end: noop }; + if (sessionId === "") return { report: noop, end: noop }; const label = intent === "create" ? "Creating session" : "Opening session"; return { - report: (phase) => { this.publishStartupProgress(sessionId, cwd, label, "active", this.startupDetail(phase)); }, + report: (phase) => { this.publishStartupProgress(sessionId, startupToken, label, "active", this.startupDetail(phase)); }, end: () => { // A real activity published during the window (an extension error, say) // is the truth about this session and must survive the clear. if (this.activities.has(sessionId)) return; - this.publishStartupProgress(sessionId, cwd, "idle", "idle", undefined); + this.publishStartupProgress(sessionId, startupToken, "idle", "idle", undefined); }, }; } @@ -3027,17 +3033,17 @@ export class PiSessionService implements SessionRouteService { } /** - * Report startup progress on the global channel only, keyed by `cwd` so a - * browser row that has no session id yet can find it. + * Report startup progress on the global channel only, echoing the caller's + * correlation token so a waiting browser row recognises its own construction. * * Unlike {@link publishActivity} this deliberately records nothing: no * `activities` entry, no workspace activity, no unread observation. There is * no session to own that state, and a failed creation would leave it stranded. */ - private publishStartupProgress(sessionId: string, cwd: string, label: string, phase: "active" | "idle", detail: string | undefined): void { + private publishStartupProgress(sessionId: string, startupToken: string | undefined, label: string, phase: "active" | "idle", detail: string | undefined): void { const at = new Date().toISOString(); const activity = detail === undefined ? { sessionId, phase, label, at } : { sessionId, phase, label, detail, at }; - this.events.publishGlobal({ type: "session.startup", cwd, activity }); + this.events.publishGlobal(startupToken === undefined ? { type: "session.startup", activity } : { type: "session.startup", startupToken, activity }); } private publishActivity(session: PiAgentSession, label: string, phase: "active" | "idle" | "error", detail?: string): void { diff --git a/src/server/sessions/sessionRoutes.test.ts b/src/server/sessions/sessionRoutes.test.ts index caddfa9..e1e01aa 100644 --- a/src/server/sessions/sessionRoutes.test.ts +++ b/src/server/sessions/sessionRoutes.test.ts @@ -26,6 +26,7 @@ import { PiSessionService, type PiSessionManagerGateway } from "./piSessionServi import { testModelRuntime } from "./piSessionService.testSupport.js"; import { SessionNotificationStore } from "./sessionNotificationStore.js"; import type { SessionRouteLookup, SessionRouteService } from "./sessionService.js"; +import type { ClientSession } from "../types.js"; import { registerSessionRoutes } from "./sessionRoutes.js"; import type { NormalizedSessionCleanupRequest } from "./sessionCleanup.js"; @@ -666,6 +667,35 @@ describe("session routes", () => { } }); + it("forwards a create's optional correlation token alongside the normalized cwd", async () => { + const routeApp = Fastify({ logger: false }); + await routeApp.register(fastifyWebsocket); + const eventHub = new SessionEventHub(); + const routeService = new CapturingRouteSessionService(); + registerSessionRoutes(routeApp, routeService, eventHub); + + try { + const requestCwd = resolve("/repo"); + const withToken = await routeApp.inject({ method: "POST", url: "/sessions", payload: { cwd: requestCwd, startupToken: "pending-session-3-k2x9" } }); + const withoutToken = await routeApp.inject({ method: "POST", url: "/sessions", payload: { cwd: requestCwd } }); + // An older browser, or any non-browser caller, sends no token; and a + // malformed one must not reach the service as a label it would echo. + const malformedToken = await routeApp.inject({ method: "POST", url: "/sessions", payload: { cwd: requestCwd, startupToken: 7 } }); + + expect(withToken.statusCode).toBe(200); + expect(withoutToken.statusCode).toBe(200); + expect(malformedToken.statusCode).toBe(400); + expect(malformedToken.json()).toEqual({ error: "startupToken field must be a string" }); + expect(routeService.startCalls).toEqual([ + { cwd: requestCwd, startupToken: "pending-session-3-k2x9" }, + { cwd: requestCwd, startupToken: undefined }, + ]); + } finally { + await routeService.dispose(); + await routeApp.close(); + } + }); + it("rejects malformed bulk mutation bodies before calling the service", async () => { const routeApp = Fastify({ logger: false }); await routeApp.register(fastifyWebsocket); @@ -706,6 +736,7 @@ class CapturingRouteSessionService implements SessionRouteService { readonly bulkArchiveCalls: SessionBulkMutationRef[][] = []; readonly bulkDeleteCalls: SessionBulkMutationRef[][] = []; readonly navigateTreeCalls: { lookup: SessionRouteLookup; request: SessionTreeNavigateRequest }[] = []; + readonly startCalls: { cwd: string; startupToken: string | undefined }[] = []; reloadError: Error | undefined; clearQueueError: Error | undefined; @@ -772,7 +803,11 @@ class CapturingRouteSessionService implements SessionRouteService { } list(): never { throw unusedRouteMethod("list"); } - start(): never { throw unusedRouteMethod("start"); } + + start(cwd: string, options?: { startupToken?: string }): Promise { + this.startCalls.push({ cwd, startupToken: options?.startupToken }); + return Promise.resolve({ id: "session-1", path: "/tmp/session-1.jsonl", cwd, created: "2026-06-25T00:00:00.000Z", modified: "2026-06-25T00:00:00.000Z", messageCount: 0, firstMessage: "" }); + } dismissWarning(lookup: SessionRouteLookup, dismissId: string): Promise { this.dismissWarningCalls.push({ lookup, dismissId }); diff --git a/src/server/sessions/sessionRoutes.ts b/src/server/sessions/sessionRoutes.ts index da73425..783c411 100644 --- a/src/server/sessions/sessionRoutes.ts +++ b/src/server/sessions/sessionRoutes.ts @@ -45,10 +45,14 @@ export function registerSessionRoutes(app: FastifyInstance, sessions: SessionRou } }); - app.post<{ Body: { cwd?: unknown } | undefined }>(`${prefix}/sessions`, async (request, reply) => { + app.post<{ Body: { cwd?: unknown; startupToken?: unknown } | undefined }>(`${prefix}/sessions`, async (request, reply) => { try { const body = requireRecord(request.body); - return await sessions.start(normalizeRequestCwd(requireString(body, "cwd"))); + // An opaque label the caller uses to recognise its own construction's + // startup reports. Optional: only a browser row waiting for a session id + // has anything to correlate. + const startupToken = body["startupToken"] === undefined ? undefined : requireNonEmptyString(body, "startupToken"); + return await sessions.start(normalizeRequestCwd(requireString(body, "cwd")), optionalField("startupToken", startupToken)); } catch (error) { return reply.code(400).send({ error: errorMessage(error) }); } diff --git a/src/server/sessions/sessionService.ts b/src/server/sessions/sessionService.ts index 7e64414..844d01f 100644 --- a/src/server/sessions/sessionService.ts +++ b/src/server/sessions/sessionService.ts @@ -40,7 +40,12 @@ export type SessionRouteLookup = string | SessionRouteRef; */ export interface SessionRouteService { list(cwd: string): Promise; - start(cwd: string): Promise; + /** + * Create a session. `startupToken` is an opaque label the caller supplies so + * it can recognise this construction's startup progress reports; the service + * echoes it and never interprets it. + */ + start(cwd: string, options?: { startupToken?: string }): Promise; messages(ref: SessionRouteLookup, page?: { before?: number; limit?: number }): Promise; status(ref: SessionRouteLookup): Promise; streamSnapshot(ref: SessionRouteLookup): Promise; diff --git a/src/shared/apiTypes.ts b/src/shared/apiTypes.ts index e934fe4..308f15f 100644 --- a/src/shared/apiTypes.ts +++ b/src/shared/apiTypes.ts @@ -438,18 +438,18 @@ export interface QueuedSessionMessage { * constructing the agent session and no `PiAgentSession` exists yet, so * `activity.update` cannot be published for it. * - * `cwd` is the routing key for a browser row that is still waiting for a - * session id: a client-invented pending start knows its workspace path but not - * the daemon's session id. `activity.sessionId` carries the daemon's real id, so - * the same event also serves the case where the browser already knows it (an - * open of an existing session). + * `startupToken` is the opaque label a create request supplied, echoed back so a + * browser row still waiting for a session id recognises its own construction. + * The daemon never interprets it and it never becomes the session id: + * `activity.sessionId` always carries the real id, which is how an *open* of a + * session the browser already knows is routed instead. * * `activity.phase === "idle"` means the startup window ended with nothing left * to report, so a browser that substituted its own text should restore it. */ export interface SessionStartupProgressEvent { type: "session.startup"; - cwd: string; + startupToken?: string; activity: SessionActivity; }