Files
shannon/apps/worker/src/services/git-manager.ts
T
ezl-keygraph 5ff40f8c6f feat(worker): migrate agent runtime from Claude Agent SDK to pi harness (#389)
* feat(worker): migrate agent runtime from Claude Agent SDK to pi harness

* feat: remove Google Vertex AI provider support

* fix(worker): route Bedrock and custom-base-URL providers from env

* feat(prompts): instruct agents to call submit_exploitation_queue and submit_auth_result

* fix(worker): count sub-agent cost and surface compaction failures

* refactor(worker): rename claude-executor to pi-executor

* feat(worker): pi-event-driven output formatting

* fix(worker): gate adaptive thinking to Opus models, drop CLAUDE_THINKING_LEVEL

* fix(worker): restore minLength/minItems on vuln-collector schemas

* feat(worker): give task sub-agent write+bash, align tool descriptions

* feat(worker): add glob custom tool and route code_path globs to it

* refactor(prompts): use pi tool names (task, todo_write, read, bash, glob)

* refactor(prompts): drop stale MCP terminology for collector tools

* refactor(prompts): drop collector server names from deliverable instructions

* fix(worker): restore minLength/minItems on pre-recon and exploit collector schemas

* feat(worker): load playwright-cli skill via pi resource loader

* refactor(cli): remove CLAUDE_CODE_MAX_OUTPUT_TOKENS config

* build: drop @anthropic-ai/claude-code from worker image

* docs: remove vertex references from llms context

* docs(worker): update stale sdk comments

* refactor(worker): unify provider precedence between preflight and executor

* feat(worker): enforce bounded bash timeouts via pi extension

* ci: bump the beta release line to 2.0.0 (#356)

* fix(cli): pin npx command hints to beta tag

* fix: render agent deliverables before the success commit so resume preserves them (#377)

* feat(cli): restructure run folder and improve terminal UX (#383)

* feat: surface report at run root and nest run internals under .shannon

* feat: use plain-language wording in user-facing terminal messages

* feat(cli): guide users to watch scan progress and surface report path on start

* docs: sync run-folder layout and CLI wording across docs and comments

* feat(cli): add version command reporting package version or git SHA

* feat(cli): detect TTY for interactive prompts, color, and progress output

* docs: document --yes flag, version command, and tty module

* fix(cli): FORCE_COLOR precedence and plain uninstall --yes output

* fix(cli): respect empty NO_COLOR

* fix(cli): let NO_COLOR take precedence over FORCE_COLOR

* docs: mark claude-code-router integration as removed

* refactor(worker): converge shared core with shannon-oss (#388)

* fix(worker): port keygraph shared-core correctness fixes

* refactor(worker): adopt collectors/ and ai/pi/ layout; add task budget cap and cancellation

* refactor(worker): drop inconsistent Collector "Server" suffix

* refactor(worker): drop unused providerConfig/apiKey seams, resolve credentials from env only

* refactor(worker): port oss code_path pattern expansion + external_directory allow

* fix(worker): preserve dotfile paths in code_path avoid patterns (.env no longer stripped to env)

* feat(worker): render Unprocessed Vulnerabilities section in exploit deliverable (align with oss)

* feat(worker): request set_blind_spots for all vuln classes (align auth/ssrf with production prompts)

* refactor(worker): adopt unified permissionSystem* naming and helper layout

* refactor(worker): inline blind_spots into vuln deliverable section array

* chore(worker): drop unused zod dependency (tree is typebox-native)

* fix(worker): normalize base32 TOTP secret to accept padding and whitespace

* refactor(worker): adopt shared toolResult helper and flatSchema naming in collectors

* refactor(worker): use undefined over null in queue-schema builders

* docs(worker): converge renderer/collector doc comments to current pi terminology

* refactor(worker): adopt schema.ts cleanInput/stringEnum helpers in collectors

* feat(worker): converge exploit-collector/renderer with vendored; capture and render overview for blocked findings

* refactor(worker): converge session-tools/pipeline/exploitation-checker with vendored

* refactor(worker): converge task-tool usage reporting with vendored onUsage callback

* refactor(worker): converge structured output onto a submitTool executor channel

* docs(worker): expand exploit-renderer docstring to match shannon-oss

* docs(worker): adopt richer vuln-renderer docstring from shannon-oss

* docs(worker): neutralize billing-detection wording for shannon-oss parity

* fix(worker): verify checkpoint hash in the deliverables clone being reset

* fix(worker): fail fast on malformed exploitation queue JSON

* fix(worker): honor retryable flag when classifying exploitation-queue check failures

* fix(worker): fail fast on corrupted session.json in run-scope validation

* feat(worker): propagate Temporal cancellation signal into agent and auth pi sessions

* fix(worker): mark exploit agent complete when exploitation is skipped so resume skips it

* prompts: drop scan description from executive report prompt

* refactor(worker): add createGenericSubmitTool for raw JSON-schema submit tools

* refactor(worker): gate playwright-cli skill to browser agents via skillsOverride (adopt shannon-oss mechanism)

* docs(worker): correct formatLogTime comment to UTC to match toISOString

* refactor(worker): converge queue-schemas with shannon-oss (guarded count, decl order)

* refactor(worker): converge task-tool with shannon-oss (byte-identical; modelRegistry optional)

* fix(worker): use replaceLiteral for all prompt value insertions to prevent $-mangling

* fix(worker): classify agent execution failures by error type instead of hardcoding validation

* fix(worker): cap auth-failure detail at 250 chars to match shannon-oss

* style(worker): apply biome formatting

* refactor(worker): remove per-session task delegation cap from task tool

* style(cli): collapse usage hint now that the beta tag is gone

* chore: mark the pi harness migration as a breaking change

BREAKING CHANGE: Google Vertex AI is no longer a supported provider. The
CLAUDE_CODE_USE_VERTEX, ANTHROPIC_VERTEX_PROJECT, CLOUD_ML_REGION, and
GOOGLE_APPLICATION_CREDENTIALS environment variables, along with the
use_vertex, vertex_project, and cloud_ml_region config.toml keys, are
removed. Vertex users must switch to Anthropic, AWS Bedrock, or a custom
Anthropic-compatible base URL.

The CLAUDE_CODE_MAX_OUTPUT_TOKENS environment variable and the
max_output_tokens config.toml key are also removed.
2026-07-16 19:13:13 +05:30

401 lines
12 KiB
TypeScript

// Copyright (C) 2025 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.
import { AsyncLocalStorage } from 'node:async_hooks';
import { $ } from 'zx';
import type { ActivityLogger } from '../types/activity-logger.js';
import { ErrorCode } from '../types/errors.js';
import { PentestError } from './error-handling.js';
/**
* Check if a directory is a git repository.
* Returns true if the directory contains a .git folder or is inside a git repo.
*/
export async function isGitRepository(dir: string): Promise<boolean> {
try {
await $`cd ${dir} && git rev-parse --git-dir`.quiet();
return true;
} catch {
return false;
}
}
interface GitOperationResult {
success: boolean;
hadChanges?: boolean;
changes?: string[];
commitHash?: string;
error?: Error;
}
/**
* Get list of changed files from git status --porcelain -z output.
* When paths is provided, the status query is scoped to those paths.
*/
async function getChangedFiles(
sourceDir: string,
operationDescription: string,
paths?: readonly string[],
): Promise<string[]> {
const args = ['git', 'status', '--porcelain', '-z'];
if (paths && paths.length > 0) {
args.push('--', ...paths);
}
const status = await executeGitCommandWithRetry(args, sourceDir, operationDescription);
return parsePorcelainZ(status.stdout);
}
/**
* Parse `git status --porcelain -z` output.
*
* -z uses NUL separators and raw (unquoted) byte paths, sidestepping the
* fragile whitespace/quote handling of the default porcelain v1 format.
* Each entry is `XY<space>PATH\0`; renames/copies (X = 'R' or 'C') emit an
* additional `ORIG\0` token immediately after the entry, which we skip.
*/
export function parsePorcelainZ(raw: string): string[] {
if (raw.length === 0) {
return [];
}
const tokens = raw.split('\0');
const entries: string[] = [];
for (let i = 0; i < tokens.length; i++) {
const tok = tokens[i];
if (!tok || tok.length < 4) {
continue;
}
entries.push(tok);
const x = tok[0];
if (x === 'R' || x === 'C') {
i++;
}
}
return entries;
}
function changedPathFromStatus(entry: string): string {
return entry.slice(3);
}
async function stageChanges(sourceDir: string, description: string, paths?: readonly string[]): Promise<string[]> {
const changes = await getChangedFiles(sourceDir, description, paths);
if (paths && paths.length > 0) {
const changedPaths = [...new Set(changes.map(changedPathFromStatus).filter((p) => p.length > 0))];
if (changedPaths.length > 0) {
await executeGitCommandWithRetry(['git', 'add', '-A', '--', ...changedPaths], sourceDir, description);
}
return changes;
}
await executeGitCommandWithRetry(['git', 'add', '-A'], sourceDir, description);
return changes;
}
/**
* Log a summary of changed files with truncation for long lists
*/
function logChangeSummary(
changes: string[],
messageWithChanges: string,
messageWithoutChanges: string,
logger: ActivityLogger,
level: 'info' | 'warn' = 'info',
maxToShow: number = 5,
): void {
if (changes.length > 0) {
const msg = messageWithChanges.replace('{count}', String(changes.length));
const fileList = changes
.slice(0, maxToShow)
.map((c) => ` ${c}`)
.join(', ');
const suffix = changes.length > maxToShow ? ` ... and ${changes.length - maxToShow} more files` : '';
logger[level](`${msg} ${fileList}${suffix}`);
} else {
logger[level](messageWithoutChanges);
}
}
/**
* Convert unknown error to GitOperationResult
*/
function toErrorResult(error: unknown): GitOperationResult {
const errMsg = error instanceof Error ? error.message : String(error);
return {
success: false,
error: error instanceof Error ? error : new Error(errMsg),
};
}
// Serializes git operations to prevent index.lock conflicts during parallel agent execution
class GitSemaphore {
private queue: Array<() => void> = [];
private running: boolean = false;
async acquire(): Promise<void> {
return new Promise((resolve) => {
this.queue.push(resolve);
this.process();
});
}
release(): void {
this.running = false;
this.process();
}
private process(): void {
if (!this.running && this.queue.length > 0) {
this.running = true;
const resolve = this.queue.shift();
resolve?.();
}
}
}
const gitSemaphore = new GitSemaphore();
// Tracks whether the current async context already holds the repo lock, so a
// composite operation (e.g. status → add → commit) can call nested git helpers
// without re-acquiring the semaphore and deadlocking on itself.
const gitLockContext = new AsyncLocalStorage<boolean>();
/**
* Run an operation while holding the repo-wide git lock. Reentrant: a nested
* call inside an already-locked context runs immediately instead of blocking.
*/
export async function withGitRepoLock<T>(operation: () => Promise<T>): Promise<T> {
if (gitLockContext.getStore()) {
return operation();
}
await gitSemaphore.acquire();
try {
return await gitLockContext.run(true, operation);
} finally {
gitSemaphore.release();
}
}
const GIT_LOCK_ERROR_PATTERNS = [
'index.lock',
'unable to lock',
'Another git process',
'fatal: Unable to create',
'fatal: index file',
];
function isGitLockError(errorMessage: string): boolean {
return GIT_LOCK_ERROR_PATTERNS.some((pattern) => errorMessage.includes(pattern));
}
// Retries git commands on lock conflicts with exponential backoff
export async function executeGitCommandWithRetry(
commandArgs: string[],
sourceDir: string,
description: string,
maxRetries: number = 5,
): Promise<{ stdout: string; stderr: string }> {
if (!gitLockContext.getStore()) {
return withGitRepoLock(() => executeGitCommandWithRetry(commandArgs, sourceDir, description, maxRetries));
}
for (let attempt = 1; attempt <= maxRetries; attempt++) {
try {
const [cmd, ...args] = commandArgs;
const result = await $`cd ${sourceDir} && ${cmd} ${args}`;
return result;
} catch (error) {
const errMsg = error instanceof Error ? error.message : String(error);
if (isGitLockError(errMsg) && attempt < maxRetries) {
const delay = 2 ** (attempt - 1) * 1000;
// executeGitCommandWithRetry is also called outside activity context
// (e.g., from resume logic), so we use console.warn as a fallback here
console.warn(
`Git lock conflict during ${description} (attempt ${attempt}/${maxRetries}). Retrying in ${delay}ms...`,
);
await new Promise((resolve) => setTimeout(resolve, delay));
continue;
}
throw error;
}
}
throw new PentestError(
`Git command failed after ${maxRetries} retries`,
'filesystem',
true, // Retryable - transient git lock issues
{ maxRetries, description },
ErrorCode.GIT_CHECKPOINT_FAILED,
);
}
// Two-phase reset: hard reset (tracked files) + clean (untracked files).
// When paths is provided, the untracked clean is scoped to those paths so a
// failing agent's rollback can't delete a concurrent sibling agent's scratch.
export async function rollbackGitWorkspace(
sourceDir: string,
reason: string = 'retry preparation',
logger: ActivityLogger,
paths?: readonly string[],
): Promise<GitOperationResult> {
// Skip git operations if not a git repository
if (!(await isGitRepository(sourceDir))) {
logger.info('Skipping git rollback (not a git repository)');
return { success: true };
}
logger.info(`Rolling back workspace for ${reason}`);
try {
const changes = await withGitRepoLock(async () => {
const pendingChanges = await getChangedFiles(sourceDir, 'status check for rollback');
await executeGitCommandWithRetry(['git', 'reset', '--hard', 'HEAD'], sourceDir, 'hard reset for rollback');
const cleanArgs = paths && paths.length > 0 ? ['git', 'clean', '-fd', '--', ...paths] : ['git', 'clean', '-fd'];
await executeGitCommandWithRetry(cleanArgs, sourceDir, 'cleaning untracked files for rollback');
return pendingChanges;
});
logChangeSummary(
changes,
'Rollback completed - removed {count} contaminated changes:',
'Rollback completed - no changes to remove',
logger,
'info',
3,
);
return { success: true };
} catch (error) {
const errMsg = error instanceof Error ? error.message : String(error);
logger.error(`Rollback failed after retries: ${errMsg}`);
return {
success: false,
error: new PentestError(
`Git rollback failed: ${errMsg}`,
'filesystem',
false, // Non-retryable - rollback is best-effort cleanup
{ sourceDir, reason },
ErrorCode.GIT_ROLLBACK_FAILED,
),
};
}
}
// Creates checkpoint before each attempt. First attempt preserves workspace; retries clean it.
export async function createGitCheckpoint(
sourceDir: string,
description: string,
attempt: number,
logger: ActivityLogger,
paths?: readonly string[],
): Promise<GitOperationResult> {
// Skip git operations if not a git repository
if (!(await isGitRepository(sourceDir))) {
logger.info('Skipping git checkpoint (not a git repository)');
return { success: true };
}
logger.info(`Creating checkpoint for ${description} (attempt ${attempt})`);
try {
const result = await withGitRepoLock(async (): Promise<GitOperationResult> => {
// 1. On retries, clean workspace to prevent pollution from previous attempt
if (attempt > 1) {
const cleanResult = await rollbackGitWorkspace(sourceDir, `${description} (retry cleanup)`, logger, paths);
if (!cleanResult.success) {
return cleanResult;
}
}
// 2. Stage scoped changes and commit checkpoint
const changes = await stageChanges(sourceDir, 'staging changes', paths);
const hasChanges = changes.length > 0;
await executeGitCommandWithRetry(
['git', 'commit', '-m', `📍 Checkpoint: ${description} (attempt ${attempt})`, '--allow-empty'],
sourceDir,
'creating commit',
);
const commitHash = await getGitCommitHash(sourceDir);
return { success: true, hadChanges: hasChanges, changes, ...(commitHash && { commitHash }) };
});
if (result.success) {
if (result.hadChanges) {
logger.info('Checkpoint created with scoped changes staged');
} else {
logger.info('Empty checkpoint created (no scoped workspace changes)');
}
}
return result;
} catch (error) {
const result = toErrorResult(error);
logger.warn(`Checkpoint creation failed after retries: ${result.error?.message}`);
return result;
}
}
export async function commitGitSuccess(
sourceDir: string,
description: string,
logger: ActivityLogger,
paths?: readonly string[],
): Promise<GitOperationResult> {
// Skip git operations if not a git repository
if (!(await isGitRepository(sourceDir))) {
logger.info('Skipping git commit (not a git repository)');
return { success: true };
}
logger.info(`Committing successful results for ${description}`);
try {
const result = await withGitRepoLock(async (): Promise<GitOperationResult> => {
const changes = await stageChanges(sourceDir, 'staging changes for success commit', paths);
await executeGitCommandWithRetry(
['git', 'commit', '-m', `✅ ${description}: completed successfully`, '--allow-empty'],
sourceDir,
'creating success commit',
);
const commitHash = await getGitCommitHash(sourceDir);
return {
success: true,
hadChanges: changes.length > 0,
changes,
...(commitHash && { commitHash }),
};
});
logChangeSummary(
result.changes ?? [],
'Success commit created with {count} file changes:',
'Empty success commit created (agent made no file changes)',
logger,
);
return result;
} catch (error) {
const result = toErrorResult(error);
logger.warn(`Success commit failed after retries: ${result.error?.message}`);
return result;
}
}
/**
* Get current git commit hash.
* Returns null if not a git repository.
*/
export async function getGitCommitHash(sourceDir: string): Promise<string | null> {
if (!(await isGitRepository(sourceDir))) {
return null;
}
try {
const result = await executeGitCommandWithRetry(['git', 'rev-parse', 'HEAD'], sourceDir, 'read HEAD commit');
return result.stdout.trim();
} catch {
return null;
}
}