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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,11 @@ import * as Scope from "effect/Scope";
import * as Stream from "effect/Stream";

import * as CheckpointStore from "../../checkpointing/CheckpointStore.ts";
import { ProviderAdapterRequestError, ProviderValidationError } from "../../provider/Errors.ts";
import {
ProviderAdapterRequestError,
ProviderAdapterValidationError,
ProviderValidationError,
} from "../../provider/Errors.ts";
import {
ProviderService,
type ProviderServiceShape,
Expand Down Expand Up @@ -80,6 +84,7 @@ type MutableState = {
failRestoreAfterEffect: boolean;
failCompletionAfterCommit: boolean;
terminalRollbackFailure: boolean;
unavailableRetainedBoundary: boolean;
interruptAcknowledgementHangs: boolean;
readonly order: string[];
readonly interruptedTurnIds: Array<TurnId | undefined>;
Expand Down Expand Up @@ -119,6 +124,7 @@ function makeState(providerSendState: ProjectionTurnRetraction["providerSendStat
failRestoreAfterEffect: false,
failCompletionAfterCommit: false,
terminalRollbackFailure: false,
unavailableRetainedBoundary: false,
interruptAcknowledgementHangs: false,
order: [],
interruptedTurnIds: [],
Expand Down Expand Up @@ -333,6 +339,16 @@ async function startHarness(
},
}),
rollbackConversation: () => unsupported(),
validateRollbackConversationTo: ({ retainedTurnCount }) =>
state.unavailableRetainedBoundary
? Effect.fail(
new ProviderAdapterValidationError({
provider: "claudeAgent",
operation: "rollbackThreadTo",
issue: `Provider history has 3 turns, below retained boundary ${retainedTurnCount}.`,
}),
)
: Effect.void,
rollbackConversationTo: ({ retainedTurnCount, targetTurnId }) =>
Effect.gen(function* () {
state.order.push("rollback");
Expand Down Expand Up @@ -734,6 +750,36 @@ it("marks terminal provider rollback failure with the correlated activity shape"
await stopHarness(harness);
});

it("silently rejects an unavailable retained boundary before interrupting the turn", async () => {
const state = makeState("claimed");
state.unavailableRetainedBoundary = true;
const harness = await startHarness(state);

expect(state.interruptedTurnIds).toEqual([]);
expect(state.order).toEqual([]);
expect(state.sessionStatus).toBe("running");
expect(state.row.status).toBe("failed");
expect(
state.dispatched.find(
(command) =>
command.type === "thread.activity.append" &&
command.activity.kind === "turn.retract.failed",
),
).toMatchObject({
type: "thread.activity.append",
activity: {
payload: {
requestId: REQUEST_ID,
stage: "provider-rollback",
retryable: false,
silent: true,
},
},
});

await stopHarness(harness);
});

it.layer(NodeServices.layer)("first-message completion integration", (it) => {
it.effect("produces reverted and deleted atomically through the WO4a decider", () =>
Effect.gen(function* () {
Expand Down
26 changes: 26 additions & 0 deletions apps/server/src/orchestration/Layers/TurnRetractionReactor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -58,6 +58,7 @@ type StageFailure = {
readonly stage: RetractionStage;
readonly retryable: boolean;
readonly detail: string;
readonly silent?: boolean;
};

const terminalProviderErrorSchemas = [
Expand All @@ -75,6 +76,12 @@ const isTerminalProviderError = (error: unknown): boolean =>
const failureDetail = (error: unknown): string =>
error instanceof Error ? error.message : String(error);

const isProviderAdapterValidationError = Schema.is(ProviderAdapterValidationError);
const isUnavailableRetainedBoundary = (error: unknown): boolean =>
isProviderAdapterValidationError(error) &&
error.operation === "rollbackThreadTo" &&
/^Provider history has \d+ turns, below retained boundary \d+\.$/.test(error.issue);

export class TurnRetractionRetryTicks extends Context.Reference<Stream.Stream<void>>(
"t3/orchestration/Layers/TurnRetractionReactor/TurnRetractionRetryTicks",
{
Expand Down Expand Up @@ -193,6 +200,7 @@ export const makeTurnRetractionReactor = Effect.gen(function* () {
stage: failure.stage,
retryable: failure.retryable,
detail: failure.detail,
...(failure.silent ? { silent: true } : {}),
},
turnId: row.targetTurnId,
createdAt,
Expand Down Expand Up @@ -462,6 +470,23 @@ export const makeTurnRetractionReactor = Effect.gen(function* () {
return;
}

if (providerService.validateRollbackConversationTo) {
yield* providerService
.validateRollbackConversationTo({
threadId: row.threadId,
retainedTurnCount: row.baselineTurnCount,
targetTurnId,
})
.pipe(
Effect.mapError((error) => ({
stage: "provider-rollback" as const,
retryable: !isTerminalProviderError(error),
detail: failureDetail(error),
...(isUnavailableRetainedBoundary(error) ? { silent: true } : {}),
})),
);
}

const nowMillis = DateTime.toEpochMillis(yield* DateTime.now);
const priorAttempt = readInterruptAttempt(row.requestId, targetTurnId);
const retryCadenceMillis = Duration.toMillis(interruptRetryCadence);
Expand Down Expand Up @@ -531,6 +556,7 @@ export const makeTurnRetractionReactor = Effect.gen(function* () {
stage: "provider-rollback" as const,
retryable: !isTerminalProviderError(error),
detail: failureDetail(error),
...(isUnavailableRetainedBoundary(error) ? { silent: true } : {}),
})),
);
yield* restoreFilesystem(row, false);
Expand Down
8 changes: 8 additions & 0 deletions apps/server/src/provider/Layers/ClaudeAdapter.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3921,6 +3921,14 @@ describe("ClaudeAdapterLive", () => {
const repeated = yield* adapter.rollbackThreadTo(session.threadId, 1);
assert.equal(repeated.turns.length, 1);

assert.isDefined(adapter.validateRollbackThreadTo);
const validation = yield* adapter.validateRollbackThreadTo!(session.threadId, 2).pipe(
Effect.result,
);
assert.equal(validation._tag, "Failure");
const afterValidation = yield* adapter.readThread(session.threadId);
assert.equal(afterValidation.turns.length, 1);

const shorterThanTarget = yield* adapter
.rollbackThreadTo(session.threadId, 2)
.pipe(Effect.result);
Expand Down
19 changes: 16 additions & 3 deletions apps/server/src/provider/Layers/ClaudeAdapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4765,9 +4765,10 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* (
},
);

const rollbackThreadTo: NonNullable<ClaudeAdapterShape["rollbackThreadTo"]> = Effect.fn(
"rollbackThreadTo",
)(function* (threadId, retainedTurnCount) {
const validateRollbackBoundary = Effect.fn("validateClaudeRollbackBoundary")(function* (
threadId: ThreadId,
retainedTurnCount: number,
) {
const context = yield* requireSession(threadId);
if (!Number.isInteger(retainedTurnCount) || retainedTurnCount < 0) {
return yield* new ProviderAdapterValidationError({
Expand All @@ -4784,6 +4785,17 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* (
issue: `Provider history has ${lifetimeTurnCount} turns, below retained boundary ${retainedTurnCount}.`,
});
}
});

const validateRollbackThreadTo: NonNullable<ClaudeAdapterShape["validateRollbackThreadTo"]> =
validateRollbackBoundary;

const rollbackThreadTo: NonNullable<ClaudeAdapterShape["rollbackThreadTo"]> = Effect.fn(
"rollbackThreadTo",
)(function* (threadId, retainedTurnCount) {
yield* validateRollbackBoundary(threadId, retainedTurnCount);
const context = yield* requireSession(threadId);
const lifetimeTurnCount = context.sessionBaseTurnCount + context.turns.length;
const delta = lifetimeTurnCount - retainedTurnCount;
const sessionLocalTurnCount = context.turns.length;
const nextLength = sessionLocalTurnCount - Math.min(delta, sessionLocalTurnCount);
Expand Down Expand Up @@ -4896,6 +4908,7 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* (
interruptTurn,
readThread,
rollbackThread,
validateRollbackThreadTo,
rollbackThreadTo,
respondToRequest,
respondToUserInput,
Expand Down
23 changes: 23 additions & 0 deletions apps/server/src/provider/Layers/ProviderService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1140,6 +1140,28 @@ const makeProviderService = Effect.fn("makeProviderService")(function* (
);
});

const validateRollbackConversationTo: NonNullable<
ProviderServiceMethod<"validateRollbackConversationTo">
> = Effect.fn("validateRollbackConversationTo")(function* (rawInput) {
const input = yield* decodeInputOrValidationError({
operation: "ProviderService.validateRollbackConversationTo",
schema: ProviderRollbackConversationToInput,
payload: rawInput,
});
const routed = yield* resolveRoutableSession({
threadId: input.threadId,
operation: "ProviderService.validateRollbackConversationTo",
allowRecovery: true,
});
if (routed.adapter.validateRollbackThreadTo !== undefined) {
yield* routed.adapter.validateRollbackThreadTo(
routed.threadId,
input.retainedTurnCount,
input.targetTurnId,
);
}
});

const runStopAll = Effect.fn("runStopAll")(function* () {
const threadIds = yield* directory.listThreadIds();
const currentAdapters = yield* getAdapterEntries;
Expand Down Expand Up @@ -1211,6 +1233,7 @@ const makeProviderService = Effect.fn("makeProviderService")(function* (
getCapabilities,
getInstanceInfo,
rollbackConversation,
validateRollbackConversationTo,
rollbackConversationTo,
// Each access creates a fresh PubSub subscription so that multiple
// consumers (ProviderRuntimeIngestion, CheckpointReactor, etc.) each
Expand Down
11 changes: 11 additions & 0 deletions apps/server/src/provider/Services/ProviderAdapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -126,6 +126,17 @@ export interface ProviderAdapterShape<TError> {
targetTurnId?: TurnId,
) => Effect.Effect<ProviderThreadSnapshot, TError>;

/**
* Validate an absolute rollback boundary without changing provider state.
* Adapters that track history relative to an opaque resume cursor use this
* to reject an unavailable boundary before a live turn is interrupted.
*/
readonly validateRollbackThreadTo?: (
threadId: ThreadId,
retainedTurnCount: number,
targetTurnId?: TurnId,
) => Effect.Effect<void, TError>;

/**
* Stop all sessions owned by this adapter.
*/
Expand Down
7 changes: 7 additions & 0 deletions apps/server/src/provider/Services/ProviderService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -116,6 +116,13 @@ export interface ProviderServiceShape {
readonly targetTurnId?: TurnId;
}) => Effect.Effect<void, ProviderServiceError>;

/** Validate an absolute rollback boundary without mutating provider state. */
readonly validateRollbackConversationTo?: (input: {
readonly threadId: ThreadId;
readonly retainedTurnCount: number;
readonly targetTurnId?: TurnId;
}) => Effect.Effect<void, ProviderServiceError>;

/**
* Canonical provider runtime event stream.
*
Expand Down
7 changes: 4 additions & 3 deletions apps/web/src/components/ChatView.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -263,7 +263,7 @@ import {
} from "./chat/chatEscapeTrigger";
import { DraftHeroHeadline } from "./chat/DraftHeroHeadline";
import { shouldRenderEmptyThreadHero } from "./chat/emptyThreadHero";
import { findCorrelatedRetractionFailure } from "./chat/lastUserMessageRecovery";
import { findCorrelatedRetractionFailureInfo } from "./chat/lastUserMessageRecovery";
import {
deriveEffectiveSessionPresentation,
usePendingRetractionForThread,
Expand Down Expand Up @@ -1540,13 +1540,14 @@ function ChatViewContent(props: ChatViewProps) {
// depend on which route is mounted.
const isServerThread = activeServerThread !== null;
const activeThread = activeServerThread ?? localDraftThread;
const retractionFailureDetail =
const retractionFailure =
activeServerThread?.turnRetraction?.status === "failed"
? findCorrelatedRetractionFailure(
? findCorrelatedRetractionFailureInfo(
activeServerThread.activities,
activeServerThread.turnRetraction.requestId,
)
: null;
const retractionFailureDetail = retractionFailure?.silent ? null : retractionFailure?.detail;
const threadError = isServerThread
? (localServerError ??
retractionFailureDetail ??
Expand Down
66 changes: 66 additions & 0 deletions apps/web/src/components/chat/RetractionRecoveryHandoff.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,8 @@ import {
resolveRetractionRecoverySignal,
} from "./RetractionRecoveryHandoff";
import {
applyOptimisticRetractionRecoveryToThread,
rememberOptimisticRetractionComposer,
snapshotLastUserMessageRecovery,
useRetractionRecoveryStore,
} from "./lastUserMessageRecovery";
Expand Down Expand Up @@ -216,6 +218,70 @@ describe("retraction recovery handoff", () => {
expect(useComposerDraftStore.getState().getDraftSession(draftId)).toBeNull();
});

it("silently restores the pre-Esc composer for an ignored boundary rejection", async () => {
useComposerDraftStore.getState().setPrompt(sourceThreadRef, "existing draft");
rememberOptimisticRetractionComposer({ requestId, sourceThreadRef });
const recovery = await seedRecovery();
useRetractionRecoveryStore.getState().setOptimisticDestination(requestId, "thread");
applyOptimisticRetractionRecoveryToThread({
sourceThreadRef,
bundle: {
prompt: "preserve this message",
images: [],
modelSelection: {
instanceId: ProviderInstanceId.make("codex"),
model: "gpt-5.6",
},
runtimeMode: "full-access",
interactionMode: "default",
envMode: "worktree",
baseBranch: "main",
startFromOrigin: true,
},
});
expect(useComposerDraftStore.getState().getComposerDraft(sourceThreadRef)?.prompt).toBe(
"existing draft\n\npreserve this message",
);

const signal = resolveRetractionRecoverySignal({
recovery,
liveCompletion: null,
projectedRetraction: {
requestId,
messageId,
targetTurnId: null,
firstUserMessage: false,
status: "failed",
completedAt: null,
},
activities: [
{
id: "ignored-failure" as never,
tone: "error",
kind: "turn.retract.failed",
summary: "Message retract failed",
payload: { requestId, detail: "boundary unavailable", silent: true },
turnId: null,
createdAt: "2026-08-11T12:00:00.100Z",
},
],
threadStatus: "live",
threadDetailExists: true,
shellSnapshotReady: true,
sourceThreadInShell: true,
nowMs: Date.parse(createdAt) + 100,
});

expect(signal).toEqual({ kind: "ignored" });
expect(signal && applyRetractionRecoverySignal({ recovery, signal, navigate: vi.fn() })).toBe(
"thread-restored",
);
expect(useComposerDraftStore.getState().getComposerDraft(sourceThreadRef)?.prompt).toBe(
"existing draft",
);
expect(useComposerDraftStore.getState().getDraftSession(draftId)).toBeNull();
});

it("surfaces the recovery draft when failure activity outlives the source thread", async () => {
const recovery = await seedRecovery();
const navigate = vi.fn();
Expand Down
Loading
Loading