diff --git a/packages/daemon/src/auto-review.ts b/packages/daemon/src/auto-review.ts index 3d11d54..cbf3c83 100644 --- a/packages/daemon/src/auto-review.ts +++ b/packages/daemon/src/auto-review.ts @@ -222,7 +222,15 @@ export function selectAutoReviewCandidates( request.value.subject !== undefined && refKey(request.value.subject) === key, ).length - candidates.push({ artifact, writer, retractedCount, writerRank: Math.max(0, reviewerOrder.indexOf(writer.did)) }) + const rank = reviewerOrder.indexOf(writer.did) + candidates.push({ + artifact, + writer, + retractedCount, + // A locally loaded writer can precede its still-uningested agent profile. Treat that + // uncertainty as the end of the known order instead of letting it bypass staggering. + writerRank: rank === -1 ? reviewerOrder.length : rank, + }) } } @@ -317,6 +325,8 @@ export class AutoReviewTrigger { readonly #inFlight = new Map>() readonly #retractions = new Map>() readonly #alreadyRetracted = new Set() + /** Duplicate requests whose locally running turn has already been abandoned. */ + readonly #alreadyAbandoned = new Set() readonly #staggerTicks = new Map() /** Per-space persist tail: chained so overlapping ticks never race the tmp+rename. */ readonly #persists = new Map>() @@ -410,6 +420,9 @@ export class AutoReviewTrigger { for (const artifact of observableArtifacts(index)) { const key = refKey(artifact) if (candidateKeys.has(key) || deferred.has(key) || this.#inFlight.has(key)) continue + // The common stagger path ends when another writer's canonical request arrives. Selection + // then omits the subject, so discard its counter as soon as the subject becomes observable. + this.#staggerTicks.delete(key) ledger.subjects.add(key) } ledger.seeded = true @@ -430,10 +443,16 @@ export class AutoReviewTrigger { const byDid = new Map(actors.all.map((actor) => [actor.did, actor])) for (const view of [...index.goals, ...index.projects]) { for (const { duplicate, canonical } of duplicateReviewRequests(view, index.members).values()) { + const reason = `duplicate auto-review request; canonical is ${canonical.uri}` + // Execution and authorship are independent: any operator may be running this open request, + // while only its author can retract it. Every daemon therefore abandons first; the author + // gate below applies only to the protocol write. + if (!this.#alreadyAbandoned.has(duplicate.uri)) { + this.#alreadyAbandoned.add(duplicate.uri) + this.#deps.onAbandon?.(duplicate.uri, reason) + } const writer = byDid.get(duplicate.did) if (!writer || this.#retractions.has(duplicate.uri) || this.#alreadyRetracted.has(duplicate.uri)) continue - const reason = `duplicate auto-review request; canonical is ${canonical.uri}` - this.#deps.onAbandon?.(duplicate.uri, reason) const promise = this.#retract(writer, duplicate, reason).finally(() => this.#retractions.delete(duplicate.uri)) this.#retractions.set(duplicate.uri, promise) } diff --git a/packages/daemon/src/dispatch.ts b/packages/daemon/src/dispatch.ts index f027dfc..85c850e 100644 --- a/packages/daemon/src/dispatch.ts +++ b/packages/daemon/src/dispatch.ts @@ -349,10 +349,11 @@ export function selectDispatchable( view: GoalView | ProjectView, request: IndexedRecord, expectedScope: 'goal' | 'project', + duplicateReviews: ReturnType, ): Dispatchable | undefined => { const type = request.value.type - const duplicate = duplicateReviewRequests(view, index.members).get(request.uri) + const duplicate = duplicateReviews.get(request.uri) if (duplicate) { onReject?.( request.uri, @@ -535,16 +536,18 @@ export function selectDispatchable( // a goal any active member has ended (`GoalView.ended`, whichever way it was ended), and a goal // whose PROJECT the project author archived, which shelves everything under it. for (const goalView of activeGoals(index)) { + const duplicateReviews = duplicateReviewRequests(goalView, index.members) for (const request of goalView.openRequests) { - const dispatchable = evaluate(goalView, request, 'goal') + const dispatchable = evaluate(goalView, request, 'goal', duplicateReviews) if (dispatchable) results.push(dispatchable) } } // Project-scoped (system-artifact) requests: registry types with scope 'project' (adr, architecture) // and review requests whose subject is a project-scoped artifact. Same gate; never implementation. for (const projectView of activeProjects(index)) { + const duplicateReviews = duplicateReviewRequests(projectView, index.members) for (const request of projectView.openRequests) { - const dispatchable = evaluate(projectView, request, 'project') + const dispatchable = evaluate(projectView, request, 'project', duplicateReviews) if (dispatchable) results.push(dispatchable) } } diff --git a/packages/daemon/test/auto-review.test.mjs b/packages/daemon/test/auto-review.test.mjs index 31bf9c2..0020ddb 100644 --- a/packages/daemon/test/auto-review.test.mjs +++ b/packages/daemon/test/auto-review.test.mjs @@ -106,6 +106,7 @@ function verdict(did, rkey, subject) { function makeIndex({ artifacts = [], requests = [], + openRequests = [], reviews = [], retracted = [], members = [AGENT_A], @@ -115,8 +116,10 @@ function makeIndex({ archived = false, projectArtifacts = [], projectRequests = [], + projectOpenRequests = [], projectReviews = [], projectRetracted = [], + agents = [], } = {}) { const spaceUri = `at://${ROOT}/${COLLECTIONS.space}/space` const projectRef = { uri: `at://${ROOT}/${COLLECTIONS.project}/p`, cid: 'cid-project' } @@ -134,7 +137,9 @@ function makeIndex({ ended: closed || archived, artifacts, requests, - openRequests: [], + openRequests, + awaitingInput: [], + winningClaims: {}, retracted, reviews, }, @@ -144,12 +149,15 @@ function makeIndex({ target: { uri: projectRef.uri, cid: projectRef.cid, value: { autoReview: projectAutoReview } }, artifacts: projectArtifacts, requests: projectRequests, - openRequests: [], + openRequests: projectOpenRequests, + awaitingInput: [], + winningClaims: {}, retracted: projectRetracted, reviews: projectReviews, currentSystemArtifacts: [], }, ], + agents, // Convenience refs for callers. _goalRef: goalRef, _projectRef: projectRef, @@ -655,6 +663,72 @@ describe('AutoReviewTrigger.pump', () => { assert.equal(creates.length, 0) assert.equal(logs.filter((m) => /produces review artifacts/.test(m)).length, 1) }) + + it('stagger-defers a rank-1 writer and suppresses it when the canonical arrives', async () => { + const creates = [] + const base = makeIndex({}) + const req = request(HUMAN, 'req', { goal: base._goalRef }) + const art = artifact(AGENT_A, 'plan-1', req, { goal: base._goalRef }) + const profiles = [AGENT_A, AGENT_B].map((did) => ({ + did, + value: { artifactTypes: ['review'] }, + })) + const seed = makeIndex({ projectAutoReview: { plan: true }, members: [AGENT_A, AGENT_B], agents: profiles }) + const live = makeIndex({ + artifacts: [art], requests: [req], projectAutoReview: { plan: true }, + members: [AGENT_A, AGENT_B], agents: profiles, + }) + const canonical = reviewRequest(AGENT_A, 'canonical', art, { goal: base._goalRef }) + const satisfied = makeIndex({ + artifacts: [art], requests: [req, canonical], openRequests: [canonical], + projectAutoReview: { plan: true }, members: [AGENT_A, AGENT_B], agents: profiles, + }) + const t = trigger(stateDir()) + const onlyRankOne = registry([fakeActor(AGENT_B, ['review'], creates)]) + t.pump(seed, onlyRankOne); await t.drain() + t.pump(live, onlyRankOne); await t.drain() + assert.equal(creates.length, 0, 'rank 1 waits during its stagger window') + t.pump(satisfied, onlyRankOne); await t.drain() + t.pump(satisfied, onlyRankOne); await t.drain() + assert.equal(creates.length, 0, 'the observed canonical suppresses the deferred write') + }) + + it('abandons a losing turn on every operator but retracts only from the author', async () => { + const base = makeIndex({}) + const subjectRequest = request(HUMAN, 'req', { goal: base._goalRef }) + const art = artifact(AGENT_A, 'plan-1', subjectRequest, { goal: base._goalRef }) + const canonical = reviewRequest(HUMAN, 'human-review', art, { goal: base._goalRef }) + const duplicate = reviewRequest(AGENT_A, 'auto-review', art, { goal: base._goalRef }) + const index = makeIndex({ + artifacts: [art], requests: [subjectRequest, canonical, duplicate], + openRequests: [canonical, duplicate], + members: [{ did: HUMAN, kind: 'human' }, AGENT_A, AGENT_B], + }) + + const executorWrites = [] + const executorAbandons = [] + const executor = new AutoReviewTrigger({ + stateDir: stateDir(), + onAbandon: (uri, reason) => executorAbandons.push({ uri, reason }), + }) + executor.pump(index, registry([fakeActor(AGENT_B, ['review'], executorWrites)])) + await executor.drain() + assert.deepEqual(executorAbandons.map((entry) => entry.uri), [duplicate.uri]) + assert.equal(executorWrites.length, 0, 'a non-author cannot retract the request') + + const authorWrites = [] + const authorAbandons = [] + const author = new AutoReviewTrigger({ + stateDir: stateDir(), + onAbandon: (uri) => authorAbandons.push(uri), + }) + author.pump(index, registry([fakeActor(AGENT_A, ['review'], authorWrites)])) + await author.drain() + assert.deepEqual(authorAbandons, [duplicate.uri]) + assert.equal(authorWrites.length, 1) + assert.equal(authorWrites[0].collection, COLLECTIONS.retractRequest) + assert.deepEqual(authorWrites[0].value.request, { uri: duplicate.uri, cid: duplicate.cid }) + }) }) // --- Round-trip through the REAL core materializer: the written request must land in the target @@ -835,6 +909,44 @@ it('the emitted open review is claimable across operators and only the confirmed } }) +it('duplicate review requests are neither claimable nor dispatchable', () => { + const { records, space, art } = buildStore({ scope: 'goal' }) + const canonical = stored(HUMAN, 'human-review', { + $type: COLLECTIONS.artifactRequest, + goal: art.value.goal, + type: 'review', + subject: strongref(art), + basedOn: [strongref(art)], + createdAt: '2026-01-01T00:05:00Z', + }, 20, 'cid-human-review') + const duplicate = stored(AGENT_A, 'auto-review', { + $type: COLLECTIONS.artifactRequest, + goal: art.value.goal, + type: 'review', + subject: strongref(art), + basedOn: [strongref(art)], + createdAt: '2026-01-01T00:04:00Z', + }, 21, 'cid-auto-review') + const store = new MemoryRecordStore() + for (const record of [...records, canonical, duplicate]) store.put(record) + const index = materialize(store, { spaceUri: space.uri }) + const actors = registry([fakeActor(AGENT_A, ['review'], [])]) + const claims = new ClaimLedger() + const turns = new TurnLedger() + try { + assert.deepEqual( + selectClaimable(index, actors, claims, turns, { budget: 5 }).map((candidate) => candidate.request.uri), + [canonical.uri], + ) + const rejected = [] + selectDispatchable(index, actors, turns, undefined, (uri, reason) => rejected.push({ uri, reason })) + assert.match(rejected.find((entry) => entry.uri === duplicate.uri).reason, /duplicate review request.*canonical/) + } finally { + claims.close() + turns.close() + } +}) + // --- Write/ledger boundary (cross-model review findings 1-5). --- describe('AutoReviewTrigger.pump — write/ledger boundary', () => {