fix(process): bind Validate resume to the current pending round

Issue: #775
User-Visible: no
This commit is contained in:
Matysh
2026-10-01 16:44:28 +00:00
committed by claude[bot]
parent dd86bf955a
commit 09fb02cf15
7 changed files with 374 additions and 63 deletions
+44 -4
View File
@@ -10328,8 +10328,8 @@ const MUTANT_DEFINITIONS = [
+ 'it would spend the model a second time (#636)',
patches: [{
file: 'scripts/process-resume.mjs',
find: " if (!pending) return { action: 'noop', reason: 'latest process run left no pending marker — it did not wait for Validate' };",
replace: " if (false && !pending) return { action: 'noop', reason: 'latest process run left no pending marker — it did not wait for Validate' };",
find: " if (!pending) return { action: 'noop', recheck: true, reason: 'latest process run left no pending marker — it did not wait for Validate' };",
replace: " if (false && !pending) return { action: 'noop', recheck: true, reason: 'latest process run left no pending marker — it did not wait for Validate' };",
}],
},
{
@@ -10338,8 +10338,48 @@ const MUTANT_DEFINITIONS = [
because: 'relabelling while a process run is active queues a second round for the same request (#636)',
patches: [{
file: 'scripts/process-resume.mjs',
find: " if (mine.some((run) => ACTIVE_RUN_STATES.has(run.status))) return { action: 'noop', reason: 'a process run for this issue is already active' };",
replace: " if (false && mine.some((run) => ACTIVE_RUN_STATES.has(run.status))) return { action: 'noop', reason: 'a process run for this issue is already active' };",
find: " if (mine.some((run) => ACTIVE_RUN_STATES.has(run.status))) return { action: 'noop', recheck: true, reason: 'a process run for this issue is already active' };",
replace: " if (false && mine.some((run) => ACTIVE_RUN_STATES.has(run.status))) return { action: 'noop', recheck: true, reason: 'a process run for this issue is already active' };",
}],
},
{
id: 'pending-history-latest-200-only',
guard: 'node --test test/process-pending-round.test.mjs',
because: '#775: unrelated label traffic must not hide the pending review beyond two pages',
patches: [{
file: 'scripts/process-reconcile.mjs',
find: ' for (let page = 1; ; page++) {',
replace: ' for (let page = 1; page <= 2; page++) {',
}],
},
{
id: 'pending-revives-previous-request',
guard: 'node --test test/process-pending-round.test.mjs',
because: '#775: a re-applied label starts a new request even one second after the preceding run',
patches: [{
file: 'scripts/process-reconcile.mjs',
find: ' && at(run.createdAt) >= at(request.at))',
replace: ' && at(run.createdAt) >= at(request.at) - 120_000)',
}],
},
{
id: 'pending-resumes-another-validate',
guard: 'node --test test/process-pending-round.test.mjs',
because: '#775: the same SHA is not proof that the pending round awaited this dispatch',
patches: [{
file: 'scripts/process-resume.mjs',
find: ' if (pending.branch !== branch || String(pending.validate_run_id) !== String(validateRun.id)) {',
replace: ' if (pending.branch !== branch) {',
}],
},
{
id: 'pending-no-fresh-snapshot-before-write',
guard: 'node --test test/process-pending-round.test.mjs',
because: '#775: a stopped or already restarted round must not be woken from an old snapshot',
patches: [{
file: 'scripts/process-resume.mjs',
find: ' const second = inspect();',
replace: ' const second = first;',
}],
},
{
+63 -22
View File
@@ -72,6 +72,32 @@ export function latestReviewRequest(events = [], label = null) {
return requests.at(-1) || null;
}
/** Runs belonging to this request only. Never revive a previous label round. */
export function reviewRunsForRequest(runs, issue, request) {
if (!request || !Number.isFinite(at(request.at))) return [];
const matching = runs.map((run) => run.issue ? run : parseProcessRun(run)).filter(Boolean)
.filter((run) => run.issue === Number(issue) && run.label === request.label
&& at(run.createdAt) >= at(request.at))
.sort((a, b) => at(b.createdAt) - at(a.createdAt) || Number(b.id) - Number(a.id));
// Fully skipped workflows ran no review. All other conclusions remain barriers:
// searching past a failed or successful non-pending run risks a second model call.
const attempted = matching.filter((run) => run.conclusion !== 'skipped');
return attempted.length ? attempted : matching;
}
export function pendingEvidenceError(pending, { issue, run }) {
if (pending.schema !== 1 || String(pending.issue) !== String(issue)
|| String(pending.run_id) !== String(run.id) || String(pending.run_attempt) !== String(run.attempt)
|| pending.stage !== run.stage || !String(pending.branch || '').startsWith(`issue/${issue}-`)) {
return 'pending marker belongs to another issue/stage/run attempt/branch';
}
if (!/^[0-9a-f]{40}$/i.test(String(pending.material_sha || ''))
|| !/^[1-9][0-9]*$/.test(String(pending.validate_run_id || ''))) {
return 'pending marker has incomplete material SHA or Validate run ID';
}
return null;
}
export function preparedEvidenceError(prepared, { issue, stage, run }) {
if (!prepared) return null;
const badIdentity = prepared.schema !== 1
@@ -112,13 +138,8 @@ export function decideReconciliation({
}
const requestAt = at(request.at);
const matching = runs
.map((run) => run.issue ? run : parseProcessRun(run))
.filter(Boolean)
.filter((run) => run.issue === Number(issue.number) && run.label === label
&& Number.isFinite(at(run.createdAt)) && at(run.createdAt) >= requestAt - 120_000)
.sort((a, b) => at(b.createdAt) - at(a.createdAt) || Number(b.id) - Number(a.id));
const run = matching[0] || null;
const matching = reviewRunsForRequest(runs, issue.number, request);
const run = matching.find((candidate) => ACTIVE_RUN_STATES.has(candidate.status)) || matching[0] || null;
if (!run) {
if (now - requestAt < graceMs) return result('wait', 'label event is still within delivery grace', { label, stage });
@@ -147,6 +168,9 @@ export function decideReconciliation({
if (Number.isFinite(settledAt) && now - settledAt < graceMs) {
return result('wait', 'completed run is still within label-application grace', { label, stage, run });
}
if (run.resultArtifact) {
return result('escalate', 'sealed model result exists but was not integrated; automatic rerun would spend the model twice', { label, stage, run });
}
if (run.conclusion === 'success' && run.pending) {
// #636: успешный прогон без вердикта — это не потеря, а осознанный выход
// подготовки: Validate с мутантами на материале ещё шёл. Пока он идёт —
@@ -160,9 +184,6 @@ export function decideReconciliation({
if (run.conclusion === 'success') {
return result('escalate', 'successful run did not move the review label', { label, stage, run });
}
if (run.resultArtifact) {
return result('escalate', 'sealed model result exists but was not integrated; automatic rerun would spend the model twice', { label, stage, run });
}
if (RETRYABLE_CONCLUSIONS.has(run.conclusion)) {
return result('retry', `transient process run conclusion: ${run.conclusion}`, { label, stage, run });
}
@@ -201,7 +222,7 @@ function issueView(repo, number) {
return ghJson(['issue', 'view', String(number), '--repo', repo, '--json', 'number,title,labels,comments,updatedAt']);
}
function issueEvents(repo, number) {
export function issueEvents(repo, number) {
const pages = ghJson(['api', '--paginate', '--slurp', `repos/${repo}/issues/${number}/events?per_page=100`]);
return Array.isArray(pages?.[0]) ? pages.flat() : (Array.isArray(pages) ? pages : []);
}
@@ -216,12 +237,34 @@ function openReviewIssues(repo) {
.sort((a, b) => a.number - b.number);
}
export function processRuns(repo, issues = []) {
const pages = [1, 2].flatMap((page) => {
const response = ghJson(['api', `repos/${repo}/actions/workflows/process.yml/runs?event=issues&per_page=100&page=${page}`]);
return response.workflow_runs || [];
});
return pages.map((raw) => {
export function processRuns(repo, issues = [], { since = null, getJson = ghJson } = {}) {
// #775: the former latest-200 window could silently hide the pending round.
// Bound history by the current label request, not unrelated workflow volume.
const starts = since ? [at(since)] : issues.map((issue) => {
const label = REVIEW_LABELS.find((candidate) => labelsOf(issue).includes(candidate));
return at(latestReviewRequest(issueEvents(repo, issue.number), label)?.at);
}).filter(Number.isFinite);
if (!starts.length) return [];
const oldest = Math.min(...starts);
if (!Number.isFinite(oldest)) throw new Error('invalid process history start');
const rows = new Map();
for (let page = 1; ; page++) {
// Do not use event/created search filters: filtered Actions queries cap at
// 1000 results. Filter issue events locally after complete pagination.
const response = getJson(['api', `repos/${repo}/actions/workflows/process.yml/runs?per_page=100&page=${page}`]);
const batch = response.workflow_runs;
if (!Array.isArray(batch)) throw new Error('invalid process history response');
if (!batch.length) break;
let added = 0;
for (const raw of batch) {
if (!raw.id || !Number.isFinite(at(raw.created_at))) throw new Error('invalid process history run');
const key = `${raw.id}/${raw.run_attempt || 1}`;
if (!rows.has(key)) { rows.set(key, raw); added++; }
}
if (!added) throw new Error('process history pagination made no progress');
if (batch.length < 100 || batch.every((raw) => at(raw.created_at) < oldest)) break;
}
return [...rows.values()].filter((raw) => raw.event === 'issues' && at(raw.created_at) >= oldest).map((raw) => {
const stable = parseProcessRun(raw);
if (stable) return stable;
// Runs created before #555 used the issue title as display_title. Accept
@@ -303,9 +346,8 @@ function hydrateRunEvidence(repo, issue, run) {
return { ...run, evidenceError: 'duplicate prepared/result/pending artifacts' };
}
const pending = pendingArtifacts.length === 1 ? loadSealedArtifact(repo, run, pendingArtifacts[0], 'pending.json') : null;
if (pending && (String(pending.issue) !== String(issue.number) || String(pending.run_id) !== String(run.id))) {
return { ...run, evidenceError: 'pending marker belongs to another issue/run' };
}
const pendingError = pending && pendingEvidenceError(pending, { issue: issue.number, run });
if (pendingError) return { ...run, evidenceError: pendingError };
return {
...run,
preparedArtifact: preparedArtifacts.length === 1,
@@ -370,8 +412,7 @@ async function snapshot(repo, baseRuns, issue) {
const labels = labelsOf(fresh);
const label = REVIEW_LABELS.find((candidate) => labels.includes(candidate)) || null;
const request = latestReviewRequest(events, label);
const candidates = baseRuns.filter((run) => run.issue === issue.number && run.label === label)
.sort((a, b) => at(b.createdAt) - at(a.createdAt));
const candidates = reviewRunsForRequest(baseRuns, issue.number, request);
const hydrated = candidates.length ? [hydrateRunEvidence(repo, fresh, candidates[0]), ...candidates.slice(1)] : candidates;
return { issue: fresh, request, runs: hydrated };
}
+71 -32
View File
@@ -16,10 +16,11 @@
// модели. Страховка на потерянное событие — process-reconcile (тот же маркер).
import { execFileSync } from 'node:child_process';
import { appendFileSync } from 'node:fs';
import { setTimeout as delay } from 'node:timers/promises';
import { isMainModule } from './spawn-portable.mjs';
import {
ACTIVE_RUN_STATES, artifactNames, loadSealedArtifact, parseProcessRun, pendingArtifactName,
processRuns, relabel,
ACTIVE_RUN_STATES, artifactNames, issueEvents, latestReviewRequest, loadSealedArtifact,
pendingArtifactName, pendingEvidenceError, processRuns, relabel, reviewRunsForRequest,
} from './process-reconcile.mjs';
export const REVIEW_LABEL = 'S7-code-review';
@@ -31,42 +32,49 @@ export function issueNumberFromBranch(branch) {
return match ? Number(match[1]) : null;
}
const at = (value) => {
const parsed = Date.parse(String(value || ''));
return Number.isFinite(parsed) ? parsed : NaN;
};
/**
* Чистое решение: будить раунд или нет.
*
* @param {object} p
* @param {string[]} p.labels метки issue
* @param {object} p.validateRun завершённый прогон Validate: { event, status, headSha }
* @param {object} p.validateRun завершённый прогон Validate: { id, event, status, headSha }
* @param {number} p.issue номер issue
* @param {object} p.request последняя постановка S7 из timeline
* @param {string} p.branch ветка материала
* @param {string} p.headSha текущая вершина ветки
* @param {string} p.sha SHA материала, на котором завершился Validate
* @param {object[]} p.runs прогоны конвейера этой issue (parseProcessRun-совместимые)
* @param {(run) => object|null} p.pendingOf маркер ожидания прогона либо null
* @returns {{action:'resume'|'noop', reason:string, run?:object}}
* @returns {{action:'resume'|'noop', reason:string, run?:object, recheck?:boolean}}
*/
export function decideResume({ labels = [], validateRun, sha, runs = [], pendingOf = () => null }) {
if (validateRun?.event !== 'workflow_dispatch') return { action: 'noop', reason: 'not a dispatch run — push runs carry no mutants' };
export function decideResume({ labels = [], issue, request, branch, headSha, validateRun, sha, runs = [], pendingOf = () => null }) {
if (validateRun?.event !== 'workflow_dispatch') return { action: 'noop', reason: 'not a dispatch run — only the awaited dispatch wakes the round' };
if (validateRun?.status !== 'completed') return { action: 'noop', reason: 'validate run is not completed' };
if (validateRun?.headSha && sha && validateRun.headSha !== sha) return { action: 'noop', reason: 'validate run head differs from the material' };
if (!labels.includes(REVIEW_LABEL)) return { action: 'noop', reason: 'issue is not awaiting code review' };
if (labels.includes('S4-spec-review')) return { action: 'noop', reason: 'ambiguous review status labels' };
const stop = STOP_LABELS.find((label) => labels.includes(label));
if (stop) return { action: 'noop', reason: `owner stopped the review (${stop})` };
const mine = runs.map((run) => run.issue ? run : parseProcessRun(run)).filter(Boolean)
.filter((run) => run.label === REVIEW_LABEL)
.sort((a, b) => at(b.createdAt) - at(a.createdAt) || Number(b.id) - Number(a.id));
if (mine.some((run) => ACTIVE_RUN_STATES.has(run.status))) return { action: 'noop', reason: 'a process run for this issue is already active' };
if (request?.label !== REVIEW_LABEL || !Number.isFinite(Date.parse(request.at))) {
return { action: 'noop', reason: 'current review label has no timeline request' };
}
if (!sha || headSha !== sha) return { action: 'noop', reason: 'branch head differs from the completed Validate material' };
const mine = reviewRunsForRequest(runs, issue, request);
if (mine.some((run) => ACTIVE_RUN_STATES.has(run.status))) return { action: 'noop', recheck: true, reason: 'a process run for this issue is already active' };
const latest = mine[0];
if (!latest) return { action: 'noop', reason: 'no process run for this issue — reconcile owns lost requests' };
if (!latest) return { action: 'noop', recheck: true, reason: 'no process run for this issue — reconcile owns lost requests' };
if (latest.status !== 'completed' || latest.conclusion !== 'success') {
return { action: 'noop', reason: `latest process run is ${latest.status}/${latest.conclusion || 'none'} — nothing was left pending` };
}
const pending = pendingOf(latest);
if (!pending) return { action: 'noop', reason: 'latest process run left no pending marker — it did not wait for Validate' };
if (!pending) return { action: 'noop', recheck: true, reason: 'latest process run left no pending marker — it did not wait for Validate' };
const error = pendingEvidenceError(pending, { issue, run: latest });
if (error) return { action: 'noop', reason: error };
if (pending.branch !== branch || String(pending.validate_run_id) !== String(validateRun.id)) {
return { action: 'noop', reason: 'pending marker waits for another branch or Validate run' };
}
if (String(pending.material_sha) !== String(sha)) {
return { action: 'noop', reason: `pending marker waits for ${String(pending.material_sha).slice(0, 8)}, not ${String(sha).slice(0, 8)}` };
return { action: 'noop', recheck: true, reason: `pending marker waits for ${String(pending.material_sha).slice(0, 8)}, not ${String(sha).slice(0, 8)}` };
}
return { action: 'resume', reason: 'the round was waiting for exactly this Validate run', run: latest };
}
@@ -75,25 +83,56 @@ function gh(args) {
return execFileSync('gh', args, { encoding: 'utf8' });
}
export async function resume({ repo, branch, sha, validateRun, apply = true, ops = null }) {
export async function resume({ repo, branch, sha, validateRun, apply = true, ops = null, sleep = delay }) {
const issue = issueNumberFromBranch(branch);
if (!issue) return { action: 'noop', reason: `branch ${branch} is not an issue branch`, issue: null };
const io = ops || {
labels: () => JSON.parse(gh(['issue', 'view', String(issue), '--repo', repo, '--json', 'labels'])).labels.map((label) => label.name),
runs: () => processRuns(repo, []),
state: () => {
const current = JSON.parse(gh(['issue', 'view', String(issue), '--repo', repo, '--json', 'labels,state']));
const labels = current.state === 'OPEN' ? current.labels.map((label) => label.name) : [];
// A deleted branch after merge/closure is normal, not a failed recovery.
if (!labels.includes(REVIEW_LABEL) || labels.includes('S4-spec-review') || STOP_LABELS.some((label) => labels.includes(label))) {
return { labels, request: null, headSha: null };
}
return { labels, request: latestReviewRequest(issueEvents(repo, issue), REVIEW_LABEL),
headSha: JSON.parse(gh(['api', `repos/${repo}/git/ref/heads/${branch}`])).object.sha };
},
runs: (request) => request ? processRuns(repo, [], { since: request.at }) : [],
pendingOf: (run) => {
const name = pendingArtifactName(issue, run);
const artifact = artifactNames(repo, run).find((item) => item.name === name && !item.expired);
if (!artifact) return null;
const pending = loadSealedArtifact(repo, run, artifact, 'pending.json');
return String(pending.issue) === String(issue) && String(pending.run_id) === String(run.id) ? pending : null;
const artifacts = artifactNames(repo, run);
if (artifacts.some((item) => item.name === `review-result-${issue}-${run.id}-${run.attempt}` && !item.expired)) {
throw new Error('sealed model result exists — refusing to wake a completed review');
}
const matches = artifacts.filter((item) => item.name === name && !item.expired);
if (matches.length > 1) throw new Error('duplicate pending artifacts');
return matches.length ? loadSealedArtifact(repo, run, matches[0], 'pending.json') : null;
},
relabel: () => relabel(repo, { number: issue }, REVIEW_LABEL),
};
const labels = io.labels();
const runs = io.runs().filter((run) => run.issue === issue);
const decision = decideResume({ labels, validateRun, sha, runs, pendingOf: io.pendingOf });
if (decision.action === 'resume' && apply) io.relabel();
const inspect = () => {
const state = io.state();
const decision = decideResume({ ...state, issue, branch, validateRun, sha,
runs: io.runs(state.request), pendingOf: io.pendingOf });
return { state, decision };
};
let first = inspect();
// A completed Validate can arrive before the process list/artifact is visible.
// Three fresh snapshots at most; never search behind a completed review barrier.
// After this short delivery grace, the scheduled reconciler remains the fallback.
for (let attempt = 1; first.decision.recheck && attempt < 3; attempt++) {
await sleep(5_000);
first = inspect();
}
const decision = first.decision;
if (decision.action === 'resume' && apply) {
const second = inspect();
if (second.decision.action !== 'resume' || first.state.request.id !== second.state.request?.id
|| decision.run.id !== second.decision.run?.id || decision.run.attempt !== second.decision.run?.attempt) {
return { action: 'noop', reason: `state changed before write: ${second.decision.reason}`, issue, applied: false };
}
io.relabel();
}
return { ...decision, issue, applied: decision.action === 'resume' && apply };
}
@@ -102,13 +141,13 @@ if (isMainModule(import.meta.url)) {
const repo = arg('repo') || process.env.GITHUB_REPOSITORY;
const branch = arg('branch');
const sha = arg('sha');
if (!repo || !branch || !sha) {
console.error('usage: process-resume.mjs --repo=<owner/repo> --branch=<issue/NN-slug> --sha=<sha> --event=<event> --status=<status> [--apply=false]');
if (!repo || !branch || !sha || !arg('run-id')) {
console.error('usage: process-resume.mjs --repo=<owner/repo> --branch=<issue/NN-slug> --sha=<sha> --run-id=<Validate id> --event=<event> --status=<status> [--apply=false]');
process.exit(2);
}
const outcome = await resume({
repo, branch, sha,
validateRun: { event: arg('event'), status: arg('status') || 'completed', headSha: sha },
validateRun: { id: arg('run-id'), event: arg('event'), status: arg('status') || 'completed', headSha: sha },
apply: arg('apply') !== 'false',
});
const lines = [`action=${outcome.action}`, `issue=${outcome.issue ?? ''}`, `reason=${outcome.reason}`, `applied=${outcome.applied ? 'true' : 'false'}`];