From 98edd35ce0fab4ed5aaa665ab008c6f9a5da3f3e Mon Sep 17 00:00:00 2001 From: Adrian Webb Date: Mon, 28 Sep 2026 22:18:22 -0400 Subject: [PATCH 01/10] Scope governance review history to proposal review nodes --- .../admission/living-allocation-inputs.ts | 5 ++- .../execution/assignment-attempt-builder.ts | 1 + .../execution/living-execution-assignment.ts | 3 +- .../admission/allocation-inputs.test.ts | 36 ++++++++++++++++++- 4 files changed, 42 insertions(+), 3 deletions(-) diff --git a/src/api/capacity/services/capacity/assignments/admission/living-allocation-inputs.ts b/src/api/capacity/services/capacity/assignments/admission/living-allocation-inputs.ts index 313743fa..3916c7c6 100644 --- a/src/api/capacity/services/capacity/assignments/admission/living-allocation-inputs.ts +++ b/src/api/capacity/services/capacity/assignments/admission/living-allocation-inputs.ts @@ -11,7 +11,7 @@ export type LivingAllocationInputs = Record { const result: LivingAllocationInputs = {}; for (const provider of input.providers) { @@ -73,6 +73,9 @@ export async function livingAllocationInputs(store: CapacityGovernanceDatabase, AND assignment.assignment_attempt_json::jsonb->'provider'->>'modelConfigurationId'=? AND assignment.assignment_attempt_json::jsonb->'provider'->>'executionCapabilityId'=? AND assignment.assignment_attempt_json::jsonb->'effectiveProfile'->>'activity'=? + AND ${input.proposalGovernanceReview + ? "node.kind='reviewing' AND node.pair_role IS NULL AND node.source_ref_json::jsonb->>'model'='proposal'" + : "NOT (node.kind='reviewing' AND node.pair_role IS NULL AND node.source_ref_json::jsonb->>'model'='proposal')"} AND usage.accounting_mode='aggregate' AND ((assignment.status='completed' AND (node.pair_role IS DISTINCT FROM 'actor' OR (node.status='completed' AND assignment.execution_node_revision=node.node_revision))) diff --git a/src/api/capacity/services/capacity/assignments/planning/execution/assignment-attempt-builder.ts b/src/api/capacity/services/capacity/assignments/planning/execution/assignment-attempt-builder.ts index c530db9d..758827d9 100644 --- a/src/api/capacity/services/capacity/assignments/planning/execution/assignment-attempt-builder.ts +++ b/src/api/capacity/services/capacity/assignments/planning/execution/assignment-attempt-builder.ts @@ -233,6 +233,7 @@ export function buildAssignmentAttempt(input: { // as the node's actual viable minimum still fits. const allocation = calculateAssignmentAllocation({ estimate: allocationEstimate, measurements: planningTurn ? [] : allocationInputs.measurements, + observedViabilityFloor: isProposalGovernanceReview(candidate.node), constraints: [{ id: 'execution-window', remainingSeconds: availableSeconds }, { id: 'utc-day-window', remainingSeconds: Math.max(0, (utcDayEnd - Date.parse(input.now)) / 1000 - preparationSeconds) }, { id: 'model-day', remainingSeconds: remaining(limits.dailyActiveSecondsLimit, observation.modelUsage) }, diff --git a/src/api/capacity/services/capacity/assignments/planning/execution/living-execution-assignment.ts b/src/api/capacity/services/capacity/assignments/planning/execution/living-execution-assignment.ts index 98741a41..254b7f2b 100644 --- a/src/api/capacity/services/capacity/assignments/planning/execution/living-execution-assignment.ts +++ b/src/api/capacity/services/capacity/assignments/planning/execution/living-execution-assignment.ts @@ -203,7 +203,8 @@ export async function assignNextReadyExecutionNode( try { const allocationInputs = await livingAllocationInputs(store, { run, runs, providers: executionProviders, capacityProviderId: principal.capacityProviderId, capabilityId: candidate.node.requiredCapabilities?.[0] ?? '', - agentClass: candidate.node.agentClass!, activity: candidate.effectiveProfile.activity, now }); + agentClass: candidate.node.agentClass!, activity: candidate.effectiveProfile.activity, + proposalGovernanceReview: isProposalGovernanceReview(candidate.node), now }); const selected = buildAssignmentAttempt({ candidate, run, principal, providerSessionId, providers: executionProviders, allocationInputs, attempt: priorAttempts + 1, now, diff --git a/tests/unit/control-plane/capacity/execution/admission/allocation-inputs.test.ts b/tests/unit/control-plane/capacity/execution/admission/allocation-inputs.test.ts index 039f015a..0348917e 100644 --- a/tests/unit/control-plane/capacity/execution/admission/allocation-inputs.test.ts +++ b/tests/unit/control-plane/capacity/execution/admission/allocation-inputs.test.ts @@ -104,7 +104,8 @@ describe('live allocation ledger inputs', () => { it('does not learn a short success from an Actor attempt rejected by review', async () => { const db = new PGlite(); try { - await db.exec(`CREATE TABLE execution_nodes (id text, team_id text, agent_class text, pair_role text, status text, node_revision integer); + await db.exec(`CREATE TABLE execution_nodes (id text, team_id text, agent_class text, pair_role text, status text, node_revision integer, + kind text DEFAULT 'acting', source_ref_json jsonb DEFAULT '{}'); CREATE TABLE capacity_provider_assignments (id text, team_id text, execution_node_id text, execution_node_revision integer, capacity_provider_id text, execution_provider_id text, status text, lifecycle_code text, assignment_attempt_json jsonb); CREATE TABLE capacity_usage_actuals (id text, assignment_id text, created_at text, active_seconds integer, accounting_mode text); @@ -134,6 +135,39 @@ describe('live allocation ledger inputs', () => { expect(result['codex-implementation']?.measurements.map(({ id }) => id)).toEqual(['expired', 'accepted']); } finally { await db.close(); } }, 15_000); + it('uses only proposal-governance history for proposal reviews, never paired-review duration', async () => { + const db = new PGlite(); + try { + await db.exec(`CREATE TABLE execution_nodes (id text, team_id text, agent_class text, pair_role text, + kind text, source_ref_json jsonb, status text, node_revision integer); + CREATE TABLE capacity_provider_assignments (id text, team_id text, execution_node_id text, + capacity_provider_id text, execution_provider_id text, status text, lifecycle_code text, + assignment_attempt_json jsonb, execution_node_revision integer); + CREATE TABLE capacity_usage_actuals (id text, assignment_id text, created_at text, active_seconds integer, accounting_mode text); + INSERT INTO execution_nodes VALUES + ('governance','team','reviewer',NULL,'reviewing','{"model":"proposal"}','completed',1), + ('paired','team','reviewer','reviewer','reviewing','{"model":"proposal"}','completed',1); + INSERT INTO capacity_provider_assignments VALUES + ('governance','team','governance','provider','codex-implementation','completed',NULL, + '{"estimate":{"expectedSeconds":250},"limits":{"maximumSeconds":165},"provider":{"modelConfigurationId":"terra-medium","executionCapabilityId":"implementation"},"effectiveProfile":{"activity":"reviewing"}}',1), + ('paired','team','paired','provider','codex-implementation','completed',NULL, + '{"estimate":{"expectedSeconds":250},"limits":{"maximumSeconds":165},"provider":{"modelConfigurationId":"terra-medium","executionCapabilityId":"implementation"},"effectiveProfile":{"activity":"reviewing"}}',1); + INSERT INTO capacity_usage_actuals VALUES + ('governance','governance','2026-09-16T12:01:00Z',84,'aggregate'), + ('paired','paired','2026-09-16T12:02:00Z',15,'aggregate');`); + const store = { all: async (sql: string, values: unknown[]) => { + if (!sql.includes('capacity_usage_actuals')) return []; + let index = 0; + return (await db.query(sql.replace(/\?/gu, () => `$${++index}`), values)).rows; + }, first: async () => ({ ready_count: 1 }) }; + const base = { run: run as never, runs: [run as never], providers: [provider as never], + capacityProviderId: 'provider', capabilityId: 'implementation', agentClass: 'reviewer', activity: 'reviewing', now }; + expect((await livingAllocationInputs(store as never, { ...base, proposalGovernanceReview: true })) + ['codex-implementation']?.measurements.map(({ id }) => id)).toEqual(['governance']); + expect((await livingAllocationInputs(store as never, { ...base, proposalGovernanceReview: false })) + ['codex-implementation']?.measurements.map(({ id }) => id)).toEqual(['paired']); + } finally { await db.close(); } + }, 15_000); it('does not exempt closing workdays from shared supply and weighted allocation', async () => { const closingRun = { ...run, parameters: { appliedPlan: { ...plan, state: 'closing' } } }; const store = { all: vi.fn(async () => []), first: vi.fn(async () => ({ ready_count: 1 })) }; From 3e0ad3e3ea543768149b80b39215b4c28ea6b5b1 Mon Sep 17 00:00:00 2001 From: Adrian Webb Date: Mon, 28 Sep 2026 22:20:46 -0400 Subject: [PATCH 02/10] Prove governance viability is wired into assignment admission --- .../assignment-attempt-builder.test.ts | 25 +++++++++++++++++++ 1 file changed, 25 insertions(+) diff --git a/tests/unit/control-plane/capacity/execution/assignment-attempt-builder.test.ts b/tests/unit/control-plane/capacity/execution/assignment-attempt-builder.test.ts index 70c1ac29..32aeabe6 100644 --- a/tests/unit/control-plane/capacity/execution/assignment-attempt-builder.test.ts +++ b/tests/unit/control-plane/capacity/execution/assignment-attempt-builder.test.ts @@ -47,6 +47,31 @@ const run = { id: 'workday', executionMode: 'simulation', parameters: { appliedP } } } as never; describe('immutable assignment-attempt construction', () => { + it('uses observed proposal-review viability without shortening an ordinary paired review', () => { + const review = structuredClone(candidate); + review.node.kind = 'reviewing'; + review.node.pairRole = null as never; + review.node.agentClass = 'reviewer'; + review.node.workspace = 'read-only'; + review.node.estimate = { minimumSeconds: 60, expectedSeconds: 100, maximumSeconds: 165 }; + review.node.requestedPermissions = { content: { read: ['proposal'], write: [] }, + tools: ['source.read', 'verification'] } as never; + review.effectiveProfile.permissionCeiling = review.node.requestedPermissions; + review.effectiveProfile.activity = 'reviewing'; + review.effectiveProfile.handler = 'reviewer'; + const measurements = Array.from({ length: 20 }, (_, index) => ({ id: `review-${index}`, + completedAt: new Date(Date.parse('2026-09-13T11:30:00.000Z') + index * 1000).toISOString(), + expectedSeconds: 250, allocatedSeconds: 200, activeSeconds: 84, outcome: 'completed' as const })); + const build = (item: typeof review) => buildAssignmentAttempt({ candidate: item as never, run, + principal: { teamId: 'team', capacityProviderId: 'provider' } as never, + allocationInputs: { codex: { measurements, constraints: [] } }, providerSessionId: 'session', + providers: [provider] as never, attempt: 1, now: '2026-09-13T12:01:00.000Z' }); + expect(build(review).allocation).toMatchObject({ admitted: true, minimumSeconds: 105, + observedMinimumSeconds: 105, allocatedSeconds: 105 }); + review.node.pairRole = 'reviewer' as never; + expect(build(review).allocation).toMatchObject({ admitted: true, minimumSeconds: 60, + observedMinimumSeconds: 0, allocatedSeconds: 60 }); + }); it('gives planning discussion the policy turn ceiling without changing acting chat allocation', () => { const discussion = structuredClone(candidate); discussion.node.kind = 'communication'; From e5bc86204e667e26f070cfd34f322462ca1fc67c Mon Sep 17 00:00:00 2001 From: Adrian Webb Date: Mon, 28 Sep 2026 22:23:36 -0400 Subject: [PATCH 03/10] Keep allocation regression within test file limits --- .../assignment-attempt-builder.test.ts | 29 ++++--------------- 1 file changed, 6 insertions(+), 23 deletions(-) diff --git a/tests/unit/control-plane/capacity/execution/assignment-attempt-builder.test.ts b/tests/unit/control-plane/capacity/execution/assignment-attempt-builder.test.ts index 32aeabe6..a8b011b3 100644 --- a/tests/unit/control-plane/capacity/execution/assignment-attempt-builder.test.ts +++ b/tests/unit/control-plane/capacity/execution/assignment-attempt-builder.test.ts @@ -48,29 +48,12 @@ const run = { id: 'workday', executionMode: 'simulation', parameters: { appliedP describe('immutable assignment-attempt construction', () => { it('uses observed proposal-review viability without shortening an ordinary paired review', () => { - const review = structuredClone(candidate); - review.node.kind = 'reviewing'; - review.node.pairRole = null as never; - review.node.agentClass = 'reviewer'; - review.node.workspace = 'read-only'; - review.node.estimate = { minimumSeconds: 60, expectedSeconds: 100, maximumSeconds: 165 }; - review.node.requestedPermissions = { content: { read: ['proposal'], write: [] }, - tools: ['source.read', 'verification'] } as never; - review.effectiveProfile.permissionCeiling = review.node.requestedPermissions; - review.effectiveProfile.activity = 'reviewing'; - review.effectiveProfile.handler = 'reviewer'; - const measurements = Array.from({ length: 20 }, (_, index) => ({ id: `review-${index}`, - completedAt: new Date(Date.parse('2026-09-13T11:30:00.000Z') + index * 1000).toISOString(), - expectedSeconds: 250, allocatedSeconds: 200, activeSeconds: 84, outcome: 'completed' as const })); - const build = (item: typeof review) => buildAssignmentAttempt({ candidate: item as never, run, - principal: { teamId: 'team', capacityProviderId: 'provider' } as never, - allocationInputs: { codex: { measurements, constraints: [] } }, providerSessionId: 'session', - providers: [provider] as never, attempt: 1, now: '2026-09-13T12:01:00.000Z' }); - expect(build(review).allocation).toMatchObject({ admitted: true, minimumSeconds: 105, - observedMinimumSeconds: 105, allocatedSeconds: 105 }); - review.node.pairRole = 'reviewer' as never; - expect(build(review).allocation).toMatchObject({ admitted: true, minimumSeconds: 60, - observedMinimumSeconds: 0, allocatedSeconds: 60 }); + const review = structuredClone(candidate); Object.assign(review.node, { kind: 'reviewing', pairRole: null, agentClass: 'reviewer', workspace: 'read-only', estimate: { minimumSeconds: 60, expectedSeconds: 100, maximumSeconds: 165 }, requestedPermissions: { content: { read: ['proposal'], write: [] }, tools: ['source.read', 'verification'] } }); + Object.assign(review.effectiveProfile, { activity: 'reviewing', handler: 'reviewer', permissionCeiling: review.node.requestedPermissions }); + const measurements = Array.from({ length: 20 }, (_, index) => ({ id: `review-${index}`, completedAt: new Date(Date.parse('2026-09-13T11:30:00.000Z') + index * 1000).toISOString(), expectedSeconds: 250, allocatedSeconds: 200, activeSeconds: 84, outcome: 'completed' as const })); + const build = (item: typeof review) => buildAssignmentAttempt({ candidate: item as never, run, principal: { teamId: 'team', capacityProviderId: 'provider' } as never, allocationInputs: { codex: { measurements, constraints: [] } }, providerSessionId: 'session', providers: [provider] as never, attempt: 1, now: '2026-09-13T12:01:00.000Z' }); + expect(build(review).allocation).toMatchObject({ admitted: true, minimumSeconds: 105, observedMinimumSeconds: 105, allocatedSeconds: 105 }); review.node.pairRole = 'reviewer' as never; + expect(build(review).allocation).toMatchObject({ admitted: true, minimumSeconds: 60, observedMinimumSeconds: 0, allocatedSeconds: 60 }); }); it('gives planning discussion the policy turn ceiling without changing acting chat allocation', () => { const discussion = structuredClone(candidate); From 3181bc617162036c11ad27586c96cd373a0eb9ee Mon Sep 17 00:00:00 2001 From: Adrian Webb Date: Mon, 28 Sep 2026 23:03:28 -0400 Subject: [PATCH 04/10] Keep chat reconciliation independent of unrelated proposals --- .../workday-admission-fence-service.ts | 5 ++-- .../scheduling/workday-scheduling-service.ts | 5 ++-- .../scheduling/workday-tick-service.ts | 6 ++-- .../execution/execution-graph-service.ts | 28 +++++++++++++------ .../communication-execution-projector.test.ts | 26 ++++++++++++++++- 5 files changed, 55 insertions(+), 15 deletions(-) diff --git a/src/api/capacity/services/capacity/workdays/lifecycle/workday-admission-fence-service.ts b/src/api/capacity/services/capacity/workdays/lifecycle/workday-admission-fence-service.ts index b6cc492b..b2ff7c7e 100644 --- a/src/api/capacity/services/capacity/workdays/lifecycle/workday-admission-fence-service.ts +++ b/src/api/capacity/services/capacity/workdays/lifecycle/workday-admission-fence-service.ts @@ -3,7 +3,7 @@ import { CapacityGovernanceError } from '../../../../database.ts'; import { assignmentContentIntegrationReadySql,CONTENT_INTEGRATED_EVENT,CONTENT_INTEGRATION_REQUIRED_EVENT } from '../../assignments/lifecycle/assignment-content-integration-requirement.ts'; import { CapacityWorkdayRunRepository } from '../../../../repositories/capacity/workdays/workday-run.ts'; import { advanceLivingWorkday } from './living-workday-lifecycle.ts'; -import { reconcileExecutionGraph } from '../../../../../control-plane/repositories/capacity/execution/execution-graph-service.ts'; +import { reconcileCommunicationExecutionGraph, reconcileExecutionGraph } from '../../../../../control-plane/repositories/capacity/execution/execution-graph-service.ts'; type WorkdayAdmissionFenceStore = Parameters[0]; @@ -26,7 +26,8 @@ export async function fenceCapacityWorkdayAdmission( 'capacity_workday_plan_missing', 'Running workday has no authoritative applied plan to close.', 409, { runId }); if (run.status === 'running' && run.parameters.appliedPlan) { await advanceLivingWorkday(store, run, new Date().toISOString(), true); - await reconcileExecutionGraph(store, teamId); + if (run.executionKind === 'conversation') await reconcileCommunicationExecutionGraph(store, teamId); + else await reconcileExecutionGraph(store, teamId); } const assignment = await store.first( `SELECT COUNT(*) AS total, diff --git a/src/api/capacity/services/capacity/workdays/scheduling/workday-scheduling-service.ts b/src/api/capacity/services/capacity/workdays/scheduling/workday-scheduling-service.ts index c2676d82..aa01db15 100644 --- a/src/api/capacity/services/capacity/workdays/scheduling/workday-scheduling-service.ts +++ b/src/api/capacity/services/capacity/workdays/scheduling/workday-scheduling-service.ts @@ -12,7 +12,7 @@ type WorkdayProject, } from '../policy/workday-project-policy.ts'; import { resolveWorkdayAgentProfileSnapshot } from '../policy/workday-agent-profile-policy.ts'; import { reconcileTreeDxRefSignals } from '../../../treedx/repositories/treedx-ref-signal-reconciler.ts'; -import { reconcileExecutionGraph } from '../../../../../control-plane/repositories/capacity/execution/execution-graph-service.ts'; +import { reconcileCommunicationExecutionGraph, reconcileExecutionGraph } from '../../../../../control-plane/repositories/capacity/execution/execution-graph-service.ts'; import { workdayParticipants } from '../../../../policy/execution/workday-participants.ts'; import { readExactProposal } from '../../../../../governance/executable-proposal.ts'; @@ -253,7 +253,8 @@ export async function scheduleCapacityWorkdayRun( if (!updated) { throw new CapacityGovernanceError('capacity_workday_run_update_failed', 'Scheduled workday run could not be updated.', 500, { runId: run.id }); } - await reconcileExecutionGraph(store, run.teamId); + if (run.executionKind === 'conversation') await reconcileCommunicationExecutionGraph(store, run.teamId); + else await reconcileExecutionGraph(store, run.teamId); await recordRequiredEvent(store, run.teamId, run.id, { eventType: 'assignment.polling_ready', status: 'recorded', title: 'Workday is ready for authenticated provider polling', diff --git a/src/api/capacity/services/capacity/workdays/scheduling/workday-tick-service.ts b/src/api/capacity/services/capacity/workdays/scheduling/workday-tick-service.ts index db2a9744..5491a9ce 100644 --- a/src/api/capacity/services/capacity/workdays/scheduling/workday-tick-service.ts +++ b/src/api/capacity/services/capacity/workdays/scheduling/workday-tick-service.ts @@ -2,7 +2,7 @@ import { createHash } from 'node:crypto'; import { CapacityGovernanceError } from '../../../../database.ts'; import { decodeDurableJsonObject } from '../../../../durable-json.ts'; import { CapacityWorkdayRunRepository } from '../../../../repositories/capacity/workdays/workday-run.ts'; -import { reconcileExecutionGraph } from '../../../../../control-plane/repositories/capacity/execution/execution-graph-service.ts'; +import { reconcileCommunicationExecutionGraph, reconcileExecutionGraph } from '../../../../../control-plane/repositories/capacity/execution/execution-graph-service.ts'; import { CapacityWorkdayEventService } from '../content/workday-event-service.ts'; import { advanceLivingWorkday } from '../lifecycle/living-workday-lifecycle.ts'; import { appliedWorkdaySchema, workdayPhase } from '@treeseed/sdk/agent-capacity'; @@ -52,7 +52,9 @@ export async function tickCapacityWorkdayRun( WHERE team_id=? AND workday_id=? AND kind IN ('planning','estimating') AND status IN ('ready','blocked')`, [now, teamId, runId]); } const lifecycle = await advanceLivingWorkday(store, run, now); - const executionGraph = await reconcileExecutionGraph(store, teamId, {}, `workday-tick:${runId}:${eventId ?? now}`); + const executionGraph = run.executionKind === 'conversation' + ? await reconcileCommunicationExecutionGraph(store, teamId) + : await reconcileExecutionGraph(store, teamId, {}, `workday-tick:${runId}:${eventId ?? now}`); const result = { runId, tickedAt: now, lifecycle, executionGraph }; if (eventId) await new CapacityWorkdayEventService(store).create(teamId, runId, { id: eventId, eventType: 'workday.tick', status: 'recorded', title: 'Workday execution-graph tick', diff --git a/src/api/control-plane/repositories/capacity/execution/execution-graph-service.ts b/src/api/control-plane/repositories/capacity/execution/execution-graph-service.ts index 828c28fb..c58c62e4 100644 --- a/src/api/control-plane/repositories/capacity/execution/execution-graph-service.ts +++ b/src/api/control-plane/repositories/capacity/execution/execution-graph-service.ts @@ -277,12 +277,12 @@ export async function persistExecutionGraph(store: any, graph: TeamGraph, curren return revisionRecord; } -async function reconcileExecutionGraphOnce(store: any, teamId: string, body: Row = {}) { +async function reconcileExecutionGraphOnce(store: any, teamId: string, body: Row = {}, scope: 'team' | 'communication' = 'team') { const current = await loadGraphSource('projection', () => readGraph(store, teamId)); const [sources, profiles, workdays, communications, activeAssignmentRows, terminalAssignmentRows, reviewCycleRows] = await Promise.all([ - loadGraphSource('governance', () => loadTeamExecutableProposalSources(store, teamId)), + scope === 'communication' ? Promise.resolve([]) : loadGraphSource('governance', () => loadTeamExecutableProposalSources(store, teamId)), loadGraphSource('agent_profiles', () => loadProfiles(store, teamId)), - loadGraphSource('workdays', () => loadActiveWorkdays(store, teamId)), + scope === 'communication' ? Promise.resolve([]) : loadGraphSource('workdays', () => loadActiveWorkdays(store, teamId)), loadGraphSource('communications', () => loadCommunicationInvocations(store, teamId)), loadGraphSource('assignments', () => store.all(`SELECT DISTINCT execution_node_id,work_day_id FROM capacity_provider_assignments WHERE team_id=? AND execution_node_id IS NOT NULL AND status IN ('pending','leased','running','returned')`, [teamId])), @@ -316,11 +316,14 @@ async function reconcileExecutionGraphOnce(store: any, teamId: string, body: Row const workdayProjection = projectActiveWorkdays({ teamId, revision, sources: workdays, profiles, decisionNodes: proposalProjection?.nodes ?? [] }); const communicationProjection = projectCommunicationInvocations({ teamId, revision, sources: communications, profiles }); - const nodes = [...(proposalProjection?.nodes ?? []), ...workdayProjection.nodes, ...communicationProjection.nodes]; - const edges = [...(proposalProjection?.edges ?? []), ...workdayProjection.edges]; + const nodes = scope === 'communication' + ? [...current.nodes.filter((node) => node.kind !== 'communication'), ...communicationProjection.nodes] + : [...(proposalProjection?.nodes ?? []), ...workdayProjection.nodes, ...communicationProjection.nodes]; + const edges = scope === 'communication' ? current.edges : [...(proposalProjection?.edges ?? []), ...workdayProjection.edges]; const changedSourceRefs = [...new Map([...(proposalProjection?.revision.changedSourceRefs ?? []), ...workdayProjection.changedSourceRefs, ...communicationProjection.changedSourceRefs, - ...(!sources.length && !workdays.length && !communications.length ? current.nodes.map((node) => node.sourceRef) : [])] + ...(scope === 'team' && !sources.length && !workdays.length && !communications.length + ? current.nodes.map((node) => node.sourceRef) : [])] .map((reference) => [stable(reference), reference])).values()]; const base: TeamGraph = { teamId, revision, digest: digest({ teamId, nodes, edges }), nodes, edges }; const nodeById = new Map(base.nodes.map((node) => [node.id, node])); @@ -423,10 +426,10 @@ async function reconcileExecutionGraphOnce(store: any, teamId: string, body: Row } /** Concurrent source changes converge by rereading the winning graph revision. */ -export async function reconcileExecutionGraph(store: any, teamId: string, body: Row = {}, ..._trace: unknown[]) { +async function reconcileGraphScope(store: any, teamId: string, body: Row, scope: 'team' | 'communication') { for (let attempt = 1; attempt <= 4; attempt += 1) { try { - return await reconcileExecutionGraphOnce(store, teamId, body); + return await reconcileExecutionGraphOnce(store, teamId, body, scope); } catch (error) { if (!(error instanceof CapacityOperationError) || error.code !== 'execution_graph_revision_conflict' || attempt === 4) throw error; } @@ -434,6 +437,15 @@ export async function reconcileExecutionGraph(store: any, teamId: string, body: throw new CapacityOperationError(409, 'execution_graph_revision_conflict', 'The execution graph changed concurrently; reconcile again.'); } +export async function reconcileExecutionGraph(store: any, teamId: string, body: Row = {}, ..._trace: unknown[]) { + return reconcileGraphScope(store, teamId, body, 'team'); +} + +/** Reconcile conversation demand in the same graph without reinterpreting unrelated accepted proposals. */ +export async function reconcileCommunicationExecutionGraph(store: any, teamId: string, body: Row = {}) { + return reconcileGraphScope(store, teamId, body, 'communication'); +} + export function createExecutionGraphService(store: any) { return { async show(principal: CapacityPrincipal, teamId: string, query: Row) { diff --git a/tests/unit/control-plane/capacity/execution/graph/communication-execution-projector.test.ts b/tests/unit/control-plane/capacity/execution/graph/communication-execution-projector.test.ts index 2c155bc1..51b28b43 100644 --- a/tests/unit/control-plane/capacity/execution/graph/communication-execution-projector.test.ts +++ b/tests/unit/control-plane/capacity/execution/graph/communication-execution-projector.test.ts @@ -1,6 +1,7 @@ -import { describe, expect, it } from 'vitest'; +import { describe, expect, it, vi } from 'vitest'; import { calculateAssignmentAllocation } from '@treeseed/sdk/agent-capacity'; import { projectCommunicationInvocations } from '../../../../../../src/api/capacity/policy/execution/communication-execution-projector.ts'; +import { reconcileCommunicationExecutionGraph } from '../../../../../../src/api/control-plane/repositories/capacity/execution/execution-graph-service.ts'; const definition = { schemaVersion: 'treeseed.agent/v1' as const, id: 'sdk/architect', name: 'SDK Architect', agentClass: 'architect', @@ -12,6 +13,29 @@ const definition = { }; describe('communication living-graph projection', () => { + it('projects chat in the one team graph without validating unrelated accepted proposal content', async () => { + const acceptedRef = { store: 'treedx', model: 'proposal', id: 'legacy-accepted', revision: 1, + digest: `sha256:${'b'.repeat(64)}`, repository: 'treeseed-ai/sdk-library', commit: 'a'.repeat(40), path: 'proposals/legacy.mdx' }; + const all = vi.fn(async (sql: string) => { + if (sql.includes('governance_proposals')) throw new Error('Unrelated accepted proposal is not executable'); + if (sql.includes('FROM execution_nodes')) return [{ id: 'accepted-condition', team_id: 'team', project_id: 'sdk', + kind: 'condition', pair_role: null, source_ref_json: acceptedRef, authority_refs_json: [], rule_revision: 1, + node_revision: 1, status: 'blocked', condition_json: { conditionType: 'authority', subjectRef: acceptedRef, + expectedState: 'accepted' }, graph_revision_created: 1, graph_revision_updated: 1 }]; + if (sql.includes('FROM project_agent_classes')) return [{ project_id: 'sdk', handler_refs_json: { agents: [definition] } }]; + if (sql.includes('FROM agent_invocation_requests')) return [{ id: 'invocation', team_id: 'team', project_id: 'sdk', + agent_id: 'architect', execution_id: 'conversation-invocation', repository_id: 'treeseed-ai/sdk-library', + metadata_json: { sourceMessagePath: 'discussion-messages/smoke/request.mdx', sourceCommit: 'a'.repeat(40), + productiveSeconds: 180 }, content_refs_json: [] }]; + return []; + }); + const store = { all, first: vi.fn(async () => ({ revision: 1, graph_digest: `sha256:${'c'.repeat(64)}` })) }; + const planned = await reconcileCommunicationExecutionGraph(store, 'team', { plan: true }); + expect(planned).toMatchObject({ baseRevision: 1, changes: { added: ['communication:invocation:conversation-invocation'], + stale: [], completed: [] } }); + expect(all.mock.calls.some(([sql]) => sql.includes('governance_proposals'))).toBe(false); + expect(all.mock.calls.some(([sql]) => sql.includes('FROM capacity_workday_runs') && !sql.includes('agent_invocation_requests'))).toBe(false); + }); it('projects an addressed message as one ready read-only chat assignment source', () => { const projected = projectCommunicationInvocations({ teamId: 'team', revision: 3, profiles: { 'sdk:architect': definition }, sources: [{ From 9def590f500fa36a86e0ded57dc382d3821e1935 Mon Sep 17 00:00:00 2001 From: Adrian Webb Date: Mon, 28 Sep 2026 23:24:55 -0400 Subject: [PATCH 05/10] Freeze non-runnable historical proposal components --- .../execution/executable-proposal-source.ts | 16 ++++++-- .../execution/execution-graph-service.ts | 28 ++++++------- .../graph/executable-proposal-source.test.ts | 39 +++++++++++++++++++ 3 files changed, 66 insertions(+), 17 deletions(-) diff --git a/src/api/capacity/services/capacity/execution/executable-proposal-source.ts b/src/api/capacity/services/capacity/execution/executable-proposal-source.ts index d3f995c0..374f1245 100644 --- a/src/api/capacity/services/capacity/execution/executable-proposal-source.ts +++ b/src/api/capacity/services/capacity/execution/executable-proposal-source.ts @@ -19,7 +19,8 @@ const stable = (value: unknown): string => { }; /** Load structurally complete drafts for review and accepted proposals for work. */ -export async function loadTeamExecutableProposalSources(store: any, teamId: string, projectId?: string): Promise { +export async function loadTeamExecutableProposalSources(store: any, teamId: string, projectId?: string, + onFrozenInvalid?: (source: { id: string; digest: string }) => void): Promise { const [rows, graphRows] = await Promise.all([store.all(`SELECT p.id AS proposal_id,p.project_id,p.active_version,p.active_content_hash,p.metadata_json, p.decision_id,d.id AS accepted_decision_id,d.proposal_version,d.proposal_content_hash,d.decision_record_json @@ -30,14 +31,15 @@ export async function loadTeamExecutableProposalSources(store: any, teamId: stri ${projectId ? 'AND p.project_id = ?' : ''} ORDER BY p.project_id,p.id`, projectId ? [teamId, projectId] : [teamId]), store.all('SELECT source_ref_json,status FROM execution_nodes WHERE team_id=?', [teamId])]); - const graphState = new Map(); + const graphState = new Map(); for (const graphRow of graphRows) { const source = record(graphRow.source_ref_json); if (text(source.model) !== 'proposal') continue; const key = `${text(source.id)}\u0000${text(source.digest)}`; - const state = graphState.get(key) ?? { count: 0, incomplete: 0 }; + const state = graphState.get(key) ?? { count: 0, incomplete: 0, ready: 0 }; state.count += 1; if (text(graphRow.status) !== 'completed') state.incomplete += 1; + if (text(graphRow.status) === 'ready') state.ready += 1; graphState.set(key, state); } const sources: ExecutableProposalSource[] = []; @@ -62,6 +64,14 @@ export async function loadTeamExecutableProposalSources(store: any, teamId: stri // would stale its existing graph nodes without a new decision. if (!accepted) continue; const value = error as { status?: number; code?: string }; + // An already-materialized historical component whose content no longer + // satisfies the current schema may remain visible, but cannot admit work. + // Preserve its exact graph nodes; never silently drop an unmaterialized or + // ready accepted decision, or mask a TreeDX outage/authority mismatch. + if (value.code === 'proposal_execution_plan_invalid' && state?.count && state.ready === 0 && onFrozenInvalid) { + onFrozenInvalid({ id: text(row.proposal_id), digest: `sha256:${text(row.active_content_hash)}` }); + continue; + } throw Object.assign(new CapacityOperationError(Number(value.status ?? 409), value.code ?? 'proposal_execution_plan_invalid', error instanceof Error ? error.message : 'The accepted proposal could not be read.'), { diagnostics: (error as { diagnostics?: unknown }).diagnostics }); } diff --git a/src/api/control-plane/repositories/capacity/execution/execution-graph-service.ts b/src/api/control-plane/repositories/capacity/execution/execution-graph-service.ts index c58c62e4..0c34520e 100644 --- a/src/api/control-plane/repositories/capacity/execution/execution-graph-service.ts +++ b/src/api/control-plane/repositories/capacity/execution/execution-graph-service.ts @@ -276,11 +276,12 @@ export async function persistExecutionGraph(store: any, graph: TeamGraph, curren } return revisionRecord; } - async function reconcileExecutionGraphOnce(store: any, teamId: string, body: Row = {}, scope: 'team' | 'communication' = 'team') { const current = await loadGraphSource('projection', () => readGraph(store, teamId)); + const frozenInvalid = new Set(); const [sources, profiles, workdays, communications, activeAssignmentRows, terminalAssignmentRows, reviewCycleRows] = await Promise.all([ - scope === 'communication' ? Promise.resolve([]) : loadGraphSource('governance', () => loadTeamExecutableProposalSources(store, teamId)), + scope === 'communication' ? Promise.resolve([]) : loadGraphSource('governance', () => + loadTeamExecutableProposalSources(store, teamId, undefined, (source) => frozenInvalid.add(`${source.id}\u0000${source.digest}`))), loadGraphSource('agent_profiles', () => loadProfiles(store, teamId)), scope === 'communication' ? Promise.resolve([]) : loadGraphSource('workdays', () => loadActiveWorkdays(store, teamId)), loadGraphSource('communications', () => loadCommunicationInvocations(store, teamId)), @@ -316,10 +317,13 @@ async function reconcileExecutionGraphOnce(store: any, teamId: string, body: Row const workdayProjection = projectActiveWorkdays({ teamId, revision, sources: workdays, profiles, decisionNodes: proposalProjection?.nodes ?? [] }); const communicationProjection = projectCommunicationInvocations({ teamId, revision, sources: communications, profiles }); + const frozenNodes = current.nodes.filter((node) => frozenInvalid.has(`${node.sourceRef.id}\u0000${node.sourceRef.digest}`)); + const frozenNodeIds = new Set(frozenNodes.map((node) => node.id)); const nodes = scope === 'communication' ? [...current.nodes.filter((node) => node.kind !== 'communication'), ...communicationProjection.nodes] - : [...(proposalProjection?.nodes ?? []), ...workdayProjection.nodes, ...communicationProjection.nodes]; - const edges = scope === 'communication' ? current.edges : [...(proposalProjection?.edges ?? []), ...workdayProjection.edges]; + : [...(proposalProjection?.nodes ?? []), ...workdayProjection.nodes, ...communicationProjection.nodes, ...frozenNodes]; + const edges = scope === 'communication' ? current.edges : [...(proposalProjection?.edges ?? []), ...workdayProjection.edges, ...current.edges.filter( + (edge) => frozenNodeIds.has(edge.fromNodeId) || frozenNodeIds.has(edge.toNodeId))]; const changedSourceRefs = [...new Map([...(proposalProjection?.revision.changedSourceRefs ?? []), ...workdayProjection.changedSourceRefs, ...communicationProjection.changedSourceRefs, ...(scope === 'team' && !sources.length && !workdays.length && !communications.length @@ -347,10 +351,8 @@ async function reconcileExecutionGraphOnce(store: any, teamId: string, body: Row const disposition = text(record(record(row.lifecycle_output_json).activityCompletion).reviewDisposition); const status = text(row.status); const nodeId = text(row.execution_node_id); - // A failed Actor revision is terminal graph evidence. Dropping it here lets - // the projection restore the Actor to blocked/ready and can cause its paired - // Reviewer to re-review an older successful candidate. Preserve the exact - // failed/cancelled revision so the review edge remains fail-closed. + // Preserve failed Actor revisions or a Reviewer could re-review an old candidate. + // Failed/cancelled revisions remain terminal, never restored to ready. if (status !== 'completed') return [[nodeId, { status: status === 'cancelled' ? 'cancelled' as const : 'failed' as const, nodeRevision: integer(row.execution_node_revision), @@ -409,6 +411,8 @@ async function reconcileExecutionGraphOnce(store: any, teamId: string, body: Row }).map((node) => node.id)); recoverInterruptedGovernanceReviews(desired, eligible, revision); } + for (const node of desired.nodes) if (frozenNodeIds.has(node.id) && node.status === 'ready') node.status = 'blocked'; + if (frozenNodeIds.size) desired.digest = digest({ teamId, nodes: desired.nodes, edges: desired.edges }); const changes = graphChanges(current, desired); if (body.plan === true) return { teamId, baseRevision: current.revision, desiredRevision: desired.revision, desiredDigest: desired.digest, changes }; if (!hasChanges(changes)) return current.revision @@ -437,14 +441,10 @@ async function reconcileGraphScope(store: any, teamId: string, body: Row, scope: throw new CapacityOperationError(409, 'execution_graph_revision_conflict', 'The execution graph changed concurrently; reconcile again.'); } -export async function reconcileExecutionGraph(store: any, teamId: string, body: Row = {}, ..._trace: unknown[]) { - return reconcileGraphScope(store, teamId, body, 'team'); -} +export async function reconcileExecutionGraph(store: any, teamId: string, body: Row = {}, ..._trace: unknown[]) { return reconcileGraphScope(store, teamId, body, 'team'); } /** Reconcile conversation demand in the same graph without reinterpreting unrelated accepted proposals. */ -export async function reconcileCommunicationExecutionGraph(store: any, teamId: string, body: Row = {}) { - return reconcileGraphScope(store, teamId, body, 'communication'); -} +export async function reconcileCommunicationExecutionGraph(store: any, teamId: string, body: Row = {}) { return reconcileGraphScope(store, teamId, body, 'communication'); } export function createExecutionGraphService(store: any) { return { diff --git a/tests/unit/control-plane/capacity/execution/graph/executable-proposal-source.test.ts b/tests/unit/control-plane/capacity/execution/graph/executable-proposal-source.test.ts index bcfe9b91..2e8525c7 100644 --- a/tests/unit/control-plane/capacity/execution/graph/executable-proposal-source.test.ts +++ b/tests/unit/control-plane/capacity/execution/graph/executable-proposal-source.test.ts @@ -8,6 +8,7 @@ vi.mock('../../../../../../src/api/governance/executable-proposal.ts', async (im import { loadTeamExecutableProposalSources } from '../../../../../../src/api/capacity/services/capacity/execution/executable-proposal-source.ts'; import { loadProposalBlockingFeedback } from '../../../../../../src/api/capacity/services/capacity/execution/proposal-planning-source.ts'; +import { reconcileExecutionGraph } from '../../../../../../src/api/control-plane/repositories/capacity/execution/execution-graph-service.ts'; const questionRef = { store: 'treedx', model: 'question', id: 'question', revision: 1, digest: `sha256:${'c'.repeat(64)}`, repository: 'repository', commit: 'b'.repeat(40), path: 'questions/question.mdx' }; @@ -123,4 +124,42 @@ describe('executable proposal source selection', () => { code: 'proposal_execution_plan_invalid', diagnostics: [{ path: 'executionPlan' }], }); }); + it('freezes an already-materialized invalid accepted component with no ready work', async () => { + const all = vi.fn(async (query: string) => query.includes('FROM execution_nodes') ? [{ + source_ref_json: { model: 'proposal', id: 'historical', digest: `sha256:${'a'.repeat(64)}` }, + status: 'blocked', + }] : [{ + proposal_id: 'historical', project_id: 'project', active_version: 1, + active_content_hash: 'a'.repeat(64), decision_id: 'decision', accepted_decision_id: 'decision', + decision_record_json: { proposalRef: { id: 'historical' } }, + }]); + exactProposal.mockRejectedValueOnce(Object.assign(new Error('Legacy proposal schema is invalid.'), { + status: 422, code: 'proposal_execution_plan_invalid', + })); + const frozen: Array<{ id: string; digest: string }> = []; + await expect(loadTeamExecutableProposalSources({ all }, 'team', undefined, (source) => frozen.push(source))) + .resolves.toEqual([]); + expect(frozen).toEqual([{ id: 'historical', digest: `sha256:${'a'.repeat(64)}` }]); + }); + it('retains a blocked historical component during ordinary team reconciliation', async () => { + const source = { store: 'treedx', model: 'proposal', id: 'historical', revision: 1, + digest: `sha256:${'a'.repeat(64)}`, repository: 'library', commit: 'b'.repeat(40), path: 'proposals/old.mdx' }; + const all = vi.fn(async (query: string) => { + if (query.includes('FROM governance_proposals')) return [{ proposal_id: 'historical', project_id: 'project', + active_version: 1, active_content_hash: 'a'.repeat(64), accepted_decision_id: 'decision', + decision_record_json: { proposalRef: { id: 'historical' } } }]; + if (query.includes('FROM execution_nodes')) return [{ id: 'old-condition', team_id: 'team', project_id: 'project', + kind: 'condition', pair_role: null, source_ref_json: source, authority_refs_json: [], rule_revision: 1, + node_revision: 1, status: 'blocked', condition_json: { conditionType: 'authority', subjectRef: source, + expectedState: 'accepted' }, graph_revision_created: 1, graph_revision_updated: 1 }]; + return []; + }); + exactProposal.mockRejectedValueOnce(Object.assign(new Error('Legacy proposal schema is invalid.'), { + status: 422, code: 'proposal_execution_plan_invalid', + })); + const store = { all, first: vi.fn(async () => ({ revision: 1, graph_digest: `sha256:${'b'.repeat(64)}` })) }; + const planned = await reconcileExecutionGraph(store, 'team', { plan: true }); + expect(planned).toMatchObject({ changes: { stale: [], added: [] } }); + expect(all).toHaveBeenCalledWith(expect.stringContaining('FROM governance_proposals'), ['team']); + }); }); From 238ab71d3b9f1ce17fb66e84b84f141bee427ea3 Mon Sep 17 00:00:00 2001 From: Adrian Webb Date: Mon, 28 Sep 2026 23:32:28 -0400 Subject: [PATCH 06/10] Retire ready demand from cancelled historical simulations --- .../execution/executable-proposal-source.ts | 7 +++-- .../graph/executable-proposal-source.test.ts | 29 +++++++++++++++++++ 2 files changed, 34 insertions(+), 2 deletions(-) diff --git a/src/api/capacity/services/capacity/execution/executable-proposal-source.ts b/src/api/capacity/services/capacity/execution/executable-proposal-source.ts index 374f1245..365f8712 100644 --- a/src/api/capacity/services/capacity/execution/executable-proposal-source.ts +++ b/src/api/capacity/services/capacity/execution/executable-proposal-source.ts @@ -30,7 +30,9 @@ export async function loadTeamExecutableProposalSources(store: any, teamId: stri WHERE p.team_id = ? AND ((p.decision_id IS NULL AND p.status IN ('draft','submitted','open','voting')) OR d.id IS NOT NULL) ${projectId ? 'AND p.project_id = ?' : ''} ORDER BY p.project_id,p.id`, projectId ? [teamId, projectId] : [teamId]), - store.all('SELECT source_ref_json,status FROM execution_nodes WHERE team_id=?', [teamId])]); + store.all(`SELECT node.source_ref_json,node.status,node.workday_id,run.status AS workday_status + FROM execution_nodes node LEFT JOIN capacity_workday_runs run + ON run.team_id=node.team_id AND run.id=node.workday_id WHERE node.team_id=?`, [teamId])]); const graphState = new Map(); for (const graphRow of graphRows) { const source = record(graphRow.source_ref_json); @@ -39,7 +41,8 @@ export async function loadTeamExecutableProposalSources(store: any, teamId: stri const state = graphState.get(key) ?? { count: 0, incomplete: 0, ready: 0 }; state.count += 1; if (text(graphRow.status) !== 'completed') state.incomplete += 1; - if (text(graphRow.status) === 'ready') state.ready += 1; + if (text(graphRow.status) === 'ready' && !(text(graphRow.workday_id) + && ['cancelled', 'completed', 'failed'].includes(text(graphRow.workday_status)))) state.ready += 1; graphState.set(key, state); } const sources: ExecutableProposalSource[] = []; diff --git a/tests/unit/control-plane/capacity/execution/graph/executable-proposal-source.test.ts b/tests/unit/control-plane/capacity/execution/graph/executable-proposal-source.test.ts index 2e8525c7..0d715f71 100644 --- a/tests/unit/control-plane/capacity/execution/graph/executable-proposal-source.test.ts +++ b/tests/unit/control-plane/capacity/execution/graph/executable-proposal-source.test.ts @@ -141,6 +141,35 @@ describe('executable proposal source selection', () => { .resolves.toEqual([]); expect(frozen).toEqual([{ id: 'historical', digest: `sha256:${'a'.repeat(64)}` }]); }); + it('treats ready nodes from a cancelled simulation as non-admissible historical state', async () => { + const row = { proposal_id: 'old', project_id: 'project', active_version: 1, + active_content_hash: 'a'.repeat(64), accepted_decision_id: 'decision', + decision_record_json: { proposalRef: { id: 'old' } } }; + const all = vi.fn(async (query: string) => query.includes('FROM execution_nodes') ? [{ + source_ref_json: { model: 'proposal', id: 'old', digest: `sha256:${'a'.repeat(64)}` }, + status: 'ready', workday_id: 'stopped', workday_status: 'cancelled', + }] : [row]); + exactProposal.mockRejectedValueOnce(Object.assign(new Error('Invalid old content.'), { + status: 422, code: 'proposal_execution_plan_invalid', + })); + const frozen: string[] = []; + await expect(loadTeamExecutableProposalSources({ all }, 'team', undefined, (source) => frozen.push(source.id))) + .resolves.toEqual([]); + expect(frozen).toEqual(['old']); + }); + it('refuses to freeze invalid accepted content with an active ready assignment', async () => { + const all = vi.fn(async (query: string) => query.includes('FROM execution_nodes') ? [{ + source_ref_json: { model: 'proposal', id: 'active', digest: `sha256:${'a'.repeat(64)}` }, + status: 'ready', workday_id: 'running', workday_status: 'running', + }] : [{ proposal_id: 'active', project_id: 'project', active_version: 1, + active_content_hash: 'a'.repeat(64), accepted_decision_id: 'decision', + decision_record_json: { proposalRef: { id: 'active' } } }]); + exactProposal.mockRejectedValueOnce(Object.assign(new Error('Invalid active content.'), { + status: 422, code: 'proposal_execution_plan_invalid', + })); + await expect(loadTeamExecutableProposalSources({ all }, 'team', undefined, () => undefined)) + .rejects.toMatchObject({ code: 'proposal_execution_plan_invalid' }); + }); it('retains a blocked historical component during ordinary team reconciliation', async () => { const source = { store: 'treedx', model: 'proposal', id: 'historical', revision: 1, digest: `sha256:${'a'.repeat(64)}`, repository: 'library', commit: 'b'.repeat(40), path: 'proposals/old.mdx' }; From a8117d6cba03f172dcfc5bc0c7281b1ed7829f1d Mon Sep 17 00:00:00 2001 From: Adrian Webb Date: Mon, 28 Sep 2026 23:39:50 -0400 Subject: [PATCH 07/10] Quarantine invalid historical demand without active assignments --- .../execution/executable-proposal-source.ts | 21 ++++++++++--------- .../graph/executable-proposal-source.test.ts | 8 +++---- 2 files changed, 15 insertions(+), 14 deletions(-) diff --git a/src/api/capacity/services/capacity/execution/executable-proposal-source.ts b/src/api/capacity/services/capacity/execution/executable-proposal-source.ts index 365f8712..0ff9c5d7 100644 --- a/src/api/capacity/services/capacity/execution/executable-proposal-source.ts +++ b/src/api/capacity/services/capacity/execution/executable-proposal-source.ts @@ -30,19 +30,20 @@ export async function loadTeamExecutableProposalSources(store: any, teamId: stri WHERE p.team_id = ? AND ((p.decision_id IS NULL AND p.status IN ('draft','submitted','open','voting')) OR d.id IS NOT NULL) ${projectId ? 'AND p.project_id = ?' : ''} ORDER BY p.project_id,p.id`, projectId ? [teamId, projectId] : [teamId]), - store.all(`SELECT node.source_ref_json,node.status,node.workday_id,run.status AS workday_status - FROM execution_nodes node LEFT JOIN capacity_workday_runs run - ON run.team_id=node.team_id AND run.id=node.workday_id WHERE node.team_id=?`, [teamId])]); - const graphState = new Map(); + store.all(`SELECT node.source_ref_json,node.status, + EXISTS (SELECT 1 FROM capacity_provider_assignments assignment + WHERE assignment.team_id=node.team_id AND assignment.execution_node_id=node.id + AND assignment.status IN ('pending','leased','running','returned')) AS active_assignment + FROM execution_nodes node WHERE node.team_id=?`, [teamId])]); + const graphState = new Map(); for (const graphRow of graphRows) { const source = record(graphRow.source_ref_json); if (text(source.model) !== 'proposal') continue; const key = `${text(source.id)}\u0000${text(source.digest)}`; - const state = graphState.get(key) ?? { count: 0, incomplete: 0, ready: 0 }; + const state = graphState.get(key) ?? { count: 0, incomplete: 0, active: 0 }; state.count += 1; if (text(graphRow.status) !== 'completed') state.incomplete += 1; - if (text(graphRow.status) === 'ready' && !(text(graphRow.workday_id) - && ['cancelled', 'completed', 'failed'].includes(text(graphRow.workday_status)))) state.ready += 1; + if (graphRow.active_assignment === true) state.active += 1; graphState.set(key, state); } const sources: ExecutableProposalSource[] = []; @@ -69,9 +70,9 @@ export async function loadTeamExecutableProposalSources(store: any, teamId: stri const value = error as { status?: number; code?: string }; // An already-materialized historical component whose content no longer // satisfies the current schema may remain visible, but cannot admit work. - // Preserve its exact graph nodes; never silently drop an unmaterialized or - // ready accepted decision, or mask a TreeDX outage/authority mismatch. - if (value.code === 'proposal_execution_plan_invalid' && state?.count && state.ready === 0 && onFrozenInvalid) { + // Preserve its exact graph nodes; never silently drop an unmaterialized + // decision, interrupt an active assignment, or mask a TreeDX outage. + if (value.code === 'proposal_execution_plan_invalid' && state?.count && state.active === 0 && onFrozenInvalid) { onFrozenInvalid({ id: text(row.proposal_id), digest: `sha256:${text(row.active_content_hash)}` }); continue; } diff --git a/tests/unit/control-plane/capacity/execution/graph/executable-proposal-source.test.ts b/tests/unit/control-plane/capacity/execution/graph/executable-proposal-source.test.ts index 0d715f71..37bac733 100644 --- a/tests/unit/control-plane/capacity/execution/graph/executable-proposal-source.test.ts +++ b/tests/unit/control-plane/capacity/execution/graph/executable-proposal-source.test.ts @@ -141,13 +141,13 @@ describe('executable proposal source selection', () => { .resolves.toEqual([]); expect(frozen).toEqual([{ id: 'historical', digest: `sha256:${'a'.repeat(64)}` }]); }); - it('treats ready nodes from a cancelled simulation as non-admissible historical state', async () => { + it('treats ready nodes without an active assignment as non-admissible historical state', async () => { const row = { proposal_id: 'old', project_id: 'project', active_version: 1, active_content_hash: 'a'.repeat(64), accepted_decision_id: 'decision', decision_record_json: { proposalRef: { id: 'old' } } }; const all = vi.fn(async (query: string) => query.includes('FROM execution_nodes') ? [{ source_ref_json: { model: 'proposal', id: 'old', digest: `sha256:${'a'.repeat(64)}` }, - status: 'ready', workday_id: 'stopped', workday_status: 'cancelled', + status: 'ready', active_assignment: false, }] : [row]); exactProposal.mockRejectedValueOnce(Object.assign(new Error('Invalid old content.'), { status: 422, code: 'proposal_execution_plan_invalid', @@ -157,10 +157,10 @@ describe('executable proposal source selection', () => { .resolves.toEqual([]); expect(frozen).toEqual(['old']); }); - it('refuses to freeze invalid accepted content with an active ready assignment', async () => { + it('refuses to freeze invalid accepted content with an active assignment', async () => { const all = vi.fn(async (query: string) => query.includes('FROM execution_nodes') ? [{ source_ref_json: { model: 'proposal', id: 'active', digest: `sha256:${'a'.repeat(64)}` }, - status: 'ready', workday_id: 'running', workday_status: 'running', + status: 'ready', active_assignment: true, }] : [{ proposal_id: 'active', project_id: 'project', active_version: 1, active_content_hash: 'a'.repeat(64), accepted_decision_id: 'decision', decision_record_json: { proposalRef: { id: 'active' } } }]); From 387fc86721d78f8b95d212f6b06bdcb9ccf6aa72 Mon Sep 17 00:00:00 2001 From: Adrian Webb Date: Mon, 28 Sep 2026 23:46:01 -0400 Subject: [PATCH 08/10] Read active assignment state once for historical graph safety --- .../execution/executable-proposal-source.ts | 13 ++++++------- .../graph/executable-proposal-source.test.ts | 6 +++--- 2 files changed, 9 insertions(+), 10 deletions(-) diff --git a/src/api/capacity/services/capacity/execution/executable-proposal-source.ts b/src/api/capacity/services/capacity/execution/executable-proposal-source.ts index 0ff9c5d7..3ec4aad3 100644 --- a/src/api/capacity/services/capacity/execution/executable-proposal-source.ts +++ b/src/api/capacity/services/capacity/execution/executable-proposal-source.ts @@ -21,7 +21,7 @@ const stable = (value: unknown): string => { /** Load structurally complete drafts for review and accepted proposals for work. */ export async function loadTeamExecutableProposalSources(store: any, teamId: string, projectId?: string, onFrozenInvalid?: (source: { id: string; digest: string }) => void): Promise { - const [rows, graphRows] = await Promise.all([store.all(`SELECT + const [rows, graphRows, activeRows] = await Promise.all([store.all(`SELECT p.id AS proposal_id,p.project_id,p.active_version,p.active_content_hash,p.metadata_json, p.decision_id,d.id AS accepted_decision_id,d.proposal_version,d.proposal_content_hash,d.decision_record_json FROM governance_proposals p @@ -30,11 +30,10 @@ export async function loadTeamExecutableProposalSources(store: any, teamId: stri WHERE p.team_id = ? AND ((p.decision_id IS NULL AND p.status IN ('draft','submitted','open','voting')) OR d.id IS NOT NULL) ${projectId ? 'AND p.project_id = ?' : ''} ORDER BY p.project_id,p.id`, projectId ? [teamId, projectId] : [teamId]), - store.all(`SELECT node.source_ref_json,node.status, - EXISTS (SELECT 1 FROM capacity_provider_assignments assignment - WHERE assignment.team_id=node.team_id AND assignment.execution_node_id=node.id - AND assignment.status IN ('pending','leased','running','returned')) AS active_assignment - FROM execution_nodes node WHERE node.team_id=?`, [teamId])]); + store.all('SELECT id,source_ref_json,status FROM execution_nodes WHERE team_id=?', [teamId]), + store.all(`SELECT DISTINCT execution_node_id FROM capacity_provider_assignments + WHERE team_id=? AND execution_node_id IS NOT NULL AND status IN ('pending','leased','running','returned')`, [teamId])]); + const activeNodeIds = new Set(activeRows.map((row: Row) => text(row.execution_node_id)).filter(Boolean)); const graphState = new Map(); for (const graphRow of graphRows) { const source = record(graphRow.source_ref_json); @@ -43,7 +42,7 @@ export async function loadTeamExecutableProposalSources(store: any, teamId: stri const state = graphState.get(key) ?? { count: 0, incomplete: 0, active: 0 }; state.count += 1; if (text(graphRow.status) !== 'completed') state.incomplete += 1; - if (graphRow.active_assignment === true) state.active += 1; + if (activeNodeIds.has(text(graphRow.id))) state.active += 1; graphState.set(key, state); } const sources: ExecutableProposalSource[] = []; diff --git a/tests/unit/control-plane/capacity/execution/graph/executable-proposal-source.test.ts b/tests/unit/control-plane/capacity/execution/graph/executable-proposal-source.test.ts index 37bac733..27fca85c 100644 --- a/tests/unit/control-plane/capacity/execution/graph/executable-proposal-source.test.ts +++ b/tests/unit/control-plane/capacity/execution/graph/executable-proposal-source.test.ts @@ -159,9 +159,9 @@ describe('executable proposal source selection', () => { }); it('refuses to freeze invalid accepted content with an active assignment', async () => { const all = vi.fn(async (query: string) => query.includes('FROM execution_nodes') ? [{ - source_ref_json: { model: 'proposal', id: 'active', digest: `sha256:${'a'.repeat(64)}` }, - status: 'ready', active_assignment: true, - }] : [{ proposal_id: 'active', project_id: 'project', active_version: 1, + id: 'active-node', source_ref_json: { model: 'proposal', id: 'active', digest: `sha256:${'a'.repeat(64)}` }, + status: 'ready', + }] : query.includes('FROM capacity_provider_assignments') ? [{ execution_node_id: 'active-node' }] : [{ proposal_id: 'active', project_id: 'project', active_version: 1, active_content_hash: 'a'.repeat(64), accepted_decision_id: 'decision', decision_record_json: { proposalRef: { id: 'active' } } }]); exactProposal.mockRejectedValueOnce(Object.assign(new Error('Invalid active content.'), { From 5361171a69580c55eff6089419be1a28c84acc1a Mon Sep 17 00:00:00 2001 From: Adrian Webb Date: Tue, 29 Sep 2026 00:16:28 -0400 Subject: [PATCH 09/10] Retry concurrent graph projection after rolled-back revision conflict --- .../lifecycle/assignment-lifecycle-service.ts | 35 ++++++++++++++-- .../planning/estimates/integration.ts | 4 +- .../assignment-graph-transition-retry.test.ts | 41 +++++++++++++++++++ 3 files changed, 76 insertions(+), 4 deletions(-) create mode 100644 tests/unit/control-plane/capacity/assignments/assignment-graph-transition-retry.test.ts diff --git a/src/api/capacity/services/capacity/assignments/lifecycle/assignment-lifecycle-service.ts b/src/api/capacity/services/capacity/assignments/lifecycle/assignment-lifecycle-service.ts index 29d6f1c4..a7242123 100644 --- a/src/api/capacity/services/capacity/assignments/lifecycle/assignment-lifecycle-service.ts +++ b/src/api/capacity/services/capacity/assignments/lifecycle/assignment-lifecycle-service.ts @@ -34,6 +34,25 @@ interface ProviderAssignmentLifecycleStore extends CapacityGovernanceDatabase { export interface ProviderAssignmentLifecycleMutationResult { assignment: DurableProviderAssignment; leaseToken: string | null; leaseSeconds: number | null; } + +export async function batchAssignmentGraphTransition(input: { + operations: Array<{ query: string; params?: unknown[] }>; + graphOperationCount: number; + batch: (operations: Array<{ query: string; params?: unknown[] }>) => Promise; + rebuildGraphOperations: () => Promise>; +}): Promise { + for (let attempt = 0; ; attempt += 1) { + try { await input.batch(input.operations); return; } + catch (error) { + const conflict = error as { code?: unknown; constraint?: unknown }; + if (attempt >= 3 || conflict.code !== '23505' + || conflict.constraint !== 'execution_graph_revisions_pkey') throw error; + const refreshed = await input.rebuildGraphOperations(); + input.operations.splice(1, input.graphOperationCount, ...refreshed); + input.graphOperationCount = refreshed.length; + } + } +} async function assertRequiredSignals(database: CapacityGovernanceDatabase, assignment: DurableProviderAssignment) { const required = Array.isArray(record(assignment.allowedOutputs).publishedSignals) ? [...new Set((record(assignment.allowedOutputs).publishedSignals as unknown[]).map(String).map((value) => value.replace(/_/gu, '-')).filter(Boolean))] : []; @@ -433,10 +452,11 @@ export class ProviderAssignmentLifecycleService { WHERE id = ? AND team_id = ? AND capacity_provider_id = ? AND membership_id = ? AND state_version = ? AND status = 'leased' AND lease_state = 'leased' AND lease_token = ? ${options.allowExpiredLease ? '' : 'AND (lease_expires_at IS NULL OR lease_expires_at > ?)'} `, params: [options.status, ...params] }]; - operations.push(...await livingExecutionLifecycleOperations({ store: this.store, assignment, + const graphOperations = await livingExecutionLifecycleOperations({ store: this.store, assignment, status: options.status, now, result: options.assignmentResult, returnCode: options.status === 'returned' ? input.code : undefined, - reviewDisposition: options.reviewDisposition ?? null })); + reviewDisposition: options.reviewDisposition ?? null }); + operations.push(...graphOperations); if (['completed','failed','cancelled'].includes(options.status)) { const terminalWorkspace = terminalAssignmentAuthority(assignment, now); operations.push({ @@ -455,7 +475,16 @@ export class ProviderAssignmentLifecycleService { params: [options.status==='cancelled'?'cancelled':'failed',assignment.id, now, JSON.stringify({ code: input.code ?? options.defaultCode, reason: input.reason ?? input.message ?? options.defaultReason }), now, assignment.invocationId, assignment.teamId], }); } - await this.store.batch(operations); + // Parallel assignment completions can observe the same graph head. The + // losing transaction is rolled back by PostgreSQL; recompute its graph + // projection against the committed head, preserving the same assignment + // transition and exactly-once result. + await batchAssignmentGraphTransition({ operations, graphOperationCount: graphOperations.length, + batch: (statements) => this.store.batch(statements), + rebuildGraphOperations: () => livingExecutionLifecycleOperations({ store: this.store, assignment, + status: options.status, now, result: options.assignmentResult, + returnCode: options.status === 'returned' ? input.code : undefined, + reviewDisposition: options.reviewDisposition ?? null }) }); const transitioned = await this.store.getProviderAssignment(principal.teamId, assignment.id); if (!transitioned || transitioned.stateVersion !== assignment.stateVersion + 1 || transitioned.status !== options.status) return null; if (assignment.operationHandoffId && (options.status === 'completed' || options.status === 'failed')) await terminalizeOperationHandoff(this.store, assignment.operationHandoffId, assignment.id, options.status, now); diff --git a/src/api/capacity/services/capacity/assignments/planning/estimates/integration.ts b/src/api/capacity/services/capacity/assignments/planning/estimates/integration.ts index dec5306f..5149a89f 100644 --- a/src/api/capacity/services/capacity/assignments/planning/estimates/integration.ts +++ b/src/api/capacity/services/capacity/assignments/planning/estimates/integration.ts @@ -129,7 +129,9 @@ async function integrateAssignmentEstimateOnce( 'assignment_estimate_content_missing', 'Estimator proposal result has no content.', 409); const parsed = validatePortableContentData('proposal', record(file.frontmatter)); if (!parsed.ok || !parsed.data) throw new CapacityGovernanceError('assignment_estimate_content_invalid', - 'Estimator proposal content is invalid.', 409, { diagnostics: parsed.diagnostics }); + `Estimator proposal content is invalid: ${parsed.diagnostics.slice(0, 5) + .map((diagnostic) => `${diagnostic.field ?? 'document'}:${diagnostic.code}`).join(', ') || 'missing data'}.`, + 409, { diagnostics: parsed.diagnostics }); const candidate = parsed.data as Row; const merged = mergeAssignmentEstimate({ frozen: frozen.definition, candidate, current: current.definition, agentClass: attempt.agentClass }); diff --git a/tests/unit/control-plane/capacity/assignments/assignment-graph-transition-retry.test.ts b/tests/unit/control-plane/capacity/assignments/assignment-graph-transition-retry.test.ts new file mode 100644 index 00000000..957c70a1 --- /dev/null +++ b/tests/unit/control-plane/capacity/assignments/assignment-graph-transition-retry.test.ts @@ -0,0 +1,41 @@ +import { describe, expect, it, vi } from 'vitest'; +import { batchAssignmentGraphTransition } from '../../../../../src/api/capacity/services/capacity/assignments/lifecycle/assignment-lifecycle-service.ts'; + +const statement = (query: string) => ({ query }); +const revisionConflict = () => Object.assign(new Error('concurrent graph revision'), { + code: '23505', constraint: 'execution_graph_revisions_pkey', +}); + +describe('concurrent assignment graph completion', () => { + it('reprojects only the rolled-back graph operations against the committed revision', async () => { + const batches: string[][] = []; + const batch = vi.fn(async (operations: Array<{ query: string }>) => { + batches.push(operations.map((operation) => operation.query)); + if (batches.length === 1) throw revisionConflict(); + }); + const rebuildGraphOperations = vi.fn(async () => [statement('graph revision 12')]); + await batchAssignmentGraphTransition({ + operations: [statement('same assignment transition'), statement('graph revision 11'), statement('same teardown')], + graphOperationCount: 1, batch, rebuildGraphOperations, + }); + expect(batches).toEqual([ + ['same assignment transition', 'graph revision 11', 'same teardown'], + ['same assignment transition', 'graph revision 12', 'same teardown'], + ]); + expect(rebuildGraphOperations).toHaveBeenCalledOnce(); + }); + + it('fails closed on unrelated uniqueness errors and bounded repeated contention', async () => { + const unrelated = Object.assign(new Error('other unique key'), { code: '23505', constraint: 'other_pkey' }); + const rebuildGraphOperations = vi.fn(async () => [statement('new graph')]); + await expect(batchAssignmentGraphTransition({ operations: [statement('assignment')], graphOperationCount: 0, + batch: async () => { throw unrelated; }, rebuildGraphOperations })).rejects.toBe(unrelated); + expect(rebuildGraphOperations).not.toHaveBeenCalled(); + let attempts = 0; + await expect(batchAssignmentGraphTransition({ operations: [statement('assignment')], graphOperationCount: 0, + batch: async () => { attempts++; throw revisionConflict(); }, rebuildGraphOperations })).rejects.toMatchObject({ + constraint: 'execution_graph_revisions_pkey', + }); + expect(attempts).toBe(4); + }); +}); From 9c20113ea1466ed1af32771893253a99f4db342e Mon Sep 17 00:00:00 2001 From: Adrian Webb Date: Tue, 29 Sep 2026 01:09:10 -0400 Subject: [PATCH 10/10] Guard invocation binding in conversation admission --- .../admission/living-execution-admission.ts | 26 ++++++++++++++++++- .../living-execution-admission.test.ts | 19 +++++++++++++- 2 files changed, 43 insertions(+), 2 deletions(-) diff --git a/src/api/capacity/services/capacity/assignments/admission/living-execution-admission.ts b/src/api/capacity/services/capacity/assignments/admission/living-execution-admission.ts index 088f3e39..29813786 100644 --- a/src/api/capacity/services/capacity/assignments/admission/living-execution-admission.ts +++ b/src/api/capacity/services/capacity/assignments/admission/living-execution-admission.ts @@ -100,6 +100,8 @@ export async function admitLivingExecutionAssignment(store: Store, input: { // concurrency. The count in the reservation INSERT is therefore atomic. { query: `SELECT id FROM capacity_workday_runs WHERE team_id=? AND id=? AND status='running' FOR UPDATE`, params: [assignment.teamId, assignment.workdayId] }, + ...(input.invocationId ? [{ query: `SELECT id FROM agent_invocation_requests WHERE id=? AND team_id=? FOR UPDATE`, + params: [input.invocationId, assignment.teamId] }] : []), ...initializeCapabilityCounters(assignment, claims, input.now), { query: `SELECT node.id FROM execution_nodes node WHERE node.team_id=? AND node.id=? AND node.node_revision=? AND node.status='ready' @@ -132,6 +134,12 @@ export async function admitLivingExecutionAssignment(store: Store, input: { AND ${claims.map(() => `EXISTS (SELECT 1 FROM capacity_admission_counters WHERE id=? AND committed_amount+?<=LEAST(hard_limit,?))`).join(' AND ')} AND EXISTS (SELECT 1 FROM capacity_workday_runs run WHERE run.team_id=? AND run.id=? AND run.status='running') + ${input.invocationId ? `AND EXISTS (SELECT 1 FROM agent_invocation_requests invocation + WHERE invocation.id=? AND invocation.team_id=? AND invocation.status IN ('admitted','running') + AND (invocation.assignment_id IS NULL OR invocation.assignment_id=? OR EXISTS ( + SELECT 1 FROM capacity_provider_assignments prior WHERE prior.id=invocation.assignment_id + AND prior.team_id=invocation.team_id AND prior.invocation_id=invocation.id + AND prior.status IN ('returned','failed','cancelled'))))` : ''} AND (SELECT COUNT(*) FROM capacity_provider_assignments active WHERE active.team_id=? AND active.work_day_id=? AND active.execution_kind=? AND active.status IN ('pending','leased','running')) < ? @@ -145,6 +153,7 @@ export async function admitLivingExecutionAssignment(store: Store, input: { nodeRevision: assignment.nodeRevision, graphRevision: assignment.graphRevision }),input.now,input.now,admissionToken,...common, ...claims.flatMap(claim => [claim.id, assignment.limits.maximumSeconds, claim.hardLimit]), assignment.teamId, assignment.workdayId, + ...(input.invocationId ? [input.invocationId, assignment.teamId, assignment.id] : []), assignment.teamId, assignment.workdayId, input.executionKind, input.workdayConcurrencyLimit, assignment.teamId, principal.capacityProviderId, input.laneId, providerConcurrencyLimit] }, ...commitCapabilityCounters(assignment, claims, admissionToken, input.now), @@ -178,7 +187,10 @@ export async function admitLivingExecutionAssignment(store: Store, input: { predecessorResults: input.predecessorResults, authorizedContext, treedxProxyHandle: input.treedxProxyHandle }),input.now, assignment.id,assignment.teamId,assignment.reservationId] }, ...(input.invocationId ? [{ query: `UPDATE agent_invocation_requests SET assignment_id=?, status='running', updated_at=? - WHERE id=? AND team_id=? AND status IN ('admitted','running') AND (assignment_id IS NULL OR assignment_id=?)`, + WHERE id=? AND team_id=? AND status IN ('admitted','running') AND (assignment_id IS NULL OR assignment_id=? OR EXISTS ( + SELECT 1 FROM capacity_provider_assignments prior WHERE prior.id=agent_invocation_requests.assignment_id + AND prior.team_id=agent_invocation_requests.team_id AND prior.invocation_id=agent_invocation_requests.id + AND prior.status IN ('returned','failed','cancelled')))`, params: [assignment.id,input.now,input.invocationId,assignment.teamId,assignment.id] }] : []), { query: `INSERT INTO treedx_proxy_handles ( id,team_id,project_id,assignment_id,repository_id,workspace_id,status,scopes_json, @@ -201,6 +213,18 @@ export async function admitLivingExecutionAssignment(store: Store, input: { assignment.nodeId,assignment.nodeRevision] }, ]); const committed = await store.getProviderAssignment(assignment.teamId, assignment.id); + if (input.invocationId) { + const invocation = await store.first('SELECT status,assignment_id FROM agent_invocation_requests WHERE id=? AND team_id=?', + [input.invocationId, assignment.teamId]); + if (committed && invocation?.assignment_id !== assignment.id) throw new CapacityGovernanceError( + 'communication_invocation_binding_failed', 'Conversation assignment was not bound to its authoritative invocation.', 409, + { invocationId: input.invocationId, assignmentId: assignment.id, + observedStatus: invocation?.status ?? null, observedAssignmentId: invocation?.assignment_id ?? null }); + if (!committed && (!invocation || !['admitted','running'].includes(String(invocation.status)))) throw new CapacityGovernanceError( + 'communication_invocation_not_admissible', 'Conversation invocation is not available for assignment admission.', 409, + { invocationId: input.invocationId, observedStatus: invocation?.status ?? null, + observedAssignmentId: invocation?.assignment_id ?? null }); + } if (!committed) { const active = await store.first(`SELECT COUNT(*) AS active_count FROM capacity_provider_assignments WHERE team_id=? AND work_day_id=? AND execution_kind=? AND status IN ('pending','leased','running')`, diff --git a/tests/unit/control-plane/capacity/execution/living-execution-admission.test.ts b/tests/unit/control-plane/capacity/execution/living-execution-admission.test.ts index a4825049..e349fd35 100644 --- a/tests/unit/control-plane/capacity/execution/living-execution-admission.test.ts +++ b/tests/unit/control-plane/capacity/execution/living-execution-admission.test.ts @@ -104,7 +104,8 @@ describe('living execution admission', () => { it('binds a conversation invocation to the assignment in the admission transaction', async () => { const committed = { id: assignment.id, executionNodeId: 'node', executionNodeRevision: 1 }; - const store = { getProviderAssignment: vi.fn().mockResolvedValueOnce(null).mockResolvedValueOnce(committed), batch: vi.fn(async () => []) }; + const store = { getProviderAssignment: vi.fn().mockResolvedValueOnce(null).mockResolvedValueOnce(committed), + batch: vi.fn(async () => []), first: vi.fn(async () => ({ status: 'running', assignment_id: assignment.id })) }; await admitLivingExecutionAssignment(store as never, { principal: { teamId: 'team', capacityProviderId: 'provider', membershipId: 'membership' } as never, accountingLimits, assignment: assignment as never, allocation, projectAgentClassId: 'class', providerSessionId: 'session', executionProviderId: 'runtime', laneId: 'communication', lanePurpose: 'communication', executionKind: 'conversation', workdayConcurrencyLimit: 2, invocationId: 'invocation-1', predecessorResults: [], treedxProxyHandle: { id: 'tdx_assignment', status: 'issued', @@ -113,5 +114,21 @@ describe('living execution admission', () => { expect(binding.params).toEqual(['assignment', assignment.createdAt, 'invocation-1', 'team', 'assignment']); const reservation = store.batch.mock.calls[0]![0].find((operation: { query: string }) => operation.query.includes('INSERT INTO capacity_reservations'))!; expect(reservation.params.slice(-8)).toEqual(['team', 'workday', 'conversation', 2, 'team', 'provider', 'communication', 1]); + expect(reservation.query).toContain("invocation.status IN ('admitted','running')"); + expect(reservation.query).toContain("prior.status IN ('returned','failed','cancelled')"); + const operations = store.batch.mock.calls[0]![0] as Array<{ query: string; params: unknown[] }>; + expect(operations.findIndex((operation) => operation.query.includes('agent_invocation_requests WHERE id=? AND team_id=? FOR UPDATE'))) + .toBeLessThan(operations.findIndex((operation) => operation.query.includes('INSERT INTO capacity_reservations'))); + for (const operation of operations) expect((operation.query.match(/\?/gu) ?? []).length).toBe(operation.params.length); + }); + it('fails closed when a committed conversation assignment lacks the exact invocation binding', async () => { + const committed = { id: assignment.id, executionNodeId: 'node', executionNodeRevision: 1 }; + const store = { getProviderAssignment: vi.fn().mockResolvedValueOnce(null).mockResolvedValueOnce(committed), + batch: vi.fn(async () => []), first: vi.fn(async () => ({ status: 'running', assignment_id: 'other-active-assignment' })) }; + await expect(admitLivingExecutionAssignment(store as never, { principal: { teamId: 'team', capacityProviderId: 'provider', membershipId: 'membership' } as never, + accountingLimits, assignment: assignment as never, allocation, projectAgentClassId: 'class', providerSessionId: 'session', executionProviderId: 'runtime', laneId: 'communication', + lanePurpose: 'communication', executionKind: 'conversation', workdayConcurrencyLimit: 2, invocationId: 'invocation-1', predecessorResults: [], treedxProxyHandle: { id: 'tdx_assignment', status: 'issued', + allowedPaths: [], allowedReadPaths: [], allowedWritePaths: [], scopes: [], allowedOperations: [] }, now: assignment.createdAt })) + .rejects.toMatchObject({ code: 'communication_invocation_binding_failed', details: { observedAssignmentId: 'other-active-assignment' } }); }); });