Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ export type LivingAllocationInputs = Record<string, { measurements: AllocationMe
/** Read existing graph/reservation/usage authority; no performance or allocation store. */
export async function livingAllocationInputs(store: CapacityGovernanceDatabase, input: {
run: DurableCapacityWorkdayRun; runs: DurableCapacityWorkdayRun[]; providers: ProviderSynthesisExecutionProvider[]; capacityProviderId: string;
capabilityId: string; agentClass: string; activity: string; now: string;
capabilityId: string; agentClass: string; activity: string; proposalGovernanceReview?: boolean; now: string;
}): Promise<LivingAllocationInputs> {
const result: LivingAllocationInputs = {};
for (const provider of input.providers) {
Expand Down Expand Up @@ -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)))
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand Down Expand Up @@ -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')) < ?
Expand All @@ -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),
Expand Down Expand Up @@ -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,
Expand All @@ -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')`,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<unknown>;
rebuildGraphOperations: () => Promise<Array<{ query: string; params?: unknown[] }>>;
}): Promise<void> {
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))] : [];
Expand Down Expand Up @@ -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({
Expand All @@ -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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 });
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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) },
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,8 +19,9 @@ 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<ExecutableProposalSource[]> {
const [rows, graphRows] = await Promise.all([store.all(`SELECT
export async function loadTeamExecutableProposalSources(store: any, teamId: string, projectId?: string,
onFrozenInvalid?: (source: { id: string; digest: string }) => void): Promise<ExecutableProposalSource[]> {
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
Expand All @@ -29,15 +30,19 @@ 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])]);
const graphState = new Map<string, { count: number; incomplete: number }>();
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<string, { count: number; incomplete: number; active: number }>();
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, active: 0 };
state.count += 1;
if (text(graphRow.status) !== 'completed') state.incomplete += 1;
if (activeNodeIds.has(text(graphRow.id))) state.active += 1;
graphState.set(key, state);
}
const sources: ExecutableProposalSource[] = [];
Expand All @@ -62,6 +67,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
// 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;
}
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 });
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<typeof advanceLivingWorkday>[0];

Expand All @@ -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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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';

Expand Down Expand Up @@ -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',
Expand Down
Loading
Loading