feat(ql3): terminalize copilot diagnosis failures

This commit is contained in:
whyour
2026-08-15 23:16:26 +08:00
parent 5fc70010f2
commit 0a31c3364b
26 changed files with 3217 additions and 127 deletions
@@ -3,6 +3,7 @@ import type { SecurityPrincipal } from '@qinglong/runtime-core/security';
import type { CopilotFailureDiagnosisAdmissionReceipt } from '../admission/contracts';
import type { CopilotFailureDiagnosisModelExecutionResult } from '../model-execution/coordinator';
import type { CopilotFailureDiagnosisToolExecutionResult } from '../tool-execution/contracts';
import type { CopilotFailureDiagnosisPreModelTerminalizationReceipt } from '../terminalization/contracts';
export const MAX_ACTIVE_COPILOT_FAILURE_DIAGNOSIS_APPLICATION_REQUESTS = 64;
@@ -17,8 +18,9 @@ export interface ExecuteCopilotFailureDiagnosisApplicationCommand {
export interface ExecuteCopilotFailureDiagnosisApplicationResult {
readonly admissionStatus: 'created' | 'existing';
readonly admission: Readonly<CopilotFailureDiagnosisAdmissionReceipt>;
readonly tool: Readonly<CopilotFailureDiagnosisToolExecutionResult>;
readonly tool: Readonly<CopilotFailureDiagnosisToolExecutionResult> | null;
readonly model: Readonly<CopilotFailureDiagnosisModelExecutionResult> | null;
readonly terminalization: Readonly<CopilotFailureDiagnosisPreModelTerminalizationReceipt> | null;
readonly terminalizationRequired: boolean;
}
@@ -7,9 +7,7 @@ import {
} from '@qinglong/runtime-core/builtin-run-log-excerpt-tool';
import type { RunRepositoryReader } from '@qinglong/runtime-core/run-repository';
import type { ToolPolicyAuthorizer } from '@qinglong/runtime-core/tool-registry';
import {
prepareToolInvocation,
} from '@qinglong/runtime-core/tool-registry';
import { prepareToolInvocation } from '@qinglong/runtime-core/tool-registry';
import type { ProjectToolDefinitionSnapshotRepository } from '@qinglong/runtime-core/project-tool-definition-snapshot';
import { projectToolDefinitionRegistry } from '@qinglong/runtime-core/project-tool-definition-snapshot';
import {
@@ -25,9 +23,7 @@ import {
} from '@qinglong/runtime-core/trusted-tool-invocation';
import { MAX_MODEL_INVOCATION_MS } from '../../../model-gateway/model';
import {
prepareCopilotFailureDiagnosisExecution,
} from '../admission/plan';
import { prepareCopilotFailureDiagnosisExecution } from '../admission/plan';
import type {
CopilotFailureDiagnosisAdmissionRepository,
CopilotFailureDiagnosisExecutionPlan,
@@ -37,10 +33,20 @@ import {
executeCopilotFailureDiagnosisTool,
type CopilotFailureDiagnosisToolExecutionDependencies,
} from '../tool-execution/coordinator';
import {
CopilotFailureDiagnosisToolExecutionDeadlineExceededError,
type CopilotFailureDiagnosisToolExecutionResult,
} from '../tool-execution/contracts';
import {
executeCopilotFailureDiagnosisModel,
type CopilotFailureDiagnosisModelExecutionDependencies,
} from '../model-execution/coordinator';
import {
terminalizeCopilotFailureDiagnosisBeforeModel,
type CopilotFailureDiagnosisPreModelTerminalizationDependencies,
type CopilotFailureDiagnosisPreModelTerminalizationTrigger,
} from '../terminalization/coordinator';
import { CopilotFailureDiagnosisPreModelTerminalizationNotReadyError } from '../terminalization/contracts';
import {
CopilotFailureDiagnosisApplicationBusyError,
CopilotFailureDiagnosisApplicationConflictError,
@@ -64,7 +70,10 @@ const IDENTITY_DOMAIN = Buffer.from(
export interface CopilotFailureDiagnosisApplicationDependencies {
readonly admissions: CopilotFailureDiagnosisAdmissionRepository;
readonly snapshots: Pick<ProjectToolDefinitionSnapshotRepository, 'findCurrent'>;
readonly snapshots: Pick<
ProjectToolDefinitionSnapshotRepository,
'findCurrent'
>;
readonly runs: Pick<
RunRepositoryReader,
'findRunById' | 'findLatestAttemptByRunId'
@@ -77,19 +86,23 @@ export interface CopilotFailureDiagnosisApplicationDependencies {
readonly authorizer: ToolPolicyAuthorizer;
readonly tool: CopilotFailureDiagnosisToolExecutionDependencies;
readonly model: CopilotFailureDiagnosisModelExecutionDependencies;
readonly terminalizations: CopilotFailureDiagnosisPreModelTerminalizationDependencies;
readonly executeTool: typeof executeCopilotFailureDiagnosisTool;
readonly executeModel: typeof executeCopilotFailureDiagnosisModel;
readonly terminalizeBeforeModel: typeof terminalizeCopilotFailureDiagnosisBeforeModel;
readonly modelIntent: Readonly<PrepareCopilotFailureDiagnosisModelIntent>;
readonly executionTimeoutMs: number;
readonly now?: () => number;
readonly nonceFactory?: (input: Readonly<{
key: Uint8Array;
keyId: string;
requestId: string;
projectId: string;
sourceRunId: string;
invocationActionDigest: string;
}>) => Uint8Array;
readonly nonceFactory?: (
input: Readonly<{
key: Uint8Array;
keyId: string;
requestId: string;
projectId: string;
sourceRunId: string;
invocationActionDigest: string;
}>,
) => Uint8Array;
}
interface ActiveRequest {
@@ -146,14 +159,16 @@ function preview(
});
}
function defaultNonce(input: Readonly<{
key: Uint8Array;
keyId: string;
requestId: string;
projectId: string;
sourceRunId: string;
invocationActionDigest: string;
}>): Uint8Array {
function defaultNonce(
input: Readonly<{
key: Uint8Array;
keyId: string;
requestId: string;
projectId: string;
sourceRunId: string;
invocationActionDigest: string;
}>,
): Uint8Array {
const derived = createHmac('sha256', Buffer.from(input.key))
.update(NONCE_DOMAIN)
.update(
@@ -224,6 +239,10 @@ function assertDependencies(
typeof value.executeModel !== 'function' ||
!value.tool ||
!value.model ||
typeof value.terminalizations?.repository?.findByRequestId !== 'function' ||
typeof value.terminalizations?.repository?.readAuthority !== 'function' ||
typeof value.terminalizations?.repository?.commit !== 'function' ||
typeof value.terminalizeBeforeModel !== 'function' ||
!value.modelIntent ||
typeof value.modelIntent !== 'object' ||
!Number.isSafeInteger(value.executionTimeoutMs) ||
@@ -359,21 +378,91 @@ export class CopilotFailureDiagnosisApplicationService {
});
}
}
const tool = await this.#dependencies.executeTool(
{
requestId: plan.requestId,
principal: command.principal,
authorizer: this.#dependencies.authorizer,
},
this.#dependencies.tool,
);
const boundary = await this.#preModel(plan.requestId, { kind: 'boundary' });
if (boundary) {
return Object.freeze({
admissionStatus,
admission,
tool: null,
model: null,
terminalization: boundary,
terminalizationRequired: false,
});
}
let tool: Readonly<CopilotFailureDiagnosisToolExecutionResult>;
try {
tool = await this.#dependencies.executeTool(
{
requestId: plan.requestId,
principal: command.principal,
authorizer: this.#dependencies.authorizer,
},
this.#dependencies.tool,
);
} catch (cause) {
if (
!(
cause instanceof
CopilotFailureDiagnosisToolExecutionDeadlineExceededError
)
) {
throw cause;
}
const terminalization = await this.#dependencies.terminalizeBeforeModel(
plan.requestId,
{ kind: 'boundary' },
this.#dependencies.terminalizations,
);
return Object.freeze({
admissionStatus,
admission,
tool: null,
model: null,
terminalization: terminalization.receipt,
terminalizationRequired: false,
});
}
if (tool.outcome !== 'succeeded') {
const terminalization = await this.#dependencies.terminalizeBeforeModel(
plan.requestId,
{ kind: 'tool_failure', completion: tool.completion },
this.#dependencies.terminalizations,
);
return Object.freeze({
admissionStatus,
admission,
tool,
model: null,
terminalizationRequired: true,
terminalization: terminalization.receipt,
terminalizationRequired: false,
});
}
const projection = await this.#preModel(plan.requestId, {
kind: 'tool_projection',
completion: tool.completion,
output: tool.output,
});
if (projection) {
return Object.freeze({
admissionStatus,
admission,
tool,
model: null,
terminalization: projection,
terminalizationRequired: false,
});
}
const afterToolBoundary = await this.#preModel(plan.requestId, {
kind: 'boundary',
});
if (afterToolBoundary) {
return Object.freeze({
admissionStatus,
admission,
tool,
model: null,
terminalization: afterToolBoundary,
terminalizationRequired: false,
});
}
const model = await this.#dependencies.executeModel(
@@ -385,10 +474,33 @@ export class CopilotFailureDiagnosisApplicationService {
admission,
tool,
model,
terminalization: null,
terminalizationRequired: false,
});
}
async #preModel(
requestId: string,
trigger: CopilotFailureDiagnosisPreModelTerminalizationTrigger,
) {
try {
const result = await this.#dependencies.terminalizeBeforeModel(
requestId,
trigger,
this.#dependencies.terminalizations,
);
return result.receipt;
} catch (cause) {
if (
cause instanceof
CopilotFailureDiagnosisPreModelTerminalizationNotReadyError
) {
return null;
}
throw cause;
}
}
async #prepare(
command: Readonly<ExecuteCopilotFailureDiagnosisApplicationCommand>,
) {
@@ -447,9 +559,7 @@ export class CopilotFailureDiagnosisApplicationService {
const key = await this.#dependencies.invocationKeys.active();
try {
const baseIdentity = identity('cda', command.requestId);
const nonce = (
this.#dependencies.nonceFactory ?? defaultNonce
)({
const nonce = (this.#dependencies.nonceFactory ?? defaultNonce)({
key: key.key,
keyId: key.keyId,
requestId: command.requestId,
@@ -507,9 +617,9 @@ export class CopilotFailureDiagnosisApplicationService {
return unavailable();
}
try {
const nonce = (
this.#dependencies.nonceFactory ?? defaultNonce
)(nonceInput(plan, material.key));
const nonce = (this.#dependencies.nonceFactory ?? defaultNonce)(
nonceInput(plan, material.key),
);
const inputArtifact = createToolInvocationInputArtifact(
{
artifactId: plan.tool.invocationArtifact.artifactId,
@@ -564,7 +674,10 @@ export class CopilotFailureDiagnosisApplicationService {
return unavailable(cause);
}
if (!Number.isSafeInteger(value) || value < 0) return invalid('clock');
if (value + this.#dependencies.executionTimeoutMs > Number.MAX_SAFE_INTEGER) {
if (
value + this.#dependencies.executionTimeoutMs >
Number.MAX_SAFE_INTEGER
) {
return invalid('deadline overflows');
}
return value;
@@ -34,8 +34,7 @@ import {
} from './finalization';
const TABLE = '"ql3_ai"."copilot_failure_diagnosis_model_outputs"';
const FINALIZATION_TABLE =
'"ql3_ai"."copilot_failure_diagnosis_finalizations"';
const FINALIZATION_TABLE = '"ql3_ai"."copilot_failure_diagnosis_finalizations"';
interface OutputRow extends Record<string, unknown> {
readonly artifactJson: unknown;
@@ -79,7 +78,9 @@ function unavailable(cause?: unknown): never {
});
}
function parse(row: OutputRow): Readonly<CopilotFailureDiagnosisOutputArtifact> {
function parse(
row: OutputRow,
): Readonly<CopilotFailureDiagnosisOutputArtifact> {
try {
return normalizeCopilotFailureDiagnosisOutputArtifact(
row.artifactJson as CopilotFailureDiagnosisOutputArtifact,
@@ -119,7 +120,8 @@ async function put(
client: PostgresClient,
artifactValue: CopilotFailureDiagnosisOutputArtifact,
): Promise<Readonly<CopilotFailureDiagnosisOutputArtifact>> {
const artifact = normalizeCopilotFailureDiagnosisOutputArtifact(artifactValue);
const artifact =
normalizeCopilotFailureDiagnosisOutputArtifact(artifactValue);
const existing = await read(client, artifact.artifactId);
if (existing) {
if (JSON.stringify(existing) !== JSON.stringify(artifact)) {
@@ -239,7 +241,9 @@ export class PostgresCopilotFailureDiagnosisModelRepository
if (!result.rows[0]) return null;
const row = result.rows[0];
const receipt = normalizeCopilotFailureDiagnosisFinalizationReceipt(
object(row.receiptJson) as unknown as CopilotFailureDiagnosisFinalizationReceipt,
object(
row.receiptJson,
) as unknown as CopilotFailureDiagnosisFinalizationReceipt,
);
if (
receipt.requestId !== string(row.requestId) ||
@@ -268,7 +272,9 @@ export class PostgresCopilotFailureDiagnosisModelRepository
event.step_run_id AS "eventStepRunId",
event.payload, event.created_at_ms AS "eventCreatedAtMs",
completion.completion_digest AS "completionDigest",
completion.outcome AS "completionOutcome"
completion.outcome AS "completionOutcome",
resolution.decision AS "resolutionDecision",
resolution.resolved_at_ms AS "resolvedAtMs"
FROM "ql3"."runs" AS run
JOIN "ql3"."step_runs" AS step
ON step.run_id = run.id AND step.id = $1
@@ -276,6 +282,8 @@ export class PostgresCopilotFailureDiagnosisModelRepository
ON event.run_id = run.id AND event.id = $2
JOIN "ql3_ai"."model_invocation_completions" AS completion
ON completion.invocation_id = $3
LEFT JOIN "ql3_ai"."model_invocation_resolutions" AS resolution
ON resolution.invocation_id = completion.invocation_id
WHERE run.id = $4`,
[
receipt.modelStepRunId,
@@ -301,7 +309,15 @@ export class PostgresCopilotFailureDiagnosisModelRepository
string(proof.eventStepRunId) !== receipt.modelStepRunId ||
integer(proof.eventCreatedAtMs) !== receipt.finalizedAtMs ||
string(proof.completionDigest) !== receipt.completionDigest ||
string(proof.completionOutcome) !== receipt.outcome ||
!(
string(proof.completionOutcome) === receipt.outcome ||
(string(proof.completionOutcome) === 'outcome_unknown' &&
((proof.resolutionDecision === 'fail' &&
receipt.outcome === 'failed') ||
(proof.resolutionDecision === 'cancel' &&
receipt.outcome === 'cancelled')) &&
integer(proof.resolvedAtMs) === receipt.finalizedAtMs)
) ||
payload.requestId !== receipt.requestId ||
payload.planDigest !== receipt.planDigest ||
payload.invocationId !== receipt.invocationId ||
@@ -320,10 +336,12 @@ export class PostgresCopilotFailureDiagnosisModelRepository
}
}
async finalize(requestId: string): Promise<Readonly<{
status: 'created' | 'existing';
receipt: Readonly<CopilotFailureDiagnosisFinalizationReceipt>;
}>> {
async finalize(requestId: string): Promise<
Readonly<{
status: 'created' | 'existing';
receipt: Readonly<CopilotFailureDiagnosisFinalizationReceipt>;
}>
> {
const existing = await this.findFinalization(requestId);
if (existing) {
return Object.freeze({ status: 'existing' as const, receipt: existing });
@@ -333,7 +351,9 @@ export class PostgresCopilotFailureDiagnosisModelRepository
try {
client = await this.#pool.connect();
} catch (cause) {
throw new CopilotFailureDiagnosisFinalizationUnavailableError({ cause });
throw new CopilotFailureDiagnosisFinalizationUnavailableError({
cause,
});
}
let began = false;
try {
@@ -373,13 +393,17 @@ export class PostgresCopilotFailureDiagnosisModelRepository
}
if (
cause instanceof CopilotFailureDiagnosisFinalizationConflictError ||
cause instanceof CopilotFailureDiagnosisModelExecutionInProgressError ||
cause instanceof CopilotFailureDiagnosisModelResolutionRequiredError ||
cause instanceof
CopilotFailureDiagnosisModelExecutionInProgressError ||
cause instanceof
CopilotFailureDiagnosisModelResolutionRequiredError ||
cause instanceof CopilotFailureDiagnosisFinalizationUnavailableError
) {
throw cause;
}
throw new CopilotFailureDiagnosisFinalizationUnavailableError({ cause });
throw new CopilotFailureDiagnosisFinalizationUnavailableError({
cause,
});
} finally {
client.release();
}
@@ -390,10 +414,12 @@ export class PostgresCopilotFailureDiagnosisModelRepository
async #finalizeInTransaction(
client: PostgresClient,
requestId: string,
): Promise<Readonly<{
status: 'created' | 'existing';
receipt: Readonly<CopilotFailureDiagnosisFinalizationReceipt>;
}>> {
): Promise<
Readonly<{
status: 'created' | 'existing';
receipt: Readonly<CopilotFailureDiagnosisFinalizationReceipt>;
}>
> {
const admission = await client.query<Record<string, unknown>>(
`SELECT plan_json AS "planJson"
FROM "ql3_ai"."copilot_failure_diagnosis_admissions"
@@ -404,19 +430,29 @@ export class PostgresCopilotFailureDiagnosisModelRepository
throw new CopilotFailureDiagnosisFinalizationConflictError();
}
const plan = normalizeCopilotFailureDiagnosisExecutionPlan(
object(admission.rows[0]!.planJson) as unknown as CopilotFailureDiagnosisExecutionPlan,
object(
admission.rows[0]!.planJson,
) as unknown as CopilotFailureDiagnosisExecutionPlan,
);
const durable = await client.query<Record<string, unknown>>(
`SELECT run.status AS "runStatus", run.version AS "runVersion",
run.event_sequence AS "runEventSequence",
step.status AS "stepStatus",
step.step_run_digest AS "stepRunDigest",
completion.record_json AS "completionJson"
completion.record_json AS "completionJson",
resolution.decision AS "resolutionDecision",
resolution.completion_digest AS "resolutionCompletionDigest",
resolution.resolved_at_ms AS "resolvedAtMs",
resolution_mutation.step_run_digest AS "resolutionStepRunDigest"
FROM "ql3"."runs" AS run
JOIN "ql3"."step_runs" AS step
ON step.run_id = run.id AND step.id = $1
LEFT JOIN "ql3_ai"."model_invocation_completions" AS completion
ON completion.invocation_id = $2
LEFT JOIN "ql3_ai"."model_invocation_resolutions" AS resolution
ON resolution.invocation_id = $2
LEFT JOIN "ql3"."step_run_mutations" AS resolution_mutation
ON resolution_mutation.mutation_id = resolution.mutation_id
WHERE run.id = $3
FOR UPDATE OF run, step`,
[plan.modelStepRunId, plan.modelInvocationId, plan.runId],
@@ -438,14 +474,41 @@ export class PostgresCopilotFailureDiagnosisModelRepository
completion.stepRunId !== plan.modelStepRunId ||
completion.traceId !== plan.traceId ||
string(row.runStatus) !== 'running' ||
string(row.stepRunDigest) !== completion.completedStepRunDigest
(completion.outcome !== 'outcome_unknown' &&
string(row.stepRunDigest) !== completion.completedStepRunDigest)
) {
throw new CopilotFailureDiagnosisFinalizationConflictError();
}
let outcome: CopilotFailureDiagnosisFinalOutcome;
let finalizedAtMs = completion.completedAtMs;
if (completion.outcome === 'outcome_unknown') {
throw new CopilotFailureDiagnosisModelResolutionRequiredError();
if (
row.resolutionDecision === null ||
row.resolutionDecision === undefined
) {
throw new CopilotFailureDiagnosisModelResolutionRequiredError();
}
if (row.resolutionDecision === 'retry') {
throw new CopilotFailureDiagnosisModelExecutionInProgressError();
}
if (
row.resolutionDecision !== 'fail' &&
row.resolutionDecision !== 'cancel'
) {
throw new CopilotFailureDiagnosisFinalizationConflictError();
}
if (
string(row.resolutionCompletionDigest) !==
completion.completionDigest ||
string(row.resolutionStepRunDigest) !== string(row.stepRunDigest)
) {
throw new CopilotFailureDiagnosisFinalizationConflictError();
}
outcome = row.resolutionDecision === 'fail' ? 'failed' : 'cancelled';
finalizedAtMs = integer(row.resolvedAtMs);
} else {
outcome = completion.outcome;
}
const outcome: CopilotFailureDiagnosisFinalOutcome = completion.outcome;
if (string(row.stepStatus) !== outcome) {
throw new CopilotFailureDiagnosisFinalizationConflictError();
}
@@ -485,19 +548,25 @@ export class PostgresCopilotFailureDiagnosisModelRepository
outputArtifactId: output?.artifactId ?? null,
finalRunVersion: runVersion + 1,
finalRunEventSequence: eventSequence + 1,
finalizedAtMs: completion.completedAtMs,
finalizedAtMs,
});
const failure = outcome === 'succeeded'
? { code: null, summary: null }
: outcome === 'timed_out'
? {
code: 'COPILOT_FAILURE_DIAGNOSIS_TIMED_OUT',
summary: 'Copilot failure diagnosis timed out',
}
: {
code: 'COPILOT_FAILURE_DIAGNOSIS_FAILED',
summary: 'Copilot failure diagnosis failed',
};
const failure =
outcome === 'succeeded'
? { code: null, summary: null }
: outcome === 'cancelled'
? {
code: 'COPILOT_FAILURE_DIAGNOSIS_CANCELLED',
summary: 'Copilot failure diagnosis was cancelled',
}
: outcome === 'timed_out'
? {
code: 'COPILOT_FAILURE_DIAGNOSIS_TIMED_OUT',
summary: 'Copilot failure diagnosis timed out',
}
: {
code: 'COPILOT_FAILURE_DIAGNOSIS_FAILED',
summary: 'Copilot failure diagnosis failed',
};
const updated = await client.query(
`UPDATE "ql3"."runs"
SET status = $1, version = $2, event_sequence = $3,
@@ -0,0 +1,4 @@
export * from './terminalization/contracts';
export * from './terminalization/protocol';
export * from './terminalization/coordinator';
export * from './terminalization/postgresRepository';
@@ -0,0 +1,160 @@
import type { RunRecord } from '@qinglong/runtime-core';
import type {
StepRunMutation,
StepRunRecord,
} from '@qinglong/runtime-core/step-run';
import type { CopilotFailureDiagnosisExecutionPlan } from '../admission/contracts';
export const COPILOT_FAILURE_DIAGNOSIS_PRE_MODEL_TERMINALIZATION_SCHEMA =
'qinglong/copilot-failure-diagnosis-pre-model-terminalization@v1' as const;
export const COPILOT_FAILURE_DIAGNOSIS_PRE_MODEL_TERMINALIZATION_COMMAND_SCHEMA =
'qinglong/copilot-failure-diagnosis-pre-model-terminalization-command@v1' as const;
export const COPILOT_FAILURE_DIAGNOSIS_PRE_MODEL_TERMINALIZATION_REASONS = [
'tool_failed',
'tool_timed_out',
'log_not_found',
'log_pending',
'log_missing',
'log_retired',
'tool_budget_exhausted',
'deadline_exceeded',
'cancellation_requested',
] as const;
export type CopilotFailureDiagnosisPreModelTerminalizationReason =
(typeof COPILOT_FAILURE_DIAGNOSIS_PRE_MODEL_TERMINALIZATION_REASONS)[number];
export type CopilotFailureDiagnosisPreModelTerminalizationStage =
| 'tool'
| 'log'
| 'deadline'
| 'cancellation';
export type CopilotFailureDiagnosisPreModelTerminalizationOutcome =
| 'failed'
| 'timed_out'
| 'cancelled';
export interface CopilotFailureDiagnosisTerminalStepReference {
readonly stepRunId: string;
readonly status: 'failed' | 'timed_out' | 'cancelled';
readonly version: number;
readonly mutationId: string;
readonly mutationDigest: string;
readonly eventId: string;
}
export interface CopilotFailureDiagnosisPreModelTerminalizationReceipt {
readonly schema: typeof COPILOT_FAILURE_DIAGNOSIS_PRE_MODEL_TERMINALIZATION_SCHEMA;
readonly requestId: string;
readonly planDigest: string;
readonly runId: string;
readonly stage: CopilotFailureDiagnosisPreModelTerminalizationStage;
readonly reason: CopilotFailureDiagnosisPreModelTerminalizationReason;
readonly outcome: CopilotFailureDiagnosisPreModelTerminalizationOutcome;
readonly evidenceDigest: string;
readonly toolStartId: string | null;
readonly toolCompletionDigest: string | null;
readonly terminalSteps: readonly Readonly<CopilotFailureDiagnosisTerminalStepReference>[];
readonly finalRunVersion: number;
readonly finalRunEventSequence: number;
readonly runEventId: string;
readonly finalizedAtMs: number;
readonly receiptDigest: string;
}
export interface CopilotFailureDiagnosisPreModelTerminalizationCommand {
readonly schema: typeof COPILOT_FAILURE_DIAGNOSIS_PRE_MODEL_TERMINALIZATION_COMMAND_SCHEMA;
readonly plan: Readonly<CopilotFailureDiagnosisExecutionPlan>;
readonly expectedRunVersion: number;
readonly expectedRunEventSequence: number;
readonly stepMutations: readonly Readonly<StepRunMutation>[];
readonly receipt: Readonly<CopilotFailureDiagnosisPreModelTerminalizationReceipt>;
readonly commandDigest: string;
}
export interface CopilotFailureDiagnosisPreModelTerminalizationAuthority {
readonly plan: Readonly<CopilotFailureDiagnosisExecutionPlan>;
readonly run: Readonly<
Pick<
RunRecord,
| 'id'
| 'projectId'
| 'status'
| 'version'
| 'eventSequence'
| 'cancelRequestedAtMs'
| 'cancelReason'
>
>;
readonly toolStep: Readonly<StepRunRecord>;
readonly modelStep: Readonly<StepRunRecord>;
readonly modelStartExists: boolean;
readonly observedAtMs: number;
}
export interface CopilotFailureDiagnosisPreModelTerminalizationRepository {
findByRequestId(
requestId: string,
): Promise<Readonly<CopilotFailureDiagnosisPreModelTerminalizationReceipt> | null>;
readAuthority(
requestId: string,
): Promise<Readonly<CopilotFailureDiagnosisPreModelTerminalizationAuthority>>;
commit(
command: Readonly<CopilotFailureDiagnosisPreModelTerminalizationCommand>,
): Promise<
Readonly<{
status: 'created' | 'existing';
receipt: Readonly<CopilotFailureDiagnosisPreModelTerminalizationReceipt>;
}>
>;
}
export class InvalidCopilotFailureDiagnosisPreModelTerminalizationError extends TypeError {
readonly code = 'COPILOT_FAILURE_DIAGNOSIS_PRE_MODEL_TERMINALIZATION_INVALID';
constructor(message: string) {
super(
`Copilot failure diagnosis pre-Model terminalization is invalid: ${message}`,
);
this.name = 'InvalidCopilotFailureDiagnosisPreModelTerminalizationError';
}
}
export class CopilotFailureDiagnosisPreModelTerminalizationConflictError extends Error {
readonly code =
'COPILOT_FAILURE_DIAGNOSIS_PRE_MODEL_TERMINALIZATION_CONFLICT';
constructor(message = 'durable terminalization authority changed') {
super(
`Copilot failure diagnosis pre-Model terminalization conflicts: ${message}`,
);
this.name = 'CopilotFailureDiagnosisPreModelTerminalizationConflictError';
}
}
export class CopilotFailureDiagnosisPreModelTerminalizationNotReadyError extends Error {
readonly code =
'COPILOT_FAILURE_DIAGNOSIS_PRE_MODEL_TERMINALIZATION_NOT_READY';
constructor() {
super('Copilot failure diagnosis has no terminal pre-Model condition');
this.name = 'CopilotFailureDiagnosisPreModelTerminalizationNotReadyError';
}
}
export class CopilotFailureDiagnosisPreModelTerminalizationUnavailableError extends Error {
readonly code =
'COPILOT_FAILURE_DIAGNOSIS_PRE_MODEL_TERMINALIZATION_UNAVAILABLE';
constructor(options?: ErrorOptions) {
super(
'Copilot failure diagnosis pre-Model terminalization is unavailable',
options,
);
this.name =
'CopilotFailureDiagnosisPreModelTerminalizationUnavailableError';
}
}
@@ -0,0 +1,405 @@
import type { ToolJsonValue } from '@qinglong/runtime-core/tool-registry';
import type { ToolExecutionCompletionRecord } from '@qinglong/runtime-core/tool-execution-completion';
import type { ToolExecutionFailureCompletionRecord } from '@qinglong/runtime-core/tool-execution-failure-completion';
import { BUILTIN_RUN_LOG_EXCERPT_TIMEOUT_SECONDS } from '@qinglong/runtime-core/builtin-run-log-excerpt-tool';
import {
STEP_RUN_TERMINAL_STATUSES,
transitionStepRunMutation,
type StepRunMutation,
type StepRunRecord,
type StepRunStatus,
} from '@qinglong/runtime-core/step-run';
import {
CopilotFailureDiagnosisPreModelTerminalizationConflictError,
CopilotFailureDiagnosisPreModelTerminalizationNotReadyError,
CopilotFailureDiagnosisPreModelTerminalizationUnavailableError,
InvalidCopilotFailureDiagnosisPreModelTerminalizationError,
type CopilotFailureDiagnosisPreModelTerminalizationReason,
type CopilotFailureDiagnosisPreModelTerminalizationRepository,
type CopilotFailureDiagnosisPreModelTerminalizationReceipt,
} from './contracts';
import {
copilotFailureDiagnosisPreModelEvidenceDigest,
copilotFailureDiagnosisPreModelTerminalizationMapping,
copilotFailureDiagnosisTerminalizationIdentity,
createCopilotFailureDiagnosisPreModelTerminalizationCommand,
createCopilotFailureDiagnosisPreModelTerminalizationReceipt,
terminalStepReference,
} from './protocol';
export type CopilotFailureDiagnosisPreModelTerminalizationTrigger =
| Readonly<{
kind: 'tool_failure';
completion: Readonly<ToolExecutionFailureCompletionRecord>;
}>
| Readonly<{
kind: 'tool_projection';
completion: Readonly<ToolExecutionCompletionRecord>;
output: ToolJsonValue;
}>
| Readonly<{ kind: 'boundary' }>;
export interface CopilotFailureDiagnosisPreModelTerminalizationDependencies {
readonly repository: CopilotFailureDiagnosisPreModelTerminalizationRepository;
}
export interface CopilotFailureDiagnosisPreModelTerminalizationResult {
readonly status: 'created' | 'existing';
readonly receipt: Readonly<CopilotFailureDiagnosisPreModelTerminalizationReceipt>;
}
function invalid(message: string): never {
throw new InvalidCopilotFailureDiagnosisPreModelTerminalizationError(message);
}
function conflict(message: string): never {
throw new CopilotFailureDiagnosisPreModelTerminalizationConflictError(
message,
);
}
function unavailable(cause?: unknown): never {
throw new CopilotFailureDiagnosisPreModelTerminalizationUnavailableError({
cause: cause instanceof Error ? cause : undefined,
});
}
function assertDependencies(
value: CopilotFailureDiagnosisPreModelTerminalizationDependencies,
): void {
if (
!value ||
typeof value !== 'object' ||
typeof value.repository?.findByRequestId !== 'function' ||
typeof value.repository?.readAuthority !== 'function' ||
typeof value.repository?.commit !== 'function'
) {
return invalid('dependencies are invalid');
}
}
function statusFromProjection(
output: ToolJsonValue,
runId: string,
attemptId: string,
): 'not_found' | 'pending' | 'missing' | 'retired' | null {
if (!output || typeof output !== 'object' || Array.isArray(output)) {
return conflict('Tool projection is not an object');
}
const record = output as Readonly<Record<string, ToolJsonValue>>;
if (
record.runId !== runId ||
record.attemptId !== attemptId ||
record.profile !== 'cluster-control' ||
typeof record.status !== 'string' ||
!['available', 'not_found', 'pending', 'missing', 'retired'].includes(
record.status,
)
) {
return conflict('Tool projection identity changed');
}
return record.status === 'available'
? null
: (record.status as 'not_found' | 'pending' | 'missing' | 'retired');
}
function triggerEvidence(
trigger: CopilotFailureDiagnosisPreModelTerminalizationTrigger,
authority: Awaited<
ReturnType<
CopilotFailureDiagnosisPreModelTerminalizationRepository['readAuthority']
>
>,
): Readonly<{
reason: CopilotFailureDiagnosisPreModelTerminalizationReason;
evidenceDigest: string;
toolStartId: string | null;
toolCompletionDigest: string | null;
}> {
const { plan, run, observedAtMs } = authority;
if (trigger.kind === 'tool_failure') {
const completion = trigger.completion;
if (
completion.runId !== plan.runId ||
completion.stepRunId !== plan.toolStepRunId ||
completion.outcome !== authority.toolStep.status ||
completion.completedStepRunDigest !== authority.toolStep.stepRunDigest
)
return conflict('Tool failure completion changed');
const reason =
completion.outcome === 'timed_out' ? 'tool_timed_out' : 'tool_failed';
return Object.freeze({
reason,
evidenceDigest: copilotFailureDiagnosisPreModelEvidenceDigest({
reason,
planDigest: plan.planDigest,
toolCompletionDigest: completion.completionDigest,
}),
toolStartId: completion.startId,
toolCompletionDigest: completion.completionDigest,
});
}
if (trigger.kind === 'tool_projection') {
const completion = trigger.completion;
if (
completion.runId !== plan.runId ||
completion.stepRunId !== plan.toolStepRunId ||
authority.toolStep.status !== 'succeeded' ||
completion.completedStepRunDigest !== authority.toolStep.stepRunDigest
)
return conflict('Tool success completion changed');
const status = statusFromProjection(
trigger.output,
plan.source.runId,
plan.source.attemptId,
);
if (status === null)
throw new CopilotFailureDiagnosisPreModelTerminalizationNotReadyError();
const reason =
`log_${status}` as CopilotFailureDiagnosisPreModelTerminalizationReason;
return Object.freeze({
reason,
evidenceDigest: copilotFailureDiagnosisPreModelEvidenceDigest({
reason,
planDigest: plan.planDigest,
toolCompletionDigest: completion.completionDigest,
sourceStatus: status,
}),
toolStartId: completion.startId,
toolCompletionDigest: completion.completionDigest,
});
}
if (run.cancelRequestedAtMs !== undefined && run.cancelReason !== undefined) {
const reason = 'cancellation_requested' as const;
return Object.freeze({
reason,
evidenceDigest: copilotFailureDiagnosisPreModelEvidenceDigest({
reason,
planDigest: plan.planDigest,
cancelRequestedAtMs: run.cancelRequestedAtMs,
cancelReason: run.cancelReason,
}),
toolStartId: null,
toolCompletionDigest: null,
});
}
if (observedAtMs >= plan.deadlineAtMs) {
const reason = 'deadline_exceeded' as const;
return Object.freeze({
reason,
evidenceDigest: copilotFailureDiagnosisPreModelEvidenceDigest({
reason,
planDigest: plan.planDigest,
deadlineAtMs: plan.deadlineAtMs,
}),
toolStartId: null,
toolCompletionDigest: null,
});
}
const requiredToolBudgetMs = BUILTIN_RUN_LOG_EXCERPT_TIMEOUT_SECONDS * 1_000;
if (
authority.toolStep.status === 'ready' &&
observedAtMs + requiredToolBudgetMs > plan.deadlineAtMs
) {
const reason = 'tool_budget_exhausted' as const;
return Object.freeze({
reason,
evidenceDigest: copilotFailureDiagnosisPreModelEvidenceDigest({
reason,
planDigest: plan.planDigest,
deadlineAtMs: plan.deadlineAtMs,
observedAtMs,
requiredToolBudgetMs,
}),
toolStartId: null,
toolCompletionDigest: null,
});
}
throw new CopilotFailureDiagnosisPreModelTerminalizationNotReadyError();
}
function targetFor(
step: Readonly<StepRunRecord>,
reason: CopilotFailureDiagnosisPreModelTerminalizationReason,
outcome: 'failed' | 'timed_out' | 'cancelled',
): StepRunStatus | null {
if (STEP_RUN_TERMINAL_STATUSES.includes(step.status as never)) return null;
if (reason === 'tool_failed' || reason === 'tool_timed_out') {
return step.kind === 'model'
? 'cancelled'
: conflict('Tool Step is not terminal');
}
if (reason.startsWith('log_')) {
return step.kind === 'model'
? 'failed'
: conflict('log evidence Tool is not terminal');
}
if (step.status === 'pending') return 'cancelled';
return outcome;
}
function transitionFacts(
target: StepRunStatus,
reason: CopilotFailureDiagnosisPreModelTerminalizationReason,
): Readonly<{ resultCode: string; errorSummary?: string }> {
if (target === 'failed') {
return Object.freeze({
resultCode: 'copilot_log_unavailable',
errorSummary: 'Failure diagnosis source log is unavailable',
});
}
if (target === 'timed_out') {
return Object.freeze({
resultCode: 'copilot_deadline_exceeded',
errorSummary: 'Copilot failure diagnosis deadline exceeded',
});
}
return Object.freeze({
resultCode: reason.startsWith('tool_')
? 'copilot_tool_dependency_failed'
: 'copilot_cancelled',
});
}
function mutations(
authority: Awaited<
ReturnType<
CopilotFailureDiagnosisPreModelTerminalizationRepository['readAuthority']
>
>,
reason: CopilotFailureDiagnosisPreModelTerminalizationReason,
outcome: 'failed' | 'timed_out' | 'cancelled',
): readonly Readonly<StepRunMutation>[] {
let runVersion = authority.run.version;
let eventSequence = authority.run.eventSequence;
const result: StepRunMutation[] = [];
for (const step of [authority.toolStep, authority.modelStep]) {
const target = targetFor(step, reason, outcome);
if (target === null) continue;
const facts = transitionFacts(target, reason);
const mutationId = copilotFailureDiagnosisTerminalizationIdentity(
'mutation',
authority.plan.planDigest,
reason,
step.id,
);
const eventId = copilotFailureDiagnosisTerminalizationIdentity(
'step-event',
authority.plan.planDigest,
reason,
step.id,
);
const mutation = transitionStepRunMutation(
step,
{
expectedVersion: step.version,
expectedDigest: step.stepRunDigest,
mutationId,
to: target,
atMs: authority.observedAtMs,
resultCode: facts.resultCode,
...(facts.errorSummary === undefined
? {}
: { errorSummary: facts.errorSummary }),
},
{
expectedRunVersion: runVersion,
expectedRunEventSequence: eventSequence,
eventId,
dedupeKey: eventId,
actor: { type: 'system' },
},
);
result.push(mutation);
runVersion += 1;
eventSequence += 1;
}
if (result.length < 1) return conflict('no non-terminal Step remains');
return Object.freeze(result);
}
export async function terminalizeCopilotFailureDiagnosisBeforeModel(
requestId: string,
trigger: CopilotFailureDiagnosisPreModelTerminalizationTrigger,
dependencies: CopilotFailureDiagnosisPreModelTerminalizationDependencies,
): Promise<Readonly<CopilotFailureDiagnosisPreModelTerminalizationResult>> {
if (!/^[A-Za-z0-9][A-Za-z0-9._:-]{0,127}$/.test(requestId))
return invalid('request id is invalid');
if (!trigger || typeof trigger !== 'object' || Array.isArray(trigger))
return invalid('trigger is invalid');
assertDependencies(dependencies);
try {
const existing = await dependencies.repository.findByRequestId(requestId);
if (existing)
return Object.freeze({ status: 'existing' as const, receipt: existing });
const authority = await dependencies.repository.readAuthority(requestId);
if (
authority.plan.requestId !== requestId ||
authority.run.id !== authority.plan.runId ||
authority.run.status !== 'running' ||
authority.run.version !== authority.run.eventSequence ||
authority.toolStep.id !== authority.plan.toolStepRunId ||
authority.toolStep.runId !== authority.plan.runId ||
authority.toolStep.kind !== 'tool' ||
authority.modelStep.id !== authority.plan.modelStepRunId ||
authority.modelStep.runId !== authority.plan.runId ||
authority.modelStep.kind !== 'model' ||
authority.modelStartExists ||
!Number.isSafeInteger(authority.observedAtMs) ||
authority.observedAtMs < authority.plan.plannedAtMs
)
return conflict('pre-Model authority is invalid');
const evidence = triggerEvidence(trigger, authority);
const mapping = copilotFailureDiagnosisPreModelTerminalizationMapping(
evidence.reason,
authority.run.cancelReason,
);
const stepMutations = mutations(
authority,
evidence.reason,
mapping.outcome,
);
const receipt = createCopilotFailureDiagnosisPreModelTerminalizationReceipt(
{
requestId,
planDigest: authority.plan.planDigest,
runId: authority.plan.runId,
stage: mapping.stage,
reason: evidence.reason,
outcome: mapping.outcome,
evidenceDigest: evidence.evidenceDigest,
toolStartId: evidence.toolStartId,
toolCompletionDigest: evidence.toolCompletionDigest,
terminalSteps: stepMutations.map(terminalStepReference),
finalRunVersion: authority.run.version + stepMutations.length + 1,
finalRunEventSequence:
authority.run.eventSequence + stepMutations.length + 1,
finalizedAtMs: authority.observedAtMs,
},
);
const command = createCopilotFailureDiagnosisPreModelTerminalizationCommand(
{
plan: authority.plan,
expectedRunVersion: authority.run.version,
expectedRunEventSequence: authority.run.eventSequence,
stepMutations,
receipt,
},
);
return dependencies.repository.commit(command);
} catch (cause) {
if (
cause instanceof
InvalidCopilotFailureDiagnosisPreModelTerminalizationError ||
cause instanceof
CopilotFailureDiagnosisPreModelTerminalizationConflictError ||
cause instanceof
CopilotFailureDiagnosisPreModelTerminalizationNotReadyError ||
cause instanceof
CopilotFailureDiagnosisPreModelTerminalizationUnavailableError
)
throw cause;
return unavailable(cause);
}
}
@@ -0,0 +1,585 @@
import type { PostgresClient, PostgresPool } from '@qinglong/runtime-core';
import { BUILTIN_RUN_LOG_EXCERPT_TIMEOUT_SECONDS } from '@qinglong/runtime-core/builtin-run-log-excerpt-tool';
import {
normalizeStepRunMutation,
normalizeStepRunRecord,
type StepRunMutation,
type StepRunRecord,
} from '@qinglong/runtime-core/step-run';
import { normalizeCopilotFailureDiagnosisExecutionPlan } from '../admission/plan';
import type { CopilotFailureDiagnosisExecutionPlan } from '../admission/contracts';
import {
CopilotFailureDiagnosisPreModelTerminalizationConflictError,
CopilotFailureDiagnosisPreModelTerminalizationUnavailableError,
type CopilotFailureDiagnosisPreModelTerminalizationAuthority,
type CopilotFailureDiagnosisPreModelTerminalizationCommand,
type CopilotFailureDiagnosisPreModelTerminalizationReceipt,
type CopilotFailureDiagnosisPreModelTerminalizationRepository,
} from './contracts';
import {
normalizeCopilotFailureDiagnosisPreModelTerminalizationCommand,
normalizeCopilotFailureDiagnosisPreModelTerminalizationReceipt,
} from './protocol';
const TABLE = '"ql3_ai"."copilot_failure_diagnosis_pre_model_terminalizations"';
const REQUEST_ID_PATTERN = /^[A-Za-z0-9][A-Za-z0-9._:-]{0,127}$/;
function conflict(message?: string): never {
throw new CopilotFailureDiagnosisPreModelTerminalizationConflictError(
message,
);
}
function unavailable(cause?: unknown): never {
throw new CopilotFailureDiagnosisPreModelTerminalizationUnavailableError({
cause: cause instanceof Error ? cause : undefined,
});
}
function object(value: unknown): Record<string, unknown> {
if (value && typeof value === 'object' && !Array.isArray(value)) {
return value as Record<string, unknown>;
}
if (typeof value === 'string') {
try {
const parsed = JSON.parse(value) as unknown;
if (parsed && typeof parsed === 'object' && !Array.isArray(parsed)) {
return parsed as Record<string, unknown>;
}
} catch {
// Mapped to the generic durable conflict below.
}
}
return conflict();
}
function integer(value: unknown): number {
if (typeof value === 'number' && Number.isSafeInteger(value) && value >= 0) {
return value;
}
if (typeof value === 'string' && /^(0|[1-9]\d*)$/.test(value)) {
const parsed = Number(value);
if (Number.isSafeInteger(parsed)) return parsed;
}
return conflict();
}
function string(value: unknown): string {
return typeof value === 'string' ? value : conflict();
}
function optionalString(value: unknown): string | undefined {
return value === null || value === undefined ? undefined : string(value);
}
function same(left: unknown, right: unknown): boolean {
return JSON.stringify(left) === JSON.stringify(right);
}
async function authority(
queryable: Pick<PostgresPool, 'query'> | Pick<PostgresClient, 'query'>,
requestId: string,
lock: boolean,
): Promise<Readonly<CopilotFailureDiagnosisPreModelTerminalizationAuthority>> {
const result = await queryable.query<Record<string, unknown>>(
`SELECT admission.plan_json AS "planJson",
run.id AS "runId", run.project_id AS "projectId",
run.status AS "runStatus", run.version AS "runVersion",
run.event_sequence AS "runEventSequence",
run.cancel_requested_at_ms AS "cancelRequestedAtMs",
run.cancel_reason AS "cancelReason",
tool_step.step_run_json AS "toolStepJson",
model_step.step_run_json AS "modelStepJson",
(model_start.invocation_id IS NOT NULL) AS "modelStartExists",
floor(extract(epoch FROM clock_timestamp()) * 1000)::bigint
AS "observedAtMs"
FROM "ql3_ai"."copilot_failure_diagnosis_admissions" AS admission
JOIN "ql3"."runs" AS run ON run.id = admission.run_id
JOIN "ql3"."step_runs" AS tool_step
ON tool_step.run_id = run.id
AND tool_step.id = admission.tool_step_run_id
JOIN "ql3"."step_runs" AS model_step
ON model_step.run_id = run.id
AND model_step.id = admission.model_step_run_id
LEFT JOIN "ql3_ai"."model_invocation_starts" AS model_start
ON model_start.invocation_id = admission.plan_json->>'modelInvocationId'
WHERE admission.request_id = $1
${lock ? 'FOR UPDATE OF run, tool_step, model_step' : ''}`,
[requestId],
);
if (result.rows.length !== 1)
return conflict('terminalization authority is absent');
const row = result.rows[0]!;
const plan = normalizeCopilotFailureDiagnosisExecutionPlan(
object(row.planJson) as unknown as CopilotFailureDiagnosisExecutionPlan,
);
const toolStep = normalizeStepRunRecord(
object(row.toolStepJson) as unknown as StepRunRecord,
);
const modelStep = normalizeStepRunRecord(
object(row.modelStepJson) as unknown as StepRunRecord,
);
return Object.freeze({
plan,
run: Object.freeze({
id: string(row.runId),
projectId: string(row.projectId),
status: string(row.runStatus) as never,
version: integer(row.runVersion),
eventSequence: integer(row.runEventSequence),
...(row.cancelRequestedAtMs === null
? {}
: { cancelRequestedAtMs: integer(row.cancelRequestedAtMs) }),
...(row.cancelReason === null
? {}
: { cancelReason: string(row.cancelReason) as never }),
}),
toolStep,
modelStep,
modelStartExists: row.modelStartExists === true,
observedAtMs: integer(row.observedAtMs),
});
}
async function updateStep(
client: PostgresClient,
mutationValue: Readonly<StepRunMutation>,
): Promise<void> {
const mutation = normalizeStepRunMutation(mutationValue);
const step = mutation.stepRun;
const updated = await client.query(
`UPDATE "ql3"."step_runs"
SET status = $1, version = $2, attempt_count = $3,
output_ref = $4, approval_request_id = $5, ready_at_ms = $6,
started_at_ms = $7, finished_at_ms = $8, result_code = $9,
error_summary = $10, updated_at_ms = $11,
last_mutation_id = $12, step_run_digest = $13,
step_run_json = $14::jsonb
WHERE id = $15 AND run_id = $16 AND version = $17
AND step_run_digest = $18 AND status = $19`,
[
step.status,
step.version,
step.attemptCount,
step.outputRef,
step.approvalRequestId,
step.readyAtMs,
step.startedAtMs,
step.finishedAtMs,
step.resultCode,
step.errorSummary,
step.updatedAtMs,
step.lastMutationId,
step.stepRunDigest,
JSON.stringify(step),
step.id,
step.runId,
mutation.expectedStepRunVersion,
mutation.expectedStepRunDigest,
mutation.previousStatus,
],
);
if ((updated.rowCount ?? updated.rows.length) !== 1) return conflict();
const run = await client.query(
`UPDATE "ql3"."runs"
SET version = version + 1, event_sequence = event_sequence + 1
WHERE id = $1 AND status = 'running'
AND version = $2 AND event_sequence = $3`,
[
mutation.runId,
mutation.expectedRunVersion,
mutation.expectedRunEventSequence,
],
);
if ((run.rowCount ?? run.rows.length) !== 1) return conflict();
const event = mutation.event;
await client.query(
`INSERT INTO "ql3"."run_events" (
id, run_id, sequence, type, dedupe_key, actor_type, actor_id,
attempt_id, step_run_id, payload, created_at_ms
) VALUES ($1, $2, $3, $4, $5, $6, $7, NULL, $8, $9::jsonb, $10)`,
[
event.id,
event.runId,
event.sequence,
event.type,
event.dedupeKey,
event.actorType,
event.actorId ?? null,
step.id,
JSON.stringify(event.payload),
event.createdAtMs,
],
);
await client.query(
`INSERT INTO "ql3"."step_run_mutations" (
mutation_id, mutation_digest, run_id, step_run_id, step_run_digest,
event_id, event_sequence, run_version, step_run_json, committed_at_ms
) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9::jsonb,
floor(extract(epoch FROM transaction_timestamp()) * 1000)::bigint)`,
[
mutation.mutationId,
mutation.mutationDigest,
mutation.runId,
step.id,
step.stepRunDigest,
event.id,
event.sequence,
mutation.expectedRunVersion + 1,
JSON.stringify(step),
],
);
}
function failure(
receipt: Readonly<CopilotFailureDiagnosisPreModelTerminalizationReceipt>,
): Readonly<{
code: string;
summary: string;
}> {
if (receipt.reason === 'tool_failed') {
return Object.freeze({
code: 'COPILOT_DIAGNOSIS_TOOL_FAILED',
summary: 'Copilot diagnosis Tool failed',
});
}
if (receipt.reason === 'tool_timed_out' || receipt.outcome === 'timed_out') {
return Object.freeze({
code: 'COPILOT_DIAGNOSIS_TIMED_OUT',
summary: 'Copilot diagnosis timed out',
});
}
if (receipt.stage === 'log') {
return Object.freeze({
code: 'COPILOT_DIAGNOSIS_LOG_UNAVAILABLE',
summary: 'Copilot diagnosis log is unavailable',
});
}
return Object.freeze({
code: 'COPILOT_DIAGNOSIS_CANCELLED',
summary: 'Copilot diagnosis was cancelled',
});
}
export class PostgresCopilotFailureDiagnosisPreModelTerminalizationRepository
implements CopilotFailureDiagnosisPreModelTerminalizationRepository
{
constructor(private readonly pool: PostgresPool) {
if (
!pool ||
typeof pool.query !== 'function' ||
typeof pool.connect !== 'function'
) {
return unavailable();
}
}
async findByRequestId(
requestId: string,
): Promise<Readonly<CopilotFailureDiagnosisPreModelTerminalizationReceipt> | null> {
if (!REQUEST_ID_PATTERN.test(requestId))
return conflict('request id is invalid');
try {
const result = await this.pool.query<Record<string, unknown>>(
`SELECT receipt_json AS "receiptJson" FROM ${TABLE} WHERE request_id = $1`,
[requestId],
);
if (result.rows.length > 1) return conflict();
if (!result.rows[0]) return null;
const receipt =
normalizeCopilotFailureDiagnosisPreModelTerminalizationReceipt(
object(
result.rows[0].receiptJson,
) as unknown as CopilotFailureDiagnosisPreModelTerminalizationReceipt,
);
const durable = await this.pool.query<Record<string, unknown>>(
`SELECT run.status AS "runStatus", run.version AS "runVersion",
run.event_sequence AS "runEventSequence",
run.finished_at_ms AS "finishedAtMs",
event.type AS "eventType", event.payload,
event.created_at_ms AS "eventCreatedAtMs"
FROM "ql3"."runs" AS run
JOIN "ql3"."run_events" AS event
ON event.run_id = run.id AND event.id = $1
WHERE run.id = $2`,
[receipt.runEventId, receipt.runId],
);
if (durable.rows.length !== 1) return conflict();
const proof = durable.rows[0]!;
const payload = object(proof.payload);
if (
string(proof.runStatus) !== receipt.outcome ||
integer(proof.runVersion) !== receipt.finalRunVersion ||
integer(proof.runEventSequence) !== receipt.finalRunEventSequence ||
integer(proof.finishedAtMs) !== receipt.finalizedAtMs ||
string(proof.eventType) !== `copilot.diagnosis.${receipt.outcome}` ||
integer(proof.eventCreatedAtMs) !== receipt.finalizedAtMs ||
payload.requestId !== receipt.requestId ||
payload.planDigest !== receipt.planDigest ||
payload.reason !== receipt.reason ||
payload.evidenceDigest !== receipt.evidenceDigest
)
return conflict();
for (const step of receipt.terminalSteps) {
const stepProof = await this.pool.query<Record<string, unknown>>(
`SELECT step.status, step.version,
mutation.mutation_digest AS "mutationDigest",
mutation.event_id AS "eventId"
FROM "ql3"."step_runs" AS step
JOIN "ql3"."step_run_mutations" AS mutation
ON mutation.mutation_id = $1 AND mutation.step_run_id = step.id
WHERE step.id = $2 AND step.run_id = $3`,
[step.mutationId, step.stepRunId, receipt.runId],
);
if (
stepProof.rows.length !== 1 ||
string(stepProof.rows[0]!.status) !== step.status ||
integer(stepProof.rows[0]!.version) !== step.version ||
string(stepProof.rows[0]!.mutationDigest) !== step.mutationDigest ||
string(stepProof.rows[0]!.eventId) !== step.eventId
)
return conflict();
}
return receipt;
} catch (cause) {
if (
cause instanceof
CopilotFailureDiagnosisPreModelTerminalizationConflictError
)
throw cause;
return unavailable(cause);
}
}
async readAuthority(
requestId: string,
): Promise<
Readonly<CopilotFailureDiagnosisPreModelTerminalizationAuthority>
> {
if (!REQUEST_ID_PATTERN.test(requestId))
return conflict('request id is invalid');
try {
return await authority(this.pool, requestId, false);
} catch (cause) {
if (
cause instanceof
CopilotFailureDiagnosisPreModelTerminalizationConflictError
)
throw cause;
return unavailable(cause);
}
}
async commit(
commandValue: Readonly<CopilotFailureDiagnosisPreModelTerminalizationCommand>,
): Promise<
Readonly<{
status: 'created' | 'existing';
receipt: Readonly<CopilotFailureDiagnosisPreModelTerminalizationReceipt>;
}>
> {
const command =
normalizeCopilotFailureDiagnosisPreModelTerminalizationCommand(
commandValue,
);
const existing = await this.findByRequestId(command.plan.requestId);
if (existing) {
if (existing.receiptDigest !== command.receipt.receiptDigest)
return conflict();
return Object.freeze({ status: 'existing' as const, receipt: existing });
}
for (let attempt = 0; attempt < 3; attempt += 1) {
let client: PostgresClient;
try {
client = await this.pool.connect();
} catch (cause) {
return unavailable(cause);
}
let began = false;
try {
await client.query('BEGIN ISOLATION LEVEL SERIALIZABLE');
began = true;
await client.query(`SELECT set_config('statement_timeout', $1, true)`, [
'5s',
]);
await client.query(`SELECT set_config('lock_timeout', $1, true)`, [
'2s',
]);
const current = await authority(client, command.plan.requestId, true);
if (
!same(current.plan, command.plan) ||
current.run.status !== 'running' ||
current.run.version !== command.expectedRunVersion ||
current.run.eventSequence !== command.expectedRunEventSequence ||
current.modelStartExists ||
current.observedAtMs < command.receipt.finalizedAtMs ||
(command.receipt.reason === 'deadline_exceeded' &&
current.observedAtMs < current.plan.deadlineAtMs) ||
(command.receipt.reason === 'tool_budget_exhausted' &&
(current.toolStep.status !== 'ready' ||
current.observedAtMs +
BUILTIN_RUN_LOG_EXCERPT_TIMEOUT_SECONDS * 1_000 <=
current.plan.deadlineAtMs)) ||
(command.receipt.reason === 'cancellation_requested' &&
(current.run.cancelRequestedAtMs === undefined ||
current.run.cancelReason === undefined))
)
return conflict();
if (command.receipt.stage === 'tool') {
const proof = await client.query<Record<string, unknown>>(
`SELECT outcome, completion_digest AS "completionDigest"
FROM "ql3"."tool_execution_failure_completions"
WHERE start_id = $1 AND run_id = $2 AND step_run_id = $3`,
[
command.receipt.toolStartId,
command.plan.runId,
command.plan.toolStepRunId,
],
);
if (
proof.rows.length !== 1 ||
string(proof.rows[0]!.completionDigest) !==
command.receipt.toolCompletionDigest ||
string(proof.rows[0]!.outcome) !== command.receipt.outcome
)
return conflict();
} else if (command.receipt.stage === 'log') {
const proof = await client.query<Record<string, unknown>>(
`SELECT completion.completion_digest AS "completionDigest",
unlock.tool_completion_digest AS "unlockDigest"
FROM "ql3"."tool_execution_completions" AS completion
JOIN "ql3_ai"."copilot_failure_diagnosis_tool_unlocks" AS unlock
ON unlock.request_id = $1 AND unlock.start_id = completion.start_id
WHERE completion.start_id = $2 AND completion.run_id = $3
AND completion.step_run_id = $4`,
[
command.plan.requestId,
command.receipt.toolStartId,
command.plan.runId,
command.plan.toolStepRunId,
],
);
if (
proof.rows.length !== 1 ||
string(proof.rows[0]!.completionDigest) !==
command.receipt.toolCompletionDigest ||
string(proof.rows[0]!.unlockDigest) !==
command.receipt.toolCompletionDigest
)
return conflict();
}
for (const mutation of command.stepMutations)
await updateStep(client, mutation);
const failureFact = failure(command.receipt);
const updated = await client.query(
`UPDATE "ql3"."runs"
SET status = $1, version = $2, event_sequence = $3,
output_ref = NULL, finished_at_ms = $4,
error_code = $5, error_summary = $6
WHERE id = $7 AND status = 'running'
AND version = $8 AND event_sequence = $9`,
[
command.receipt.outcome,
command.receipt.finalRunVersion,
command.receipt.finalRunEventSequence,
command.receipt.finalizedAtMs,
failureFact.code,
failureFact.summary,
command.plan.runId,
command.receipt.finalRunVersion - 1,
command.receipt.finalRunEventSequence - 1,
],
);
if ((updated.rowCount ?? updated.rows.length) !== 1) return conflict();
const payload = JSON.stringify({
requestId: command.receipt.requestId,
planDigest: command.receipt.planDigest,
reason: command.receipt.reason,
evidenceDigest: command.receipt.evidenceDigest,
outcome: command.receipt.outcome,
});
await client.query(
`INSERT INTO "ql3"."run_events" (
id, run_id, sequence, type, dedupe_key, actor_type, actor_id,
attempt_id, step_run_id, payload, created_at_ms
) VALUES ($1, $2, $3, $4, $1, 'system', NULL, NULL, NULL,
$5::jsonb, $6)`,
[
command.receipt.runEventId,
command.receipt.runId,
command.receipt.finalRunEventSequence,
`copilot.diagnosis.${command.receipt.outcome}`,
payload,
command.receipt.finalizedAtMs,
],
);
await client.query(
`INSERT INTO ${TABLE} (
request_id, plan_digest, run_id, stage, reason, outcome,
evidence_digest, tool_start_id, tool_completion_digest,
terminal_steps_json, final_run_version,
final_run_event_sequence, run_event_id, finalized_at_ms,
receipt_digest, receipt_json
) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10::jsonb,
$11, $12, $13, $14, $15, $16::jsonb)`,
[
command.receipt.requestId,
command.receipt.planDigest,
command.receipt.runId,
command.receipt.stage,
command.receipt.reason,
command.receipt.outcome,
command.receipt.evidenceDigest,
command.receipt.toolStartId,
command.receipt.toolCompletionDigest,
JSON.stringify(command.receipt.terminalSteps),
command.receipt.finalRunVersion,
command.receipt.finalRunEventSequence,
command.receipt.runEventId,
command.receipt.finalizedAtMs,
command.receipt.receiptDigest,
JSON.stringify(command.receipt),
],
);
await client.query('COMMIT');
began = false;
return Object.freeze({
status: 'created' as const,
receipt: command.receipt,
});
} catch (cause) {
if (began) {
try {
await client.query('ROLLBACK');
} catch {
/* preserve */
}
}
const state =
cause && typeof cause === 'object' && 'code' in cause
? String(cause.code)
: '';
if ((state === '40001' || state === '40P01') && attempt < 2) continue;
const recovered = await this.findByRequestId(command.plan.requestId);
if (recovered) {
if (recovered.receiptDigest !== command.receipt.receiptDigest)
return conflict();
return Object.freeze({
status: 'existing' as const,
receipt: recovered,
});
}
if (
cause instanceof
CopilotFailureDiagnosisPreModelTerminalizationConflictError
)
throw cause;
return unavailable(cause);
} finally {
client.release();
}
}
return unavailable();
}
}
@@ -0,0 +1,471 @@
import { createHash } from 'node:crypto';
import {
normalizeStepRunMutation,
type StepRunMutation,
} from '@qinglong/runtime-core/step-run';
import { normalizeCopilotFailureDiagnosisExecutionPlan } from '../admission/plan';
import {
COPILOT_FAILURE_DIAGNOSIS_PRE_MODEL_TERMINALIZATION_COMMAND_SCHEMA,
COPILOT_FAILURE_DIAGNOSIS_PRE_MODEL_TERMINALIZATION_REASONS,
COPILOT_FAILURE_DIAGNOSIS_PRE_MODEL_TERMINALIZATION_SCHEMA,
InvalidCopilotFailureDiagnosisPreModelTerminalizationError,
type CopilotFailureDiagnosisPreModelTerminalizationCommand,
type CopilotFailureDiagnosisPreModelTerminalizationOutcome,
type CopilotFailureDiagnosisPreModelTerminalizationReason,
type CopilotFailureDiagnosisPreModelTerminalizationReceipt,
type CopilotFailureDiagnosisPreModelTerminalizationStage,
type CopilotFailureDiagnosisTerminalStepReference,
} from './contracts';
const ID_PATTERN = /^[A-Za-z0-9][A-Za-z0-9._:/-]{0,127}$/;
const RUN_ID_PATTERN = /^[A-Za-z0-9][A-Za-z0-9._:-]{0,35}$/;
const DIGEST_PATTERN = /^[0-9a-f]{64}$/;
const RECEIPT_DOMAIN =
'qinglong/copilot-failure-diagnosis-pre-model-terminalization-receipt@v1\0';
const COMMAND_DOMAIN =
'qinglong/copilot-failure-diagnosis-pre-model-terminalization-command@v1\0';
const IDENTITY_DOMAIN =
'qinglong/copilot-failure-diagnosis-pre-model-terminalization-identity@v1\0';
function invalid(message: string): never {
throw new InvalidCopilotFailureDiagnosisPreModelTerminalizationError(message);
}
function exact(
value: object,
expected: readonly string[],
label: string,
): void {
const keys = Reflect.ownKeys(value);
if (
keys.length !== expected.length ||
keys.some((key) => typeof key !== 'string' || !expected.includes(key))
) {
return invalid(`${label} shape is invalid`);
}
}
function text(value: unknown, pattern: RegExp, label: string): string {
if (typeof value !== 'string' || !pattern.test(value)) {
return invalid(`${label} is invalid`);
}
return value;
}
function integer(value: unknown, minimum: number, label: string): number {
if (!Number.isSafeInteger(value) || (value as number) < minimum) {
return invalid(`${label} is invalid`);
}
return value as number;
}
function hash(domain: string, value: unknown): string {
return createHash('sha256')
.update(domain)
.update(JSON.stringify(value))
.digest('hex');
}
export function copilotFailureDiagnosisPreModelTerminalizationMapping(
reason: CopilotFailureDiagnosisPreModelTerminalizationReason,
cancellationReason?: string,
): Readonly<{
stage: CopilotFailureDiagnosisPreModelTerminalizationStage;
outcome: CopilotFailureDiagnosisPreModelTerminalizationOutcome;
}> {
if (
!COPILOT_FAILURE_DIAGNOSIS_PRE_MODEL_TERMINALIZATION_REASONS.includes(
reason,
)
) {
return invalid('reason is invalid');
}
if (reason === 'tool_failed')
return Object.freeze({ stage: 'tool', outcome: 'failed' });
if (reason === 'tool_timed_out') {
return Object.freeze({ stage: 'tool', outcome: 'timed_out' });
}
if (reason.startsWith('log_'))
return Object.freeze({ stage: 'log', outcome: 'failed' });
if (reason === 'tool_budget_exhausted' || reason === 'deadline_exceeded') {
return Object.freeze({ stage: 'deadline', outcome: 'timed_out' });
}
return Object.freeze({
stage: 'cancellation',
outcome: cancellationReason === 'timeout' ? 'timed_out' : 'cancelled',
});
}
export function copilotFailureDiagnosisTerminalizationIdentity(
prefix: 'mutation' | 'step-event' | 'run-event',
planDigest: string,
reason: CopilotFailureDiagnosisPreModelTerminalizationReason,
target: string,
): string {
const digest = hash(IDENTITY_DOMAIN, { prefix, planDigest, reason, target });
if (prefix !== 'run-event' && prefix !== 'step-event') {
return `cdx:${digest.slice(0, 31)}`;
}
const hex = digest.slice(0, 32).split('');
hex[12] = '4';
hex[16] = '8';
const value = hex.join('');
return `${value.slice(0, 8)}-${value.slice(8, 12)}-${value.slice(
12,
16,
)}-${value.slice(16, 20)}-${value.slice(20)}`;
}
function terminalStep(
value: CopilotFailureDiagnosisTerminalStepReference,
): Readonly<CopilotFailureDiagnosisTerminalStepReference> {
exact(
value,
[
'eventId',
'mutationDigest',
'mutationId',
'status',
'stepRunId',
'version',
],
'terminal Step',
);
if (!['failed', 'timed_out', 'cancelled'].includes(value.status)) {
return invalid('terminal Step status is invalid');
}
return Object.freeze({
stepRunId: text(value.stepRunId, ID_PATTERN, 'StepRun id'),
status: value.status,
version: integer(value.version, 1, 'StepRun version'),
mutationId: text(value.mutationId, ID_PATTERN, 'mutation id'),
mutationDigest: text(
value.mutationDigest,
DIGEST_PATTERN,
'mutation digest',
),
eventId: text(value.eventId, RUN_ID_PATTERN, 'event id'),
});
}
function unsignedReceipt(
value: Omit<
CopilotFailureDiagnosisPreModelTerminalizationReceipt,
'receiptDigest'
>,
): object {
return {
schema: value.schema,
requestId: value.requestId,
planDigest: value.planDigest,
runId: value.runId,
stage: value.stage,
reason: value.reason,
outcome: value.outcome,
evidenceDigest: value.evidenceDigest,
toolStartId: value.toolStartId,
toolCompletionDigest: value.toolCompletionDigest,
terminalSteps: value.terminalSteps,
finalRunVersion: value.finalRunVersion,
finalRunEventSequence: value.finalRunEventSequence,
runEventId: value.runEventId,
finalizedAtMs: value.finalizedAtMs,
};
}
export function copilotFailureDiagnosisPreModelTerminalizationReceiptDigest(
value: Omit<
CopilotFailureDiagnosisPreModelTerminalizationReceipt,
'receiptDigest'
>,
): string {
return hash(RECEIPT_DOMAIN, unsignedReceipt(value));
}
export function normalizeCopilotFailureDiagnosisPreModelTerminalizationReceipt(
value: CopilotFailureDiagnosisPreModelTerminalizationReceipt,
): Readonly<CopilotFailureDiagnosisPreModelTerminalizationReceipt> {
if (!value || typeof value !== 'object' || Array.isArray(value))
return invalid('receipt');
exact(
value,
[
'evidenceDigest',
'finalRunEventSequence',
'finalRunVersion',
'finalizedAtMs',
'outcome',
'planDigest',
'reason',
'receiptDigest',
'requestId',
'runEventId',
'runId',
'schema',
'stage',
'terminalSteps',
'toolCompletionDigest',
'toolStartId',
],
'receipt',
);
if (
value.schema !== COPILOT_FAILURE_DIAGNOSIS_PRE_MODEL_TERMINALIZATION_SCHEMA
) {
return invalid('receipt schema is invalid');
}
const mapping = copilotFailureDiagnosisPreModelTerminalizationMapping(
value.reason,
);
if (
value.stage !== mapping.stage ||
(value.reason !== 'cancellation_requested' &&
value.outcome !== mapping.outcome) ||
(value.reason === 'cancellation_requested' &&
!['cancelled', 'timed_out'].includes(value.outcome)) ||
!Array.isArray(value.terminalSteps) ||
value.terminalSteps.length < 1 ||
value.terminalSteps.length > 2
) {
return invalid('receipt state is invalid');
}
const terminalSteps = Object.freeze(value.terminalSteps.map(terminalStep));
if (
new Set(terminalSteps.map((item) => item.stepRunId)).size !==
terminalSteps.length
) {
return invalid('terminal Step identities are duplicated');
}
const unsigned = Object.freeze({
schema: value.schema,
requestId: text(value.requestId, ID_PATTERN, 'request id'),
planDigest: text(value.planDigest, DIGEST_PATTERN, 'plan digest'),
runId: text(value.runId, RUN_ID_PATTERN, 'Run id'),
stage: value.stage,
reason: value.reason,
outcome: value.outcome,
evidenceDigest: text(
value.evidenceDigest,
DIGEST_PATTERN,
'evidence digest',
),
toolStartId:
value.toolStartId === null
? null
: text(value.toolStartId, ID_PATTERN, 'Tool start id'),
toolCompletionDigest:
value.toolCompletionDigest === null
? null
: text(
value.toolCompletionDigest,
DIGEST_PATTERN,
'Tool completion digest',
),
terminalSteps,
finalRunVersion: integer(value.finalRunVersion, 1, 'final Run version'),
finalRunEventSequence: integer(
value.finalRunEventSequence,
1,
'final Run event sequence',
),
runEventId: text(value.runEventId, RUN_ID_PATTERN, 'Run event id'),
finalizedAtMs: integer(value.finalizedAtMs, 0, 'finalized time'),
} satisfies Omit<CopilotFailureDiagnosisPreModelTerminalizationReceipt, 'receiptDigest'>);
if (
unsigned.finalRunVersion !== unsigned.finalRunEventSequence ||
unsigned.runEventId !==
copilotFailureDiagnosisTerminalizationIdentity(
'run-event',
unsigned.planDigest,
unsigned.reason,
unsigned.runId,
) ||
(unsigned.stage === 'tool' || unsigned.stage === 'log') !==
(unsigned.toolStartId !== null &&
unsigned.toolCompletionDigest !== null) ||
text(value.receiptDigest, DIGEST_PATTERN, 'receipt digest') !==
copilotFailureDiagnosisPreModelTerminalizationReceiptDigest(unsigned)
) {
return invalid('receipt evidence is invalid');
}
return Object.freeze({ ...unsigned, receiptDigest: value.receiptDigest });
}
export function createCopilotFailureDiagnosisPreModelTerminalizationReceipt(
value: Omit<
CopilotFailureDiagnosisPreModelTerminalizationReceipt,
'schema' | 'runEventId' | 'receiptDigest'
>,
): Readonly<CopilotFailureDiagnosisPreModelTerminalizationReceipt> {
const unsigned = Object.freeze({
schema: COPILOT_FAILURE_DIAGNOSIS_PRE_MODEL_TERMINALIZATION_SCHEMA,
...value,
runEventId: copilotFailureDiagnosisTerminalizationIdentity(
'run-event',
value.planDigest,
value.reason,
value.runId,
),
});
return normalizeCopilotFailureDiagnosisPreModelTerminalizationReceipt({
...unsigned,
receiptDigest:
copilotFailureDiagnosisPreModelTerminalizationReceiptDigest(unsigned),
});
}
function commandUnsigned(
value: Omit<
CopilotFailureDiagnosisPreModelTerminalizationCommand,
'commandDigest'
>,
): object {
return {
schema: value.schema,
plan: value.plan,
expectedRunVersion: value.expectedRunVersion,
expectedRunEventSequence: value.expectedRunEventSequence,
stepMutations: value.stepMutations,
receipt: value.receipt,
};
}
export function normalizeCopilotFailureDiagnosisPreModelTerminalizationCommand(
value: CopilotFailureDiagnosisPreModelTerminalizationCommand,
): Readonly<CopilotFailureDiagnosisPreModelTerminalizationCommand> {
if (!value || typeof value !== 'object' || Array.isArray(value))
return invalid('command');
exact(
value,
[
'commandDigest',
'expectedRunEventSequence',
'expectedRunVersion',
'plan',
'receipt',
'schema',
'stepMutations',
],
'command',
);
if (
value.schema !==
COPILOT_FAILURE_DIAGNOSIS_PRE_MODEL_TERMINALIZATION_COMMAND_SCHEMA
) {
return invalid('command schema is invalid');
}
const plan = normalizeCopilotFailureDiagnosisExecutionPlan(value.plan);
const receipt =
normalizeCopilotFailureDiagnosisPreModelTerminalizationReceipt(
value.receipt,
);
if (
!Array.isArray(value.stepMutations) ||
value.stepMutations.length !== receipt.terminalSteps.length
) {
return invalid('Step mutation count is invalid');
}
const stepMutations = Object.freeze(
value.stepMutations.map(normalizeStepRunMutation),
);
for (let index = 0; index < stepMutations.length; index += 1) {
const mutation = stepMutations[index]!;
const reference = receipt.terminalSteps[index]!;
if (
mutation.runId !== plan.runId ||
mutation.stepRun.id !== reference.stepRunId ||
mutation.stepRun.status !== reference.status ||
mutation.stepRun.version !== reference.version ||
mutation.mutationId !== reference.mutationId ||
mutation.mutationDigest !== reference.mutationDigest ||
mutation.event.id !== reference.eventId
)
return invalid('Step mutation evidence changed');
}
const expectedRunVersion = integer(
value.expectedRunVersion,
1,
'expected Run version',
);
const expectedRunEventSequence = integer(
value.expectedRunEventSequence,
1,
'expected Run event sequence',
);
if (
expectedRunVersion !== expectedRunEventSequence ||
receipt.requestId !== plan.requestId ||
receipt.planDigest !== plan.planDigest ||
receipt.runId !== plan.runId ||
receipt.finalRunVersion !== expectedRunVersion + stepMutations.length + 1
)
return invalid('command aggregate fence is invalid');
const unsigned = Object.freeze({
schema: value.schema,
plan,
expectedRunVersion,
expectedRunEventSequence,
stepMutations,
receipt,
});
if (
text(value.commandDigest, DIGEST_PATTERN, 'command digest') !==
hash(COMMAND_DOMAIN, commandUnsigned(unsigned))
) {
return invalid('command digest does not match');
}
return Object.freeze({ ...unsigned, commandDigest: value.commandDigest });
}
export function createCopilotFailureDiagnosisPreModelTerminalizationCommand(
value: Omit<
CopilotFailureDiagnosisPreModelTerminalizationCommand,
'schema' | 'commandDigest'
>,
): Readonly<CopilotFailureDiagnosisPreModelTerminalizationCommand> {
const unsigned = Object.freeze({
schema: COPILOT_FAILURE_DIAGNOSIS_PRE_MODEL_TERMINALIZATION_COMMAND_SCHEMA,
...value,
});
return normalizeCopilotFailureDiagnosisPreModelTerminalizationCommand({
...unsigned,
commandDigest: hash(COMMAND_DOMAIN, commandUnsigned(unsigned)),
});
}
export function copilotFailureDiagnosisPreModelEvidenceDigest(
value: Readonly<{
reason: CopilotFailureDiagnosisPreModelTerminalizationReason;
planDigest: string;
toolCompletionDigest?: string;
sourceStatus?: string;
deadlineAtMs?: number;
observedAtMs?: number;
requiredToolBudgetMs?: number;
cancelRequestedAtMs?: number;
cancelReason?: string;
}>,
): string {
return hash(`${IDENTITY_DOMAIN}evidence\0`, value);
}
export function terminalStepReference(
mutation: Readonly<StepRunMutation>,
): Readonly<CopilotFailureDiagnosisTerminalStepReference> {
const normalized = normalizeStepRunMutation(mutation);
if (
!['failed', 'timed_out', 'cancelled'].includes(normalized.stepRun.status)
) {
return invalid('Step mutation is not terminal');
}
return terminalStep({
stepRunId: normalized.stepRun.id,
status: normalized.stepRun.status as 'failed' | 'timed_out' | 'cancelled',
version: normalized.stepRun.version,
mutationId: normalized.mutationId,
mutationDigest: normalized.mutationDigest,
eventId: normalized.event.id,
});
}
@@ -5,6 +5,8 @@ import type {
ToolExecutionCompletionRecord,
ToolExecutionResultArtifactReference,
} from '@qinglong/runtime-core/tool-execution-completion';
import type { ToolJsonValue } from '@qinglong/runtime-core/tool-registry';
import type { ToolExecutionFailureCompletionRecord } from '@qinglong/runtime-core/tool-execution-failure-completion';
import type { ToolPolicyAuthorizer } from '@qinglong/runtime-core/tool-registry';
import type {
@@ -76,12 +78,14 @@ export type CopilotFailureDiagnosisToolExecutionResult =
completionStatus: 'created' | 'existing';
unlockStatus: 'created' | 'existing';
completion: Readonly<ToolExecutionCompletionRecord>;
output: ToolJsonValue;
unlock: Readonly<CopilotFailureDiagnosisToolUnlockReceipt>;
}>
| Readonly<{
outcome: 'failed' | 'timed_out';
completionStatus: 'created' | 'existing';
unlockStatus: null;
completion: Readonly<ToolExecutionFailureCompletionRecord>;
}>;
export class InvalidCopilotFailureDiagnosisToolExecutionError extends TypeError {
@@ -586,6 +586,7 @@ export async function executeCopilotFailureDiagnosisTool(
outcome: completed.outcome,
completionStatus: completed.status,
unlockStatus: null,
completion: completed.completion,
});
}
const unlock = await unlockModel(plan, completed.completion, dependencies);
@@ -594,6 +595,7 @@ export async function executeCopilotFailureDiagnosisTool(
completionStatus: completed.status,
unlockStatus: unlock.status,
completion: completed.completion,
output: completed.output,
unlock: unlock.receipt,
});
}
@@ -67,6 +67,8 @@ export const POSTGRES_COPILOT_FAILURE_DIAGNOSIS_TOOL_UNLOCK_MIGRATION_ID =
'pg-9019-ai-copilot-failure-diagnosis-tool-unlocks';
export const POSTGRES_COPILOT_FAILURE_DIAGNOSIS_MODEL_EXECUTION_MIGRATION_ID =
'pg-9020-ai-copilot-failure-diagnosis-model-executions';
export const POSTGRES_COPILOT_FAILURE_DIAGNOSIS_PRE_MODEL_TERMINALIZATION_MIGRATION_ID =
'pg-9021-ai-copilot-failure-diagnosis-pre-model-terminalizations';
export const POSTGRES_MODEL_INVOCATION_SCHEMA = 'ql3_ai';
export const LOCAL_MODEL_INVOCATION_MIGRATION_HISTORY_TABLE =
'QingLong3AiSchemaMigrations';
@@ -30,6 +30,7 @@ import {
POSTGRES_COPILOT_FAILURE_DIAGNOSIS_ADMISSION_MIGRATION_ID,
POSTGRES_COPILOT_FAILURE_DIAGNOSIS_TOOL_UNLOCK_MIGRATION_ID,
POSTGRES_COPILOT_FAILURE_DIAGNOSIS_MODEL_EXECUTION_MIGRATION_ID,
POSTGRES_COPILOT_FAILURE_DIAGNOSIS_PRE_MODEL_TERMINALIZATION_MIGRATION_ID,
POSTGRES_MODEL_INVOCATION_SCHEMA,
POSTGRES_MODEL_INVOCATION_MIGRATION_HISTORY_TABLE,
} from './identities';
@@ -67,6 +68,7 @@ const POSTGRES_HISTORY_IDENTITY = Object.freeze({
POSTGRES_COPILOT_FAILURE_DIAGNOSIS_ADMISSION_MIGRATION_ID,
POSTGRES_COPILOT_FAILURE_DIAGNOSIS_TOOL_UNLOCK_MIGRATION_ID,
POSTGRES_COPILOT_FAILURE_DIAGNOSIS_MODEL_EXECUTION_MIGRATION_ID,
POSTGRES_COPILOT_FAILURE_DIAGNOSIS_PRE_MODEL_TERMINALIZATION_MIGRATION_ID,
]),
streamId: POSTGRES_MODEL_INVOCATION_MIGRATION_STREAM_ID,
dialect: 'postgresql' as const,
@@ -4,6 +4,7 @@ import {
POSTGRES_COPILOT_FAILURE_DIAGNOSIS_ADMISSION_MIGRATION_ID,
POSTGRES_COPILOT_FAILURE_DIAGNOSIS_TOOL_UNLOCK_MIGRATION_ID,
POSTGRES_COPILOT_FAILURE_DIAGNOSIS_MODEL_EXECUTION_MIGRATION_ID,
POSTGRES_COPILOT_FAILURE_DIAGNOSIS_PRE_MODEL_TERMINALIZATION_MIGRATION_ID,
POSTGRES_MODEL_INVOCATION_SCHEMA,
} from '../identities';
import { defineSqlMigration } from '../shared';
@@ -14,6 +15,8 @@ const SOURCE_SNAPSHOT_FUNCTION =
const TOOL_UNLOCK_TABLE = 'copilot_failure_diagnosis_tool_unlocks';
const MODEL_OUTPUT_TABLE = 'copilot_failure_diagnosis_model_outputs';
const FINALIZATION_TABLE = 'copilot_failure_diagnosis_finalizations';
const PRE_MODEL_TERMINALIZATION_TABLE =
'copilot_failure_diagnosis_pre_model_terminalizations';
const POSTGRES_COPILOT_FAILURE_DIAGNOSIS_ADMISSION_TABLE_SQL = `
CREATE TABLE "${POSTGRES_MODEL_INVOCATION_SCHEMA}"."${ADMISSION_TABLE}" (
@@ -480,8 +483,102 @@ const postgresCopilotFailureDiagnosisModelExecutionMigration =
(context, statement) => context.query(statement).then(() => undefined),
);
const POSTGRES_COPILOT_FAILURE_DIAGNOSIS_PRE_MODEL_TERMINALIZATION_TABLE_SQL = `
CREATE TABLE "${POSTGRES_MODEL_INVOCATION_SCHEMA}"."${PRE_MODEL_TERMINALIZATION_TABLE}" (
request_id varchar(128) PRIMARY KEY,
plan_digest char(64) NOT NULL UNIQUE,
run_id varchar(36) NOT NULL UNIQUE,
stage varchar(16) NOT NULL,
reason varchar(32) NOT NULL,
outcome varchar(16) NOT NULL,
evidence_digest char(64) NOT NULL,
tool_start_id varchar(128),
tool_completion_digest char(64),
terminal_steps_json jsonb NOT NULL,
final_run_version integer NOT NULL,
final_run_event_sequence integer NOT NULL,
run_event_id varchar(36) NOT NULL UNIQUE,
finalized_at_ms bigint NOT NULL,
receipt_digest char(64) NOT NULL UNIQUE,
receipt_json jsonb NOT NULL,
CONSTRAINT ql3_ai_copilot_pre_model_terminalization_admission_fk
FOREIGN KEY (request_id)
REFERENCES "${POSTGRES_MODEL_INVOCATION_SCHEMA}"."${ADMISSION_TABLE}"
(request_id) ON DELETE RESTRICT,
CONSTRAINT ql3_ai_copilot_pre_model_terminalization_run_fk
FOREIGN KEY (run_id) REFERENCES "ql3"."runs" (id) ON DELETE RESTRICT,
CONSTRAINT ql3_ai_copilot_pre_model_terminalization_event_fk
FOREIGN KEY (run_event_id) REFERENCES "ql3"."run_events" (id) ON DELETE RESTRICT,
CONSTRAINT ql3_ai_copilot_pre_model_terminalization_state_check CHECK (
stage IN ('tool', 'log', 'deadline', 'cancellation') AND
reason IN (
'tool_failed', 'tool_timed_out', 'log_not_found', 'log_pending',
'log_missing', 'log_retired', 'tool_budget_exhausted',
'deadline_exceeded',
'cancellation_requested'
) AND
outcome IN ('failed', 'timed_out', 'cancelled') AND
((stage = 'tool' AND reason IN ('tool_failed', 'tool_timed_out')) OR
(stage = 'log' AND reason IN (
'log_not_found', 'log_pending', 'log_missing', 'log_retired'
)) OR
(stage = 'deadline' AND reason IN (
'tool_budget_exhausted', 'deadline_exceeded'
)) OR
(stage = 'cancellation' AND reason = 'cancellation_requested')) AND
((stage IN ('tool', 'log') AND tool_start_id IS NOT NULL AND
tool_completion_digest IS NOT NULL) OR
(stage IN ('deadline', 'cancellation') AND tool_start_id IS NULL AND
tool_completion_digest IS NULL)) AND
final_run_version >= 1 AND
final_run_event_sequence = final_run_version AND
finalized_at_ms >= 0
),
CONSTRAINT ql3_ai_copilot_pre_model_terminalization_digest_check CHECK (
plan_digest ~ '^[0-9a-f]{64}$' AND
evidence_digest ~ '^[0-9a-f]{64}$' AND
(tool_completion_digest IS NULL OR
tool_completion_digest ~ '^[0-9a-f]{64}$') AND
receipt_digest ~ '^[0-9a-f]{64}$'
),
CONSTRAINT ql3_ai_copilot_pre_model_terminalization_json_check CHECK (
jsonb_typeof(terminal_steps_json) = 'array' AND
jsonb_array_length(terminal_steps_json) BETWEEN 1 AND 2 AND
jsonb_typeof(receipt_json) = 'object' AND
octet_length(receipt_json::text) BETWEEN 2 AND 32768 AND
receipt_json @> jsonb_build_object(
'schema',
'qinglong/copilot-failure-diagnosis-pre-model-terminalization@v1',
'requestId', request_id, 'planDigest', plan_digest,
'runId', run_id, 'stage', stage, 'reason', reason,
'outcome', outcome, 'evidenceDigest', evidence_digest,
'terminalSteps', terminal_steps_json,
'finalRunVersion', final_run_version,
'finalRunEventSequence', final_run_event_sequence,
'runEventId', run_event_id, 'finalizedAtMs', finalized_at_ms,
'receiptDigest', receipt_digest
)
)
)`;
const postgresCopilotFailureDiagnosisPreModelTerminalizationMigration =
defineSqlMigration<PostgresQueryable>(
POSTGRES_COPILOT_FAILURE_DIAGNOSIS_PRE_MODEL_TERMINALIZATION_MIGRATION_ID,
[
POSTGRES_COPILOT_FAILURE_DIAGNOSIS_PRE_MODEL_TERMINALIZATION_TABLE_SQL,
`REVOKE ALL ON TABLE
"${POSTGRES_MODEL_INVOCATION_SCHEMA}"."${PRE_MODEL_TERMINALIZATION_TABLE}"
FROM PUBLIC`,
`GRANT SELECT, INSERT ON TABLE
"${POSTGRES_MODEL_INVOCATION_SCHEMA}"."${PRE_MODEL_TERMINALIZATION_TABLE}"
TO ql3_runtime`,
],
(context, statement) => context.query(statement).then(() => undefined),
);
export const postgresCopilotMigrations = Object.freeze([
postgresCopilotFailureDiagnosisAdmissionMigration,
postgresCopilotFailureDiagnosisToolUnlockMigration,
postgresCopilotFailureDiagnosisModelExecutionMigration,
postgresCopilotFailureDiagnosisPreModelTerminalizationMigration,
]);