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
58 changes: 56 additions & 2 deletions deploy/cloudflare-memory/runLedger.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -701,11 +701,25 @@ describe("run ledger — the coordinator's event and the key (items 47–48)", (
expect(sent).toEqual([]);
// The interrupted close written outside `finish`: one send, the record's status on it.
expect(
await post("/runs/put", { storeKey: key, record: { ...childRecord("p1", "slack:C1:3.0", "interrupted") } }),
await post("/runs/put", {
storeKey: key,
record: {
...childRecord("p1", "slack:C1:3.0", "interrupted"),
events: [
{
type: "coordinator_tag",
parentInstanceId: "ship_acme_api_1",
transportWorkflowId: "recovery-review-1",
at: 1,
},
],
eventCount: 1,
},
}),
).toMatchObject({ status: 200, data: { ok: true, stored: true } });
expect(sent).toEqual([
{
instance: "ship_acme_api_1",
instance: "recovery-review-1",
type: "run-finished-p1",
payload: expect.objectContaining({ runId: "p1", status: "interrupted", parentInstanceId: "ship_acme_api_1" }),
},
Expand Down Expand Up @@ -1105,6 +1119,46 @@ describe("run ledger — the coordinator's unit rows (item 50)", () => {
});
});

it("lists only active recovery rows in SQL, ignores terminal history, and fails closed on malformed candidate JSON", async () => {
const key = storeKey();
const active = unit("U12", {
pr: { number: 7, url: "https://github.com/acme/api/pull/7" },
recovery: {
kind: "review",
round: 2,
expectedHeadSha: "a".repeat(40),
remainingMs: 60_000,
claimedAt: 2_000,
step: "U12/recovery/2/review",
reviewRunId: "review-1",
previousEnding: { kind: "aborted", report: "recoverable", at: 1_000 },
workflowId: "recovery-review-1",
deadlineAt: 62_000,
reviewKey: `${INSTANCE_ID}:U12/1/review`,
},
});
const terminal = unit("U13", {
ending: { kind: "merged", report: "done", at: 3_000 },
});
expect((await post("/runs/coordinator/units/put", { storeKey: key, units: [active, terminal] })).status).toBe(200);
const stub = env.RUNS.get(env.RUNS.idFromName(key));

await expect(runInDurableObject(stub, (inst: RunHistoryDO) => inst.listActiveRecoveries())).resolves.toEqual([
active,
]);

await runInDurableObject(stub, async (_inst, state) => {
state.storage.sql.exec(
`INSERT INTO coordinator_units (instance_id, unit, json, updated_at) VALUES (?, ?, ?, ?)`,
INSTANCE_ID,
"U14",
`{"recovery":`,
4_000,
);
});
await expect(runInDurableObject(stub, (inst: RunHistoryDO) => inst.listActiveRecoveries())).rejects.toThrow();
});

it("the validating wake boundary accepts a first-segment resume without inventing a renewal segment row", async () => {
const key = storeKey();
const row = unit("U12", {
Expand Down
52 changes: 50 additions & 2 deletions deploy/cloudflare-memory/worker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -117,6 +117,7 @@ import {
sendRunFinished,
STEP_NAME_PATTERN,
UNIT_PATTERN,
unitOfIdempotencyKey,
type CoordinatorInstance,
type CoordinatorUnit,
type RunFinishedSend,
Expand Down Expand Up @@ -2771,6 +2772,37 @@ export class RunHistoryDO extends DurableObject<Env> {
.map((r) => JSON.parse(r.json) as CoordinatorUnit);
}

async listActiveRecoveries(): Promise<CoordinatorUnit[]> {
return this.sql
.exec<{ json: string }>(
`SELECT json FROM coordinator_units
WHERE json_valid(json) = 0 OR json_type(json, '$.recovery') IS NOT NULL
ORDER BY rowid`,
)
.toArray()
.map((r) => {
const unit: unknown = JSON.parse(r.json);
if (!isCoordinatorUnit(unit) || unit.recovery === undefined)
throw new Error("active recovery index contains a malformed coordinator unit");
return unit;
});
}

private recoveryTransport(parentInstanceId: string, idempotencyKey: string | undefined): string | undefined {
const unit = idempotencyKey === undefined ? undefined : unitOfIdempotencyKey(idempotencyKey);
if (unit === undefined) return undefined;
const row = this.sql
.exec<{ json: string }>(
`SELECT json FROM coordinator_units WHERE instance_id = ? AND unit = ?`,
parentInstanceId,
unit,
)
.toArray()[0];
if (row === undefined) return undefined;
const parsed = JSON.parse(row.json) as CoordinatorUnit;
return parsed.recovery?.workflowId;
}

// ---- the thread events of a unit-owned thread (record 0051's reply-as-event rule) --------------

/** The next sequence assigned in one transaction, the per-event cap applied
Expand Down Expand Up @@ -3025,9 +3057,11 @@ export class RunHistoryDO extends DurableObject<Env> {
// tells the waiting parent the child resumed — best effort, beside the
// bot's own announcement; a duplicate is consumed and re-armed, harmless.
if (req.meta.restartOf !== undefined && req.meta.parentInstanceId !== undefined) {
const recoveryTransport = this.recoveryTransport(req.meta.parentInstanceId, req.meta.idempotencyKey);
const sent = await sendChildSignal(this.env.SHIP_COORDINATOR, {
runId: req.runId,
parentInstanceId: req.meta.parentInstanceId,
...(recoveryTransport !== undefined ? { transportWorkflowId: recoveryTransport } : {}),
kind: "resumed",
reason: `restarted from run ${req.meta.restartOf}`,
at: now,
Expand Down Expand Up @@ -3318,7 +3352,14 @@ export class RunHistoryDO extends DurableObject<Env> {
if ((await this.ctx.storage.getAlarm()) === null)
await this.ctx.storage.setAlarm(systemClock() + RUN_SWEEP_INTERVAL_MS);
await this.refreshSessionBytes(record.session?.key);
const event = await sendRunFinished(this.env.SHIP_COORDINATOR, record);
const transportWorkflowId =
record.parentInstanceId === undefined
? undefined
: this.recoveryTransport(record.parentInstanceId, record.idempotencyKey);
const event = await sendRunFinished(this.env.SHIP_COORDINATOR, {
...record,
...(transportWorkflowId !== undefined ? { transportWorkflowId } : {}),
});
if (event.kind === "failed")
console.warn(`[runs/finish] ${runId} → ${event.type} not delivered to ${event.instance}: ${event.reason}`);
return { ...out, event: event.kind };
Expand Down Expand Up @@ -3654,7 +3695,11 @@ export class RunHistoryDO extends DurableObject<Env> {
// `finishedAt` equals `startedAt`); the parent confirms by `read-record`
// before it acts, so a duplicate send is harmless.
if (result.stored && record.parentInstanceId !== undefined && record.finishedAt > record.startedAt) {
const event = await sendRunFinished(this.env.SHIP_COORDINATOR, record);
const transportWorkflowId = record.events.find((event) => event.type === "coordinator_tag")?.transportWorkflowId;
const event = await sendRunFinished(this.env.SHIP_COORDINATOR, {
...record,
...(transportWorkflowId !== undefined ? { transportWorkflowId } : {}),
});
if (event.kind === "failed")
console.warn(`[runs/put] ${record.id} → ${event.type} not delivered to ${event.instance}: ${event.reason}`);
}
Expand Down Expand Up @@ -5220,6 +5265,7 @@ const LEDGER_ROUTES = new Set([
"/runs/coordinator/stop",
"/runs/coordinator/units/put",
"/runs/coordinator/units/claim-legacy-continuation",
"/runs/coordinator/units/list-active-recoveries",
"/runs/coordinator/units/list",
"/runs/coordinator/events/append",
"/runs/coordinator/events/list",
Expand Down Expand Up @@ -5935,6 +5981,8 @@ async function handleLedger(pathname: string, body: unknown, env: Env): Promise<
return json({ error: "instanceId must be a Workflow instance id" }, 400);
return json({ units: await stub.listUnits(b.instanceId) });
}
if (pathname === "/runs/coordinator/units/list-active-recoveries")
return json({ units: await stub.listActiveRecoveries() });
if (pathname === "/runs/coordinator/wake") {
if (!isCoordinatorUnit(b.unit)) return json({ error: "unit must be a coordinator unit row" }, 400);
if (typeof b.waitId !== "string" || !STEP_NAME_PATTERN.test(b.waitId))
Expand Down
6 changes: 5 additions & 1 deletion deploy/cloudflare/coordinator.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -81,7 +81,11 @@ describe("the coordinator holds no credential", () => {
);
expect(source).not.toMatch(/https?:\/\/(?!switchboard-keepalive\.internal)/);
});
it("runs the plan runner's driver and nothing of its own: `run()` is one `runPlan` over the platform's step and the container bot", () => {
it("selects the recovery driver only from typed recovery params and otherwise runs the plan driver", () => {
expect(source).toContain('if (event.payload.kind === "recover-original-unit")');
expect(source).toMatch(
/return runOriginalUnitRecovery\(workflowSteps\(step\), containerBot\(this\.env\), event\.instanceId, params\);/,
);
expect(source).toMatch(/return runPlan\(workflowSteps\(step\), containerBot\(this\.env\), event\.instanceId\);/);
});
});
Expand Down
19 changes: 18 additions & 1 deletion deploy/cloudflare/coordinator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,9 @@ import { NonRetryableError } from "cloudflare:workflows";
import { COORDINATOR_IDENTITY, COORDINATOR_STEP_PATH_PREFIX } from "../../src/core/coordinator/contract.ts";
import {
readBotAnswer,
runOriginalUnitRecovery,
runPlan,
type OriginalUnitRecoveryParams,
type CoordinatorBot,
type PlanRunSummary,
type StepRunner,
Expand All @@ -38,7 +40,17 @@ export type CoordinatorEnv = Pick<Env, "SWITCHBOARD" | "SWITCHBOARD_INGRESS_TOKE
/** What an instance is created with: nothing the driver reads — the instance
* id names the plan, and the bot's rows are the input (the coordinator's
* contract: ids only, never a task's text or a thread's contents). */
export type ShipCoordinatorParams = Record<string, unknown>;
export type ShipCoordinatorParams = Record<string, unknown> | OriginalUnitRecoveryParams;

function originalUnitRecoveryParams(value: ShipCoordinatorParams): OriginalUnitRecoveryParams | undefined {
if (
value.kind !== "recover-original-unit" ||
typeof value.parentInstanceId !== "string" ||
typeof value.unit !== "string"
)
return undefined;
return { kind: value.kind, parentInstanceId: value.parentInstanceId, unit: value.unit };
}

/** The bot behind the container binding: the reply as the wire carried it.
* The transport's failures throw for the step's retry; the door's own refusal
Expand Down Expand Up @@ -81,6 +93,11 @@ function workflowSteps(step: WorkflowStep): StepRunner {

export class ShipCoordinator extends WorkflowEntrypoint<CoordinatorEnv, ShipCoordinatorParams> {
async run(event: Readonly<WorkflowEvent<ShipCoordinatorParams>>, step: WorkflowStep): Promise<PlanRunSummary> {
if (event.payload.kind === "recover-original-unit") {
const params = originalUnitRecoveryParams(event.payload);
if (params === undefined) throw new NonRetryableError("the original-unit recovery params are malformed");
return runOriginalUnitRecovery(workflowSteps(step), containerBot(this.env), event.instanceId, params);
}
return runPlan(workflowSteps(step), containerBot(this.env), event.instanceId);
}
}
5 changes: 4 additions & 1 deletion deploy/cloudflare/worker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -325,7 +325,10 @@ async function handleCoordinatorInstances(request: Request, env: Env): Promise<R
const existing = await (await env.SHIP_COORDINATOR.get(parsed.id)).status();
outcome = { kind: "duplicate", id: parsed.id, status: existing.status };
} catch {
outcome = { kind: "failed", id: parsed.id, reason };
// Creation may have committed before its response was lost. When the
// status read is also unavailable, do not call that a definite failure:
// the bot must retain its same-id claim for a later duplicate replay.
return json(503, { ok: false, error: "create_unanswered", message: reason });
}
}
console.log(`[coordinator] ${auth.subject} → instance ${parsed.id}: ${outcome.kind}`);
Expand Down
Loading