mirror of
https://github.com/Ark0N/Codeman.git
synced 2026-09-30 12:39:42 +02:00
chore: bump version to 0.1583
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
This commit is contained in:
@@ -35,7 +35,7 @@ When user says "COM":
|
||||
1. Increment version in BOTH `package.json` AND `CLAUDE.md` (verify they match with `grep version package.json && grep Version CLAUDE.md`)
|
||||
2. Run: `git add -A && git commit -m "chore: bump version to X.XXXX" && git push && npm run build && systemctl --user restart claudeman-web`
|
||||
|
||||
**Version**: 0.1582 (must match `package.json` for npm publish)
|
||||
**Version**: 0.1583 (must match `package.json` for npm publish)
|
||||
|
||||
## Project Overview
|
||||
|
||||
|
||||
+1
-1
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "claudeman",
|
||||
"version": "0.1582",
|
||||
"version": "0.1583",
|
||||
"description": "The missing control plane for Claude Code - run 20 autonomous agents with real-time monitoring and session persistence",
|
||||
"type": "module",
|
||||
"main": "dist/index.js",
|
||||
|
||||
+83
-18
@@ -33,6 +33,7 @@ export interface SubagentInfo {
|
||||
modelShort?: 'haiku' | 'sonnet' | 'opus'; // Short model identifier
|
||||
totalInputTokens?: number; // Running total of input tokens
|
||||
totalOutputTokens?: number; // Running total of output tokens
|
||||
pid?: number; // Cached process ID for fast liveness checks
|
||||
}
|
||||
|
||||
export interface SubagentToolCall {
|
||||
@@ -126,6 +127,7 @@ const IDLE_TIMEOUT_MS = 30000; // Consider agent idle after 30s of no activity
|
||||
const POLL_INTERVAL_MS = 1000; // Base poll interval (lightweight checks)
|
||||
const FULL_SCAN_EVERY_N_POLLS = 5; // Full directory traversal every 5th poll (5s)
|
||||
const LIVENESS_CHECK_MS = 10000; // Check if subagent processes are still alive every 10s
|
||||
const FILE_ALIVE_THRESHOLD_MS = 30000; // File mtime within 30s = agent alive (primary check)
|
||||
const STALE_COMPLETED_MAX_AGE_MS = 60 * 60 * 1000; // Remove completed agents older than 1 hour
|
||||
const STALE_IDLE_MAX_AGE_MS = 4 * 60 * 60 * 1000; // Remove idle agents older than 4 hours
|
||||
const STARTUP_MAX_FILE_AGE_MS = 4 * 60 * 60 * 1000; // Only load files modified in last 4 hours on startup
|
||||
@@ -235,7 +237,12 @@ export class SubagentWatcher extends EventEmitter {
|
||||
|
||||
/**
|
||||
* Start periodic liveness checker
|
||||
* Detects when subagent processes have exited but status is still active/idle
|
||||
* Detects when subagent processes have exited but status is still active/idle.
|
||||
*
|
||||
* Uses a 3-tier check to minimize cost:
|
||||
* 1. File mtime (stat ~0.3ms/agent) — if transcript modified recently, agent is alive
|
||||
* 2. Cached PID (/proc/{pid}/stat ~0.1ms) — if stored PID still exists, alive
|
||||
* 3. Full pgrep scan (expensive, ~500ms) — only for agents that fail tiers 1+2
|
||||
*/
|
||||
private startLivenessChecker(): void {
|
||||
if (this.livenessInterval) return;
|
||||
@@ -246,26 +253,40 @@ export class SubagentWatcher extends EventEmitter {
|
||||
this._isCheckingLiveness = true;
|
||||
|
||||
try {
|
||||
// Run pgrep ONCE and collect all claude process info
|
||||
const pidMap = await this.getClaudePids();
|
||||
// Collect agents that need the expensive pgrep scan
|
||||
const needsFullScan: SubagentInfo[] = [];
|
||||
|
||||
for (const [_agentId, info] of this.agentInfo) {
|
||||
// Re-check status in case another check completed this agent
|
||||
if (info.status === 'active' || info.status === 'idle') {
|
||||
if (info.status !== 'active' && info.status !== 'idle') continue;
|
||||
|
||||
// Tier 1: File mtime check (~0.3ms per agent)
|
||||
if (await this.checkSubagentFileAlive(info)) continue;
|
||||
|
||||
// Tier 2: Cached PID check (~0.1ms per agent)
|
||||
if (info.pid && await this.checkPidAlive(info.pid)) continue;
|
||||
|
||||
// Tiers 1+2 failed — need expensive scan for this agent
|
||||
needsFullScan.push(info);
|
||||
}
|
||||
|
||||
// Tier 3: Full pgrep scan — only if any agents failed cheap checks
|
||||
if (needsFullScan.length > 0) {
|
||||
const pidMap = await this.getClaudePids();
|
||||
|
||||
for (const info of needsFullScan) {
|
||||
// Re-check status in case another check completed this agent
|
||||
if (info.status !== 'active' && info.status !== 'idle') continue;
|
||||
|
||||
const alive = this.checkSubagentAliveFromPidMap(info, pidMap);
|
||||
// If process not found, fall back to file mtime check
|
||||
if (!alive) {
|
||||
const fileAlive = await this.checkSubagentFileAlive(info);
|
||||
if (!fileAlive && (info.status === 'active' || info.status === 'idle')) {
|
||||
// Double-check status after async call to prevent race
|
||||
info.status = 'completed';
|
||||
// Clean up pendingToolCalls for this agent to prevent memory leak
|
||||
this.pendingToolCalls.delete(info.agentId);
|
||||
this.emit('subagent:completed', info);
|
||||
}
|
||||
info.pid = undefined;
|
||||
info.status = 'completed';
|
||||
this.pendingToolCalls.delete(info.agentId);
|
||||
this.emit('subagent:completed', info);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Periodically clean up stale completed agents (older than 24 hours)
|
||||
this.cleanupStaleAgents();
|
||||
} finally {
|
||||
@@ -274,10 +295,23 @@ export class SubagentWatcher extends EventEmitter {
|
||||
}, LIVENESS_CHECK_MS);
|
||||
}
|
||||
|
||||
/**
|
||||
* Check if a PID is still alive via /proc/{pid}/stat (single file read, ~0.1ms).
|
||||
*/
|
||||
private async checkPidAlive(pid: number): Promise<boolean> {
|
||||
try {
|
||||
await statAsync(`/proc/${pid}`);
|
||||
return true;
|
||||
} catch {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Run pgrep once and read /proc info for all Claude PIDs in parallel.
|
||||
* Returns a Map of pid -> { environ, cmdline } for subagent processes only.
|
||||
* Excludes main Claudeman-managed Claude processes (CLAUDEMAN_MUX=1).
|
||||
* Also updates cached PIDs on tracked agents when a match is found.
|
||||
*/
|
||||
private async getClaudePids(): Promise<Map<number, { environ: string; cmdline: string }>> {
|
||||
const result = new Map<number, { environ: string; cmdline: string }>();
|
||||
@@ -302,6 +336,17 @@ export class SubagentWatcher extends EventEmitter {
|
||||
result.set(pid, { environ, cmdline });
|
||||
}
|
||||
}));
|
||||
|
||||
// Update cached PIDs on tracked agents
|
||||
for (const [pid, procInfo] of result) {
|
||||
for (const [_agentId, info] of this.agentInfo) {
|
||||
if (info.status !== 'active' && info.status !== 'idle') continue;
|
||||
if (procInfo.environ.includes(info.sessionId) || procInfo.cmdline.includes(info.sessionId)) {
|
||||
info.pid = pid;
|
||||
break; // Each PID belongs to at most one agent
|
||||
}
|
||||
}
|
||||
}
|
||||
} catch {
|
||||
// pgrep returns non-zero if no matches
|
||||
}
|
||||
@@ -324,14 +369,15 @@ export class SubagentWatcher extends EventEmitter {
|
||||
}
|
||||
|
||||
/**
|
||||
* Check if a subagent's transcript file was recently modified (fallback for process check).
|
||||
* Check if a subagent's transcript file was recently modified.
|
||||
* Primary liveness signal — transcript files are written to continuously while agent is active.
|
||||
*/
|
||||
private async checkSubagentFileAlive(info: SubagentInfo): Promise<boolean> {
|
||||
try {
|
||||
const fileStat = await statAsync(info.filePath);
|
||||
const mtime = fileStat.mtime.getTime();
|
||||
const now = Date.now();
|
||||
if (now - mtime < 60000) {
|
||||
if (now - mtime < FILE_ALIVE_THRESHOLD_MS) {
|
||||
return true;
|
||||
}
|
||||
} catch {
|
||||
@@ -597,7 +643,7 @@ export class SubagentWatcher extends EventEmitter {
|
||||
|
||||
/**
|
||||
* Kill a subagent by its agent ID
|
||||
* Finds the Claude process and sends SIGTERM
|
||||
* Uses cached PID first, falls back to findSubagentProcess if needed.
|
||||
*/
|
||||
async killSubagent(agentId: string): Promise<boolean> {
|
||||
const info = this.agentInfo.get(agentId);
|
||||
@@ -607,10 +653,12 @@ export class SubagentWatcher extends EventEmitter {
|
||||
if (info.status === 'completed') return false;
|
||||
|
||||
try {
|
||||
// Find Claude process with matching session ID
|
||||
// Always use findSubagentProcess for kill — it verifies environ/cmdline,
|
||||
// preventing PID reuse attacks (cached PID may have been recycled by OS)
|
||||
const pid = await this.findSubagentProcess(info.sessionId);
|
||||
if (pid) {
|
||||
process.kill(pid, 'SIGTERM');
|
||||
info.pid = undefined;
|
||||
info.status = 'completed';
|
||||
this.emit('subagent:completed', info);
|
||||
return true;
|
||||
@@ -620,6 +668,7 @@ export class SubagentWatcher extends EventEmitter {
|
||||
}
|
||||
|
||||
// Mark as completed even if we couldn't find the process
|
||||
info.pid = undefined;
|
||||
info.status = 'completed';
|
||||
this.emit('subagent:completed', info);
|
||||
return true;
|
||||
@@ -646,6 +695,7 @@ export class SubagentWatcher extends EventEmitter {
|
||||
* Find the process ID of a Claude subagent by its session ID.
|
||||
* Searches /proc for claude processes with matching session ID in environment.
|
||||
* Skips main Claudeman-managed Claude processes (identified by CLAUDEMAN_MUX=1).
|
||||
* Caches the discovered PID on the matching agent info for future fast checks.
|
||||
*/
|
||||
private async findSubagentProcess(sessionId: string): Promise<number | null> {
|
||||
try {
|
||||
@@ -674,12 +724,15 @@ export class SubagentWatcher extends EventEmitter {
|
||||
if (environ.includes('CLAUDEMAN_MUX=1')) continue;
|
||||
|
||||
if (environ.includes(sessionId)) {
|
||||
// Cache PID on the matching agent
|
||||
this.cacheAgentPid(sessionId, pid);
|
||||
return pid;
|
||||
}
|
||||
|
||||
try {
|
||||
const cmdline = await readFile(`/proc/${pid}/cmdline`, 'utf8');
|
||||
if (cmdline.includes(sessionId)) {
|
||||
this.cacheAgentPid(sessionId, pid);
|
||||
return pid;
|
||||
}
|
||||
} catch {
|
||||
@@ -692,6 +745,18 @@ export class SubagentWatcher extends EventEmitter {
|
||||
return null;
|
||||
}
|
||||
|
||||
/**
|
||||
* Store a discovered PID on the agent info with matching sessionId.
|
||||
*/
|
||||
private cacheAgentPid(sessionId: string, pid: number): void {
|
||||
for (const [_agentId, info] of this.agentInfo) {
|
||||
if (info.sessionId === sessionId && (info.status === 'active' || info.status === 'idle')) {
|
||||
info.pid = pid;
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Get transcript for a subagent (optionally limited to last N entries)
|
||||
*/
|
||||
|
||||
+114
-37
@@ -504,55 +504,132 @@ describe('SubagentWatcher performance', () => {
|
||||
});
|
||||
});
|
||||
|
||||
describe('liveness checker: serial pgrep amplification', () => {
|
||||
it('should measure total event loop impact of serial liveness checks', async () => {
|
||||
// The liveness checker iterates all active agents SERIALLY with await:
|
||||
// for (const [agentId, info] of this.agentInfo) {
|
||||
// const alive = await this.checkSubagentAlive(agentId);
|
||||
// }
|
||||
// Each checkSubagentAlive calls findSubagentProcess which runs:
|
||||
// 1. pgrep -f claude
|
||||
// 2. For each PID: readFile(/proc/{pid}/environ) + readFile(/proc/{pid}/cmdline)
|
||||
//
|
||||
// With 20 active agents, the entire for-loop holds the event loop hostage
|
||||
// because Node's microtask queue won't process other work between awaits
|
||||
// in a tight loop.
|
||||
|
||||
describe('liveness checker: tiered optimization', () => {
|
||||
it('should measure old approach: full pgrep scan for 20 agents', async () => {
|
||||
const { execFile } = await import('node:child_process');
|
||||
const { readFile } = await import('node:fs/promises');
|
||||
|
||||
// Simulate the liveness check loop for 20 agents
|
||||
const lagPromise = measureEventLoopLag(3000);
|
||||
|
||||
const start = performance.now();
|
||||
for (let agent = 0; agent < 20; agent++) {
|
||||
// Step 1: pgrep (what findSubagentProcess does)
|
||||
const pids = await new Promise<string[]>((resolve) => {
|
||||
execFile('pgrep', ['-f', 'node'], { encoding: 'utf8' }, (_err, stdout) => {
|
||||
resolve((stdout || '').trim().split('\n').filter(Boolean));
|
||||
});
|
||||
// Old approach: pgrep once + /proc reads for all PIDs, then iterate all agents
|
||||
const pids = await new Promise<string[]>((resolve) => {
|
||||
execFile('pgrep', ['-f', 'node'], { encoding: 'utf8' }, (_err, stdout) => {
|
||||
resolve((stdout || '').trim().split('\n').filter(Boolean));
|
||||
});
|
||||
|
||||
// Step 2: Read /proc for each PID (what findSubagentProcess does per PID)
|
||||
for (const pidStr of pids.slice(0, 5)) { // limit to 5 PIDs for test sanity
|
||||
try {
|
||||
await readFile(`/proc/${pidStr}/environ`, 'utf8');
|
||||
} catch { /* expected for many PIDs */ }
|
||||
try {
|
||||
await readFile(`/proc/${pidStr}/cmdline`, 'utf8');
|
||||
} catch { /* expected */ }
|
||||
}
|
||||
});
|
||||
for (const pidStr of pids.slice(0, 10)) {
|
||||
try { await readFile(`/proc/${pidStr}/environ`, 'utf8'); } catch { /* */ }
|
||||
try { await readFile(`/proc/${pidStr}/cmdline`, 'utf8'); } catch { /* */ }
|
||||
}
|
||||
const elapsed = performance.now() - start;
|
||||
|
||||
const lag = await lagPromise;
|
||||
|
||||
console.log(`[liveness × 20 agents] total: ${elapsed.toFixed(0)}ms`);
|
||||
console.log(`[liveness × 20 agents] max event loop lag: ${lag.maxLagMs.toFixed(1)}ms`);
|
||||
console.log(`[liveness × 20 agents] avg event loop lag: ${lag.avgLagMs.toFixed(1)}ms`);
|
||||
console.log(`[old liveness] pgrep + /proc for ${pids.length} PIDs: ${elapsed.toFixed(0)}ms`);
|
||||
console.log(`[old liveness] max event loop lag: ${lag.maxLagMs.toFixed(1)}ms`);
|
||||
});
|
||||
|
||||
// Document the cost: this runs every 10 seconds and takes N×hundreds of ms
|
||||
// During this time, SSE broadcasts, terminal data, and API responses are delayed
|
||||
it('should measure new tier-1: file stat for 20 agents (fast path)', async () => {
|
||||
// Create 20 agent files to stat
|
||||
const agentDir = join(tmpDir, 'tier1-agents');
|
||||
mkdirSync(agentDir, { recursive: true });
|
||||
const files: string[] = [];
|
||||
for (let i = 0; i < 20; i++) {
|
||||
const f = join(agentDir, `agent-${i}.jsonl`);
|
||||
writeFileSync(f, generateAgentTranscript(50));
|
||||
files.push(f);
|
||||
}
|
||||
|
||||
const { stat: statFn } = await import('node:fs/promises');
|
||||
const lagPromise = measureEventLoopLag(500);
|
||||
|
||||
const start = performance.now();
|
||||
let aliveCount = 0;
|
||||
for (const f of files) {
|
||||
try {
|
||||
const s = await statFn(f);
|
||||
if (Date.now() - s.mtime.getTime() < 30000) aliveCount++;
|
||||
} catch { /* */ }
|
||||
}
|
||||
const elapsed = performance.now() - start;
|
||||
const lag = await lagPromise;
|
||||
|
||||
console.log(`[tier-1 file stat × 20 agents] ${elapsed.toFixed(1)}ms, all alive: ${aliveCount === 20}`);
|
||||
console.log(`[tier-1 file stat] max event loop lag: ${lag.maxLagMs.toFixed(1)}ms`);
|
||||
|
||||
// File stat for 20 agents should be well under 5ms total
|
||||
expect(elapsed).toBeLessThan(20);
|
||||
expect(aliveCount).toBe(20);
|
||||
});
|
||||
|
||||
it('should measure new tier-2: /proc/{pid}/stat for 20 cached PIDs', async () => {
|
||||
const { stat: statFn } = await import('node:fs/promises');
|
||||
const ourPid = process.pid;
|
||||
|
||||
const lagPromise = measureEventLoopLag(500);
|
||||
|
||||
const start = performance.now();
|
||||
let aliveCount = 0;
|
||||
// Simulate checking 20 cached PIDs (all point to our own PID for testing)
|
||||
for (let i = 0; i < 20; i++) {
|
||||
try {
|
||||
await statFn(`/proc/${ourPid}`);
|
||||
aliveCount++;
|
||||
} catch { /* */ }
|
||||
}
|
||||
const elapsed = performance.now() - start;
|
||||
const lag = await lagPromise;
|
||||
|
||||
console.log(`[tier-2 /proc/pid × 20 agents] ${elapsed.toFixed(1)}ms, alive: ${aliveCount}`);
|
||||
console.log(`[tier-2 /proc/pid] max event loop lag: ${lag.maxLagMs.toFixed(1)}ms`);
|
||||
|
||||
// /proc/pid stat for 20 agents should be well under 5ms total
|
||||
expect(elapsed).toBeLessThan(20);
|
||||
});
|
||||
|
||||
it('should show tier-1+2 is orders of magnitude faster than full pgrep scan', async () => {
|
||||
const { stat: statFn } = await import('node:fs/promises');
|
||||
const { execFile: execFileFn } = await import('node:child_process');
|
||||
|
||||
// Create 20 agent files
|
||||
const agentDir = join(tmpDir, 'comparison-agents');
|
||||
mkdirSync(agentDir, { recursive: true });
|
||||
const files: string[] = [];
|
||||
for (let i = 0; i < 20; i++) {
|
||||
const f = join(agentDir, `agent-${i}.jsonl`);
|
||||
writeFileSync(f, 'x'.repeat(100));
|
||||
files.push(f);
|
||||
}
|
||||
|
||||
// Tier 1+2 approach: stat files + stat /proc/pid
|
||||
const startFast = performance.now();
|
||||
for (const f of files) {
|
||||
await statFn(f); // tier 1
|
||||
}
|
||||
for (let i = 0; i < 20; i++) {
|
||||
try { await statFn(`/proc/${process.pid}`); } catch { /* */ } // tier 2
|
||||
}
|
||||
const fastElapsed = performance.now() - startFast;
|
||||
|
||||
// Old approach: pgrep + /proc reads
|
||||
const startSlow = performance.now();
|
||||
const { readFile } = await import('node:fs/promises');
|
||||
const pids = await new Promise<string[]>((resolve) => {
|
||||
execFileFn('pgrep', ['-f', 'node'], { encoding: 'utf8' }, (_err, stdout) => {
|
||||
resolve((stdout || '').trim().split('\n').filter(Boolean));
|
||||
});
|
||||
});
|
||||
for (const pidStr of pids.slice(0, 10)) {
|
||||
try { await readFile(`/proc/${pidStr}/environ`, 'utf8'); } catch { /* */ }
|
||||
try { await readFile(`/proc/${pidStr}/cmdline`, 'utf8'); } catch { /* */ }
|
||||
}
|
||||
const slowElapsed = performance.now() - startSlow;
|
||||
|
||||
const speedup = slowElapsed / Math.max(fastElapsed, 0.01);
|
||||
console.log(`[comparison] tier-1+2: ${fastElapsed.toFixed(1)}ms, old pgrep: ${slowElapsed.toFixed(1)}ms, speedup: ${speedup.toFixed(0)}x`);
|
||||
|
||||
// Tiered approach should be significantly faster
|
||||
expect(fastElapsed).toBeLessThan(slowElapsed);
|
||||
});
|
||||
});
|
||||
|
||||
|
||||
Reference in New Issue
Block a user