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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
211 changes: 211 additions & 0 deletions packages/api/mcp/operationLifecycle.ts
Original file line number Diff line number Diff line change
@@ -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<LifecycleOutcome>(['completed', 'failed', 'cancelled']);

function record(value: unknown): Record<string, unknown> | undefined {
return value !== null && typeof value === 'object' && !Array.isArray(value) ? value as Record<string, unknown> : 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<string, unknown>,
continuation: Record<string, unknown>,
result: Record<string, unknown>,
targetIssues: Record<string, unknown>[],
): 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<string, unknown>,
continuation: Record<string, unknown>,
target: Record<string, unknown>,
targetIssues: Record<string, unknown>[],
): number[] {
const issueNumbers = new Set<number>();
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<Operation, 'repository'>, receipt: Record<string, unknown>): Record<string, unknown> {
const result = record(receipt.result) ?? {};
const continuation = record(result.continuation) ?? {};
const target = record(receipt.targetState) ?? {};
const targetIssues: Record<string, unknown>[] = Array.isArray(target.issues)
? target.issues.map(record).filter((value): value is Record<string, unknown> => !!value) : [];
const artifacts: Record<string, unknown> = {};

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<string, unknown>): McpErrorEnvelope {
return { code, message: redactSecrets(message), stage: 'internal', retryable: false, status: 500, ...(details ? { details } : {}) };
}

function resultFailure(result: Record<string, unknown> | 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<string, unknown>): 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<string, unknown> => 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<string, unknown> | 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<string, unknown>,
): Promise<void> {
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<string, unknown> | undefined,
result: Record<string, unknown> | 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<string, unknown>,
): Promise<void> {
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);
}
75 changes: 63 additions & 12 deletions packages/api/mcp/operationTracking.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, unknown>,
): Promise<void> {
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<string, unknown>,
result: ExecutionResult,
task: TrackingContext['task'] | undefined,
): void {
const target = result.targetState
?? (receipt.targetState as Record<string, unknown> | 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<string, unknown>): Promise<void> {
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<string, unknown>;
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<string, unknown> | 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<Operation>('mcp_operations').where({ id: row.id }).first();
if (current && terminalStates.includes(current.state)) {
const currentResult = current.result ? JSON.parse(current.result) as ExecutionResult & Record<string, unknown> : {};
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<string, unknown>;
reviewResults?: Array<{ success: boolean; commentId?: number; commentUrl?: string }>;
loop?: { completionStatus?: string | null } & Record<string, unknown>;
}
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<string, unknown>;
}

Expand Down Expand Up @@ -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<void> {
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';
Expand Down
Loading
Loading