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
105 changes: 105 additions & 0 deletions apps/server/src/orchestration-v2/ProjectionStore.test.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import { assert, it, vi } from "@effect/vitest";
import {
EventId,
CommandId,
CheckpointId,
CheckpointRef,
CheckpointScopeId,
Expand Down Expand Up @@ -2003,6 +2004,29 @@ it.layer(TestLayer)("ProjectionStoreV2", (it) => {
);
}
assert.isNull(shells.threads.find((row) => row.id === otherThreadId)!.lastErrorClass);
const candidates = yield* store.getLimitRecoveryCandidates({
now,
autoResume: true,
snooze: false,
});
const candidate = candidates.find((row) => row.id === threadId);
if (lastErrorClass === "usage_limit") {
assert.deepEqual(candidate, {
id: sqlShell.id,
status: sqlShell.status,
lastErrorClass: sqlShell.lastErrorClass,
usageLimitResetAt: sqlShell.usageLimitResetAt,
latestRunId: sqlShell.latestRunId,
latestRunCompletedAt: sqlShell.latestRunCompletedAt,
updatedAt: sqlShell.updatedAt,
archivedAt: sqlShell.archivedAt,
settledOverride: sqlShell.settledOverride,
pendingRuntimeRequest: null,
limitRecovery: sqlShell.limitRecovery,
snoozedUntil: sqlShell.snoozedUntil,
});
} else assert.isUndefined(candidate);
assert.isUndefined(candidates.find((row) => row.id === otherThreadId));
});
yield* store.apply({
id: EventId.make("event:limit-shell:error"),
Expand All @@ -2013,6 +2037,87 @@ it.layer(TestLayer)("ProjectionStoreV2", (it) => {
});
yield* applyRun("failed");
yield* assertSummary("Plan limit reached.", "usage_limit");
const sql = yield* SqlClient.SqlClient;
const [originalRow] = yield* sql<{
payload_json: string;
}>`SELECT payload_json FROM orchestration_v2_projection_threads WHERE thread_id = ${threadId}`;
for (const [field, value] of [
["archivedAt", DateTime.formatIso(now)],
["settledOverride", "settled"],
]) {
yield* sql`UPDATE orchestration_v2_projection_threads
SET payload_json = json_set(payload_json, ${`$.${field}`}, ${value})
WHERE thread_id = ${threadId}`;
assert.isUndefined(
(yield* store.getLimitRecoveryCandidates({ now, autoResume: true, snooze: false })).find(
(row) => row.id === threadId,
),
);
yield* sql`UPDATE orchestration_v2_projection_threads
SET payload_json = ${originalRow!.payload_json} WHERE thread_id = ${threadId}`;
}
yield* sql`UPDATE orchestration_v2_projection_threads SET deleted_at = ${DateTime.formatIso(now)} WHERE thread_id = ${threadId}`;
assert.isUndefined(
(yield* store.getLimitRecoveryCandidates({ now, autoResume: true, snooze: false })).find(
(row) => row.id === threadId,
),
);
yield* sql`UPDATE orchestration_v2_projection_threads SET deleted_at = NULL WHERE thread_id = ${threadId}`;
const recoveryOptions = { now, autoResume: false, snooze: false };
assert.isUndefined(
(yield* store.getLimitRecoveryCandidates(recoveryOptions)).find(
(row) => row.id === threadId,
),
);
const reset = DateTime.makeUnsafe(limitItem.failure.resetAt);
const recovery = {
runId: original.id,
resetAt: limitItem.failure.resetAt,
autoResume: true,
requestId: CommandId.make("recovery:choice"),
};
yield* sql`UPDATE orchestration_v2_projection_threads
SET payload_json = json_set(payload_json, '$.limitRecovery', json(${encodeUnknownJsonString(recovery)}))
WHERE thread_id = ${threadId}`;
// Armed future retries need no state decoding until they become due.
assert.isUndefined(
(yield* store.getLimitRecoveryCandidates({ ...recoveryOptions, autoResume: true })).find(
(row) => row.id === threadId,
),
);
const due = (yield* store.getLimitRecoveryCandidates({
...recoveryOptions,
now: reset,
})).find((row) => row.id === threadId)!;
assert.deepEqual(due.limitRecovery, recovery);
yield* sql`INSERT INTO orchestration_v2_projection_runtime_requests
(runtime_request_id, thread_id, node_id, kind, status, created_at, payload_json)
VALUES ('limit-shell:pending-request', ${threadId}, ${original.rootNodeId}, 'approval', 'pending', ${DateTime.formatIso(now)}, '{}')`;
assert.isUndefined(
(yield* store.getLimitRecoveryCandidates({ ...recoveryOptions, now: reset })).find(
(row) => row.id === threadId,
),
);
yield* sql`DELETE FROM orchestration_v2_projection_runtime_requests WHERE runtime_request_id = 'limit-shell:pending-request'`;
yield* sql`UPDATE orchestration_v2_projection_threads
SET payload_json = json_set(payload_json, '$.snoozedUntil', ${DateTime.formatIso(DateTime.add(reset, { minutes: 1 }))})
WHERE thread_id = ${threadId}`;
assert.isUndefined(
(yield* store.getLimitRecoveryCandidates({ ...recoveryOptions, now: reset })).find(
(row) => row.id === threadId,
),
);
yield* sql`UPDATE orchestration_v2_projection_threads
SET payload_json = json_set(payload_json, '$.snoozedUntil', NULL, '$.limitRecovery.autoResume', json('false'))
WHERE thread_id = ${threadId}`;
assert.isUndefined(
(yield* store.getLimitRecoveryCandidates({
...recoveryOptions,
now: reset,
autoResume: true,
})).find((row) => row.id === threadId),
);
yield* sql`UPDATE orchestration_v2_projection_threads SET payload_json = ${originalRow!.payload_json} WHERE thread_id = ${threadId}`;
const session = {
id: ProviderSessionId.make("session:limit-shell:shared"),
driver,
Expand Down
166 changes: 166 additions & 0 deletions apps/server/src/orchestration-v2/ProjectionStore.ts
Original file line number Diff line number Diff line change
Expand Up @@ -130,6 +130,23 @@ export type ProjectionRecoveryKind =
| "subagent-results"
| "delegated-completions";

/** Persisted state needed for limit recovery, without transcript or fork history. */
export type ProjectionLimitRecoveryCandidate = Pick<
OrchestrationV2ThreadShell,
| "id"
| "status"
| "lastErrorClass"
| "latestRunId"
| "usageLimitResetAt"
| "archivedAt"
| "settledOverride"
| "pendingRuntimeRequest"
| "latestRunCompletedAt"
| "updatedAt"
| "limitRecovery"
| "snoozedUntil"
>;

/** Thread activity needed by settlement, without transcript or fork history. */
export type ProjectionSettlementCandidate = Pick<
OrchestrationV2ThreadShell,
Expand Down Expand Up @@ -255,6 +272,11 @@ export interface ProjectionStoreV2Shape {
readonly getThread: (
threadId: ThreadId,
) => Effect.Effect<OrchestrationV2AppThread, ProjectionStoreV2Error>;
readonly getLimitRecoveryCandidates: (options: {
readonly now: DateTime.Utc;
readonly autoResume: boolean;
readonly snooze: boolean;
}) => Effect.Effect<ReadonlyArray<ProjectionLimitRecoveryCandidate>, ProjectionStoreV2Error>;
readonly getSettlementCandidates: () => Effect.Effect<
ReadonlyArray<ProjectionSettlementCandidate>,
ProjectionStoreV2Error
Expand Down Expand Up @@ -3054,6 +3076,107 @@ export const layer: Layer.Layer<ProjectionStoreV2, never, SqlClient.SqlClient> =
),
);

const getLimitRecoveryCandidates = Effect.fn("ProjectionStore.getLimitRecoveryCandidates")(
function* (options: Parameters<ProjectionStoreV2Shape["getLimitRecoveryCandidates"]>[0]) {
// Indexed latest-run and root-error lookups avoid reading run histories,
// counting transcript items, or walking fork ancestors on scheduler ticks.
const rows = yield* sql<{
readonly payload_json: string;
readonly run_id: string;
readonly completed_at: string | null;
readonly failure_payload_json: string;
readonly last_error: string | null;
}>`
SELECT t.payload_json, r.run_id, r.completed_at,
item.payload_json AS failure_payload_json,
(
SELECT json_extract(session.payload_json, '$.lastError')
FROM orchestration_v2_projection_provider_sessions session
INNER JOIN orchestration_v2_projection_provider_session_bindings binding
ON binding.provider_session_id = session.provider_session_id
WHERE binding.thread_id = t.thread_id
AND session.provider_instance_id = t.provider_instance_id
ORDER BY session.updated_at DESC, session.provider_session_id DESC
LIMIT 1
) AS last_error
FROM orchestration_v2_projection_threads t
INNER JOIN orchestration_v2_projection_runs r ON r.run_id = (
SELECT latest.run_id FROM orchestration_v2_projection_runs latest
WHERE latest.thread_id = t.thread_id
ORDER BY latest.ordinal DESC, latest.run_id DESC LIMIT 1
) AND r.status = 'failed'
INNER JOIN orchestration_v2_projection_turn_items item ON item.turn_item_id = (
SELECT error.turn_item_id FROM orchestration_v2_projection_turn_items error
WHERE error.thread_id = t.thread_id AND error.run_id = r.run_id
AND error.type = 'error' AND error.status = 'failed'
AND error.node_id IS json_extract(r.payload_json, '$.rootNodeId')
ORDER BY error.updated_at DESC, error.ordinal DESC, error.turn_item_id DESC
LIMIT 1
)
WHERE t.deleted_at IS NULL
AND json_extract(t.payload_json, '$.archivedAt') IS NULL
AND json_extract(t.payload_json, '$.settledOverride') IS NOT 'settled'
AND json_extract(item.payload_json, '$.failure.class') = 'usage_limit'
AND json_extract(item.payload_json, '$.failure.resetAt') IS NOT NULL
AND julianday(json_extract(item.payload_json, '$.failure.resetAt')) > julianday(COALESCE(r.completed_at, json_extract(t.payload_json, '$.updatedAt')))
AND (
(
json_extract(t.payload_json, '$.limitRecovery.runId') IS r.run_id
AND json_extract(t.payload_json, '$.limitRecovery.resetAt') IS json_extract(item.payload_json, '$.failure.resetAt')
AND json_extract(t.payload_json, '$.limitRecovery.autoResume') = 1
AND julianday(json_extract(item.payload_json, '$.failure.resetAt')) <= julianday(${DateTime.formatIso(options.now)})
AND (
json_extract(t.payload_json, '$.snoozedUntil') IS NULL
OR julianday(json_extract(t.payload_json, '$.snoozedUntil')) <= julianday(${DateTime.formatIso(options.now)})
)
)
OR (
(
json_extract(t.payload_json, '$.limitRecovery.runId') IS NOT r.run_id
OR json_extract(t.payload_json, '$.limitRecovery.resetAt') IS NOT json_extract(item.payload_json, '$.failure.resetAt')
)
AND (
${options.autoResume}
OR (${options.snooze} AND julianday(json_extract(item.payload_json, '$.failure.resetAt')) > julianday(${DateTime.formatIso(options.now)}))
)
)
)
AND NOT EXISTS (
SELECT 1 FROM orchestration_v2_projection_runtime_requests request
WHERE request.thread_id = t.thread_id AND request.status = 'pending'
)
ORDER BY t.thread_id
`;
const candidates: Array<ProjectionLimitRecoveryCandidate> = [];
for (const row of rows) {
const thread = yield* decodeThreadPayload(row.payload_json);
const item = yield* decodeTurnItemPayload(row.failure_payload_json);
const summary = threadErrorSummary(
item.type === "error" ? item.failure : null,
row.last_error,
);
if (summary.lastErrorClass !== "usage_limit") continue;
candidates.push({
id: thread.id,
status: "failed",
lastErrorClass: summary.lastErrorClass,
usageLimitResetAt: summary.usageLimitResetAt,
latestRunId: RunId.make(row.run_id),
latestRunCompletedAt:
row.completed_at === null ? null : DateTime.makeUnsafe(row.completed_at),
updatedAt: thread.updatedAt,
archivedAt: thread.archivedAt,
settledOverride: thread.settledOverride,
pendingRuntimeRequest: null,
limitRecovery: thread.limitRecovery ?? null,
snoozedUntil: thread.snoozedUntil ?? null,
});
}
return candidates;
},
Effect.mapError((cause) => new ProjectionStoreSetupError({ cause })),
);

const getRecoveryThreadIds = Effect.fn("ProjectionStore.getRecoveryThreadIds")(
function* (kind: ProjectionRecoveryKind) {
const candidates = (() => {
Expand Down Expand Up @@ -4813,6 +4936,7 @@ export const layer: Layer.Layer<ProjectionStoreV2, never, SqlClient.SqlClient> =
getRuntimeRequest,
getPlan,
getProviderControlContext,
getLimitRecoveryCandidates,
getRecoveryThreadIds,
getUnreadableThreadIds,
getThreadSnapshot,
Expand Down Expand Up @@ -4911,6 +5035,48 @@ export const layerMemory: Layer.Layer<ProjectionStoreV2> = Layer.effect(
left.id.localeCompare(right.id),
);
}),
getLimitRecoveryCandidates: (options) =>
Ref.get(replayState).pipe(
Effect.map((state) =>
[...state.projections.values()]
.filter(
({ thread }) =>
thread.deletedAt === null &&
thread.archivedAt === null &&
thread.settledOverride !== "settled",
)
.map(threadShellFromProjection)
.filter(
(thread) =>
thread.status === "failed" &&
thread.lastErrorClass === "usage_limit" &&
thread.usageLimitResetAt !== null &&
thread.pendingRuntimeRequest === null,
)
.filter((thread) => {
const resetMs = Date.parse(thread.usageLimitResetAt!);
const nowMs = DateTime.toEpochMillis(options.now);
if (
!Number.isFinite(resetMs) ||
resetMs <= DateTime.toEpochMillis(thread.latestRunCompletedAt ?? thread.updatedAt)
)
return false;
if (
thread.limitRecovery?.runId !== thread.latestRunId ||
thread.limitRecovery.resetAt !== thread.usageLimitResetAt
) {
return options.autoResume || (options.snooze && resetMs > nowMs);
}
return (
thread.limitRecovery.autoResume &&
resetMs <= nowMs &&
(thread.snoozedUntil == null ||
DateTime.toEpochMillis(thread.snoozedUntil) <= nowMs)
);
})
.toSorted((left, right) => left.id.localeCompare(right.id)),
),
),
getRecoveryThreadIds: (kind) =>
Effect.gen(function* () {
const projections = (yield* Ref.get(replayState)).projections;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -214,6 +214,7 @@ it.effect(
ProjectionStoreV2,
ProjectionStoreV2.of({
apply: () => Effect.void,
getLimitRecoveryCandidates: () => Effect.die("unused getLimitRecoveryCandidates"),
getShellSnapshot: () => Effect.die("unused getShellSnapshot"),
getThreadShell: () => Effect.die("unused getThreadShell"),
getThread: () => Ref.get(projection).pipe(Effect.map((state) => state.thread)),
Expand Down
3 changes: 2 additions & 1 deletion apps/server/src/orchestration-v2/ThreadLaunchService.test.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import * as Scheduler from "../scheduling/Scheduler.ts";
import * as WorktreeSetupTracker from "../project/WorktreeSetupTracker.ts";
import * as ProjectCloneTracker from "../project/ProjectCloneTracker.ts";
import * as TerminalManager from "../terminal/Manager.ts";
Expand Down Expand Up @@ -267,7 +268,7 @@ for (const target of ["new", "existing"] as const) {
() => {
const harness = makeHarness();
const scheduledTasks = ScheduledTasks.layer.pipe(
Layer.provide(Layer.mergeAll(harness.layer, NodeCrypto.layer)),
Layer.provide(Layer.mergeAll(harness.layer, NodeCrypto.layer, Scheduler.layer)),
);
return Effect.gen(function* () {
const tasks = yield* ScheduledTasks.ScheduledTaskService;
Expand Down
Loading
Loading