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
11 changes: 10 additions & 1 deletion bin/jury.js
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@ import { buildPrompt } from "../lib/prompt.js";
import { serve } from "../lib/server.js";
import { appendEvent, writeRun, readEvents, foldEvents, slugFor, attemptStamp, runsDir, writeArtifact, openArtifact } from "../lib/store.js";
import { findingsIn, gate, settledList, VERDICTS } from "../lib/findings.js";
import { MAX_TURNS, turnsFor, outstanding, deadlocked, refreshSettled, replyRound, record, sessionsIn, openFindings, currentJudge } from "../lib/loop.js";
import { MAX_TURNS, turnsFor, outstanding, deadlocked, refreshSettled, replyRound, record, sessionsIn, judgeSessionIn, openFindings, currentJudge } from "../lib/loop.js";
import { triageOne } from "../lib/triage.js";
import { assertPrCheckout, repositoryFromPrUrl, resolvePrCheckout } from "../lib/repository.js";
import { resolveJuryDirectory } from "../lib/directories.js";
Expand Down Expand Up @@ -612,6 +612,9 @@ async function cmdAgent(argv) {
}
}
const first = Math.max(0, ...prior.filter((e) => e.t === "round.start").map((e) => e.n)) + 1;
// Keep one judge conversation for every finding and every round. The event
// makes a resumed process pick up the same provider session after a crash.
let judgeSessionId = judgeSessionIn(prior, judge.name);

if (roles && !values.resume) {
const names = knownAgents(cfg).filter(a => a.name === roles.judge || roles.reviewers.includes(a.name)).map(a => a.name);
Expand Down Expand Up @@ -752,12 +755,18 @@ async function cmdAgent(argv) {
try {
v = await triageOne(judge, f, {
worktree, trunk, context: group ? groupContext(group.targets) : "", stopToken: cfg.stopToken, dryRun: values["dry-run"],
sessionId: judgeSessionId,
onLog: (m) => console.log(` ${st.muted(m)}`),
});
} finally {
clearInterval(heart);
}

if (v.sessionId && v.sessionId !== judgeSessionId) {
judgeSessionId = v.sessionId;
await appendEvent(dir, { t: "judge.session", agent: judge.name, sessionId: judgeSessionId, round });
}

if (!v.verdict) {
// Left open on purpose: an unjudged finding must be raised again rather
// than silently disappearing into a round that reports convergence.
Expand Down
19 changes: 11 additions & 8 deletions lib/agents.js
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@ function subst(argv, vars) {
* dead reviewer is a result, not a crash, and the round should still report the
* others.
*/
export async function runAgent(agent, { worktree, prompt, stopToken, dryRun, onLog, onChunk, timeoutSeconds }) {
export async function runAgent(agent, { worktree, prompt, stopToken, dryRun, onLog, onChunk, timeoutSeconds, sessionId = null }) {
const started = Date.now();

if (dryRun) {
Expand All @@ -49,15 +49,18 @@ export async function runAgent(agent, { worktree, prompt, stopToken, dryRun, onL
await writeFile(promptFile, prompt, "utf8");
}

// Resuming uses the exact conversation that judged the previous finding.
// An agent that takes its session id as input gets one generated here, so the
// reply can name the exact conversation instead of asking for "the last one"
// and hoping nothing else ran in between.
const assigned = agent.newSession && agent.argv.some(arg => arg.includes("{{sessionId}}")) ? randomUUID() : null;
const resumable = Boolean(sessionId && agent.resume?.supported && agent.resume?.argv?.length);
const commandArgv = resumable ? agent.resume.argv : agent.argv;
const assigned = !resumable && agent.newSession && commandArgv.some(arg => arg.includes("{{sessionId}}")) ? randomUUID() : null;
const vars = {
worktree, promptFile: promptFile ?? "", promptText: prompt,
sessionId: assigned ?? "", packageDir: PACKAGE_DIR,
sessionId: sessionId ?? assigned ?? "", packageDir: PACKAGE_DIR,
};
const [cmd, ...args] = subst(agent.argv, vars);
const [cmd, ...args] = subst(commandArgv, vars);
const cwd = agent.cwd === "worktree" ? worktree : process.cwd();

// The caller already prints the agent's name at the head of this line, so
Expand Down Expand Up @@ -119,13 +122,13 @@ export async function runAgent(agent, { worktree, prompt, stopToken, dryRun, onL
// A resumable agent that prints a session id gets it captured here. Without
// it a reply can only say "resume the last session", which is the review only
// if nothing else ran in the meantime.
let sessionId = assigned;
if (!sessionId && agent.resume?.idFrom) {
let returnedSessionId = sessionId ?? assigned;
if (!returnedSessionId && agent.resume?.idFrom) {
try {
const re = agent.resume.idFrom instanceof RegExp
? agent.resume.idFrom
: new RegExp(agent.resume.idFrom, "i");
sessionId = out.stdout.match(re)?.[1] ?? null;
returnedSessionId = out.stdout.match(re)?.[1] ?? null;
} catch { /* a bad pattern must not fail the review */ }
}

Expand All @@ -152,7 +155,7 @@ export async function runAgent(agent, { worktree, prompt, stopToken, dryRun, onL
contradicted: said && findings.length > 0,
findings,
report,
sessionId,
sessionId: returnedSessionId,
raw: out.stdout,
};
}
Expand Down
10 changes: 10 additions & 0 deletions lib/loop.js
Original file line number Diff line number Diff line change
Expand Up @@ -429,3 +429,13 @@ export function sessionsIn(events) {
}
return byAgent;
}

/** The judge's exact conversation, persisted separately from reviewer threads. */
export function judgeSessionIn(events, judge) {
let sessionId = null;
for (const e of events) {
if (e.t !== "judge.session" || (judge && e.agent !== judge)) continue;
sessionId = e.sessionId ?? null;
}
return sessionId;
}
11 changes: 6 additions & 5 deletions lib/triage.js
Original file line number Diff line number Diff line change
Expand Up @@ -125,21 +125,22 @@ export function parseVerdict(text) {
* other verdicts, so a failure is a null verdict and the finding stays open for
* the next round to raise again.
*/
export async function triageOne(main, finding, { worktree, trunk, stopToken, dryRun, onLog, onChunk, context }) {
export async function triageOne(main, finding, { worktree, trunk, stopToken, dryRun, onLog, onChunk, context, sessionId = null }) {
const prompt = triagePrompt({ finding, trunk, stopToken, context });
if (dryRun) {
return { verdict: "rejected", reproduced: null, reason: "dry run — not triaged", test: null, seconds: 0 };
return { verdict: "rejected", reproduced: null, reason: "dry run — not triaged", test: null, seconds: 0, sessionId };
}
const r = await runAgent(main, {
worktree, prompt, stopToken, onLog, onChunk,
// Triage is slower than review: it reads, reproduces, edits and runs tests.
// Floor as well as multiple: a misconfigured 0 must not mean 'kill it now'.
timeoutSeconds: Math.max(900, (main.expectSeconds || 900) * 3),
sessionId,
});
if (!r.ok) return { verdict: null, failed: true, report: r.report, seconds: r.seconds };
if (!r.ok) return { verdict: null, failed: true, report: r.report, seconds: r.seconds, sessionId: r.sessionId ?? sessionId };
const v = parseVerdict(r.report) ?? parseVerdict(r.raw);
if (!v?.verdict) {
return { verdict: null, unparsed: true, report: r.report, seconds: r.seconds };
return { verdict: null, unparsed: true, report: r.report, seconds: r.seconds, sessionId: r.sessionId ?? sessionId };
}
return { ...v, report: r.report, seconds: r.seconds };
return { ...v, report: r.report, seconds: r.seconds, sessionId: r.sessionId ?? sessionId };
}
42 changes: 42 additions & 0 deletions test/judge-session.test.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,42 @@
import test from "node:test";
import assert from "node:assert/strict";
import { mkdtemp, writeFile, readFile, rm, chmod } from "node:fs/promises";
import { tmpdir } from "node:os";
import path from "node:path";
import { triageOne } from "../lib/triage.js";

test("judge triage resumes one provider session across findings", async () => {
const dir = await mkdtemp(path.join(tmpdir(), "jury-judge-session-"));
const calls = path.join(dir, "calls.log");
const fake = path.join(dir, "fake-judge.js");
await writeFile(fake, `#!/usr/bin/env node
const fs = require("node:fs");
fs.appendFileSync(${JSON.stringify(calls)}, (process.argv.includes("--resume") ? "resume " + process.argv[process.argv.indexOf("--resume") + 1] : "first") + "\\n");
console.log("SESSION=judge-session-1");
console.log(JSON.stringify({reproduced:null, verdict:"rejected", reason:"not reproducible", test:null}));
`);
await chmod(fake, 0o755);
const agent = {
name: "fake-judge", promptDelivery: "argv", cwd: "worktree",
argv: [fake, "{{promptText}}"],
resume: { supported: true, argv: [fake, "--resume", "{{sessionId}}", "{{promptText}}"], idFrom: "SESSION=([a-z0-9-]+)" },
report: "whole", expectSeconds: 1,
};
const finding = { id: "F1", agent: "reviewer", claim: "something is wrong", loc: "a.txt:1" };
try {
const first = await triageOne(agent, finding, { worktree: dir, trunk: "master", stopToken: "NO NEW FINDINGS", context: "" });
assert.equal(first.verdict, "rejected");
assert.equal(first.sessionId, "judge-session-1");
const second = await triageOne(agent, { ...finding, id: "F2" }, {
worktree: dir, trunk: "master", stopToken: "NO NEW FINDINGS", context: "", sessionId: first.sessionId,
});
assert.equal(second.verdict, "rejected");
assert.equal(second.sessionId, first.sessionId);
const lines = (await readFile(calls, "utf8")).trim().split("\n");
assert.equal(lines.length, 2);
assert.equal(lines[0], "first");
assert.equal(lines[1], "resume judge-session-1");
} finally {
await rm(dir, { recursive: true, force: true });
}
});
3 changes: 2 additions & 1 deletion test/kimi.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -12,10 +12,11 @@ import { mkdtemp, mkdir, writeFile, readFile, readdir, rm, access } from "node:f
import { tmpdir } from "node:os";
import path from "node:path";
import { execFileSync, spawnSync } from "node:child_process";
import { realpathSync } from "node:fs";
import { fileURLToPath, pathToFileURL } from "node:url";

const cli = process.env.JURY_TEST_CLI ?? fileURLToPath(new URL("../bin/jury.js", import.meta.url));
const root = path.dirname(path.dirname(cli));
const root = realpathSync(path.dirname(path.dirname(cli)));
const { DEFAULTS, loadConfig, judgeAgent, knownAgents } = await import(pathToFileURL(path.join(root, "lib/config.js")));
const { runAgent, hasStopToken } = await import(pathToFileURL(path.join(root, "lib/agents.js")));
const { replyArgv } = await import(pathToFileURL(path.join(root, "lib/reply.js")));
Expand Down
Loading