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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 8 additions & 1 deletion src/adapters/processRegistry.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,14 @@ function fakeChild(pid: number): EventEmitter & { pid: number } {
function register(pid: number): EventEmitter & { pid: number } {
const proc = fakeChild(pid);
registerProcess(
{ pid, adapter: 'codex', role: 'worker', command: 'codex', spawnedAt: Date.now() } as never,
{
pid,
taskId: 't',
stage: 'worker',
projectPath: '/tmp',
spawnedAt: Date.now(),
lastActivityAt: Date.now(),
},
proc as never,
);
return proc;
Expand Down
111 changes: 67 additions & 44 deletions src/adapters/processRegistry.ts
Original file line number Diff line number Diff line change
Expand Up @@ -50,13 +50,13 @@ export function registerProcess(info: ProcessInfo, proc: ChildProcess): void {
},
});

// Track activity from stdout/stderr
const updateActivity = () => {
// Activity tracking
const updateActivity = (): void => {
const entry = registry.get(info.pid);
if (!entry) return;
entry.lastActivityAt = Date.now();
// Ownership check: only update if this exact identity still owns the PID
if (!entry || entry.taskId !== info.taskId || entry.spawnedAt !== info.spawnedAt) return;

// Throttled broadcast
entry.lastActivityAt = Date.now();
const lastBroadcast = activityThrottle.get(info.pid) ?? 0;
if (Date.now() - lastBroadcast >= ACTIVITY_THROTTLE_MS) {
activityThrottle.set(info.pid, Date.now());
Expand All @@ -67,10 +67,18 @@ export function registerProcess(info: ProcessInfo, proc: ChildProcess): void {
proc.stdout?.on('data', updateActivity);
proc.stderr?.on('data', updateActivity);

// Cleanup on close
// Cleanup on close — verify ownership before acting so a reused PID does not
// cause this handler to clean up a newer process that happens to share the same
// numeric PID (PID reuse is real on busy systems).
proc.on('close', (code, signal) => {
const entry = registry.get(info.pid);
const durationMs = entry ? Date.now() - entry.spawnedAt : 0;
// Ownership check: the entry must match this exact process identity (taskId +
// spawnedAt), not just the PID. A reused PID would have a different spawnedAt
// or taskId.
if (!entry || entry.taskId !== info.taskId || entry.spawnedAt !== info.spawnedAt) {
return;
}
const durationMs = Date.now() - entry.spawnedAt;
registry.delete(info.pid);
processHandles.delete(info.pid);
activityThrottle.delete(info.pid);
Expand Down Expand Up @@ -105,61 +113,76 @@ export function getAllProcesses(): ProcessInfo[] {
return Array.from(registry.values());
}

const KILL_ESCALATE_MS = 5_000;

/**
* Kill a tracked process. Sends SIGTERM first, then SIGKILL after 5s.
* Kill a tracked process by PID.
* Signals through the retained ChildProcess handle — never by raw PID — so a
* recycled numeric PID cannot be escalated into after the original child exits.
* Soft kill schedules SIGKILL escalation after {@link KILL_ESCALATE_MS} only
* while this exact registry identity still owns the handle; the caller does not
* wait for that timer (so a recycled PID in those seconds is never signalled).
*/
export async function killProcess(pid: number, force = false): Promise<boolean> {
const entry = registry.get(pid);
if (!entry) return false;
const info = registry.get(pid);
const proc = processHandles.get(pid);
if (!info || !proc) return false;

const stillOurs = (): boolean => {
const current = registry.get(pid);
return (
!!current
&& current.taskId === info.taskId
&& current.spawnedAt === info.spawnedAt
&& processHandles.get(pid) === proc
);
};

if (!stillOurs()) return false;

try {
const proc = processHandles.get(pid);
if (proc) {
if (force) {
terminateCliProcessTree(proc);
} else {
signalCliProcessTree(proc, 'SIGTERM');
// Escalate the same process group. The npm wrapper may exit before its
// native CLI/MCP descendants, so checking only proc.exitCode is unsafe.
const escalation = setTimeout(() => terminateCliProcessTree(proc), 5000);
escalation.unref();
}
if (force) {
if (!stillOurs()) return false;
terminateCliProcessTree(proc);
return true;
}
// No handle means we cannot prove this PID is still the child we spawned.
// Signalling it anyway is the ownership hazard this registry exists to
// remove: PIDs are recycled, and a delayed escalation in particular would
// land on whatever unrelated process inherited the number. Registration and
// handle are written and cleared together, so this is a defensive branch —
// report it rather than guessing.
console.warn(`[ProcessRegistry] No process handle for pid ${pid}; refusing to signal by PID`);
registry.delete(pid);
activityThrottle.delete(pid);
return false;

signalCliProcessTree(proc, 'SIGTERM');
const timer = setTimeout(() => {
if (stillOurs()) terminateCliProcessTree(proc);
}, KILL_ESCALATE_MS);
proc.once('close', () => {
clearTimeout(timer);
});
return true;
} catch {
// Process already gone
registry.delete(pid);
processHandles.delete(pid);
activityThrottle.delete(pid);
return false;
}
}

/**
* Start periodic health checker that verifies processes are still alive.
* Removes stale entries from the registry.
* Start periodic health checker that removes stale entries
* where the process handle has already exited.
*
* Identity-safe: re-reads the registry entry before acting so a reused PID
* does not cause this handler to clean up a newer process.
*/
export function startHealthChecker(intervalMs = 30000): void {
if (healthCheckTimer) return;

healthCheckTimer = setInterval(() => {
const now = Date.now();
for (const [pid, info] of registry) {
try {
process.kill(pid, 0); // No-op signal — just checks if alive
} catch {
// Process is dead but wasn't cleaned up
const durationMs = Date.now() - info.spawnedAt;
const proc = processHandles.get(pid);
if (!proc) {
// Ownership check: re-read the entry to confirm it hasn't been replaced
// by a reused PID with a different identity (taskId + spawnedAt).
const current = registry.get(pid);
if (!current || current.taskId !== info.taskId || current.spawnedAt !== info.spawnedAt) {
continue;
}
const durationMs = now - info.spawnedAt;
registry.delete(pid);
processHandles.delete(pid);
activityThrottle.delete(pid);
broadcastEvent({
type: 'process:exit',
Expand Down Expand Up @@ -188,4 +211,4 @@ export function stopHealthChecker(): void {
clearInterval(healthCheckTimer);
healthCheckTimer = null;
}
}
}
35 changes: 21 additions & 14 deletions src/agents/pairPipeline.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
// Worker → Reviewer → Tester → Documenter pipeline
// ============================================
import { EventEmitter } from 'node:events';
import { AsyncLocalStorage } from 'node:async_hooks';
import { taskEventKey, type TaskItem } from '../orchestration/decisionEngine.js';
import { enforcedFileScope } from '../orchestration/writeScope.js';
import type { WorkerResult, ReviewResult } from './agentPair.js';
Expand Down Expand Up @@ -77,6 +78,7 @@ export { buildTaskPrefix } from './pipelineTaskPrefix.js';
export { stageTimeoutMs } from './stageTimeouts.js';
import { stageTimeoutMs } from './stageTimeouts.js';

type RunControl = { signal?: AbortSignal; stuck: StuckDetector };

/**
* Resolve the coordination identity for one stage of a task.
Expand All @@ -89,14 +91,18 @@ import { stageTimeoutMs } from './stageTimeouts.js';
*/
export class PairPipeline extends EventEmitter {
private config: PipelineConfig;
/** Fallback for callers outside run(); each run() installs a fresh detector via AsyncLocalStorage. */
private stuckDetector: StuckDetector;
/** Set per run() — aborts the pipeline + in-flight adapter call on cancel/disable. */
private abortSignal?: AbortSignal;
/** Per-run abort + stuck controls — never share across concurrent run() calls. */
private static runControl = new AsyncLocalStorage<RunControl>();
/** Cache of adapter default models (heavy: OAuth + live catalog) keyed by adapter name. (INT-2393) */
private defaultModelCache = new Map<string, Promise<string | undefined>>();
/** Throw if this run has been cancelled. Called at iteration/stage boundaries. */
private throwIfAborted(): void {
if (this.abortSignal?.aborted) throw new PipelineCancelledError();
if (PairPipeline.runControl.getStore()?.signal?.aborted) throw new PipelineCancelledError();
}
private getStuck(): StuckDetector {
return PairPipeline.runControl.getStore()?.stuck ?? this.stuckDetector;
}

constructor(config: PipelineConfig) {
Expand Down Expand Up @@ -129,13 +135,11 @@ export class PairPipeline extends EventEmitter {
* On failure at any stage, returns to Worker (up to maxIterations)
*/
async run(task: TaskItem, projectPath: string, opts?: { signal?: AbortSignal }): Promise<PipelineResult> {
const stuckDetector = createStuckDetector({ sameErrorRepeat: 3, revisionLoop: 4 });
return PairPipeline.runControl.run({ signal: opts?.signal, stuck: stuckDetector }, async () => {
const startTime = Date.now();
const stages: StageResult[] = [];
const maxIterations = this.config.maxIterations ?? 3;
this.abortSignal = opts?.signal;

// Reset stuck detector (new pipeline run)
this.stuckDetector.reset();

// Ensure repo graph snapshot exists (first-time scan if needed)
if (!hasRepoSnapshot(projectPath)) {
Expand Down Expand Up @@ -168,6 +172,8 @@ export class PairPipeline extends EventEmitter {
currentIteration: 0,
taskPrefix,
reflection: createReflectionState(),
abortSignal: opts?.signal,
stuckDetector,
};
try {
if (this.config.verify?.enabled) try {
Expand Down Expand Up @@ -215,7 +221,7 @@ export class PairPipeline extends EventEmitter {
} catch (error) {
// Cancellation (project disable / manual stop) is not a failure — surface it
// as 'cancelled' so the scheduler doesn't count it failed or trigger a retry.
const cancelled = error instanceof PipelineCancelledError || !!this.abortSignal?.aborted;
const cancelled = error instanceof PipelineCancelledError || !!PairPipeline.runControl.getStore()?.signal?.aborted;
// A 429/usage-limit propagates up here from any stage (worker/reviewer/…).
// Surface it as its own finalStatus so the runner pauses until quota resets
// instead of counting a failure and spamming Linear comments. (INT-1906)
Expand Down Expand Up @@ -260,6 +266,7 @@ export class PairPipeline extends EventEmitter {
},
};
}
});
}
/**
* Worker에 주입할 코드 컨텍스트 수집
Expand All @@ -274,7 +281,7 @@ export class PairPipeline extends EventEmitter {
private async runPostSuccessStage(stage: PipelineStage, context: PipelineContext, stages: StageResult[]): Promise<void> {
try { stages.push(await this.runStage(stage, context)); }
catch (err) {
if (err instanceof PipelineCancelledError || this.abortSignal?.aborted) throw err;
if (err instanceof PipelineCancelledError || PairPipeline.runControl.getStore()?.signal?.aborted) throw err;
safeConsole.warn(`[${context.taskPrefix}] ${stage} skipped (non-blocking failure): ${err instanceof Error ? err.message : String(err)}`);
}
}
Expand Down Expand Up @@ -415,7 +422,7 @@ export class PairPipeline extends EventEmitter {
onLog,
processContext: { taskId: taskEventKey(context.task), stage: 'worker' },
workerContext,
signal: this.abortSignal,
signal: PairPipeline.runControl.getStore()?.signal,
instructionCapsule: this.config.instructionCapsule,
mcpTools: this.config.roleMcpTools?.worker,
adapterRouting: this.config.adapterRouting,
Expand Down Expand Up @@ -512,7 +519,7 @@ export class PairPipeline extends EventEmitter {
type: 'log',
data: { taskId: taskEventKey(context.task), stage: 'reviewer', line: `[${prefix}] ${line}` },
}),
signal: this.abortSignal,
signal: PairPipeline.runControl.getStore()?.signal,
instructionCapsule: this.config.instructionCapsule,
mcpTools: this.config.roleMcpTools?.reviewer,
coordinationContext: coordinationContextFor(context, 'reviewer'),
Expand Down Expand Up @@ -777,7 +784,7 @@ export class PairPipeline extends EventEmitter {
context.currentIteration++;

// Stuck detection check (before iteration starts)
const stuckCheck = this.stuckDetector.check();
const stuckCheck = this.getStuck().check();
if (stuckCheck.isStuck) {
context.stuckReason = stuckCheck.reason;
safeConsole.error(`[${context.taskPrefix}] STUCK DETECTED: ${stuckCheck.reason}`);
Expand Down Expand Up @@ -833,7 +840,7 @@ export class PairPipeline extends EventEmitter {
stages.push(workerResult);

// Record Worker result in stuck detector
this.stuckDetector.addEntry({
this.getStuck().addEntry({
stage: 'worker',
success: workerResult.success,
output: (workerResult.result as WorkerResult).summary,
Expand Down Expand Up @@ -1143,7 +1150,7 @@ export class PairPipeline extends EventEmitter {
const decision = (reviewerResult.result as ReviewResult).decision;

// Record Reviewer result in stuck detector
this.stuckDetector.addEntry({
this.getStuck().addEntry({
stage: 'reviewer',
success: reviewerResult.success,
decision: decision,
Expand Down
3 changes: 3 additions & 0 deletions src/agents/pairPipelineTypes.ts
Original file line number Diff line number Diff line change
Expand Up @@ -179,6 +179,9 @@ export interface PipelineContext {
newSecurityFindings?: import('../verify/securityAudit.js').SecurityFinding[];
/** Pair-level stagnation detector reason, preserved so the scheduler does not rerun the same loop. */
stuckReason?: string;
abortSignal?: AbortSignal;
/** Per-run stuck detector — never share across concurrent run() calls. */
stuckDetector?: import('../support/stuckDetector.js').StuckDetector;
}

export type PipelineEventType = 'stage:start' | 'stage:complete' | 'stage:fail' | 'iteration:start' | 'iteration:complete' | 'iteration:fail' | 'pipeline:complete' | 'pipeline:fail' | 'fanout:gate' | 'halt';
Loading