mirror of
https://github.com/KeygraphHQ/shannon.git
synced 2026-10-06 16:26:56 +02:00
fix(reconciliation): run task formation on live repo, not a copy (#480)
This commit is contained in:
1 parent
0ab7c0b41b
commit
84212f376d
3 files changed
+65
-404
No files matched your search
@@ -1,295 +0,0 @@
|
||||
// Copyright (C) 2026 Keygraph, Inc.
|
||||
//
|
||||
// This program is free software: you can redistribute it and/or modify
|
||||
// it under the terms of the GNU Affero General Public License version 3
|
||||
// as published by the Free Software Foundation.
|
||||
|
||||
/** Attempt-local working-tree copy used by the task-formation model boundary. */
|
||||
|
||||
import type { Dirent, Stats } from 'node:fs';
|
||||
import { cp, lstat, mkdir, mkdtemp, readdir, realpath, rm } from 'node:fs/promises';
|
||||
import os from 'node:os';
|
||||
import path from 'node:path';
|
||||
import { ArtifactIntegrityError, ReconciliationIoError } from '../reconciliation/artifact-store.js';
|
||||
|
||||
const JAIL_PREFIX = 'shannon-task-formation-';
|
||||
// Never copied into the model-readable jail: `.git` carries deliverables history, `.shannon` holds
|
||||
// scan internals, and `.pi` holds provider credentials. Any of these reaching the jail would expose
|
||||
// them to the tools the model drives. The post-copy verification re-checks their absence by name.
|
||||
const ALWAYS_EXCLUDED_NAMES = Object.freeze(['.git', '.shannon', '.pi'] as const);
|
||||
|
||||
export interface SourceJailOptions {
|
||||
readonly sourceRoot: string;
|
||||
readonly deliverablesPath: string;
|
||||
readonly reconciliationWorkspacePath: string;
|
||||
readonly signal?: AbortSignal;
|
||||
/** Test-only filesystem selector. Production uses `os.tmpdir()`. */
|
||||
readonly tempRoot?: string;
|
||||
}
|
||||
|
||||
/** One source-only jail plus the immutable deny rules used by its live tool gate. */
|
||||
export interface SourceJail {
|
||||
readonly dir: string;
|
||||
readonly deniedPaths: readonly string[];
|
||||
cleanup(): Promise<void>;
|
||||
}
|
||||
|
||||
function isErrno(error: unknown, code: string): boolean {
|
||||
return error instanceof Error && (error as NodeJS.ErrnoException).code === code;
|
||||
}
|
||||
|
||||
function cancellationError(signal: AbortSignal): Error {
|
||||
if (signal.reason instanceof Error) return signal.reason;
|
||||
return new DOMException('Task formation was cancelled.', 'AbortError');
|
||||
}
|
||||
|
||||
function checkCancellation(signal: AbortSignal | undefined): void {
|
||||
if (signal?.aborted === true) throw cancellationError(signal);
|
||||
}
|
||||
|
||||
// Path-confinement predicate: true only when `candidate` is `root` itself or lies beneath it.
|
||||
// A relative path that escapes upward (`..`) or is absolute means the candidate is outside the root.
|
||||
function isWithin(root: string, candidate: string): boolean {
|
||||
const relativePath = path.relative(root, candidate);
|
||||
return (
|
||||
relativePath === '' ||
|
||||
(!relativePath.startsWith(`..${path.sep}`) && relativePath !== '..' && !path.isAbsolute(relativePath))
|
||||
);
|
||||
}
|
||||
|
||||
async function relativeExclusion(
|
||||
sourceRoot: string,
|
||||
lexicalSourceRoot: string,
|
||||
candidate: string,
|
||||
): Promise<string | undefined> {
|
||||
const resolved = path.resolve(candidate);
|
||||
let relativePath: string | undefined;
|
||||
if (isWithin(sourceRoot, resolved)) {
|
||||
relativePath = path.relative(sourceRoot, resolved);
|
||||
} else if (isWithin(lexicalSourceRoot, resolved)) {
|
||||
relativePath = path.relative(lexicalSourceRoot, resolved);
|
||||
} else {
|
||||
try {
|
||||
const canonicalCandidate = await realpath(resolved);
|
||||
if (isWithin(sourceRoot, canonicalCandidate)) {
|
||||
relativePath = path.relative(sourceRoot, canonicalCandidate);
|
||||
}
|
||||
} catch {
|
||||
return undefined;
|
||||
}
|
||||
}
|
||||
if (relativePath === undefined) return undefined;
|
||||
|
||||
if (relativePath === '') {
|
||||
// An exclusion that resolves to the whole root would empty the jail. Fail closed rather than
|
||||
// copy nothing and hand the model an empty tree.
|
||||
throw new ArtifactIntegrityError('A task-formation exclusion resolves to the complete source root');
|
||||
}
|
||||
return relativePath;
|
||||
}
|
||||
|
||||
async function buildDynamicExclusions(
|
||||
options: SourceJailOptions,
|
||||
sourceRoot: string,
|
||||
lexicalSourceRoot: string,
|
||||
): Promise<readonly string[]> {
|
||||
const exclusions = (
|
||||
await Promise.all([
|
||||
relativeExclusion(sourceRoot, lexicalSourceRoot, options.deliverablesPath),
|
||||
relativeExclusion(sourceRoot, lexicalSourceRoot, options.reconciliationWorkspacePath),
|
||||
])
|
||||
).filter((value): value is string => value !== undefined);
|
||||
return Object.freeze([...new Set(exclusions)]);
|
||||
}
|
||||
|
||||
function pathHasAlwaysExcludedName(relativePath: string): boolean {
|
||||
const segments = relativePath.split(path.sep);
|
||||
return segments.some((segment) => (ALWAYS_EXCLUDED_NAMES as readonly string[]).includes(segment));
|
||||
}
|
||||
|
||||
function pathIsDynamicallyExcluded(relativePath: string, exclusions: readonly string[]): boolean {
|
||||
return exclusions.some((excluded) => relativePath === excluded || relativePath.startsWith(`${excluded}${path.sep}`));
|
||||
}
|
||||
|
||||
async function copySourceTree(
|
||||
sourceRoot: string,
|
||||
destination: string,
|
||||
dynamicExclusions: readonly string[],
|
||||
signal: AbortSignal | undefined,
|
||||
): Promise<void> {
|
||||
let entries: Dirent[];
|
||||
try {
|
||||
entries = (await readdir(sourceRoot, { withFileTypes: true })).sort((left, right) =>
|
||||
left.name.localeCompare(right.name),
|
||||
);
|
||||
} catch {
|
||||
throw new ReconciliationIoError('Unable to enumerate the task-formation source tree');
|
||||
}
|
||||
|
||||
// Cancellation is checked before every top-level entry and inside the copy filter so an aborted
|
||||
// scan stops promptly instead of copying a whole large tree first.
|
||||
for (const entry of entries) {
|
||||
checkCancellation(signal);
|
||||
const source = path.join(sourceRoot, entry.name);
|
||||
const destinationEntry = path.join(destination, entry.name);
|
||||
try {
|
||||
// verbatimSymlinks copies links as links rather than following them, so a link pointing
|
||||
// outside the tree cannot pull external content in; the filter then drops any path that
|
||||
// resolves outside the root, plus the always- and dynamically-excluded paths.
|
||||
await cp(source, destinationEntry, {
|
||||
recursive: true,
|
||||
verbatimSymlinks: true,
|
||||
errorOnExist: true,
|
||||
force: false,
|
||||
async filter(candidate) {
|
||||
checkCancellation(signal);
|
||||
const relativePath = path.relative(sourceRoot, candidate);
|
||||
if (relativePath === '' || !isWithin(sourceRoot, path.resolve(candidate))) return false;
|
||||
if (pathHasAlwaysExcludedName(relativePath)) return false;
|
||||
return !pathIsDynamicallyExcluded(relativePath, dynamicExclusions);
|
||||
},
|
||||
});
|
||||
} catch (error) {
|
||||
if (signal?.aborted === true) throw cancellationError(signal);
|
||||
if (error instanceof ArtifactIntegrityError) throw error;
|
||||
throw new ReconciliationIoError('Unable to copy the task-formation source tree');
|
||||
}
|
||||
}
|
||||
checkCancellation(signal);
|
||||
}
|
||||
|
||||
async function assertAlwaysExcludedNamesAbsent(directory: string, signal: AbortSignal | undefined): Promise<void> {
|
||||
checkCancellation(signal);
|
||||
let entries: Dirent[];
|
||||
try {
|
||||
entries = await readdir(directory, { withFileTypes: true });
|
||||
} catch {
|
||||
throw new ReconciliationIoError('Unable to verify the task-formation source jail');
|
||||
}
|
||||
|
||||
for (const entry of entries) {
|
||||
checkCancellation(signal);
|
||||
if ((ALWAYS_EXCLUDED_NAMES as readonly string[]).includes(entry.name)) {
|
||||
throw new ArtifactIntegrityError('The task-formation source jail contains an excluded entry');
|
||||
}
|
||||
if (entry.isDirectory() && !entry.isSymbolicLink()) {
|
||||
await assertAlwaysExcludedNamesAbsent(path.join(directory, entry.name), signal);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async function assertDynamicExclusionsAbsent(
|
||||
directory: string,
|
||||
exclusions: readonly string[],
|
||||
signal: AbortSignal | undefined,
|
||||
): Promise<void> {
|
||||
for (const excluded of exclusions) {
|
||||
checkCancellation(signal);
|
||||
try {
|
||||
await lstat(path.join(directory, excluded));
|
||||
} catch (error) {
|
||||
if (isErrno(error, 'ENOENT')) continue;
|
||||
throw new ReconciliationIoError('Unable to verify a task-formation jail exclusion');
|
||||
}
|
||||
throw new ArtifactIntegrityError('The task-formation source jail contains a protected workspace entry');
|
||||
}
|
||||
}
|
||||
|
||||
// Re-verify the copied tree independently of the copy filter: the jail root must be a real
|
||||
// directory (not a symlink), and no excluded name or protected workspace path may survive. This
|
||||
// catches a filter gap or a race during the copy before the model is allowed to read the tree.
|
||||
async function verifyJail(
|
||||
directory: string,
|
||||
dynamicExclusions: readonly string[],
|
||||
signal: AbortSignal | undefined,
|
||||
): Promise<void> {
|
||||
checkCancellation(signal);
|
||||
let stats: Stats;
|
||||
try {
|
||||
stats = await lstat(directory);
|
||||
} catch {
|
||||
throw new ReconciliationIoError('Unable to inspect the task-formation source jail');
|
||||
}
|
||||
if (stats.isSymbolicLink() || !stats.isDirectory()) {
|
||||
throw new ArtifactIntegrityError('The task-formation source jail is not a real directory');
|
||||
}
|
||||
await assertAlwaysExcludedNamesAbsent(directory, signal);
|
||||
await assertDynamicExclusionsAbsent(directory, dynamicExclusions, signal);
|
||||
checkCancellation(signal);
|
||||
}
|
||||
|
||||
async function removeJail(directory: string): Promise<void> {
|
||||
try {
|
||||
await rm(directory, { recursive: true, force: true });
|
||||
} catch {
|
||||
throw new ReconciliationIoError('Unable to remove the task-formation source jail');
|
||||
}
|
||||
|
||||
try {
|
||||
await lstat(directory);
|
||||
} catch (error) {
|
||||
if (isErrno(error, 'ENOENT')) return;
|
||||
throw new ReconciliationIoError('Unable to verify task-formation source-jail cleanup');
|
||||
}
|
||||
throw new ReconciliationIoError('Task-formation source-jail cleanup left the jail on disk');
|
||||
}
|
||||
|
||||
/**
|
||||
* Copy the scanned working tree into an isolated temporary directory without following symlinks.
|
||||
* Every failure removes the attempt-local directory before it propagates.
|
||||
*/
|
||||
export async function materializeSourceJail(options: SourceJailOptions): Promise<SourceJail> {
|
||||
checkCancellation(options.signal);
|
||||
|
||||
const lexicalSourceRoot = path.resolve(options.sourceRoot);
|
||||
let sourceRoot: string;
|
||||
try {
|
||||
sourceRoot = await realpath(options.sourceRoot);
|
||||
const sourceStats = await lstat(sourceRoot);
|
||||
if (sourceStats.isSymbolicLink() || !sourceStats.isDirectory()) {
|
||||
throw new ArtifactIntegrityError('The task-formation source root is not a real directory');
|
||||
}
|
||||
} catch (error) {
|
||||
if (error instanceof ArtifactIntegrityError) throw error;
|
||||
throw new ReconciliationIoError('Unable to resolve the task-formation source root');
|
||||
}
|
||||
|
||||
let tempRoot: string;
|
||||
try {
|
||||
const configuredTempRoot = options.tempRoot ?? os.tmpdir();
|
||||
await mkdir(configuredTempRoot, { recursive: true });
|
||||
tempRoot = await realpath(configuredTempRoot);
|
||||
} catch {
|
||||
throw new ReconciliationIoError('Unable to resolve the task-formation temporary root');
|
||||
}
|
||||
// A temp root inside the source tree would make the copy try to copy the jail into itself.
|
||||
if (isWithin(sourceRoot, tempRoot)) {
|
||||
throw new ArtifactIntegrityError('The task-formation temporary root cannot be inside the source tree');
|
||||
}
|
||||
|
||||
const dynamicExclusions = await buildDynamicExclusions(options, sourceRoot, lexicalSourceRoot);
|
||||
let directory: string;
|
||||
try {
|
||||
directory = await mkdtemp(path.join(tempRoot, JAIL_PREFIX));
|
||||
} catch {
|
||||
throw new ReconciliationIoError('Unable to create the task-formation source jail');
|
||||
}
|
||||
|
||||
let cleaned = false;
|
||||
const cleanup = async (): Promise<void> => {
|
||||
if (cleaned) return;
|
||||
await removeJail(directory);
|
||||
cleaned = true;
|
||||
};
|
||||
|
||||
try {
|
||||
await copySourceTree(sourceRoot, directory, dynamicExclusions, options.signal);
|
||||
await verifyJail(directory, dynamicExclusions, options.signal);
|
||||
} catch (error) {
|
||||
await cleanup().catch(() => undefined);
|
||||
throw error;
|
||||
}
|
||||
|
||||
const deniedPaths = Object.freeze([...ALWAYS_EXCLUDED_NAMES, ...dynamicExclusions]);
|
||||
return Object.freeze({ dir: directory, deniedPaths, cleanup });
|
||||
}
|
||||
@@ -38,10 +38,9 @@ const MAX_TURNS = 128;
|
||||
const MAX_LIST_RESULTS = 500;
|
||||
const DEFAULT_LIST_RESULTS = 200;
|
||||
const MAX_OUTPUT_BYTES = 64 * 1024;
|
||||
// The live-tool-side counterpart of the source jail's copy-time exclusion (source-jail.ts): even if
|
||||
// one of these somehow existed in the jailed tree, the read/grep/find/ls/glob tools built below must
|
||||
// still refuse to serve it. `.git` is deliverables history, `.shannon` is scan internals, `.pi` is
|
||||
// provider credentials.
|
||||
// Task formation reads the live repository, so these read-only tools are the sole barrier keeping the
|
||||
// model out of `.git` (source history), `.shannon` (scan internals, incl. the deliverables Git repo),
|
||||
// and `.pi` (credentials). Always denied, whatever extra denies a caller passes.
|
||||
const ALWAYS_DENIED_PATHS = Object.freeze(['.git', '.shannon', '.pi'] as const);
|
||||
const TRANSIENT_IO_CODES = new Set([
|
||||
'EAGAIN',
|
||||
@@ -288,9 +287,9 @@ function createGlobTool(confinement: RepositoryConfinement): ToolDefinition {
|
||||
return defineTool({
|
||||
name: 'glob',
|
||||
label: 'Glob source files',
|
||||
description: 'Match bounded file globs from the source-jail root without following symlinks.',
|
||||
promptSnippet: 'glob: match source files from the jail root',
|
||||
promptGuidelines: ['Patterns are always rooted in the source jail.'],
|
||||
description: 'Match bounded file globs from the repository root without following symlinks.',
|
||||
promptSnippet: 'glob: match source files from the repository root',
|
||||
promptGuidelines: ['Patterns are always rooted in the repository.'],
|
||||
parameters: Type.Object(
|
||||
{
|
||||
pattern: Type.String({ minLength: 1, maxLength: 256 }),
|
||||
@@ -322,7 +321,7 @@ function createGlobTool(confinement: RepositoryConfinement): ToolDefinition {
|
||||
});
|
||||
}
|
||||
|
||||
/** Create the five code-owned source tools that share one canonical jail policy. */
|
||||
/** Create the five code-owned source tools that share one canonical deny policy. */
|
||||
export async function createTaskFormationSourceTools(options: ToolFactoryOptions): Promise<readonly ToolDefinition[]> {
|
||||
const deniedPaths = uniqueDeniedPaths(options.deniedPaths);
|
||||
const capellaTools = await createCapellaRepositoryTools({
|
||||
|
||||
@@ -6,12 +6,10 @@
|
||||
|
||||
/** Pass 1 task formation over current observations only. */
|
||||
|
||||
import path from 'node:path';
|
||||
import { DEFAULT_DELIVERABLES_SUBDIR, WORKSPACES_DIR } from '../../paths.js';
|
||||
import { WORKSPACES_DIR } from '../../paths.js';
|
||||
import { loadPrompt } from '../../services/prompt-manager.js';
|
||||
import type { ActivityLogger } from '../../types/activity-logger.js';
|
||||
import type { ReconciliationClass } from '../../types/reconciliation.js';
|
||||
import { materializeSourceJail } from '../pi/source-jail.js';
|
||||
import {
|
||||
isTaskFormationFallbackReason,
|
||||
type TaskFormationExecutionContext,
|
||||
@@ -61,7 +59,6 @@ export interface FormClassExploitTasksInput {
|
||||
readonly repositoryPath: string;
|
||||
readonly producerRef: ArtifactRef<'producer-observations'>;
|
||||
readonly supplementalRef: ArtifactRef<'supplemental-observations'>;
|
||||
readonly deliverablesSubdir?: string;
|
||||
readonly webUrl?: string;
|
||||
}
|
||||
|
||||
@@ -181,14 +178,14 @@ function taskFormationPromptName(vulnerabilityClass: ReconciliationClass): strin
|
||||
// prompt itself is missing or unreadable content, which Temporal should not spend retries on.
|
||||
async function loadClassPolicy(
|
||||
vulnerabilityClass: ReconciliationClass,
|
||||
jailPath: string,
|
||||
repoPath: string,
|
||||
webUrl: string,
|
||||
logger: ActivityLogger,
|
||||
): Promise<string> {
|
||||
try {
|
||||
return await loadPrompt(
|
||||
taskFormationPromptName(vulnerabilityClass),
|
||||
{ webUrl, repoPath: jailPath, AUTH_STATE_FILE: '' },
|
||||
{ webUrl, repoPath, AUTH_STATE_FILE: '' },
|
||||
null,
|
||||
false,
|
||||
logger,
|
||||
@@ -332,109 +329,69 @@ export function createFormClassExploitTasks(
|
||||
const submitTool = createValidatingSubmitTool(buildTaskFormationSchema([...labelSet]), (parameters) =>
|
||||
findTaskFormationProblems(parameters, labelSet),
|
||||
);
|
||||
const deliverablesPath = path.resolve(
|
||||
const classPolicy = await loadClassPolicy(
|
||||
input.vulnerabilityClass,
|
||||
input.repositoryPath,
|
||||
input.deliverablesSubdir ?? DEFAULT_DELIVERABLES_SUBDIR,
|
||||
input.webUrl ?? 'https://not-applicable.invalid',
|
||||
logger,
|
||||
);
|
||||
const reconciliationWorkspacePath = path.resolve(workspacesDir, input.sessionId, '.shannon', 'reconciliation');
|
||||
// Task formation runs against a disposable copy of the source tree rather than the live
|
||||
// repository or the deliverables directory, so the model's tool calls during this stage cannot
|
||||
// read or modify anything outside what it was actually given to reason about.
|
||||
const jail = await materializeSourceJail({
|
||||
sourceRoot: input.repositoryPath,
|
||||
deliverablesPath,
|
||||
reconciliationWorkspacePath,
|
||||
...(signal !== undefined && { signal }),
|
||||
});
|
||||
checkCancellation(signal);
|
||||
|
||||
let formation: FormClassExploitTasksResult;
|
||||
let modelResult: TaskFormationExecutorResult;
|
||||
try {
|
||||
const classPolicy = await loadClassPolicy(
|
||||
input.vulnerabilityClass,
|
||||
jail.dir,
|
||||
input.webUrl ?? 'https://not-applicable.invalid',
|
||||
logger,
|
||||
);
|
||||
checkCancellation(signal);
|
||||
|
||||
let modelResult: TaskFormationExecutorResult;
|
||||
try {
|
||||
const executorTimeoutMs = deps.executorTimeoutMsFor?.();
|
||||
modelResult = await executor.run({
|
||||
cwd: jail.dir,
|
||||
systemPrompt: classPolicy,
|
||||
modelContext: modelInput.serialized,
|
||||
deniedPaths: jail.deniedPaths,
|
||||
submitTool,
|
||||
signal: signal ?? new AbortController().signal,
|
||||
...(executorTimeoutMs !== undefined && { timeoutMs: executorTimeoutMs }),
|
||||
correlation: {
|
||||
...deps.executionContextFor?.(),
|
||||
stage: 'task-formation',
|
||||
vulnerabilityClass: input.vulnerabilityClass,
|
||||
},
|
||||
});
|
||||
} catch (error) {
|
||||
if (!(error instanceof TaskFormationExecutorError)) throw error;
|
||||
if (error.failureKind === 'infrastructure') {
|
||||
throw new ReconciliationIoError(
|
||||
'Task-formation executor setup encountered a retryable infrastructure failure',
|
||||
);
|
||||
}
|
||||
if (error.failureKind !== 'model') throw error;
|
||||
const metrics = metricsFromUsage(error.usage, error.modelCalls);
|
||||
deps.onMetrics?.(metrics);
|
||||
throw new TaskFormationModelError({
|
||||
message: error.message,
|
||||
retryable: error.retryable,
|
||||
...(error.fallbackReason !== undefined && { fallbackReason: error.fallbackReason }),
|
||||
metrics,
|
||||
});
|
||||
}
|
||||
|
||||
const metrics = metricsFromUsage(modelResult.usage, modelResult.modelCalls);
|
||||
deps.onMetrics?.(metrics);
|
||||
checkCancellation(signal);
|
||||
const accepted = acceptTaskGroups(modelResult.output, labelSet);
|
||||
const groups = accepted.groups.map((group) => ({
|
||||
producer_ids: group.queue_labels.map((label) => {
|
||||
const producerId = modelInput.labelToProducerId.get(label);
|
||||
if (producerId === undefined) {
|
||||
throw new ArtifactIntegrityError('An accepted task-formation label has no observation mapping');
|
||||
}
|
||||
return producerId;
|
||||
}),
|
||||
reasoning: group.reasoning,
|
||||
}));
|
||||
const body: TaskFormationBody = {
|
||||
model_ran: true,
|
||||
groups,
|
||||
rejected_group_count: accepted.rejectedGroupCount,
|
||||
dropped_unknown_label_count: accepted.droppedUnknownLabelCount,
|
||||
};
|
||||
const ref = await writeFormationArtifact(input, workspacesDir, body);
|
||||
formation = { ref, metrics, model: `${modelResult.providerId}:${modelResult.modelId}` };
|
||||
const executorTimeoutMs = deps.executorTimeoutMsFor?.();
|
||||
modelResult = await executor.run({
|
||||
cwd: input.repositoryPath,
|
||||
systemPrompt: classPolicy,
|
||||
modelContext: modelInput.serialized,
|
||||
deniedPaths: [],
|
||||
submitTool,
|
||||
signal: signal ?? new AbortController().signal,
|
||||
...(executorTimeoutMs !== undefined && { timeoutMs: executorTimeoutMs }),
|
||||
correlation: {
|
||||
...deps.executionContextFor?.(),
|
||||
stage: 'task-formation',
|
||||
vulnerabilityClass: input.vulnerabilityClass,
|
||||
},
|
||||
});
|
||||
} catch (error) {
|
||||
// A primary error — including cancellation — already owns the outcome, so a cleanup failure
|
||||
// is logged and swallowed rather than replacing that error's type or cause chain.
|
||||
try {
|
||||
await jail.cleanup();
|
||||
} catch {
|
||||
logger.error(
|
||||
'A temporary copy of your source code could not be removed after analysis. It is inside the scan workspace and is safe to delete.',
|
||||
{
|
||||
stage: 'task-formation',
|
||||
vulnerabilityClass: input.vulnerabilityClass,
|
||||
},
|
||||
);
|
||||
if (!(error instanceof TaskFormationExecutorError)) throw error;
|
||||
if (error.failureKind === 'infrastructure') {
|
||||
throw new ReconciliationIoError('Task-formation executor setup encountered a retryable infrastructure failure');
|
||||
}
|
||||
throw error;
|
||||
if (error.failureKind !== 'model') throw error;
|
||||
const metrics = metricsFromUsage(error.usage, error.modelCalls);
|
||||
deps.onMetrics?.(metrics);
|
||||
throw new TaskFormationModelError({
|
||||
message: error.message,
|
||||
retryable: error.retryable,
|
||||
...(error.fallbackReason !== undefined && { fallbackReason: error.fallbackReason }),
|
||||
metrics,
|
||||
});
|
||||
}
|
||||
|
||||
// Nothing else is in flight after a successful formation, so an unremoved or unverifiable jail
|
||||
// is the stage's outcome: it leaves a full copy of the scanned tree on disk and fails here.
|
||||
await jail.cleanup();
|
||||
return formation;
|
||||
const metrics = metricsFromUsage(modelResult.usage, modelResult.modelCalls);
|
||||
deps.onMetrics?.(metrics);
|
||||
checkCancellation(signal);
|
||||
const accepted = acceptTaskGroups(modelResult.output, labelSet);
|
||||
const groups = accepted.groups.map((group) => ({
|
||||
producer_ids: group.queue_labels.map((label) => {
|
||||
const producerId = modelInput.labelToProducerId.get(label);
|
||||
if (producerId === undefined) {
|
||||
throw new ArtifactIntegrityError('An accepted task-formation label has no observation mapping');
|
||||
}
|
||||
return producerId;
|
||||
}),
|
||||
reasoning: group.reasoning,
|
||||
}));
|
||||
const body: TaskFormationBody = {
|
||||
model_ran: true,
|
||||
groups,
|
||||
rejected_group_count: accepted.rejectedGroupCount,
|
||||
dropped_unknown_label_count: accepted.droppedUnknownLabelCount,
|
||||
};
|
||||
const ref = await writeFormationArtifact(input, workspacesDir, body);
|
||||
return { ref, metrics, model: `${modelResult.providerId}:${modelResult.modelId}` };
|
||||
};
|
||||
}
|
||||
|
||||
|
||||
Reference in new issue
Block a user