From b4e4f6251e26c31c0131de9aeb6055f88302567d Mon Sep 17 00:00:00 2001 From: murashit Date: Fri, 3 Jul 2026 16:44:58 +0900 Subject: [PATCH] Separate chat runtime notification events --- ...-chat-application-outer-layer-imports.grit | 4 +- .../threads}/rename-lifecycle.ts | 0 .../app-server/inbound/notification-plan.ts | 582 ++++-------------- .../inbound/notification-routing.ts | 10 +- .../chat/app-server/inbound/runtime-events.ts | 213 +++++++ .../conversation/runtime-event-plan.ts | 233 +++++++ .../conversation/runtime-events.ts | 94 +++ .../chat/application/state/ui-state.ts | 2 +- src/features/threads-view/state.ts | 2 +- src/features/threads/list/row-projection.ts | 2 +- .../threads}/rename-lifecycle.test.ts | 2 +- .../chat/app-server/inbound/routing.test.ts | 15 + .../app-server/inbound/runtime-events.test.ts | 56 ++ .../conversation/runtime-event-plan.test.ts | 102 +++ tests/scripts/grit-policy.test.mjs | 15 +- 15 files changed, 851 insertions(+), 481 deletions(-) rename src/{features/threads/list => domain/threads}/rename-lifecycle.ts (100%) create mode 100644 src/features/chat/app-server/inbound/runtime-events.ts create mode 100644 src/features/chat/application/conversation/runtime-event-plan.ts create mode 100644 src/features/chat/application/conversation/runtime-events.ts rename tests/{features/threads/list => domain/threads}/rename-lifecycle.test.ts (97%) create mode 100644 tests/features/chat/app-server/inbound/runtime-events.test.ts create mode 100644 tests/features/chat/application/conversation/runtime-event-plan.test.ts diff --git a/scripts/grit/import-boundaries/no-chat-application-outer-layer-imports.grit b/scripts/grit/import-boundaries/no-chat-application-outer-layer-imports.grit index c366b2a5..de773b7b 100644 --- a/scripts/grit/import-boundaries/no-chat-application-outer-layer-imports.grit +++ b/scripts/grit/import-boundaries/no-chat-application-outer-layer-imports.grit @@ -9,7 +9,7 @@ private pattern js_module_reference() { js_module_reference() as $stmt where { $stmt <: contains `$source` where { - $source <: r"^[\"'](?:(?:\.\./)+(?:app-server|host|panel|presentation|ui)|src/(?:app-server|features/chat/(?:app-server|host|panel|presentation|ui)))(?:/.*)?[\"']$" + $source <: r"^[\"'](?:(?:\.\./)+(?:app-server|host|panel|presentation|ui)|(?:\.\./){2,}(?:selection-rewrite|thread-picker|threads|threads-view|turn-diff)|src/(?:app-server|features/(?:chat/(?:app-server|host|panel|presentation|ui)|selection-rewrite|thread-picker|threads|threads-view|turn-diff)))(?:/.*)?[\"']$" }, - register_diagnostic(span=$stmt, message="Chat application modules must not import app-server, host, panel, presentation, or UI layers; expose state and workflow contracts instead.", severity="error") + register_diagnostic(span=$stmt, message="Chat application modules must not import app-server, sibling feature, host, panel, presentation, or UI layers; expose state and workflow contracts instead.", severity="error") } diff --git a/src/features/threads/list/rename-lifecycle.ts b/src/domain/threads/rename-lifecycle.ts similarity index 100% rename from src/features/threads/list/rename-lifecycle.ts rename to src/domain/threads/rename-lifecycle.ts diff --git a/src/features/chat/app-server/inbound/notification-plan.ts b/src/features/chat/app-server/inbound/notification-plan.ts index 36827e3d..50812cfc 100644 --- a/src/features/chat/app-server/inbound/notification-plan.ts +++ b/src/features/chat/app-server/inbound/notification-plan.ts @@ -1,59 +1,20 @@ import type { ServerNotification } from "../../../../app-server/connection/rpc-messages"; import { threadFromAppServerRecord } from "../../../../app-server/services/threads"; -import { jsonPreview } from "../../../../domain/display/json-preview"; import { threadTokenUsageFromRuntimeUsage } from "../../../../domain/runtime/metrics"; import { normalizeExplicitThreadName } from "../../../../domain/threads/model"; -import type { ThreadConversationSummary } from "../../../../domain/threads/transcript"; import type { ThreadCatalogEvent } from "../../../threads/catalog/thread-catalog"; -import { activeTurnId, pendingTurnStart as pendingTurnStartForState } from "../../application/conversation/turn-state"; +import { type ConversationRuntimeEffect, planConversationRuntimeEvents } from "../../application/conversation/runtime-event-plan"; import { activeThreadSettingsAppliedAction } from "../../application/state/actions"; -import { messageStreamItems } from "../../application/state/message-stream"; import type { ChatAction, ChatState } from "../../application/state/root-reducer"; -import { reconcileCompletedTurnItems } from "../../domain/message-stream/completed-turn-reconciliation"; import { goalChangeItem } from "../../domain/message-stream/factories/goal-items"; -import { - STREAMED_COMMAND_RUNNING_TEXT, - STREAMED_FILE_CHANGE_IN_PROGRESS_TEXT, - STREAMED_MCP_PROGRESS_LABEL, -} from "../../domain/message-stream/factories/streaming-items"; -import { createSystemItem } from "../../domain/message-stream/factories/system-items"; -import type { MessageStreamItem, MessageStreamItemKind } from "../../domain/message-stream/items"; -import { attachHookRunsToTurn, completeReasoningItems, upsertMessageStreamItemById } from "../../domain/message-stream/updates"; import type { AppServerResourceEvent } from "../actions/metadata"; -import { - type AppServerFileChange, - normalizeFileChanges, - streamingFileChangeMessageStreamItem, -} from "../mappers/message-stream/file-changes"; -import { hookRunMessageStreamItem } from "../mappers/message-stream/hook-run-items"; -import { createAutoReviewResultItem, createReviewResultItem } from "../mappers/message-stream/review-result-items"; -import { taskProgressMessageStreamItem } from "../mappers/message-stream/task-progress"; -import { - type AppServerTurnItem, - completedConversationSummaryFromAppServerTurn, - messageStreamItemFromTurnItem, - messageStreamItemsFromTurns, - shouldSuppressLifecycleItem, -} from "../mappers/message-stream/turn-items"; -import { - type DiagnosticStatusNotification, - type DiagnosticStatusNotificationMethod, - routeServerNotification, - type StreamUpdateNotification, - type StreamUpdateNotificationMethod, - type ThreadLifecycleNotification, - type ThreadLifecycleNotificationMethod, - type TurnLifecycleNotification, - type TurnLifecycleNotificationMethod, - type UserVisibleNoticeNotification, - type UserVisibleNoticeNotificationMethod, -} from "./notification-routing"; +import { type DiagnosticStatusNotification, routeServerNotification, type ThreadLifecycleNotification } from "./notification-routing"; +import { conversationRuntimeEventsFromNotification } from "./runtime-events"; export type ChatNotificationEffect = - | { type: "refresh-threads" } + | Exclude | { type: "refresh-server-diagnostics"; forceResourceProbes?: boolean } | { type: "apply-app-server-resource-event"; event: AppServerResourceEvent } - | { type: "maybe-name-thread"; threadId: string; turnId: string; completedSummary: ThreadConversationSummary | null } | { type: "apply-thread-catalog-event"; event: ThreadCatalogEvent }; export interface ChatNotificationPlan { @@ -64,257 +25,6 @@ export interface ChatNotificationPlan { export type LocalItemIdProvider = (prefix: string) => string; const EMPTY_PLAN: ChatNotificationPlan = { actions: [], effects: [] }; -const MESSAGE_CONTEXT_COMPACTED = "Context compacted."; - -type ServerNotificationPlanner = ( - notification: Extract, -) => ChatNotificationPlan; -type ServerNotificationPlannerMap = { [Method in M]: ServerNotificationPlanner }; -type ServerNotificationLocalPlanner = ( - notification: Extract, - localItemId: LocalItemIdProvider, -) => ChatNotificationPlan; -type ServerNotificationLocalPlannerMap = { - [Method in M]: ServerNotificationLocalPlanner; -}; -type ServerNotificationStatePlanner = ( - state: ChatState, - notification: Extract, - localItemId: LocalItemIdProvider, -) => ChatNotificationPlan; -type ServerNotificationStatePlannerMap = { - [Method in M]: ServerNotificationStatePlanner; -}; - -const DIAGNOSTIC_STATUS_PLANNERS = { - "thread/tokenUsage/updated": (notification) => - actionPlan({ - type: "active-thread/token-usage-set", - tokenUsage: threadTokenUsageFromRuntimeUsage(notification.params.tokenUsage), - }), - "account/rateLimits/updated": () => ({ - actions: [], - effects: [{ type: "apply-app-server-resource-event", event: { type: "rate-limits-updated", preserveExistingOnFailure: true } }], - }), - "skills/changed": () => ({ - actions: [], - effects: [{ type: "apply-app-server-resource-event", event: { type: "skills-changed", forceReload: true } }], - }), - "app/list/updated": () => ({ - actions: [], - effects: [{ type: "refresh-server-diagnostics" }], - }), - "mcpServer/oauthLogin/completed": () => ({ - actions: [], - effects: [{ type: "refresh-server-diagnostics", forceResourceProbes: true }], - }), - "mcpServer/startupStatus/updated": (notification) => ({ - actions: [], - effects: [ - { - type: "apply-app-server-resource-event", - event: { - type: "mcp-startup-status-updated", - name: notification.params.name, - status: notification.params.status, - message: notification.params.error, - }, - }, - ], - }), -} satisfies ServerNotificationPlannerMap; - -const USER_VISIBLE_NOTICE_PLANNERS = { - "thread/compacted": (_notification, localItemId) => systemMessagePlan({ id: localItemId("system"), text: MESSAGE_CONTEXT_COMPACTED }), - "model/rerouted": jsonNoticePlan, - deprecationNotice: jsonNoticePlan, - error: jsonNoticePlan, - warning: jsonNoticePlan, - configWarning: jsonNoticePlan, - "windows/worldWritableWarning": jsonNoticePlan, - "windowsSandbox/setupCompleted": (notification, localItemId) => - notification.params.success ? EMPTY_PLAN : jsonNoticePlan(notification, localItemId), -} satisfies ServerNotificationLocalPlannerMap; - -const STREAM_UPDATE_PLANNERS = { - "item/agentMessage/delta": (_state, notification) => { - const { params } = notification; - return actionPlan({ - type: "message-stream/assistant-delta-appended", - itemId: params.itemId, - turnId: params.turnId, - delta: params.delta, - completeReasoning: true, - }); - }, - "item/plan/delta": (_state, notification) => { - const { params } = notification; - return actionPlan({ - type: "message-stream/plan-delta-appended", - itemId: params.itemId, - turnId: params.turnId, - delta: params.delta, - }); - }, - "turn/plan/updated": (_state, notification) => - actionPlan({ - type: "message-stream/item-upserted", - item: taskProgressMessageStreamItem(notification.params.turnId, notification.params.explanation, notification.params.plan), - }), - "item/reasoning/summaryTextDelta": (state, notification) => - appendToolTextPlan(state, notification.params.itemId, notification.params.turnId, "reasoning", notification.params.delta, "reasoning"), - "item/reasoning/textDelta": (state, notification) => - appendToolTextPlan(state, notification.params.itemId, notification.params.turnId, "reasoning", notification.params.delta, "reasoning"), - "item/reasoning/summaryPartAdded": (state, notification) => - appendToolTextPlan(state, notification.params.itemId, notification.params.turnId, "reasoning", "", "reasoning"), - "item/started": (_state, notification) => startedItemPlan(notification.params.item, notification.params.turnId), - "item/completed": (_state, notification) => completedItemPlan(notification.params.item, notification.params.turnId), - "item/commandExecution/outputDelta": (_state, notification) => - actionPlan({ - type: "message-stream/item-output-appended", - itemId: notification.params.itemId, - turnId: notification.params.turnId, - delta: notification.params.delta, - kind: "command", - fallbackText: STREAMED_COMMAND_RUNNING_TEXT, - }), - "item/fileChange/patchUpdated": (_state, notification) => - fileChangePlan(notification.params.itemId, notification.params.turnId, notification.params.changes, "inProgress"), - "item/fileChange/outputDelta": (_state, notification) => - actionPlan({ - type: "message-stream/item-output-appended", - itemId: notification.params.itemId, - turnId: notification.params.turnId, - delta: notification.params.delta, - kind: "fileChange", - fallbackText: STREAMED_FILE_CHANGE_IN_PROGRESS_TEXT, - }), - "turn/diff/updated": (_state, notification) => - actionPlan({ type: "message-stream/turn-diff-updated", turnId: notification.params.turnId, diff: notification.params.diff }), - "hook/started": (state, notification) => hookRunPlan(state, notification.params.run, notification.params.turnId, "running"), - "hook/completed": (state, notification) => - hookRunPlan(state, notification.params.run, notification.params.turnId, notification.params.run.status), - "item/mcpToolCall/progress": (_state, notification) => - actionPlan({ - type: "message-stream/tool-output-appended", - itemId: notification.params.itemId, - turnId: notification.params.turnId, - delta: notification.params.message, - fallbackLabel: STREAMED_MCP_PROGRESS_LABEL, - }), - "item/autoApprovalReview/started": autoApprovalReviewPlan, - "item/autoApprovalReview/completed": autoApprovalReviewPlan, - guardianWarning: (state, notification, localItemId) => { - const item = createReviewResultItem(localItemId("review"), notification.params.message); - if ( - isUnstructuredAutoReviewWarning(item) && - hasStructuredAutoReviewResult(messageStreamItems(state.messageStream), activeTurnId(state)) - ) { - return EMPTY_PLAN; - } - return actionPlan({ type: "message-stream/item-upserted", item }); - }, -} satisfies ServerNotificationStatePlannerMap; - -const TURN_LIFECYCLE_PLANNERS = { - "turn/started": (state, notification) => ({ - actions: [ - { - type: "turn/started", - threadId: notification.params.threadId, - turnId: notification.params.turn.id, - items: messageStreamItemsWithPendingPromptSubmitHooks(state, notification.params.turn.id), - }, - ], - effects: [ - { - type: "apply-thread-catalog-event", - event: { type: "thread-touched", threadId: notification.params.threadId, recencyAt: notification.params.turn.startedAt }, - }, - ], - }), - "turn/completed": (state, notification) => { - if (activeTurnId(state) !== notification.params.turn.id) return EMPTY_PLAN; - return { - actions: [ - { - type: "turn/completed", - turnId: notification.params.turn.id, - status: notification.params.turn.status, - items: completeReasoningItems( - reconcileCompletedTurnItems({ - currentItems: messageStreamItems(state.messageStream), - completedTurnId: notification.params.turn.id, - turnItems: messageStreamItemsFromTurns([notification.params.turn]), - }), - notification.params.turn.id, - ), - }, - ], - effects: [ - { - type: "maybe-name-thread", - threadId: notification.params.threadId, - turnId: notification.params.turn.id, - completedSummary: completedConversationSummaryFromAppServerTurn(notification.params.turn), - }, - { type: "refresh-threads" }, - ], - }; - }, -} satisfies ServerNotificationStatePlannerMap; - -const THREAD_LIFECYCLE_PLANNERS = { - "thread/started": (state, notification) => { - const effects: ChatNotificationEffect[] = [ - { - type: "apply-thread-catalog-event", - event: { type: "thread-started", thread: threadFromAppServerRecord(notification.params.thread) }, - }, - ]; - if (!state.activeThread.id || state.activeThread.id === notification.params.thread.id) { - return { actions: [{ type: "active-thread/cwd-set", cwd: notification.params.thread.cwd }], effects }; - } - return { actions: [], effects }; - }, - "thread/archived": (_state, notification) => ({ - actions: [], - effects: [{ type: "apply-thread-catalog-event", event: { type: "thread-archived", threadId: notification.params.threadId } }], - }), - "thread/deleted": (_state, notification) => ({ - actions: [], - effects: [{ type: "apply-thread-catalog-event", event: { type: "thread-deleted", threadId: notification.params.threadId } }], - }), - "thread/unarchived": (_state, notification) => ({ - actions: [], - effects: [{ type: "apply-thread-catalog-event", event: { type: "thread-unarchived", threadId: notification.params.threadId } }], - }), - "thread/name/updated": (_state, notification) => { - const name = normalizeExplicitThreadName(notification.params.threadName); - return { - actions: [], - effects: [{ type: "apply-thread-catalog-event", event: { type: "thread-renamed", threadId: notification.params.threadId, name } }], - }; - }, - "thread/settings/updated": (state, notification) => { - if (state.activeThread.id !== notification.params.threadId) return EMPTY_PLAN; - return actionPlan(activeThreadSettingsAppliedAction(notification.params.threadSettings)); - }, - "thread/goal/updated": (state, notification, localItemId) => { - if (state.activeThread.id !== notification.params.threadId) return EMPTY_PLAN; - const actions: ChatAction[] = [{ type: "active-thread/goal-set", goal: notification.params.goal }]; - const item = goalChangeItem(localItemId("goal"), state.activeThread.goal, notification.params.goal); - if (item) actions.push({ type: "message-stream/item-upserted", item }); - return { actions, effects: [] }; - }, - "thread/goal/cleared": (state, notification, localItemId) => { - if (state.activeThread.id !== notification.params.threadId) return EMPTY_PLAN; - const actions: ChatAction[] = [{ type: "active-thread/goal-set", goal: null }]; - const item = goalChangeItem(localItemId("goal"), state.activeThread.goal, null); - if (item) actions.push({ type: "message-stream/item-upserted", item }); - return { actions, effects: [] }; - }, -} satisfies ServerNotificationStatePlannerMap; export function planChatNotification( state: ChatState, @@ -323,7 +33,7 @@ export function planChatNotification( ): ChatNotificationPlan { const route = routeServerNotification(notification, { activeThreadId: state.activeThread.id, - activeTurnId: activeTurnId(state), + activeTurnId: activeTurnIdForState(state), }); switch (route.kind) { case "inactive": @@ -331,37 +41,68 @@ export function planChatNotification( case "unhandled": return EMPTY_PLAN; case "streamUpdate": - return planStreamUpdate(state, route.notification, localItemId); case "turnLifecycle": - return planTurnLifecycle(state, route.notification, localItemId); + case "requestResolved": + case "userVisibleNotice": + return runtimeEventsPlan(state, route.notification, localItemId); case "threadLifecycle": return planThreadLifecycle(state, route.notification, localItemId); - case "requestResolved": - return { - actions: [{ type: "request/resolved", requestId: route.notification.params.requestId }], - effects: [], - }; case "diagnosticStatus": return planDiagnosticStatus(route.notification); - case "userVisibleNotice": - return planUserVisibleNotice(route.notification, localItemId); } } -function planStreamUpdate( +function runtimeEventsPlan( state: ChatState, - notification: StreamUpdateNotification, + notification: Parameters[0], localItemId: LocalItemIdProvider, ): ChatNotificationPlan { - return planNotificationWithStateByMethod(state, notification, STREAM_UPDATE_PLANNERS, localItemId); + const plan = planConversationRuntimeEvents(state, conversationRuntimeEventsFromNotification(notification, localItemId)); + return { actions: plan.actions, effects: plan.effects.map(chatNotificationEffectFromConversationRuntimeEffect) }; } -function planTurnLifecycle( - state: ChatState, - notification: TurnLifecycleNotification, - localItemId: LocalItemIdProvider, -): ChatNotificationPlan { - return planNotificationWithStateByMethod(state, notification, TURN_LIFECYCLE_PLANNERS, localItemId); +function chatNotificationEffectFromConversationRuntimeEffect(effect: ConversationRuntimeEffect): ChatNotificationEffect { + switch (effect.type) { + case "thread-recency-touched": + return { + type: "apply-thread-catalog-event", + event: { type: "thread-touched", threadId: effect.threadId, recencyAt: effect.recencyAt }, + }; + case "refresh-threads": + case "maybe-name-thread": + return effect; + } +} + +function planDiagnosticStatus(notification: DiagnosticStatusNotification): ChatNotificationPlan { + switch (notification.method) { + case "thread/tokenUsage/updated": + return actionPlan({ + type: "active-thread/token-usage-set", + tokenUsage: threadTokenUsageFromRuntimeUsage(notification.params.tokenUsage), + }); + case "account/rateLimits/updated": + return effectPlan({ + type: "apply-app-server-resource-event", + event: { type: "rate-limits-updated", preserveExistingOnFailure: true }, + }); + case "skills/changed": + return effectPlan({ type: "apply-app-server-resource-event", event: { type: "skills-changed", forceReload: true } }); + case "app/list/updated": + return effectPlan({ type: "refresh-server-diagnostics" }); + case "mcpServer/oauthLogin/completed": + return effectPlan({ type: "refresh-server-diagnostics", forceResourceProbes: true }); + case "mcpServer/startupStatus/updated": + return effectPlan({ + type: "apply-app-server-resource-event", + event: { + type: "mcp-startup-status-updated", + name: notification.params.name, + status: notification.params.status, + message: notification.params.error, + }, + }); + } } function planThreadLifecycle( @@ -369,172 +110,75 @@ function planThreadLifecycle( notification: ThreadLifecycleNotification, localItemId: LocalItemIdProvider, ): ChatNotificationPlan { - return planNotificationWithStateByMethod(state, notification, THREAD_LIFECYCLE_PLANNERS, localItemId); -} - -function planDiagnosticStatus(notification: DiagnosticStatusNotification): ChatNotificationPlan { - return planNotificationByMethod(notification, DIAGNOSTIC_STATUS_PLANNERS); -} - -function planUserVisibleNotice(notification: UserVisibleNoticeNotification, localItemId: LocalItemIdProvider): ChatNotificationPlan { - return planNotificationWithLocalItemIdByMethod(notification, USER_VISIBLE_NOTICE_PLANNERS, localItemId); -} - -function planNotificationByMethod( - notification: Extract, - planners: ServerNotificationPlannerMap, -): ChatNotificationPlan { - const planner = planners[notification.method]; - return planner(notification); -} - -function planNotificationWithLocalItemIdByMethod( - notification: Extract, - planners: ServerNotificationLocalPlannerMap, - localItemId: LocalItemIdProvider, -): ChatNotificationPlan { - const planner = planners[notification.method]; - return planner(notification, localItemId); -} - -function planNotificationWithStateByMethod( - state: ChatState, - notification: Extract, - planners: ServerNotificationStatePlannerMap, - localItemId: LocalItemIdProvider, -): ChatNotificationPlan { - const planner = planners[notification.method]; - return planner(state, notification, localItemId); -} - -function jsonNoticePlan( - notification: Extract }>, - localItemId: LocalItemIdProvider, -): ChatNotificationPlan { - return systemMessagePlan({ id: localItemId("system"), text: `${notification.method}: ${jsonPreview(notification.params)}` }); -} - -function autoApprovalReviewPlan( - state: ChatState, - notification: Extract, -): ChatNotificationPlan { - const reviewItem = createAutoReviewResultItem(notification.params); - return actionPlan({ - type: "message-stream/items-replaced", - items: upsertMessageStreamItemById( - messageStreamItems(state.messageStream).filter((item) => !isUnstructuredAutoReviewWarning(item)), - reviewItem, - ), - }); -} - -function startedItemPlan(item: AppServerTurnItem, turnId: string): ChatNotificationPlan { - if (shouldSuppressLifecycleItem(item)) return EMPTY_PLAN; - const streamItem = messageStreamItemFromTurnItem(item, turnId); - return streamItem ? actionPlan({ type: "message-stream/item-upserted", item: streamItem }) : EMPTY_PLAN; -} - -function completedItemPlan(item: AppServerTurnItem, turnId: string): ChatNotificationPlan { - if (item.type === "userMessage") return EMPTY_PLAN; - const streamItem = messageStreamItemFromTurnItem(item, turnId); - if (!streamItem) return EMPTY_PLAN; - return { - actions: [ - { type: "message-stream/item-upserted", item: streamItem }, - ...(streamItem.kind === "reasoning" ? ([{ type: "message-stream/reasoning-completed", turnId }] satisfies ChatAction[]) : []), - ], - effects: [], - }; -} - -function fileChangePlan(itemId: string, turnId: string, changes: readonly AppServerFileChange[], status: string): ChatNotificationPlan { - return actionPlan({ - type: "message-stream/item-upserted", - item: streamingFileChangeMessageStreamItem(itemId, turnId, normalizeFileChanges(changes), status), - }); -} - -function appendToolTextPlan( - _state: ChatState, - itemId: string, - turnId: string, - label: string, - delta: string, - kind: Extract = "tool", -): ChatNotificationPlan { - return actionPlan({ - type: "message-stream/item-text-appended", - itemId, - turnId, - label, - delta, - kind, - }); -} - -function hookRunPlan( - state: ChatState, - run: Extract["params"]["run"], - turnId: string | null, - status: string, -): ChatNotificationPlan { - const resolvedTurnId = hookRunTurnId(state, run, turnId); - const item = hookRunMessageStreamItem(run, resolvedTurnId, status); - if (!item) return EMPTY_PLAN; - const currentPendingTurnStart = pendingTurnStartForState(state); - let pendingTurnStart = currentPendingTurnStart; - if (!resolvedTurnId && currentPendingTurnStart && run.eventName === "userPromptSubmit") { - const hookIds = currentPendingTurnStart.promptSubmitHookItemIds; - pendingTurnStart = hookIds.includes(item.id) - ? currentPendingTurnStart - : { ...currentPendingTurnStart, promptSubmitHookItemIds: [...hookIds, item.id] }; + switch (notification.method) { + case "thread/started": + return threadStartedPlan(state, notification); + case "thread/archived": + return effectPlan({ type: "apply-thread-catalog-event", event: { type: "thread-archived", threadId: notification.params.threadId } }); + case "thread/deleted": + return effectPlan({ type: "apply-thread-catalog-event", event: { type: "thread-deleted", threadId: notification.params.threadId } }); + case "thread/unarchived": + return effectPlan({ + type: "apply-thread-catalog-event", + event: { type: "thread-unarchived", threadId: notification.params.threadId }, + }); + case "thread/name/updated": + return effectPlan({ + type: "apply-thread-catalog-event", + event: { + type: "thread-renamed", + threadId: notification.params.threadId, + name: normalizeExplicitThreadName(notification.params.threadName), + }, + }); + case "thread/settings/updated": + if (state.activeThread.id !== notification.params.threadId) return EMPTY_PLAN; + return actionPlan(activeThreadSettingsAppliedAction(notification.params.threadSettings)); + case "thread/goal/updated": + return threadGoalPlan(state, notification.params.threadId, notification.params.goal, localItemId); + case "thread/goal/cleared": + return threadGoalPlan(state, notification.params.threadId, null, localItemId); } - return actionPlan({ - type: "turn/pending-start-hook-upserted", - item, - pendingTurnStart, - }); } -function hookRunTurnId( +function threadStartedPlan( state: ChatState, - run: Extract["params"]["run"], - turnId: string | null, -): string | null { - if (turnId) return turnId; - if (run.eventName === "userPromptSubmit" && !pendingTurnStartForState(state)) return activeTurnId(state); - return null; + notification: Extract, +): ChatNotificationPlan { + const effects: ChatNotificationEffect[] = [ + { + type: "apply-thread-catalog-event", + event: { type: "thread-started", thread: threadFromAppServerRecord(notification.params.thread) }, + }, + ]; + if (!state.activeThread.id || state.activeThread.id === notification.params.thread.id) { + return { actions: [{ type: "active-thread/cwd-set", cwd: notification.params.thread.cwd }], effects }; + } + return { actions: [], effects }; } -function messageStreamItemsWithPendingPromptSubmitHooks(state: ChatState, turnId: string): readonly MessageStreamItem[] { - const pending = pendingTurnStartForState(state); - const items = messageStreamItems(state.messageStream); - if (!pending) return items; - return attachHookRunsToTurn(items, turnId, pending.promptSubmitHookItemIds, pending.anchorItemId); +function threadGoalPlan( + state: ChatState, + threadId: string, + goal: Extract["params"]["goal"] | null, + localItemId: LocalItemIdProvider, +): ChatNotificationPlan { + if (state.activeThread.id !== threadId) return EMPTY_PLAN; + const actions: ChatAction[] = [{ type: "active-thread/goal-set", goal }]; + const item = goalChangeItem(localItemId("goal"), state.activeThread.goal, goal); + if (item) actions.push({ type: "message-stream/item-upserted", item }); + return { actions, effects: [] }; } -function hasStructuredAutoReviewResult(items: readonly MessageStreamItem[], activeTurnId: string | null): boolean { - return items.some( - (item) => - item.kind === "reviewResult" && - Boolean(item.turnId) && - (!activeTurnId || item.turnId === activeTurnId) && - isAutoReviewText(item.text), - ); -} - -function isUnstructuredAutoReviewWarning(item: MessageStreamItem): boolean { - return item.kind === "reviewResult" && !item.turnId && isAutoReviewText(item.text); -} - -function isAutoReviewText(text: string): boolean { - return /^Auto-review\b/i.test(text.trim()); -} - -function systemMessagePlan(message: { id: string; text: string }): ChatNotificationPlan { - return actionPlan({ type: "message-stream/system-item-added", item: createSystemItem(message.id, message.text) }); +function activeTurnIdForState(state: ChatState): string | null { + const lifecycle = state.turn.lifecycle; + return lifecycle.kind === "running" ? lifecycle.turnId : null; } function actionPlan(action: ChatAction): ChatNotificationPlan { return { actions: [action], effects: [] }; } + +function effectPlan(effect: ChatNotificationEffect): ChatNotificationPlan { + return { actions: [], effects: [effect] }; +} diff --git a/src/features/chat/app-server/inbound/notification-routing.ts b/src/features/chat/app-server/inbound/notification-routing.ts index fea7ab34..11fc561b 100644 --- a/src/features/chat/app-server/inbound/notification-routing.ts +++ b/src/features/chat/app-server/inbound/notification-routing.ts @@ -60,11 +60,11 @@ const STREAM_UPDATE_NOTIFICATION_METHODS = [ "guardianWarning", ] as const; -export type StreamUpdateNotificationMethod = (typeof STREAM_UPDATE_NOTIFICATION_METHODS)[number]; +type StreamUpdateNotificationMethod = (typeof STREAM_UPDATE_NOTIFICATION_METHODS)[number]; const TURN_LIFECYCLE_NOTIFICATION_METHODS = ["turn/started", "turn/completed"] as const; -export type TurnLifecycleNotificationMethod = (typeof TURN_LIFECYCLE_NOTIFICATION_METHODS)[number]; +type TurnLifecycleNotificationMethod = (typeof TURN_LIFECYCLE_NOTIFICATION_METHODS)[number]; const THREAD_LIFECYCLE_NOTIFICATION_METHODS = [ "thread/started", @@ -77,7 +77,7 @@ const THREAD_LIFECYCLE_NOTIFICATION_METHODS = [ "thread/settings/updated", ] as const; -export type ThreadLifecycleNotificationMethod = (typeof THREAD_LIFECYCLE_NOTIFICATION_METHODS)[number]; +type ThreadLifecycleNotificationMethod = (typeof THREAD_LIFECYCLE_NOTIFICATION_METHODS)[number]; const DIAGNOSTIC_STATUS_NOTIFICATION_METHODS = [ "thread/tokenUsage/updated", @@ -88,7 +88,7 @@ const DIAGNOSTIC_STATUS_NOTIFICATION_METHODS = [ "mcpServer/startupStatus/updated", ] as const; -export type DiagnosticStatusNotificationMethod = (typeof DIAGNOSTIC_STATUS_NOTIFICATION_METHODS)[number]; +type DiagnosticStatusNotificationMethod = (typeof DIAGNOSTIC_STATUS_NOTIFICATION_METHODS)[number]; const USER_VISIBLE_NOTICE_NOTIFICATION_METHODS = [ "thread/compacted", @@ -101,7 +101,7 @@ const USER_VISIBLE_NOTICE_NOTIFICATION_METHODS = [ "windowsSandbox/setupCompleted", ] as const; -export type UserVisibleNoticeNotificationMethod = (typeof USER_VISIBLE_NOTICE_NOTIFICATION_METHODS)[number]; +type UserVisibleNoticeNotificationMethod = (typeof USER_VISIBLE_NOTICE_NOTIFICATION_METHODS)[number]; const IGNORED_SERVER_NOTIFICATION_METHODS = [ "thread/status/changed", diff --git a/src/features/chat/app-server/inbound/runtime-events.ts b/src/features/chat/app-server/inbound/runtime-events.ts new file mode 100644 index 00000000..905c93ab --- /dev/null +++ b/src/features/chat/app-server/inbound/runtime-events.ts @@ -0,0 +1,213 @@ +import type { ServerNotification } from "../../../../app-server/connection/rpc-messages"; +import { jsonPreview } from "../../../../domain/display/json-preview"; +import type { ConversationRuntimeEvent } from "../../application/conversation/runtime-events"; +import { + STREAMED_COMMAND_RUNNING_TEXT, + STREAMED_FILE_CHANGE_IN_PROGRESS_TEXT, + STREAMED_MCP_PROGRESS_LABEL, +} from "../../domain/message-stream/factories/streaming-items"; +import { createSystemItem } from "../../domain/message-stream/factories/system-items"; +import type { MessageStreamItem } from "../../domain/message-stream/items"; +import { + type AppServerFileChange, + normalizeFileChanges, + streamingFileChangeMessageStreamItem, +} from "../mappers/message-stream/file-changes"; +import { hookRunMessageStreamItem } from "../mappers/message-stream/hook-run-items"; +import { createAutoReviewResultItem, createReviewResultItem } from "../mappers/message-stream/review-result-items"; +import { taskProgressMessageStreamItem } from "../mappers/message-stream/task-progress"; +import { + type AppServerTurnItem, + completedConversationSummaryFromAppServerTurn, + messageStreamItemFromTurnItem, + messageStreamItemsFromTurns, + shouldSuppressLifecycleItem, +} from "../mappers/message-stream/turn-items"; +import type { StreamUpdateNotification, TurnLifecycleNotification, UserVisibleNoticeNotification } from "./notification-routing"; + +const MESSAGE_CONTEXT_COMPACTED = "Context compacted."; + +type RuntimeEventSource = + | StreamUpdateNotification + | TurnLifecycleNotification + | Extract + | UserVisibleNoticeNotification; + +export function conversationRuntimeEventsFromNotification( + notification: RuntimeEventSource, + localItemId: (prefix: string) => string, +): readonly ConversationRuntimeEvent[] { + switch (notification.method) { + case "item/agentMessage/delta": + return [ + { + type: "assistantDelta", + itemId: notification.params.itemId, + runId: notification.params.turnId, + delta: notification.params.delta, + completeReasoning: true, + }, + ]; + case "item/plan/delta": + return [ + { + type: "planDelta", + itemId: notification.params.itemId, + runId: notification.params.turnId, + delta: notification.params.delta, + }, + ]; + case "turn/plan/updated": + return [ + { + type: "itemUpserted", + item: taskProgressMessageStreamItem(notification.params.turnId, notification.params.explanation, notification.params.plan), + }, + ]; + case "item/reasoning/summaryTextDelta": + case "item/reasoning/textDelta": + return [ + { + type: "textDelta", + itemId: notification.params.itemId, + runId: notification.params.turnId, + label: "reasoning", + delta: notification.params.delta, + kind: "reasoning", + }, + ]; + case "item/reasoning/summaryPartAdded": + return [ + { + type: "textDelta", + itemId: notification.params.itemId, + runId: notification.params.turnId, + label: "reasoning", + delta: "", + kind: "reasoning", + }, + ]; + case "item/started": + return startedItemEvents(notification.params.item, notification.params.turnId); + case "item/completed": + return completedItemEvents(notification.params.item, notification.params.turnId); + case "item/commandExecution/outputDelta": + return [ + { + type: "itemOutputDelta", + itemId: notification.params.itemId, + runId: notification.params.turnId, + delta: notification.params.delta, + kind: "command", + fallbackText: STREAMED_COMMAND_RUNNING_TEXT, + }, + ]; + case "item/fileChange/patchUpdated": + return [ + { + type: "itemUpserted", + item: fileChangeItem(notification.params.itemId, notification.params.turnId, notification.params.changes, "inProgress"), + }, + ]; + case "item/fileChange/outputDelta": + return [ + { + type: "itemOutputDelta", + itemId: notification.params.itemId, + runId: notification.params.turnId, + delta: notification.params.delta, + kind: "fileChange", + fallbackText: STREAMED_FILE_CHANGE_IN_PROGRESS_TEXT, + }, + ]; + case "turn/diff/updated": + return [{ type: "turnDiffUpdated", runId: notification.params.turnId, diff: notification.params.diff }]; + case "hook/started": + return hookRunEvents(notification.params.run, notification.params.turnId, "running"); + case "hook/completed": + return hookRunEvents(notification.params.run, notification.params.turnId, notification.params.run.status); + case "item/mcpToolCall/progress": + return [ + { + type: "toolOutputDelta", + itemId: notification.params.itemId, + runId: notification.params.turnId, + delta: notification.params.message, + fallbackLabel: STREAMED_MCP_PROGRESS_LABEL, + }, + ]; + case "item/autoApprovalReview/started": + case "item/autoApprovalReview/completed": + return [{ type: "autoReviewUpdated", item: createAutoReviewResultItem(notification.params) }]; + case "guardianWarning": + return [{ type: "reviewWarning", item: createReviewResultItem(localItemId("review"), notification.params.message) }]; + case "turn/started": + return [ + { + type: "runStarted", + threadId: notification.params.threadId, + runId: notification.params.turn.id, + recencyAt: notification.params.turn.startedAt, + }, + ]; + case "turn/completed": + return [ + { + type: "runCompleted", + threadId: notification.params.threadId, + runId: notification.params.turn.id, + status: notification.params.turn.status, + completedItems: messageStreamItemsFromTurns([notification.params.turn]), + completedSummary: completedConversationSummaryFromAppServerTurn(notification.params.turn), + }, + ]; + case "serverRequest/resolved": + return [{ type: "requestResolved", requestId: notification.params.requestId }]; + case "thread/compacted": + return [{ type: "systemNotice", item: createSystemItem(localItemId("system"), MESSAGE_CONTEXT_COMPACTED) }]; + case "model/rerouted": + case "deprecationNotice": + case "error": + case "warning": + case "configWarning": + case "windows/worldWritableWarning": + return [jsonNoticeEvent(notification, localItemId)]; + case "windowsSandbox/setupCompleted": + return notification.params.success ? [] : [jsonNoticeEvent(notification, localItemId)]; + } +} + +function startedItemEvents(item: AppServerTurnItem, runId: string): readonly ConversationRuntimeEvent[] { + if (shouldSuppressLifecycleItem(item)) return []; + const streamItem = messageStreamItemFromTurnItem(item, runId); + return streamItem ? [{ type: "itemUpserted", item: streamItem }] : []; +} + +function completedItemEvents(item: AppServerTurnItem, runId: string): readonly ConversationRuntimeEvent[] { + if (item.type === "userMessage") return []; + const streamItem = messageStreamItemFromTurnItem(item, runId); + return streamItem ? [{ type: "itemCompleted", runId, item: streamItem }] : []; +} + +function fileChangeItem(itemId: string, runId: string, changes: readonly AppServerFileChange[], status: string): MessageStreamItem { + return streamingFileChangeMessageStreamItem(itemId, runId, normalizeFileChanges(changes), status); +} + +function hookRunEvents( + run: Extract["params"]["run"], + runId: string | null, + status: string, +): readonly ConversationRuntimeEvent[] { + const item = hookRunMessageStreamItem(run, runId, status); + return item ? [{ type: "hookRunObserved", item, runId, eventName: run.eventName }] : []; +} + +function jsonNoticeEvent( + notification: Extract }>, + localItemId: (prefix: string) => string, +): ConversationRuntimeEvent { + return { + type: "systemNotice", + item: createSystemItem(localItemId("system"), `${notification.method}: ${jsonPreview(notification.params)}`), + }; +} diff --git a/src/features/chat/application/conversation/runtime-event-plan.ts b/src/features/chat/application/conversation/runtime-event-plan.ts new file mode 100644 index 00000000..563de83b --- /dev/null +++ b/src/features/chat/application/conversation/runtime-event-plan.ts @@ -0,0 +1,233 @@ +import { reconcileCompletedTurnItems } from "../../domain/message-stream/completed-turn-reconciliation"; +import type { MessageStreamItem } from "../../domain/message-stream/items"; +import { attachHookRunsToTurn, completeReasoningItems, upsertMessageStreamItemById } from "../../domain/message-stream/updates"; +import { messageStreamItems } from "../state/message-stream"; +import { type ChatAction, type ChatState, chatReducer } from "../state/root-reducer"; +import type { ConversationRuntimeEvent } from "./runtime-events"; +import { activeTurnId, pendingTurnStart as pendingTurnStartForState } from "./turn-state"; + +export type ConversationRuntimeEffect = + | { type: "refresh-threads" } + | { type: "maybe-name-thread"; threadId: string; turnId: string; completedSummary: ConversationRuntimeEventCompletedSummary } + | { type: "thread-recency-touched"; threadId: string; recencyAt: number | null }; + +type ConversationRuntimeEventCompletedSummary = Extract["completedSummary"]; + +export interface ConversationRuntimePlan { + actions: readonly ChatAction[]; + effects: readonly ConversationRuntimeEffect[]; +} + +const EMPTY_PLAN: ConversationRuntimePlan = { actions: [], effects: [] }; + +export function planConversationRuntimeEvents(state: ChatState, events: readonly ConversationRuntimeEvent[]): ConversationRuntimePlan { + let currentState = state; + const actions: ChatAction[] = []; + const effects: ConversationRuntimeEffect[] = []; + for (const event of events) { + const plan = planConversationRuntimeEvent(currentState, event); + actions.push(...plan.actions); + effects.push(...plan.effects); + currentState = reducePlannedActions(currentState, plan.actions); + } + return actions.length === 0 && effects.length === 0 ? EMPTY_PLAN : { actions, effects }; +} + +function planConversationRuntimeEvent(state: ChatState, event: ConversationRuntimeEvent): ConversationRuntimePlan { + switch (event.type) { + case "assistantDelta": + return actionPlan({ + type: "message-stream/assistant-delta-appended", + itemId: event.itemId, + turnId: event.runId, + delta: event.delta, + completeReasoning: event.completeReasoning, + }); + case "planDelta": + return actionPlan({ + type: "message-stream/plan-delta-appended", + itemId: event.itemId, + turnId: event.runId, + delta: event.delta, + }); + case "textDelta": + return actionPlan({ + type: "message-stream/item-text-appended", + itemId: event.itemId, + turnId: event.runId, + label: event.label, + delta: event.delta, + kind: event.kind, + }); + case "toolOutputDelta": + return actionPlan({ + type: "message-stream/tool-output-appended", + itemId: event.itemId, + turnId: event.runId, + delta: event.delta, + fallbackLabel: event.fallbackLabel, + }); + case "itemOutputDelta": + return actionPlan({ + type: "message-stream/item-output-appended", + itemId: event.itemId, + turnId: event.runId, + delta: event.delta, + kind: event.kind, + fallbackText: event.fallbackText, + }); + case "itemUpserted": + return actionPlan({ type: "message-stream/item-upserted", item: event.item }); + case "itemCompleted": + return completedItemPlan(event.item, event.runId); + case "autoReviewUpdated": + return autoReviewUpdatedPlan(state, event.item); + case "runStarted": + return runStartedPlan(state, event); + case "runCompleted": + return runCompletedPlan(state, event); + case "turnDiffUpdated": + return actionPlan({ type: "message-stream/turn-diff-updated", turnId: event.runId, diff: event.diff }); + case "hookRunObserved": + return hookRunPlan(state, event); + case "requestResolved": + return actionPlan({ type: "request/resolved", requestId: event.requestId }); + case "reviewWarning": + return reviewWarningPlan(state, event.item); + case "systemNotice": + return actionPlan({ type: "message-stream/system-item-added", item: event.item }); + } +} + +function runStartedPlan(state: ChatState, event: Extract): ConversationRuntimePlan { + return { + actions: [ + { + type: "turn/started", + threadId: event.threadId, + turnId: event.runId, + items: messageStreamItemsWithPendingPromptSubmitHooks(state, event.runId), + }, + ], + effects: [ + { + type: "thread-recency-touched", + threadId: event.threadId, + recencyAt: event.recencyAt, + }, + ], + }; +} + +function runCompletedPlan(state: ChatState, event: Extract): ConversationRuntimePlan { + if (activeTurnId(state) !== event.runId) return EMPTY_PLAN; + return { + actions: [ + { + type: "turn/completed", + turnId: event.runId, + status: event.status, + items: completeReasoningItems( + reconcileCompletedTurnItems({ + currentItems: messageStreamItems(state.messageStream), + completedTurnId: event.runId, + turnItems: event.completedItems, + }), + event.runId, + ), + }, + ], + effects: [ + { type: "maybe-name-thread", threadId: event.threadId, turnId: event.runId, completedSummary: event.completedSummary }, + { type: "refresh-threads" }, + ], + }; +} + +function completedItemPlan(item: MessageStreamItem, runId: string): ConversationRuntimePlan { + return { + actions: [ + { type: "message-stream/item-upserted", item }, + ...(item.kind === "reasoning" ? ([{ type: "message-stream/reasoning-completed", turnId: runId }] satisfies ChatAction[]) : []), + ], + effects: [], + }; +} + +function hookRunPlan(state: ChatState, event: Extract): ConversationRuntimePlan { + const resolvedRunId = hookRunId(state, event); + const item = resolvedRunId ? { ...event.item, turnId: resolvedRunId } : event.item; + const currentPendingTurnStart = pendingTurnStartForState(state); + let pendingTurnStart = currentPendingTurnStart; + if (!resolvedRunId && currentPendingTurnStart && event.eventName === "userPromptSubmit") { + const hookIds = currentPendingTurnStart.promptSubmitHookItemIds; + pendingTurnStart = hookIds.includes(item.id) + ? currentPendingTurnStart + : { ...currentPendingTurnStart, promptSubmitHookItemIds: [...hookIds, item.id] }; + } + return actionPlan({ + type: "turn/pending-start-hook-upserted", + item, + pendingTurnStart, + }); +} + +function hookRunId(state: ChatState, event: Extract): string | null { + if (event.runId) return event.runId; + if (event.eventName === "userPromptSubmit" && !pendingTurnStartForState(state)) return activeTurnId(state); + return null; +} + +function reviewWarningPlan(state: ChatState, item: MessageStreamItem): ConversationRuntimePlan { + if ( + isUnstructuredAutoReviewWarning(item) && + hasStructuredAutoReviewResult(messageStreamItems(state.messageStream), activeTurnId(state)) + ) { + return EMPTY_PLAN; + } + return actionPlan({ type: "message-stream/item-upserted", item }); +} + +function autoReviewUpdatedPlan(state: ChatState, item: MessageStreamItem): ConversationRuntimePlan { + return actionPlan({ + type: "message-stream/items-replaced", + items: upsertMessageStreamItemById( + messageStreamItems(state.messageStream).filter((currentItem) => !isUnstructuredAutoReviewWarning(currentItem)), + item, + ), + }); +} + +function messageStreamItemsWithPendingPromptSubmitHooks(state: ChatState, runId: string): readonly MessageStreamItem[] { + const pending = pendingTurnStartForState(state); + const items = messageStreamItems(state.messageStream); + if (!pending) return items; + return attachHookRunsToTurn(items, runId, pending.promptSubmitHookItemIds, pending.anchorItemId); +} + +function hasStructuredAutoReviewResult(items: readonly MessageStreamItem[], activeRunId: string | null): boolean { + return items.some( + (item) => + item.kind === "reviewResult" && Boolean(item.turnId) && (!activeRunId || item.turnId === activeRunId) && isAutoReviewText(item.text), + ); +} + +function isUnstructuredAutoReviewWarning(item: MessageStreamItem): boolean { + return item.kind === "reviewResult" && !item.turnId && isAutoReviewText(item.text); +} + +function isAutoReviewText(text: string): boolean { + return /^Auto-review\b/i.test(text.trim()); +} + +function reducePlannedActions(state: ChatState, actions: readonly ChatAction[]): ChatState { + return actions.reduce(reducePlannedAction, state); +} + +function reducePlannedAction(state: ChatState, action: ChatAction): ChatState { + return chatReducer(state, action); +} + +function actionPlan(action: ChatAction): ConversationRuntimePlan { + return { actions: [action], effects: [] }; +} diff --git a/src/features/chat/application/conversation/runtime-events.ts b/src/features/chat/application/conversation/runtime-events.ts new file mode 100644 index 00000000..23d900ac --- /dev/null +++ b/src/features/chat/application/conversation/runtime-events.ts @@ -0,0 +1,94 @@ +import type { PendingRequestId } from "../../../../domain/pending-requests/model"; +import type { ThreadConversationSummary } from "../../../../domain/threads/transcript"; +import type { MessageStreamItem } from "../../domain/message-stream/items"; + +type ConversationRuntimeTextItemKind = "tool" | "hook" | "reasoning"; +type ConversationRuntimeOutputItemKind = "command" | "fileChange"; + +export type ConversationRuntimeEvent = + | { + type: "assistantDelta"; + runId: string; + itemId: string; + delta: string; + completeReasoning: boolean; + } + | { + type: "planDelta"; + runId: string; + itemId: string; + delta: string; + } + | { + type: "textDelta"; + runId: string; + itemId: string; + label: string; + delta: string; + kind: ConversationRuntimeTextItemKind; + } + | { + type: "toolOutputDelta"; + runId: string; + itemId: string; + delta: string; + fallbackLabel: string; + } + | { + type: "itemOutputDelta"; + runId: string; + itemId: string; + delta: string; + kind: ConversationRuntimeOutputItemKind; + fallbackText: string; + } + | { + type: "itemUpserted"; + item: MessageStreamItem; + } + | { + type: "itemCompleted"; + runId: string; + item: MessageStreamItem; + } + | { + type: "autoReviewUpdated"; + item: MessageStreamItem; + } + | { + type: "runStarted"; + threadId: string; + runId: string; + recencyAt: number | null; + } + | { + type: "runCompleted"; + threadId: string; + runId: string; + status: string; + completedItems: readonly MessageStreamItem[]; + completedSummary: ThreadConversationSummary | null; + } + | { + type: "turnDiffUpdated"; + runId: string; + diff: string; + } + | { + type: "hookRunObserved"; + item: MessageStreamItem; + runId: string | null; + eventName: string; + } + | { + type: "requestResolved"; + requestId: PendingRequestId; + } + | { + type: "reviewWarning"; + item: MessageStreamItem; + } + | { + type: "systemNotice"; + item: MessageStreamItem; + }; diff --git a/src/features/chat/application/state/ui-state.ts b/src/features/chat/application/state/ui-state.ts index f3f174f6..b913a0ec 100644 --- a/src/features/chat/application/state/ui-state.ts +++ b/src/features/chat/application/state/ui-state.ts @@ -7,7 +7,7 @@ import { type ThreadRenameLifecycleState, threadRenameGenerationStillActive, transitionThreadRenameLifecycleState, -} from "../../../threads/list/rename-lifecycle"; +} from "../../../../domain/threads/rename-lifecycle"; import type { DisclosureSetAction } from "./actions"; import { patchObject } from "./patch"; diff --git a/src/features/threads-view/state.ts b/src/features/threads-view/state.ts index 79b4b123..aa97bef0 100644 --- a/src/features/threads-view/state.ts +++ b/src/features/threads-view/state.ts @@ -6,7 +6,7 @@ import { type ThreadRenameActiveState, type ThreadRenameGeneratingState, transitionThreadRenameLifecycleState, -} from "../threads/list/rename-lifecycle"; +} from "../../domain/threads/rename-lifecycle"; import { type ThreadRowCoreProjection, threadRowCoreProjection } from "../threads/list/row-projection"; type ThreadsLiveStatus = "pending" | "running" | "open"; diff --git a/src/features/threads/list/row-projection.ts b/src/features/threads/list/row-projection.ts index f3e9f4d9..03f1298a 100644 --- a/src/features/threads/list/row-projection.ts +++ b/src/features/threads/list/row-projection.ts @@ -1,6 +1,6 @@ import type { Thread } from "../../../domain/threads/model"; +import type { ThreadRenameActiveState } from "../../../domain/threads/rename-lifecycle"; import { threadDisplayTitle, threadRenameDraftTitle } from "../../../domain/threads/title"; -import type { ThreadRenameActiveState } from "./rename-lifecycle"; interface ThreadRowCoreRenameProjection { readonly active: boolean; diff --git a/tests/features/threads/list/rename-lifecycle.test.ts b/tests/domain/threads/rename-lifecycle.test.ts similarity index 97% rename from tests/features/threads/list/rename-lifecycle.test.ts rename to tests/domain/threads/rename-lifecycle.test.ts index 68e4e4c5..fce9e488 100644 --- a/tests/features/threads/list/rename-lifecycle.test.ts +++ b/tests/domain/threads/rename-lifecycle.test.ts @@ -5,7 +5,7 @@ import { type ThreadRenameGeneratingState, type ThreadRenameLifecycleState, transitionThreadRenameLifecycleState, -} from "../../../../src/features/threads/list/rename-lifecycle"; +} from "../../../src/domain/threads/rename-lifecycle"; describe("thread rename lifecycle", () => { it("keeps late generation callbacks scoped to the active unchanged generation", () => { diff --git a/tests/features/chat/app-server/inbound/routing.test.ts b/tests/features/chat/app-server/inbound/routing.test.ts index d03166cc..edebf403 100644 --- a/tests/features/chat/app-server/inbound/routing.test.ts +++ b/tests/features/chat/app-server/inbound/routing.test.ts @@ -70,6 +70,21 @@ describe("chat inbound routing", () => { expectNotificationRouteKind(notification, "threadLifecycle", { activeThreadId: "thread-other", activeTurnId: "turn-active" }); }); + it("translates run recency runtime effects to thread catalog events at the inbound boundary", () => { + let state = chatStateFixture(); + state = chatStateWith(state, { activeThread: { id: "thread-active" } }); + state = chatStateWith(state, { turn: { lifecycle: { kind: "running", turnId: "turn-active" } } }); + + const plan = planChatNotification(state, turnStartedNotification(), (prefix) => `${prefix}-1`); + + expect(plan.effects).toEqual([ + { + type: "apply-thread-catalog-event", + event: { type: "thread-touched", threadId: "thread-active", recencyAt: null }, + }, + ]); + }); + it.each([ { name: "command approval", request: commandApprovalRequest(), kind: "approval" }, { name: "file change approval", request: fileChangeApprovalRequest(), kind: "approval" }, diff --git a/tests/features/chat/app-server/inbound/runtime-events.test.ts b/tests/features/chat/app-server/inbound/runtime-events.test.ts new file mode 100644 index 00000000..c5139c6b --- /dev/null +++ b/tests/features/chat/app-server/inbound/runtime-events.test.ts @@ -0,0 +1,56 @@ +import { describe, expect, it } from "vitest"; +import type { ServerNotification } from "../../../../../src/app-server/connection/rpc-messages"; +import { conversationRuntimeEventsFromNotification } from "../../../../../src/features/chat/app-server/inbound/runtime-events"; + +describe("app-server conversation runtime event mapping", () => { + it("maps assistant deltas to panel-owned runtime events", () => { + const notification = { + method: "item/agentMessage/delta", + params: { threadId: "thread-active", turnId: "turn-active", itemId: "a1", delta: "hello" }, + } satisfies Extract; + + expect(conversationRuntimeEventsFromNotification(notification, (prefix) => `${prefix}-1`)).toEqual([ + { type: "assistantDelta", runId: "turn-active", itemId: "a1", delta: "hello", completeReasoning: true }, + ]); + }); + + it("maps completed turns to completed run snapshots", () => { + const notification = { + method: "turn/completed", + params: { + threadId: "thread-active", + turn: { + id: "turn-active", + status: "completed", + error: null, + startedAt: 1, + completedAt: 2, + durationMs: 1, + itemsView: "full", + items: [ + { type: "userMessage", id: "u1", clientId: null, content: [{ type: "text", text: "hello", text_elements: [] }] }, + { type: "agentMessage", id: "a1", text: "done", phase: "final_answer", memoryCitation: null }, + ], + }, + }, + } satisfies Extract; + + const events = conversationRuntimeEventsFromNotification(notification, (prefix) => `${prefix}-1`); + + expect(events).toEqual([ + expect.objectContaining({ + type: "runCompleted", + threadId: "thread-active", + runId: "turn-active", + status: "completed", + completedSummary: { userText: "hello", assistantText: "done" }, + }), + ]); + expect(events[0]).toMatchObject({ + completedItems: [ + expect.objectContaining({ id: "u1", kind: "message", role: "user", text: "hello" }), + expect.objectContaining({ id: "a1", kind: "message", role: "assistant", text: "done" }), + ], + }); + }); +}); diff --git a/tests/features/chat/application/conversation/runtime-event-plan.test.ts b/tests/features/chat/application/conversation/runtime-event-plan.test.ts new file mode 100644 index 00000000..8519f4ba --- /dev/null +++ b/tests/features/chat/application/conversation/runtime-event-plan.test.ts @@ -0,0 +1,102 @@ +import { describe, expect, it } from "vitest"; +import { planConversationRuntimeEvents } from "../../../../../src/features/chat/application/conversation/runtime-event-plan"; +import type { ConversationRuntimeEvent } from "../../../../../src/features/chat/application/conversation/runtime-events"; +import { type ChatAction, type ChatState, chatReducer } from "../../../../../src/features/chat/application/state/root-reducer"; +import type { MessageStreamItem } from "../../../../../src/features/chat/domain/message-stream/items"; +import { chatStateMessageStreamItems, withChatStateMessageStreamItems } from "../../support/message-stream"; +import { chatStateFixture, chatStateWith } from "../../support/state"; + +function activeRunningState(): ChatState { + let state = chatStateFixture(); + state = chatStateWith(state, { activeThread: { id: "thread-active" } }); + return chatStateWith(state, { turn: { lifecycle: { kind: "running", turnId: "turn-active" } } }); +} + +function applyActions(state: ChatState, actions: readonly ChatAction[]): ChatState { + return actions.reduce(chatReducer, state); +} + +describe("ConversationRuntimeEvent planner", () => { + it("keeps run recency updates as conversation-owned effects", () => { + const state = chatStateWith(chatStateFixture(), { activeThread: { id: "thread-active" } }); + + const plan = planConversationRuntimeEvents(state, [ + { type: "runStarted", threadId: "thread-active", runId: "turn-active", recencyAt: 123 }, + ]); + + expect(plan.effects).toEqual([{ type: "thread-recency-touched", threadId: "thread-active", recencyAt: 123 }]); + }); + + it("reconciles completed run snapshots with optimistic local user messages", () => { + let state = activeRunningState(); + state = withChatStateMessageStreamItems(state, [ + { id: "local-user-1", kind: "message", messageKind: "user", role: "user", text: "hello", turnId: "turn-active" }, + ]); + const events: ConversationRuntimeEvent[] = [ + { + type: "runCompleted", + threadId: "thread-active", + runId: "turn-active", + status: "completed", + completedSummary: { userText: "hello", assistantText: "done" }, + completedItems: [ + { + id: "u1", + sourceItemId: "u1", + kind: "message", + messageKind: "user", + role: "user", + text: "hello", + clientId: "local-user-1", + turnId: "turn-active", + }, + { + id: "a1", + sourceItemId: "a1", + kind: "message", + messageKind: "assistantResponse", + role: "assistant", + text: "done", + messageState: "completed", + turnId: "turn-active", + }, + ], + }, + ]; + + const plan = planConversationRuntimeEvents(state, events); + const next = applyActions(state, plan.actions); + + expect(chatStateMessageStreamItems(next).map((item) => item.id)).toEqual(["u1", "a1"]); + expect(plan.effects).toEqual([ + { + type: "maybe-name-thread", + threadId: "thread-active", + turnId: "turn-active", + completedSummary: { userText: "hello", assistantText: "done" }, + }, + { type: "refresh-threads" }, + ]); + }); + + it("upserts structured auto-review results without dropping unrelated stream items", () => { + let state = activeRunningState(); + state = withChatStateMessageStreamItems(state, [ + { id: "m1", kind: "message", messageKind: "assistantResponse", role: "assistant", text: "working", messageState: "completed" }, + { id: "warning-1", kind: "reviewResult", role: "tool", text: "Auto-review warning", executionState: "completed" }, + ]); + const item: MessageStreamItem = { + id: "review-1", + kind: "reviewResult", + role: "tool", + text: "Auto-review approved", + turnId: "turn-active", + executionState: "completed", + }; + + const plan = planConversationRuntimeEvents(state, [{ type: "autoReviewUpdated", item }]); + const next = applyActions(state, plan.actions); + + expect(chatStateMessageStreamItems(next).map((streamItem) => streamItem.id)).toEqual(["m1", "review-1"]); + }); +}); diff --git a/tests/scripts/grit-policy.test.mjs b/tests/scripts/grit-policy.test.mjs index bca138b0..d9845a8e 100644 --- a/tests/scripts/grit-policy.test.mjs +++ b/tests/scripts/grit-policy.test.mjs @@ -22,7 +22,7 @@ const RESPONSIBILITY_ROOT_MODULE_FILE_MESSAGE = const APP_SERVER_SUBFOLDER_ROOT_IMPORT_MESSAGE = "App-server subfolders must not import sibling root modules; move the dependency into a responsibility subfolder."; const CHAT_APPLICATION_OUTER_LAYER_MESSAGE = - "Chat application modules must not import app-server, host, panel, presentation, or UI layers; expose state and workflow contracts instead."; + "Chat application modules must not import app-server, sibling feature, host, panel, presentation, or UI layers; expose state and workflow contracts instead."; const CHAT_APP_SERVER_OUTER_LAYER_MESSAGE = "Chat app-server adapters must not import chat host, panel, presentation, or UI layers."; const CHAT_WORKSPACE_BOUNDARY_MESSAGE = "Chat modules must not import workspace modules; pass workspace capabilities through chat host contracts."; @@ -501,6 +501,15 @@ export type Item = MessageStreamItem; import type { AppServerClient } from "../../../app-server/connection/client"; export type Escape = AppServerClient; +`.trimStart(), + ); + await writeFile( + path.join(cwd, "src/features/chat/application/sibling-feature.ts"), + ` +import type { ThreadRenameLifecycleState } from "../../../threads/list/rename-lifecycle"; +import type { ThreadPickerItem } from "../../../thread-picker/model"; + +export type Escape = ThreadRenameLifecycleState | ThreadPickerItem; `.trimStart(), ); await writeFile( @@ -695,6 +704,7 @@ export const value = statusText; "src/features/chat/application/outer.ts", "src/features/chat/application/allowed.ts", "src/features/chat/application/root-app-server.ts", + "src/features/chat/application/sibling-feature.ts", "src/features/chat/app-server/outer.ts", "src/features/chat/app-server/allowed.ts", "src/features/chat/host/workspace-escape.ts", @@ -725,6 +735,9 @@ export const value = statusText; ); expect(pluginDiagnostics(report, "src/features/chat/application/allowed.ts")).toEqual([]); expect(pluginMessages(report, "src/features/chat/application/root-app-server.ts")).toEqual([CHAT_APPLICATION_OUTER_LAYER_MESSAGE]); + expect(pluginMessages(report, "src/features/chat/application/sibling-feature.ts")).toEqual( + Array.from({ length: 2 }, () => CHAT_APPLICATION_OUTER_LAYER_MESSAGE), + ); expect(pluginMessages(report, "src/features/chat/app-server/outer.ts")).toEqual( Array.from({ length: 4 }, () => CHAT_APP_SERVER_OUTER_LAYER_MESSAGE), );