Skip to content
Merged
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
6 changes: 6 additions & 0 deletions src/adapters/codexResponses.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -664,6 +664,12 @@ describe('401 refresh scope', () => {
// refreshAndRetry expires the token through the store rather than writing
// a snapshot back, so the fake has to offer the same door. (INT-2961)
expireProfile: (_k: string) => { profile = { ...profile, expires: 0 }; return true; },
// The refresh path now runs under the store lock: it re-reads the profile
// from disk under that lock and persists without re-entering it
// (AuthProfileStore.saveUnlocked/setProfileUnlocked/reloadProfileFromDisk).
// The fake is in-memory, so both simply read/write the one profile.
reloadProfileFromDisk: (_k: string) => profile,
setProfileUnlocked: (_k: string, p: typeof profile) => { profile = p; },
};
}

Expand Down
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
128 changes: 85 additions & 43 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,13 +67,26 @@ 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);
const pending = escalationTimers.get(info.pid);
if (pending) {
clearTimeout(pending);
escalationTimers.delete(info.pid);
}

broadcastEvent({
type: 'process:exit',
Expand Down Expand Up @@ -105,61 +118,90 @@ export function getAllProcesses(): ProcessInfo[] {
return Array.from(registry.values());
}

const KILL_ESCALATE_MS = 5_000;

/** Pending SIGKILL escalations per pid, so a force kill can cancel one. */
const escalationTimers = new Map<number, NodeJS.Timeout>();

/**
* 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;
// A soft kill's pending escalation must not outlive this call: it would
// signal again 5s later, and stillOurs() alone cannot see that.
if (escalationTimers.get(pid)) {
clearTimeout(escalationTimers.get(pid));
escalationTimers.delete(pid);
}
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(() => {
escalationTimers.delete(pid);
if (stillOurs()) terminateCliProcessTree(proc);
}, KILL_ESCALATE_MS);
// Unref: a pending escalation must not keep the event loop alive.
timer.unref();
escalationTimers.set(pid, timer);
proc.once('close', () => {
clearTimeout(timer);
escalationTimers.delete(pid);
});
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 +230,4 @@ export function stopHealthChecker(): void {
clearInterval(healthCheckTimer);
healthCheckTimer = null;
}
}
}
94 changes: 94 additions & 0 deletions src/agents/pairPipeline.cancelIsolation.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,94 @@
// Regression: the per-run abort/stuck controls live in AsyncLocalStorage, so a
// run() call can no longer observe or overwrite another run's cancellation
// state. Before the fix both were instance fields, and the second run's already
// aborted signal made the FIRST run cancel itself at its next stage boundary.
import { describe, it, expect, vi, beforeEach, afterEach } from 'vitest';
import type { TaskItem } from '../orchestration/decisionEngine.js';

const runWorker = vi.fn();
const broadcastEvent = vi.fn();
const getDefaultModel = vi.fn();
const hasRepoSnapshot = vi.fn();
const scanAndCache = vi.fn();
const analyzeIssue = vi.fn();
const recallRepoKnowledge = vi.fn();

vi.mock('./worker.js', async () => {
const actual = await vi.importActual<typeof import('./worker.js')>('./worker.js');
return { ...actual, runWorker };
});
vi.mock('./tester.js', async () => {
const actual = await vi.importActual<typeof import('./tester.js')>('./tester.js');
return { ...actual, runTester: vi.fn() };
});
vi.mock('../knowledge/index.js', () => ({ hasRepoSnapshot, scanAndCache, analyzeIssue, recallRepoKnowledge }));
vi.mock('../core/eventHub.js', () => ({ broadcastEvent }));
vi.mock('../adapters/index.js', async () => {
const actual = await vi.importActual<typeof import('../adapters/index.js')>('../adapters/index.js');
return { ...actual, getAdapter: () => ({ getDefaultModel }) };
});

describe('PairPipeline concurrent run isolation', () => {
beforeEach(() => {
vi.spyOn(console, 'log').mockImplementation(() => undefined);
vi.spyOn(console, 'warn').mockImplementation(() => undefined);
vi.spyOn(console, 'error').mockImplementation(() => undefined);
hasRepoSnapshot.mockReturnValue(true);
scanAndCache.mockResolvedValue(undefined);
analyzeIssue.mockResolvedValue(undefined);
recallRepoKnowledge.mockResolvedValue([]);
getDefaultModel.mockResolvedValue('model');
runWorker.mockResolvedValue({
success: true,
summary: 'done',
filesChanged: ['src/a.ts'],
commands: [],
output: '',
confidencePercent: 95,
});
});
afterEach(() => { vi.restoreAllMocks(); vi.clearAllMocks(); });

function task(id: string): TaskItem {
return { id, source: 'linear', title: id, description: 'd', priority: 1, createdAt: Date.now(), estimatedMinutes: 30 };
}

it('does not cancel a run just because a concurrent run on the same instance is aborted', async () => {
const { PairPipeline } = await import('./pairPipeline.js');
const pipeline = new PairPipeline({
stages: ['worker'],
maxIterations: 2,
roles: { worker: { enabled: true, timeoutMs: 0 } },
});

let releaseFirstCall: () => void = () => {};
const firstCallGate = new Promise<void>((resolve) => { releaseFirstCall = resolve; });
let workerCalls = 0;
runWorker.mockImplementation(async () => {
workerCalls++;
if (workerCalls === 1) await firstCallGate;
return {
success: true,
summary: 'done',
filesChanged: ['src/a.ts'],
commands: [],
output: '',
confidencePercent: 95,
};
});

const first = pipeline.run(task('TASK-A'), process.cwd());

// A second run on the SAME instance that is already cancelled. It must not
// install its aborted signal anywhere the first run can read it.
const controller = new AbortController();
controller.abort();
const second = await pipeline.run(task('TASK-B'), process.cwd(), { signal: controller.signal });
expect(second.finalStatus).toBe('cancelled');

releaseFirstCall();
const firstResult = await first;
expect(firstResult.finalStatus).not.toBe('cancelled');
expect(firstResult.success).toBe(true);
});
});
Loading
Loading