Skip to content
197 changes: 197 additions & 0 deletions apps/server/src/orchestration-v2/FoundationPersistence.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ import {
ProviderThreadId,
RunAttemptId,
RunId,
RuntimeRequestId,
ThreadId,
TurnItemId,
} from "@t3tools/contracts";
Expand Down Expand Up @@ -338,6 +339,202 @@ it.effect("keeps other database work runnable while discovering compaction candi
).pipe(Effect.provide(TestLayer)),
);

it.effect.each([
{ trigger: "startup", runStatus: null },
{ trigger: "shutdown", runStatus: null },
{ trigger: "startup", runStatus: "completed" },
{ trigger: "shutdown", runStatus: "completed" },
] as const)(
"recovers request transcripts through SQLite on $trigger with run status $runStatus",
({ trigger, runStatus }) =>
Effect.gen(function* () {
const eventSink = yield* EventSink.EventSinkV2;
const projections = yield* ProjectionStore.ProjectionStoreV2;
const now = yield* DateTime.now;
const threadId = ThreadId.make(`thread:request-transcript:${trigger}:${runStatus}`);
const runId = runStatus === null ? null : RunId.make(`run:${threadId}`);
const events: Array<OrchestrationV2DomainEvent> = [
threadCreatedEvent({
id: `event:${threadId}:thread`,
thread: makeThread(threadId, now),
now,
}),
];
if (runId !== null) {
events.push({
id: EventId.make(`event:${threadId}:run`),
type: "run.created",
threadId,
runId,
occurredAt: now,
payload: {
id: runId,
threadId,
ordinal: 1,
providerInstanceId,
modelSelection,
providerThreadId: null,
userMessageId: MessageId.make(`message:${threadId}`),
rootNodeId: null,
activeAttemptId: null,
status: "completed",
requestedAt: now,
startedAt: now,
completedAt: now,
checkpointId: null,
contextHandoffId: null,
},
});
}
for (const kind of ["approval", "question", "durable-question"] as const) {
const requestId = RuntimeRequestId.make(`request:${threadId}:${kind}`);
const nodeId = NodeId.make(`node:${threadId}:${kind}`);
const itemType = kind === "approval" ? "approval_request" : "user_input_request";
const itemBase = {
id: TurnItemId.make(`item:${threadId}:${kind}`),
threadId,
runId,
nodeId,
providerThreadId: null,
providerTurnId: null,
nativeItemRef: null,
parentItemId: null,
ordinal: events.length,
status: "waiting" as const,
title: null,
startedAt: now,
completedAt: null,
updatedAt: now,
requestId,
};
const item: OrchestrationV2TurnItem =
kind === "approval"
? {
...itemBase,
type: "approval_request",
requestKind: "command",
prompt: "Run command?",
}
: { ...itemBase, type: "user_input_request", questions: [] };
events.push(
{
id: EventId.make(`event:${nodeId}`),
type: "node.updated",
threadId,
nodeId,
occurredAt: now,
payload: {
id: nodeId,
threadId,
runId,
parentNodeId: null,
rootNodeId: nodeId,
kind: itemType,
status: "waiting",
countsForRun: false,
providerThreadId: null,
providerTurnId: null,
nativeItemRef: null,
runtimeRequestId: requestId,
checkpointScopeId: null,
startedAt: now,
completedAt: null,
},
},
{
id: EventId.make(`event:${requestId}`),
type: "runtime-request.updated",
threadId,
occurredAt: now,
payload: {
id: requestId,
nodeId,
providerTurnId: null,
nativeRequestRef: null,
kind: kind === "approval" ? "command" : "user_input",
status: "pending",
responseCapability:
kind === "durable-question"
? { type: "message" }
: {
type: "live",
providerSessionId: ProviderSessionId.make(`session:${threadId}`),
},
createdAt: now,
resolvedAt: null,
},
},
{
id: EventId.make(`event:${item.id}`),
type: "turn-item.updated",
threadId,
occurredAt: now,
payload: item,
},
);
if (kind === "approval") {
events.push({
id: EventId.make(`event:${item.id}:terminal`),
type: "turn-item.updated",
threadId,
occurredAt: now,
payload: {
...item,
id: TurnItemId.make(`item:${threadId}:terminal`),
status: "completed",
completedAt: now,
},
});
}
}
yield* eventSink.commitCommand({
commandId: CommandId.make(`command:${threadId}:seed`),
threadId,
commandType: "foundation.request-transcript",
acceptedAt: now,
events,
effects: [],
});
const selected = yield* projections.getRuntimeRecoveryProjection(threadId);
assert.sameMembers(
selected.nodes.map((node) => node.id),
[NodeId.make(`node:${threadId}:approval`), NodeId.make(`node:${threadId}:question`)],
);
assert.sameMembers(
selected.turnItems.map((item) => item.id),
[
TurnItemId.make(`item:${threadId}:approval`),
TurnItemId.make(`item:${threadId}:question`),
],
);
const recovery = yield* ProviderRuntimeRecovery.make.pipe(
Effect.provide(ServerSettings.layerTest()),
);
assert.equal((yield* recovery.reconcile(trigger)).closedRequests, 2);
const final = yield* projections.getThreadProjection(threadId);
for (const kind of ["approval", "question", "durable-question"] as const) {
const durable = kind === "durable-question";
assert.equal(
final.runtimeRequests.find((request) => request.id === `request:${threadId}:${kind}`)
?.status,
durable ? "pending" : trigger === "startup" ? "expired" : "cancelled",
);
assert.equal(
final.nodes.find((node) => node.id === `node:${threadId}:${kind}`)?.status,
durable ? "waiting" : "cancelled",
);
assert.equal(
final.turnItems.find((item) => item.id === `item:${threadId}:${kind}`)?.status,
durable ? "waiting" : "cancelled",
);
}
assert.equal(
final.turnItems.find((item) => item.id === `item:${threadId}:terminal`)?.status,
"completed",
);
}).pipe(Effect.provide(Layer.fresh(TestLayer))),
);

it.layer(TestLayer)("orchestration V2 foundation persistence", (it) => {
it.effect("projects oversized tool bodies before both replay and live RPC retention", () =>
Effect.scoped(
Expand Down
14 changes: 14 additions & 0 deletions apps/server/src/orchestration-v2/ProjectionStore.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3888,6 +3888,11 @@ export const layer: Layer.Layer<ProjectionStoreV2, never, SqlClient.SqlClient> =
WHERE item.thread_id = ${threadId} AND item.type = 'subagent'
AND item.status IN ('pending', 'running', 'waiting')
)
OR node.node_id IN (
SELECT request.node_id FROM orchestration_v2_projection_runtime_requests AS request
WHERE request.thread_id = ${threadId} AND request.status = 'pending'
AND json_extract(request.payload_json, '$.responseCapability.type') <> 'message'
)
)
ORDER BY COALESCE(node.started_at, ''), node.node_id ASC
`,
Expand Down Expand Up @@ -4008,6 +4013,15 @@ export const layer: Layer.Layer<ProjectionStoreV2, never, SqlClient.SqlClient> =
AND status IN ('queued', 'preparing', 'starting', 'running', 'waiting')
)
OR item.type IN ('command_execution', 'dynamic_tool', 'subagent')
OR (
item.type IN ('approval_request', 'user_input_request')
AND json_extract(item.payload_json, '$.requestId') IN (
SELECT request.runtime_request_id
FROM orchestration_v2_projection_runtime_requests AS request
WHERE request.thread_id = ${threadId} AND request.status = 'pending'
AND json_extract(request.payload_json, '$.responseCapability.type') <> 'message'
)
)
OR (
item.run_id IS NULL
AND item.node_id IN (
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -180,6 +180,117 @@ it.effect("expires orphaned runtime requests before command readiness", () => {
}).pipe(Effect.provide(layer));
});

it.effect.each([
{ trigger: "startup", requestType: "approval_request", runStatus: null },
{ trigger: "startup", requestType: "user_input_request", runStatus: null },
{ trigger: "shutdown", requestType: "approval_request", runStatus: null },
{ trigger: "shutdown", requestType: "user_input_request", runStatus: null },
{ trigger: "startup", requestType: "approval_request", runStatus: "running" },
{ trigger: "startup", requestType: "approval_request", runStatus: "completed" },
] as const)(
"closes $requestType transcript entities on $trigger with run status $runStatus",
({ trigger, requestType, runStatus }) => {
const threadId = ThreadId.make(`thread_recovery_${trigger}_${requestType}`);
const requestId = RuntimeRequestId.make("request_runless");
const nodeId = NodeId.make("node_runless");
const itemId = TurnItemId.make("item_runless");
const terminalItemId = TurnItemId.make("item_already_terminal");
const messageNodeId = NodeId.make("node_message");
const runId = runStatus === null ? null : RunId.make("run_recovery_request");
let committedInput: Parameters<EventSink.EventSinkV2["Service"]["commitCommand"]>[0] | null =
null;
const projection = {
thread: { id: threadId },
runtimeRequests: [
{ id: requestId, nodeId, status: "pending", responseCapability: { type: "live" } },
{
id: RuntimeRequestId.make("request_message"),
nodeId: messageNodeId,
status: "pending",
responseCapability: { type: "message" },
},
],
providerSessions: [],
providerThreads: [],
providerTurns: [],
runs:
runStatus === null
? []
: [
{
id: runId,
status: runStatus,
providerInstanceId: ProviderInstanceId.make("codex"),
},
],
attempts: [],
subagents: [],
messages: [],
nodes: [
{ id: nodeId, runId, kind: requestType, status: "waiting" },
{ id: messageNodeId, runId, kind: "user_input_request", status: "waiting" },
],
turnItems: [
{ id: itemId, runId, nodeId, type: requestType, requestId, status: "waiting" },
{
id: terminalItemId,
runId,
nodeId,
type: requestType,
requestId,
status: "completed",
},
{
id: TurnItemId.make("item_message"),
runId,
nodeId: messageNodeId,
type: "user_input_request",
requestId: RuntimeRequestId.make("request_message"),
status: "waiting",
},
],
} as unknown as OrchestrationV2ThreadProjection;
const layer = ProviderRuntimeRecovery.layer.pipe(
Layer.provide(ServerSettings.layerTest()),
Layer.provide(
Layer.mergeAll(
Layer.mock(ProjectionStore.ProjectionStoreV2)({
getRecoveryThreadIds: () => Effect.succeed([threadId]),
getRuntimeRecoveryProjection: () => Effect.succeed(projection),
}),
Layer.mock(EventSink.EventSinkV2)({
commitCommand: (input) => {
committedInput = input;
return Effect.succeed({ committed: true, cancelledEffectCount: 0 } as never);
},
}),
IdAllocator.layer,
Layer.mock(EffectOutbox.EffectOutboxV2)({
reconcileAfterProcessLoss: Effect.succeed({ requeued: 0, cancelled: 0 }),
}),
),
),
);
return Effect.gen(function* () {
const summary =
yield* (yield* ProviderRuntimeRecovery.ProviderRuntimeRecoveryService).reconcile(trigger);
assert.equal(summary.closedRequests, 1);
const events = committedInput?.events ?? [];
assert.equal(events.length, runStatus === "running" ? 4 : 3);
const requestEvent = events.find((event) => event.type === "runtime-request.updated");
assert.equal(requestEvent?.payload.status, trigger === "startup" ? "expired" : "cancelled");
const nodeEvent = events.find((event) => event.type === "node.updated");
const itemEvent = events.find((event) => event.type === "turn-item.updated");
assert.equal(nodeEvent?.payload.id, nodeId);
assert.equal(nodeEvent?.payload.status, "cancelled");
assert.isNotNull(nodeEvent?.payload.completedAt);
assert.equal(itemEvent?.payload.id, itemId);
assert.equal(itemEvent?.payload.status, "cancelled");
assert.isNotNull(itemEvent?.payload.completedAt);
}).pipe(Effect.provide(layer));
},
);

it.effect("preserves async questions across startup and shutdown", () => {
const threadId = ThreadId.make("async-recovery-thread");
const nodeId = NodeId.make("async-recovery-node");
Expand Down
Loading
Loading