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
558 changes: 558 additions & 0 deletions apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.test.ts

Large diffs are not rendered by default.

123 changes: 103 additions & 20 deletions apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,7 @@ import type {
ThreadId,
} from "@t3tools/contracts";
import * as CodexClient from "effect-codex-app-server/client";
import * as CodexErrors from "effect-codex-app-server/errors";
import * as CodexSchema from "effect-codex-app-server/schema";
import * as Context from "effect/Context";
import * as DateTime from "effect/DateTime";
Expand Down Expand Up @@ -831,6 +832,58 @@ const resolveCodexForkRollbackTurnCount = Effect.fn("CodexAdapterV2.resolveForkR
},
);

/**
* Prefer a native `thread/fork` turn boundary over the fork-then-rollback
* fallback. `lastTurnId` is inclusive on the Codex side, so a fork requested at
* the selected turn omits every later turn atomically. That matters for
* paginated threads, which reject `thread/rollback` entirely. The count
* fallback remains only for source turns that predate native turn references.
*/
export const resolveCodexForkBoundary = Effect.fn("CodexAdapterV2.resolveForkBoundary")(function* (
input: ProviderAdapterV2ForkThreadInput,
) {
const rollbackTurnCount = yield* resolveCodexForkRollbackTurnCount(input);
if (input.providerTurnId === undefined || input.sourceProviderTurns === undefined) {
return { lastTurnId: undefined, rollbackTurnCount };
}

const boundaryTurn = providerTurnsForThread(
input.sourceProviderTurns,
input.sourceProviderThread,
).find((turn) => turn.id === input.providerTurnId);
const nativeTurnId = boundaryTurn?.nativeTurnRef?.nativeId;
if (nativeTurnId === null || nativeTurnId === undefined) {
return { lastTurnId: undefined, rollbackTurnCount };
}

return { lastTurnId: nativeTurnId, rollbackTurnCount: 0 };
});

/**
* The generated `thread/read` response schema does not surface `historyMode`,
* so the probe goes through the raw request channel with a permissive decode
* (mirrors the V1 session runtime's paginated-history detection).
*/
const CodexThreadHistoryMetadata = Schema.Struct({
thread: Schema.Struct({
historyMode: Schema.optionalKey(Schema.Literals(["legacy", "paginated"])),
}),
});
const decodeCodexThreadHistoryMetadata = Schema.decodeUnknownEffect(CodexThreadHistoryMetadata);

const readCodexThreadHistoryMode = Effect.fn("CodexAdapterV2.readThreadHistoryMode")(function* (
raw: Pick<CodexClient.CodexAppServerClient["Service"]["raw"], "request">,
threadId: string,
) {
const response = yield* raw.request("thread/read", { threadId, includeTurns: false });
const metadata = yield* decodeCodexThreadHistoryMetadata(response).pipe(
Effect.mapError((error) =>
CodexErrors.CodexAppServerRequestError.invalidPayload("thread/read", "decode-payload", error),
),
);
return metadata.thread.historyMode;
});

export const resolveCodexRollbackTurnCount = Effect.fn("CodexAdapterV2.resolveRollbackTurnCount")(
function* (input: ProviderAdapterV2RollbackThreadInput) {
const providerTurns = input.providerThreadTurns;
Expand Down Expand Up @@ -5583,6 +5636,19 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi
runtimeRequests: [],
};
}
// thread/rollback only exists for legacy-history threads; Codex
// rejects it on paginated threads and this adapter has no
// paginated rollback path yet, so surface that honestly.
const historyMode = yield* ensureInitialized.pipe(
Effect.andThen(readCodexThreadHistoryMode(client.raw, threadId)),
);
if (historyMode === "paginated") {
return yield* new ProviderAdapterRollbackThreadError({
driver: CODEX_PROVIDER,
providerThreadId: threadInput.providerThread.id,
cause: `Cannot roll back Codex thread ${threadId}: the thread uses paginated history, which rejects thread/rollback, and this adapter does not implement paginated conversation rollback.`,
});
}
const response = yield* ensureInitialized.pipe(
Effect.andThen(client.request("thread/rollback", { threadId, numTurns })),
);
Expand Down Expand Up @@ -5617,10 +5683,14 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi
forkThread: (threadInput) =>
Effect.gen(function* () {
const threadId = yield* getNativeThreadId(threadInput.sourceProviderThread);
const boundary = yield* resolveCodexForkBoundary(threadInput);
const response = yield* ensureInitialized.pipe(
Effect.andThen(
client.request("thread/fork", {
threadId,
...(boundary.lastTurnId === undefined
? {}
: { lastTurnId: boundary.lastTurnId }),
...codexThreadRuntimeParams({
threadId: threadInput.targetThreadId,
...(threadInput.modelSelection === undefined
Expand All @@ -5641,26 +5711,39 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi
}),
),
);
const rollbackTurnCount = yield* resolveCodexForkRollbackTurnCount(threadInput);
const forkedThread =
rollbackTurnCount === 0
? response.thread
: (yield* ensureInitialized.pipe(
Effect.andThen(
client.request("thread/rollback", {
threadId: response.thread.id,
numTurns: rollbackTurnCount,
}),
),
Effect.mapError(
(cause) =>
new ProviderAdapterForkThreadError({
driver: CODEX_PROVIDER,
providerThreadId: threadInput.sourceProviderThread.id,
cause: normalizeCodexCause(cause),
}),
),
)).thread;
let forkedThread = response.thread;
if (boundary.rollbackTurnCount > 0) {
// Reached only when the selected source turn has no native
// turn reference, so the fork had to be taken at head and then
// trimmed. thread/rollback is legacy-history only; on a
// paginated fork the boundary cannot be honored at all.
const historyMode = yield* ensureInitialized.pipe(
Effect.andThen(readCodexThreadHistoryMode(client.raw, response.thread.id)),
);
if (historyMode === "paginated") {
return yield* new ProviderAdapterForkThreadError({
driver: CODEX_PROVIDER,
providerThreadId: threadInput.sourceProviderThread.id,
cause: `Cannot fork Codex thread ${threadId} at provider turn ${threadInput.providerTurnId}: the source turn has no native Codex turn reference, and the forked thread uses paginated history which rejects thread/rollback.`,
});
}
forkedThread = (yield* ensureInitialized.pipe(
Effect.andThen(
client.request("thread/rollback", {
threadId: response.thread.id,
numTurns: boundary.rollbackTurnCount,
}),
),
Effect.mapError(
(cause) =>
new ProviderAdapterForkThreadError({
driver: CODEX_PROVIDER,
providerThreadId: threadInput.sourceProviderThread.id,
cause: normalizeCodexCause(cause),
}),
),
)).thread;
}
return providerThreadFromCodexThread({
appThreadId: threadInput.targetThreadId,
idAllocator,
Expand Down
8 changes: 6 additions & 2 deletions apps/server/src/orchestration-v2/TODO.md
Original file line number Diff line number Diff line change
Expand Up @@ -12,8 +12,12 @@ implementation checklist for `apps/server/src/orchestration-v2`.
- Checkpoint rollback is currently a full revert: filesystem checkpoint restore, provider thread
rollback, stale checkpoint marking, and later run/node `rolled_back` projection state.
- Codex same-provider fork is lazy: `thread.fork` records lineage and pending transfer, and first
dispatch resolves native Codex fork. Earlier source-point forks use native `thread/fork` followed
by fork-local `thread/rollback`.
dispatch resolves native Codex fork. Earlier source-point forks pass the source turn's native id
to `thread/fork` as `lastTurnId`; the fork-then-rollback fallback remains only for source turns
without a native turn reference and only on legacy-history threads.
- Codex provider conversation rollback supports only legacy-history threads. The adapter probes
`historyMode` and fails explicitly on paginated threads; the V1 session runtime's
`thread/turns/list` + `thread/revert` path is the reference for closing this gap.
- Native Codex fork-from-earlier-run has a real replay-backed test fixture:
`testkit/fixtures/thread_fork_native_prior_turn`.
- Merge-back from a fork into its source thread records a `merge_back` context transfer, materializes
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -272,7 +272,14 @@ const scenarioExpectations = {
approvalRequestCount: 0,
},
thread_rollback: {
outgoing: ["initialize", "initialized", "thread/start", "turn/start", "thread/rollback"],
outgoing: [
"initialize",
"initialized",
"thread/start",
"turn/start",
"thread/read",
"thread/rollback",
],
incoming: ["turn/started", "turn/completed", "item/agentMessage/delta"],
turnStartCount: 3,
turnCompletedCount: 3,
Expand Down
Loading
Loading