From 041f49b77c5a3143dbd7619594121f852ac88066 Mon Sep 17 00:00:00 2001 From: Adam Firestone Date: Thu, 13 Aug 2026 11:36:42 -0500 Subject: [PATCH 1/4] fix(server,web): hide first-message retractions immediately --- .../Layers/TurnRetractionReactor.test.ts | 38 +++++++++- .../Layers/TurnRetractionReactor.ts | 27 +++++++ apps/web/src/components/CommandPalette.tsx | 5 +- apps/web/src/components/LegacySidebar.tsx | 16 ++-- apps/web/src/components/Sidebar.tsx | 5 +- .../src/components/chat/DraftHeroHeadline.tsx | 5 +- .../chat/useDiscoverableThreadShells.test.ts | 75 +++++++++++++++++++ .../chat/useDiscoverableThreadShells.ts | 45 +++++++++++ .../settings/ProjectSettingsPanel.tsx | 5 +- apps/web/src/routes/_chat.index.tsx | 9 +-- 10 files changed, 206 insertions(+), 24 deletions(-) create mode 100644 apps/web/src/components/chat/useDiscoverableThreadShells.test.ts create mode 100644 apps/web/src/components/chat/useDiscoverableThreadShells.ts diff --git a/apps/server/src/orchestration/Layers/TurnRetractionReactor.test.ts b/apps/server/src/orchestration/Layers/TurnRetractionReactor.test.ts index 8df6cf98cd8b..061c10f842f1 100644 --- a/apps/server/src/orchestration/Layers/TurnRetractionReactor.test.ts +++ b/apps/server/src/orchestration/Layers/TurnRetractionReactor.test.ts @@ -79,6 +79,7 @@ type MutableState = { activeTurnId: TurnId | null; historyTurnCount: number; readonly rollbackTargetTurnIds: Array; + providerStopped: boolean; filesystemRestored: boolean; failRollbackAfterEffect: boolean; failRestoreAfterEffect: boolean; @@ -119,6 +120,7 @@ function makeState(providerSendState: ProjectionTurnRetraction["providerSendStat activeTurnId: providerSendState === "claimed" ? TURN_ID : null, historyTurnCount: 2, rollbackTargetTurnIds: [], + providerStopped: false, filesystemRestored: false, failRollbackAfterEffect: false, failRestoreAfterEffect: false, @@ -324,7 +326,13 @@ async function startHarness( }).pipe(Effect.andThen(state.interruptAcknowledgementHangs ? Effect.never : Effect.void)), respondToRequest: () => unsupported(), respondToUserInput: () => unsupported(), - stopSession: () => unsupported(), + stopSession: () => + Effect.sync(() => { + state.order.push("stop"); + state.providerStopped = true; + state.sessionStatus = "ready"; + state.activeTurnId = null; + }), listSessions: () => Effect.succeed([]), getCapabilities: () => Effect.succeed({ sessionModelSwitch: "in-session" }), getInstanceInfo: (instanceId) => @@ -451,6 +459,34 @@ it("completes a cancelled provider-send path after filesystem convergence", asyn await stopHarness(harness); }); +it("stops, restores, and completes a claimed first-message retraction without rollback", async () => { + const state = makeState("claimed"); + state.row = pendingRow("claimed", true); + state.sessionStatus = "starting"; + state.activeTurnId = null; + const harness = await startHarness(state); + + expect(state.providerStopped).toBe(true); + expect(state.filesystemRestored).toBe(true); + expect(state.row.status).toBe("completed"); + expect(state.order).toEqual(["stop", "restore", "complete"]); + expect(state.interruptedTurnIds).toEqual([]); + expect(state.rollbackTargetTurnIds).toEqual([]); + await stopHarness(harness); +}); + +it("restores and completes a cancelled first-message send without stopping a provider", async () => { + const state = makeState("cancelled"); + state.row = pendingRow("cancelled", true); + const harness = await startHarness(state); + + expect(state.providerStopped).toBe(false); + expect(state.filesystemRestored).toBe(true); + expect(state.row.status).toBe("completed"); + expect(state.order).toEqual(["restore", "complete"]); + await stopHarness(harness); +}); + it("drives claimed convergence from interrupt through a settlement event", async () => { const state = makeState("claimed"); const harness = await startHarness(state); diff --git a/apps/server/src/orchestration/Layers/TurnRetractionReactor.ts b/apps/server/src/orchestration/Layers/TurnRetractionReactor.ts index beb5a07b58ea..c4a8fab77ae2 100644 --- a/apps/server/src/orchestration/Layers/TurnRetractionReactor.ts +++ b/apps/server/src/orchestration/Layers/TurnRetractionReactor.ts @@ -50,6 +50,7 @@ type RetractionStage = | "eligibility" | "interrupt" | "settlement" + | "provider-stop" | "provider-rollback" | "checkpoint-restore" | "cleanup"; @@ -415,6 +416,32 @@ export const makeTurnRetractionReactor = Effect.gen(function* () { })), ); + // A first-message retraction deletes the thread, so there is no provider + // conversation to retain or resume. Stop the provider before restoring the + // baseline so it cannot write into the workspace after restoration, then + // let completion emit the existing reverted + deleted event pair. + if (row.firstUserMessage) { + if (row.providerSendState === "claimed") { + yield* providerService.stopSession({ threadId: row.threadId }).pipe( + Effect.mapError((error) => ({ + stage: "provider-stop" as const, + retryable: !isTerminalProviderError(error), + detail: failureDetail(error), + })), + ); + } + yield* restoreFilesystem(row, row.providerSendState === "cancelled"); + yield* dispatchCompletion(row, targetTurnId); + clearIssuedInterrupts(row.requestId); + yield* logConvergence(row, "cleanup", "completed", { + action: + row.providerSendState === "claimed" + ? "stop-provider-restore-and-delete-first-message-thread" + : "restore-and-delete-cancelled-first-message-thread", + }); + return; + } + if (row.providerSendState === "cancelled") { yield* restoreFilesystem(row, true); yield* dispatchCompletion(row, targetTurnId); diff --git a/apps/web/src/components/CommandPalette.tsx b/apps/web/src/components/CommandPalette.tsx index ad9099629681..352b0f99b94e 100644 --- a/apps/web/src/components/CommandPalette.tsx +++ b/apps/web/src/components/CommandPalette.tsx @@ -68,7 +68,7 @@ import { sourceControlEnvironment } from "../state/sourceControl"; import { useAtomCommand } from "../state/use-atom-command"; import { useAtomQueryRunner } from "../state/use-atom-query-runner"; import { useEnvironments, usePrimaryEnvironmentId } from "../state/environments"; -import { useProjects, useThreadShells } from "../state/entities"; +import { useProjects } from "../state/entities"; import { useThreadSearch } from "../state/queries"; import { resolveThreadActionProjectRef, startNewThreadFromContext } from "../lib/chatThreadActions"; import { @@ -90,6 +90,7 @@ import { getLatestThreadForProject, sortThreads } from "../lib/threadSort"; import { cn, isMacPlatform, isWindowsPlatform, newProjectId } from "../lib/utils"; import { selectThreadTerminalUiState, useTerminalUiStateStore } from "../terminalUiStateStore"; import { buildThreadRouteParams, resolveThreadRouteTarget } from "../threadRoutes"; +import { useDiscoverableThreadShells } from "./chat/useDiscoverableThreadShells"; import { applyWslEnvironmentConfiguration, parseWslUncPath, @@ -580,7 +581,7 @@ function OpenCommandPaletteDialog(props: { useHandleNewThread(); const projects = useProjects(); const projectOrder = useUiStateStore((store) => store.projectOrder); - const threads = useThreadShells(); + const threads = useDiscoverableThreadShells(); const keybindings = useAtomValue(primaryServerKeybindingsAtom); const { theme, themeHalves, resolvedTheme } = useTheme(); const providers = useAtomValue(primaryServerProvidersAtom); diff --git a/apps/web/src/components/LegacySidebar.tsx b/apps/web/src/components/LegacySidebar.tsx index 916e3214ac16..8da34c0a9f7c 100644 --- a/apps/web/src/components/LegacySidebar.tsx +++ b/apps/web/src/components/LegacySidebar.tsx @@ -76,13 +76,7 @@ import { isElectron } from "../env"; import { useOpenPrLink } from "../lib/openPullRequestLink"; import { isTerminalFocused } from "../lib/terminalFocus"; import { isMacPlatform } from "../lib/utils"; -import { - readThreadShell, - useProject, - useProjects, - useThreadShells, - useThreadShellsForProjectRefs, -} from "../state/entities"; +import { readThreadShell, useProject, useProjects } from "../state/entities"; import { selectThreadTerminalUiState, useTerminalUiStateStore } from "../terminalUiStateStore"; import { useThreadRunningTerminalIds } from "../state/terminalSessions"; import { useThreadDiscoveredPorts } from "../portDiscoveryState"; @@ -108,6 +102,10 @@ import { ensureLocalApi, readLocalApi } from "../localApi"; import { useComposerDraftStore } from "../composerDraftStore"; import { useNewThreadHandler } from "../hooks/useHandleNewThread"; import { useDesktopUpdateState } from "../state/desktopUpdate"; +import { + useDiscoverableThreadShells, + useDiscoverableThreadShellsForProjectRefs, +} from "./chat/useDiscoverableThreadShells"; import { useThreadActions } from "../hooks/useThreadActions"; import { projectEnvironment } from "../state/projects"; @@ -1182,7 +1180,7 @@ const SidebarProjectItem = memo(function SidebarProjectItem(props: SidebarProjec }, }); const openPrLink = useOpenPrLink(); - const sidebarThreads = useThreadShellsForProjectRefs(project.memberProjectRefs); + const sidebarThreads = useDiscoverableThreadShellsForProjectRefs(project.memberProjectRefs); const sidebarThreadByKey = useMemo( () => new Map( @@ -3021,7 +3019,7 @@ const SidebarProjectsContent = memo(function SidebarProjectsContent( export default function LegacySidebar() { const projects = useProjects(); - const sidebarThreads = useThreadShells(); + const sidebarThreads = useDiscoverableThreadShells(); const projectExpandedById = useUiStateStore((store) => store.projectExpandedById); const projectOrder = useUiStateStore((store) => store.projectOrder); const reorderProjects = useUiStateStore((store) => store.reorderProjects); diff --git a/apps/web/src/components/Sidebar.tsx b/apps/web/src/components/Sidebar.tsx index b7a32070b21f..38324233b5f6 100644 --- a/apps/web/src/components/Sidebar.tsx +++ b/apps/web/src/components/Sidebar.tsx @@ -104,7 +104,7 @@ import { useCopyToClipboard } from "../hooks/useCopyToClipboard"; import { useLocalStorage } from "../hooks/useLocalStorage"; import { useNowMinute } from "../hooks/useNowMinute"; import { useEnvironments, usePrimaryEnvironmentId } from "../state/environments"; -import { useProjects, useThreadShells } from "../state/entities"; +import { useProjects } from "../state/entities"; import { environmentServerConfigsAtom, primaryServerKeybindingsAtom } from "../state/server"; import { vcsEnvironment } from "../state/vcs"; import { threadEnvironment } from "../state/threads"; @@ -139,6 +139,7 @@ import { sortThreadsForSidebar, } from "./Sidebar.logic"; import { useRetractedTurnPresentationSuppressed } from "./chat/retractedTurnPresentation"; +import { useDiscoverableThreadShells } from "./chat/useDiscoverableThreadShells"; import { resolveLocalCheckoutBranchMismatch } from "./BranchToolbar.logic"; import { ThreadWorktreeIndicator, @@ -1598,7 +1599,7 @@ const SidebarSearchResultRow = memo(function SidebarSearchResultRow(props: { export default function Sidebar() { const projects = useProjects(); const projectOrder = useUiStateStore((store) => store.projectOrder); - const threads = useThreadShells(); + const threads = useDiscoverableThreadShells(); const router = useRouter(); const { isMobile, setOpenMobile } = useSidebar(); const keybindings = useAtomValue(primaryServerKeybindingsAtom); diff --git a/apps/web/src/components/chat/DraftHeroHeadline.tsx b/apps/web/src/components/chat/DraftHeroHeadline.tsx index 0fb187a54621..dad3438bd658 100644 --- a/apps/web/src/components/chat/DraftHeroHeadline.tsx +++ b/apps/web/src/components/chat/DraftHeroHeadline.tsx @@ -11,9 +11,10 @@ import { buildSidebarProjectPickerEntries, buildSidebarProjectSnapshots, } from "~/sidebarProjectGrouping"; -import { useProjects, useThreadShells } from "~/state/entities"; +import { useProjects } from "~/state/entities"; import { useEnvironments, usePrimaryEnvironmentId } from "~/state/environments"; import { sortLogicalProjectsForSidebar } from "../Sidebar.logic"; +import { useDiscoverableThreadShells } from "./useDiscoverableThreadShells"; import { Menu, MenuItem, @@ -34,7 +35,7 @@ export function DraftHeroHeadline({ activeProjectTitle, }: DraftHeroHeadlineProps) { const projects = useProjects(); - const threads = useThreadShells(); + const threads = useDiscoverableThreadShells(); const { environments } = useEnvironments(); const primaryEnvironmentId = usePrimaryEnvironmentId(); const projectGroupingSettings = useClientSettings(selectProjectGroupingSettings); diff --git a/apps/web/src/components/chat/useDiscoverableThreadShells.test.ts b/apps/web/src/components/chat/useDiscoverableThreadShells.test.ts new file mode 100644 index 000000000000..4edaa6376c9f --- /dev/null +++ b/apps/web/src/components/chat/useDiscoverableThreadShells.test.ts @@ -0,0 +1,75 @@ +import type { EnvironmentThreadShell } from "@t3tools/client-runtime/state/models"; +import { scopeProjectRef, scopeThreadRef } from "@t3tools/client-runtime/environment"; +import { + CommandId, + EnvironmentId, + MessageId, + ProjectId, + ProviderInstanceId, + ThreadId, +} from "@t3tools/contracts"; +import { describe, expect, it } from "vite-plus/test"; + +import { DraftId } from "../../composerDraftStore"; +import { filterDiscoverableThreadShells } from "./useDiscoverableThreadShells"; + +const environmentId = EnvironmentId.make("environment-1"); +const projectId = ProjectId.make("project-1"); + +function thread(id: string): EnvironmentThreadShell { + return { + environmentId, + id: ThreadId.make(id), + projectId, + title: id, + modelSelection: { instanceId: ProviderInstanceId.make("codex"), model: "gpt-5.6" }, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + latestTurn: null, + createdAt: "2026-08-13T12:00:00.000Z", + updatedAt: "2026-08-13T12:00:00.000Z", + archivedAt: null, + settledOverride: null, + settledAt: null, + session: null, + latestUserMessageAt: "2026-08-13T12:00:00.000Z", + hasPendingApprovals: false, + hasPendingUserInput: false, + hasActionableProposedPlan: false, + }; +} + +function recovery(threadId: ThreadId, firstUserMessage: boolean) { + return { + requestId: CommandId.make(`request-${threadId}`), + messageId: MessageId.make(`message-${threadId}`), + sourceThreadRef: scopeThreadRef(environmentId, threadId), + projectRef: scopeProjectRef(environmentId, projectId), + draftId: DraftId.make(`draft-${threadId}`), + createdAt: "2026-08-13T12:00:00.000Z", + firstUserMessage, + optimisticDestination: "thread" as const, + }; +} + +describe("discoverable thread shells", () => { + it("hides only threads with a pending first-message recovery", () => { + const first = thread("first-message"); + const middle = thread("mid-thread"); + const ordinary = thread("ordinary"); + + expect( + filterDiscoverableThreadShells([first, middle, ordinary], { + first: recovery(first.id, true), + middle: recovery(middle.id, false), + }), + ).toEqual([middle, ordinary]); + }); + + it("returns the original list when no first-message recovery is pending", () => { + const threads = [thread("ordinary")]; + expect(filterDiscoverableThreadShells(threads, {})).toBe(threads); + }); +}); diff --git a/apps/web/src/components/chat/useDiscoverableThreadShells.ts b/apps/web/src/components/chat/useDiscoverableThreadShells.ts new file mode 100644 index 000000000000..35b11c5fd27d --- /dev/null +++ b/apps/web/src/components/chat/useDiscoverableThreadShells.ts @@ -0,0 +1,45 @@ +import type { EnvironmentThreadShell } from "@t3tools/client-runtime/state/models"; +import type { ScopedProjectRef } from "@t3tools/contracts"; +import { scopedThreadKey } from "@t3tools/client-runtime/environment"; +import { useMemo } from "react"; + +import { useThreadShells, useThreadShellsForProjectRefs } from "../../state/entities"; +import { + type PendingRetractionRecovery, + useRetractionRecoveryStore, +} from "./lastUserMessageRecovery"; + +export function filterDiscoverableThreadShells( + threads: ReadonlyArray, + recoveries: Readonly>, +): ReadonlyArray { + const hiddenThreadKeys = new Set( + Object.values(recoveries).flatMap((recovery) => + recovery.firstUserMessage === true ? [scopedThreadKey(recovery.sourceThreadRef)] : [], + ), + ); + if (hiddenThreadKeys.size === 0) return threads; + return threads.filter( + (thread) => + !hiddenThreadKeys.has( + scopedThreadKey({ environmentId: thread.environmentId, threadId: thread.id }), + ), + ); +} + +function useDiscoverableThreads( + threads: ReadonlyArray, +): ReadonlyArray { + const recoveries = useRetractionRecoveryStore((state) => state.byRequestId); + return useMemo(() => filterDiscoverableThreadShells(threads, recoveries), [recoveries, threads]); +} + +export function useDiscoverableThreadShells(): ReadonlyArray { + return useDiscoverableThreads(useThreadShells()); +} + +export function useDiscoverableThreadShellsForProjectRefs( + refs: ReadonlyArray, +): ReadonlyArray { + return useDiscoverableThreads(useThreadShellsForProjectRefs(refs)); +} diff --git a/apps/web/src/components/settings/ProjectSettingsPanel.tsx b/apps/web/src/components/settings/ProjectSettingsPanel.tsx index 50cf9c318040..8cb9b4874bd7 100644 --- a/apps/web/src/components/settings/ProjectSettingsPanel.tsx +++ b/apps/web/src/components/settings/ProjectSettingsPanel.tsx @@ -66,11 +66,12 @@ import { type SidebarProjectSnapshot, } from "../../sidebarProjectGrouping"; import { useEnvironments, usePrimaryEnvironmentId } from "../../state/environments"; -import { useProjects, useThreadShells } from "../../state/entities"; +import { useProjects } from "../../state/entities"; import { projectEnvironment } from "../../state/projects"; import { primaryServerProvidersAtom, serverEnvironment } from "../../state/server"; import { useAtomCommand } from "../../state/use-atom-command"; import { ProviderModelPicker } from "../chat/ProviderModelPicker"; +import { useDiscoverableThreadShells } from "../chat/useDiscoverableThreadShells"; import { TraitsPicker } from "../chat/TraitsPicker"; import { ProjectFavicon } from "../ProjectFavicon"; import { @@ -307,7 +308,7 @@ function ProjectDetail({ group }: { group: SidebarProjectSnapshot }) { const updateClientSettings = useUpdateClientSettings(); const projectGroupingSettings = useClientSettings(selectProjectGroupingSettings); const serverProviders = useAtomValue(primaryServerProvidersAtom); - const threads = useThreadShells(); + const threads = useDiscoverableThreadShells(); const updateProject = useAtomCommand(projectEnvironment.update, { reportFailure: false }); const deleteProject = useAtomCommand(projectEnvironment.delete, { reportFailure: false }); const upsertKeybinding = useAtomCommand(serverEnvironment.upsertKeybinding, { diff --git a/apps/web/src/routes/_chat.index.tsx b/apps/web/src/routes/_chat.index.tsx index 12bbdf666c43..9316f893794d 100644 --- a/apps/web/src/routes/_chat.index.tsx +++ b/apps/web/src/routes/_chat.index.tsx @@ -5,15 +5,12 @@ import { useCallback, useEffect, useMemo, useRef, useState } from "react"; import { openCommandPalette } from "../commandPaletteBus"; import { sortScopedProjectsForSidebar } from "../components/Sidebar.logic"; +import { useDiscoverableThreadShells } from "../components/chat/useDiscoverableThreadShells"; import { Button } from "../components/ui/button"; import { Empty, EmptyDescription, EmptyHeader, EmptyTitle } from "../components/ui/empty"; import { SidebarInset } from "../components/ui/sidebar"; import { useNewThreadHandler } from "../hooks/useHandleNewThread"; -import { - useAllEnvironmentShellsBootstrapped, - useProjects, - useThreadShells, -} from "../state/entities"; +import { useAllEnvironmentShellsBootstrapped, useProjects } from "../state/entities"; import { useEnvironments } from "../state/environments"; import { APP_DISPLAY_NAME } from "~/branding"; import { hasCloudPublicConfig } from "~/cloud/publicConfig"; @@ -38,7 +35,7 @@ function ChatIndexRouteView() { */ function IndexDraftLanding() { const projects = useProjects(); - const threads = useThreadShells(); + const threads = useDiscoverableThreadShells(); const bootstrapped = useAllEnvironmentShellsBootstrapped(); const handleNewThread = useNewThreadHandler(); const startingRef = useRef(false); From 2ab6edb6892d57e8eef831cace9cb08570ac8d49 Mon Sep 17 00:00:00 2001 From: Adam Firestone Date: Thu, 13 Aug 2026 11:49:04 -0500 Subject: [PATCH 2/4] fix(server): settle claimed first turns before deletion --- .../Layers/TurnRetractionReactor.test.ts | 41 +++++++++++------- .../Layers/TurnRetractionReactor.ts | 42 +++++++------------ 2 files changed, 42 insertions(+), 41 deletions(-) diff --git a/apps/server/src/orchestration/Layers/TurnRetractionReactor.test.ts b/apps/server/src/orchestration/Layers/TurnRetractionReactor.test.ts index 061c10f842f1..8c4e229621d5 100644 --- a/apps/server/src/orchestration/Layers/TurnRetractionReactor.test.ts +++ b/apps/server/src/orchestration/Layers/TurnRetractionReactor.test.ts @@ -79,7 +79,6 @@ type MutableState = { activeTurnId: TurnId | null; historyTurnCount: number; readonly rollbackTargetTurnIds: Array; - providerStopped: boolean; filesystemRestored: boolean; failRollbackAfterEffect: boolean; failRestoreAfterEffect: boolean; @@ -120,7 +119,6 @@ function makeState(providerSendState: ProjectionTurnRetraction["providerSendStat activeTurnId: providerSendState === "claimed" ? TURN_ID : null, historyTurnCount: 2, rollbackTargetTurnIds: [], - providerStopped: false, filesystemRestored: false, failRollbackAfterEffect: false, failRestoreAfterEffect: false, @@ -326,13 +324,7 @@ async function startHarness( }).pipe(Effect.andThen(state.interruptAcknowledgementHangs ? Effect.never : Effect.void)), respondToRequest: () => unsupported(), respondToUserInput: () => unsupported(), - stopSession: () => - Effect.sync(() => { - state.order.push("stop"); - state.providerStopped = true; - state.sessionStatus = "ready"; - state.activeTurnId = null; - }), + stopSession: () => unsupported(), listSessions: () => Effect.succeed([]), getCapabilities: () => Effect.succeed({ sessionModelSwitch: "in-session" }), getInstanceInfo: (instanceId) => @@ -459,18 +451,40 @@ it("completes a cancelled provider-send path after filesystem convergence", asyn await stopHarness(harness); }); -it("stops, restores, and completes a claimed first-message retraction without rollback", async () => { +it("settles a claimed first-message turn before restoring and completing without rollback", async () => { const state = makeState("claimed"); state.row = pendingRow("claimed", true); state.sessionStatus = "starting"; state.activeTurnId = null; const harness = await startHarness(state); - expect(state.providerStopped).toBe(true); + expect(state.order).toEqual([]); + expect(state.row.status).toBe("requested"); + + state.sessionStatus = "running"; + state.activeTurnId = TURN_ID; + await harness.emitRuntime({ + type: "turn.started", + eventId: EventId.make("evt-first-message-turn-started"), + provider: ProviderDriverKind.make("codex"), + createdAt: NOW, + threadId: THREAD_ID, + turnId: TURN_ID, + payload: {}, + }); + await harness.runtime.runPromise(Effect.yieldNow); + await harness.runtime.runPromise(harness.reactor.drain); + expect(state.order).toEqual(["interrupt"]); + + state.sessionStatus = "ready"; + state.activeTurnId = null; + await harness.retryTick(); + await harness.runtime.runPromise(harness.reactor.drain); + expect(state.filesystemRestored).toBe(true); expect(state.row.status).toBe("completed"); - expect(state.order).toEqual(["stop", "restore", "complete"]); - expect(state.interruptedTurnIds).toEqual([]); + expect(state.order).toEqual(["interrupt", "restore", "complete"]); + expect(state.interruptedTurnIds).toEqual([TURN_ID]); expect(state.rollbackTargetTurnIds).toEqual([]); await stopHarness(harness); }); @@ -480,7 +494,6 @@ it("restores and completes a cancelled first-message send without stopping a pro state.row = pendingRow("cancelled", true); const harness = await startHarness(state); - expect(state.providerStopped).toBe(false); expect(state.filesystemRestored).toBe(true); expect(state.row.status).toBe("completed"); expect(state.order).toEqual(["restore", "complete"]); diff --git a/apps/server/src/orchestration/Layers/TurnRetractionReactor.ts b/apps/server/src/orchestration/Layers/TurnRetractionReactor.ts index c4a8fab77ae2..1d15d3d6ce73 100644 --- a/apps/server/src/orchestration/Layers/TurnRetractionReactor.ts +++ b/apps/server/src/orchestration/Layers/TurnRetractionReactor.ts @@ -50,7 +50,6 @@ type RetractionStage = | "eligibility" | "interrupt" | "settlement" - | "provider-stop" | "provider-rollback" | "checkpoint-restore" | "cleanup"; @@ -416,32 +415,6 @@ export const makeTurnRetractionReactor = Effect.gen(function* () { })), ); - // A first-message retraction deletes the thread, so there is no provider - // conversation to retain or resume. Stop the provider before restoring the - // baseline so it cannot write into the workspace after restoration, then - // let completion emit the existing reverted + deleted event pair. - if (row.firstUserMessage) { - if (row.providerSendState === "claimed") { - yield* providerService.stopSession({ threadId: row.threadId }).pipe( - Effect.mapError((error) => ({ - stage: "provider-stop" as const, - retryable: !isTerminalProviderError(error), - detail: failureDetail(error), - })), - ); - } - yield* restoreFilesystem(row, row.providerSendState === "cancelled"); - yield* dispatchCompletion(row, targetTurnId); - clearIssuedInterrupts(row.requestId); - yield* logConvergence(row, "cleanup", "completed", { - action: - row.providerSendState === "claimed" - ? "stop-provider-restore-and-delete-first-message-thread" - : "restore-and-delete-cancelled-first-message-thread", - }); - return; - } - if (row.providerSendState === "cancelled") { yield* restoreFilesystem(row, true); yield* dispatchCompletion(row, targetTurnId); @@ -572,6 +545,21 @@ export const makeTurnRetractionReactor = Effect.gen(function* () { } } + // A first-message retraction deletes the thread, so its settled provider + // conversation does not need a rollback boundary. Do not stop the session + // while ProviderCommandReactor may still be starting the send: the normal + // start/interrupt settlement above lets that durable worker finish. The + // correlated thread.deleted event owns final provider-session cleanup. + if (row.firstUserMessage) { + yield* restoreFilesystem(row, false); + yield* dispatchCompletion(row, targetTurnId); + clearIssuedInterrupts(row.requestId); + yield* logConvergence(row, "cleanup", "completed", { + action: "restore-and-delete-settled-first-message-thread", + }); + return; + } + yield* providerService .rollbackConversationTo({ threadId: row.threadId, From c47359278eb20fb743d25c3583f8f1d440ffac1b Mon Sep 17 00:00:00 2001 From: Adam Firestone Date: Thu, 13 Aug 2026 12:25:38 -0500 Subject: [PATCH 3/4] fix(codex): discard retracted transient threads --- .../Layers/CheckpointReactor.test.ts | 1 + .../Layers/ProviderCommandReactor.test.ts | 1 + .../Layers/ProviderRuntimeIngestion.test.ts | 1 + .../Layers/ThreadDeletionReactor.test.ts | 8 ++++ .../Layers/ThreadDeletionReactor.ts | 14 ++++++ .../Layers/TurnRetractionReactor.test.ts | 1 + .../src/provider/Layers/CodexAdapter.test.ts | 25 +++++++++++ .../src/provider/Layers/CodexAdapter.ts | 13 ++++++ .../Layers/CodexSessionRuntime.test.ts | 24 +++++++++++ .../provider/Layers/CodexSessionRuntime.ts | 20 +++++++++ .../provider/Layers/ProviderService.test.ts | 43 +++++++++++++++++++ .../src/provider/Layers/ProviderService.ts | 29 +++++++++++++ .../Layers/ProviderSessionReaper.test.ts | 1 + .../src/provider/Services/ProviderAdapter.ts | 7 +++ .../src/provider/Services/ProviderService.ts | 9 ++++ 15 files changed, 197 insertions(+) diff --git a/apps/server/src/orchestration/Layers/CheckpointReactor.test.ts b/apps/server/src/orchestration/Layers/CheckpointReactor.test.ts index b6bbe0b5c2d2..1b9518da9f34 100644 --- a/apps/server/src/orchestration/Layers/CheckpointReactor.test.ts +++ b/apps/server/src/orchestration/Layers/CheckpointReactor.test.ts @@ -110,6 +110,7 @@ function createProviderServiceHarness( interruptTurn: () => unsupported(), respondToRequest: () => unsupported(), respondToUserInput: () => unsupported(), + discardTransientThread: () => unsupported(), stopSession: () => unsupported(), listSessions, getCapabilities: () => Effect.succeed({ sessionModelSwitch: "in-session" }), diff --git a/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts b/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts index 57c1e6720b8d..0b0bfff5ffed 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts @@ -321,6 +321,7 @@ describe("ProviderCommandReactor", () => { interruptTurn: interruptTurn as ProviderServiceShape["interruptTurn"], respondToRequest: respondToRequest as ProviderServiceShape["respondToRequest"], respondToUserInput: respondToUserInput as ProviderServiceShape["respondToUserInput"], + discardTransientThread: () => unsupported(), stopSession: stopSession as ProviderServiceShape["stopSession"], listSessions: () => Effect.succeed(runtimeSessions), getCapabilities: (_provider) => diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index fdc9f9301e1c..2dfc4eeee886 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -185,6 +185,7 @@ function createProviderServiceHarness() { interruptTurn: () => unsupported(), respondToRequest: () => unsupported(), respondToUserInput: () => unsupported(), + discardTransientThread: () => unsupported(), stopSession: () => unsupported(), listSessions: () => Effect.succeed([...runtimeSessions]), getCapabilities: () => Effect.succeed({ sessionModelSwitch: "in-session" }), diff --git a/apps/server/src/orchestration/Layers/ThreadDeletionReactor.test.ts b/apps/server/src/orchestration/Layers/ThreadDeletionReactor.test.ts index 6623cef15d57..46a25744c906 100644 --- a/apps/server/src/orchestration/Layers/ThreadDeletionReactor.test.ts +++ b/apps/server/src/orchestration/Layers/ThreadDeletionReactor.test.ts @@ -16,6 +16,7 @@ import type { ProjectionThread } from "../../persistence/Services/ProjectionThre import { logCleanupCauseUnlessInterrupted, managedWorktreeCleanupTarget, + shouldDiscardTransientProviderThread, } from "./ThreadDeletionReactor.ts"; const threadId = ThreadId.make("thread-deletion-reactor-test"); @@ -160,3 +161,10 @@ describe("managedWorktreeCleanupTarget", () => { ).toBeNull(); }); }); + +describe("shouldDiscardTransientProviderThread", () => { + it("selects only durable first-message retraction deletions", () => { + expect(shouldDiscardTransientProviderThread(deletedEvent())).toBe(true); + expect(shouldDiscardTransientProviderThread(deletedEvent(false))).toBe(false); + }); +}); diff --git a/apps/server/src/orchestration/Layers/ThreadDeletionReactor.ts b/apps/server/src/orchestration/Layers/ThreadDeletionReactor.ts index 9a687be3afa4..f7aee2483651 100644 --- a/apps/server/src/orchestration/Layers/ThreadDeletionReactor.ts +++ b/apps/server/src/orchestration/Layers/ThreadDeletionReactor.ts @@ -22,6 +22,10 @@ import { forkParked } from "../../serverActivation.ts"; type ThreadDeletedEvent = Extract; +export function shouldDiscardTransientProviderThread(event: ThreadDeletedEvent): boolean { + return event.payload.retraction?.firstUserMessage === true; +} + export function managedWorktreeCleanupTarget(input: { readonly event: ThreadDeletedEvent; readonly thread: ProjectionThread; @@ -77,6 +81,13 @@ const make = Effect.gen(function* () { threadId, }); + const discardTransientProviderThread = (threadId: ThreadDeletedEvent["payload"]["threadId"]) => + logCleanupCauseUnlessInterrupted({ + effect: providerService.discardTransientThread({ threadId }), + message: "thread retraction cleanup skipped transient provider thread discard", + threadId, + }); + const closeThreadTerminals = (threadId: ThreadDeletedEvent["payload"]["threadId"]) => logCleanupCauseUnlessInterrupted({ effect: terminalManager.close({ threadId, deleteHistory: true }), @@ -113,6 +124,9 @@ const make = Effect.gen(function* () { event: ThreadDeletedEvent, ) { const { threadId } = event.payload; + if (shouldDiscardTransientProviderThread(event)) { + yield* discardTransientProviderThread(threadId); + } yield* stopProviderSession(threadId); yield* closeThreadTerminals(threadId); yield* removeRetractedManagedWorktree(event); diff --git a/apps/server/src/orchestration/Layers/TurnRetractionReactor.test.ts b/apps/server/src/orchestration/Layers/TurnRetractionReactor.test.ts index 8c4e229621d5..970224d236a9 100644 --- a/apps/server/src/orchestration/Layers/TurnRetractionReactor.test.ts +++ b/apps/server/src/orchestration/Layers/TurnRetractionReactor.test.ts @@ -324,6 +324,7 @@ async function startHarness( }).pipe(Effect.andThen(state.interruptAcknowledgementHangs ? Effect.never : Effect.void)), respondToRequest: () => unsupported(), respondToUserInput: () => unsupported(), + discardTransientThread: () => unsupported(), stopSession: () => unsupported(), listSessions: () => Effect.succeed([]), getCapabilities: () => Effect.succeed({ sessionModelSwitch: "in-session" }), diff --git a/apps/server/src/provider/Layers/CodexAdapter.test.ts b/apps/server/src/provider/Layers/CodexAdapter.test.ts index 2c3c7b4946f3..5f36ddea4c4b 100644 --- a/apps/server/src/provider/Layers/CodexAdapter.test.ts +++ b/apps/server/src/provider/Layers/CodexAdapter.test.ts @@ -103,6 +103,8 @@ class FakeCodexRuntime implements CodexSessionRuntimeShape { }), ); + public readonly deleteThreadImpl = vi.fn((): Promise => Promise.resolve(undefined)); + public readonly respondToRequestImpl = vi.fn( (_requestId: ApprovalRequestId, _decision: ProviderApprovalDecision): Promise => Promise.resolve(undefined), @@ -141,6 +143,8 @@ class FakeCodexRuntime implements CodexSessionRuntimeShape { return Effect.promise(() => this.rollbackThreadImpl(numTurns)); } + deleteThread = Effect.promise(() => this.deleteThreadImpl()); + respondToRequest(requestId: ApprovalRequestId, decision: ProviderApprovalDecision) { return Effect.promise(() => this.respondToRequestImpl(requestId, decision)); } @@ -309,6 +313,27 @@ const sessionErrorLayer = it.layer( ); sessionErrorLayer("CodexAdapterLive session errors", (it) => { + it.effect("discards the active provider-owned thread without stopping its session", () => + Effect.gen(function* () { + const adapter = yield* CodexAdapter; + const threadId = asThreadId("discard-transient-thread"); + yield* adapter.startSession({ + provider: ProviderDriverKind.make("codex"), + threadId, + runtimeMode: "full-access", + }); + const runtime = sessionRuntimeFactory.lastRuntime; + NodeAssert.ok(runtime); + NodeAssert.ok(adapter.discardTransientThread); + + yield* adapter.discardTransientThread(threadId); + + NodeAssert.equal(runtime.deleteThreadImpl.mock.calls.length, 1); + NodeAssert.equal(runtime.closeImpl.mock.calls.length, 0); + NodeAssert.equal(yield* adapter.hasSession(threadId), true); + }), + ); + it.effect("computes the remaining absolute rollback delta and is idempotent", () => Effect.gen(function* () { const adapter = yield* CodexAdapter; diff --git a/apps/server/src/provider/Layers/CodexAdapter.ts b/apps/server/src/provider/Layers/CodexAdapter.ts index 3697dc29ec12..fe8ca8285838 100644 --- a/apps/server/src/provider/Layers/CodexAdapter.ts +++ b/apps/server/src/provider/Layers/CodexAdapter.ts @@ -2028,6 +2028,18 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* ( yield* stopSessionInternal(session); }); + const discardTransientThread: NonNullable = ( + threadId, + ) => + requireSession(threadId).pipe( + Effect.flatMap((session) => session.runtime.deleteThread), + Effect.mapError((cause) => + cause._tag === "ProviderAdapterSessionNotFoundError" + ? cause + : mapCodexRuntimeError(threadId, "thread/delete", cause), + ), + ); + const listSessions: CodexAdapterShape["listSessions"] = () => Effect.forEach( Array.from(sessions.values()).filter((session) => !session.stopped), @@ -2066,6 +2078,7 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* ( respondToRequest, respondToUserInput, stopSession, + discardTransientThread, listSessions, hasSession, stopAll, diff --git a/apps/server/src/provider/Layers/CodexSessionRuntime.test.ts b/apps/server/src/provider/Layers/CodexSessionRuntime.test.ts index d7346a0e0dbe..7719951cf872 100644 --- a/apps/server/src/provider/Layers/CodexSessionRuntime.test.ts +++ b/apps/server/src/provider/Layers/CodexSessionRuntime.test.ts @@ -16,6 +16,7 @@ import { import { codexSessionAppServerArgs } from "./codexLaunchArgs.ts"; import { buildTurnStartParams, + deleteCodexThread, hasConfiguredMcpServer, isRecoverableThreadResumeError, openCodexThread, @@ -469,3 +470,26 @@ describe("openCodexThread", () => { }), ); }); + +describe("deleteCodexThread", () => { + it.effect("uses the provider thread id with thread/delete", () => + Effect.gen(function* () { + const calls: Array<{ method: string; payload: unknown }> = []; + const client = { + request: (method: "thread/delete", payload: { readonly threadId: string }) => { + calls.push({ method, payload }); + return Effect.succeed({}); + }, + }; + + yield* deleteCodexThread(client, "provider-thread-transient"); + + NodeAssert.deepStrictEqual(calls, [ + { + method: "thread/delete", + payload: { threadId: "provider-thread-transient" }, + }, + ]); + }), + ); +}); diff --git a/apps/server/src/provider/Layers/CodexSessionRuntime.ts b/apps/server/src/provider/Layers/CodexSessionRuntime.ts index 58c012bd63ea..94592b9e03e4 100644 --- a/apps/server/src/provider/Layers/CodexSessionRuntime.ts +++ b/apps/server/src/provider/Layers/CodexSessionRuntime.ts @@ -141,6 +141,7 @@ export interface CodexSessionRuntimeShape { readonly rollbackThread: ( numTurns: number, ) => Effect.Effect; + readonly deleteThread: Effect.Effect; readonly respondToRequest: ( requestId: ApprovalRequestId, decision: ProviderApprovalDecision, @@ -454,6 +455,22 @@ interface CodexThreadOpenClient { ) => Effect.Effect; } +interface CodexThreadDeleteClient { + readonly request: ( + method: "thread/delete", + payload: CodexRpc.ClientRequestParamsByMethod["thread/delete"], + ) => Effect.Effect< + CodexRpc.ClientRequestResponsesByMethod["thread/delete"], + CodexErrors.CodexAppServerError + >; +} + +export const deleteCodexThread = ( + client: CodexThreadDeleteClient, + threadId: string, +): Effect.Effect => + client.request("thread/delete", { threadId }).pipe(Effect.asVoid); + export const openCodexThread = (input: { readonly client: CodexThreadOpenClient; readonly threadId: ThreadId; @@ -1853,6 +1870,9 @@ export const makeCodexSessionRuntime = ( }); return parseThreadSnapshot(response); }), + deleteThread: Effect.flatMap(readProviderThreadId, (providerThreadId) => + deleteCodexThread(client, providerThreadId), + ), respondToRequest: (requestId, decision) => Effect.gen(function* () { const pending = (yield* Ref.get(pendingApprovalsRef)).get(requestId); diff --git a/apps/server/src/provider/Layers/ProviderService.test.ts b/apps/server/src/provider/Layers/ProviderService.test.ts index 9f2bc19fa23c..ccd49ad5c1b6 100644 --- a/apps/server/src/provider/Layers/ProviderService.test.ts +++ b/apps/server/src/provider/Layers/ProviderService.test.ts @@ -163,6 +163,10 @@ function makeFakeCodexAdapter(provider: ProviderDriverKind = CODEX_DRIVER) { }), ); + const discardTransientThread = vi.fn( + (_threadId: ThreadId): Effect.Effect => Effect.void, + ); + const listSessions = vi.fn( (): Effect.Effect> => Effect.sync(() => Array.from(sessions.values())), @@ -223,6 +227,7 @@ function makeFakeCodexAdapter(provider: ProviderDriverKind = CODEX_DRIVER) { respondToRequest, respondToUserInput, stopSession, + ...(provider === CODEX_DRIVER ? { discardTransientThread } : {}), listSessions, hasSession, readThread, @@ -259,6 +264,7 @@ function makeFakeCodexAdapter(provider: ProviderDriverKind = CODEX_DRIVER) { respondToRequest, respondToUserInput, stopSession, + discardTransientThread, listSessions, hasSession, readThread, @@ -921,6 +927,43 @@ routing.layer("ProviderServiceLive routing", (it) => { }), ); + it.effect("discards only provider threads whose adapter explicitly supports it", () => + Effect.gen(function* () { + const provider = yield* ProviderService.ProviderService; + const codexThreadId = asThreadId("transient-codex-thread"); + const claudeThreadId = asThreadId("transient-claude-thread"); + yield* provider.startSession(codexThreadId, { + provider: CODEX_DRIVER, + providerInstanceId: codexInstanceId, + threadId: codexThreadId, + runtimeMode: "full-access", + }); + yield* provider.startSession(claudeThreadId, { + provider: CLAUDE_AGENT_DRIVER, + providerInstanceId: claudeAgentInstanceId, + threadId: claudeThreadId, + runtimeMode: "full-access", + }); + routing.codex.discardTransientThread.mockClear(); + routing.claude.discardTransientThread.mockClear(); + + yield* provider.discardTransientThread({ threadId: codexThreadId }); + yield* provider.discardTransientThread({ threadId: claudeThreadId }); + + assert.deepEqual(routing.codex.discardTransientThread.mock.calls, [[codexThreadId]]); + assert.equal(routing.claude.discardTransientThread.mock.calls.length, 0); + + yield* provider.stopSession({ threadId: codexThreadId }); + yield* provider.stopSession({ threadId: claudeThreadId }); + routing.codex.startSession.mockClear(); + routing.codex.stopSession.mockClear(); + routing.codex.discardTransientThread.mockClear(); + routing.claude.startSession.mockClear(); + routing.claude.stopSession.mockClear(); + routing.claude.discardTransientThread.mockClear(); + }), + ); + it.effect("routes provider operations and rollback conversation", () => Effect.gen(function* () { const provider = yield* ProviderService.ProviderService; diff --git a/apps/server/src/provider/Layers/ProviderService.ts b/apps/server/src/provider/Layers/ProviderService.ts index 0414d05d67ca..8fe2c7ac3461 100644 --- a/apps/server/src/provider/Layers/ProviderService.ts +++ b/apps/server/src/provider/Layers/ProviderService.ts @@ -938,6 +938,34 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( }, ); + const discardTransientThread: ProviderServiceMethod<"discardTransientThread"> = Effect.fn( + "discardTransientThread", + )(function* (rawInput) { + const input = yield* decodeInputOrValidationError({ + operation: "ProviderService.discardTransientThread", + schema: ProviderStopSessionInput, + payload: rawInput, + }); + const routed = yield* resolveRoutableSession({ + threadId: input.threadId, + operation: "ProviderService.discardTransientThread", + allowRecovery: false, + }); + yield* Effect.annotateCurrentSpan({ + "provider.operation": "discard-transient-thread", + "provider.kind": routed.adapter.provider, + "provider.thread_id": input.threadId, + }); + if (!routed.isActive || routed.adapter.discardTransientThread === undefined) { + return; + } + yield* routed.adapter.discardTransientThread(routed.threadId); + yield* analytics.record("provider.thread.discarded", { + provider: routed.adapter.provider, + reason: "transient", + }); + }); + const listSessions: ProviderServiceMethod<"listSessions"> = Effect.fn("listSessions")( function* () { const currentAdapters = yield* getAdapterEntries; @@ -1247,6 +1275,7 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( interruptTurn, respondToRequest, respondToUserInput, + discardTransientThread, stopSession, listSessions, getCapabilities, diff --git a/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts b/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts index b424d91bee00..ebdb65721848 100644 --- a/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts +++ b/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts @@ -167,6 +167,7 @@ describe("ProviderSessionReaper", () => { interruptTurn: () => unsupported(), respondToRequest: () => unsupported(), respondToUserInput: () => unsupported(), + discardTransientThread: () => unsupported(), stopSession, listSessions: () => Effect.succeed([]), getCapabilities: () => Effect.succeed({ sessionModelSwitch: "in-session" }), diff --git a/apps/server/src/provider/Services/ProviderAdapter.ts b/apps/server/src/provider/Services/ProviderAdapter.ts index 293d24559792..ee712aed6e71 100644 --- a/apps/server/src/provider/Services/ProviderAdapter.ts +++ b/apps/server/src/provider/Services/ProviderAdapter.ts @@ -91,6 +91,13 @@ export interface ProviderAdapterShape { */ readonly stopSession: (threadId: ThreadId) => Effect.Effect; + /** + * Permanently discard a provider-owned thread that T3 has already decided + * is transient. This is intentionally optional: callers must have their own + * durable proof that the conversation should not be retained. + */ + readonly discardTransientThread?: (threadId: ThreadId) => Effect.Effect; + /** * List currently active provider sessions for this adapter. */ diff --git a/apps/server/src/provider/Services/ProviderService.ts b/apps/server/src/provider/Services/ProviderService.ts index 7092b83a39b7..4b53021bdb76 100644 --- a/apps/server/src/provider/Services/ProviderService.ts +++ b/apps/server/src/provider/Services/ProviderService.ts @@ -80,6 +80,15 @@ export interface ProviderServiceShape { input: ProviderStopSessionInput, ) => Effect.Effect; + /** + * Permanently discard the provider-owned thread for a conversation that T3 + * has already durably identified as transient. Unsupported providers and + * inactive sessions are left untouched. + */ + readonly discardTransientThread: ( + input: ProviderStopSessionInput, + ) => Effect.Effect; + /** * List active provider sessions. * From 4e7125f75e87bd3464337650e94850cec3b1db98 Mon Sep 17 00:00:00 2001 From: Adam Firestone Date: Thu, 13 Aug 2026 12:42:38 -0500 Subject: [PATCH 4/4] test(server): validate persisted rollback payload --- apps/server/src/provider/Layers/ProviderService.test.ts | 8 +++++++- 1 file changed, 7 insertions(+), 1 deletion(-) diff --git a/apps/server/src/provider/Layers/ProviderService.test.ts b/apps/server/src/provider/Layers/ProviderService.test.ts index ccd49ad5c1b6..8ac49fcd008a 100644 --- a/apps/server/src/provider/Layers/ProviderService.test.ts +++ b/apps/server/src/provider/Layers/ProviderService.test.ts @@ -910,8 +910,14 @@ routing.layer("ProviderServiceLive routing", (it) => { assert.equal(Option.isSome(persisted), true); if (Option.isSome(persisted)) { assert.deepEqual(persisted.value.resumeCursor, rolledBackCursor); + const runtimePayload = persisted.value.runtimePayload; assert.equal( - persisted.value.runtimePayload.lastRuntimeEvent, + runtimePayload !== null && + typeof runtimePayload === "object" && + !Array.isArray(runtimePayload) && + "lastRuntimeEvent" in runtimePayload + ? runtimePayload.lastRuntimeEvent + : undefined, "provider.rollbackConversationTo", ); }