mirror of
https://github.com/Ark0N/Codeman.git
synced 2026-10-09 16:59:43 +02:00
feat: add long-running session support and process management
- Add buffer management for 12-24+ hour sessions - Terminal buffer: 5MB max with auto-trim to 4MB - Text output: 2MB max with auto-trim to 1.5MB - Messages: 1000 max, keeps recent 800 - Add DELETE /api/sessions endpoint to kill all sessions - Enhanced session.stop() with SIGKILL fallback - Kill process groups to terminate child processes - Clean up respawn controller on session delete - Add terminal batching at 60fps for performance - Add bufferStats to session details for monitoring Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
This commit is contained in:
+264
-69
@@ -1,8 +1,18 @@
|
|||||||
import { spawn, ChildProcess } from 'node:child_process';
|
|
||||||
import { EventEmitter } from 'node:events';
|
import { EventEmitter } from 'node:events';
|
||||||
import { v4 as uuidv4 } from 'uuid';
|
import { v4 as uuidv4 } from 'uuid';
|
||||||
|
import * as pty from 'node-pty';
|
||||||
import { SessionState, SessionStatus, SessionConfig } from './types.js';
|
import { SessionState, SessionStatus, SessionConfig } from './types.js';
|
||||||
|
|
||||||
|
// Maximum terminal buffer size in characters (default 5MB of text)
|
||||||
|
const MAX_TERMINAL_BUFFER_SIZE = 5 * 1024 * 1024;
|
||||||
|
// When trimming, keep the most recent portion (4MB)
|
||||||
|
const TERMINAL_BUFFER_TRIM_SIZE = 4 * 1024 * 1024;
|
||||||
|
// Maximum text output buffer size (2MB)
|
||||||
|
const MAX_TEXT_OUTPUT_SIZE = 2 * 1024 * 1024;
|
||||||
|
const TEXT_OUTPUT_TRIM_SIZE = 1.5 * 1024 * 1024;
|
||||||
|
// Maximum number of Claude messages to keep in memory
|
||||||
|
const MAX_MESSAGES = 1000;
|
||||||
|
|
||||||
export interface ClaudeMessage {
|
export interface ClaudeMessage {
|
||||||
type: 'system' | 'assistant' | 'user' | 'result';
|
type: 'system' | 'assistant' | 'user' | 'result';
|
||||||
subtype?: string;
|
subtype?: string;
|
||||||
@@ -26,6 +36,7 @@ export interface SessionEvents {
|
|||||||
error: (data: string) => void;
|
error: (data: string) => void;
|
||||||
exit: (code: number | null) => void;
|
exit: (code: number | null) => void;
|
||||||
completion: (result: string, cost: number) => void;
|
completion: (result: string, cost: number) => void;
|
||||||
|
terminal: (data: string) => void; // Raw terminal data
|
||||||
}
|
}
|
||||||
|
|
||||||
export class Session extends EventEmitter {
|
export class Session extends EventEmitter {
|
||||||
@@ -33,9 +44,11 @@ export class Session extends EventEmitter {
|
|||||||
readonly workingDir: string;
|
readonly workingDir: string;
|
||||||
readonly createdAt: number;
|
readonly createdAt: number;
|
||||||
|
|
||||||
private process: ChildProcess | null = null;
|
private ptyProcess: pty.IPty | null = null;
|
||||||
|
private _pid: number | null = null;
|
||||||
private _status: SessionStatus = 'idle';
|
private _status: SessionStatus = 'idle';
|
||||||
private _currentTaskId: string | null = null;
|
private _currentTaskId: string | null = null;
|
||||||
|
private _terminalBuffer: string = ''; // Raw terminal output
|
||||||
private _outputBuffer: string = '';
|
private _outputBuffer: string = '';
|
||||||
private _textOutput: string = '';
|
private _textOutput: string = '';
|
||||||
private _errorBuffer: string = '';
|
private _errorBuffer: string = '';
|
||||||
@@ -44,6 +57,11 @@ export class Session extends EventEmitter {
|
|||||||
private _totalCost: number = 0;
|
private _totalCost: number = 0;
|
||||||
private _messages: ClaudeMessage[] = [];
|
private _messages: ClaudeMessage[] = [];
|
||||||
private _lineBuffer: string = '';
|
private _lineBuffer: string = '';
|
||||||
|
private resolvePromise: ((value: { result: string; cost: number }) => void) | null = null;
|
||||||
|
private rejectPromise: ((reason: Error) => void) | null = null;
|
||||||
|
private _isWorking: boolean = false;
|
||||||
|
private _lastPromptTime: number = 0;
|
||||||
|
private activityTimeout: NodeJS.Timeout | null = null;
|
||||||
|
|
||||||
constructor(config: Partial<SessionConfig> & { workingDir: string }) {
|
constructor(config: Partial<SessionConfig> & { workingDir: string }) {
|
||||||
super();
|
super();
|
||||||
@@ -62,7 +80,11 @@ export class Session extends EventEmitter {
|
|||||||
}
|
}
|
||||||
|
|
||||||
get pid(): number | null {
|
get pid(): number | null {
|
||||||
return this.process?.pid ?? null;
|
return this._pid;
|
||||||
|
}
|
||||||
|
|
||||||
|
get terminalBuffer(): string {
|
||||||
|
return this._terminalBuffer;
|
||||||
}
|
}
|
||||||
|
|
||||||
get outputBuffer(): string {
|
get outputBuffer(): string {
|
||||||
@@ -93,6 +115,14 @@ export class Session extends EventEmitter {
|
|||||||
return this._messages;
|
return this._messages;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
get isWorking(): boolean {
|
||||||
|
return this._isWorking;
|
||||||
|
}
|
||||||
|
|
||||||
|
get lastPromptTime(): number {
|
||||||
|
return this._lastPromptTime;
|
||||||
|
}
|
||||||
|
|
||||||
isIdle(): boolean {
|
isIdle(): boolean {
|
||||||
return this._status === 'idle';
|
return this._status === 'idle';
|
||||||
}
|
}
|
||||||
@@ -123,18 +153,110 @@ export class Session extends EventEmitter {
|
|||||||
claudeSessionId: this._claudeSessionId,
|
claudeSessionId: this._claudeSessionId,
|
||||||
totalCost: this._totalCost,
|
totalCost: this._totalCost,
|
||||||
textOutput: this._textOutput,
|
textOutput: this._textOutput,
|
||||||
|
terminalBuffer: this._terminalBuffer,
|
||||||
messageCount: this._messages.length,
|
messageCount: this._messages.length,
|
||||||
|
isWorking: this._isWorking,
|
||||||
|
lastPromptTime: this._lastPromptTime,
|
||||||
|
// Buffer statistics for monitoring long-running sessions
|
||||||
|
bufferStats: {
|
||||||
|
terminalBufferSize: this._terminalBuffer.length,
|
||||||
|
textOutputSize: this._textOutput.length,
|
||||||
|
messageCount: this._messages.length,
|
||||||
|
maxTerminalBuffer: MAX_TERMINAL_BUFFER_SIZE,
|
||||||
|
maxTextOutput: MAX_TEXT_OUTPUT_SIZE,
|
||||||
|
maxMessages: MAX_MESSAGES,
|
||||||
|
},
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Start an interactive Claude Code session (full terminal)
|
||||||
|
async startInteractive(): Promise<void> {
|
||||||
|
if (this.ptyProcess) {
|
||||||
|
throw new Error('Session already has a running process');
|
||||||
|
}
|
||||||
|
|
||||||
|
this._status = 'busy';
|
||||||
|
this._terminalBuffer = '';
|
||||||
|
this._outputBuffer = '';
|
||||||
|
this._textOutput = '';
|
||||||
|
this._errorBuffer = '';
|
||||||
|
this._messages = [];
|
||||||
|
this._lineBuffer = '';
|
||||||
|
this._lastActivityAt = Date.now();
|
||||||
|
|
||||||
|
console.log('[Session] Starting interactive Claude session');
|
||||||
|
|
||||||
|
this.ptyProcess = pty.spawn('claude', [
|
||||||
|
'--dangerously-skip-permissions'
|
||||||
|
], {
|
||||||
|
name: 'xterm-256color',
|
||||||
|
cols: 120,
|
||||||
|
rows: 40,
|
||||||
|
cwd: this.workingDir,
|
||||||
|
env: { ...process.env, TERM: 'xterm-256color' },
|
||||||
|
});
|
||||||
|
|
||||||
|
this._pid = this.ptyProcess.pid;
|
||||||
|
console.log('[Session] Interactive PTY spawned with PID:', this._pid);
|
||||||
|
|
||||||
|
this.ptyProcess.onData((data: string) => {
|
||||||
|
this._terminalBuffer += data;
|
||||||
|
this._lastActivityAt = Date.now();
|
||||||
|
|
||||||
|
// Trim buffer if it exceeds max size to prevent memory issues
|
||||||
|
if (this._terminalBuffer.length > MAX_TERMINAL_BUFFER_SIZE) {
|
||||||
|
this._terminalBuffer = this._terminalBuffer.slice(-TERMINAL_BUFFER_TRIM_SIZE);
|
||||||
|
}
|
||||||
|
|
||||||
|
this.emit('terminal', data);
|
||||||
|
this.emit('output', data);
|
||||||
|
|
||||||
|
// Detect if Claude is working or at prompt
|
||||||
|
// The prompt line contains "❯" when waiting for input
|
||||||
|
if (data.includes('❯') || data.includes('\u276f')) {
|
||||||
|
// Reset activity timeout - if no activity for 2 seconds after prompt, Claude is idle
|
||||||
|
if (this.activityTimeout) clearTimeout(this.activityTimeout);
|
||||||
|
this.activityTimeout = setTimeout(() => {
|
||||||
|
if (this._isWorking) {
|
||||||
|
this._isWorking = false;
|
||||||
|
this._lastPromptTime = Date.now();
|
||||||
|
this.emit('idle');
|
||||||
|
}
|
||||||
|
}, 2000);
|
||||||
|
}
|
||||||
|
|
||||||
|
// Detect when Claude starts working (thinking, writing, etc)
|
||||||
|
if (data.includes('Thinking') || data.includes('Writing') || data.includes('Reading') ||
|
||||||
|
data.includes('Running') || data.includes('⠋') || data.includes('⠙') ||
|
||||||
|
data.includes('⠹') || data.includes('⠸') || data.includes('⠼') ||
|
||||||
|
data.includes('⠴') || data.includes('⠦') || data.includes('⠧')) {
|
||||||
|
if (!this._isWorking) {
|
||||||
|
this._isWorking = true;
|
||||||
|
this.emit('working');
|
||||||
|
}
|
||||||
|
// Reset timeout since Claude is active
|
||||||
|
if (this.activityTimeout) clearTimeout(this.activityTimeout);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
this.ptyProcess.onExit(({ exitCode }) => {
|
||||||
|
console.log('[Session] Interactive PTY exited with code:', exitCode);
|
||||||
|
this.ptyProcess = null;
|
||||||
|
this._pid = null;
|
||||||
|
this._status = 'idle';
|
||||||
|
this.emit('exit', exitCode);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
async runPrompt(prompt: string): Promise<{ result: string; cost: number }> {
|
async runPrompt(prompt: string): Promise<{ result: string; cost: number }> {
|
||||||
return new Promise((resolve, reject) => {
|
return new Promise((resolve, reject) => {
|
||||||
if (this.process) {
|
if (this.ptyProcess) {
|
||||||
reject(new Error('Session already has a running process'));
|
reject(new Error('Session already has a running process'));
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
this._status = 'busy';
|
this._status = 'busy';
|
||||||
|
this._terminalBuffer = '';
|
||||||
this._outputBuffer = '';
|
this._outputBuffer = '';
|
||||||
this._textOutput = '';
|
this._textOutput = '';
|
||||||
this._errorBuffer = '';
|
this._errorBuffer = '';
|
||||||
@@ -142,43 +264,52 @@ export class Session extends EventEmitter {
|
|||||||
this._lineBuffer = '';
|
this._lineBuffer = '';
|
||||||
this._lastActivityAt = Date.now();
|
this._lastActivityAt = Date.now();
|
||||||
|
|
||||||
|
this.resolvePromise = resolve;
|
||||||
|
this.rejectPromise = reject;
|
||||||
|
|
||||||
try {
|
try {
|
||||||
// Spawn claude with streaming JSON output
|
// Spawn claude in a real PTY
|
||||||
this.process = spawn('claude', [
|
console.log('[Session] Spawning PTY for claude with prompt:', prompt.substring(0, 50));
|
||||||
|
|
||||||
|
this.ptyProcess = pty.spawn('claude', [
|
||||||
'-p',
|
'-p',
|
||||||
'--output-format', 'stream-json',
|
'--dangerously-skip-permissions',
|
||||||
prompt
|
prompt
|
||||||
], {
|
], {
|
||||||
|
name: 'xterm-256color',
|
||||||
|
cols: 120,
|
||||||
|
rows: 40,
|
||||||
cwd: this.workingDir,
|
cwd: this.workingDir,
|
||||||
stdio: ['pipe', 'pipe', 'pipe'],
|
env: { ...process.env, TERM: 'xterm-256color' },
|
||||||
env: { ...process.env },
|
|
||||||
});
|
});
|
||||||
|
|
||||||
this.process.stdout?.on('data', (data: Buffer) => {
|
this._pid = this.ptyProcess.pid;
|
||||||
const text = data.toString();
|
console.log('[Session] PTY spawned with PID:', this._pid);
|
||||||
this._outputBuffer += text;
|
|
||||||
|
// Handle terminal data
|
||||||
|
this.ptyProcess.onData((data: string) => {
|
||||||
|
this._terminalBuffer += data;
|
||||||
this._lastActivityAt = Date.now();
|
this._lastActivityAt = Date.now();
|
||||||
this.emit('output', text);
|
|
||||||
this.processJsonLines(text);
|
// Trim buffer if it exceeds max size to prevent memory issues
|
||||||
|
if (this._terminalBuffer.length > MAX_TERMINAL_BUFFER_SIZE) {
|
||||||
|
this._terminalBuffer = this._terminalBuffer.slice(-TERMINAL_BUFFER_TRIM_SIZE);
|
||||||
|
}
|
||||||
|
|
||||||
|
this.emit('terminal', data);
|
||||||
|
this.emit('output', data);
|
||||||
|
|
||||||
|
// Also try to parse JSON lines for structured data
|
||||||
|
this.processOutput(data);
|
||||||
});
|
});
|
||||||
|
|
||||||
this.process.stderr?.on('data', (data: Buffer) => {
|
// Handle exit
|
||||||
const text = data.toString();
|
this.ptyProcess.onExit(({ exitCode }) => {
|
||||||
this._errorBuffer += text;
|
console.log('[Session] PTY exited with code:', exitCode);
|
||||||
this._lastActivityAt = Date.now();
|
this.ptyProcess = null;
|
||||||
this.emit('error', text);
|
this._pid = null;
|
||||||
});
|
|
||||||
|
|
||||||
this.process.on('error', (err) => {
|
// Find result from parsed messages or use text output
|
||||||
this._status = 'error';
|
|
||||||
this.process = null;
|
|
||||||
reject(err);
|
|
||||||
});
|
|
||||||
|
|
||||||
this.process.on('exit', (code) => {
|
|
||||||
this.process = null;
|
|
||||||
|
|
||||||
// Find the result message
|
|
||||||
const resultMsg = this._messages.find(m => m.type === 'result');
|
const resultMsg = this._messages.find(m => m.type === 'result');
|
||||||
|
|
||||||
if (resultMsg && !resultMsg.is_error) {
|
if (resultMsg && !resultMsg.is_error) {
|
||||||
@@ -186,17 +317,26 @@ export class Session extends EventEmitter {
|
|||||||
const cost = resultMsg.total_cost_usd || 0;
|
const cost = resultMsg.total_cost_usd || 0;
|
||||||
this._totalCost += cost;
|
this._totalCost += cost;
|
||||||
this.emit('completion', resultMsg.result || '', cost);
|
this.emit('completion', resultMsg.result || '', cost);
|
||||||
resolve({ result: resultMsg.result || '', cost });
|
if (this.resolvePromise) {
|
||||||
} else if (code !== 0 || (resultMsg && resultMsg.is_error)) {
|
this.resolvePromise({ result: resultMsg.result || '', cost });
|
||||||
|
}
|
||||||
|
} else if (exitCode !== 0 || (resultMsg && resultMsg.is_error)) {
|
||||||
this._status = 'error';
|
this._status = 'error';
|
||||||
reject(new Error(this._errorBuffer || 'Process exited with error'));
|
if (this.rejectPromise) {
|
||||||
|
this.rejectPromise(new Error(this._errorBuffer || this._textOutput || 'Process exited with error'));
|
||||||
|
}
|
||||||
} else {
|
} else {
|
||||||
this._status = 'idle';
|
this._status = 'idle';
|
||||||
resolve({ result: this._textOutput, cost: 0 });
|
if (this.resolvePromise) {
|
||||||
|
this.resolvePromise({ result: this._textOutput || this._terminalBuffer, cost: this._totalCost });
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
this.emit('exit', code);
|
this.resolvePromise = null;
|
||||||
|
this.rejectPromise = null;
|
||||||
|
this.emit('exit', exitCode);
|
||||||
});
|
});
|
||||||
|
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
this._status = 'error';
|
this._status = 'error';
|
||||||
reject(err);
|
reject(err);
|
||||||
@@ -204,21 +344,28 @@ export class Session extends EventEmitter {
|
|||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
private processJsonLines(chunk: string): void {
|
private processOutput(data: string): void {
|
||||||
this._lineBuffer += chunk;
|
// Try to extract JSON from output (Claude may output JSON in stream mode)
|
||||||
|
this._lineBuffer += data;
|
||||||
const lines = this._lineBuffer.split('\n');
|
const lines = this._lineBuffer.split('\n');
|
||||||
|
|
||||||
// Keep incomplete line in buffer
|
|
||||||
this._lineBuffer = lines.pop() || '';
|
this._lineBuffer = lines.pop() || '';
|
||||||
|
|
||||||
for (const line of lines) {
|
for (const line of lines) {
|
||||||
if (line.trim()) {
|
const trimmed = line.trim();
|
||||||
|
// Remove ANSI escape codes for JSON parsing
|
||||||
|
const cleanLine = trimmed.replace(/\x1b\[[0-9;]*m/g, '');
|
||||||
|
|
||||||
|
if (cleanLine.startsWith('{') && cleanLine.endsWith('}')) {
|
||||||
try {
|
try {
|
||||||
const msg = JSON.parse(line) as ClaudeMessage;
|
const msg = JSON.parse(cleanLine) as ClaudeMessage;
|
||||||
this._messages.push(msg);
|
this._messages.push(msg);
|
||||||
this.emit('message', msg);
|
this.emit('message', msg);
|
||||||
|
|
||||||
// Extract text content
|
// Trim messages array for long-running sessions
|
||||||
|
if (this._messages.length > MAX_MESSAGES) {
|
||||||
|
this._messages = this._messages.slice(-Math.floor(MAX_MESSAGES * 0.8));
|
||||||
|
}
|
||||||
|
|
||||||
if (msg.type === 'system' && msg.session_id) {
|
if (msg.type === 'system' && msg.session_id) {
|
||||||
this._claudeSessionId = msg.session_id;
|
this._claudeSessionId = msg.session_id;
|
||||||
}
|
}
|
||||||
@@ -231,22 +378,40 @@ export class Session extends EventEmitter {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
if (msg.type === 'result') {
|
if (msg.type === 'result' && msg.total_cost_usd) {
|
||||||
if (msg.total_cost_usd) {
|
this._totalCost = msg.total_cost_usd;
|
||||||
this._totalCost += msg.total_cost_usd;
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
} catch {
|
} catch {
|
||||||
// Not JSON, treat as raw text
|
// Not JSON, just regular output
|
||||||
this._textOutput += line + '\n';
|
this._textOutput += line + '\n';
|
||||||
}
|
}
|
||||||
|
} else if (trimmed) {
|
||||||
|
this._textOutput += line + '\n';
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Trim text output buffer for long-running sessions
|
||||||
|
if (this._textOutput.length > MAX_TEXT_OUTPUT_SIZE) {
|
||||||
|
this._textOutput = this._textOutput.slice(-TEXT_OUTPUT_TRIM_SIZE);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Send input to the PTY (for interactive sessions)
|
||||||
|
write(data: string): void {
|
||||||
|
if (this.ptyProcess) {
|
||||||
|
this.ptyProcess.write(data);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Resize the PTY
|
||||||
|
resize(cols: number, rows: number): void {
|
||||||
|
if (this.ptyProcess) {
|
||||||
|
this.ptyProcess.resize(cols, rows);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// Legacy method for compatibility with session-manager
|
// Legacy method for compatibility with session-manager
|
||||||
async start(): Promise<void> {
|
async start(): Promise<void> {
|
||||||
// Session is ready by default, actual process starts with runPrompt
|
|
||||||
this._status = 'idle';
|
this._status = 'idle';
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -254,41 +419,66 @@ export class Session extends EventEmitter {
|
|||||||
async sendInput(input: string): Promise<void> {
|
async sendInput(input: string): Promise<void> {
|
||||||
this._status = 'busy';
|
this._status = 'busy';
|
||||||
this._lastActivityAt = Date.now();
|
this._lastActivityAt = Date.now();
|
||||||
// Run the prompt asynchronously
|
|
||||||
this.runPrompt(input).catch(err => {
|
this.runPrompt(input).catch(err => {
|
||||||
this.emit('error', err.message);
|
this.emit('error', err.message);
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
stop(): Promise<void> {
|
async stop(): Promise<void> {
|
||||||
return new Promise((resolve) => {
|
// Clear activity timeout to prevent memory leak
|
||||||
if (!this.process) {
|
if (this.activityTimeout) {
|
||||||
this._status = 'stopped';
|
clearTimeout(this.activityTimeout);
|
||||||
resolve();
|
this.activityTimeout = null;
|
||||||
return;
|
}
|
||||||
|
|
||||||
|
if (this.ptyProcess) {
|
||||||
|
const pid = this.ptyProcess.pid;
|
||||||
|
|
||||||
|
// First try graceful SIGTERM
|
||||||
|
try {
|
||||||
|
this.ptyProcess.kill();
|
||||||
|
} catch {
|
||||||
|
// Process may already be dead
|
||||||
}
|
}
|
||||||
|
|
||||||
const cleanup = () => {
|
// Give it a moment to terminate gracefully
|
||||||
this.process = null;
|
await new Promise(resolve => setTimeout(resolve, 100));
|
||||||
this._status = 'stopped';
|
|
||||||
this._currentTaskId = null;
|
|
||||||
resolve();
|
|
||||||
};
|
|
||||||
|
|
||||||
this.process.once('exit', cleanup);
|
// Force kill with SIGKILL if still alive
|
||||||
this.process.kill('SIGTERM');
|
try {
|
||||||
|
if (pid) {
|
||||||
setTimeout(() => {
|
process.kill(pid, 'SIGKILL');
|
||||||
if (this.process && !this.process.killed) {
|
|
||||||
this.process.kill('SIGKILL');
|
|
||||||
}
|
}
|
||||||
}, 5000);
|
} catch {
|
||||||
});
|
// Process already terminated
|
||||||
|
}
|
||||||
|
|
||||||
|
// Also try to kill any child processes in the process group
|
||||||
|
try {
|
||||||
|
if (pid) {
|
||||||
|
process.kill(-pid, 'SIGKILL');
|
||||||
|
}
|
||||||
|
} catch {
|
||||||
|
// Process group may not exist or already terminated
|
||||||
|
}
|
||||||
|
|
||||||
|
this.ptyProcess = null;
|
||||||
|
}
|
||||||
|
this._pid = null;
|
||||||
|
this._status = 'stopped';
|
||||||
|
this._currentTaskId = null;
|
||||||
|
|
||||||
|
if (this.rejectPromise) {
|
||||||
|
this.rejectPromise(new Error('Session stopped'));
|
||||||
|
this.resolvePromise = null;
|
||||||
|
this.rejectPromise = null;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
assignTask(taskId: string): void {
|
assignTask(taskId: string): void {
|
||||||
this._currentTaskId = taskId;
|
this._currentTaskId = taskId;
|
||||||
this._status = 'busy';
|
this._status = 'busy';
|
||||||
|
this._terminalBuffer = '';
|
||||||
this._outputBuffer = '';
|
this._outputBuffer = '';
|
||||||
this._textOutput = '';
|
this._textOutput = '';
|
||||||
this._errorBuffer = '';
|
this._errorBuffer = '';
|
||||||
@@ -310,7 +500,12 @@ export class Session extends EventEmitter {
|
|||||||
return this._errorBuffer;
|
return this._errorBuffer;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
getTerminalBuffer(): string {
|
||||||
|
return this._terminalBuffer;
|
||||||
|
}
|
||||||
|
|
||||||
clearBuffers(): void {
|
clearBuffers(): void {
|
||||||
|
this._terminalBuffer = '';
|
||||||
this._outputBuffer = '';
|
this._outputBuffer = '';
|
||||||
this._textOutput = '';
|
this._textOutput = '';
|
||||||
this._errorBuffer = '';
|
this._errorBuffer = '';
|
||||||
|
|||||||
+110
-1
@@ -92,6 +92,115 @@ export interface TaskAssignment {
|
|||||||
assignedAt: number;
|
assignedAt: number;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Error codes for consistent error handling
|
||||||
|
export enum ApiErrorCode {
|
||||||
|
NOT_FOUND = 'NOT_FOUND',
|
||||||
|
INVALID_INPUT = 'INVALID_INPUT',
|
||||||
|
SESSION_BUSY = 'SESSION_BUSY',
|
||||||
|
OPERATION_FAILED = 'OPERATION_FAILED',
|
||||||
|
ALREADY_EXISTS = 'ALREADY_EXISTS',
|
||||||
|
INTERNAL_ERROR = 'INTERNAL_ERROR',
|
||||||
|
}
|
||||||
|
|
||||||
|
// Mapping of error codes to user-friendly messages
|
||||||
|
export const ErrorMessages: Record<ApiErrorCode, string> = {
|
||||||
|
[ApiErrorCode.NOT_FOUND]: 'The requested resource was not found',
|
||||||
|
[ApiErrorCode.INVALID_INPUT]: 'Invalid input provided',
|
||||||
|
[ApiErrorCode.SESSION_BUSY]: 'Session is currently busy',
|
||||||
|
[ApiErrorCode.OPERATION_FAILED]: 'The operation failed',
|
||||||
|
[ApiErrorCode.ALREADY_EXISTS]: 'Resource already exists',
|
||||||
|
[ApiErrorCode.INTERNAL_ERROR]: 'An internal error occurred',
|
||||||
|
};
|
||||||
|
|
||||||
|
// API Request/Response types for type safety
|
||||||
|
export interface CreateSessionRequest {
|
||||||
|
workingDir?: string;
|
||||||
|
}
|
||||||
|
|
||||||
|
export interface RunPromptRequest {
|
||||||
|
prompt: string;
|
||||||
|
}
|
||||||
|
|
||||||
|
export interface SessionInputRequest {
|
||||||
|
input: string;
|
||||||
|
}
|
||||||
|
|
||||||
|
export interface ResizeRequest {
|
||||||
|
cols: number;
|
||||||
|
rows: number;
|
||||||
|
}
|
||||||
|
|
||||||
|
export interface CreateCaseRequest {
|
||||||
|
name: string;
|
||||||
|
description?: string;
|
||||||
|
}
|
||||||
|
|
||||||
|
export interface QuickStartRequest {
|
||||||
|
caseName?: string;
|
||||||
|
}
|
||||||
|
|
||||||
|
export interface CreateScheduledRunRequest {
|
||||||
|
prompt: string;
|
||||||
|
workingDir?: string;
|
||||||
|
durationMinutes: number;
|
||||||
|
}
|
||||||
|
|
||||||
|
export interface QuickRunRequest {
|
||||||
|
prompt: string;
|
||||||
|
workingDir?: string;
|
||||||
|
}
|
||||||
|
|
||||||
|
export interface ApiResponse<T = unknown> {
|
||||||
|
success: boolean;
|
||||||
|
error?: string;
|
||||||
|
errorCode?: ApiErrorCode;
|
||||||
|
data?: T;
|
||||||
|
}
|
||||||
|
|
||||||
|
// Helper functions for creating consistent error responses
|
||||||
|
export function createErrorResponse(code: ApiErrorCode, details?: string): ApiResponse {
|
||||||
|
return {
|
||||||
|
success: false,
|
||||||
|
error: details || ErrorMessages[code],
|
||||||
|
errorCode: code,
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
export function createSuccessResponse<T>(data?: T): ApiResponse<T> {
|
||||||
|
return {
|
||||||
|
success: true,
|
||||||
|
data,
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
export interface SessionResponse {
|
||||||
|
success: boolean;
|
||||||
|
session?: SessionState & {
|
||||||
|
claudeSessionId: string | null;
|
||||||
|
totalCost: number;
|
||||||
|
textOutput: string;
|
||||||
|
terminalBuffer: string;
|
||||||
|
messageCount: number;
|
||||||
|
isWorking: boolean;
|
||||||
|
lastPromptTime: number;
|
||||||
|
};
|
||||||
|
error?: string;
|
||||||
|
}
|
||||||
|
|
||||||
|
export interface QuickStartResponse {
|
||||||
|
success: boolean;
|
||||||
|
sessionId?: string;
|
||||||
|
casePath?: string;
|
||||||
|
caseName?: string;
|
||||||
|
error?: string;
|
||||||
|
}
|
||||||
|
|
||||||
|
export interface CaseInfo {
|
||||||
|
name: string;
|
||||||
|
path: string;
|
||||||
|
hasClaudeMd?: boolean;
|
||||||
|
}
|
||||||
|
|
||||||
export const DEFAULT_CONFIG: AppConfig = {
|
export const DEFAULT_CONFIG: AppConfig = {
|
||||||
pollIntervalMs: 1000,
|
pollIntervalMs: 1000,
|
||||||
defaultTimeoutMs: 300000, // 5 minutes
|
defaultTimeoutMs: 300000, // 5 minutes
|
||||||
@@ -101,7 +210,7 @@ export const DEFAULT_CONFIG: AppConfig = {
|
|||||||
idleTimeoutMs: 5000, // 5 seconds of no activity after prompt
|
idleTimeoutMs: 5000, // 5 seconds of no activity after prompt
|
||||||
updatePrompt: 'update all the docs and CLAUDE.md',
|
updatePrompt: 'update all the docs and CLAUDE.md',
|
||||||
interStepDelayMs: 1000, // 1 second between steps
|
interStepDelayMs: 1000, // 1 second between steps
|
||||||
enabled: true,
|
enabled: false, // disabled by default
|
||||||
},
|
},
|
||||||
};
|
};
|
||||||
|
|
||||||
|
|||||||
+105
-24
@@ -10,6 +10,20 @@ import { RespawnController, RespawnConfig, RespawnState } from '../respawn-contr
|
|||||||
import { getStore } from '../state-store.js';
|
import { getStore } from '../state-store.js';
|
||||||
import { generateClaudeMd } from '../templates/claude-md.js';
|
import { generateClaudeMd } from '../templates/claude-md.js';
|
||||||
import { v4 as uuidv4 } from 'uuid';
|
import { v4 as uuidv4 } from 'uuid';
|
||||||
|
import type {
|
||||||
|
CreateSessionRequest,
|
||||||
|
RunPromptRequest,
|
||||||
|
SessionInputRequest,
|
||||||
|
ResizeRequest,
|
||||||
|
CreateCaseRequest,
|
||||||
|
QuickStartRequest,
|
||||||
|
CreateScheduledRunRequest,
|
||||||
|
QuickRunRequest,
|
||||||
|
ApiResponse,
|
||||||
|
SessionResponse,
|
||||||
|
QuickStartResponse,
|
||||||
|
CaseInfo,
|
||||||
|
} from '../types.js';
|
||||||
|
|
||||||
const __dirname = dirname(fileURLToPath(import.meta.url));
|
const __dirname = dirname(fileURLToPath(import.meta.url));
|
||||||
|
|
||||||
@@ -27,6 +41,9 @@ interface ScheduledRun {
|
|||||||
logs: string[];
|
logs: string[];
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Batch terminal data for performance - collect for 16ms (60fps) before sending
|
||||||
|
const TERMINAL_BATCH_INTERVAL = 16;
|
||||||
|
|
||||||
export class WebServer extends EventEmitter {
|
export class WebServer extends EventEmitter {
|
||||||
private app: FastifyInstance;
|
private app: FastifyInstance;
|
||||||
private sessions: Map<string, Session> = new Map();
|
private sessions: Map<string, Session> = new Map();
|
||||||
@@ -35,6 +52,9 @@ export class WebServer extends EventEmitter {
|
|||||||
private sseClients: Set<FastifyReply> = new Set();
|
private sseClients: Set<FastifyReply> = new Set();
|
||||||
private store = getStore();
|
private store = getStore();
|
||||||
private port: number;
|
private port: number;
|
||||||
|
// Terminal batching for performance
|
||||||
|
private terminalBatches: Map<string, string> = new Map();
|
||||||
|
private terminalBatchTimer: NodeJS.Timeout | null = null;
|
||||||
|
|
||||||
constructor(port: number = 3000) {
|
constructor(port: number = 3000) {
|
||||||
super();
|
super();
|
||||||
@@ -74,8 +94,8 @@ export class WebServer extends EventEmitter {
|
|||||||
// Session management
|
// Session management
|
||||||
this.app.get('/api/sessions', async () => this.getSessionsState());
|
this.app.get('/api/sessions', async () => this.getSessionsState());
|
||||||
|
|
||||||
this.app.post('/api/sessions', async (req) => {
|
this.app.post('/api/sessions', async (req): Promise<SessionResponse> => {
|
||||||
const body = req.body as { workingDir?: string };
|
const body = req.body as CreateSessionRequest;
|
||||||
const workingDir = body.workingDir || process.cwd();
|
const workingDir = body.workingDir || process.cwd();
|
||||||
const session = new Session({ workingDir });
|
const session = new Session({ workingDir });
|
||||||
|
|
||||||
@@ -86,7 +106,7 @@ export class WebServer extends EventEmitter {
|
|||||||
return { success: true, session: session.toDetailedState() };
|
return { success: true, session: session.toDetailedState() };
|
||||||
});
|
});
|
||||||
|
|
||||||
this.app.delete('/api/sessions/:id', async (req) => {
|
this.app.delete('/api/sessions/:id', async (req): Promise<ApiResponse> => {
|
||||||
const { id } = req.params as { id: string };
|
const { id } = req.params as { id: string };
|
||||||
const session = this.sessions.get(id);
|
const session = this.sessions.get(id);
|
||||||
|
|
||||||
@@ -94,12 +114,46 @@ export class WebServer extends EventEmitter {
|
|||||||
return { success: false, error: 'Session not found' };
|
return { success: false, error: 'Session not found' };
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Stop respawn controller first
|
||||||
|
const controller = this.respawnControllers.get(id);
|
||||||
|
if (controller) {
|
||||||
|
controller.stop();
|
||||||
|
this.respawnControllers.delete(id);
|
||||||
|
}
|
||||||
|
|
||||||
await session.stop();
|
await session.stop();
|
||||||
this.sessions.delete(id);
|
this.sessions.delete(id);
|
||||||
|
this.terminalBatches.delete(id);
|
||||||
this.broadcast('session:deleted', { id });
|
this.broadcast('session:deleted', { id });
|
||||||
return { success: true };
|
return { success: true };
|
||||||
});
|
});
|
||||||
|
|
||||||
|
// Kill all sessions at once
|
||||||
|
this.app.delete('/api/sessions', async (): Promise<ApiResponse<{ killed: number }>> => {
|
||||||
|
const sessionIds = Array.from(this.sessions.keys());
|
||||||
|
let killed = 0;
|
||||||
|
|
||||||
|
for (const id of sessionIds) {
|
||||||
|
const session = this.sessions.get(id);
|
||||||
|
if (session) {
|
||||||
|
// Stop respawn controller first
|
||||||
|
const controller = this.respawnControllers.get(id);
|
||||||
|
if (controller) {
|
||||||
|
controller.stop();
|
||||||
|
this.respawnControllers.delete(id);
|
||||||
|
}
|
||||||
|
|
||||||
|
await session.stop();
|
||||||
|
this.sessions.delete(id);
|
||||||
|
this.terminalBatches.delete(id);
|
||||||
|
this.broadcast('session:deleted', { id });
|
||||||
|
killed++;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
return { success: true, data: { killed } };
|
||||||
|
});
|
||||||
|
|
||||||
this.app.get('/api/sessions/:id', async (req) => {
|
this.app.get('/api/sessions/:id', async (req) => {
|
||||||
const { id } = req.params as { id: string };
|
const { id } = req.params as { id: string };
|
||||||
const session = this.sessions.get(id);
|
const session = this.sessions.get(id);
|
||||||
@@ -127,9 +181,9 @@ export class WebServer extends EventEmitter {
|
|||||||
});
|
});
|
||||||
|
|
||||||
// Run prompt in session
|
// Run prompt in session
|
||||||
this.app.post('/api/sessions/:id/run', async (req) => {
|
this.app.post('/api/sessions/:id/run', async (req): Promise<{ success?: boolean; message?: string; error?: string }> => {
|
||||||
const { id } = req.params as { id: string };
|
const { id } = req.params as { id: string };
|
||||||
const { prompt } = req.body as { prompt: string };
|
const { prompt } = req.body as RunPromptRequest;
|
||||||
const session = this.sessions.get(id);
|
const session = this.sessions.get(id);
|
||||||
|
|
||||||
if (!session) {
|
if (!session) {
|
||||||
@@ -172,13 +226,13 @@ export class WebServer extends EventEmitter {
|
|||||||
});
|
});
|
||||||
|
|
||||||
// Send input to interactive session
|
// Send input to interactive session
|
||||||
this.app.post('/api/sessions/:id/input', async (req) => {
|
this.app.post('/api/sessions/:id/input', async (req): Promise<ApiResponse> => {
|
||||||
const { id } = req.params as { id: string };
|
const { id } = req.params as { id: string };
|
||||||
const { input } = req.body as { input: string };
|
const { input } = req.body as SessionInputRequest;
|
||||||
const session = this.sessions.get(id);
|
const session = this.sessions.get(id);
|
||||||
|
|
||||||
if (!session) {
|
if (!session) {
|
||||||
return { error: 'Session not found' };
|
return { success: false, error: 'Session not found' };
|
||||||
}
|
}
|
||||||
|
|
||||||
session.write(input);
|
session.write(input);
|
||||||
@@ -186,13 +240,13 @@ export class WebServer extends EventEmitter {
|
|||||||
});
|
});
|
||||||
|
|
||||||
// Resize session terminal
|
// Resize session terminal
|
||||||
this.app.post('/api/sessions/:id/resize', async (req) => {
|
this.app.post('/api/sessions/:id/resize', async (req): Promise<ApiResponse> => {
|
||||||
const { id } = req.params as { id: string };
|
const { id } = req.params as { id: string };
|
||||||
const { cols, rows } = req.body as { cols: number; rows: number };
|
const { cols, rows } = req.body as ResizeRequest;
|
||||||
const session = this.sessions.get(id);
|
const session = this.sessions.get(id);
|
||||||
|
|
||||||
if (!session) {
|
if (!session) {
|
||||||
return { error: 'Session not found' };
|
return { success: false, error: 'Session not found' };
|
||||||
}
|
}
|
||||||
|
|
||||||
session.resize(cols, rows);
|
session.resize(cols, rows);
|
||||||
@@ -327,7 +381,7 @@ export class WebServer extends EventEmitter {
|
|||||||
|
|
||||||
// Quick run (create session, run prompt, return result)
|
// Quick run (create session, run prompt, return result)
|
||||||
this.app.post('/api/run', async (req) => {
|
this.app.post('/api/run', async (req) => {
|
||||||
const { prompt, workingDir } = req.body as { prompt: string; workingDir?: string };
|
const { prompt, workingDir } = req.body as QuickRunRequest;
|
||||||
const dir = workingDir || process.cwd();
|
const dir = workingDir || process.cwd();
|
||||||
|
|
||||||
const session = new Session({ workingDir: dir });
|
const session = new Session({ workingDir: dir });
|
||||||
@@ -349,12 +403,8 @@ export class WebServer extends EventEmitter {
|
|||||||
return Array.from(this.scheduledRuns.values());
|
return Array.from(this.scheduledRuns.values());
|
||||||
});
|
});
|
||||||
|
|
||||||
this.app.post('/api/scheduled', async (req) => {
|
this.app.post('/api/scheduled', async (req): Promise<{ success: boolean; run: ScheduledRun }> => {
|
||||||
const { prompt, workingDir, durationMinutes } = req.body as {
|
const { prompt, workingDir, durationMinutes } = req.body as CreateScheduledRunRequest;
|
||||||
prompt: string;
|
|
||||||
workingDir?: string;
|
|
||||||
durationMinutes: number;
|
|
||||||
};
|
|
||||||
|
|
||||||
const run = await this.startScheduledRun(prompt, workingDir || process.cwd(), durationMinutes);
|
const run = await this.startScheduledRun(prompt, workingDir || process.cwd(), durationMinutes);
|
||||||
return { success: true, run };
|
return { success: true, run };
|
||||||
@@ -386,7 +436,7 @@ export class WebServer extends EventEmitter {
|
|||||||
// Case management
|
// Case management
|
||||||
const casesDir = join(homedir(), 'claudeman-cases');
|
const casesDir = join(homedir(), 'claudeman-cases');
|
||||||
|
|
||||||
this.app.get('/api/cases', async () => {
|
this.app.get('/api/cases', async (): Promise<CaseInfo[]> => {
|
||||||
if (!existsSync(casesDir)) {
|
if (!existsSync(casesDir)) {
|
||||||
return [];
|
return [];
|
||||||
}
|
}
|
||||||
@@ -400,8 +450,8 @@ export class WebServer extends EventEmitter {
|
|||||||
}));
|
}));
|
||||||
});
|
});
|
||||||
|
|
||||||
this.app.post('/api/cases', async (req) => {
|
this.app.post('/api/cases', async (req): Promise<{ success: boolean; case?: { name: string; path: string }; error?: string }> => {
|
||||||
const { name, description } = req.body as { name: string; description?: string };
|
const { name, description } = req.body as CreateCaseRequest;
|
||||||
|
|
||||||
if (!name || !/^[a-zA-Z0-9_-]+$/.test(name)) {
|
if (!name || !/^[a-zA-Z0-9_-]+$/.test(name)) {
|
||||||
return { success: false, error: 'Invalid case name. Use only letters, numbers, hyphens, underscores.' };
|
return { success: false, error: 'Invalid case name. Use only letters, numbers, hyphens, underscores.' };
|
||||||
@@ -444,8 +494,8 @@ export class WebServer extends EventEmitter {
|
|||||||
});
|
});
|
||||||
|
|
||||||
// Quick Start: Create case (if needed) and start interactive session in one click
|
// Quick Start: Create case (if needed) and start interactive session in one click
|
||||||
this.app.post('/api/quick-start', async (req) => {
|
this.app.post('/api/quick-start', async (req): Promise<QuickStartResponse> => {
|
||||||
const { caseName = 'testcase' } = req.body as { caseName?: string };
|
const { caseName = 'testcase' } = req.body as QuickStartRequest;
|
||||||
|
|
||||||
// Validate case name
|
// Validate case name
|
||||||
if (!/^[a-zA-Z0-9_-]+$/.test(caseName)) {
|
if (!/^[a-zA-Z0-9_-]+$/.test(caseName)) {
|
||||||
@@ -498,7 +548,8 @@ export class WebServer extends EventEmitter {
|
|||||||
});
|
});
|
||||||
|
|
||||||
session.on('terminal', (data) => {
|
session.on('terminal', (data) => {
|
||||||
this.broadcast('session:terminal', { id: session.id, data });
|
// Use batching for better performance at high throughput
|
||||||
|
this.batchTerminalData(session.id, data);
|
||||||
});
|
});
|
||||||
|
|
||||||
session.on('message', (msg: ClaudeMessage) => {
|
session.on('message', (msg: ClaudeMessage) => {
|
||||||
@@ -692,6 +743,29 @@ export class WebServer extends EventEmitter {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Batch terminal data for better performance (60fps)
|
||||||
|
private batchTerminalData(sessionId: string, data: string): void {
|
||||||
|
const existing = this.terminalBatches.get(sessionId) || '';
|
||||||
|
this.terminalBatches.set(sessionId, existing + data);
|
||||||
|
|
||||||
|
// Start batch timer if not already running
|
||||||
|
if (!this.terminalBatchTimer) {
|
||||||
|
this.terminalBatchTimer = setTimeout(() => {
|
||||||
|
this.flushTerminalBatches();
|
||||||
|
this.terminalBatchTimer = null;
|
||||||
|
}, TERMINAL_BATCH_INTERVAL);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private flushTerminalBatches(): void {
|
||||||
|
for (const [sessionId, data] of this.terminalBatches) {
|
||||||
|
if (data.length > 0) {
|
||||||
|
this.broadcast('session:terminal', { id: sessionId, data });
|
||||||
|
}
|
||||||
|
}
|
||||||
|
this.terminalBatches.clear();
|
||||||
|
}
|
||||||
|
|
||||||
async start(): Promise<void> {
|
async start(): Promise<void> {
|
||||||
await this.setupRoutes();
|
await this.setupRoutes();
|
||||||
await this.app.listen({ port: this.port, host: '0.0.0.0' });
|
await this.app.listen({ port: this.port, host: '0.0.0.0' });
|
||||||
@@ -699,6 +773,13 @@ export class WebServer extends EventEmitter {
|
|||||||
}
|
}
|
||||||
|
|
||||||
async stop(): Promise<void> {
|
async stop(): Promise<void> {
|
||||||
|
// Clear batch timer
|
||||||
|
if (this.terminalBatchTimer) {
|
||||||
|
clearTimeout(this.terminalBatchTimer);
|
||||||
|
this.terminalBatchTimer = null;
|
||||||
|
}
|
||||||
|
this.terminalBatches.clear();
|
||||||
|
|
||||||
// Stop all respawn controllers
|
// Stop all respawn controllers
|
||||||
for (const controller of this.respawnControllers.values()) {
|
for (const controller of this.respawnControllers.values()) {
|
||||||
controller.stop();
|
controller.stop();
|
||||||
|
|||||||
Reference in New Issue
Block a user