mirror of
https://github.com/KeygraphHQ/shannon.git
synced 2026-10-04 23:36:54 +02:00
feat: Shannon 3.0 Agentic SAST (#433)
* feat(worker): add agentic static analysis Add the ten-stage Agentic SAST pipeline, confined repository tools, model runtime, prompt templates, and SARIF export. Make retries, repair sessions, reduced coverage, usage accounting, and model-output drift durable across Temporal replay and resume. Keep retry diagnostics in their actionable closed vocabulary. Package the Mantis-derived license material with the prompts that require it. * feat(worker): deduplicate static and runtime findings before exploitation Parse Agentic SAST SARIF into typed observations, enrich and route those observations, and reconcile them with pentest findings before exploitation. Publish deterministic exploitation queues with stable lineage, exact-path Git commits, retry-safe manifests, named drop reasons, and confined task formation. Reject duplicate producer IDs before commit and adopt either legal provenance shape after a lost acknowledgement. * feat(config)!: replace vuln_classes with agentic_sast Wire Agentic SAST and reconciliation into the main pipeline, persist their durable state, and add the Miscellaneous finding and exploitation lane. Make scan completion, cancellation, partial outcomes, resume identity, and report recovery use the integrated final workflow contract. Introduce the atomic finalization, ordering, renumbering, compaction, and output services that workflow calls. Keep completed Miscellaneous work and report drafts idempotent across resume, preserve public main's default-on exploit SARIF behavior, and describe stage-fallback candidates without claiming they were exported. BREAKING CHANGE: `vuln_classes` has been removed. Configs containing it now fail validation, and all five core pentest classes run on every scan. Workspaces created by Shannon 2.x cannot be resumed. Finish or discard in-flight scans before upgrading, then start a new workspace name. * perf: overlap static analysis and the Miscellaneous lane with the pentest Run Agentic SAST alongside vulnerability analysis and run Miscellaneous exploitation alongside the specialist exploitation lanes. Keep reconciliation dependent on the completed static-analysis result while preserving parallel work everywhere that has no data dependency. * feat(cli)!: default the scan target and add a JSON error contract List local scans, resolve the active or most recent workspace automatically, and make logs, status, and stop use one canonical scan identity. Add stable machine-readable failures, richer status output, explicit help errors, and seven-day Temporal retention. Treat absent Temporal pending-activity failures as absent whether the decoder represents them as `null` or missing. BREAKING CHANGE: `status --json` now returns a fixed `failureMessage`. Read `partialReasons`, `agenticSast`, and `workflow.log` for diagnostic detail. * feat(logging): trace tool calls and write a log per agent Record complete tool-call arguments in the workflow log and project each agent's events into its own durable log. Add agent listing and agent-specific log tailing while preserving byte-exact output and draining log handles before activities return. * feat(worker): standardize severity and reporting guidance in exploit prompts Give every exploit agent the same status, confidence, severity-reasoning, report-writing, credential-handling, and scope contract. Apply the same task-formation and SAST-enrichment procedure to the Miscellaneous lane. * feat(worker): disclose scan coverage and make reporting auditable Build on the retry-safe finalization foundation to preserve correct identities, source locations, scan dates, partial-coverage limitations, and consistent report JSON, Markdown, SARIF, and PDF output. Report Agentic SAST, reconciliation wall-clock time, stage usage, retry spend, and background work without duplicate or hardcoded totals. Keep report findings canonical, drop cross-class restatements, name enrichment losses, and render the executive-summary narrative in the PDF. * chore(license): attribute Mantis and Pi and refresh the docs Add the final Mantis and Pi notices, license copies, acknowledgements, and residual copyright updates. Update the README, maintained documentation, contributor guidance, and hand-maintained mirrors to describe Agentic SAST, reconciliation, the Miscellaneous lane, current CLI behavior, and the final release contract. Correct stale workspace and container guidance and annotate long-standing internals for maintainers. * fix(logging): treat a slash as a word separator in agent labels * feat(cli)!: rebuild scan status around model work - show Capella stages beneath the concurrent Agentic SAST phase - attach reconciliation time to the class row it feeds - hide completed bookkeeping and the duplicate miscellaneous wrapper - carry validated child-workflow progress into durable parent state - derive the terminal tree and status JSON from the same phase shape BREAKING CHANGE: `status --json` replaces phase `parallel` with `children` and `meta`, adds phase summaries and notes plus agent attachment fields, and removes the `analysis-engines` and `operational-work` phases. * fix(report): drop the empty Critical Findings section from the PDF summary * fix(sast): align Capella export with the submit-time code-path contract The export gate required every code_paths entry to be file:line, but submit only requires the primary sink to be file:line and accepts bare trace steps. A single malformed trace step therefore dropped an otherwise-valid finding at export. - add isValidPrimaryCodePath as the one shared primary-sink contract - validate only the primary at export; buildResult already drops unusable steps - route the submit-time validator through the same helper so the two cannot drift * feat(sast): tolerate hygiene-only Capella reductions instead of going partial A reduction only makes a run partial when it loses real coverage or a whole finding. Malformed model output, salvaged turn-limit work, and rejected duplicate verdicts are recorded as evidence but no longer flip the run to partial. - add reductionIsTolerable: partial only when genuine-loss counts are nonzero - drive runCapella's partial reasons and display coverage off non-tolerable ones - keep every reduction in agenticSast.reductions so nothing is lost as evidence * feat(logging): record the provider reason for a failed agent turn A failed provider turn collapsed to AGENT_EXECUTION_FAILED/unknown with the underlying reason discarded, so a model-side rejection or safeguard was indistinguishable from a transport fault in the error log. - add safeProviderTurnDetails: write bounded, non-sensitive fields (provider, model, responseId, stop reason, tool-in-flight, category, retryable) to error.log - gate a sanitized errorMessage snippet behind SHANNON_DEBUG_PROVIDER_ERRORS, off by default - forward SHANNON_DEBUG_PROVIDER_ERRORS from the CLI into the worker container * fix(cli): keep shannon logs tailing through a Temporal blip - End the interactive tail on the log's own terminal marker or Ctrl-C, so a transient Temporal outage no longer aborts the command with exit 1. - Rebuild the memoized Temporal client after a failed poll: a wedged gRPC channel was cached forever, so "retrying…" could never reconnect. - Keep start --follow (CI) bounded — a genuinely dead Temporal still fails the run instead of hanging. * fix(worker): correct PDF finding reporting - Render OWASP category, authentication state, and remediation - Omit the redundant per-finding exploited status - Preserve canonical category and field ordering across report modes - Continue Proof of Impact numbering across embedded code blocks - Wrap long PDF code lines without changing canonical report content * fix: attribute a reconciliation failure to exploitation only - Stop marking a class's vulnerability-analysis agent failed when that agent succeeded and only reconciliation failed; the status tree now renders the analysis row completed and the exploitation row failed - Consume the worker's failedReconciliations signal in the CLI, which the mirrored PipelineState already declared but never read - Correct the class_reconciliation_failed message, which claimed the class's analysis results were still in the report when the class is excluded from it * fix(pi): give each task sub-session its own resource loader to prevent stale extension ctx * fix(prompts): scope exploit agents to in-band proof, mark OOB-only findings blocked * fix(cli): reject a shell credential that shadows a gateway config.toml key * fix(cli): make scan shutdown verifiable - preselect and persist workflow identity before worker launch - cancel first, then verify bounded Temporal termination - reconcile Docker workers with Temporal open workflows - fail closed on stale images and unavailable lifecycle state - mark cancellation only after confirmed shutdown * feat(cli): prompt for setup on a bare npx invocation with no credentials * fix(cli): don't blame anthropic when no credentials are configured at all * chore(release): bump beta base version to 3.0.0 * feat(cli): show a 'start your first scan' box in help on a TTY * docs: refresh README and platform overview for Shannon 3.0 - lead with the 3.0 launch note and rewrite key capabilities around security code analysis, the rebuilt terminal experience, native CI/CD, and PDF/SARIF - recast the editions table as Shannon Open Source against the Keygraph Enterprise Platform, stating open source is not a trial edition - rewrite the platform overview around exhaustive agentic SAST, canonical findings, automated remediation, targeted verification, and governance - add five product screenshots under assets/keygraph-platform/, referenced relative to docs/ * docs: add the Shannon naming section and swap in the 3.0 demo GIF - explain the Claude Shannon information-theory origin under "What is Shannon?" - point "Shannon in Action" at the 3.0 recording in assets/Shannon3GIF.gif Both taken from the README half of #438. * docs: document CI/CD integrations and the reconciled analysis pipeline - add a CI/CD Integrations section covering the official GitHub Action and GitLab component, pipeline artifacts, and exploit-only severity gates - redraw the architecture section as a Mermaid flow: agentic code analysis and recon feed finding reconciliation, then exploitation and reporting - describe open-source code analysis as a multi-stage agentic workflow and reserve parsed-code CPGs and exhaustive verification for Enterprise - sharpen the privacy wording: results stay local, but model requests carry source context to whichever endpoint you configure - drop the "not recommended" framing on local models and add a section on why Shannon complements rather than replaces human pentesters - regenerate llms-full.txt from the updated README and docs * docs: add the Photoview benchmark across three models - Add a "Shannon in Action" table for Photoview 2.4.0 runs on DeepSeek v4 Flash, Grok 4.6, and Claude Opus 5, each linking its PDF report and SARIF output - Store the per-model reports under benchmark/ - Link the (forthcoming) benchmark writeup from the section intro * docs: add the Shannon vs XBOW/Aikido Photoview benchmark writeup - Add docs/shannon-xbow-aikido-benchmark.md with methodology, per-model cost/coverage tables, and links to each model's report and SARIF - Link the writeup from the README "Shannon in Action" section * docs: link the benchmark announcement discussion from the README * fix(readme): restore theme-aware banner, badge, and buttons * feat!: trigger the Shannon 3.0 major release --------- Co-authored-by: ezl-keygraph <ezhil@keygraph.io>
This commit is contained in:
1 parent
6108de3cfc
commit
9767ebe633
267 files changed
+127389
-3268
No files matched your search
File diff suppressed because it is too large.
Load diff
@@ -1,4 +1,4 @@
|
||||
// Copyright (C) 2025 Keygraph, Inc.
|
||||
// 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
|
||||
|
||||
@@ -7,7 +7,12 @@
|
||||
|
||||
export type { ActivityInput } from './activities.js';
|
||||
export type {
|
||||
AgenticSastInput,
|
||||
AgenticSastState,
|
||||
AgentMetrics,
|
||||
NonFatalFailure,
|
||||
OperationalMetrics,
|
||||
OperationalStageState,
|
||||
PipelineInput,
|
||||
PipelineState,
|
||||
PipelineSummary,
|
||||
|
||||
@@ -0,0 +1,637 @@
|
||||
// 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.
|
||||
|
||||
/** Temporal boundary for the standalone reconciliation stages. */
|
||||
|
||||
import { ApplicationFailure, CancelledFailure, Context, heartbeat } from '@temporalio/activity';
|
||||
import { type ModelHost, modelHost } from '../ai/model-host.js';
|
||||
import { createPiStructuredGenerationPort } from '../ai/pi/structured-generation.js';
|
||||
import {
|
||||
createTaskFormationExecutor,
|
||||
TASK_FORMATION_FALLBACK_REASONS,
|
||||
type TaskFormationExecutionContext,
|
||||
TaskFormationExecutorError,
|
||||
type TaskFormationFallbackReason,
|
||||
} from '../ai/pi/task-formation-executor.js';
|
||||
import { ReconciliationError } from '../ai/reconciliation/artifact-store.js';
|
||||
import {
|
||||
createEnrichClassSastObservations,
|
||||
type EnrichClassSastObservationsInput,
|
||||
SastEnrichmentModelError,
|
||||
} from '../ai/reconciliation/enrich.js';
|
||||
import {
|
||||
createFormClassExploitTasks,
|
||||
type FormClassExploitTasksInput,
|
||||
TaskFormationModelError,
|
||||
} from '../ai/reconciliation/form.js';
|
||||
import {
|
||||
type MaterializeClassExploitTasksArgs,
|
||||
materializeClassExploitTasks as materializeClassExploitTasksStage,
|
||||
} from '../ai/reconciliation/materialize.js';
|
||||
import {
|
||||
type PrepareClassReconciliationArgs,
|
||||
prepareClassReconciliation as prepareClassReconciliationStage,
|
||||
} from '../ai/reconciliation/prepare.js';
|
||||
import {
|
||||
type PublishClassReconciliationOssArgs,
|
||||
publicationContractForClass,
|
||||
publishClassReconciliationOss as publishClassReconciliationOssStage,
|
||||
} from '../ai/reconciliation/publish.js';
|
||||
import {
|
||||
type SeedEmptyProducerQueueArgs,
|
||||
seedEmptyProducerQueue as seedEmptyProducerQueueStage,
|
||||
} from '../ai/reconciliation/seed-miscellaneous.js';
|
||||
import type {
|
||||
EnrichSuccess,
|
||||
FormSuccess,
|
||||
MaterializeResult,
|
||||
PrepareResult,
|
||||
} from '../ai/reconciliation/stage-contracts.js';
|
||||
import type { ActivityLogger } from '../types/activity-logger.js';
|
||||
import type { ReconciliationClass } from '../types/reconciliation.js';
|
||||
import { renderSafeMessage } from '../types/run-state.js';
|
||||
import { createActivityLogger } from './activity-logger.js';
|
||||
import {
|
||||
ACCEPTED_TASK_FORMATION_FALLBACK_REASONS,
|
||||
type EnrichClassSastObservationsActivityInput,
|
||||
type FormClassExploitTasksActivityInput,
|
||||
type FormClassExploitTasksActivityResult,
|
||||
type MaterializeClassExploitTasksActivityInput,
|
||||
type PrepareClassReconciliationActivityInput,
|
||||
type PublishClassReconciliationActivityInput,
|
||||
RECONCILIATION_ACTIVITY_NAMES,
|
||||
RECONCILIATION_ACTIVITY_PROFILES,
|
||||
RECONCILIATION_STABLE_FAILURE_TYPES,
|
||||
type ReconciliationActivityName,
|
||||
type ReconciliationActivityRegistry,
|
||||
type ReconciliationClassActivityName,
|
||||
type ReconciliationStableFailureType,
|
||||
resolveReconciliationActivityBudget,
|
||||
type SeedEmptyProducerQueueActivityInput,
|
||||
TASK_FORMATION_EXECUTOR_TIMEOUT_MARGIN_MS,
|
||||
} from './reconcile-activity-types.js';
|
||||
|
||||
const STABLE_FAILURE_TYPES: ReadonlySet<string> = new Set(RECONCILIATION_STABLE_FAILURE_TYPES);
|
||||
|
||||
// The workflow validates fallback reasons against its bundle-safe mirror; fail fast at worker
|
||||
// startup if the mirror ever drifts from the executor's authoritative closed set.
|
||||
{
|
||||
const mirror = [...ACCEPTED_TASK_FORMATION_FALLBACK_REASONS].sort();
|
||||
const authoritative = [...TASK_FORMATION_FALLBACK_REASONS].sort();
|
||||
if (mirror.length !== authoritative.length || mirror.some((reason, index) => reason !== authoritative[index])) {
|
||||
throw new Error('The workflow fallback-reason mirror does not match the task-formation executor contract');
|
||||
}
|
||||
}
|
||||
|
||||
const DEFAULT_RETRYABILITY: Readonly<Record<ReconciliationStableFailureType, boolean>> = Object.freeze({
|
||||
TaskFormationModelError: true,
|
||||
SastEnrichmentModelError: true,
|
||||
ReconciliationArtifactNotFound: true,
|
||||
ReconciliationIoError: true,
|
||||
ConfigurationError: false,
|
||||
SastEnrichmentInputError: false,
|
||||
ArtifactIntegrityError: false,
|
||||
PublicationConflict: false,
|
||||
UnmappableSurvivor: false,
|
||||
KeySetDivergence: false,
|
||||
});
|
||||
|
||||
/**
|
||||
* One sentence per stable failure type, written for the reader rather than for the
|
||||
* reconciliation design. `{Class}` and `{class}` are substituted from the failing class,
|
||||
* which every reconciliation activity carries in its input.
|
||||
*/
|
||||
const SAFE_FAILURE_MESSAGES: Readonly<Record<ReconciliationStableFailureType, string>> = Object.freeze({
|
||||
TaskFormationModelError: 'Shannon could not group {class} findings into test cases.',
|
||||
SastEnrichmentModelError: 'Shannon could not add code context to the {class} findings from static analysis.',
|
||||
ReconciliationArtifactNotFound:
|
||||
'A saved {class} result could not be read back. Re-running this workspace retries it.',
|
||||
ReconciliationIoError: 'A reconciliation filesystem or Git operation failed.',
|
||||
ConfigurationError: 'Reconciliation activity configuration is invalid.',
|
||||
SastEnrichmentInputError: 'The supplied SAST reference is invalid.',
|
||||
ArtifactIntegrityError: 'A saved {class} result failed its integrity check and was not used.',
|
||||
PublicationConflict:
|
||||
"{Class} results were already published by an earlier run, and this run's results differ. Nothing was overwritten.",
|
||||
UnmappableSurvivor:
|
||||
'Shannon could not match a finding in the report back to the test case it came from. {Class} results were not published.',
|
||||
KeySetDivergence:
|
||||
'Shannon found two disagreeing sets of findings for {class} and stopped rather than publish either.',
|
||||
});
|
||||
|
||||
interface ReconciliationHeartbeatDetails {
|
||||
readonly stage: ReconciliationActivityName;
|
||||
readonly attempt: number;
|
||||
readonly elapsedSeconds: number;
|
||||
readonly classDeadlineMs: number;
|
||||
}
|
||||
|
||||
export interface ReconciliationActivityRuntime {
|
||||
readonly attempt: number;
|
||||
readonly cancellationSignal: AbortSignal;
|
||||
readonly logger: ActivityLogger;
|
||||
/** Temporal's granted per-attempt execution budget, from the activity info. */
|
||||
readonly startToCloseTimeoutMs?: number;
|
||||
/** Bounded per-attempt correlation identifier (run id + activity id). */
|
||||
readonly executionKey?: string;
|
||||
heartbeat(details: ReconciliationHeartbeatDetails): void;
|
||||
}
|
||||
|
||||
interface ReconciliationStageRuntime {
|
||||
readonly signal: AbortSignal;
|
||||
readonly logger: ActivityLogger;
|
||||
readonly modelHost: ModelHost;
|
||||
/** Remaining granted budget minus the deterministic margin, evaluated at call time. */
|
||||
readonly executorTimeoutMsFor?: () => number | undefined;
|
||||
readonly executionContextFor?: () => TaskFormationExecutionContext | undefined;
|
||||
}
|
||||
|
||||
export interface ReconciliationStageBindings {
|
||||
seedEmptyProducerQueue(args: SeedEmptyProducerQueueArgs): Promise<{
|
||||
alreadySeeded: boolean;
|
||||
alreadyPublished: boolean;
|
||||
commitHash: string;
|
||||
}>;
|
||||
prepareClassReconciliation(args: PrepareClassReconciliationArgs): Promise<PrepareResult>;
|
||||
enrichClassSastObservations(
|
||||
input: EnrichClassSastObservationsInput,
|
||||
runtime: ReconciliationStageRuntime,
|
||||
): Promise<EnrichSuccess>;
|
||||
formClassExploitTasks(input: FormClassExploitTasksInput, runtime: ReconciliationStageRuntime): Promise<FormSuccess>;
|
||||
materializeClassExploitTasks(args: MaterializeClassExploitTasksArgs): Promise<MaterializeResult>;
|
||||
publishClassReconciliationOss(args: PublishClassReconciliationOssArgs): Promise<{
|
||||
alreadyPublished: boolean;
|
||||
manifestSha256: string;
|
||||
commitHash: string;
|
||||
}>;
|
||||
}
|
||||
|
||||
export interface ReconciliationActivityBindings {
|
||||
/** Worker-local paths are bound here and never enter Temporal activity arguments. */
|
||||
readonly repositoryPath: string;
|
||||
readonly deliverablesDir: string;
|
||||
readonly workspacesDir: string;
|
||||
readonly webUrl?: string;
|
||||
readonly modelHost?: ModelHost;
|
||||
readonly now?: () => number;
|
||||
readonly runtime?: () => ReconciliationActivityRuntime;
|
||||
readonly stages?: Partial<ReconciliationStageBindings>;
|
||||
}
|
||||
|
||||
interface FailureMetrics {
|
||||
readonly costUsd: number;
|
||||
readonly modelCalls: number;
|
||||
readonly inputTokens: number;
|
||||
readonly outputTokens: number;
|
||||
}
|
||||
|
||||
interface StableFailureDetails {
|
||||
readonly metrics?: FailureMetrics;
|
||||
readonly fallbackReason?: TaskFormationFallbackReason;
|
||||
}
|
||||
|
||||
function defaultRuntime(): ReconciliationActivityRuntime {
|
||||
const context = Context.current();
|
||||
return {
|
||||
attempt: context.info.attempt,
|
||||
cancellationSignal: context.cancellationSignal,
|
||||
logger: createActivityLogger(),
|
||||
startToCloseTimeoutMs: context.info.startToCloseTimeoutMs,
|
||||
executionKey: `${context.info.workflowExecution.runId}:${context.info.activityId}`,
|
||||
heartbeat,
|
||||
};
|
||||
}
|
||||
|
||||
function isStableFailureType(value: string): value is ReconciliationStableFailureType {
|
||||
return STABLE_FAILURE_TYPES.has(value);
|
||||
}
|
||||
|
||||
function failureMetrics(error: TaskFormationModelError | SastEnrichmentModelError): FailureMetrics {
|
||||
return {
|
||||
costUsd: error.metrics.costUsd,
|
||||
modelCalls: error.metrics.modelCalls,
|
||||
inputTokens: error.metrics.inputTokens,
|
||||
outputTokens: error.metrics.outputTokens,
|
||||
};
|
||||
}
|
||||
|
||||
/** Build the one ApplicationFailure shape every reconciliation stage failure normalizes into. */
|
||||
function applicationFailure(
|
||||
type: ReconciliationStableFailureType,
|
||||
retryable: boolean,
|
||||
stage: ReconciliationActivityName,
|
||||
vulnerabilityClass: ReconciliationClass,
|
||||
details: StableFailureDetails = {},
|
||||
): ApplicationFailure {
|
||||
return ApplicationFailure.create({
|
||||
message: renderSafeMessage(SAFE_FAILURE_MESSAGES[type], { vulnerabilityClass }),
|
||||
type,
|
||||
nonRetryable: !retryable,
|
||||
details: [
|
||||
{
|
||||
stage,
|
||||
...(details.metrics !== undefined && { metrics: details.metrics }),
|
||||
...(details.fallbackReason !== undefined && { fallbackReason: details.fallbackReason }),
|
||||
},
|
||||
],
|
||||
});
|
||||
}
|
||||
|
||||
const CANCELLATION_CHAIN_DEPTH = 8;
|
||||
|
||||
/**
|
||||
* A failure counts as cancellation only when the activity signal is aborted AND its bounded
|
||||
* cause chain carries a real cancellation (the signal's own reason, a `CancelledFailure`, or
|
||||
* a cancellation-named abort raised under the aborted signal). A provider timeout, an
|
||||
* abort-shaped provider error with the signal unset, a cleanup failure, or any infrastructure
|
||||
* fault therefore stays an ordinary typed failure and is never manufactured into cancellation.
|
||||
*/
|
||||
function chainContainsRealCancellation(error: unknown, signal: AbortSignal): boolean {
|
||||
let current: unknown = error;
|
||||
const seen = new Set<unknown>();
|
||||
for (let depth = 0; depth < CANCELLATION_CHAIN_DEPTH; depth++) {
|
||||
if (current === undefined || current === null || seen.has(current)) return false;
|
||||
if (current === signal.reason) return true;
|
||||
if (current instanceof CancelledFailure) return true;
|
||||
if (current instanceof Error && (current.name === 'CancelledFailure' || current.name === 'AbortError')) return true;
|
||||
seen.add(current);
|
||||
current = current instanceof Error ? current.cause : undefined;
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
function cancellationFrom(error: unknown, signal: AbortSignal): CancelledFailure | undefined {
|
||||
if (!signal.aborted) return undefined;
|
||||
// A proactive check before any stage work has an aborted signal and no failure to inspect.
|
||||
if (error !== undefined && error !== null && !chainContainsRealCancellation(error, signal)) return undefined;
|
||||
|
||||
const reason = signal.reason;
|
||||
if (reason instanceof CancelledFailure) return reason;
|
||||
return new CancelledFailure('Reconciliation activity cancelled');
|
||||
}
|
||||
|
||||
/**
|
||||
* Map every shape a reconciliation stage can throw (a model-call error, a wrapped executor
|
||||
* error, an artifact-store error, an already-classified ApplicationFailure, or an unrecognized
|
||||
* error) onto the closed set of stable failure types. Cancellation is checked first and
|
||||
* always wins, since a stage aborted for cancellation is not a stage that failed.
|
||||
*/
|
||||
function normalizeFailure(
|
||||
error: unknown,
|
||||
stage: ReconciliationActivityName,
|
||||
vulnerabilityClass: ReconciliationClass,
|
||||
signal: AbortSignal,
|
||||
): never {
|
||||
const cancellation = cancellationFrom(error, signal);
|
||||
if (cancellation !== undefined) throw cancellation;
|
||||
|
||||
if (error instanceof TaskFormationModelError) {
|
||||
throw applicationFailure('TaskFormationModelError', error.retryable, stage, vulnerabilityClass, {
|
||||
metrics: failureMetrics(error),
|
||||
...(error.fallbackReason !== undefined && { fallbackReason: error.fallbackReason }),
|
||||
});
|
||||
}
|
||||
if (error instanceof SastEnrichmentModelError) {
|
||||
throw applicationFailure('SastEnrichmentModelError', error.retryable, stage, vulnerabilityClass, {
|
||||
metrics: failureMetrics(error),
|
||||
});
|
||||
}
|
||||
if (error instanceof ReconciliationError) {
|
||||
throw applicationFailure(error.failureType, error.retryable, stage, vulnerabilityClass);
|
||||
}
|
||||
if (error instanceof TaskFormationExecutorError) {
|
||||
if (error.failureKind === 'model') {
|
||||
throw applicationFailure('TaskFormationModelError', error.retryable, stage, vulnerabilityClass, {
|
||||
metrics: {
|
||||
costUsd: error.usage.costUsd,
|
||||
modelCalls: error.modelCalls,
|
||||
inputTokens: error.usage.inputTokens,
|
||||
outputTokens: error.usage.outputTokens,
|
||||
},
|
||||
...(error.fallbackReason !== undefined && { fallbackReason: error.fallbackReason }),
|
||||
});
|
||||
}
|
||||
// Retryable executor infrastructure faults (session setup, transient IO) must stay
|
||||
// retryable IO at the boundary instead of colliding with terminal ConfigurationError.
|
||||
if (error.failureKind === 'infrastructure') {
|
||||
throw applicationFailure('ReconciliationIoError', error.retryable, stage, vulnerabilityClass);
|
||||
}
|
||||
const type = error.failureKind === 'confinement' ? 'ArtifactIntegrityError' : 'ConfigurationError';
|
||||
throw applicationFailure(type, error.retryable, stage, vulnerabilityClass);
|
||||
}
|
||||
if (error instanceof ApplicationFailure) {
|
||||
const errorType = error.type;
|
||||
if (typeof errorType === 'string' && isStableFailureType(errorType)) {
|
||||
throw applicationFailure(errorType, !error.nonRetryable, stage, vulnerabilityClass);
|
||||
}
|
||||
throw applicationFailure('ReconciliationIoError', true, stage, vulnerabilityClass);
|
||||
}
|
||||
if (error instanceof Error && isStableFailureType(error.name)) {
|
||||
const retryable =
|
||||
'retryable' in error && typeof error.retryable === 'boolean' ? error.retryable : DEFAULT_RETRYABILITY[error.name];
|
||||
throw applicationFailure(error.name, retryable, stage, vulnerabilityClass);
|
||||
}
|
||||
|
||||
// Unknown failures remain retryable. A generic error name is not evidence that the fault is terminal.
|
||||
throw applicationFailure('ReconciliationIoError', true, stage, vulnerabilityClass);
|
||||
}
|
||||
|
||||
/** Refuse to schedule a class's remaining reconciliation stages once its 12-hour budget is spent. */
|
||||
function assertActivityCanRun(
|
||||
activityName: ReconciliationClassActivityName,
|
||||
classDeadlineMs: number,
|
||||
vulnerabilityClass: ReconciliationClass,
|
||||
nowMs: number,
|
||||
): ReturnType<typeof resolveReconciliationActivityBudget> {
|
||||
try {
|
||||
const budget = resolveReconciliationActivityBudget(activityName, classDeadlineMs, nowMs);
|
||||
if (!budget.shouldSchedule) {
|
||||
throw applicationFailure('ConfigurationError', false, activityName, vulnerabilityClass);
|
||||
}
|
||||
return budget;
|
||||
} catch (error) {
|
||||
if (error instanceof ApplicationFailure) throw error;
|
||||
throw applicationFailure('ConfigurationError', false, activityName, vulnerabilityClass);
|
||||
}
|
||||
}
|
||||
|
||||
async function runReconciliationStage<T>(
|
||||
activityName: ReconciliationClassActivityName,
|
||||
classDeadlineMs: number,
|
||||
vulnerabilityClass: ReconciliationClass,
|
||||
runtime: ReconciliationActivityRuntime,
|
||||
now: () => number,
|
||||
stage: (runtime: ReconciliationStageRuntime) => Promise<T>,
|
||||
activityModelHost: ModelHost,
|
||||
): Promise<T> {
|
||||
const cancellation = cancellationFrom(undefined, runtime.cancellationSignal);
|
||||
if (cancellation !== undefined) throw cancellation;
|
||||
|
||||
const budget = assertActivityCanRun(activityName, classDeadlineMs, vulnerabilityClass, now());
|
||||
const profile = RECONCILIATION_ACTIVITY_PROFILES[activityName];
|
||||
const startedAt = now();
|
||||
let heartbeatInterval: ReturnType<typeof setInterval> | undefined;
|
||||
|
||||
if (profile.profile === 'model' && budget.heartbeatIntervalMs !== null) {
|
||||
runtime.heartbeat({ stage: activityName, attempt: runtime.attempt, elapsedSeconds: 0, classDeadlineMs });
|
||||
heartbeatInterval = setInterval(() => {
|
||||
runtime.heartbeat({
|
||||
stage: activityName,
|
||||
attempt: runtime.attempt,
|
||||
elapsedSeconds: Math.max(0, Math.floor((now() - startedAt) / 1_000)),
|
||||
classDeadlineMs,
|
||||
});
|
||||
}, budget.heartbeatIntervalMs);
|
||||
}
|
||||
|
||||
// The executor's own timer must expire before Temporal's activity timeout, so the
|
||||
// metrics-bearing model-stage-timeout failure stays reachable. Evaluate the remaining
|
||||
// granted budget at call time because jail materialization can consume minutes first.
|
||||
const grantedBudgetMs = runtime.startToCloseTimeoutMs;
|
||||
const executorTimeoutMsFor = (): number | undefined => {
|
||||
if (grantedBudgetMs === undefined || grantedBudgetMs <= 0) return undefined;
|
||||
const remainingMs = startedAt + grantedBudgetMs - now();
|
||||
return Math.max(1_000, remainingMs - TASK_FORMATION_EXECUTOR_TIMEOUT_MARGIN_MS);
|
||||
};
|
||||
const executionContextFor = (): TaskFormationExecutionContext | undefined => ({
|
||||
attempt: runtime.attempt,
|
||||
...(runtime.executionKey !== undefined && { executionKey: runtime.executionKey }),
|
||||
});
|
||||
|
||||
try {
|
||||
return await stage({
|
||||
signal: runtime.cancellationSignal,
|
||||
logger: runtime.logger,
|
||||
modelHost: activityModelHost,
|
||||
executorTimeoutMsFor,
|
||||
executionContextFor,
|
||||
});
|
||||
} catch (error) {
|
||||
return normalizeFailure(error, activityName, vulnerabilityClass, runtime.cancellationSignal);
|
||||
} finally {
|
||||
if (heartbeatInterval !== undefined) clearInterval(heartbeatInterval);
|
||||
}
|
||||
}
|
||||
|
||||
async function runSeedStage<T>(runtime: ReconciliationActivityRuntime, stage: () => Promise<T>): Promise<T> {
|
||||
const cancellation = cancellationFrom(undefined, runtime.cancellationSignal);
|
||||
if (cancellation !== undefined) throw cancellation;
|
||||
|
||||
try {
|
||||
return await stage();
|
||||
} catch (error) {
|
||||
return normalizeFailure(error, 'seedEmptyProducerQueue', 'miscellaneous', runtime.cancellationSignal);
|
||||
}
|
||||
}
|
||||
|
||||
function defaultStages(workspacesDir: string): ReconciliationStageBindings {
|
||||
return {
|
||||
seedEmptyProducerQueue: seedEmptyProducerQueueStage,
|
||||
prepareClassReconciliation: prepareClassReconciliationStage,
|
||||
enrichClassSastObservations: (input, runtime) =>
|
||||
createEnrichClassSastObservations({
|
||||
generation: createPiStructuredGenerationPort(runtime.modelHost),
|
||||
modelContextFor: () => undefined,
|
||||
workspacesDir,
|
||||
signalFor: () => runtime.signal,
|
||||
logger: runtime.logger,
|
||||
})(input),
|
||||
formClassExploitTasks: (input, runtime) =>
|
||||
createFormClassExploitTasks({
|
||||
executor: createTaskFormationExecutor(runtime.modelHost),
|
||||
workspacesDir,
|
||||
signalFor: () => runtime.signal,
|
||||
logger: runtime.logger,
|
||||
...(runtime.executorTimeoutMsFor !== undefined && { executorTimeoutMsFor: runtime.executorTimeoutMsFor }),
|
||||
...(runtime.executionContextFor !== undefined && { executionContextFor: runtime.executionContextFor }),
|
||||
})(input),
|
||||
materializeClassExploitTasks: materializeClassExploitTasksStage,
|
||||
publishClassReconciliationOss: publishClassReconciliationOssStage,
|
||||
};
|
||||
}
|
||||
|
||||
function bindStages(
|
||||
workspacesDir: string,
|
||||
overrides: Partial<ReconciliationStageBindings> | undefined,
|
||||
): ReconciliationStageBindings {
|
||||
return { ...defaultStages(workspacesDir), ...overrides };
|
||||
}
|
||||
|
||||
function assertRegistryNames(registry: ReconciliationActivityRegistry): void {
|
||||
const actualNames = Object.keys(registry).sort();
|
||||
const expectedNames = [...RECONCILIATION_ACTIVITY_NAMES].sort();
|
||||
const expectedNamesAreUnique = new Set(expectedNames).size === expectedNames.length;
|
||||
if (
|
||||
!expectedNamesAreUnique ||
|
||||
actualNames.length !== expectedNames.length ||
|
||||
actualNames.some((name, index) => name !== expectedNames[index])
|
||||
) {
|
||||
throw new Error('Reconciliation activity registry does not match its frozen six-name contract');
|
||||
}
|
||||
}
|
||||
|
||||
/** Bind worker-local filesystem/model dependencies and return the frozen six-activity registry. */
|
||||
export function createReconciliationActivityRegistry(
|
||||
bindings: ReconciliationActivityBindings,
|
||||
): Readonly<ReconciliationActivityRegistry> {
|
||||
const now = bindings.now ?? Date.now;
|
||||
const runtimeFor = bindings.runtime ?? defaultRuntime;
|
||||
const activityModelHost = bindings.modelHost ?? modelHost;
|
||||
const stages = bindStages(bindings.workspacesDir, bindings.stages);
|
||||
|
||||
async function seedEmptyProducerQueue(input: SeedEmptyProducerQueueActivityInput) {
|
||||
const runtime = runtimeFor();
|
||||
const result = await runSeedStage(runtime, () =>
|
||||
stages.seedEmptyProducerQueue({
|
||||
deliverablesDir: bindings.deliverablesDir,
|
||||
sessionId: input.sessionId,
|
||||
logger: runtime.logger,
|
||||
}),
|
||||
);
|
||||
return {
|
||||
alreadySeeded: result.alreadySeeded,
|
||||
alreadyPublished: result.alreadyPublished,
|
||||
commitHash: result.commitHash,
|
||||
};
|
||||
}
|
||||
|
||||
async function prepareClassReconciliation(input: PrepareClassReconciliationActivityInput) {
|
||||
const runtime = runtimeFor();
|
||||
const result = await runReconciliationStage(
|
||||
'prepareClassReconciliation',
|
||||
input.classDeadlineMs,
|
||||
input.vulnerabilityClass,
|
||||
runtime,
|
||||
now,
|
||||
() =>
|
||||
stages.prepareClassReconciliation({
|
||||
deliverablesDir: bindings.deliverablesDir,
|
||||
sessionId: input.sessionId,
|
||||
vulnerabilityClass: input.vulnerabilityClass,
|
||||
// The contract fixes which fields the eventual published queue may carry for this
|
||||
// class. Passing it through unmodified is what keeps internal producer and
|
||||
// reconciliation identifiers out of the exploitation queue a downstream exploit
|
||||
// agent reads; widening it here would leak those identifiers into model-facing input.
|
||||
contract: publicationContractForClass(input.vulnerabilityClass, input.includeSastProvenance),
|
||||
workspacesDir: bindings.workspacesDir,
|
||||
}),
|
||||
activityModelHost,
|
||||
);
|
||||
if (result.outcome === 'already_published') {
|
||||
return { outcome: result.outcome, manifestSha256: result.manifestSha256 };
|
||||
}
|
||||
return { outcome: result.outcome, ref: result.ref };
|
||||
}
|
||||
|
||||
async function enrichClassSastObservations(input: EnrichClassSastObservationsActivityInput) {
|
||||
const runtime = runtimeFor();
|
||||
const result = await runReconciliationStage(
|
||||
'enrichClassSastObservations',
|
||||
input.classDeadlineMs,
|
||||
input.vulnerabilityClass,
|
||||
runtime,
|
||||
now,
|
||||
(stageRuntime) =>
|
||||
stages.enrichClassSastObservations(
|
||||
{
|
||||
sessionId: input.sessionId,
|
||||
vulnerabilityClass: input.vulnerabilityClass,
|
||||
...(input.sarif !== undefined && { sarif: input.sarif }),
|
||||
},
|
||||
stageRuntime,
|
||||
),
|
||||
activityModelHost,
|
||||
);
|
||||
return { ref: result.ref, metrics: result.metrics };
|
||||
}
|
||||
|
||||
async function formClassExploitTasks(
|
||||
input: FormClassExploitTasksActivityInput,
|
||||
): Promise<FormClassExploitTasksActivityResult> {
|
||||
const runtime = runtimeFor();
|
||||
const result = await runReconciliationStage(
|
||||
'formClassExploitTasks',
|
||||
input.classDeadlineMs,
|
||||
input.vulnerabilityClass,
|
||||
runtime,
|
||||
now,
|
||||
(stageRuntime) =>
|
||||
stages.formClassExploitTasks(
|
||||
{
|
||||
sessionId: input.sessionId,
|
||||
vulnerabilityClass: input.vulnerabilityClass,
|
||||
repositoryPath: bindings.repositoryPath,
|
||||
producerRef: input.producerRef,
|
||||
supplementalRef: input.supplementalRef,
|
||||
...(bindings.webUrl !== undefined && { webUrl: bindings.webUrl }),
|
||||
},
|
||||
stageRuntime,
|
||||
),
|
||||
activityModelHost,
|
||||
);
|
||||
return { ref: result.ref, metrics: result.metrics };
|
||||
}
|
||||
|
||||
async function materializeClassExploitTasks(input: MaterializeClassExploitTasksActivityInput) {
|
||||
const runtime = runtimeFor();
|
||||
const result = await runReconciliationStage(
|
||||
'materializeClassExploitTasks',
|
||||
input.classDeadlineMs,
|
||||
input.vulnerabilityClass,
|
||||
runtime,
|
||||
now,
|
||||
() =>
|
||||
stages.materializeClassExploitTasks({
|
||||
sessionId: input.sessionId,
|
||||
workspacesDir: bindings.workspacesDir,
|
||||
vulnerabilityClass: input.vulnerabilityClass,
|
||||
producerRef: input.producerRef,
|
||||
supplementalRef: input.supplementalRef,
|
||||
form: input.form,
|
||||
}),
|
||||
activityModelHost,
|
||||
);
|
||||
return { ref: result.ref };
|
||||
}
|
||||
|
||||
async function publishClassReconciliationOss(input: PublishClassReconciliationActivityInput) {
|
||||
const runtime = runtimeFor();
|
||||
const result = await runReconciliationStage(
|
||||
'publishClassReconciliationOss',
|
||||
input.classDeadlineMs,
|
||||
input.vulnerabilityClass,
|
||||
runtime,
|
||||
now,
|
||||
() =>
|
||||
stages.publishClassReconciliationOss({
|
||||
deliverablesDir: bindings.deliverablesDir,
|
||||
sessionId: input.sessionId,
|
||||
workspacesDir: bindings.workspacesDir,
|
||||
vulnerabilityClass: input.vulnerabilityClass,
|
||||
producerRef: input.producerRef,
|
||||
supplementalRef: input.supplementalRef,
|
||||
fixedTasksRef: input.fixedTasksRef,
|
||||
logger: runtime.logger,
|
||||
}),
|
||||
activityModelHost,
|
||||
);
|
||||
return {
|
||||
alreadyPublished: result.alreadyPublished,
|
||||
manifestSha256: result.manifestSha256,
|
||||
commitHash: result.commitHash,
|
||||
};
|
||||
}
|
||||
|
||||
const registry = {
|
||||
seedEmptyProducerQueue,
|
||||
prepareClassReconciliation,
|
||||
enrichClassSastObservations,
|
||||
formClassExploitTasks,
|
||||
materializeClassExploitTasks,
|
||||
publishClassReconciliationOss,
|
||||
} satisfies ReconciliationActivityRegistry;
|
||||
assertRegistryNames(registry);
|
||||
return Object.freeze(registry);
|
||||
}
|
||||
@@ -0,0 +1,278 @@
|
||||
// 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.
|
||||
|
||||
/** Workflow-safe reconciliation activity signatures and scheduling policy. */
|
||||
|
||||
import type { TaskFormationFallbackReason } from '../ai/pi/task-formation-executor.js';
|
||||
import type { ArtifactRef } from '../ai/reconciliation/contracts.js';
|
||||
import type { StageMetrics } from '../ai/reconciliation/stage-contracts.js';
|
||||
import type { SarifRef } from '../ai/sast/types.js';
|
||||
import type { ReconciliationClass } from '../types/reconciliation.js';
|
||||
|
||||
const MINUTE_MS = 60 * 1_000;
|
||||
const HOUR_MS = 60 * MINUTE_MS;
|
||||
|
||||
export const RECONCILIATION_CLASS_BUDGET_MS = 12 * HOUR_MS;
|
||||
export const RECONCILIATION_LATER_STAGE_RESERVE_MS = 5 * MINUTE_MS;
|
||||
|
||||
/**
|
||||
* Deterministic safety margin subtracted from the granted activity budget before it is passed
|
||||
* to the Pass 1 executor timer, so the executor's own timeout always fires before Temporal's
|
||||
* activity timeout and the metrics-bearing model-stage-timeout path stays reachable.
|
||||
*/
|
||||
export const TASK_FORMATION_EXECUTOR_TIMEOUT_MARGIN_MS = MINUTE_MS;
|
||||
|
||||
/**
|
||||
* Workflow-safe mirror of Agent A's closed fallback-reason set. The executor module itself is
|
||||
* not bundle-safe, so the workflow validates deserialized failure details against this frozen
|
||||
* copy; the `satisfies` clause and the exhaustiveness check keep the two sets identical at
|
||||
* compile time, and the activity boundary re-asserts equality at module load.
|
||||
*/
|
||||
export const ACCEPTED_TASK_FORMATION_FALLBACK_REASONS = Object.freeze([
|
||||
'retryable_model_failure',
|
||||
'missing_accepted_submission',
|
||||
'model_stage_timeout',
|
||||
] as const satisfies readonly TaskFormationFallbackReason[]);
|
||||
|
||||
type UnlistedFallbackReason = Exclude<
|
||||
TaskFormationFallbackReason,
|
||||
(typeof ACCEPTED_TASK_FORMATION_FALLBACK_REASONS)[number]
|
||||
>;
|
||||
const _everyFallbackReasonIsListed: UnlistedFallbackReason extends never ? true : never = true;
|
||||
void _everyFallbackReasonIsListed;
|
||||
|
||||
/** Validate one deserialized fallback reason against the closed set. */
|
||||
export function isAcceptedTaskFormationFallbackReason(value: unknown): value is TaskFormationFallbackReason {
|
||||
return (ACCEPTED_TASK_FORMATION_FALLBACK_REASONS as readonly unknown[]).includes(value);
|
||||
}
|
||||
|
||||
export interface ReconciliationActivityDeadline {
|
||||
/** Fixed workflow-derived deadline for this class, measured as Unix epoch milliseconds. */
|
||||
readonly classDeadlineMs: number;
|
||||
}
|
||||
|
||||
export interface ReconciliationActivityBaseInput extends ReconciliationActivityDeadline {
|
||||
readonly sessionId: string;
|
||||
readonly vulnerabilityClass: ReconciliationClass;
|
||||
}
|
||||
|
||||
export interface SeedEmptyProducerQueueActivityInput {
|
||||
readonly sessionId: string;
|
||||
}
|
||||
|
||||
export interface PrepareClassReconciliationActivityInput extends ReconciliationActivityBaseInput {
|
||||
readonly includeSastProvenance: boolean;
|
||||
}
|
||||
|
||||
export interface EnrichClassSastObservationsActivityInput extends ReconciliationActivityBaseInput {
|
||||
readonly sarif?: SarifRef;
|
||||
}
|
||||
|
||||
export interface FormClassExploitTasksActivityInput extends ReconciliationActivityBaseInput {
|
||||
readonly producerRef: ArtifactRef<'producer-observations'>;
|
||||
readonly supplementalRef: ArtifactRef<'supplemental-observations'>;
|
||||
}
|
||||
|
||||
export interface FormClassExploitTasksActivityResult {
|
||||
readonly ref: ArtifactRef<'task-formation'>;
|
||||
readonly metrics: StageMetrics;
|
||||
}
|
||||
|
||||
export type MaterializationFormationResult = FormClassExploitTasksActivityResult | 'singleton_fallback';
|
||||
|
||||
export interface MaterializeClassExploitTasksActivityInput extends ReconciliationActivityBaseInput {
|
||||
readonly producerRef: ArtifactRef<'producer-observations'>;
|
||||
readonly supplementalRef: ArtifactRef<'supplemental-observations'>;
|
||||
readonly form: MaterializationFormationResult;
|
||||
}
|
||||
|
||||
export interface PublishClassReconciliationActivityInput extends ReconciliationActivityBaseInput {
|
||||
readonly producerRef: ArtifactRef<'producer-observations'>;
|
||||
readonly supplementalRef: ArtifactRef<'supplemental-observations'>;
|
||||
readonly fixedTasksRef: ArtifactRef<'fixed-tasks'>;
|
||||
}
|
||||
|
||||
export interface SeedEmptyProducerQueueActivityResult {
|
||||
readonly alreadySeeded: boolean;
|
||||
readonly alreadyPublished: boolean;
|
||||
readonly commitHash: string;
|
||||
}
|
||||
|
||||
export type PrepareClassReconciliationActivityResult =
|
||||
| {
|
||||
readonly outcome: 'already_published';
|
||||
readonly manifestSha256: string;
|
||||
}
|
||||
| {
|
||||
readonly outcome: 'pending';
|
||||
readonly ref: ArtifactRef<'producer-observations'>;
|
||||
};
|
||||
|
||||
export interface EnrichClassSastObservationsActivityResult {
|
||||
readonly ref: ArtifactRef<'supplemental-observations'>;
|
||||
readonly metrics: StageMetrics;
|
||||
}
|
||||
|
||||
export interface MaterializeClassExploitTasksActivityResult {
|
||||
readonly ref: ArtifactRef<'fixed-tasks'>;
|
||||
}
|
||||
|
||||
export interface PublishClassReconciliationActivityResult {
|
||||
readonly alreadyPublished: boolean;
|
||||
readonly manifestSha256: string;
|
||||
readonly commitHash: string;
|
||||
}
|
||||
|
||||
export interface ReconciliationActivityRegistry {
|
||||
readonly seedEmptyProducerQueue: (
|
||||
input: SeedEmptyProducerQueueActivityInput,
|
||||
) => Promise<SeedEmptyProducerQueueActivityResult>;
|
||||
readonly prepareClassReconciliation: (
|
||||
input: PrepareClassReconciliationActivityInput,
|
||||
) => Promise<PrepareClassReconciliationActivityResult>;
|
||||
readonly enrichClassSastObservations: (
|
||||
input: EnrichClassSastObservationsActivityInput,
|
||||
) => Promise<EnrichClassSastObservationsActivityResult>;
|
||||
readonly formClassExploitTasks: (
|
||||
input: FormClassExploitTasksActivityInput,
|
||||
) => Promise<FormClassExploitTasksActivityResult>;
|
||||
readonly materializeClassExploitTasks: (
|
||||
input: MaterializeClassExploitTasksActivityInput,
|
||||
) => Promise<MaterializeClassExploitTasksActivityResult>;
|
||||
readonly publishClassReconciliationOss: (
|
||||
input: PublishClassReconciliationActivityInput,
|
||||
) => Promise<PublishClassReconciliationActivityResult>;
|
||||
}
|
||||
|
||||
export const RECONCILIATION_ACTIVITY_NAMES = Object.freeze([
|
||||
'seedEmptyProducerQueue',
|
||||
'prepareClassReconciliation',
|
||||
'enrichClassSastObservations',
|
||||
'formClassExploitTasks',
|
||||
'materializeClassExploitTasks',
|
||||
'publishClassReconciliationOss',
|
||||
] as const satisfies readonly (keyof ReconciliationActivityRegistry)[]);
|
||||
|
||||
export type ReconciliationActivityName = (typeof RECONCILIATION_ACTIVITY_NAMES)[number];
|
||||
export type ReconciliationClassActivityName = Exclude<ReconciliationActivityName, 'seedEmptyProducerQueue'>;
|
||||
|
||||
export type ReconciliationActivityProfileName = 'deterministic' | 'model';
|
||||
|
||||
export interface ReconciliationActivityProfile {
|
||||
readonly profile: ReconciliationActivityProfileName;
|
||||
readonly startToCloseTimeoutMs: number;
|
||||
readonly heartbeatTimeoutMs: number | null;
|
||||
readonly maximumAttempts: number;
|
||||
readonly retryInitialIntervalMs: number;
|
||||
readonly retryBackoffCoefficient: number;
|
||||
}
|
||||
|
||||
function profile(
|
||||
name: ReconciliationActivityProfileName,
|
||||
startToCloseTimeoutMs: number,
|
||||
heartbeatTimeoutMs: number | null,
|
||||
maximumAttempts: number,
|
||||
): Readonly<ReconciliationActivityProfile> {
|
||||
return Object.freeze({
|
||||
profile: name,
|
||||
startToCloseTimeoutMs,
|
||||
heartbeatTimeoutMs,
|
||||
maximumAttempts,
|
||||
retryInitialIntervalMs: 1_000,
|
||||
retryBackoffCoefficient: 2,
|
||||
});
|
||||
}
|
||||
|
||||
// Deterministic stages (filesystem/git only) get a short timeout, no heartbeat, and more
|
||||
// attempts, since a transient IO failure is cheap to retry. Model-backed stages get a long
|
||||
// timeout, a heartbeat so a wedged model call is detected before the full timeout elapses, and
|
||||
// fewer attempts, since each attempt can itself cost real time and money.
|
||||
const DETERMINISTIC_PROFILE = profile('deterministic', 2 * MINUTE_MS, null, 5);
|
||||
const MODEL_PROFILE = profile('model', 30 * MINUTE_MS, 5 * MINUTE_MS, 3);
|
||||
|
||||
export const RECONCILIATION_ACTIVITY_PROFILES = Object.freeze({
|
||||
seedEmptyProducerQueue: DETERMINISTIC_PROFILE,
|
||||
prepareClassReconciliation: DETERMINISTIC_PROFILE,
|
||||
enrichClassSastObservations: MODEL_PROFILE,
|
||||
formClassExploitTasks: MODEL_PROFILE,
|
||||
materializeClassExploitTasks: DETERMINISTIC_PROFILE,
|
||||
publishClassReconciliationOss: DETERMINISTIC_PROFILE,
|
||||
} as const satisfies Readonly<Record<ReconciliationActivityName, Readonly<ReconciliationActivityProfile>>>);
|
||||
|
||||
export interface ReconciliationActivityBudget {
|
||||
/** False once the class's 12-hour deadline leaves no time for this stage; the caller must not schedule it. */
|
||||
readonly shouldSchedule: boolean;
|
||||
readonly remainingClassBudgetMs: number;
|
||||
readonly reservedForLaterStagesMs: number;
|
||||
readonly scheduleToCloseTimeoutMs: number;
|
||||
readonly startToCloseTimeoutMs: number;
|
||||
readonly heartbeatTimeoutMs: number | null;
|
||||
readonly heartbeatIntervalMs: number | null;
|
||||
}
|
||||
|
||||
function assertEpochMilliseconds(value: number, label: string): void {
|
||||
if (!Number.isSafeInteger(value) || value < 0) {
|
||||
throw new Error(`${label} must be non-negative safe-integer epoch milliseconds`);
|
||||
}
|
||||
}
|
||||
|
||||
/** Derive the fixed per-class deadline immediately before the caller enters prepare. */
|
||||
export function reconciliationClassDeadlineFrom(startedAtMs: number): number {
|
||||
assertEpochMilliseconds(startedAtMs, 'Reconciliation class start');
|
||||
const deadlineMs = startedAtMs + RECONCILIATION_CLASS_BUDGET_MS;
|
||||
assertEpochMilliseconds(deadlineMs, 'Reconciliation class deadline');
|
||||
return deadlineMs;
|
||||
}
|
||||
|
||||
/**
|
||||
* Cap one activity's whole retry window, per-attempt timeout, and heartbeat against the class deadline.
|
||||
* A zero timeout means the caller must not schedule the activity.
|
||||
*/
|
||||
export function resolveReconciliationActivityBudget(
|
||||
activityName: ReconciliationClassActivityName,
|
||||
classDeadlineMs: number,
|
||||
nowMs: number,
|
||||
): Readonly<ReconciliationActivityBudget> {
|
||||
assertEpochMilliseconds(classDeadlineMs, 'Reconciliation class deadline');
|
||||
assertEpochMilliseconds(nowMs, 'Reconciliation budget time');
|
||||
|
||||
const profile = RECONCILIATION_ACTIVITY_PROFILES[activityName];
|
||||
const remainingClassBudgetMs = Math.max(0, classDeadlineMs - nowMs);
|
||||
const reservedForLaterStagesMs = Math.min(remainingClassBudgetMs, RECONCILIATION_LATER_STAGE_RESERVE_MS);
|
||||
const scheduleToCloseTimeoutMs = Math.max(0, remainingClassBudgetMs - RECONCILIATION_LATER_STAGE_RESERVE_MS);
|
||||
const startToCloseTimeoutMs = Math.min(profile.startToCloseTimeoutMs, scheduleToCloseTimeoutMs);
|
||||
const heartbeatTimeoutMs =
|
||||
profile.heartbeatTimeoutMs === null || startToCloseTimeoutMs === 0
|
||||
? null
|
||||
: Math.min(profile.heartbeatTimeoutMs, startToCloseTimeoutMs);
|
||||
const heartbeatIntervalMs =
|
||||
heartbeatTimeoutMs === null ? null : Math.max(1, Math.min(100_000, Math.floor(heartbeatTimeoutMs / 3)));
|
||||
|
||||
return Object.freeze({
|
||||
shouldSchedule: scheduleToCloseTimeoutMs > 0,
|
||||
remainingClassBudgetMs,
|
||||
reservedForLaterStagesMs,
|
||||
scheduleToCloseTimeoutMs,
|
||||
startToCloseTimeoutMs,
|
||||
heartbeatTimeoutMs,
|
||||
heartbeatIntervalMs,
|
||||
});
|
||||
}
|
||||
|
||||
export const RECONCILIATION_STABLE_FAILURE_TYPES = Object.freeze([
|
||||
'TaskFormationModelError',
|
||||
'SastEnrichmentModelError',
|
||||
'ReconciliationArtifactNotFound',
|
||||
'ReconciliationIoError',
|
||||
'ConfigurationError',
|
||||
'SastEnrichmentInputError',
|
||||
'ArtifactIntegrityError',
|
||||
'PublicationConflict',
|
||||
'UnmappableSurvivor',
|
||||
'KeySetDivergence',
|
||||
] as const);
|
||||
|
||||
export type ReconciliationStableFailureType = (typeof RECONCILIATION_STABLE_FAILURE_TYPES)[number];
|
||||
@@ -1,16 +1,97 @@
|
||||
import { defineQuery } from '@temporalio/workflow';
|
||||
import { defineQuery, defineSignal } from '@temporalio/workflow';
|
||||
|
||||
export type { AgentMetrics } from '../types/metrics.js';
|
||||
|
||||
import type { DistributedConfig, VulnClass } from '../types/config.js';
|
||||
import type {
|
||||
AgenticSastReduction,
|
||||
CapellaFailurePoint,
|
||||
CapellaRecoveredFailure,
|
||||
CapellaStage,
|
||||
SarifRef,
|
||||
} from '../ai/sast/types.js';
|
||||
import type { VulnClass } from '../types/config.js';
|
||||
import type { ErrorCode } from '../types/errors.js';
|
||||
import type { AgentMetrics } from '../types/metrics.js';
|
||||
import type { ReconciliationClass } from '../types/reconciliation.js';
|
||||
import type {
|
||||
MiscellaneousOutcome,
|
||||
PartialReasonView,
|
||||
ReportProgress,
|
||||
ReportSarifDisposition,
|
||||
StoredPdfProvenance,
|
||||
} from '../types/run-state.js';
|
||||
|
||||
/**
|
||||
* The serializable slice of Capella's configuration passed across the Temporal workflow
|
||||
* boundary into the child workflow input. Everything the workflow needs from the parsed
|
||||
* config or the model spec must be flattened into plain data here; the workflow sandbox
|
||||
* cannot carry functions or class instances across that boundary.
|
||||
*/
|
||||
export interface AgenticSastInput {
|
||||
readonly codePathAvoids: readonly string[];
|
||||
readonly codePathFocus: readonly string[];
|
||||
readonly modelSpec: string;
|
||||
readonly capellaFormatVersion: string;
|
||||
readonly promptSetVersion: string;
|
||||
}
|
||||
|
||||
/**
|
||||
* The agentic SAST lifecycle as seen from the pentest workflow: not configured, running as a
|
||||
* child workflow, or one of two terminal outcomes. This is what the live `getProgress` query
|
||||
* and the terminal `PipelineState` both report, so a caller never needs to inspect the Capella
|
||||
* child workflow's own result type directly.
|
||||
*/
|
||||
export type AgenticSastState =
|
||||
| { readonly status: 'disabled' }
|
||||
| { readonly status: 'running'; readonly startedAt: number }
|
||||
| {
|
||||
readonly status: 'succeeded';
|
||||
readonly findingCount: number;
|
||||
readonly sarifSha256: string;
|
||||
readonly coverage: 'complete' | 'reduced';
|
||||
readonly warnings: readonly string[];
|
||||
readonly durationMs: number;
|
||||
readonly reductions?: readonly AgenticSastReduction[];
|
||||
readonly recoveredFailure?: CapellaRecoveredFailure;
|
||||
}
|
||||
| {
|
||||
readonly status: 'failed';
|
||||
readonly failedStage: CapellaFailurePoint;
|
||||
/** Reader-facing name of `failedStage`, projected once so no surface renders the slug. */
|
||||
readonly failedStageLabel: string;
|
||||
readonly error: string;
|
||||
/** Bounded machine code preserved from the failing Capella activity, when one crossed the child. */
|
||||
readonly errorCode?: string;
|
||||
readonly completedStages: readonly CapellaStage[];
|
||||
readonly warnings: readonly string[];
|
||||
readonly durationMs: number;
|
||||
};
|
||||
|
||||
export type OperationalStageStatus = 'pending' | 'running' | 'completed' | 'failed' | 'skipped';
|
||||
|
||||
export interface OperationalStageState {
|
||||
readonly key: string;
|
||||
readonly label: string;
|
||||
readonly status: OperationalStageStatus;
|
||||
readonly startedAt?: number;
|
||||
readonly durationMs?: number;
|
||||
readonly error?: string;
|
||||
}
|
||||
|
||||
export interface OperationalMetrics extends AgentMetrics {
|
||||
readonly usageComplete?: boolean;
|
||||
}
|
||||
|
||||
/** A degradation the scan recorded and continued past, kept for the terminal summary log rather than for control flow. */
|
||||
export interface NonFatalFailure {
|
||||
readonly phase: string;
|
||||
readonly error: string;
|
||||
}
|
||||
|
||||
export interface PipelineInput {
|
||||
webUrl: string;
|
||||
repoPath: string;
|
||||
configPath?: string;
|
||||
outputPath?: string;
|
||||
pipelineTestingMode?: boolean;
|
||||
workflowId?: string; // Used for audit correlation
|
||||
sessionId?: string; // Workspace directory name (distinct from workflowId for named workspaces)
|
||||
@@ -19,44 +100,114 @@ export interface PipelineInput {
|
||||
|
||||
// Config fields — serializable, flow through to ActivityInput → getOrCreateContainer()
|
||||
configYAML?: string; // Raw YAML string (parsed in activity, not workflow — workflow sandbox can't use Node.js)
|
||||
configData?: DistributedConfig; // Pre-parsed config (bypasses file loading)
|
||||
deliverablesSubdir?: string; // Override deliverables path (default: '.shannon/deliverables')
|
||||
auditDir?: string; // Override audit log directory (default: './workspaces')
|
||||
promptDir?: string; // Override prompt template directory
|
||||
sastSarifPath?: string; // Optional path for consumer-supplied findings input
|
||||
agenticSast?: AgenticSastInput;
|
||||
sastSarif?: SarifRef;
|
||||
customerOutputPath?: string; // Stable mounted path for final customer copies only
|
||||
checkpointsEnabled?: boolean; // Enable checkpoint activities (default: false)
|
||||
vulnClasses?: VulnClass[]; // omitted = all five
|
||||
exploit?: boolean; // false skips the exploitation phase
|
||||
}
|
||||
|
||||
/** What `loadResumeState` reconstructs from a prior workspace: independently verified, never assumed from session.json alone. */
|
||||
export interface ResumeState {
|
||||
workspaceName: string;
|
||||
originalUrl: string;
|
||||
completedAgents: string[];
|
||||
checkpointHash: string;
|
||||
originalWorkflowId: string;
|
||||
expectedAgents: string[];
|
||||
participatingClasses: ReconciliationClass[];
|
||||
exploit: boolean;
|
||||
reportProgress?: ReportProgress;
|
||||
miscellaneousOutcome?: MiscellaneousOutcome;
|
||||
}
|
||||
|
||||
/** The narrow view of the durable scan-state record the workflow needs to keep its own queryable state in sync. */
|
||||
export interface DurableStateSummary {
|
||||
readonly exploit: boolean;
|
||||
readonly expectedAgents: readonly string[];
|
||||
readonly participatingClasses: readonly ReconciliationClass[];
|
||||
readonly reportStage: ReportProgress['stage'] | 'uninitialized';
|
||||
readonly miscellaneousOutcome?: MiscellaneousOutcome;
|
||||
}
|
||||
|
||||
/** Common result shape for the deterministic report-processing activities (renumber, compaction). */
|
||||
export interface ReconciliationActivityResult {
|
||||
readonly vulnerabilityClass?: ReconciliationClass;
|
||||
readonly skipped: boolean;
|
||||
readonly changedPathCount: number;
|
||||
readonly checkpoint?: string;
|
||||
readonly alreadyCommitted?: boolean;
|
||||
}
|
||||
|
||||
export interface FinalizeReportActivityResult {
|
||||
readonly checkpoint: string;
|
||||
readonly manifestSha256: string;
|
||||
readonly changedPathCount: number;
|
||||
readonly alreadyCommitted: boolean;
|
||||
/** Adopted-or-produced SARIF disposition from the committed finalization manifest. */
|
||||
readonly sarifDisposition: ReportSarifDisposition;
|
||||
readonly pdfGenerated: boolean;
|
||||
/** Verified provenance for the current PDF bytes, or null when no trustworthy PDF exists. */
|
||||
readonly pdfProvenance: StoredPdfProvenance | null;
|
||||
readonly warningCount: number;
|
||||
}
|
||||
|
||||
export interface AssembleReportActivityResult {
|
||||
/** Classes whose findings could not be included in the assembled report inputs. */
|
||||
readonly failedClasses: readonly ReconciliationClass[];
|
||||
}
|
||||
|
||||
export interface SurfaceReportActivityResult {
|
||||
readonly surfaced: readonly string[];
|
||||
readonly removedStale: readonly string[];
|
||||
readonly warningCount: number;
|
||||
}
|
||||
|
||||
export interface PipelineSummary {
|
||||
totalCostUsd: number;
|
||||
totalDurationMs: number; // Wall-clock time (end - start)
|
||||
totalTurns: number;
|
||||
/** Total resolved agents: those that ran plus those that were skipped. */
|
||||
agentCount: number;
|
||||
/** False when operational (Capella/reconciliation) spend is known to be incomplete. */
|
||||
usageAccountingComplete: boolean;
|
||||
}
|
||||
|
||||
/**
|
||||
* The workflow's whole queryable and terminal state. The CLI cannot import this package, so
|
||||
* `apps/cli/src/scan/pipeline.ts` mirrors this shape (along with AgentMetrics and the
|
||||
* activity-name-to-agent map) by hand; a field added, renamed, or removed here needs the same
|
||||
* change there, or the CLI's status rendering silently falls out of sync with a running scan.
|
||||
*/
|
||||
export interface PipelineState {
|
||||
status: 'running' | 'completed' | 'failed' | 'cancelled' | 'partial';
|
||||
currentPhase: string | null;
|
||||
currentAgent: string | null;
|
||||
/** Agents that actually ran. Mutually exclusive from `skippedAgents`. */
|
||||
completedAgents: string[];
|
||||
/** Expected agents that never ran because their class had nothing to exploit. */
|
||||
skippedAgents: string[];
|
||||
expectedAgents: string[];
|
||||
participatingClasses: ReconciliationClass[];
|
||||
// Vuln classes whose pipeline failed while at least one other succeeded. Drives the
|
||||
// partial terminal status so a crashed class isn't reported as if it fully passed.
|
||||
failedPipelines: { vulnType: VulnClass; error: string }[];
|
||||
failedReconciliations: { vulnerabilityClass: ReconciliationClass; error: string }[];
|
||||
failedAgent: string | null;
|
||||
error: string | null;
|
||||
errorCode?: ErrorCode;
|
||||
startTime: number;
|
||||
agentMetrics: Record<string, AgentMetrics>;
|
||||
operationalMetrics: Record<string, OperationalMetrics>;
|
||||
operationalStages: Record<string, OperationalStageState>;
|
||||
agenticSast: AgenticSastState;
|
||||
nonFatalFailures: NonFatalFailure[];
|
||||
/** Ordered durable degradation reasons with derived safe messages; empty for a full success. */
|
||||
partialReasons: PartialReasonView[];
|
||||
reportProgress?: ReportProgress;
|
||||
summary: PipelineSummary | null;
|
||||
}
|
||||
|
||||
@@ -79,3 +230,20 @@ export interface VulnExploitPipelineResult {
|
||||
}
|
||||
|
||||
export const getProgress = defineQuery<PipelineProgress>('getProgress');
|
||||
|
||||
/**
|
||||
* One Capella stage transition, reported by the SAST child workflow to its parent.
|
||||
*
|
||||
* Capella runs as a child workflow, so its activities never appear in the parent's
|
||||
* pending activities and the CLI cannot observe them. This signal is how per-stage
|
||||
* progress reaches the parent's durable `operationalStages`, which is what both the
|
||||
* live `getProgress` query and the terminal result render from.
|
||||
*/
|
||||
export interface CapellaStageProgress {
|
||||
readonly stage: CapellaStage;
|
||||
readonly status: 'running' | 'completed' | 'failed';
|
||||
readonly startedAt: number;
|
||||
readonly durationMs?: number;
|
||||
}
|
||||
|
||||
export const capellaStageProgress = defineSignal<[CapellaStageProgress]>('capellaStageProgress');
|
||||
@@ -1,4 +1,4 @@
|
||||
// Copyright (C) 2025 Keygraph, Inc.
|
||||
// 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
|
||||
@@ -29,14 +29,59 @@ export function toWorkflowSummary(
|
||||
throw new Error('toWorkflowSummary: state.summary must be set before calling');
|
||||
}
|
||||
|
||||
// The failure detail is one of the child workflow's fixed safe sentences, so it carries no
|
||||
// provider, prompt, repository, or path content and travels with the stable code.
|
||||
const agenticSastFailure = state.agenticSast.status === 'failed' ? state.agenticSast : undefined;
|
||||
const agenticSastErrorCode = agenticSastFailure?.errorCode;
|
||||
const agenticSastFailureMessage = agenticSastFailure?.error;
|
||||
const agenticSastFailedStage = agenticSastFailure?.failedStageLabel;
|
||||
// Both terminal Capella variants carry a warnings array; a disabled or still-running run has none.
|
||||
const agenticSast = state.agenticSast;
|
||||
const usageAccountingWarnings =
|
||||
agenticSast.status === 'succeeded' || agenticSast.status === 'failed' ? agenticSast.warnings : [];
|
||||
// Carry the terminal disposition so a successful (or reduced-coverage) run is visible in the summary,
|
||||
// not only a failed one. Coverage is meaningful only on success.
|
||||
const agenticSastCoverage = agenticSast.status === 'succeeded' ? agenticSast.coverage : undefined;
|
||||
const endedAtMs = state.startTime + summary.totalDurationMs;
|
||||
return {
|
||||
status,
|
||||
startedAtMs: state.startTime,
|
||||
endedAtMs,
|
||||
totalDurationMs: summary.totalDurationMs,
|
||||
totalCostUsd: summary.totalCostUsd,
|
||||
completedAgents: state.completedAgents,
|
||||
skippedAgents: state.skippedAgents,
|
||||
agentMetrics: Object.fromEntries(
|
||||
Object.entries(state.agentMetrics).map(([name, m]) => [name, { durationMs: m.durationMs, costUsd: m.costUsd }]),
|
||||
Object.entries(state.agentMetrics).map(([name, metrics]) => [
|
||||
name,
|
||||
{ durationMs: metrics.durationMs, costUsd: metrics.costUsd },
|
||||
]),
|
||||
),
|
||||
...(state.error && { error: state.error }),
|
||||
operationalMetrics: Object.fromEntries(
|
||||
Object.entries(state.operationalMetrics).map(([name, metrics]) => [
|
||||
name,
|
||||
{ ...metrics, usageComplete: metrics.usageComplete !== false },
|
||||
]),
|
||||
),
|
||||
// The stage wall-clocks the summary reads to report each group's real elapsed time; the priced
|
||||
// metrics carry cost but no faithful duration (reconciliation stages record 0).
|
||||
operationalStages: Object.fromEntries(
|
||||
Object.entries(state.operationalStages).map(([key, stage]) => [
|
||||
key,
|
||||
{
|
||||
...(stage.startedAt !== undefined && { startedAt: stage.startedAt }),
|
||||
...(stage.durationMs !== undefined && { durationMs: stage.durationMs }),
|
||||
},
|
||||
]),
|
||||
),
|
||||
partialReasons: state.partialReasons,
|
||||
usageAccountingComplete: summary.usageAccountingComplete,
|
||||
usageAccountingWarnings: [...usageAccountingWarnings],
|
||||
agenticSastStatus: agenticSast.status,
|
||||
...(agenticSastCoverage !== undefined && { agenticSastCoverage }),
|
||||
...(agenticSastFailedStage !== undefined && { agenticSastFailedStage }),
|
||||
...(agenticSastFailureMessage !== undefined && { agenticSastFailureMessage }),
|
||||
...(agenticSastErrorCode !== undefined && { agenticSastErrorCode }),
|
||||
...(state.errorCode !== undefined && { errorCode: state.errorCode }),
|
||||
};
|
||||
}
|
||||
@@ -1,6 +1,6 @@
|
||||
#!/usr/bin/env node
|
||||
|
||||
// Copyright (C) 2025 Keygraph, Inc.
|
||||
// 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
|
||||
@@ -18,8 +18,9 @@
|
||||
*
|
||||
* Options:
|
||||
* --task-queue <name> Task queue name (required, unique per scan)
|
||||
* --workflow-id <id> Workflow ID selected by the Shannon CLI
|
||||
* --config <path> Configuration file path
|
||||
* --output <path> Output directory for workspaces
|
||||
* --output <path> Stable mounted path for final customer report copies
|
||||
* --workspace <name> Resume from existing workspace
|
||||
* --pipeline-testing Use minimal prompts for fast testing
|
||||
*
|
||||
@@ -27,24 +28,69 @@
|
||||
* TEMPORAL_ADDRESS - Temporal server address (default: localhost:7233)
|
||||
*/
|
||||
|
||||
import fs from 'node:fs';
|
||||
import path from 'node:path';
|
||||
import { fileURLToPath } from 'node:url';
|
||||
import { Client, Connection, type WorkflowHandle, WorkflowNotFoundError } from '@temporalio/client';
|
||||
import { bundleWorkflowCode, NativeConnection, Worker } from '@temporalio/worker';
|
||||
import dotenv from 'dotenv';
|
||||
import { DEFAULT_MODEL_SPEC } from '../ai/models.js';
|
||||
import { capellaTerminalStageLabel, isCapellaSafeFailureMessage } from '../ai/sast/capella/safe-failures.js';
|
||||
import { capellaActivities, mergeActivityRegistries } from '../ai/sast/capella/temporal/registry.js';
|
||||
import { CAPELLA_FORMAT_VERSION, CAPELLA_PROMPT_SET_VERSION } from '../ai/sast/capella/types.js';
|
||||
import { summarizeOperationalMetrics } from '../audit/operational-summary.js';
|
||||
import { sanitizeHostname } from '../audit/utils.js';
|
||||
import { parseConfig } from '../config-parser.js';
|
||||
import { distributeConfig, parseConfig } from '../config-parser.js';
|
||||
import { deliverablesDir, resolveSessionJsonPath } from '../paths.js';
|
||||
import { isProviderFailureCategory } from '../types/errors.js';
|
||||
import {
|
||||
ASSEMBLED_REPORT_PDF_FILENAME,
|
||||
deliverablesDir,
|
||||
FINAL_REPORT_PDF_FILENAME,
|
||||
resolveSessionJsonPath,
|
||||
} from '../paths.js';
|
||||
import type { VulnClass } from '../types/config.js';
|
||||
ACCEPTED_CAPELLA_FAILURE_STAGES,
|
||||
isPartialReason,
|
||||
projectPartialReasons,
|
||||
SAFE_RUN_STATE_MESSAGES,
|
||||
workspaceExploitMismatchMessage,
|
||||
} from '../types/run-state.js';
|
||||
import { fileExists, readJson } from '../utils/file-io.js';
|
||||
import * as activities from './activities.js';
|
||||
import type { PipelineInput, PipelineProgress, PipelineState } from './shared.js';
|
||||
import {
|
||||
assembleReportActivity,
|
||||
checkExploitationQueue,
|
||||
compactReportFindings,
|
||||
finalizeReportOutputs,
|
||||
initDeliverableGit,
|
||||
initializeDurableScanState,
|
||||
initializeReportProgress,
|
||||
loadResumeState,
|
||||
logPhaseTransition,
|
||||
logWorkflowComplete,
|
||||
persistCanonicalReportProgress,
|
||||
persistFinalizedReportProgress,
|
||||
persistMiscellaneousOutcome,
|
||||
recordResumeAttempt,
|
||||
registerResumeAttempt,
|
||||
renumberClassFindings,
|
||||
restoreGitCheckpoint,
|
||||
runAuthExploitAgent,
|
||||
runAuthenticationValidation,
|
||||
runAuthVulnAgent,
|
||||
runAuthzExploitAgent,
|
||||
runAuthzVulnAgent,
|
||||
runInjectionExploitAgent,
|
||||
runInjectionVulnAgent,
|
||||
runMiscellaneousExploitAgent,
|
||||
runPreflightValidation,
|
||||
runPreReconAgent,
|
||||
runReconAgent,
|
||||
runReportAgent,
|
||||
runSsrfExploitAgent,
|
||||
runSsrfVulnAgent,
|
||||
runXssExploitAgent,
|
||||
runXssVulnAgent,
|
||||
saveCheckpoint,
|
||||
surfaceReportOutputs,
|
||||
syncCodePathDenyRules,
|
||||
syncPlaywrightStealthConfig,
|
||||
} from './activities.js';
|
||||
import { createReconciliationActivityRegistry } from './reconcile-activities.js';
|
||||
import type { AgenticSastInput, PipelineInput, PipelineProgress, PipelineState } from './shared.js';
|
||||
|
||||
dotenv.config();
|
||||
|
||||
@@ -52,14 +98,152 @@ const __dirname = path.dirname(fileURLToPath(import.meta.url));
|
||||
|
||||
const PROGRESS_QUERY = 'getProgress';
|
||||
|
||||
/** Accept only a code shaped like a stable error code, or a known provider-failure category; anything else is treated as absent rather than printed. */
|
||||
function safeFailureCode(value: string | undefined): string | undefined {
|
||||
if (value !== undefined && (/^[A-Z][A-Z0-9_]{0,63}$/u.test(value) || isProviderFailureCategory(value))) {
|
||||
return value;
|
||||
}
|
||||
return undefined;
|
||||
}
|
||||
|
||||
/** Re-derive a printable line from the closed partial-reason projection instead of trusting the workflow-returned view directly, so a malformed reason prints nothing rather than something wrong. */
|
||||
function safePartialReasonMessage(reason: PipelineState['partialReasons'][number]): string | undefined {
|
||||
if (reason.code === 'agentic_sast_reduced') return 'Agentic SAST completed with reduced coverage.';
|
||||
const candidate = {
|
||||
code: reason.code,
|
||||
...(reason.vulnerabilityClass !== undefined && { vulnerabilityClass: reason.vulnerabilityClass }),
|
||||
...(reason.stage !== undefined && { stage: reason.stage }),
|
||||
...(reason.reductionReason !== undefined && { reductionReason: reason.reductionReason }),
|
||||
...(reason.omittedCount !== undefined && { omittedCount: reason.omittedCount }),
|
||||
...(reason.consideredCount !== undefined && { consideredCount: reason.consideredCount }),
|
||||
...(reason.classifiedCount !== undefined && { classifiedCount: reason.classifiedCount }),
|
||||
...(reason.affectedBatchCount !== undefined && { affectedBatchCount: reason.affectedBatchCount }),
|
||||
};
|
||||
if (!isPartialReason(candidate)) return undefined;
|
||||
return projectPartialReasons([candidate])[0]?.message;
|
||||
}
|
||||
|
||||
// The ordinary activity names. This frozen list is one of three that together form the
|
||||
// registered activity set the CLI status reader mirrors: the Capella names in
|
||||
// ai/sast/capella/temporal/activity-types.ts and the reconciliation names in
|
||||
// reconcile-activity-types.ts are the other two. Adding or removing an activity means
|
||||
// updating both this list and the `pentestActivities` object below, or the load-time check
|
||||
// throws.
|
||||
export const PENTEST_ACTIVITY_NAMES = Object.freeze([
|
||||
'runPreReconAgent',
|
||||
'runReconAgent',
|
||||
'runInjectionVulnAgent',
|
||||
'runXssVulnAgent',
|
||||
'runAuthVulnAgent',
|
||||
'runAuthzVulnAgent',
|
||||
'runSsrfVulnAgent',
|
||||
'runInjectionExploitAgent',
|
||||
'runXssExploitAgent',
|
||||
'runAuthExploitAgent',
|
||||
'runAuthzExploitAgent',
|
||||
'runSsrfExploitAgent',
|
||||
'runMiscellaneousExploitAgent',
|
||||
'runReportAgent',
|
||||
'runPreflightValidation',
|
||||
'runAuthenticationValidation',
|
||||
'initDeliverableGit',
|
||||
'syncPlaywrightStealthConfig',
|
||||
'syncCodePathDenyRules',
|
||||
'initializeDurableScanState',
|
||||
'persistMiscellaneousOutcome',
|
||||
'initializeReportProgress',
|
||||
'renumberClassFindings',
|
||||
'assembleReportActivity',
|
||||
'compactReportFindings',
|
||||
'persistCanonicalReportProgress',
|
||||
'finalizeReportOutputs',
|
||||
'persistFinalizedReportProgress',
|
||||
'surfaceReportOutputs',
|
||||
'checkExploitationQueue',
|
||||
'loadResumeState',
|
||||
'restoreGitCheckpoint',
|
||||
'registerResumeAttempt',
|
||||
'recordResumeAttempt',
|
||||
'logPhaseTransition',
|
||||
'logWorkflowComplete',
|
||||
'saveCheckpoint',
|
||||
] as const);
|
||||
|
||||
export const pentestActivities = Object.freeze({
|
||||
runPreReconAgent,
|
||||
runReconAgent,
|
||||
runInjectionVulnAgent,
|
||||
runXssVulnAgent,
|
||||
runAuthVulnAgent,
|
||||
runAuthzVulnAgent,
|
||||
runSsrfVulnAgent,
|
||||
runInjectionExploitAgent,
|
||||
runXssExploitAgent,
|
||||
runAuthExploitAgent,
|
||||
runAuthzExploitAgent,
|
||||
runSsrfExploitAgent,
|
||||
runMiscellaneousExploitAgent,
|
||||
runReportAgent,
|
||||
runPreflightValidation,
|
||||
runAuthenticationValidation,
|
||||
initDeliverableGit,
|
||||
syncPlaywrightStealthConfig,
|
||||
syncCodePathDenyRules,
|
||||
initializeDurableScanState,
|
||||
persistMiscellaneousOutcome,
|
||||
initializeReportProgress,
|
||||
renumberClassFindings,
|
||||
assembleReportActivity,
|
||||
compactReportFindings,
|
||||
persistCanonicalReportProgress,
|
||||
finalizeReportOutputs,
|
||||
persistFinalizedReportProgress,
|
||||
surfaceReportOutputs,
|
||||
checkExploitationQueue,
|
||||
loadResumeState,
|
||||
restoreGitCheckpoint,
|
||||
registerResumeAttempt,
|
||||
recordResumeAttempt,
|
||||
logPhaseTransition,
|
||||
logWorkflowComplete,
|
||||
saveCheckpoint,
|
||||
});
|
||||
|
||||
const registeredPentestNames = Object.keys(pentestActivities).sort();
|
||||
const expectedPentestNames = [...PENTEST_ACTIVITY_NAMES].sort();
|
||||
if (
|
||||
registeredPentestNames.length !== expectedPentestNames.length ||
|
||||
registeredPentestNames.some((name, index) => name !== expectedPentestNames[index])
|
||||
) {
|
||||
throw new Error('Pentest activity registry does not match its frozen ordinary activity contract');
|
||||
}
|
||||
|
||||
export interface ProductionActivityBindings {
|
||||
readonly repositoryPath: string;
|
||||
readonly webUrl: string;
|
||||
readonly workspacesDir: string;
|
||||
}
|
||||
|
||||
/** Compose the frozen ordinary, Capella, and reconciliation activity namespaces. */
|
||||
export function createProductionActivityRegistry(bindings: ProductionActivityBindings): Readonly<object> {
|
||||
const reconciliationActivities = createReconciliationActivityRegistry({
|
||||
repositoryPath: bindings.repositoryPath,
|
||||
deliverablesDir: deliverablesDir(bindings.repositoryPath),
|
||||
workspacesDir: bindings.workspacesDir,
|
||||
webUrl: bindings.webUrl,
|
||||
});
|
||||
return mergeActivityRegistries(pentestActivities, capellaActivities, reconciliationActivities);
|
||||
}
|
||||
|
||||
// === CLI Argument Parsing ===
|
||||
|
||||
interface CliArgs {
|
||||
webUrl: string;
|
||||
repoPath: string;
|
||||
taskQueue: string;
|
||||
workflowId?: string;
|
||||
configPath?: string;
|
||||
outputPath?: string;
|
||||
customerOutputPath?: string;
|
||||
pipelineTestingMode: boolean;
|
||||
resumeFromWorkspace?: string;
|
||||
}
|
||||
@@ -71,8 +255,10 @@ function showUsage(): void {
|
||||
console.log(' node dist/temporal/worker.js <webUrl> <repoPath> --task-queue <name> [options]\n');
|
||||
console.log('Options:');
|
||||
console.log(' --task-queue <name> Task queue name (required)');
|
||||
console.log(' --workflow-id <id> Workflow ID selected by the Shannon CLI');
|
||||
console.log(' --config <path> Configuration file path');
|
||||
console.log(' --workspace <name> Resume from existing workspace');
|
||||
console.log(' --output <path> Stable mounted path for final customer report copies');
|
||||
console.log(' --pipeline-testing Use minimal prompts for fast testing\n');
|
||||
}
|
||||
|
||||
@@ -85,8 +271,9 @@ function parseCliArgs(argv: string[]): CliArgs {
|
||||
let webUrl: string | undefined;
|
||||
let repoPath: string | undefined;
|
||||
let taskQueue: string | undefined;
|
||||
let workflowId: string | undefined;
|
||||
let configPath: string | undefined;
|
||||
let outputPath: string | undefined;
|
||||
let customerOutputPath: string | undefined;
|
||||
let pipelineTestingMode = false;
|
||||
let resumeFromWorkspace: string | undefined;
|
||||
|
||||
@@ -98,6 +285,12 @@ function parseCliArgs(argv: string[]): CliArgs {
|
||||
taskQueue = nextArg;
|
||||
i++;
|
||||
}
|
||||
} else if (arg === '--workflow-id') {
|
||||
const nextArg = argv[i + 1];
|
||||
if (nextArg && !nextArg.startsWith('-')) {
|
||||
workflowId = nextArg;
|
||||
i++;
|
||||
}
|
||||
} else if (arg === '--config') {
|
||||
const nextArg = argv[i + 1];
|
||||
if (nextArg && !nextArg.startsWith('-')) {
|
||||
@@ -107,7 +300,7 @@ function parseCliArgs(argv: string[]): CliArgs {
|
||||
} else if (arg === '--output') {
|
||||
const nextArg = argv[i + 1];
|
||||
if (nextArg && !nextArg.startsWith('-')) {
|
||||
outputPath = nextArg;
|
||||
customerOutputPath = nextArg;
|
||||
i++;
|
||||
}
|
||||
} else if (arg === '--workspace') {
|
||||
@@ -143,9 +336,10 @@ function parseCliArgs(argv: string[]): CliArgs {
|
||||
webUrl,
|
||||
repoPath,
|
||||
taskQueue,
|
||||
...(workflowId && { workflowId }),
|
||||
pipelineTestingMode,
|
||||
...(configPath && { configPath }),
|
||||
...(outputPath && { outputPath }),
|
||||
...(customerOutputPath && { customerOutputPath }),
|
||||
...(resumeFromWorkspace && { resumeFromWorkspace }),
|
||||
};
|
||||
}
|
||||
@@ -158,16 +352,32 @@ interface SessionJson {
|
||||
webUrl: string;
|
||||
originalWorkflowId?: string;
|
||||
resumeAttempts?: Array<{ workflowId: string }>;
|
||||
status?: 'in-progress' | 'completed' | 'failed' | 'cancelled' | 'partial';
|
||||
};
|
||||
metrics: {
|
||||
total_cost_usd: number;
|
||||
};
|
||||
durableScanState?: {
|
||||
schema_version?: unknown;
|
||||
exploit?: unknown;
|
||||
};
|
||||
}
|
||||
|
||||
function isValidWorkspaceName(name: string): boolean {
|
||||
return /^[a-zA-Z0-9][a-zA-Z0-9_-]{0,127}$/.test(name);
|
||||
}
|
||||
|
||||
function escapeRegExp(value: string): string {
|
||||
return value.replace(/[.*+?^${}()|[\]\\]/g, '\\$&');
|
||||
}
|
||||
|
||||
/** Accept a CLI-owned ID only when it preserves this launch branch's public naming contract. */
|
||||
function selectWorkflowId(requested: string | undefined, fallback: string, expected: RegExp): string {
|
||||
if (requested === undefined) return fallback;
|
||||
if (!expected.test(requested)) throw new Error('Invalid workflow identity supplied by the Shannon CLI');
|
||||
return requested;
|
||||
}
|
||||
|
||||
interface WorkspaceResolution {
|
||||
workflowId: string;
|
||||
sessionId: string;
|
||||
@@ -216,10 +426,15 @@ async function terminateExistingWorkflows(client: Client, workspaceName: string)
|
||||
return terminated;
|
||||
}
|
||||
|
||||
async function resolveWorkspace(client: Client, args: CliArgs): Promise<WorkspaceResolution> {
|
||||
async function resolveWorkspace(client: Client, args: CliArgs, expectedExploit: boolean): Promise<WorkspaceResolution> {
|
||||
if (!args.resumeFromWorkspace) {
|
||||
const hostname = sanitizeHostname(args.webUrl);
|
||||
const workflowId = `${hostname}_shannon-${Date.now()}`;
|
||||
const fallback = `${hostname}_shannon-${Date.now()}`;
|
||||
const workflowId = selectWorkflowId(
|
||||
args.workflowId,
|
||||
fallback,
|
||||
new RegExp(`^${escapeRegExp(hostname)}_shannon-\\d+$`),
|
||||
);
|
||||
return {
|
||||
workflowId,
|
||||
sessionId: workflowId,
|
||||
@@ -233,6 +448,19 @@ async function resolveWorkspace(client: Client, args: CliArgs): Promise<Workspac
|
||||
const workspaceExists = await fileExists(sessionPath);
|
||||
|
||||
if (workspaceExists) {
|
||||
const session = await readJson<SessionJson>(sessionPath);
|
||||
if (session.session.webUrl !== args.webUrl) {
|
||||
throw new Error(
|
||||
'This workspace was created for a different target URL, so it cannot be resumed against this one. Check -u, or start a new scan with a different -w name.',
|
||||
);
|
||||
}
|
||||
if (session.durableScanState?.schema_version !== 1 || typeof session.durableScanState.exploit !== 'boolean') {
|
||||
throw new Error(SAFE_RUN_STATE_MESSAGES.CorruptedSessionError);
|
||||
}
|
||||
if (session.durableScanState.exploit !== expectedExploit) {
|
||||
throw new Error(workspaceExploitMismatchMessage(session.durableScanState.exploit));
|
||||
}
|
||||
|
||||
console.log('=== RESUME MODE ===');
|
||||
console.log(`Workspace: ${workspace}\n`);
|
||||
|
||||
@@ -241,16 +469,9 @@ async function resolveWorkspace(client: Client, args: CliArgs): Promise<Workspac
|
||||
console.log(`Terminated ${terminatedWorkflows.length} previous scan(s)\n`);
|
||||
}
|
||||
|
||||
const session = await readJson<SessionJson>(sessionPath);
|
||||
if (session.session.webUrl !== args.webUrl) {
|
||||
console.error('ERROR: URL mismatch with workspace');
|
||||
console.error(` Workspace URL: ${session.session.webUrl}`);
|
||||
console.error(` Provided URL: ${args.webUrl}`);
|
||||
process.exit(1);
|
||||
}
|
||||
|
||||
const fallback = `${workspace}_resume_${Date.now()}`;
|
||||
return {
|
||||
workflowId: `${workspace}_resume_${Date.now()}`,
|
||||
workflowId: selectWorkflowId(args.workflowId, fallback, new RegExp(`^${escapeRegExp(workspace)}_resume_\\d+$`)),
|
||||
sessionId: workspace,
|
||||
isResume: true,
|
||||
terminatedWorkflows,
|
||||
@@ -258,7 +479,7 @@ async function resolveWorkspace(client: Client, args: CliArgs): Promise<Workspac
|
||||
}
|
||||
|
||||
if (!isValidWorkspaceName(workspace)) {
|
||||
console.error(`ERROR: Invalid workspace name: "${workspace}"`);
|
||||
console.error('ERROR: Invalid workspace name.');
|
||||
console.error(' Must be 1-128 characters, alphanumeric/hyphens/underscores, starting with alphanumeric');
|
||||
process.exit(1);
|
||||
}
|
||||
@@ -268,7 +489,12 @@ async function resolveWorkspace(client: Client, args: CliArgs): Promise<Workspac
|
||||
|
||||
// If the workspace name already looks like a CLI-generated ID
|
||||
// (ends with _shannon-<digits>), use it directly to avoid double _shannon- suffixes
|
||||
const workflowId = /_shannon-\d+$/.test(workspace) ? workspace : `${workspace}_shannon-${Date.now()}`;
|
||||
const fallback = /_shannon-\d+$/.test(workspace) ? workspace : `${workspace}_shannon-${Date.now()}`;
|
||||
const expected =
|
||||
fallback === workspace
|
||||
? new RegExp(`^${escapeRegExp(workspace)}$`)
|
||||
: new RegExp(`^${escapeRegExp(workspace)}_shannon-\\d+$`);
|
||||
const workflowId = selectWorkflowId(args.workflowId, fallback, expected);
|
||||
|
||||
return {
|
||||
workflowId,
|
||||
@@ -281,7 +507,7 @@ async function resolveWorkspace(client: Client, args: CliArgs): Promise<Workspac
|
||||
// === Pipeline Input Construction ===
|
||||
|
||||
interface OrchestrationConfig {
|
||||
vulnClasses?: VulnClass[];
|
||||
agenticSast?: AgenticSastInput;
|
||||
exploit?: boolean;
|
||||
}
|
||||
|
||||
@@ -289,16 +515,26 @@ async function loadOrchestrationConfig(configPath: string | undefined): Promise<
|
||||
if (!configPath) return {};
|
||||
try {
|
||||
const config = await parseConfig(configPath);
|
||||
const distributed = distributeConfig(config);
|
||||
const codePathAvoids = distributed.avoid.filter((rule) => rule.type === 'code_path').map((rule) => rule.value);
|
||||
const codePathFocus = distributed.focus.filter((rule) => rule.type === 'code_path').map((rule) => rule.value);
|
||||
|
||||
return {
|
||||
...(config.vuln_classes && config.vuln_classes.length > 0 && { vulnClasses: [...config.vuln_classes] }),
|
||||
...(config.exploit !== undefined && { exploit: config.exploit === 'true' }),
|
||||
...(distributed.agenticSast && {
|
||||
agenticSast: {
|
||||
codePathAvoids,
|
||||
codePathFocus,
|
||||
modelSpec: process.env.SHANNON_AI_MODEL?.trim() || DEFAULT_MODEL_SPEC,
|
||||
capellaFormatVersion: CAPELLA_FORMAT_VERSION,
|
||||
promptSetVersion: CAPELLA_PROMPT_SET_VERSION,
|
||||
},
|
||||
}),
|
||||
exploit: distributed.exploit,
|
||||
};
|
||||
} catch (error) {
|
||||
// A broken config must fail the run, not silently fall back to empty
|
||||
// defaults that quietly change scope (vuln classes, exploit, retries).
|
||||
const message = error instanceof Error ? error.message : String(error);
|
||||
console.error(`Failed to parse config ${configPath}: ${message}`);
|
||||
console.error('Worker configuration could not be loaded. Reference code: CONFIG_VALIDATION_FAILED');
|
||||
process.exit(1);
|
||||
}
|
||||
}
|
||||
@@ -317,7 +553,8 @@ function buildPipelineInput(
|
||||
...(args.pipelineTestingMode && { pipelineTestingMode: args.pipelineTestingMode }),
|
||||
...(workspace.isResume && args.resumeFromWorkspace && { resumeFromWorkspace: args.resumeFromWorkspace }),
|
||||
...(workspace.terminatedWorkflows.length > 0 && { terminatedWorkflows: workspace.terminatedWorkflows }),
|
||||
...(orchestration.vulnClasses && { vulnClasses: orchestration.vulnClasses }),
|
||||
...(args.customerOutputPath !== undefined && { customerOutputPath: args.customerOutputPath }),
|
||||
...(orchestration.agenticSast !== undefined && { agenticSast: orchestration.agenticSast }),
|
||||
...(orchestration.exploit !== undefined && { exploit: orchestration.exploit }),
|
||||
};
|
||||
}
|
||||
@@ -332,8 +569,11 @@ async function waitForWorkflowResult(
|
||||
try {
|
||||
const progress = await handle.query<PipelineProgress>(PROGRESS_QUERY);
|
||||
const elapsed = Math.floor(progress.elapsedMs / 1000);
|
||||
const expectedCount = progress.expectedAgents.length;
|
||||
// Agentic SAST runs alongside the phase above, so the line names it while it is working.
|
||||
const agenticSast = progress.agenticSast.status === 'running' ? ' | Agentic SAST: running' : '';
|
||||
console.log(
|
||||
`[${elapsed}s] Phase: ${progress.currentPhase || 'unknown'} | Agent: ${progress.currentAgent || 'none'} | Completed: ${progress.completedAgents.length}/13`,
|
||||
`[${elapsed}s] Phase: ${progress.currentPhase || 'unknown'} | Agent: ${progress.currentAgent || 'none'} | Completed: ${progress.completedAgents.length + progress.skippedAgents.length}/${expectedCount}${agenticSast}`,
|
||||
);
|
||||
} catch {
|
||||
// Workflow may have completed
|
||||
@@ -344,12 +584,56 @@ async function waitForWorkflowResult(
|
||||
const result = await handle.result();
|
||||
clearInterval(progressInterval);
|
||||
|
||||
console.log('\nPipeline completed successfully!');
|
||||
// The returned workflow state distinguishes completed, partial, and cancelled runs;
|
||||
// each prints its own terminal line so degradation is never labelled as full success.
|
||||
if (result.status === 'partial') {
|
||||
console.log('\nScan completed with gaps (partial). The reasons are listed below.');
|
||||
for (const reason of result.partialReasons) {
|
||||
const message = safePartialReasonMessage(reason);
|
||||
if (message !== undefined) console.log(` - ${message}`);
|
||||
}
|
||||
// The reason above says a class of coverage degraded; these three name the sanitized
|
||||
// agentic-SAST failure behind it, under the same labels every other surface uses.
|
||||
if (result.agenticSast.status === 'failed') {
|
||||
const stage = ACCEPTED_CAPELLA_FAILURE_STAGES.includes(result.agenticSast.failedStage)
|
||||
? capellaTerminalStageLabel(result.agenticSast.failedStage)
|
||||
: 'orchestration';
|
||||
const message = isCapellaSafeFailureMessage(result.agenticSast.error)
|
||||
? result.agenticSast.error
|
||||
: 'An agentic SAST step failed.';
|
||||
console.log(` Agentic SAST stopped at: ${stage}`);
|
||||
console.log(` What happened: ${message}`);
|
||||
const code = safeFailureCode(result.agenticSast.errorCode);
|
||||
if (code !== undefined) {
|
||||
console.log(` Reference code (for a bug report): ${code}`);
|
||||
}
|
||||
}
|
||||
} else if (result.status === 'cancelled') {
|
||||
console.log('\nScan cancelled before it finished.');
|
||||
} else {
|
||||
console.log('\nScan completed.');
|
||||
}
|
||||
if (result.summary) {
|
||||
console.log(`Duration: ${Math.floor(result.summary.totalDurationMs / 1000)}s`);
|
||||
console.log(`Agents completed: ${result.summary.agentCount}`);
|
||||
console.log(`Agents resolved: ${result.summary.agentCount}`);
|
||||
// Agentic SAST is not an agent, so it is absent from the count above; name it so its spend in
|
||||
// Run cost is accounted for. The failure detail, if any, already printed above.
|
||||
if (result.agenticSast.status === 'succeeded') {
|
||||
const sastGroup = summarizeOperationalMetrics(result.operationalMetrics).find(
|
||||
(group) => group.key === 'agentic-sast',
|
||||
);
|
||||
const cost = sastGroup === undefined || sastGroup.costUsd === null ? 'N/A' : `$${sastGroup.costUsd.toFixed(4)}`;
|
||||
const duration = sastGroup === undefined ? '0s' : `${Math.floor(sastGroup.durationMs / 1000)}s`;
|
||||
const coverage = result.agenticSast.coverage === 'reduced' ? ' — reduced coverage' : '';
|
||||
console.log(`Agentic SAST: completed (${duration}, ${cost})${coverage}`);
|
||||
} else if (result.agenticSast.status === 'failed') {
|
||||
console.log('Agentic SAST: failed');
|
||||
}
|
||||
console.log(`Total turns: ${result.summary.totalTurns}`);
|
||||
console.log(`Run cost: $${result.summary.totalCostUsd.toFixed(4)}`);
|
||||
if (result.summary.usageAccountingComplete === false) {
|
||||
console.log('Cost is incomplete — some background work is not included in this total.');
|
||||
}
|
||||
|
||||
if (workspace.isResume) {
|
||||
try {
|
||||
@@ -362,46 +646,13 @@ async function waitForWorkflowResult(
|
||||
}
|
||||
}
|
||||
}
|
||||
} catch (error) {
|
||||
} catch {
|
||||
clearInterval(progressInterval);
|
||||
console.error('\nPipeline failed:', error);
|
||||
console.error('\nScan failed. Reference code: WORKFLOW_FAILED');
|
||||
process.exit(1);
|
||||
}
|
||||
}
|
||||
|
||||
// === Deliverables Copy ===
|
||||
|
||||
function copyDeliverables(repoPath: string, outputPath: string): void {
|
||||
const outputDir = deliverablesDir(repoPath);
|
||||
if (!fs.existsSync(outputDir)) {
|
||||
console.log('No deliverables directory found, skipping copy');
|
||||
return;
|
||||
}
|
||||
|
||||
const files = fs.readdirSync(outputDir);
|
||||
if (files.length === 0) {
|
||||
console.log('No deliverables to copy');
|
||||
return;
|
||||
}
|
||||
|
||||
fs.mkdirSync(outputPath, { recursive: true });
|
||||
|
||||
for (const file of files) {
|
||||
if (file === '.git') continue;
|
||||
const src = path.join(outputDir, file);
|
||||
const dest = path.join(outputPath, file);
|
||||
fs.cpSync(src, dest, { recursive: true });
|
||||
}
|
||||
|
||||
// Surface the report under its human-facing name alongside the raw deliverables
|
||||
const assembledPdf = path.join(outputDir, ASSEMBLED_REPORT_PDF_FILENAME);
|
||||
if (fs.existsSync(assembledPdf)) {
|
||||
fs.copyFileSync(assembledPdf, path.join(outputPath, FINAL_REPORT_PDF_FILENAME));
|
||||
}
|
||||
|
||||
console.log(`Copied ${files.length} deliverable(s) to ${outputPath}`);
|
||||
}
|
||||
|
||||
// === Main Entry Point ===
|
||||
|
||||
async function run(): Promise<void> {
|
||||
@@ -417,30 +668,40 @@ async function run(): Promise<void> {
|
||||
const client = new Client({ connection: clientConnection });
|
||||
|
||||
try {
|
||||
// 3. Bundle workflows and create worker on per-invocation task queue
|
||||
// 3. Validate orchestration and resume state before terminating any workflow.
|
||||
const orchestration = await loadOrchestrationConfig(args.configPath);
|
||||
const workspace = await resolveWorkspace(client, args, orchestration.exploit ?? true);
|
||||
|
||||
// 4. Bundle workflows and create the worker with the collision-checked activity registry.
|
||||
console.log('Preparing scan...');
|
||||
const workflowBundle = await bundleWorkflowCode({
|
||||
workflowsPath: path.join(__dirname, 'workflows.js'),
|
||||
});
|
||||
|
||||
const productionActivities = createProductionActivityRegistry({
|
||||
repositoryPath: args.repoPath,
|
||||
webUrl: args.webUrl,
|
||||
workspacesDir: path.resolve('./workspaces'),
|
||||
});
|
||||
// args.taskQueue is generated fresh per scan (see resolveWorkspace), so Temporal can only
|
||||
// ever route this worker's activities to this scan's own container: an activity task from
|
||||
// an older or unrelated scan can never execute against the repo mounted here.
|
||||
const worker = await Worker.create({
|
||||
connection,
|
||||
namespace: 'default',
|
||||
workflowBundle,
|
||||
activities,
|
||||
activities: productionActivities,
|
||||
taskQueue: args.taskQueue,
|
||||
maxConcurrentActivityTaskExecutions: 25,
|
||||
});
|
||||
|
||||
// 4. Resolve workspace and build pipeline input
|
||||
const workspace = await resolveWorkspace(client, args);
|
||||
const orchestration = await loadOrchestrationConfig(args.configPath);
|
||||
// 5. Build the fixed-scope pipeline input.
|
||||
const input = buildPipelineInput(args, workspace, orchestration);
|
||||
|
||||
// 5. Start worker polling in the background
|
||||
// 6. Start worker polling in the background.
|
||||
const workerDone = worker.run();
|
||||
|
||||
// 6. Submit workflow to the same task queue
|
||||
// 7. Submit workflow to the same task queue.
|
||||
const handle = await client.workflow.start<(input: PipelineInput) => Promise<PipelineState>>(
|
||||
'pentestPipelineWorkflow',
|
||||
{
|
||||
@@ -450,15 +711,10 @@ async function run(): Promise<void> {
|
||||
},
|
||||
);
|
||||
|
||||
// 7. Wait for workflow result
|
||||
// 8. Wait for workflow result.
|
||||
await waitForWorkflowResult(handle, workspace);
|
||||
|
||||
// 8. Copy deliverables to output directory
|
||||
if (args.outputPath) {
|
||||
copyDeliverables(args.repoPath, args.outputPath);
|
||||
}
|
||||
|
||||
// 9. Shut down worker gracefully
|
||||
// 9. Shut down worker gracefully. Final customer copies are workflow-owned.
|
||||
worker.shutdown();
|
||||
await workerDone;
|
||||
} finally {
|
||||
@@ -467,7 +723,10 @@ async function run(): Promise<void> {
|
||||
}
|
||||
}
|
||||
|
||||
run().catch((err) => {
|
||||
console.error('Worker failed:', err);
|
||||
process.exit(1);
|
||||
});
|
||||
const invokedPath = process.argv[1] ? path.resolve(process.argv[1]) : undefined;
|
||||
if (invokedPath === fileURLToPath(import.meta.url)) {
|
||||
run().catch(() => {
|
||||
console.error('Worker failed. Reference code: WORKER_FAILED');
|
||||
process.exit(1);
|
||||
});
|
||||
}
|
||||
@@ -1,4 +1,4 @@
|
||||
// Copyright (C) 2025 Keygraph, Inc.
|
||||
// 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
|
||||
@@ -9,6 +9,8 @@
|
||||
* Pure functions with no side effects — safe for Temporal workflow sandbox.
|
||||
*/
|
||||
|
||||
import { WORKFLOW_PHASES } from '../audit/safe-fields.js';
|
||||
import { ALL_AGENTS } from '../types/agents.js';
|
||||
import { ErrorCode } from '../types/errors.js';
|
||||
|
||||
/**
|
||||
@@ -26,6 +28,12 @@ const ERROR_TYPE_TO_CODE: Record<string, ErrorCode> = {
|
||||
AgentExecutionError: ErrorCode.AGENT_EXECUTION_FAILED,
|
||||
GitError: ErrorCode.GIT_CHECKPOINT_FAILED,
|
||||
InvalidTargetError: ErrorCode.TARGET_UNREACHABLE,
|
||||
AuthLoginFailedError: ErrorCode.AUTH_LOGIN_FAILED,
|
||||
PipelineFailedError: ErrorCode.AGENT_EXECUTION_FAILED,
|
||||
ReportDraftError: ErrorCode.AGENT_EXECUTION_FAILED,
|
||||
ReportSarifRenderError: ErrorCode.OUTPUT_VALIDATION_FAILED,
|
||||
IncompatibleWorkspaceError: ErrorCode.CONFIG_VALIDATION_FAILED,
|
||||
WorkspaceNotFoundError: ErrorCode.CONFIG_NOT_FOUND,
|
||||
};
|
||||
|
||||
export function classifyErrorCode(error: unknown): ErrorCode | undefined {
|
||||
@@ -40,49 +48,66 @@ export function classifyErrorCode(error: unknown): ErrorCode | undefined {
|
||||
return undefined;
|
||||
}
|
||||
|
||||
/** Maps Temporal error type strings to actionable remediation hints. */
|
||||
/**
|
||||
* Maps Temporal error type strings to actionable remediation hints. A type earns an entry
|
||||
* only when the reader has a next step to take; the rest print without a hint line.
|
||||
*/
|
||||
const REMEDIATION_HINTS: Record<string, string> = {
|
||||
AuthenticationError: "Verify the selected provider's API key is valid and not expired.",
|
||||
ConfigurationError: 'Check your CONFIG file path and contents.',
|
||||
GitError: 'Check repository path and git state.',
|
||||
InvalidTargetError: 'Verify the target URL is correct and accessible.',
|
||||
IncompatibleWorkspaceError: 'start a new scan with a different -w name.',
|
||||
WorkspaceNotFoundError: 'check the -w name against: shannon scans',
|
||||
PipelineFailedError: 're-run the same -w to retry from the last checkpoint.',
|
||||
};
|
||||
|
||||
/**
|
||||
* Walk the .cause chain to find the innermost error with a .type property.
|
||||
* Temporal wraps ApplicationFailure in ActivityFailure — the useful info is inside.
|
||||
* Every message a terminal scan failure can show. Closed on purpose: an ApplicationFailure's
|
||||
* own `.message` can carry raw activity or provider detail, so it is never surfaced directly.
|
||||
* A type absent from this record falls back to one generic sentence instead.
|
||||
*/
|
||||
const SAFE_WORKFLOW_FAILURE_MESSAGES: Readonly<Record<string, string>> = {
|
||||
AuthenticationError: 'Provider authentication failed.',
|
||||
ConfigurationError: 'The scan configuration is invalid.',
|
||||
OutputValidationError: 'A scan step returned an unusable result.',
|
||||
AgentExecutionError: 'An agent could not complete its work.',
|
||||
GitError: 'The scan checkpoint could not be updated.',
|
||||
InvalidTargetError: 'The target could not be reached.',
|
||||
AuthLoginFailedError: 'The configured login could not be completed.',
|
||||
PipelineFailedError: 'The vulnerability analysis phase could not be completed.',
|
||||
ReportDraftError: 'The report could not be saved.',
|
||||
ReportSarifRenderError: 'The report SARIF output could not be rendered.',
|
||||
IncompatibleWorkspaceError: 'This workspace cannot be resumed.',
|
||||
WorkspaceNotFoundError: 'The requested workspace was not found.',
|
||||
};
|
||||
|
||||
const WORKFLOW_PHASE_SET = new Set<string>(WORKFLOW_PHASES);
|
||||
const AGENT_NAME_SET = new Set<string>(ALL_AGENTS);
|
||||
|
||||
/**
|
||||
* Walk the .cause chain to find the innermost approved failure type.
|
||||
* Temporal wraps ApplicationFailure in ActivityFailure, so classification must inspect causes.
|
||||
*
|
||||
* Uses duck-typing because workflow code cannot import @temporalio/activity types.
|
||||
*/
|
||||
function unwrapActivityError(error: unknown): {
|
||||
message: string;
|
||||
type: string | null;
|
||||
} {
|
||||
function unwrapActivityError(error: unknown): { type: string | null } {
|
||||
let current: unknown = error;
|
||||
let typed: { message: string; type: string } | null = null;
|
||||
let type: string | null = null;
|
||||
|
||||
while (current instanceof Error) {
|
||||
if ('type' in current && typeof (current as { type: unknown }).type === 'string') {
|
||||
typed = {
|
||||
message: current.message,
|
||||
type: (current as { type: string }).type,
|
||||
};
|
||||
const candidate = (current as { type: string }).type;
|
||||
if (candidate in SAFE_WORKFLOW_FAILURE_MESSAGES) type = candidate;
|
||||
}
|
||||
current = (current as { cause?: unknown }).cause;
|
||||
}
|
||||
|
||||
if (typed) {
|
||||
return typed;
|
||||
}
|
||||
|
||||
return {
|
||||
message: error instanceof Error ? error.message : String(error),
|
||||
type: null,
|
||||
};
|
||||
return { type };
|
||||
}
|
||||
|
||||
/**
|
||||
* Format a structured error string from workflow catch context.
|
||||
* Format a structured, closed-field error string from workflow catch context.
|
||||
* Segments are delimited by | for multi-line rendering by WorkflowLogger.
|
||||
*/
|
||||
export function formatWorkflowError(error: unknown, currentPhase: string | null, currentAgent: string | null): string {
|
||||
@@ -90,10 +115,12 @@ export function formatWorkflowError(error: unknown, currentPhase: string | null,
|
||||
|
||||
// Phase context (first segment)
|
||||
let phaseContext = 'Pipeline failed';
|
||||
if (currentPhase && currentAgent && currentPhase !== currentAgent) {
|
||||
phaseContext = `${currentPhase} failed (agent: ${currentAgent})`;
|
||||
} else if (currentPhase) {
|
||||
phaseContext = `${currentPhase} failed`;
|
||||
const safePhase = currentPhase !== null && WORKFLOW_PHASE_SET.has(currentPhase) ? currentPhase : null;
|
||||
const safeAgent = currentAgent !== null && AGENT_NAME_SET.has(currentAgent) ? currentAgent : null;
|
||||
if (safePhase && safeAgent && safePhase !== safeAgent) {
|
||||
phaseContext = `${safePhase} failed (agent: ${safeAgent})`;
|
||||
} else if (safePhase) {
|
||||
phaseContext = `${safePhase} failed`;
|
||||
}
|
||||
|
||||
const segments: string[] = [phaseContext];
|
||||
@@ -102,8 +129,11 @@ export function formatWorkflowError(error: unknown, currentPhase: string | null,
|
||||
segments.push(unwrapped.type);
|
||||
}
|
||||
|
||||
// Sanitize pipe characters from message to preserve delimiter format
|
||||
segments.push(unwrapped.message.replaceAll('|', '/'));
|
||||
segments.push(
|
||||
unwrapped.type === null
|
||||
? 'The scan could not be completed.'
|
||||
: (SAFE_WORKFLOW_FAILURE_MESSAGES[unwrapped.type] ?? 'The scan could not be completed.'),
|
||||
);
|
||||
|
||||
if (unwrapped.type) {
|
||||
const hint = REMEDIATION_HINTS[unwrapped.type];
|
||||
|
||||
+1172
-461
File diff suppressed because it is too large.
Load diff
Reference in new issue
Block a user