diff --git a/packages/api/mcp/operationLifecycle.ts b/packages/api/mcp/operationLifecycle.ts new file mode 100644 index 000000000..2fd36aeed --- /dev/null +++ b/packages/api/mcp/operationLifecycle.ts @@ -0,0 +1,211 @@ +import { redactSecrets, type McpErrorEnvelope } from './errorEnvelope.js'; +import type { McpOperations, LifecycleOutcome, Operation } from './operations.js'; + +const startedTaskStates = new Set(['processing', 'claude_execution', 'post_processing']); +const executedGoalTaskStates = new Set([...startedTaskStates, 'completed', 'failed']); +const terminalStates = new Set(['completed', 'failed', 'cancelled']); + +function record(value: unknown): Record | undefined { + return value !== null && typeof value === 'object' && !Array.isArray(value) ? value as Record : undefined; +} + +function positiveInteger(...values: unknown[]): number | undefined { + for (const value of values) { + const number = Number(value); + if (Number.isSafeInteger(number) && number > 0) return number; + } + return undefined; +} + +function nonEmptyString(...values: unknown[]): string | undefined { + return values.find(value => typeof value === 'string' && value.length > 0) as string | undefined; +} + +function epochMilliseconds(value: unknown): number | undefined { + if (typeof value === 'number' && Number.isFinite(value)) return value; + if (typeof value !== 'string') return undefined; + const parsed = Date.parse(value); + return Number.isNaN(parsed) ? undefined : parsed; +} + +function taskIdFromReceipt( + target: Record, + continuation: Record, + result: Record, + targetIssues: Record[], +): string | undefined { + const currentTask = record(target.currentTask) ?? {}; + return nonEmptyString(target.taskId, target.task_id, target.current_task_id, currentTask.taskId, + currentTask.task_id, continuation.taskId, result.taskId, + ...targetIssues.flatMap(issue => [issue.taskId, issue.task_id])); +} + +function issueNumbersFromReceipt( + result: Record, + continuation: Record, + target: Record, + targetIssues: Record[], +): number[] { + const issueNumbers = new Set(); + for (const value of [result.issueNumber, result.issue_number, continuation.issueNumber, target.issueNumber, target.issue_number]) { + const number = positiveInteger(value); + if (number) issueNumbers.add(number); + } + if (Array.isArray(result.issues)) for (const value of result.issues) { + const issue = record(value); + const number = positiveInteger(issue?.number, issue?.issueNumber, value); + if (number) issueNumbers.add(number); + } + for (const issue of targetIssues) { + const number = positiveInteger(issue.number, issue.issueNumber, issue.issue_number); + if (number) issueNumbers.add(number); + } + return [...issueNumbers]; +} + +/** Collect stable output handles from mutation results and tracker observations. */ +export function artifactsFromReceipt(row: Pick, receipt: Record): Record { + const result = record(receipt.result) ?? {}; + const continuation = record(result.continuation) ?? {}; + const target = record(receipt.targetState) ?? {}; + const targetIssues: Record[] = Array.isArray(target.issues) + ? target.issues.map(record).filter((value): value is Record => !!value) : []; + const artifacts: Record = {}; + + const taskId = taskIdFromReceipt(target, continuation, result, targetIssues); + if (taskId) artifacts.taskId = taskId; + + const repository = nonEmptyString(result.repository, row.repository); + const pullRequestNumber = positiveInteger( + result.pullRequest, result.pr_number, result.prNumber, + continuation.pullRequest, continuation.pr_number, + target.pullRequest, target.pr_number, target.final_pr_number, + ...targetIssues.flatMap(issue => [issue.pullRequest, issue.pr_number]), + ); + if (repository && pullRequestNumber) artifacts.pullRequest = { + repository, + number: pullRequestNumber, + url: `https://github.com/${repository}/pull/${pullRequestNumber}`, + }; + + const issueNumbers = issueNumbersFromReceipt(result, continuation, target, targetIssues); + if (repository && issueNumbers.length) artifacts.issues = issueNumbers.map(number => ({ + repository, + number, + url: `https://github.com/${repository}/issues/${number}`, + })); + + const commentId = positiveInteger(result.commentId, continuation.commentId, target.commentId); + if (commentId) artifacts.commentId = commentId; + return artifacts; +} + +function publicFailure(code: string, message: string, details?: Record): McpErrorEnvelope { + return { code, message: redactSecrets(message), stage: 'internal', retryable: false, status: 500, ...(details ? { details } : {}) }; +} + +function resultFailure(result: Record | undefined): McpErrorEnvelope | undefined { + const resultError = record(result?.error); + if (resultError && typeof resultError.code === 'string' && typeof resultError.message === 'string' + && typeof resultError.retryable === 'boolean' && typeof resultError.status === 'number') { + return resultError as unknown as McpErrorEnvelope; + } + + const loop = record(result?.loop); + return loop?.completionStatus === 'failed' + ? publicFailure('EXECUTION_FAILED', nonEmptyString(loop.completionReason) ?? 'Ultrafix loop failed.') + : undefined; +} + +/** Normalize durable backend failure evidence into the public error envelope. */ +export function failureFromReceipt(receipt: Record): McpErrorEnvelope | undefined { + const result = record(receipt.result); + const backendFailure = resultFailure(result); + if (backendFailure) return backendFailure; + + const target = record(receipt.targetState); + const reviewResults = Array.isArray(target?.reviewResults) ? target.reviewResults + : Array.isArray(result?.reviewResults) ? result.reviewResults : []; + const failedReviews = reviewResults.map(record).filter((review): review is Record => review?.success === false); + if (failedReviews.length && failedReviews.length === reviewResults.length) { + const reasons = failedReviews.flatMap(review => typeof review.error === 'string' && review.error.length ? [review.error] : []); + return publicFailure('REVIEW_FAILED', reasons.length ? reasons.join('; ') : 'Every requested review failed.', { + failedReviewCount: failedReviews.length, + }); + } + + const currentTask = record(target?.currentTask); + const reason = nonEmptyString(target?.failure_reason, currentTask?.reason, target?.reason, result?.reason); + return reason ? publicFailure('EXECUTION_FAILED', reason) : undefined; +} + +function lifecycleOutcome( + row: Operation, + target: Record | undefined, + receiptState: string, + targetState: string, +): LifecycleOutcome | undefined { + if (terminalStates.has(receiptState as LifecycleOutcome)) return receiptState as LifecycleOutcome; + if (row.tool === 'run_ultrafix' || (!target?.taskId && !target?.task_id)) return undefined; + return terminalStates.has(targetState as LifecycleOutcome) ? targetState as LifecycleOutcome : undefined; +} + +async function syncCancellation( + operations: McpOperations, + row: Operation, + receipt: Record, +): Promise { + if (row.tool !== 'cancel_operation') return; + const result = record(receipt.result); + const sourceId = nonEmptyString(result?.operationId); + if (sourceId && result?.cancellation === 'confirmed') await operations.finishCancellationSource(row, sourceId); +} + +function observedStartTimestamp( + target: Record | undefined, + result: Record | undefined, + targetState: string, +): number | null | undefined { + const loop = record(result?.loop); + const currentTask = record(target?.currentTask); + const currentTaskState = String(currentTask?.state ?? ''); + const observedStart = startedTaskStates.has(targetState) || executedGoalTaskStates.has(currentTaskState) + || target?.queueState === 'active' || loop?.active === true; + if (!observedStart) return null; + return epochMilliseconds(currentTask?.timestamp) ?? epochMilliseconds(target?.timestamp); +} + +/** Persist tracker observations without allowing stale concurrent polls to undo newer lifecycle facts. */ +export async function syncLifecycle( + operations: McpOperations, + row: Operation, + receipt: Record, +): Promise { + const artifacts = artifactsFromReceipt(row, receipt); + if (Object.keys(artifacts).length) await operations.recordArtifacts(row.id, artifacts); + + const target = record(receipt.targetState); + const targetState = String(target?.state ?? ''); + const receiptState = String(receipt.state ?? ''); + const result = record(receipt.result); + const outcome = lifecycleOutcome(row, target, receiptState, targetState); + if (target && !outcome) await operations.recordProgress(row.id, target); + + const startedAt = observedStartTimestamp(target, result, targetState); + if (startedAt !== null) await operations.markStarted(row.id, startedAt); + + if (outcome) { + // Persist the terminal snapshot in the same guarded update as the outcome. + // This closes the window where an older nonterminal poll could otherwise + // replace terminal progress between two lifecycle writes. + await operations.finish(row.id, outcome, outcome === 'failed' ? failureFromReceipt(receipt) : undefined, target); + } else if (receiptState === 'unknown') { + await operations.markUnknown(row.id); + } else if (receiptState === 'queued' && target) { + // A task or queue can appear after an earlier timeout. Resolve pre-start + // uncertainty without erasing evidence that execution had already begun. + await operations.markAccepted(row.id); + } + + await syncCancellation(operations, row, receipt); +} diff --git a/packages/api/mcp/operationTracking.ts b/packages/api/mcp/operationTracking.ts index 81dc2ac74..945b142aa 100644 --- a/packages/api/mcp/operationTracking.ts +++ b/packages/api/mcp/operationTracking.ts @@ -69,37 +69,87 @@ export async function trackCancellation(deps: ToolDeps, row: Operation, principa await deps.db('mcp_operations').where({ id: row.id }).whereNotIn('state', terminalStates).update({ state: receipt.state, result: JSON.stringify(result), updated_at: Date.now() }); } +async function refreshPullRequestContext( + row: Operation, + principal: McpPrincipal, + result: ExecutionResult & Record, +): Promise { + if (!result.pullRequest) return; + const [owner, repo] = String(row.repository).split('/'); + const { data: pr } = await principal.github.request('GET /repos/{owner}/{repo}/pulls/{pull_number}', { + owner, repo, pull_number: result.pullRequest, + }); + result.currentHead = pr.head.sha; + result.results = { + tool: 'get_pull_request_discussion', repository: row.repository, + pullRequest: result.pullRequest, taskId: result.continuation?.taskId, + }; +} + +function restoreResolvedTarget( + receipt: Record, + result: ExecutionResult, + task: TrackingContext['task'] | undefined, +): void { + const target = result.targetState + ?? (receipt.targetState as Record | undefined) + ?? {}; + if (!task && !Object.keys(target).length) return; + receipt.targetState = { + ...target, + ...(task ? { taskId: task.task_id, pr_number: task.pr_number } : {}), + }; +} + /** Resolve the execution from the actual job or the exact triggering comment. */ export async function trackExecution(deps: ToolDeps, row: Operation, principal: McpPrincipal, receipt: Record): Promise { if (!trackedTools.includes(row.tool) || !row.result) return; const { db } = deps; - const result = JSON.parse(row.result); + const result = JSON.parse(row.result) as ExecutionResult & Record; if (['create_task', 'retry_task_submission'].includes(row.tool) && !result.continuation?.taskId) return; - if (result.error || (result.executionResolved && terminalStates.includes(row.state))) return; + if (result.error) return; const task = row.tool === 'index_repository' ? undefined : await findExecutionTask(deps, row, result); + if (result.executionResolved && terminalStates.includes(row.state)) { + restoreResolvedTarget(receipt, result, task); + return; + } if (task) await trackTask(deps, row, { task, result, receipt }); else if (result.jobId) await trackQueuedJob(row, result.jobId, receipt); else if (Date.now() - Number(row.created_at) > 120000) receipt.state = 'unknown'; - if (result.pullRequest) { - const [owner, repo] = String(row.repository).split('/'); - const { data: pr } = await principal.github.request('GET /repos/{owner}/{repo}/pulls/{pull_number}', { owner, repo, pull_number: result.pullRequest }); - result.currentHead = pr.head.sha; - result.results = { tool: 'get_pull_request_discussion', repository: row.repository, pullRequest: result.pullRequest, taskId: result.continuation?.taskId }; + await refreshPullRequestContext(row, principal, result); + if (terminalStates.includes(String(receipt.state))) { + result.executionResolved = true; + result.targetState = receipt.targetState as Record | undefined; } - if (terminalStates.includes(String(receipt.state))) result.executionResolved = true; receipt.result = result; if (receipt.state === 'unknown') receipt.message = 'Execution cannot yet be confirmed. Inspect the linked comment/job; polling can still resolve it. Do not blindly resubmit.'; - await db('mcp_operations').where({ id: row.id }).update({ state: receipt.state, result: JSON.stringify(result), updated_at: Date.now() }); + const recorded = await db('mcp_operations').where({ id: row.id }).whereNotIn('state', terminalStates) + .update({ state: receipt.state, result: JSON.stringify(result), updated_at: Date.now() }); + if (!recorded) { + // Another poll persisted terminal evidence while this observation was + // awaiting external context. Return that durable receipt and let lifecycle + // synchronization use its terminal target instead of this stale snapshot. + const current = await db('mcp_operations').where({ id: row.id }).first(); + if (current && terminalStates.includes(current.state)) { + const currentResult = current.result ? JSON.parse(current.result) as ExecutionResult & Record : {}; + receipt.state = current.state; + receipt.result = currentResult; + if (currentResult.targetState) receipt.targetState = currentResult.targetState; + else delete receipt.targetState; + delete receipt.message; + } + } } interface ExecutionResult { jobId?: string; commentId?: number; pullRequest?: number; continuation?: { taskId?: string; jobId?: string; sourceTaskId?: string }; + targetState?: Record; reviewResults?: Array<{ success: boolean; commentId?: number; commentUrl?: string }>; loop?: { completionStatus?: string | null } & Record; } interface TrackingContext { - task: { task_id: string; initial_job_data: unknown }; + task: { task_id: string; pr_number: number | null; initial_job_data: unknown }; result: ExecutionResult; receipt: Record; } @@ -134,14 +184,15 @@ async function findExecutionTask(deps: ToolDeps, row: Operation, result: Executi } else { query.andWhere(builder => builder.where('task_id', result.jobId || continuation.taskId).orWhere('job_id', result.jobId || continuation.jobId)); } - return query.orderBy('created_at', 'desc').first('task_id', 'initial_job_data'); + return query.orderBy('created_at', 'desc').first('task_id', 'pr_number', 'initial_job_data'); } async function trackTask(deps: ToolDeps, row: Operation, { task, result, receipt }: TrackingContext): Promise { result.continuation = { ...result.continuation, taskId: task.task_id }; const event = await deps.db('task_history').where({ task_id: task.task_id }).orderBy('history_id', 'desc').first('state', 'timestamp', 'reason', 'metadata'); const metadata = typeof event?.metadata === 'string' ? JSON.parse(event.metadata) : event?.metadata; - receipt.targetState = { taskId: task.task_id, state: event?.state, timestamp: event?.timestamp, reason: event?.reason, reviewResults: metadata?.reviewResults }; + receipt.targetState = { taskId: task.task_id, pr_number: task.pr_number, state: event?.state, + timestamp: event?.timestamp, reason: event?.reason, reviewResults: metadata?.reviewResults }; if (metadata?.reviewResults) result.reviewResults = metadata.reviewResults; receipt.state = terminalStates.includes(event?.state) ? event.state : !event || event.state === 'pending' ? 'queued' : 'running'; if (event?.state === 'completed' && result.reviewResults?.length && result.reviewResults.every(review => !review.success)) receipt.state = 'failed'; diff --git a/packages/api/mcp/operations.ts b/packages/api/mcp/operations.ts index 827ad39bb..17cf38112 100644 --- a/packages/api/mcp/operations.ts +++ b/packages/api/mcp/operations.ts @@ -1,14 +1,21 @@ import { randomUUID } from 'node:crypto'; import type { Knex } from 'knex'; import { McpError } from './config.js'; -import { classifyError } from './errorEnvelope.js'; +import { classifyError, type McpErrorEnvelope } from './errorEnvelope.js'; import { digest } from './store.js'; import type { McpPrincipal } from './policy.js'; +import { artifactsFromReceipt, failureFromReceipt } from './operationLifecycle.js'; + +const interruptionTimeoutMs = 120_000; export interface OperationResult { status: number; data: unknown } +export type LifecycleState = 'accepted' | 'running' | 'completed' | 'failed' | 'cancelled' | 'unknown'; +export type LifecycleOutcome = 'completed' | 'failed' | 'cancelled'; export interface Operation { id: string; owner_id: string; grant_id: string; idempotency_key: string; tool: string; repository: string | null; state: string; result: string | null; created_at: number; updated_at: number; payload_hash: string; + lifecycle: LifecycleState; accepted_at: number; started_at: number | null; finished_at: number | null; + failure: string | null; artifacts: string | null; progress: string | null; } function canonical(value: unknown): string { @@ -17,13 +24,75 @@ function canonical(value: unknown): string { return JSON.stringify(value); } +function json(value: unknown): unknown { + if (typeof value !== 'string') return value ?? null; + try { return JSON.parse(value); } catch { return null; } +} + +function record(value: unknown): Record | undefined { + return value !== null && typeof value === 'object' && !Array.isArray(value) ? value as Record : undefined; +} + +function recoveryReceipt(row: Pick) { + const result = json(row.result); + const targetState = record(record(result)?.targetState); + return { state: row.state, result, ...(targetState ? { targetState } : {}) }; +} + +function confirmedCancellationSource(row: Pick, receipt: ReturnType): string | undefined { + const result = record(receipt.result); + return row.tool === 'cancel_operation' && result?.cancellation === 'confirmed' + && typeof result.operationId === 'string' && result.operationId.length > 0 ? result.operationId : undefined; +} + +function iso(value: number | null | undefined): string | null { + if (value === null || value === undefined) return null; + const date = new Date(Number(value)); + return Number.isNaN(date.getTime()) ? null : date.toISOString(); +} + +function errorEnvelope(value: unknown): McpErrorEnvelope | undefined { + if (!value || typeof value !== 'object') return undefined; + const envelope = value as Partial; + return typeof envelope.code === 'string' && typeof envelope.message === 'string' + && typeof envelope.retryable === 'boolean' && typeof envelope.status === 'number' + ? envelope as McpErrorEnvelope : undefined; +} + +function operationState(result: OperationResult): string { + const reported = (result.data as { state?: string })?.state; + if (reported === 'browser_required') return reported; + if (result.status !== 202) return 'completed'; + return ['posted', 'queued', 'unknown', 'failed'].includes(reported || '') ? reported! : 'accepted'; +} + +function invocationInterrupted(row: Operation, now = Date.now()): boolean { + const invokedAt = Number(row.accepted_at); + return row.state === 'accepted' && row.result === null && Number.isFinite(invokedAt) + && now - invokedAt > interruptionTimeoutMs; +} + +function needsProgressRecovery(targetState: unknown, lifecycleMissing: boolean, progress: string | null): boolean { + if (targetState === undefined) return false; + return lifecycleMissing || progress === null; +} + +function recoveredProgress(db: Knex, targetState: unknown, lifecycleMissing: boolean): unknown { + if (lifecycleMissing) return JSON.stringify(targetState); + return db.raw('COALESCE(progress, ?)', [JSON.stringify(targetState)]); +} + export class McpOperations { constructor(readonly db: Knex) {} async replay(principal: McpPrincipal, tool: string, args: Record): Promise | undefined> { - const previous = await this.db('mcp_operations').where({ owner_id: principal.user.id, grant_id: principal.grant.id, idempotency_key: String(args.idempotencyKey) }).first(); + const identity = { owner_id: principal.user.id, grant_id: principal.grant.id, idempotency_key: String(args.idempotencyKey) }; + let previous = await this.db('mcp_operations').where(identity).first(); if (!previous) return undefined; if (previous.payload_hash !== digest(canonical({ tool, args }))) throw new McpError('IDEMPOTENCY_CONFLICT', 'This key was already used with different arguments.', 409); + await this.reconcileTerminalLifecycles(principal, previous.id); + await this.markInterruptedInvocations(principal, previous.id); + previous = (await this.db('mcp_operations').where(identity).first())!; return this.project(previous); } @@ -33,38 +102,216 @@ export class McpOperations { const identity = { owner_id: principal.user.id, grant_id: principal.grant.id, idempotency_key: key }; const payloadHash = digest(canonical({ tool, args })); const id = randomUUID(); + const acceptedAt = Date.now(); const inserted = await this.db('mcp_operations').insert({ ...identity, id, tool, repository: repository || null, - payload_hash: payloadHash, state: 'running', created_at: Date.now(), updated_at: Date.now() }).onConflict(['owner_id', 'grant_id', 'idempotency_key']).ignore().returning('id'); + payload_hash: payloadHash, state: 'accepted', lifecycle: 'accepted', accepted_at: acceptedAt, artifacts: JSON.stringify({}), + created_at: acceptedAt, updated_at: acceptedAt }).onConflict(['owner_id', 'grant_id', 'idempotency_key']).ignore().returning('id'); if (!inserted.length) { - const previous = await this.db('mcp_operations').where(identity).first(); + let previous = await this.db('mcp_operations').where(identity).first(); if (!previous || previous.payload_hash !== payloadHash) throw new McpError('IDEMPOTENCY_CONFLICT', 'This key was already used with different arguments.', 409); + await this.reconcileTerminalLifecycles(principal, previous.id); + previous = (await this.db('mcp_operations').where(identity).first())!; return this.project(previous); } try { const result = await invoke(id); - const reported = (result.data as { state?: string })?.state; - const state = reported === 'browser_required' ? reported : result.status === 202 && ['posted', 'queued', 'unknown', 'failed'].includes(reported || '') ? reported! : result.status === 202 ? 'accepted' : 'completed'; - await this.db('mcp_operations').where({ id }).update({ state, result: JSON.stringify(result.data), updated_at: Date.now() }); + const state = operationState(result); + const recorded = await this.db('mcp_operations').where({ id }).whereNotIn('state', ['completed', 'failed', 'cancelled']) + .update({ state, result: JSON.stringify(result.data), updated_at: Date.now() }); + if (!recorded) return this.project((await this.db('mcp_operations').where({ id }).first())!); + const receipt = { state, result: result.data }; + await this.recordArtifacts(id, artifactsFromReceipt({ repository: repository || null }, receipt)); + if (['completed', 'failed', 'cancelled'].includes(state)) { + const failure = errorEnvelope((result.data as { error?: unknown } | null)?.error) ?? failureFromReceipt(receipt); + await this.finish(id, state as LifecycleOutcome, failure); + } else if (state === 'unknown') { + const failure = errorEnvelope((result.data as { error?: unknown } | null)?.error); + await this.db('mcp_operations').where({ id }).whereIn('lifecycle', ['accepted', 'unknown']) + .update({ lifecycle: 'unknown', failure: failure ? JSON.stringify(failure) : null, updated_at: Date.now() }); + } else { + // A live invocation result is authoritative if a concurrent poll had + // already classified its previously result-less receipt as interrupted. + await this.markAccepted(id); + } } catch (error) { // A transport failure can follow an external side effect. Never replay it // automatically or claim it was rolled back. The handle remains durable. const envelope = classifyError(error, { sideEffectsPossible: true }); - await this.db('mcp_operations').where({ id }).update({ state: envelope.code === 'OUTCOME_UNKNOWN' ? 'unknown' : 'failed', - result: JSON.stringify({ error: envelope }), updated_at: Date.now() }); + const state = envelope.code === 'OUTCOME_UNKNOWN' ? 'unknown' : 'failed'; + const recorded = await this.db('mcp_operations').where({ id }).whereNotIn('state', ['completed', 'failed', 'cancelled']) + .update({ state, result: JSON.stringify({ error: envelope }), updated_at: Date.now() }); + if (!recorded) return this.project((await this.db('mcp_operations').where({ id }).first())!); + if (state === 'failed') await this.finish(id, 'failed', envelope); + else await this.db('mcp_operations').where({ id }).whereIn('lifecycle', ['accepted', 'unknown']) + .update({ lifecycle: 'unknown', failure: JSON.stringify(envelope), updated_at: Date.now() }); } return this.project((await this.db('mcp_operations').where({ id }).first())!); } async get(principal: McpPrincipal, id: string): Promise { + await this.reconcileTerminalLifecycles(principal, id); + await this.markInterruptedInvocations(principal, id); const row = await this.db('mcp_operations').where({ id, owner_id: principal.user.id, grant_id: principal.grant.id }).first(); if (!row) throw new McpError('NOT_FOUND', 'Operation not found.', 404); return row; } + /** Repair a process interruption after its terminal receipt write but before lifecycle synchronization. */ + async reconcileTerminalLifecycles(principal: McpPrincipal, id?: string): Promise { + const query = this.db('mcp_operations').where({ owner_id: principal.user.id, grant_id: principal.grant.id }) + .whereIn('state', ['completed', 'failed', 'cancelled']); + if (id) query.andWhere({ id }); + const rows = await query.select(); + for (const row of rows) { + const receipt = recoveryReceipt(row); + const cancellationSourceId = confirmedCancellationSource(row, receipt); + await this.finishCancellationSource(row, cancellationSourceId); + const targetState = receipt.targetState; + const artifacts = artifactsFromReceipt(row, receipt); + const storedArtifacts = json(row.artifacts); + const artifactRecord = storedArtifacts && typeof storedArtifacts === 'object' && !Array.isArray(storedArtifacts) + ? storedArtifacts as Record : {}; + const missingArtifacts = Object.fromEntries(Object.entries(artifacts) + .filter(([key, value]) => canonical(artifactRecord[key]) !== canonical(value))); + const failure = row.state === 'failed' ? failureFromReceipt(receipt) : undefined; + const lifecycleMissing = ['accepted', 'running', 'unknown'].includes(row.lifecycle) || row.finished_at === null; + const progressNeedsRecovery = needsProgressRecovery(targetState, lifecycleMissing, row.progress); + if (!lifecycleMissing && !Object.keys(missingArtifacts).length && !(failure && row.failure === null) && !progressNeedsRecovery) continue; + + const update: Record = { updated_at: Date.now() }; + if (lifecycleMissing) { + update.lifecycle = row.state; + // The receipt update time is the strongest durable evidence of when + // the terminal outcome was persisted; recovery itself must not move it. + update.finished_at = this.db.raw('COALESCE(finished_at, ?, ?)', [row.updated_at, Date.now()]); + } + if (Object.keys(missingArtifacts).length) { + update.artifacts = this.db.raw("json_patch(COALESCE(artifacts, '{}'), ?)", [JSON.stringify(missingArtifacts)]); + } + if (failure && row.failure === null) { + update.failure = this.db.raw('COALESCE(failure, ?)', [JSON.stringify(failure)]); + } + if (progressNeedsRecovery) { + // A terminal tracker receipt is newer than any nonterminal progress + // recorded before lifecycle synchronization was interrupted. + update.progress = recoveredProgress(this.db, targetState, lifecycleMissing); + } + + // Do not attach metadata derived from a receipt that changed after the + // read. A later reconciliation will use the newer durable evidence. + const eligible = this.db('mcp_operations').where({ + id: row.id, owner_id: row.owner_id, grant_id: row.grant_id, + state: row.state, lifecycle: row.lifecycle, + }); + if (row.result === null) eligible.whereNull('result'); + else eligible.andWhere('result', row.result); + await eligible.update(update); + } + } + + async markInterruptedInvocations(principal: McpPrincipal, id?: string): Promise { + const now = Date.now(); + const query = this.db('mcp_operations').where({ owner_id: principal.user.id, grant_id: principal.grant.id }) + .where({ state: 'accepted' }).whereNull('result').whereIn('lifecycle', ['accepted', 'running', 'unknown']) + .where('accepted_at', '<', now - interruptionTimeoutMs); + if (id) query.andWhere({ id }); + await query.update({ state: 'unknown', lifecycle: 'unknown', updated_at: now }); + } + + async markStarted(id: string, at = Date.now()): Promise { + await this.db('mcp_operations').where({ id }).whereIn('lifecycle', ['accepted', 'running', 'unknown']).update({ + lifecycle: 'running', + started_at: this.db.raw('COALESCE(started_at, ?)', [at]), + updated_at: Date.now(), + }); + } + + async markUnknown(id: string): Promise { + await this.db('mcp_operations').where({ id }).whereIn('lifecycle', ['accepted', 'running', 'unknown']).update({ + lifecycle: 'unknown', + updated_at: Date.now(), + }); + } + + async markAccepted(id: string): Promise { + await this.db('mcp_operations').where({ id, lifecycle: 'unknown' }).whereNull('started_at').update({ + lifecycle: 'accepted', + updated_at: Date.now(), + }); + } + + async recordArtifacts(id: string, partial: Record): Promise { + if (!Object.keys(partial).length) return; + await this.db('mcp_operations').where({ id }).update({ + artifacts: this.db.raw("json_patch(COALESCE(artifacts, '{}'), ?)", [JSON.stringify(partial)]), + updated_at: Date.now(), + }); + } + + async recordProgress(id: string, progress: unknown): Promise { + await this.db('mcp_operations').where({ id }) + .whereIn('lifecycle', ['accepted', 'running', 'unknown']) + .whereNotIn('state', ['completed', 'failed', 'cancelled']) + .update({ progress: JSON.stringify(progress), updated_at: Date.now() }); + } + + async finish(id: string, outcome: LifecycleOutcome, failure?: McpErrorEnvelope, progress?: unknown): Promise { + const at = Date.now(); + const eligible = this.db('mcp_operations').where({ id }).andWhere(builder => { + builder.whereIn('lifecycle', ['accepted', 'running', 'unknown']); + if (outcome === 'failed' && failure) builder.orWhere(nested => nested.where({ lifecycle: 'failed' }).whereNull('failure')); + }); + const update: Record = { + lifecycle: outcome, + finished_at: this.db.raw('COALESCE(finished_at, ?)', [at]), + failure: failure ? this.db.raw('COALESCE(failure, ?)', [JSON.stringify(failure)]) : null, + updated_at: at, + }; + if (progress !== undefined) update.progress = this.db.raw(`CASE + WHEN lifecycle IN ('accepted', 'running', 'unknown') THEN ? + ELSE COALESCE(progress, ?) + END`, [JSON.stringify(progress), JSON.stringify(progress)]); + await eligible.update(update); + } + + /** Apply durable cancellation evidence only to the receipt owner's source operation. */ + async finishCancellationSource( + cancellation: Pick, + sourceId: string | undefined, + ): Promise { + if (!sourceId) return; + const durable = await this.db('mcp_operations').where({ + id: cancellation.id, owner_id: cancellation.owner_id, grant_id: cancellation.grant_id, tool: 'cancel_operation', + }).whereIn('state', ['completed', 'failed', 'cancelled']).first('result', 'updated_at'); + const result = record(json(durable?.result)); + if (result?.cancellation !== 'confirmed' || result.operationId !== sourceId) return; + + const now = Date.now(); + const confirmedAt = Number(durable?.updated_at); + await this.db('mcp_operations').where({ + id: sourceId, owner_id: cancellation.owner_id, grant_id: cancellation.grant_id, + }).whereIn('lifecycle', ['accepted', 'running', 'unknown']).update({ + lifecycle: 'cancelled', + finished_at: this.db.raw('COALESCE(finished_at, ?)', [Number.isFinite(confirmedAt) ? confirmedAt : now]), + failure: null, + updated_at: now, + }); + } + project(row: Operation): Record { - const stale = row.state === 'running' && Date.now() - Number(row.updated_at) > 120_000; - const state = stale ? 'unknown' : row.state; - return { operationId: row.id, tool: row.tool, state, result: row.result ? JSON.parse(row.result) : null, + const interrupted = invocationInterrupted(row); + const terminal = ['completed', 'failed', 'cancelled'].includes(row.lifecycle); + const stale = !terminal && interrupted; + const state = terminal ? row.lifecycle : stale ? 'unknown' : row.state; + return { operationId: row.id, tool: row.tool, state, result: json(row.result), lifecycle: { + state: interrupted && !terminal ? 'unknown' : row.lifecycle, + acceptedAt: iso(row.accepted_at), + startedAt: iso(row.started_at), + finishedAt: iso(row.finished_at), + failure: json(row.failure), + artifacts: json(row.artifacts) ?? {}, + progress: json(row.progress), + }, ...(['accepted', 'posted', 'queued', 'running'].includes(state) ? { retryAfterSeconds: 3 } : {}), ...(stale ? { message: 'Execution may have been interrupted. Inspect the target; this action will not be replayed automatically.' } : {}) }; } diff --git a/packages/api/mcp/tools.ts b/packages/api/mcp/tools.ts index a8460c46f..21e47a157 100644 --- a/packages/api/mcp/tools.ts +++ b/packages/api/mcp/tools.ts @@ -1,3 +1,4 @@ +/* eslint-disable max-lines -- MCP tool registration stays centralized so authorization and dispatch remain auditable */ import { assertPlannerCancellationIdentity, cancellationTarget, cancellationOutcome, trackCancellation, trackExecution } from './operationTracking.js'; import { z } from 'zod'; import packageInfo from '../package.json' with { type: 'json' }; @@ -20,6 +21,7 @@ import { createAgentRuntimeRoutes } from '../routes/agentRuntimeRoutes.js'; import { McpError, type McpScope } from './config.js'; import { McpPolicy, type McpPrincipal } from './policy.js'; import { McpOperations, type OperationResult, type Operation } from './operations.js'; +import { syncLifecycle } from './operationLifecycle.js'; import { callWorkflow, type WorkflowHandler } from './adapter.js'; import { addTaskSubmissionTools, trackTaskSubmission } from './toolsTaskSubmissions.js'; import { addPlanningTools } from './toolsPlanning.js'; @@ -296,21 +298,58 @@ export function createToolCatalog(deps: ToolDeps): McpTool[] { workflow(tools, { name: 'delete_task', description: 'Delete an exact inactive task and its persisted execution history. Active tasks must first be cancelled.', scope: 'execute', schema: z.object({ ...taskShape, ...mutationShape }).strict(), target: taskTarget }, tasks.deleteTask, args => ({ params: { taskId: args.taskId }, query: { force: 'false' } })); const operations = new McpOperations(db); - tools.push({ name: 'get_operation', description: 'Read a durable mutation receipt and honest acceptance/completion state. Poll no faster than retryAfterSeconds.', scope: 'read', readOnly: true, schema: z.object({ operationId: z.uuid() }).strict(), run: async ({ principal, args }) => { - const row = await operations.get(principal, args.operationId); - if (row.repository) await policy.repository(principal, row.repository, false, { includeDisabled: row.tool.endsWith('_repository_configuration'), allowUnconfigured: row.tool === 'remove_repository_configuration' }); + const operationResult = (row: Operation): Record => { + if (!row.result) return {}; + try { return JSON.parse(row.result); } catch { return {}; } + }; + const authorizeStoredTool = async ( + row: Operation, + principal: McpPrincipal, + repositories?: Map>, + ): Promise => { + if (row.repository) { + const options = { includeDisabled: row.tool.endsWith('_repository_configuration'), allowUnconfigured: row.tool === 'remove_repository_configuration' }; + const key = `${row.repository}\0${Number(options.includeDisabled)}${Number(options.allowUnconfigured)}`; + let authorization = repositories?.get(key); + if (!authorization) { + authorization = policy.repository(principal, row.repository, false, options); + repositories?.set(key, authorization); + } + await authorization; + } const original = tools.find(tool => tool.name === row.tool); if (original?.permission) policy.requirePermission(principal, original.permission); - if (row.tool === 'cancel_operation' && row.result && JSON.parse(row.result).operationId) { - const source = await operations.get(principal, JSON.parse(row.result).operationId); - if (source.repository) await policy.repository(principal, source.repository); - row.repository = source.repository; - } + }; + const authorizeOperation = async ( + row: Operation, + principal: McpPrincipal, + repositories?: Map>, + ): Promise => { + await authorizeStoredTool(row, principal, repositories); + const sourceId = row.tool === 'cancel_operation' ? operationResult(row).operationId : undefined; + if (typeof sourceId !== 'string') return; + const source = await operations.get(principal, sourceId); + await authorizeStoredTool(source, principal, repositories); + row.repository = source.repository; + }; + tools.push({ name: 'get_operation', description: 'Read a durable mutation receipt and honest lifecycle. "accepted" means the request was recorded and handed to the backend; "running" means execution was observed; the loop/receipt is only "completed" when the backend reached a terminal success state. Poll no faster than retryAfterSeconds.', scope: 'read', readOnly: true, schema: z.object({ operationId: z.uuid() }).strict(), run: async ({ principal, args }) => { + const row = await operations.get(principal, args.operationId); + await authorizeOperation(row, principal); const receipt = operations.project(row); - const result = row.result ? JSON.parse(row.result) : {}; - const continuation = result.continuation || result; + const result = operationResult(row); + const continuation = result.continuation && typeof result.continuation === 'object' && !Array.isArray(result.continuation) + ? result.continuation as Record : result; if (continuation.planId) receipt.targetState = await db('task_drafts').where({ draft_id: continuation.planId, user_id: principal.user.id }).first('status', 'paused', 'mcp_revision'); - if (continuation.goalId) receipt.targetState = await db('goals').where({ goal_id: continuation.goalId, owner_id: principal.user.id }).first('desired_state', 'result_state', 'current_task_id'); + if (continuation.goalId) { + const goal = await db('goals').where({ goal_id: continuation.goalId, owner_id: principal.user.id }) + .first('desired_state', 'result_state', 'current_task_id', 'final_pr_number', 'failure_reason'); + if (goal) { + const currentTask = typeof goal.current_task_id === 'string' + ? await db('task_history').where({ task_id: goal.current_task_id }).orderBy('history_id', 'desc').first('state', 'timestamp', 'reason') + : undefined; + receipt.targetState = { ...goal, ...(currentTask ? { currentTask: { taskId: goal.current_task_id, ...currentTask } } : {}) }; + } + } if (continuation.taskId) receipt.targetState = await db('task_history').where({ task_id: continuation.taskId }).orderBy('history_id', 'desc').first('state', 'timestamp'); updateReceiptState(row, receipt); await trackTaskSubmission(deps, row, principal, receipt); @@ -321,10 +360,58 @@ export function createToolCatalog(deps: ToolDeps): McpTool[] { receipt.targetState = { issues }; if (issues.length === result.issues.length && issues.every(issue => ['under_review', 'merged', 'closed'].includes(issue.status))) receipt.state = 'completed'; } + await syncLifecycle(operations, row, receipt); + receipt.lifecycle = operations.project(await operations.get(principal, row.id)).lifecycle; if (['accepted', 'posted', 'queued', 'running'].includes(String(receipt.state))) receipt.retryAfterSeconds = 3; else delete receipt.retryAfterSeconds; return ok(receipt); } }); + const operationLifecycleSchema = z.enum(['accepted', 'running', 'completed', 'failed', 'cancelled', 'unknown', 'active']); + tools.push({ name: 'list_operations', description: 'List durable mutation receipts newest first without refreshing backend trackers. "accepted" means the request was recorded and handed to the backend; "running" means execution was observed; the loop/receipt is only "completed" when the backend reached a terminal success state. Use refreshWith on an item when a live refresh is needed.', scope: 'read', readOnly: true, + schema: z.object({ + tool: z.string().min(1).max(128).optional().describe('Exact tool name.'), + lifecycle: operationLifecycleSchema.optional(), + sinceMinutes: z.number().int().min(1).max(10080).default(1440), + repository: repositorySchema.optional(), + offset: z.number().int().min(0).max(100000).default(0), + limit: z.number().int().min(1).max(50).default(20), + }).strict(), run: async ({ principal, args }) => { + await operations.reconcileTerminalLifecycles(principal); + await operations.markInterruptedInvocations(principal); + const query = db('mcp_operations').where({ owner_id: principal.user.id, grant_id: principal.grant.id }) + .where('accepted_at', '>=', Date.now() - args.sinceMinutes * 60_000); + if (args.tool) query.where('tool', args.tool); + if (args.repository) query.andWhere(builder => builder.where('repository', args.repository).orWhere('tool', 'cancel_operation')); + if (args.lifecycle === 'active') query.whereIn('lifecycle', ['accepted', 'running']); + else if (args.lifecycle) query.where('lifecycle', args.lifecycle); + const ordered = query.orderBy('accepted_at', 'desc').orderBy('id', 'desc'); + const authorized: Operation[] = []; + const repositoryAuthorizations = new Map>(); + const wanted = args.offset + args.limit + 1; + const batchSize = Math.max(50, args.limit); + let databaseOffset = 0; + while (authorized.length < wanted) { + const rows = await ordered.clone().offset(databaseOffset).limit(batchSize); + if (!rows.length) break; + databaseOffset += rows.length; + for (const row of rows) { + try { + await authorizeOperation(row, principal, repositoryAuthorizations); + if (args.repository && row.repository !== args.repository) continue; + authorized.push(row); + } catch (error) { + // Discovery is a filtered view: current authorization failures do + // not reveal that a matching receipt exists. + if (!(error instanceof McpError) || ![403, 404].includes(error.status)) throw error; + } + if (authorized.length >= wanted) break; + } + if (rows.length < batchSize) break; + } + const page = authorized.slice(args.offset, args.offset + args.limit); + return ok({ operations: page.map(row => ({ ...operations.project(row), refreshWith: 'get_operation' })), + nextOffset: authorized.length > args.offset + args.limit ? args.offset + args.limit : null }); + } }); tools.push({ name: 'cancel_operation', description: 'Request cancellation of an accepted plan generation, goal or task operation. Completed external effects cannot be undone.', scope: 'execute', schema: z.object({ ...mutationShape, operationId: z.uuid() }).strict(), run: async ({ principal, args }) => { const row = await operations.get(principal, args.operationId); if (row.repository) await policy.repository(principal, row.repository, true); diff --git a/packages/api/test/fixtures/mcpCancellation.ts b/packages/api/test/fixtures/mcpCancellation.ts index 0826fca1a..293d6daf4 100644 --- a/packages/api/test/fixtures/mcpCancellation.ts +++ b/packages/api/test/fixtures/mcpCancellation.ts @@ -78,9 +78,11 @@ export async function verifyCancellation({ call, client, principal, deps, agentI await db('task_drafts').where({ draft_id: id }).update({ status: tool === 'generate_plan' ? 'generating' : 'refining', [column]: JSON.stringify({ runId, steps: [] }) }); // This is the persisted 202 boundary produced by the planner start handler. const operationId = randomUUID(); + const acceptedAt = Date.now(); await db('mcp_operations').insert({ id: operationId, owner_id: principal.user.id, grant_id: principal.grant.id, idempotency_key: randomUUID(), tool, repository, payload_hash: 'fixture', state: 'accepted', - result: JSON.stringify({ runId, continuation: { planId: id } }), created_at: Date.now(), updated_at: Date.now() }); + result: JSON.stringify({ runId, continuation: { planId: id } }), lifecycle: 'accepted', accepted_at: acceptedAt, + artifacts: JSON.stringify({}), created_at: acceptedAt, updated_at: acceptedAt }); let enter!: () => void, release!: () => void; const entered = new Promise(resolve => { enter = resolve; }); const hold = new Promise(resolve => { release = resolve; }); @@ -110,9 +112,11 @@ export async function verifyCancellation({ call, client, principal, deps, agentI assert.equal(stale.result.error.code, 'NOT_CANCELLABLE'); for (const outcome of ['completed', 'failed']) { const completedRun = randomUUID(), completedOperation = randomUUID(); + const completedAcceptedAt = Date.now(); await db('mcp_operations').insert({ id: completedOperation, owner_id: principal.user.id, grant_id: principal.grant.id, idempotency_key: randomUUID(), tool, repository, payload_hash: 'fixture', state: 'accepted', - result: JSON.stringify({ runId: completedRun, continuation: { planId: id } }), created_at: Date.now(), updated_at: Date.now() }); + result: JSON.stringify({ runId: completedRun, continuation: { planId: id } }), lifecycle: 'accepted', accepted_at: completedAcceptedAt, + artifacts: JSON.stringify({}), created_at: completedAcceptedAt, updated_at: completedAcceptedAt }); await db('task_drafts').where({ draft_id: id }).update({ status: tool === 'generate_plan' && outcome === 'failed' ? 'failed' : 'review', [column]: JSON.stringify({ runId: completedRun, status: outcome }), diff --git a/packages/api/test/mcpOperations.test.ts b/packages/api/test/mcpOperations.test.ts index 112ce64a1..408662827 100644 --- a/packages/api/test/mcpOperations.test.ts +++ b/packages/api/test/mcpOperations.test.ts @@ -1,14 +1,22 @@ +/* eslint-disable max-lines -- MCP operation and migration regressions share one database fixture. */ import assert from 'node:assert/strict'; -import { test } from 'node:test'; +import { after, test } from 'node:test'; import { mkdtemp, rm } from 'node:fs/promises'; import { tmpdir } from 'node:os'; import path from 'node:path'; import knex from 'knex'; +import { closeConnection } from '@propr/core'; import { up } from '../../core/src/db/migrations/20260910220000_add_mcp.js'; -import { McpOperations } from '../mcp/operations.js'; +import { McpOperations, type Operation } from '../mcp/operations.js'; import { McpError } from '../mcp/config.js'; import { callWorkflow } from '../mcp/adapter.js'; import { parseClientMetadataDocument } from '../mcp/clients.js'; +import { createToolCatalog, type ToolDeps } from '../mcp/tools.js'; +import type { McpPrincipal } from '../mcp/policy.js'; +import { artifactsFromReceipt, failureFromReceipt, syncLifecycle } from '../mcp/operationLifecycle.js'; +import { trackCancellation, trackExecution } from '../mcp/operationTracking.js'; + +after(closeConnection); test('mutation deduplication survives concurrent callers and reopening the SQLite database', async () => { const root = await mkdtemp(path.join(tmpdir(), 'propr-mcp-')); @@ -28,6 +36,9 @@ test('mutation deduplication survives concurrent callers and reopening the SQLit const restarted = new McpOperations(db); const result = await restarted.run(principal, { tool: 'fixture_action', args, repository: 'acme/repo' }, async () => { throw new Error('Must not replay'); }); assert.equal(result.state, 'completed'); assert.deepEqual(result.result, { changed: true }); + assert.equal((result.lifecycle as { state: string }).state, 'completed'); + assert.match((result.lifecycle as { acceptedAt: string }).acceptedAt, /^\d{4}-\d\d-\d\dT/); + assert.ok((result.lifecycle as { finishedAt: string }).finishedAt); await assert.rejects(restarted.run(principal, { tool: 'fixture_action', args: { ...args, value: 'changed' }, repository: 'acme/repo' }, async () => ({ status: 200, data: {} })), /different arguments/); await assert.rejects(restarted.get({ user: { id: '999' }, grant: { id: 'grant-1' } } as never, String(result.operationId)), /not found/); const uncertain = await restarted.run(principal, { tool: 'external_action', args: { idempotencyKey: 'uncertain-key-1' }, repository: 'acme/repo' }, async () => { throw new Error('Network disconnected after possible side effect'); }); @@ -46,6 +57,7 @@ test('mutation deduplication survives concurrent callers and reopening the SQLit } }); const rejected = await restarted.run(principal, { tool: 'guarded_action', args: { idempotencyKey: 'rejected-key-1' }, repository: 'acme/repo' }, async () => { throw new McpError('STALE_HEAD', 'Head changed', 409); }); assert.equal(rejected.state, 'failed'); + assert.deepEqual((rejected.lifecycle as { failure: unknown }).failure, rejected.result && (rejected.result as { error: unknown }).error); const workflowFailure = await restarted.run(principal, { tool: 'workflow_action', args: { idempotencyKey: 'workflow-failure-1' }, repository: 'acme/repo' }, async () => callWorkflow(async (_req, res) => { res.status(500).json({ error: 'The workflow may have changed the target.' }); }, principal, {})); assert.equal(workflowFailure.state, 'unknown'); @@ -60,6 +72,820 @@ test('mutation deduplication survives concurrent callers and reopening the SQLit } finally { await db.destroy(); await rm(root, { recursive: true, force: true }); } }); +test('operation lifecycle transitions, artifacts and progress are durable and monotonic', async t => { + const db = knex({ client: 'better-sqlite3', connection: { filename: ':memory:' }, useNullAsDefault: true }); + t.after(() => db.destroy()); + await db.schema.createTable('task_drafts', table => table.string('draft_id').primary()); + await up(db); + const operations = new McpOperations(db); + const principal = { user: { id: 'alice' }, grant: { id: 'grant-a' } } as never; + const receipt = await operations.run(principal, { + tool: 'run_ultrafix', args: { idempotencyKey: 'lifecycle-key-1' }, repository: 'acme/repo', + }, async () => ({ status: 202, data: { state: 'queued', pullRequest: 42, commentId: 99 } })); + assert.equal((receipt.lifecycle as { state: string }).state, 'accepted'); + assert.equal((receipt.lifecycle as { startedAt: unknown }).startedAt, null); + assert.deepEqual((receipt.lifecycle as { artifacts: unknown }).artifacts, { + pullRequest: { repository: 'acme/repo', number: 42, url: 'https://github.com/acme/repo/pull/42' }, commentId: 99, + }); + const replay = await operations.run(principal, { + tool: 'run_ultrafix', args: { idempotencyKey: 'lifecycle-key-1' }, repository: 'acme/repo', + }, async () => { throw new Error('Must not invoke on replay'); }); + assert.deepEqual((replay.lifecycle as { artifacts: unknown }).artifacts, (receipt.lifecycle as { artifacts: unknown }).artifacts); + + const id = String(receipt.operationId); + await operations.markStarted(id, 1_800_000_000_000); + await Promise.all([ + operations.recordArtifacts(id, { taskId: 'task-1' }), + operations.recordArtifacts(id, { pullRequest: { repository: 'acme/repo', number: 42, url: 'https://github.com/acme/repo/pull/42' } }), + ]); + await operations.recordProgress(id, { taskId: 'task-1', state: 'processing' }); + await Promise.all([operations.finish(id, 'completed'), operations.markStarted(id, 1_700_000_000_000)]); + await operations.finish(id, 'failed', { code: 'LATE_FAILURE', message: 'stale', stage: null, retryable: false, status: 500 }); + + const projected = operations.project(await operations.get(principal, id)); + assert.equal((projected.lifecycle as { state: string }).state, 'completed'); + assert.equal((projected.lifecycle as { startedAt: string }).startedAt, '2027-01-15T08:00:00.000Z'); + assert.ok((projected.lifecycle as { finishedAt: string }).finishedAt); + assert.deepEqual((projected.lifecycle as { artifacts: unknown }).artifacts, { + taskId: 'task-1', pullRequest: { repository: 'acme/repo', number: 42, url: 'https://github.com/acme/repo/pull/42' }, commentId: 99, + }); + assert.deepEqual((projected.lifecycle as { progress: unknown }).progress, { taskId: 'task-1', state: 'processing' }); + assert.equal((projected.lifecycle as { failure: unknown }).failure, null); +}); + +test('a delayed invocation result cannot replace terminal evidence persisted by a poll', async t => { + const db = knex({ client: 'better-sqlite3', connection: { filename: ':memory:' }, useNullAsDefault: true }); + t.after(() => db.destroy()); + await db.schema.createTable('task_drafts', table => table.string('draft_id').primary()); + await up(db); + const operations = new McpOperations(db); + const principal = { user: { id: 'alice' }, grant: { id: 'grant-a' } } as McpPrincipal; + let invocationStarted!: () => void; + let finishInvocation!: (result: { status: number; data: unknown }) => void; + const started = new Promise(resolve => { invocationStarted = resolve; }); + const pending = operations.run(principal, { + tool: 'review_pull_request', args: { idempotencyKey: 'late-invocation-result' }, repository: 'acme/repo', + }, async () => { + invocationStarted(); + return new Promise(resolve => { finishInvocation = resolve; }); + }); + await started; + + const row = (await db('mcp_operations').where({ idempotency_key: 'late-invocation-result' }).first())!; + const targetState = { taskId: 'review-task', state: 'completed', timestamp: '2026-09-29T04:01:00.000Z' }; + await db('mcp_operations').where({ id: row.id }).update({ + state: 'completed', result: JSON.stringify({ executionResolved: true, targetState }), updated_at: Date.now(), + }); + await operations.finish(row.id, 'completed', undefined, targetState); + finishInvocation({ status: 202, data: { state: 'queued', pullRequest: 42 } }); + + const completed = await pending; + assert.equal(completed.state, 'completed'); + assert.deepEqual(completed.result, { executionResolved: true, targetState }); + assert.deepEqual((completed.lifecycle as { progress: unknown }).progress, targetState); +}); + +test('stale concurrent polls cannot replace terminal tracker receipts or lifecycle progress', async t => { + const db = knex({ client: 'better-sqlite3', connection: { filename: ':memory:' }, useNullAsDefault: true }); + t.after(() => db.destroy()); + await db.schema.createTable('task_drafts', table => table.string('draft_id').primary()); + await db.schema.createTable('tasks', table => { + table.string('task_id').primary(); table.string('repository'); table.string('task_type'); + table.integer('issue_number'); table.integer('pr_number'); table.text('initial_job_data'); + table.timestamp('created_at').defaultTo(db.fn.now()); + }); + await db.schema.createTable('task_history', table => { + table.increments('history_id').primary(); table.string('task_id'); table.string('state'); + table.timestamp('timestamp'); table.text('reason'); table.text('metadata'); + }); + await up(db); + + const operations = new McpOperations(db); + const args = { idempotencyKey: 'stale-terminal-poll-1' }; + let releaseStale!: () => void; + let staleReachedRefresh!: () => void; + const release = new Promise(resolve => { releaseStale = resolve; }); + const reachedRefresh = new Promise(resolve => { staleReachedRefresh = resolve; }); + let githubRequests = 0; + const principal = { + user: { id: 'alice' }, grant: { id: 'grant-a' }, github: { request: async () => { + githubRequests++; + if (githubRequests === 1) { staleReachedRefresh(); await release; } + return { data: { head: { sha: `head-${githubRequests}` } } }; + } }, + } as unknown as McpPrincipal; + const receipt = await operations.run(principal, { tool: 'review_pull_request', args, repository: 'acme/repo' }, async () => ({ + status: 202, data: { state: 'queued', pullRequest: 42, commentId: 99 }, + })); + await db('tasks').insert({ task_id: 'review-task', repository: 'acme/repo', task_type: 'issue', issue_number: 42, + initial_job_data: JSON.stringify({ commandCommentId: 99, commandMode: 'review' }) }); + await db('task_history').insert({ task_id: 'review-task', state: 'processing', + timestamp: '2026-09-29T04:00:00.000Z', reason: 'Review is running' }); + const deps = { db, policy: { repository: async () => {}, requirePermission: () => {}, config: {} } as never, + taskQueue: {} as never, redisClient: {} as never, runtimeBuildQueue: {} as never } as ToolDeps; + + const staleRow = (await db('mcp_operations').where({ id: receipt.operationId }).first())!; + const staleReceipt = operations.project(staleRow); + const stalePoll = trackExecution(deps, staleRow, principal, staleReceipt); + await reachedRefresh; + + const terminalAt = '2026-09-29T04:01:00.000Z'; + await db('task_history').insert({ task_id: 'review-task', state: 'completed', timestamp: terminalAt, + reason: 'Review processing completed successfully' }); + const terminalRow = (await db('mcp_operations').where({ id: receipt.operationId }).first())!; + const terminalReceipt = operations.project(terminalRow); + await trackExecution(deps, terminalRow, principal, terminalReceipt); + await syncLifecycle(operations, terminalRow, terminalReceipt); + + releaseStale(); + await stalePoll; + await syncLifecycle(operations, staleRow, staleReceipt); + + const durable = (await db('mcp_operations').where({ id: receipt.operationId }).first())!; + const durableResult = JSON.parse(durable.result!); + assert.equal(durable.state, 'completed'); + assert.equal(durable.lifecycle, 'completed'); + assert.equal(durableResult.executionResolved, true); + assert.equal(durableResult.targetState.state, 'completed'); + assert.equal(durableResult.targetState.reason, 'Review processing completed successfully'); + assert.deepEqual(JSON.parse(durable.progress!), durableResult.targetState); + assert.equal(staleReceipt.state, 'completed', 'the stale caller returns the winning durable observation'); + assert.deepEqual(staleReceipt.targetState, durableResult.targetState); + + await db('task_history').delete(); + const replayed = await operations.replay(principal, 'review_pull_request', args); + assert.equal((replayed?.lifecycle as { state: string }).state, 'completed'); + assert.deepEqual((replayed?.lifecycle as { progress: unknown }).progress, durableResult.targetState); + const list = createToolCatalog(deps).find(tool => tool.name === 'list_operations')!; + const listed = (await list.run({ principal, args: list.schema.parse({}) })).data as { operations: Array> }; + const listedReceipt = listed.operations.find(operation => operation.operationId === receipt.operationId)!; + assert.deepEqual((listedReceipt.lifecycle as { progress: unknown }).progress, durableResult.targetState); +}); + +test('terminal execution restoration preserves target state resolved by get_operation', async t => { + const db = knex({ client: 'better-sqlite3', connection: { filename: ':memory:' }, useNullAsDefault: true }); + t.after(() => db.destroy()); + await db.schema.createTable('task_drafts', table => table.string('draft_id').primary()); + await db.schema.createTable('tasks', table => { + table.string('task_id').primary(); table.string('job_id'); table.string('repository'); table.string('task_type'); + table.integer('pr_number'); table.text('initial_job_data'); table.timestamp('created_at').defaultTo(db.fn.now()); + }); + await db.schema.createTable('task_history', table => { + table.increments('history_id').primary(); table.string('task_id'); table.string('state'); table.timestamp('timestamp'); + }); + await up(db); + + const principal = { user: { id: 'alice' }, grant: { id: 'grant-a' } } as McpPrincipal; + const operations = new McpOperations(db); + const receipt = await operations.run(principal, { + tool: 'send_task_followup', args: { idempotencyKey: 'resolved-target-fallback' }, repository: 'acme/repo', + }, async () => ({ status: 202, data: { state: 'queued', continuation: { taskId: 'resolved-task' } } })); + const observedAt = '2026-09-29T05:00:00.000Z'; + await db('tasks').insert({ task_id: 'resolved-task', job_id: 'resolved-job', repository: 'acme/repo', task_type: 'issue', pr_number: 81 }); + await db('task_history').insert({ task_id: 'resolved-task', state: 'completed', timestamp: observedAt }); + await db('mcp_operations').where({ id: receipt.operationId }).update({ + state: 'completed', lifecycle: 'completed', finished_at: Date.now(), + result: JSON.stringify({ jobId: 'resolved-job', continuation: { taskId: 'resolved-task' }, executionResolved: true }), + }); + + const deps = { db, policy: { repository: async () => {}, requirePermission: () => {}, config: {} } as never, + taskQueue: {} as never, redisClient: {} as never, runtimeBuildQueue: {} as never } as ToolDeps; + const get = createToolCatalog(deps).find(tool => tool.name === 'get_operation')!; + const restored = (await get.run({ principal, args: get.schema.parse({ operationId: receipt.operationId }) })).data as Record; + assert.deepEqual(restored.targetState, { + state: 'completed', timestamp: observedAt, taskId: 'resolved-task', pr_number: 81, + }); +}); + +test('replay recovers terminal lifecycle, artifacts and failure from durable receipts', async t => { + const db = knex({ client: 'better-sqlite3', connection: { filename: ':memory:' }, useNullAsDefault: true }); + t.after(() => db.destroy()); + await db.schema.createTable('task_drafts', table => table.string('draft_id').primary()); + await up(db); + const principal = { user: { id: 'alice' }, grant: { id: 'grant-a' } } as McpPrincipal; + const operations = new McpOperations(db); + const args = { idempotencyKey: 'interrupted-success-1' }; + const receipt = await operations.run(principal, { tool: 'fixture_action', args, repository: 'acme/repo' }, + async () => ({ status: 202, data: { state: 'queued' } })); + const persistedAt = Date.now() - 500; + await db('mcp_operations').where({ id: receipt.operationId }).update({ + state: 'completed', result: JSON.stringify({ changed: true, taskId: 'task-recovered' }), updated_at: persistedAt, + }); + assert.deepEqual(await db('mcp_operations').where({ id: receipt.operationId }).first('state', 'lifecycle', 'finished_at'), { + state: 'completed', lifecycle: 'accepted', finished_at: null, + }); + + const restarted = new McpOperations(db); + const replay = await restarted.replay(principal, 'fixture_action', args); + assert.equal(replay?.state, 'completed'); + assert.deepEqual(replay?.result, { changed: true, taskId: 'task-recovered' }); + assert.deepEqual(replay?.lifecycle, { + state: 'completed', acceptedAt: (receipt.lifecycle as { acceptedAt: string }).acceptedAt, + startedAt: null, finishedAt: new Date(persistedAt).toISOString(), failure: null, + artifacts: { taskId: 'task-recovered' }, progress: null, + }); + + // Repeat entry points also repair rows left terminal by the older + // lifecycle-only reconciliation. + await db('mcp_operations').where({ id: receipt.operationId }).update({ artifacts: JSON.stringify({}) }); + let invoked = false; + const duplicate = await restarted.run(principal, { tool: 'fixture_action', args, repository: 'acme/repo' }, async () => { + invoked = true; + return { status: 200, data: {} }; + }); + assert.equal(invoked, false); + assert.equal((duplicate.lifecycle as { state: string }).state, 'completed'); + assert.deepEqual((duplicate.lifecycle as { artifacts: unknown }).artifacts, { taskId: 'task-recovered' }); + + const failedArgs = { idempotencyKey: 'interrupted-failure-1' }; + const failed = await operations.run(principal, { tool: 'fixture_action', args: failedArgs, repository: 'acme/repo' }, + async () => ({ status: 202, data: { state: 'queued' } })); + const failure = { code: 'WORKFLOW_FAILED', message: 'The durable workflow failed.', stage: 'internal', retryable: false, status: 500 }; + await db('mcp_operations').where({ id: failed.operationId }).update({ + state: 'failed', result: JSON.stringify({ error: failure }), updated_at: persistedAt, + }); + const failedReplay = await restarted.replay(principal, 'fixture_action', failedArgs); + assert.equal((failedReplay?.lifecycle as { state: string }).state, 'failed'); + assert.deepEqual((failedReplay?.lifecycle as { failure: unknown }).failure, failure); +}); + +test('list_operations recovers tracker lifecycle, artifacts and failure before filtering', async t => { + const db = knex({ client: 'better-sqlite3', connection: { filename: ':memory:' }, useNullAsDefault: true }); + t.after(() => db.destroy()); + await db.schema.createTable('task_drafts', table => table.string('draft_id').primary()); + await db.schema.createTable('tasks', table => { + table.string('task_id').primary(); table.string('job_id'); table.string('repository'); table.string('task_type'); + table.integer('pr_number'); table.text('initial_job_data'); table.timestamp('created_at').defaultTo(db.fn.now()); + }); + await db.schema.createTable('task_history', table => { + table.increments('history_id').primary(); table.string('task_id'); table.string('state'); + table.timestamp('timestamp'); table.text('reason'); table.text('metadata'); + }); + await up(db); + const principal = { user: { id: 'alice' }, grant: { id: 'grant-a' } } as McpPrincipal; + const operations = new McpOperations(db); + const failedArgs = { idempotencyKey: 'tracker-failure-02' }; + const failed = await operations.run(principal, { + tool: 'send_task_followup', args: failedArgs, repository: 'acme/repo', + }, async () => ({ status: 202, data: { state: 'queued', jobId: 'job-failed', continuation: { taskId: 'task-failed' } } })); + const deps = { db, policy: { repository: async () => {}, requirePermission: () => {}, config: {} } as never, + taskQueue: {} as never, redisClient: {} as never, runtimeBuildQueue: {} as never } as ToolDeps; + const observedAt = '2026-09-29T03:00:00.000Z'; + await db('tasks').insert({ task_id: 'task-failed', job_id: 'job-failed', repository: 'acme/repo', task_type: 'issue', pr_number: 73 }); + await db('task_history').insert({ task_id: 'task-failed', state: 'failed', timestamp: observedAt, reason: 'Agent stopped after tests failed.' }); + await operations.recordProgress(String(failed.operationId), { taskId: 'task-failed', state: 'processing' }); + const failedRow = await db('mcp_operations').where({ id: failed.operationId }).first(); + // Stop immediately after the awaited tracker write, before syncLifecycle can + // project its nested targetState into lifecycle columns. + await trackExecution(deps, failedRow!, principal, operations.project(failedRow!)); + const trackerWrite = await db('mcp_operations').where({ id: failed.operationId }).first(); + assert.equal(trackerWrite?.state, 'failed'); + assert.equal(trackerWrite?.lifecycle, 'accepted'); + assert.equal(trackerWrite?.failure, null); + assert.deepEqual(JSON.parse(trackerWrite!.progress!), { taskId: 'task-failed', state: 'processing' }); + assert.deepEqual(JSON.parse(trackerWrite!.result!).targetState, { + taskId: 'task-failed', pr_number: 73, state: 'failed', timestamp: observedAt, reason: 'Agent stopped after tests failed.', + }); + + const list = createToolCatalog(deps).find(tool => tool.name === 'list_operations')!; + const active = (await list.run({ principal, args: list.schema.parse({ lifecycle: 'active' }) })).data as { operations: Array> }; + const failures = (await list.run({ principal, args: list.schema.parse({ lifecycle: 'failed' }) })).data as { operations: Array> }; + assert.deepEqual(active.operations, []); + assert.deepEqual(failures.operations.map(row => row.operationId), [failed.operationId]); + const recoveredLifecycle = failures.operations[0].lifecycle as Record; + assert.deepEqual(recoveredLifecycle.failure, { + code: 'EXECUTION_FAILED', message: 'Agent stopped after tests failed.', stage: 'internal', retryable: false, status: 500, + }); + assert.deepEqual(recoveredLifecycle.artifacts, { + taskId: 'task-failed', + pullRequest: { repository: 'acme/repo', number: 73, url: 'https://github.com/acme/repo/pull/73' }, + }); + assert.deepEqual(recoveredLifecycle.progress, { + taskId: 'task-failed', pr_number: 73, state: 'failed', timestamp: observedAt, reason: 'Agent stopped after tests failed.', + }); + const replay = await operations.replay(principal, 'send_task_followup', failedArgs); + assert.deepEqual(replay?.lifecycle, recoveredLifecycle); +}); + +test('replay and listing recover confirmed cancellation propagation after the tracker write', async t => { + const db = knex({ client: 'better-sqlite3', connection: { filename: ':memory:' }, useNullAsDefault: true }); + t.after(() => db.destroy()); + await db.schema.createTable('task_drafts', table => table.string('draft_id').primary()); + await db.schema.createTable('tasks', table => { + table.string('task_id').primary(); table.string('repository'); table.string('task_type'); + }); + await db.schema.createTable('goals', table => { + table.string('goal_id').primary(); table.string('current_task_id'); + }); + await db.schema.createTable('task_history', table => { + table.increments('history_id').primary(); table.string('task_id'); table.string('state'); table.timestamp('timestamp'); + }); + await up(db); + const principal = { user: { id: 'alice' }, grant: { id: 'grant-a' } } as McpPrincipal; + const otherGrant = { user: { id: 'alice' }, grant: { id: 'grant-b' } } as McpPrincipal; + const operations = new McpOperations(db); + const deps = { db, policy: { repository: async () => {}, requirePermission: () => {}, config: {} } as never, + taskQueue: {} as never, redisClient: {} as never, runtimeBuildQueue: {} as never } as ToolDeps; + const createInterruptedConfirmation = async (suffix: string) => { + const taskId = `cancel-task-${suffix}`; + const source = await operations.run(principal, { + tool: 'review_pull_request', args: { idempotencyKey: `cancel-source-${suffix}` }, repository: 'acme/repo', + }, async () => ({ status: 202, data: { state: 'queued', continuation: { taskId } } })); + const args = { idempotencyKey: `cancel-receipt-${suffix}` }; + const cancellation = await operations.run(principal, { tool: 'cancel_operation', args, repository: 'acme/repo' }, async () => ({ + status: 202, data: { operationId: source.operationId, cancellation: 'requested', continuation: { taskId }, targetTool: 'review_pull_request' }, + })); + await db('tasks').insert({ task_id: taskId, repository: 'acme/repo', task_type: 'issue' }); + await db('task_history').insert({ task_id: taskId, state: 'cancelled', timestamp: new Date() }); + const row = (await db('mcp_operations').where({ id: cancellation.operationId }).first())!; + await trackCancellation(deps, row, principal, operations.project(row)); + const persisted = (await db('mcp_operations').where({ id: cancellation.operationId }).first())!; + assert.equal(persisted.state, 'completed'); + assert.equal(persisted.lifecycle, 'accepted'); + assert.equal((await db('mcp_operations').where({ id: source.operationId }).first())!.lifecycle, 'accepted'); + return { args, cancellation, persisted, source }; + }; + + const replayCase = await createInterruptedConfirmation('replay-01'); + const replayed = await operations.replay(principal, 'cancel_operation', replayCase.args); + assert.equal((replayed?.lifecycle as { state: string }).state, 'completed'); + const replayedSource = (await db('mcp_operations').where({ id: replayCase.source.operationId }).first())!; + assert.equal(replayedSource.lifecycle, 'cancelled'); + assert.equal(replayedSource.finished_at, replayCase.persisted.updated_at); + + const listCase = await createInterruptedConfirmation('listing-1'); + const foreignSource = await operations.run(otherGrant, { + tool: 'review_pull_request', args: { idempotencyKey: 'foreign-cancel-source' }, repository: 'acme/repo', + }, async () => ({ status: 202, data: { state: 'queued' } })); + await operations.run(principal, { + tool: 'cancel_operation', args: { idempotencyKey: 'foreign-cancel-proof' }, + }, async () => ({ status: 200, data: { operationId: foreignSource.operationId, cancellation: 'confirmed', executionResolved: true } })); + const completedSource = await operations.run(principal, { + tool: 'review_pull_request', args: { idempotencyKey: 'completed-cancel-source' }, repository: 'acme/repo', + }, async () => ({ status: 200, data: { changed: true } })); + await operations.run(principal, { + tool: 'cancel_operation', args: { idempotencyKey: 'completed-cancel-proof' }, repository: 'acme/repo', + }, async () => ({ status: 200, data: { operationId: completedSource.operationId, cancellation: 'confirmed', executionResolved: true } })); + + const list = createToolCatalog(deps).find(tool => tool.name === 'list_operations')!; + const listed = (await list.run({ principal, args: list.schema.parse({}) })).data as { operations: Array> }; + assert.ok(listed.operations.some(row => row.operationId === listCase.cancellation.operationId)); + const listedSource = (await db('mcp_operations').where({ id: listCase.source.operationId }).first())!; + assert.equal(listedSource.lifecycle, 'cancelled'); + assert.equal(listedSource.finished_at, listCase.persisted.updated_at); + assert.equal((await db('mcp_operations').where({ id: foreignSource.operationId }).first())!.lifecycle, 'accepted'); + assert.equal((await db('mcp_operations').where({ id: completedSource.operationId }).first())!.lifecycle, 'completed'); +}); + +test('polling an accepted wrapper does not fabricate backend execution start', async t => { + const db = knex({ client: 'better-sqlite3', connection: { filename: ':memory:' }, useNullAsDefault: true }); + t.after(() => db.destroy()); + await db.schema.createTable('task_drafts', table => table.string('draft_id').primary()); + await up(db); + const operations = new McpOperations(db); + const principal = { user: { id: 'alice' }, grant: { id: 'grant-a' } } as McpPrincipal; + let finishInvocation!: (result: { status: number; data: unknown }) => void; + let invocationStarted!: () => void; + const started = new Promise(resolve => { invocationStarted = resolve; }); + const pending = operations.run(principal, { + tool: 'review_pull_request', args: { idempotencyKey: 'pending-wrapper-1' }, repository: 'acme/repo', + }, async () => { + invocationStarted(); + return new Promise(resolve => { finishInvocation = resolve; }); + }); + await started; + + const row = await db('mcp_operations').first() as Operation; + const polled = operations.project(row); + assert.equal(polled.state, 'accepted'); + await syncLifecycle(operations, row, polled); + const unchanged = operations.project(await operations.get(principal, row.id)); + assert.equal((unchanged.lifecycle as { state: string }).state, 'accepted'); + assert.equal((unchanged.lifecycle as { startedAt: unknown }).startedAt, null); + + finishInvocation({ status: 202, data: { state: 'queued' } }); + const queued = await pending; + assert.equal(queued.state, 'queued'); + assert.equal((queued.lifecycle as { state: string }).state, 'accepted'); + assert.equal((queued.lifecycle as { startedAt: unknown }).startedAt, null); +}); + +test('interrupted accepted invocations become durable unknown while acknowledged queues remain active', async t => { + const db = knex({ client: 'better-sqlite3', connection: { filename: ':memory:' }, useNullAsDefault: true }); + t.after(() => db.destroy()); + await db.schema.createTable('task_drafts', table => table.string('draft_id').primary()); + await up(db); + const operations = new McpOperations(db); + const principal = { user: { id: 'alice' }, grant: { id: 'grant-a' } } as McpPrincipal; + let finishInvocation!: (result: { status: number; data: unknown }) => void; + let invocationStarted!: () => void; + const started = new Promise(resolve => { invocationStarted = resolve; }); + const args = { idempotencyKey: 'interrupted-wrapper-1' }; + const pending = operations.run(principal, { tool: 'review_pull_request', args, repository: 'acme/repo' }, async () => { + invocationStarted(); + return new Promise(resolve => { finishInvocation = resolve; }); + }); + await started; + + const interrupted = await db('mcp_operations').where({ idempotency_key: args.idempotencyKey }).first(); + const invokedAt = Date.now() - 120_001; + await db('mcp_operations').where({ id: interrupted!.id }).update({ accepted_at: invokedAt, created_at: invokedAt, updated_at: Date.now() }); + const projected = operations.project((await db('mcp_operations').where({ id: interrupted!.id }).first())!); + assert.equal(projected.state, 'unknown'); + assert.equal((projected.lifecycle as { state: string }).state, 'unknown'); + assert.equal(projected.retryAfterSeconds, undefined); + assert.match(String(projected.message), /may have been interrupted/); + + const queued = await operations.run(principal, { + tool: 'review_pull_request', args: { idempotencyKey: 'acknowledged-queue-1' }, repository: 'acme/repo', + }, async () => ({ status: 202, data: { state: 'queued', commentId: 42 } })); + await db('mcp_operations').where({ id: queued.operationId }).update({ accepted_at: invokedAt, created_at: invokedAt, updated_at: Date.now() }); + + const deps = { db, policy: { repository: async () => {}, requirePermission: () => {}, config: {} } as never, + taskQueue: {} as never, redisClient: {} as never, runtimeBuildQueue: {} as never } as ToolDeps; + const catalog = createToolCatalog(deps); + const list = catalog.find(tool => tool.name === 'list_operations')!; + const get = catalog.find(tool => tool.name === 'get_operation')!; + const active = (await list.run({ principal, args: list.schema.parse({ lifecycle: 'active' }) })).data as { operations: Array> }; + const unknown = (await list.run({ principal, args: list.schema.parse({ lifecycle: 'unknown' }) })).data as { operations: Array> }; + assert.deepEqual(active.operations.map(row => row.operationId), [queued.operationId]); + assert.deepEqual(unknown.operations.map(row => row.operationId), [interrupted!.id]); + + const refreshed = (await get.run({ principal, args: get.schema.parse({ operationId: interrupted!.id }) })).data as Record; + assert.equal(refreshed.state, 'unknown'); + assert.equal((refreshed.lifecycle as { state: string }).state, 'unknown'); + let replayed = false; + const replay = await operations.run(principal, { tool: 'review_pull_request', args, repository: 'acme/repo' }, async () => { + replayed = true; + return { status: 200, data: {} }; + }); + assert.equal(replayed, false); + assert.equal(replay.state, 'unknown'); + assert.equal((replay.lifecycle as { state: string }).state, 'unknown'); + + const queuedRow = await operations.get(principal, String(queued.operationId)); + assert.equal(queuedRow.state, 'queued'); + assert.equal(queuedRow.lifecycle, 'accepted'); + + // If the original process is merely slow rather than gone, its fresh result + // remains authoritative and resolves the unavoidable timeout race. + finishInvocation({ status: 202, data: { state: 'queued', commentId: 99 } }); + const resolved = await pending; + assert.equal(resolved.state, 'queued'); + assert.equal((resolved.lifecycle as { state: string }).state, 'accepted'); +}); + +test('tracker-observed running operations remain active when their projection timestamp ages', async t => { + const db = knex({ client: 'better-sqlite3', connection: { filename: ':memory:' }, useNullAsDefault: true }); + t.after(() => db.destroy()); + await db.schema.createTable('task_drafts', table => table.string('draft_id').primary()); + await db.schema.createTable('tasks', table => { + table.string('task_id').primary(); table.string('job_id'); table.string('repository'); table.string('task_type'); + table.integer('pr_number'); table.text('initial_job_data'); table.timestamp('created_at').defaultTo(db.fn.now()); + }); + await db.schema.createTable('task_history', table => { + table.increments('history_id').primary(); table.string('task_id'); table.string('state'); + table.timestamp('timestamp'); table.text('reason'); table.text('metadata'); + }); + await up(db); + + const operations = new McpOperations(db); + const principal = { user: { id: 'alice' }, grant: { id: 'grant-a' } } as McpPrincipal; + const args = { idempotencyKey: 'long-running-tracker-1' }; + const receipt = await operations.run(principal, { + tool: 'send_task_followup', args, repository: 'acme/repo', + }, async () => ({ status: 202, data: { + state: 'queued', jobId: 'long-running-job', continuation: { taskId: 'long-running-task' }, + } })); + await db('tasks').insert({ + task_id: 'long-running-task', job_id: 'long-running-job', repository: 'acme/repo', task_type: 'issue', + initial_job_data: JSON.stringify({}), + }); + await db('task_history').insert({ + task_id: 'long-running-task', state: 'processing', timestamp: new Date(), reason: 'Task is still running', + }); + const deps = { db, policy: { repository: async () => {}, requirePermission: () => {}, config: {} } as never, + taskQueue: {} as never, redisClient: {} as never, runtimeBuildQueue: {} as never } as ToolDeps; + const catalog = createToolCatalog(deps); + const get = catalog.find(tool => tool.name === 'get_operation')!; + const list = catalog.find(tool => tool.name === 'list_operations')!; + + const polled = (await get.run({ + principal, args: get.schema.parse({ operationId: receipt.operationId }), + })).data as Record; + assert.equal(polled.state, 'running'); + assert.equal((polled.lifecycle as { state: string }).state, 'running'); + + await db('mcp_operations').where({ id: receipt.operationId }).update({ updated_at: Date.now() - 120_001 }); + const assertRunning = (projected: Record | undefined) => { + assert.equal(projected?.state, 'running'); + assert.equal((projected?.lifecycle as { state: string }).state, 'running'); + assert.equal(projected?.retryAfterSeconds, 3); + assert.equal(projected?.message, undefined); + }; + + const active = (await list.run({ + principal, args: list.schema.parse({ lifecycle: 'active' }), + })).data as { operations: Array> }; + assertRunning(active.operations.find(operation => operation.operationId === receipt.operationId)); + assertRunning(await operations.replay(principal, 'send_task_followup', args)); + + let invoked = false; + const duplicate = await operations.run(principal, { + tool: 'send_task_followup', args, repository: 'acme/repo', + }, async () => { + invoked = true; + return { status: 200, data: {} }; + }); + assert.equal(invoked, false); + assertRunning(duplicate); +}); + +test('tracker uncertainty resolves from later evidence and preserves observed timestamps and terminal outcomes', async t => { + const db = knex({ client: 'better-sqlite3', connection: { filename: ':memory:' }, useNullAsDefault: true }); + t.after(() => db.destroy()); + await db.schema.createTable('task_drafts', table => table.string('draft_id').primary()); + await up(db); + const operations = new McpOperations(db); + const principal = { user: { id: 'alice' }, grant: { id: 'grant-a' } } as McpPrincipal; + const receipt = await operations.run(principal, { + tool: 'send_task_followup', args: { idempotencyKey: 'tracker-unknown-1' }, repository: 'acme/repo', + }, async () => ({ status: 202, data: { state: 'queued' } })); + const id = String(receipt.operationId); + const row = await operations.get(principal, id); + + await syncLifecycle(operations, row, { ...receipt, state: 'unknown' }); + assert.equal((operations.project(await operations.get(principal, id)).lifecycle as { state: string }).state, 'unknown'); + const deps = { db, policy: { repository: async () => {}, requirePermission: () => {}, config: {} } as never, + taskQueue: {} as never, redisClient: {} as never, runtimeBuildQueue: {} as never } as ToolDeps; + const list = createToolCatalog(deps).find(tool => tool.name === 'list_operations')!; + const unknown = (await list.run({ principal, args: list.schema.parse({ lifecycle: 'unknown' }) })).data as { operations: unknown[] }; + const active = (await list.run({ principal, args: list.schema.parse({ lifecycle: 'active' }) })).data as { operations: unknown[] }; + assert.equal(unknown.operations.length, 1); + assert.deepEqual(active.operations, []); + await syncLifecycle(operations, row, { ...receipt, state: 'queued', targetState: { taskId: 'task-1', state: 'pending' } }); + assert.equal((operations.project(await operations.get(principal, id)).lifecycle as { state: string }).state, 'accepted'); + + const observedAt = 1_800_000_000_000; + await syncLifecycle(operations, row, { ...receipt, state: 'running', targetState: { taskId: 'task-1', state: 'processing', timestamp: observedAt } }); + await syncLifecycle(operations, row, { ...receipt, state: 'unknown' }); + let lifecycle = operations.project(await operations.get(principal, id)).lifecycle as { state: string; startedAt: string | null }; + assert.equal(lifecycle.state, 'unknown'); + assert.equal(lifecycle.startedAt, '2027-01-15T08:00:00.000Z'); + await syncLifecycle(operations, row, { ...receipt, state: 'running', targetState: { taskId: 'task-1', state: 'processing' } }); + lifecycle = operations.project(await operations.get(principal, id)).lifecycle as typeof lifecycle; + assert.equal(lifecycle.state, 'running'); + assert.equal(lifecycle.startedAt, '2027-01-15T08:00:00.000Z'); + await operations.finish(id, 'completed'); + await syncLifecycle(operations, row, { ...receipt, state: 'unknown' }); + assert.equal((operations.project(await operations.get(principal, id)).lifecycle as { state: string }).state, 'completed'); +}); + +test('tracker task and review failures populate and can enrich the durable failure envelope', async t => { + const db = knex({ client: 'better-sqlite3', connection: { filename: ':memory:' }, useNullAsDefault: true }); + t.after(() => db.destroy()); + await db.schema.createTable('task_drafts', table => table.string('draft_id').primary()); + await up(db); + const operations = new McpOperations(db); + const principal = { user: { id: 'alice' }, grant: { id: 'grant-a' } } as McpPrincipal; + const create = (tool: string, key: string) => operations.run(principal, { + tool, args: { idempotencyKey: key }, repository: 'acme/repo', + }, async () => ({ status: 202, data: { state: 'queued' } })); + + const task = await create('send_task_followup', 'tracker-failure-1'); + const taskRow = await operations.get(principal, String(task.operationId)); + await syncLifecycle(operations, taskRow, { ...task, state: 'failed', targetState: { taskId: 'task-1', state: 'failed' } }); + assert.equal((operations.project(await operations.get(principal, taskRow.id)).lifecycle as { failure: unknown }).failure, null); + await syncLifecycle(operations, taskRow, { ...task, state: 'failed', targetState: { taskId: 'task-1', state: 'failed', reason: 'Agent stopped' } }); + assert.deepEqual((operations.project(await operations.get(principal, taskRow.id)).lifecycle as { failure: unknown }).failure, { + code: 'EXECUTION_FAILED', message: 'Agent stopped', stage: 'internal', retryable: false, status: 500, + }); + + const review = await create('review_pull_request', 'review-failure-01'); + const reviewRow = await operations.get(principal, String(review.operationId)); + await syncLifecycle(operations, reviewRow, { ...review, state: 'failed', targetState: { + taskId: 'review-task', state: 'completed', reason: 'Review processing completed successfully', + reviewResults: [{ success: false, error: 'Reviewer unavailable' }, { success: false, error: 'Model timed out' }], + } }); + assert.deepEqual((operations.project(await operations.get(principal, reviewRow.id)).lifecycle as { failure: unknown }).failure, { + code: 'REVIEW_FAILED', message: 'Reviewer unavailable; Model timed out', stage: 'internal', retryable: false, status: 500, + details: { failedReviewCount: 2 }, + }); + + const ultrafix = await create('run_ultrafix', 'ultrafix-failure-1'); + const ultrafixRow = await operations.get(principal, String(ultrafix.operationId)); + await syncLifecycle(operations, ultrafixRow, { ...ultrafix, state: 'failed', + targetState: { taskId: 'ultrafix-task', state: 'completed', reason: 'Task completed successfully' }, + result: { loop: { completionStatus: 'failed', completionReason: 'Maximum cycles reached without resolving the findings.' } }, + }); + assert.deepEqual((operations.project(await operations.get(principal, ultrafixRow.id)).lifecycle as { failure: unknown }).failure, { + code: 'EXECUTION_FAILED', message: 'Maximum cycles reached without resolving the findings.', + stage: 'internal', retryable: false, status: 500, + }); + assert.equal(failureFromReceipt({ + targetState: { state: 'completed', reason: 'Task completed successfully' }, + result: { loop: { completionStatus: 'failed' } }, + })?.message, 'Ultrafix loop failed.'); + assert.equal(failureFromReceipt({ targetState: { currentTask: { reason: 'Nested goal task failed' } } })?.message, + 'Nested goal task failed'); + assert.deepEqual(artifactsFromReceipt({ repository: 'acme/repo' }, { + targetState: { currentTask: { taskId: 'nested-goal-task' }, final_pr_number: 64 }, + }), { + taskId: 'nested-goal-task', + pullRequest: { repository: 'acme/repo', number: 64, url: 'https://github.com/acme/repo/pull/64' }, + }); +}); + +test('goal failures and generated task pull requests remain durable after backend history is removed', async t => { + const db = knex({ client: 'better-sqlite3', connection: { filename: ':memory:' }, useNullAsDefault: true }); + t.after(() => db.destroy()); + await db.schema.createTable('task_drafts', table => table.string('draft_id').primary()); + await db.schema.createTable('goals', table => { + table.string('goal_id').primary(); table.string('owner_id'); table.string('repository'); + table.string('desired_state'); table.string('result_state'); table.string('current_task_id'); + table.integer('final_pr_number'); table.text('failure_reason'); + }); + await db.schema.createTable('tasks', table => { + table.string('task_id').primary(); table.string('job_id'); table.string('repository'); table.string('task_type'); + table.integer('pr_number'); table.text('initial_job_data'); table.timestamp('created_at').defaultTo(db.fn.now()); + }); + await db.schema.createTable('task_history', table => { + table.increments('history_id').primary(); table.string('task_id'); table.string('state'); + table.timestamp('timestamp').defaultTo(db.fn.now()); table.text('reason'); table.text('metadata'); + }); + await db.schema.createTable('task_submissions', table => { + table.string('id').primary(); table.string('user_id'); table.string('submission_key'); table.string('payload_hash'); + table.string('repository'); table.text('payload'); table.text('attachments'); table.string('state'); + table.integer('issue_number'); table.text('issue_url'); table.string('task_id'); table.string('retry_event_id'); + table.boolean('dispatch_complete'); table.text('error'); table.timestamp('created_at').defaultTo(db.fn.now()); + }); + await up(db); + + const principal = { user: { id: 'alice' }, grant: { id: 'grant-a' } } as McpPrincipal; + const operations = new McpOperations(db); + const deps = { db, policy: { repository: async () => {}, requirePermission: () => {}, config: {} } as never, + taskQueue: {} as never, redisClient: {} as never, runtimeBuildQueue: {} as never } as ToolDeps; + const catalog = createToolCatalog(deps); + const get = catalog.find(tool => tool.name === 'get_operation')!; + const list = catalog.find(tool => tool.name === 'list_operations')!; + + const goalArgs = { idempotencyKey: 'durable-goal-output-1' }; + const goalId = 'goal-output-1'; + const goalTaskId = 'goal-task-output-1'; + const goalReceipt = await operations.run(principal, { tool: 'create_goal', args: goalArgs, repository: 'acme/repo' }, + async () => ({ status: 202, data: { state: 'accepted', continuation: { goalId } } })); + await db('tasks').insert({ task_id: goalTaskId, repository: 'acme/repo', task_type: 'goal' }); + await db('task_history').insert({ task_id: goalTaskId, state: 'failed', reason: 'Backing task stopped' }); + await db('goals').insert({ goal_id: goalId, owner_id: 'alice', repository: 'acme/repo', desired_state: 'running', + result_state: 'failed', current_task_id: goalTaskId, final_pr_number: 61, failure_reason: 'Provider exhausted its retry budget' }); + + const failedGoal = (await get.run({ principal, args: get.schema.parse({ operationId: goalReceipt.operationId }) })).data as Record; + assert.deepEqual(failedGoal.lifecycle, { + state: 'failed', acceptedAt: (goalReceipt.lifecycle as { acceptedAt: string }).acceptedAt, + startedAt: (failedGoal.lifecycle as { startedAt: string }).startedAt, + finishedAt: (failedGoal.lifecycle as { finishedAt: string }).finishedAt, + failure: { code: 'EXECUTION_FAILED', message: 'Provider exhausted its retry budget', stage: 'internal', retryable: false, status: 500 }, + artifacts: { taskId: goalTaskId, pullRequest: { repository: 'acme/repo', number: 61, url: 'https://github.com/acme/repo/pull/61' } }, + progress: (failedGoal.lifecycle as { progress: unknown }).progress, + }); + + const submissionId = 'submission-output-1'; + const taskId = 'ordinary-task-output-1'; + const taskArgs = { idempotencyKey: 'durable-task-output-1' }; + const taskReceipt = await operations.run(principal, { tool: 'create_task', args: taskArgs, repository: 'acme/repo' }, async () => ({ + status: 202, data: { id: submissionId, submissionId, submissionState: 'queued', state: 'queued', + continuation: { submissionId, taskId } }, + })); + await db('task_submissions').insert({ id: submissionId, user_id: 'alice', submission_key: 'submission-key', payload_hash: 'hash', + repository: 'acme/repo', payload: '{}', attachments: '[]', state: 'queued', task_id: taskId, dispatch_complete: true }); + await db('tasks').insert({ task_id: taskId, repository: 'acme/repo', task_type: 'issue' }); + await db('task_history').insert({ task_id: taskId, state: 'completed', reason: 'Task completed successfully' }); + + const completedWithoutPr = (await get.run({ principal, args: get.schema.parse({ operationId: taskReceipt.operationId }) })).data as Record; + assert.deepEqual((completedWithoutPr.lifecycle as { artifacts: unknown }).artifacts, { taskId }); + // Task completion is published before tasks.pr_number, so a later poll must + // still enrich the already-terminal durable receipt. + await db('tasks').where({ task_id: taskId }).update({ pr_number: 73 }); + const completedWithPr = (await get.run({ principal, args: get.schema.parse({ operationId: taskReceipt.operationId }) })).data as Record; + assert.deepEqual((completedWithPr.lifecycle as { artifacts: unknown }).artifacts, { + taskId, pullRequest: { repository: 'acme/repo', number: 73, url: 'https://github.com/acme/repo/pull/73' }, + }); + + await db('task_history').delete(); + await db('goals').delete(); + await db('tasks').delete(); + await db('task_submissions').delete(); + + const durableGoal = (await get.run({ principal, args: get.schema.parse({ operationId: goalReceipt.operationId }) })).data as Record; + assert.deepEqual((durableGoal.lifecycle as { failure: unknown }).failure, + { code: 'EXECUTION_FAILED', message: 'Provider exhausted its retry budget', stage: 'internal', retryable: false, status: 500 }); + assert.deepEqual((durableGoal.lifecycle as { artifacts: unknown }).artifacts, + { taskId: goalTaskId, pullRequest: { repository: 'acme/repo', number: 61, url: 'https://github.com/acme/repo/pull/61' } }); + const durableTask = await operations.replay(principal, 'create_task', taskArgs); + assert.deepEqual((durableTask?.lifecycle as { artifacts: unknown }).artifacts, + { taskId, pullRequest: { repository: 'acme/repo', number: 73, url: 'https://github.com/acme/repo/pull/73' } }); + const replayedGoal = await operations.replay(principal, 'create_goal', goalArgs); + assert.equal(replayedGoal?.state, 'failed'); + assert.equal(replayedGoal?.retryAfterSeconds, undefined); + const listed = (await list.run({ principal, args: list.schema.parse({}) })).data as { operations: Array> }; + assert.equal(listed.operations.length, 2); + assert.ok(listed.operations.every(operation => Object.keys((operation.lifecycle as { artifacts: object }).artifacts).length > 0)); + const listedGoal = listed.operations.find(operation => operation.operationId === goalReceipt.operationId)!; + assert.equal(listedGoal.state, 'failed'); + assert.equal(listedGoal.retryAfterSeconds, undefined); +}); + +test('list_operations filters active receipts by exact owner and grant without refreshing trackers', async t => { + const db = knex({ client: 'better-sqlite3', connection: { filename: ':memory:' }, useNullAsDefault: true }); + t.after(() => db.destroy()); + await db.schema.createTable('task_drafts', table => table.string('draft_id').primary()); + await up(db); + const operations = new McpOperations(db); + const principal = { user: { id: 'alice' }, grant: { id: 'grant-a' } } as McpPrincipal; + const otherGrant = { user: { id: 'alice' }, grant: { id: 'grant-b' } } as McpPrincipal; + const accepted = async (actor: McpPrincipal, tool: string, key: string) => operations.run(actor, { + tool, args: { idempotencyKey: key }, repository: 'acme/repo', + }, async () => ({ status: 202, data: { state: 'queued' } })); + await accepted(principal, 'run_ultrafix', 'active-ultrafix-1'); + await accepted(principal, 'review_pull_request', 'active-review-0001'); + await accepted(otherGrant, 'review_pull_request', 'other-grant-run1'); + await operations.run(principal, { tool: 'run_ultrafix', args: { idempotencyKey: 'done-ultrafix-01' }, repository: 'acme/repo' }, + async () => ({ status: 200, data: { changed: true } })); + + const deps = { db, policy: { repository: async () => {}, requirePermission: () => {}, config: {} } as never, + taskQueue: {} as never, redisClient: {} as never, runtimeBuildQueue: {} as never } as ToolDeps; + const tool = createToolCatalog(deps).find(candidate => candidate.name === 'list_operations')!; + const args = tool.schema.parse({ tool: 'run_ultrafix', lifecycle: 'active' }); + const result = (await tool.run({ principal, args })).data as { operations: Array>; nextOffset: number | null }; + assert.equal(result.operations.length, 1); + assert.equal(result.operations[0].tool, 'run_ultrafix'); + assert.equal((result.operations[0].lifecycle as { state: string }).state, 'accepted'); + assert.equal(result.operations[0].refreshWith, 'get_operation'); + const hidden = (await tool.run({ principal: otherGrant, args: tool.schema.parse({ tool: 'run_ultrafix', lifecycle: 'active' }) })).data as { operations: unknown[] }; + assert.deepEqual(hidden.operations, []); +}); + +test('list_operations filters current repository, tool permission and cancellation-source authorization before paging', async t => { + const db = knex({ client: 'better-sqlite3', connection: { filename: ':memory:' }, useNullAsDefault: true }); + t.after(() => db.destroy()); + await db.schema.createTable('task_drafts', table => table.string('draft_id').primary()); + await up(db); + const operations = new McpOperations(db); + const principal = { user: { id: 'alice' }, grant: { id: 'grant-a' }, authorization: { permissions: [] } } as unknown as McpPrincipal; + const run = (tool: string, key: string, repository?: string, data: Record = {}) => operations.run(principal, { + tool, args: { idempotencyKey: key }, repository, + }, async () => ({ status: 200, data })); + const forbiddenPermission = await run('update_execution_settings', 'permission-hidden-1', undefined, { secret: 'permission-secret' }); + const forbiddenRepository = await run('run_ultrafix', 'repository-hidden1', 'acme/forbidden', { secret: 'repository-secret' }); + const firstAllowed = await run('run_ultrafix', 'allowed-operation1', 'acme/allowed', { marker: 'first' }); + const secondAllowed = await run('review_pull_request', 'allowed-operation2', 'acme/allowed', { marker: 'second' }); + const source = await run('review_pull_request', 'cancel-source-key', 'acme/forbidden', { marker: 'source' }); + const cancellation = await run('cancel_operation', 'cancel-receipt-01', undefined, { operationId: source.operationId, cancellation: 'requested' }); + const now = Date.now(); + for (const [receipt, acceptedAt] of [[forbiddenPermission, now], [forbiddenRepository, now - 1], [firstAllowed, now - 2], [secondAllowed, now - 3], [source, now - 4], [cancellation, now - 5]] as const) { + await db('mcp_operations').where({ id: receipt.operationId }).update({ accepted_at: acceptedAt, created_at: acceptedAt }); + } + + const policy = { + config: {}, + repository: async (_actor: McpPrincipal, repository: string) => { + if (repository === 'acme/forbidden') throw new McpError('REPOSITORY_FORBIDDEN', 'No current access', 403); + }, + requirePermission: (actor: McpPrincipal, permission: string) => { + if (!actor.authorization.permissions.includes(permission as never)) throw new McpError('INSUFFICIENT_INSTANCE_PERMISSION', 'Permission removed', 403); + }, + }; + const deps = { db, policy, taskQueue: {} as never, redisClient: {} as never, runtimeBuildQueue: {} as never } as unknown as ToolDeps; + const catalog = createToolCatalog(deps); + const list = catalog.find(tool => tool.name === 'list_operations')!; + const get = catalog.find(tool => tool.name === 'get_operation')!; + + const pageOne = (await list.run({ principal, args: list.schema.parse({ limit: 1 }) })).data as { operations: Array>; nextOffset: number | null }; + assert.deepEqual(pageOne.operations.map(item => item.operationId), [firstAllowed.operationId]); + assert.equal(pageOne.nextOffset, 1); + assert.doesNotMatch(JSON.stringify(pageOne), /permission-secret|repository-secret/); + const pageTwo = (await list.run({ principal, args: list.schema.parse({ offset: 1, limit: 1 }) })).data as { operations: Array>; nextOffset: number | null }; + assert.deepEqual(pageTwo.operations.map(item => item.operationId), [secondAllowed.operationId]); + assert.equal(pageTwo.nextOffset, null); + const cancellations = (await list.run({ principal, args: list.schema.parse({ tool: 'cancel_operation' }) })).data as { operations: unknown[] }; + assert.deepEqual(cancellations.operations, []); + + const allowedSource = await run('review_pull_request', 'allowed-cancel-src', 'acme/allowed', { marker: 'allowed-source' }); + const allowedCancellation = await run('cancel_operation', 'allowed-cancel-rct', undefined, + { operationId: allowedSource.operationId, cancellation: 'requested' }); + const otherSource = await run('review_pull_request', 'other-cancel-src-1', 'acme/other', { marker: 'other-source' }); + await run('cancel_operation', 'other-cancel-rct1', undefined, + { operationId: otherSource.operationId, cancellation: 'requested' }); + const repositoryCancellations = (await list.run({ principal, args: list.schema.parse({ + tool: 'cancel_operation', repository: 'acme/allowed', limit: 1, + }) })).data as { operations: Array> }; + assert.deepEqual(repositoryCancellations.operations.map(item => item.operationId), [allowedCancellation.operationId]); + + await assert.rejects(get.run({ principal, args: get.schema.parse({ operationId: forbiddenPermission.operationId }) }), /Permission removed/); + await assert.rejects(get.run({ principal, args: get.schema.parse({ operationId: forbiddenRepository.operationId }) }), /No current access/); + await assert.rejects(get.run({ principal, args: get.schema.parse({ operationId: cancellation.operationId }) }), /No current access/); +}); + test('CIMD intersects plural supported methods with public PKCE instead of trusting a legacy preference', () => { const id = 'https://client.example/oauth/client.json'; const document = { client_id: id, client_name: 'Test', redirect_uris: ['https://client.example/callback'], token_endpoint_auth_methods_supported: ['none', 'private_key_jwt'], token_endpoint_auth_method: 'private_key_jwt' }; diff --git a/packages/api/test/mcpTaskSubmissions.test.ts b/packages/api/test/mcpTaskSubmissions.test.ts index f4a1ef3c1..0aa8b734e 100644 --- a/packages/api/test/mcpTaskSubmissions.test.ts +++ b/packages/api/test/mcpTaskSubmissions.test.ts @@ -27,6 +27,10 @@ interface Receipt { state: string; result: SubmissionData; retryAfterSeconds?: number; + lifecycle: { + state: string; acceptedAt: string; startedAt: string | null; finishedAt: string | null; + failure: unknown; artifacts: Record; progress: unknown; + }; } async function fixture() { @@ -37,7 +41,7 @@ async function fixture() { await identityMigration(db); await db.schema.createTable('tasks', table => { table.string('task_id').primary(); table.string('repository'); table.string('task_type'); - table.string('initial_job_data'); table.timestamp('created_at').defaultTo(db.fn.now()); + table.integer('pr_number'); table.string('initial_job_data'); table.timestamp('created_at').defaultTo(db.fn.now()); }); await db.schema.createTable('task_history', table => { table.increments('history_id'); table.string('task_id'); table.string('state'); @@ -90,6 +94,8 @@ test('MCP launches ordinary issue work once and follows delayed task association const first = await f.call('create_task', args); const receipt = first.data as Receipt; assert.equal(receipt.state, 'queued'); + assert.equal(receipt.lifecycle.state, 'accepted'); + assert.ok(receipt.lifecycle.acceptedAt); assert.equal(receipt.result.taskId, null); assert.equal(receipt.result.issueUrl, 'https://github.com/owner/repo/issues/42'); assert.deepEqual((await f.call('create_task', args)).data, first.data); @@ -111,16 +117,26 @@ test('MCP launches ordinary issue work once and follows delayed task association await f.db('tasks').insert({ task_id: 'ordinary-task', repository: 'owner/repo', task_type: 'issue' }); await f.db('task_submissions').update({ task_id: 'ordinary-task' }); await f.db('task_history').insert({ task_id: 'ordinary-task', state: 'processing' }); - assert.equal((await poll()).state, 'running'); + const running = await poll(); + assert.equal(running.state, 'running'); + assert.equal(running.lifecycle.state, 'running'); + assert.ok(running.lifecycle.startedAt); const status = await f.call('get_task_submission', { repository: 'owner/repo', submissionId: stored.id }); assert.equal((status.data as SubmissionData).taskId, 'ordinary-task'); assert.equal(status.links.ui, 'https://instance.example/tasks/ordinary-task'); await f.db('task_history').insert({ task_id: 'ordinary-task', state: 'completed' }); - const completed = await poll(); + const [completed, concurrent] = await Promise.all([poll(), poll()]); assert.equal(completed.state, 'completed'); + assert.equal(completed.lifecycle.state, 'completed'); + assert.equal(concurrent.lifecycle.state, 'completed'); + assert.ok(completed.lifecycle.finishedAt); + assert.equal(completed.lifecycle.artifacts.taskId, 'ordinary-task'); assert.equal(completed.result.continuation.taskId, 'ordinary-task'); assert.equal(completed.retryAfterSeconds, undefined); - assert.equal((await poll()).state, 'completed'); + await f.db('task_history').where({ task_id: 'ordinary-task' }).delete(); + const durable = await poll(); + assert.equal(durable.lifecycle.state, 'completed'); + assert.equal(durable.lifecycle.artifacts.taskId, 'ordinary-task'); } finally { await f.db.destroy(); } }); diff --git a/packages/api/test/mcpWorkflows.test.ts b/packages/api/test/mcpWorkflows.test.ts index dc195dda0..a8a244473 100644 --- a/packages/api/test/mcpWorkflows.test.ts +++ b/packages/api/test/mcpWorkflows.test.ts @@ -1,3 +1,4 @@ +/* eslint-disable max-lines -- both MCP SDK eras share one end-to-end workflow fixture */ import assert from 'node:assert/strict'; import { test, mock } from 'node:test'; import { randomBytes, randomUUID } from 'node:crypto'; @@ -246,8 +247,17 @@ test('both SDK eras drive persisted goal, TODO, notification, settings and guard const created = await call('create_goal', { repository, objective: 'Improve reliability', agentId: agent.config.id, model: 'fixture-model', launchStrategy: 'direct' }, true); assert.equal(created.state, 'accepted', JSON.stringify(created)); const goalId = created.result.continuation.goalId; - assert.equal((await db('goals').where({ goal_id: goalId }).first()).desired_state, 'running'); + const createdGoal = await db('goals').where({ goal_id: goalId }).first(); + assert.equal(createdGoal.desired_state, 'running'); assert.ok(jobs.some(job => job.goalId === goalId)); + await db('tasks').insert({ task_id: createdGoal.current_task_id, repository, task_type: 'goal' }); + await db('task_history').insert({ task_id: createdGoal.current_task_id, state: 'processing' }); + await db('task_history').insert({ task_id: createdGoal.current_task_id, state: 'completed', reason: 'Goal task completed successfully' }); + const runningGoal = await call('get_operation', { operationId: created.operationId }); + assert.equal(runningGoal.lifecycle.state, 'running'); + assert.ok(runningGoal.lifecycle.startedAt); + const runningOperations = await call('list_operations', { lifecycle: 'running' }); + assert.ok(runningOperations.operations.some((operation: { operationId: string }) => operation.operationId === created.operationId)); const input = await call('send_goal_input', { repository, goalId, message: 'Add error handling' }, true); assert.equal(input.state, 'completed', JSON.stringify(input)); assert.ok(await db('goal_inputs').where({ goal_id: goalId, message: 'Add error handling' }).first()); await call('pause_goal', { repository, goalId }, true); assert.equal((await db('goals').where({ goal_id: goalId }).first()).desired_state, 'paused'); @@ -316,19 +326,20 @@ test('both SDK eras drive persisted goal, TODO, notification, settings and guard const ultrafix = await call('run_ultrafix', pr, true); assert.equal(ultrafix.state, 'posted'); const loopTask = `ultrafix-start-${modern}`, workEpoch = modern ? 1 : 2; await pendingTask(loopTask, ultrafix.result.commentId, modern ? 'review' : 'fix', workEpoch); - await db('task_history').insert({ task_id: loopTask, state: 'completed' }); + await db('task_history').insert({ task_id: loopTask, state: 'completed', reason: 'Task completed successfully' }); const loop = { active: true, workEpoch, cycleCount: 0, completionStatus: null as string | null, completionReason: null as string | null }; redisValues.set('ultrafix:state:acme:repo:42', JSON.stringify(loop)); assert.equal((await call('get_operation', { operationId: ultrafix.operationId })).state, 'running'); - loop.active = false; loop.completionStatus = 'succeeded'; loop.completionReason = 'Goal reached'; + loop.active = false; loop.completionStatus = modern ? 'succeeded' : 'failed'; + loop.completionReason = modern ? 'Goal reached' : 'Maximum cycles reached without resolving the findings.'; redisValues.set('ultrafix:state:acme:repo:42', JSON.stringify(loop)); const completedLoop = await call('get_operation', { operationId: ultrafix.operationId }); - assert.equal(completedLoop.state, 'completed'); - assert.equal(completedLoop.result.loop.workEpoch, workEpoch); - assert.equal(completedLoop.result.loop.completionStatus, 'succeeded'); + assert.equal(completedLoop.state, modern ? 'completed' : 'failed'); + assert.equal(completedLoop.result.loop.workEpoch, workEpoch); assert.equal(completedLoop.result.loop.completionStatus, loop.completionStatus); + if (!modern) assert.equal(completedLoop.lifecycle.failure.message, loop.completionReason); assert.equal(completedLoop.result.continuation.sourceTaskId, loopTask); redisValues.set('ultrafix:state:acme:repo:42', JSON.stringify({ ...loop, workEpoch: workEpoch + 1, active: true, completionStatus: null })); - assert.equal((await call('get_operation', { operationId: ultrafix.operationId })).state, 'completed'); + assert.equal((await call('get_operation', { operationId: ultrafix.operationId })).state, modern ? 'completed' : 'failed'); const lostIntake = await call('review_pull_request', pr, true); await db('mcp_operations').where({ id: lostIntake.operationId }).update({ created_at: Date.now() - 180000 }); assert.equal((await call('get_operation', { operationId: lostIntake.operationId })).state, 'unknown'); diff --git a/packages/core/src/db/migrations/20260910220000_add_mcp.js b/packages/core/src/db/migrations/20260910220000_add_mcp.js index 2145d3a9f..cd1fb350a 100644 --- a/packages/core/src/db/migrations/20260910220000_add_mcp.js +++ b/packages/core/src/db/migrations/20260910220000_add_mcp.js @@ -19,9 +19,17 @@ export async function up(knex) { table.string('payload_hash', 64).notNullable(); table.string('state').notNullable(); table.text('result').nullable(); + table.string('lifecycle', 16).notNullable(); + table.bigInteger('accepted_at').notNullable(); + table.bigInteger('started_at').nullable(); + table.bigInteger('finished_at').nullable(); + table.json('failure').nullable(); + table.json('artifacts').notNullable(); + table.json('progress').nullable(); table.bigInteger('created_at').notNullable(); table.bigInteger('updated_at').notNullable(); table.unique(['owner_id', 'grant_id', 'idempotency_key']); + table.index(['owner_id', 'grant_id', 'accepted_at'], 'mcp_operations_owner_grant_accepted_idx'); }); await knex.schema.alterTable('task_drafts', table => { table.integer('mcp_revision').notNullable().defaultTo(0); }); // Every writer, including browser handlers and background workers, invalidates