Skip to content
Closed
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
7 changes: 4 additions & 3 deletions docs/reference/specs/agent-ship.md

Large diffs are not rendered by default.

3 changes: 2 additions & 1 deletion docs/reference/specs/run-history.md

Large diffs are not rendered by default.

147 changes: 147 additions & 0 deletions src/channels/adminCoordinator.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3419,6 +3419,134 @@ describe("the plan runner's steps — plan, unit-start, branch, round, unit-end,
expect((await h.instances.listUnits(PLAN_INSTANCE.id))[0].rounds).toHaveLength(1); // nothing malformed was appended
});

it("an operator-check boundary keeps the report, and a lost route response replays without another boundary or thread post", async () => {
const replies: string[] = [];
const h = await planHarness({
ioFor: () => ({
reply: async (text: string) => void replies.push(text),
status: async () => ({ update: () => {}, done: async () => {} }),
history: async () => [],
}),
});
await h.instances.putUnits([unitRow("U10", { threadKey: "slack:C1:2.0" })]);
const reportHead = "a".repeat(40);
const boundary = {
parentInstanceId: PLAN_INSTANCE.id,
unit: "U10",
index: 1,
agent: "review",
outcome: "blocked_by_operator_check",
};
const report = `⏸️ Blocked by an operator check at approved head \`${reportHead}\`; no coding child was started.\n\n> The deployed bot holds no OPENAI_API_KEY.`;

expect(await call(h, "round", { ...boundary, report, reportHead })).toEqual({
status: 200,
body: { ok: true, at: NOW },
});
const row = (await h.instances.listUnits(PLAN_INSTANCE.id))[0];
expect(row.rounds).toEqual([
{ index: 1, agent: "review", outcome: "blocked_by_operator_check", reportHead, at: NOW },
]);
expect(row.operatorCheckReports).toEqual({ [reportHead]: { report, deliveredAt: NOW } });
expect(replies).toEqual([report]);
// The runner may lose this route's successful HTTP response and replay its
// durable step. The same per-head identity is a read, not a second boundary
// or Slack post.
expect(await call(h, "round", { ...boundary, report, reportHead })).toEqual({
status: 200,
body: { ok: true, at: NOW },
});
expect((await h.instances.listUnits(PLAN_INSTANCE.id))[0].rounds).toHaveLength(1);
expect(replies).toEqual([report]);
expect((await call(h, "round", boundary)).status).toBe(400);
expect(
(
await call(h, "round", {
...boundary,
outcome: "approve",
report,
reportHead,
})
).status,
).toBe(400);
});

it("retries an undelivered operator-check report without appending its per-head boundary twice", async () => {
const replies: string[] = [];
let attempts = 0;
const h = await planHarness({
ioFor: () => ({
reply: async (text: string) => {
attempts++;
if (attempts === 1) throw new Error("Slack temporarily unavailable");
replies.push(text);
},
status: async () => ({ update: () => {}, done: async () => {} }),
history: async () => [],
}),
});
await h.instances.putUnits([unitRow("U10", { threadKey: "slack:C1:2.0" })]);
const reportHead = "b".repeat(40);
const report = `⏸️ Blocked by an operator check at approved head \`${reportHead}\`.\n\n> Run deploy secrets bot --only OPENAI_API_KEY.`;
const boundary = {
parentInstanceId: PLAN_INSTANCE.id,
unit: "U10",
index: 1,
agent: "review",
outcome: "blocked_by_operator_check",
report,
reportHead,
};

expect((await call(h, "round", boundary)).status).toBe(502);
let row = (await h.instances.listUnits(PLAN_INSTANCE.id))[0];
expect(row.rounds).toHaveLength(1);
expect(row.operatorCheckReports).toEqual({ [reportHead]: { report } });

expect(await call(h, "round", boundary)).toEqual({ status: 200, body: { ok: true, at: NOW } });
row = (await h.instances.listUnits(PLAN_INSTANCE.id))[0];
expect(row.rounds).toHaveLength(1);
expect(row.operatorCheckReports).toEqual({ [reportHead]: { report, deliveredAt: NOW } });
expect(attempts).toBe(2);
expect(replies).toEqual([report]);
});

it("reconciles a lost successful reply from thread history, so catch-up posts no duplicate report or boundary", async () => {
const posted: string[] = [];
let attempts = 0;
const h = await planHarness({
ioFor: () => ({
reply: async (text: string) => {
attempts++;
posted.push(text);
throw new Error("Slack accepted the post but the response was lost");
},
status: async () => ({ update: () => {}, done: async () => {} }),
history: async () => posted.map((text) => ({ role: "assistant" as const, text })),
}),
});
await h.instances.putUnits([unitRow("U10", { threadKey: "slack:C1:2.0" })]);
const reportHead = "c".repeat(40);
const report = `⏸️ Blocked by an operator check at approved head \`${reportHead}\`.\n\n> Run deploy secrets bot --only OPENAI_API_KEY.`;
const boundary = {
parentInstanceId: PLAN_INSTANCE.id,
unit: "U10",
index: 1,
agent: "review",
outcome: "blocked_by_operator_check",
report,
reportHead,
};

expect((await call(h, "round", boundary)).status).toBe(502);
expect(await call(h, "round", boundary)).toEqual({ status: 200, body: { ok: true, at: NOW } });
const row = (await h.instances.listUnits(PLAN_INSTANCE.id))[0];
expect(row.rounds).toHaveLength(1);
expect(row.operatorCheckReports).toEqual({ [reportHead]: { report, deliveredAt: NOW } });
expect(attempts).toBe(1);
expect(posted).toEqual([report]);
});

// Record 0065 / issue 1968: `ShipRoundOutcome` grew `continued` (decision 0046's
// renewal) and `idle` (record 0051) while the route's accepted list did not,
// so a renewed round 0 threw in the driver. The route now accepts the whole
Expand Down Expand Up @@ -4560,6 +4688,25 @@ describe("POST /admin/coordinator/checks — the round's checks read at the revi
expect(none.mergeWaitNotes).toEqual([{ headSha: HEAD, instanceId: INSTANCE.id, at: NOW }]);
});

it("an operator-owned red check registers the approved head so its re-run wakes the same checks step", async () => {
const operator = await checksHarness({
roundChecks: {
total: 1,
pending: [],
failed: [
{
name: "production impact",
conclusion: "failure",
operatorPrecondition: true,
output: "The deployed bot holds no OPENAI_API_KEY.",
},
],
},
});
await checks(operator);
expect(operator.mergeWaitNotes).toEqual([{ headSha: HEAD, instanceId: INSTANCE.id, at: NOW }]);
});

it("an expected check not yet reported registers the head in the merge-wait book like a pending one, and the round reader is handed the pull request's own base for the required checks (issue 2063)", async () => {
const h = await checksHarness({
roundChecks: { total: 3, pending: [], failed: [], expected: ["approve"] },
Expand Down
106 changes: 96 additions & 10 deletions src/channels/adminCoordinator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2052,6 +2052,7 @@ const ROUND_OUTCOMES = [
"request_changes",
"no_verdict",
"checks_failed",
"blocked_by_operator_check",
"checks_restarted",
"transient",
"enqueued",
Expand Down Expand Up @@ -2152,6 +2153,11 @@ function parseGate(raw: unknown): { level: AddressSeverity; findings: string[] }
return { level: g.level, findings: g.findings as string[] };
}

/** The visible, per-head identity at the start of an operator-check report.
* Slack may split a long report, so reconciliation looks for this marker
* rather than requiring one history turn to equal the whole report. */
const operatorCheckReportIdentity = (headSha: string): string => `approved head \`${headSha}\``;

/** A round boundary: appended to the unit's row and drawn on the card. */
async function round(body: Record<string, unknown>, deps: AdminCoordinatorDeps): Promise<IngressResponse> {
const id = parseInstanceId(body.parentInstanceId);
Expand All @@ -2174,6 +2180,22 @@ async function round(body: Record<string, unknown>, deps: AdminCoordinatorDeps):
ok: false,
error: "gate must be { level: blocking|major|minor|nit, findings: string[] }",
});
const report =
typeof body.report === "string" && body.report.length > 0 && body.report.length <= 20_000 ? body.report : undefined;
const reportHead = normalizeHead(body.reportHead);
const operatorBoundary = body.outcome === "blocked_by_operator_check";
if (
(operatorBoundary &&
(report === undefined ||
reportHead === undefined ||
!report.includes(operatorCheckReportIdentity(reportHead)))) ||
(!operatorBoundary && (body.report !== undefined || body.reportHead !== undefined))
)
return json(400, {
ok: false,
error:
"blocked_by_operator_check must carry its report with the approved reportHead identity, and no other outcome may",
});
const at = (deps.clock ?? systemClock)();
const instance = await deps.instances.get(id.value);
if (!instance) return json(404, { ok: false, error: "unknown_instance" });
Expand All @@ -2183,15 +2205,39 @@ async function round(body: Record<string, unknown>, deps: AdminCoordinatorDeps):
const units = await deps.instances.listUnits(instance.id);
const row = units.find((u) => u.unit === body.unit);
if (!row) return json(404, { ok: false, error: "unit_not_found", unit: body.unit });
const updated: CoordinatorUnit = {
const existingReport = reportHead !== undefined ? row.operatorCheckReports?.[reportHead] : undefined;
if (existingReport !== undefined && existingReport.report !== report)
return json(409, { ok: false, error: "operator_check_report_changed", reportHead });
const boundaryExists =
reportHead !== undefined &&
row.rounds.some((round) => round.outcome === "blocked_by_operator_check" && round.reportHead === reportHead);
let updated: CoordinatorUnit = {
...row,
rounds: [
...row.rounds,
{ index: body.index, agent: body.agent, outcome: body.outcome as string, at, ...(gate ? { gate } : {}) },
],
rounds: boundaryExists
? row.rounds
: [
...row.rounds,
{
index: body.index,
agent: body.agent,
outcome: body.outcome as string,
at,
...(gate ? { gate } : {}),
...(reportHead !== undefined ? { reportHead } : {}),
},
],
...(report !== undefined && reportHead !== undefined
? {
operatorCheckReports: {
...(row.operatorCheckReports ?? {}),
[reportHead]: existingReport ?? { report },
},
}
: {}),
};
await deps.instances.putUnits([updated]);
if (host.kind === "host") {
const stored = await deps.instances.putUnits([updated]);
if (!stored.ok) return json(503, { ok: false, error: "round_boundary_unrecorded", at });
if (host.kind === "host" && !boundaryExists) {
const thread = unitThread(instance, updated, units.length);
hostPublish(
deps,
Expand All @@ -2212,6 +2258,7 @@ async function round(body: Record<string, unknown>, deps: AdminCoordinatorDeps):
state: body.outcome as string,
...(thread.threadKey !== undefined ? { threadKey: thread.threadKey } : {}),
...(updated.pr !== undefined ? { pr: updated.pr.number } : {}),
...(report !== undefined ? { report } : {}),
at,
},
],
Expand All @@ -2229,6 +2276,40 @@ async function round(body: Record<string, unknown>, deps: AdminCoordinatorDeps):
).catch((err) =>
(deps.log ?? console.warn)(`[coordinator] ${instance.id}: the card could not be redrawn: ${describe(err)}`),
);
if (report !== undefined && reportHead !== undefined && existingReport?.deliveredAt === undefined) {
const thread = unitThread(instance, updated, units.length);
const io = deps.ioFor({ threadKey: thread.threadKey ?? instance.threadKey, userId: instance.userId });
if (!io) return json(503, { ok: false, error: "operator_check_report_undeliverable", at });
let alreadyPosted: boolean;
try {
const identity = operatorCheckReportIdentity(reportHead);
alreadyPosted = (await io.history()).some((item) => item.role === "assistant" && item.text.includes(identity));
} catch (err) {
(deps.log ?? console.warn)(
`[coordinator] ${instance.id} ${row.unit}: operator-check report history could not be reconciled: ${describe(err)}`,
);
return json(502, { ok: false, error: "operator_check_report_unreconciled", at });
}
if (!alreadyPosted) {
try {
await io.reply(report);
} catch (err) {
(deps.log ?? console.warn)(
`[coordinator] ${instance.id} ${row.unit}: operator-check report could not reach the unit thread: ${describe(err)}`,
);
return json(502, { ok: false, error: "operator_check_report_undelivered", at });
}
}
updated = {
...updated,
operatorCheckReports: {
...(updated.operatorCheckReports ?? {}),
[reportHead]: { report, deliveredAt: at },
},
};
const marked = await deps.instances.putUnits([updated]);
if (!marked.ok) return json(503, { ok: false, error: "operator_check_report_delivery_unrecorded", at });
}
return json(200, { ok: true, at });
}

Expand Down Expand Up @@ -3064,22 +3145,27 @@ async function checksStep(body: Record<string, unknown>, deps: AdminCoordinatorD
};
}
// A head still pending — a run not completed, a required check whose run
// does not exist yet, no check reported, or a draft waiting on its ready
// event — is what the machine's checks wait rides: register it so the
// does not exist yet, no check reported, an operator precondition waiting
// for its re-run, or a draft waiting on its ready event — is what the
// machine's checks wait rides: register it so the
// intake's settled event wakes it (http-ingress item 12), exactly as the
// merge step's pending answer does.
if (
checks === undefined ||
checks.pending.length > 0 ||
(checks.expected?.length ?? 0) > 0 ||
checks.failed.some((failure) => failure.operatorPrecondition === true) ||
checks.total === 0 ||
facts?.draft === true
)
deps.noteMergeWait?.(headSha, id.value, at);
if (checks !== undefined && checks.failed.length > 0)
log(
`[coordinator] ${instance.id} ${body.unit}: CI red at ${headSha.slice(0, 7)} — ${checks.failed
.map((f) => `${f.name} (${f.conclusion}${f.flakeSuspect === true ? ", suspected flake" : ""})`)
.map(
(f) =>
`${f.name} (${f.conclusion}${f.flakeSuspect === true ? ", suspected flake" : ""}${f.operatorPrecondition === true ? ", operator precondition" : ""})`,
)
.join(", ")}`,
);
return json(200, {
Expand Down
23 changes: 23 additions & 0 deletions src/core/coordinator/contract.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -301,6 +301,29 @@ describe("isCoordinatorUnit — one unit's row", () => {
expect(isCoordinatorUnit({ ...unit, resume: "7" })).toBe(false);
});

it("operator-check report delivery is keyed by a commit head and validated separately from its round boundary", () => {
const head = "a".repeat(40);
const report = "Operator action required";
expect(
isCoordinatorUnit({
...unit,
operatorCheckReports: { [head]: { report, deliveredAt: 1_100 } },
rounds: [{ index: 1, agent: "review", outcome: "blocked_by_operator_check", reportHead: head, at: 1_000 }],
}),
).toBe(true);
expect(isCoordinatorUnit({ ...unit, operatorCheckReports: { short: { report } } })).toBe(false);
expect(isCoordinatorUnit({ ...unit, operatorCheckReports: { [head]: { report: "" } } })).toBe(false);
expect(isCoordinatorUnit({ ...unit, operatorCheckReports: { [head]: { report, deliveredAt: "now" } } })).toBe(
false,
);
expect(
isCoordinatorUnit({
...unit,
rounds: [{ index: 1, agent: "review", outcome: "blocked_by_operator_check", reportHead: "short", at: 1_000 }],
}),
).toBe(false);
});

it("the review thread is a thread key with an optional link, beside the unit's own thread; a review thread without its key, with a malformed link, or as a bare string is refused", () => {
expect(isCoordinatorUnit({ ...unit, reviewThread: { threadKey: "slack:C1:3.0" } })).toBe(true);
expect(isCoordinatorUnit({ ...unit, reviewThread: undefined })).toBe(true);
Expand Down
Loading