From 638cecb991ed7281408cffd51b815fee1257f88c Mon Sep 17 00:00:00 2001 From: murashit Date: Fri, 26 Jun 2026 17:14:12 +0900 Subject: [PATCH] Fix app-server thread request boundaries --- .../chat/app-server/inbound/handler.ts | 7 +++-- .../{routing.ts => adapter.ts} | 31 ++++++++++++++++++- .../inbound/server-requests/responses.ts | 26 ---------------- .../threads/thread-management-actions.ts | 5 +-- src/workspace/thread-catalog.ts | 13 +++++++- .../chat/protocol/inbound/routing.test.ts | 2 +- .../threads/thread-management-actions.test.ts | 20 +++++++++++- tests/scripts/grit-policy.test.mjs | 6 ++-- tests/workspace/thread-catalog.test.ts | 17 ++++++++++ 9 files changed, 89 insertions(+), 38 deletions(-) rename src/features/chat/app-server/inbound/server-requests/{routing.ts => adapter.ts} (82%) delete mode 100644 src/features/chat/app-server/inbound/server-requests/responses.ts diff --git a/src/features/chat/app-server/inbound/handler.ts b/src/features/chat/app-server/inbound/handler.ts index c09feeea..327d76cb 100644 --- a/src/features/chat/app-server/inbound/handler.ts +++ b/src/features/chat/app-server/inbound/handler.ts @@ -25,11 +25,12 @@ import type { AppServerResourceEvent } from "../actions/metadata"; import { classifyAppServerLog } from "./app-server-logs"; import { type ChatNotificationEffect, planChatNotification } from "./notification-plan"; import { + routeServerRequest, serverRequestApprovalResponse, + serverRequestCurrentTimeResponse, serverRequestMcpElicitationResponse, serverRequestUserInputResponse, -} from "./server-requests/responses"; -import { routeServerRequest } from "./server-requests/routing"; +} from "./server-requests/adapter"; function cannotSendApprovalResponseMessage(): string { return "Could not send approval response because Codex app-server is not connected."; @@ -235,7 +236,7 @@ function respondToCurrentTimeRequest( context: ChatInboundHandlerContext, request: Extract, ): void { - if (!context.actions.respondToServerRequest(request.id, { currentTimeAt: Math.floor(Date.now() / 1000) })) { + if (!context.actions.respondToServerRequest(request.id, serverRequestCurrentTimeResponse(Date.now()))) { addSystemMessage(context, cannotSendCurrentTimeMessage()); } } diff --git a/src/features/chat/app-server/inbound/server-requests/routing.ts b/src/features/chat/app-server/inbound/server-requests/adapter.ts similarity index 82% rename from src/features/chat/app-server/inbound/server-requests/routing.ts rename to src/features/chat/app-server/inbound/server-requests/adapter.ts index 6a706fce..42deeeef 100644 --- a/src/features/chat/app-server/inbound/server-requests/routing.ts +++ b/src/features/chat/app-server/inbound/server-requests/adapter.ts @@ -1,10 +1,20 @@ import type { ServerRequest } from "../../../../../app-server/connection/rpc-messages"; import { appServerApprovalRequest, + appServerApprovalResponse, appServerMcpElicitationRequest, + appServerMcpElicitationResponse, appServerUserInputRequest, + appServerUserInputResponse, } from "../../../../../app-server/protocol/server-requests"; -import type { PendingApproval, PendingMcpElicitation, PendingUserInput } from "../../../../../domain/pending-requests/model"; +import type { + ApprovalAction, + McpElicitationAction, + McpElicitationContentValue, + PendingApproval, + PendingMcpElicitation, + PendingUserInput, +} from "../../../../../domain/pending-requests/model"; import { type ActiveRouteScope, fallbackMessageScope, @@ -89,6 +99,25 @@ export function routeServerRequest(request: ServerRequest, scope: ActiveRouteSco } } +export function serverRequestApprovalResponse(approval: PendingApproval, action: ApprovalAction): unknown { + return appServerApprovalResponse(approval, action); +} + +export function serverRequestUserInputResponse(questions: readonly { id: string }[], answers: Record): unknown { + return appServerUserInputResponse(questions, answers); +} + +export function serverRequestMcpElicitationResponse( + action: McpElicitationAction, + content: Record | null, +): unknown { + return appServerMcpElicitationResponse(action, content); +} + +export function serverRequestCurrentTimeResponse(currentTimeMs: number): unknown { + return { currentTimeAt: Math.floor(currentTimeMs / 1000) }; +} + function serverRequestScope(request: ServerRequest): MessageScope { if (!isServerRequest(request)) return fallbackMessageScope(request); const extractor = SERVER_REQUEST_SCOPE_EXTRACTORS[request.method] as (request: ServerRequest) => MessageScope; diff --git a/src/features/chat/app-server/inbound/server-requests/responses.ts b/src/features/chat/app-server/inbound/server-requests/responses.ts deleted file mode 100644 index ba7c12ef..00000000 --- a/src/features/chat/app-server/inbound/server-requests/responses.ts +++ /dev/null @@ -1,26 +0,0 @@ -import { - appServerApprovalResponse, - appServerMcpElicitationResponse, - appServerUserInputResponse, -} from "../../../../../app-server/protocol/server-requests"; -import type { - ApprovalAction, - McpElicitationAction, - McpElicitationContentValue, - PendingApproval, -} from "../../../../../domain/pending-requests/model"; - -export function serverRequestApprovalResponse(approval: PendingApproval, action: ApprovalAction): unknown { - return appServerApprovalResponse(approval, action); -} - -export function serverRequestUserInputResponse(questions: readonly { id: string }[], answers: Record): unknown { - return appServerUserInputResponse(questions, answers); -} - -export function serverRequestMcpElicitationResponse( - action: McpElicitationAction, - content: Record | null, -): unknown { - return appServerMcpElicitationResponse(action, content); -} diff --git a/src/features/chat/application/threads/thread-management-actions.ts b/src/features/chat/application/threads/thread-management-actions.ts index 99c1a58b..540fada1 100644 --- a/src/features/chat/application/threads/thread-management-actions.ts +++ b/src/features/chat/application/threads/thread-management-actions.ts @@ -173,12 +173,13 @@ async function forkThreadFromTurn( try { const sourceName = inheritedForkThreadName(threadId, threadManagementState(host).threadList.listedThreads); - const forkedThread = await forkThreadOnAppServer(scope.client, threadId, host.vaultPath); + let forkedThread = await forkThreadOnAppServer(scope.client, threadId, host.vaultPath); if (threadManagementScopeClientStale(host, scope)) return; const forkedThreadId = forkedThread.id; if (turnsToDrop > 0) { - await rollbackThreadOnAppServer(scope.client, forkedThreadId, turnsToDrop); + const snapshot = await rollbackThreadOnAppServer(scope.client, forkedThreadId, turnsToDrop); if (threadManagementScopeClientStale(host, scope)) return; + forkedThread = snapshot.thread; } host.applyThreadCatalogEvent({ type: "thread-forked", thread: forkedThread }); if (!threadManagementScopeStillTargetsOriginalPanel(host, scope)) return; diff --git a/src/workspace/thread-catalog.ts b/src/workspace/thread-catalog.ts index fd23e836..c88ad15f 100644 --- a/src/workspace/thread-catalog.ts +++ b/src/workspace/thread-catalog.ts @@ -123,9 +123,11 @@ function applyThreadCatalogEvent( store.setArchivedThreads(event.threads); return; case "thread-started": - case "thread-forked": upsertActiveThread(store, activeFacts, event.thread, acknowledgeByThreadId); return; + case "thread-forked": + upsertActiveThread(store, activeFacts, event.thread, acknowledgedByThreadVersion(event.thread)); + return; case "thread-touched": applyThreadTouchedEvent(store, activeFacts, event.threadId, event.recencyAt); return; @@ -329,6 +331,15 @@ function acknowledgedByName(name: string | null): (thread: Thread) => boolean { return (thread) => thread.name === name; } +function acknowledgedByThreadVersion(reference: Thread): (thread: Thread) => boolean { + return (thread) => + thread.updatedAt > reference.updatedAt || + (thread.updatedAt === reference.updatedAt && + thread.preview === reference.preview && + thread.name === reference.name && + thread.recencyAt === reference.recencyAt); +} + function acknowledgedByRecency(recencyAt: number | null): (thread: Thread) => boolean { return (thread) => thread.recencyAt === recencyAt; } diff --git a/tests/features/chat/protocol/inbound/routing.test.ts b/tests/features/chat/protocol/inbound/routing.test.ts index ba3531ed..13b6be6f 100644 --- a/tests/features/chat/protocol/inbound/routing.test.ts +++ b/tests/features/chat/protocol/inbound/routing.test.ts @@ -8,7 +8,7 @@ import { ROUTED_SERVER_NOTIFICATION_METHODS_BY_ROUTE_KIND, routeServerNotification, } from "../../../../../src/features/chat/app-server/inbound/notification-routing"; -import { routeServerRequest } from "../../../../../src/features/chat/app-server/inbound/server-requests/routing"; +import { routeServerRequest } from "../../../../../src/features/chat/app-server/inbound/server-requests/adapter"; import { chatStateFixture, chatStateWith } from "../../support/state"; const activeScope = { activeThreadId: "thread-active", activeTurnId: "turn-active" }; diff --git a/tests/features/chat/threads/thread-management-actions.test.ts b/tests/features/chat/threads/thread-management-actions.test.ts index c9d40e28..85c12895 100644 --- a/tests/features/chat/threads/thread-management-actions.test.ts +++ b/tests/features/chat/threads/thread-management-actions.test.ts @@ -164,6 +164,19 @@ describe("thread management actions", () => { it("forks from a selected turn by dropping later turns on the fork", async () => { const client = clientMock(); + client.forkThread.mockResolvedValue({ + thread: { + ...archivedThread(), + id: "forked", + sessionId: "forked", + name: "Fork before rollback", + preview: "Pre-rollback", + updatedAt: 10, + }, + }); + client.rollbackThread.mockResolvedValue({ + thread: { ...rollbackThread(), preview: "Post-rollback", updatedAt: 20 }, + }); const host = hostMock({ client, items: turnItems() }); const controller = threadManagementActions(host); @@ -173,7 +186,12 @@ describe("thread management actions", () => { expect(client.rollbackThread).toHaveBeenCalledWith("forked", 2); expect(host.applyThreadCatalogEvent).toHaveBeenCalledWith({ type: "thread-forked", - thread: expect.objectContaining({ id: "forked" }), + thread: expect.objectContaining({ + id: "forked", + name: "Rolled Back Thread", + preview: "Post-rollback", + updatedAt: 20, + }), }); expect(host.openThreadInNewView).toHaveBeenCalledWith("forked"); expect(client.archiveThread).not.toHaveBeenCalled(); diff --git a/tests/scripts/grit-policy.test.mjs b/tests/scripts/grit-policy.test.mjs index f81c7703..b41c66e6 100644 --- a/tests/scripts/grit-policy.test.mjs +++ b/tests/scripts/grit-policy.test.mjs @@ -351,7 +351,7 @@ export function timestamp(): number { APP_SERVER_PROTOCOL_BOUNDARY_MESSAGE, APP_SERVER_PROTOCOL_BOUNDARY_MESSAGE, ]); - expect(pluginMessages(report, "src/features/chat/app-server/inbound/server-requests/responses.ts")).toEqual([ + expect(pluginMessages(report, "src/features/chat/app-server/inbound/server-requests/adapter.ts")).toEqual([ APP_SERVER_PROTOCOL_BOUNDARY_MESSAGE, ]); }); @@ -661,7 +661,7 @@ export type Item = TurnItem; `.trimStart(), ); await writeFile( - path.join(cwd, "src/features/chat/app-server/inbound/server-requests/responses.ts"), + path.join(cwd, "src/features/chat/app-server/inbound/server-requests/adapter.ts"), ` import { appServerUserInputResponse } from "../../../../../app-server/protocol/server-requests"; @@ -825,7 +825,7 @@ export type AppServerThreadResumeClient = Pick; "src/features/chat/ui/protocol-leak.tsx", "src/features/chat/app-server/inbound/app-server-logs.ts", "src/features/chat/app-server/mappers/message-stream/turn-items.ts", - "src/features/chat/app-server/inbound/server-requests/responses.ts", + "src/features/chat/app-server/inbound/server-requests/adapter.ts", "src/domain/threads/model.ts", "src/domain/threads/format.ts", "src/features/chat/domain/message-stream/selectors.ts", diff --git a/tests/workspace/thread-catalog.test.ts b/tests/workspace/thread-catalog.test.ts index 32b2f90a..5bc9b3a0 100644 --- a/tests/workspace/thread-catalog.test.ts +++ b/tests/workspace/thread-catalog.test.ts @@ -176,6 +176,23 @@ describe("ThreadCatalog", () => { expect(catalog.activeSnapshot()).toEqual([thread("other")]); }); + it("keeps rollback fork metadata until active snapshots catch up to the rollback version", () => { + const { catalog } = catalogFixture(); + const forkBeforeRollback = thread("forked", false, { name: "Before", preview: "Before rollback", updatedAt: 20 }); + const forkAfterRollback = thread("forked", false, { name: "After", preview: "After rollback", updatedAt: 20 }); + const forkAfterFutureUpdate = thread("forked", false, { name: "Future", preview: "Future update", updatedAt: 21 }); + catalog.apply({ type: "active-list-snapshot-received", threads: [thread("existing")] }); + + catalog.apply({ type: "thread-forked", thread: forkAfterRollback }); + catalog.apply({ type: "active-list-snapshot-received", threads: [forkBeforeRollback, thread("existing")] }); + + expect(catalog.activeSnapshot()).toEqual([forkAfterRollback, thread("existing")]); + + catalog.apply({ type: "active-list-snapshot-received", threads: [forkAfterFutureUpdate, thread("existing")] }); + + expect(catalog.activeSnapshot()).toEqual([forkAfterFutureUpdate, thread("existing")]); + }); + it("keeps app-server rename facts when an older active list resolves later", async () => { const staleRefresh = deferred(); const fetchThreads = vi