mirror of
https://github.com/Ark0N/Codeman.git
synced 2026-10-03 14:09:42 +02:00
feat: implement Claudeman - Claude session manager with Ralph Loop
Claudeman manages multiple Claude CLI sessions as subprocesses with: - Session management (start/stop/list/logs) - Priority-based task queue with dependency support - Ralph Loop for autonomous task assignment and completion detection - Time-aware loops for extended work sessions - JSON file persistence to ~/.claudeman/state.json Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
This commit is contained in:
+452
@@ -0,0 +1,452 @@
|
||||
import { Command } from 'commander';
|
||||
import chalk from 'chalk';
|
||||
import { getSessionManager } from './session-manager.js';
|
||||
import { getTaskQueue } from './task-queue.js';
|
||||
import { getRalphLoop } from './ralph-loop.js';
|
||||
import { getStore } from './state-store.js';
|
||||
|
||||
const program = new Command();
|
||||
|
||||
program
|
||||
.name('claudeman')
|
||||
.description('Claude Code session manager with autonomous Ralph Loop')
|
||||
.version('1.0.0');
|
||||
|
||||
// ============ Session Commands ============
|
||||
|
||||
const sessionCmd = program
|
||||
.command('session')
|
||||
.alias('s')
|
||||
.description('Manage Claude sessions');
|
||||
|
||||
sessionCmd
|
||||
.command('start')
|
||||
.description('Start a new Claude session')
|
||||
.option('-d, --dir <path>', 'Working directory', process.cwd())
|
||||
.action(async (options) => {
|
||||
try {
|
||||
const manager = getSessionManager();
|
||||
const session = await manager.createSession(options.dir);
|
||||
console.log(chalk.green(`✓ Session started: ${session.id}`));
|
||||
console.log(` Working directory: ${session.workingDir}`);
|
||||
console.log(` PID: ${session.pid}`);
|
||||
} catch (err) {
|
||||
console.error(chalk.red(`✗ Failed to start session: ${(err as Error).message}`));
|
||||
process.exit(1);
|
||||
}
|
||||
});
|
||||
|
||||
sessionCmd
|
||||
.command('stop <id>')
|
||||
.description('Stop a session')
|
||||
.action(async (id) => {
|
||||
try {
|
||||
const manager = getSessionManager();
|
||||
await manager.stopSession(id);
|
||||
console.log(chalk.green(`✓ Session stopped: ${id}`));
|
||||
} catch (err) {
|
||||
console.error(chalk.red(`✗ Failed to stop session: ${(err as Error).message}`));
|
||||
process.exit(1);
|
||||
}
|
||||
});
|
||||
|
||||
sessionCmd
|
||||
.command('list')
|
||||
.alias('ls')
|
||||
.description('List all sessions')
|
||||
.action(() => {
|
||||
const manager = getSessionManager();
|
||||
const sessions = manager.getAllSessions();
|
||||
const stored = manager.getStoredSessions();
|
||||
|
||||
if (sessions.length === 0 && Object.keys(stored).length === 0) {
|
||||
console.log(chalk.yellow('No sessions found'));
|
||||
return;
|
||||
}
|
||||
|
||||
console.log(chalk.bold('\nActive Sessions:'));
|
||||
if (sessions.length === 0) {
|
||||
console.log(' (none)');
|
||||
} else {
|
||||
for (const session of sessions) {
|
||||
const status = session.status === 'idle'
|
||||
? chalk.green('idle')
|
||||
: session.status === 'busy'
|
||||
? chalk.yellow('busy')
|
||||
: chalk.red(session.status);
|
||||
console.log(` ${chalk.cyan(session.id.slice(0, 8))} ${status} ${session.workingDir}`);
|
||||
}
|
||||
}
|
||||
|
||||
const stoppedSessions = Object.values(stored).filter((s) => s.status === 'stopped');
|
||||
if (stoppedSessions.length > 0) {
|
||||
console.log(chalk.bold('\nStopped Sessions:'));
|
||||
for (const session of stoppedSessions) {
|
||||
console.log(` ${chalk.gray(session.id.slice(0, 8))} ${chalk.gray('stopped')} ${session.workingDir}`);
|
||||
}
|
||||
}
|
||||
console.log('');
|
||||
});
|
||||
|
||||
sessionCmd
|
||||
.command('logs <id>')
|
||||
.description('View session output')
|
||||
.option('-e, --errors', 'Show stderr instead of stdout')
|
||||
.action((id, options) => {
|
||||
const manager = getSessionManager();
|
||||
const output = options.errors
|
||||
? manager.getSessionError(id)
|
||||
: manager.getSessionOutput(id);
|
||||
|
||||
if (output === null) {
|
||||
console.log(chalk.yellow(`Session ${id} not found or not active`));
|
||||
return;
|
||||
}
|
||||
|
||||
if (output === '') {
|
||||
console.log(chalk.gray('(no output)'));
|
||||
return;
|
||||
}
|
||||
|
||||
console.log(output);
|
||||
});
|
||||
|
||||
// ============ Task Commands ============
|
||||
|
||||
const taskCmd = program
|
||||
.command('task')
|
||||
.alias('t')
|
||||
.description('Manage tasks');
|
||||
|
||||
taskCmd
|
||||
.command('add <prompt>')
|
||||
.description('Add a new task')
|
||||
.option('-d, --dir <path>', 'Working directory', process.cwd())
|
||||
.option('-p, --priority <n>', 'Priority (higher = first)', '0')
|
||||
.option('-c, --completion <phrase>', 'Completion phrase to detect')
|
||||
.option('--timeout <ms>', 'Timeout in milliseconds')
|
||||
.action((prompt, options) => {
|
||||
const queue = getTaskQueue();
|
||||
const task = queue.addTask({
|
||||
prompt,
|
||||
workingDir: options.dir,
|
||||
priority: parseInt(options.priority, 10),
|
||||
completionPhrase: options.completion,
|
||||
timeoutMs: options.timeout ? parseInt(options.timeout, 10) : undefined,
|
||||
});
|
||||
console.log(chalk.green(`✓ Task added: ${task.id}`));
|
||||
console.log(` Prompt: ${prompt.slice(0, 50)}${prompt.length > 50 ? '...' : ''}`);
|
||||
console.log(` Priority: ${task.priority}`);
|
||||
});
|
||||
|
||||
taskCmd
|
||||
.command('list')
|
||||
.alias('ls')
|
||||
.description('List all tasks')
|
||||
.option('-s, --status <status>', 'Filter by status (pending, running, completed, failed)')
|
||||
.action((options) => {
|
||||
const queue = getTaskQueue();
|
||||
let tasks = queue.getAllTasks();
|
||||
|
||||
if (options.status) {
|
||||
tasks = tasks.filter((t) => t.status === options.status);
|
||||
}
|
||||
|
||||
if (tasks.length === 0) {
|
||||
console.log(chalk.yellow('No tasks found'));
|
||||
return;
|
||||
}
|
||||
|
||||
const statusColors = {
|
||||
pending: chalk.gray,
|
||||
running: chalk.yellow,
|
||||
completed: chalk.green,
|
||||
failed: chalk.red,
|
||||
};
|
||||
|
||||
console.log(chalk.bold('\nTasks:'));
|
||||
for (const task of tasks) {
|
||||
const color = statusColors[task.status];
|
||||
const prompt = task.prompt.slice(0, 40) + (task.prompt.length > 40 ? '...' : '');
|
||||
console.log(` ${chalk.cyan(task.id.slice(0, 8))} ${color(task.status.padEnd(10))} [${task.priority}] ${prompt}`);
|
||||
}
|
||||
|
||||
const counts = queue.getCount();
|
||||
console.log(chalk.bold('\nSummary:'));
|
||||
console.log(` Pending: ${counts.pending}, Running: ${counts.running}, Completed: ${counts.completed}, Failed: ${counts.failed}`);
|
||||
console.log('');
|
||||
});
|
||||
|
||||
taskCmd
|
||||
.command('status <id>')
|
||||
.description('Show task details')
|
||||
.action((id) => {
|
||||
const queue = getTaskQueue();
|
||||
const task = queue.getTask(id);
|
||||
|
||||
if (!task) {
|
||||
console.log(chalk.red(`Task ${id} not found`));
|
||||
return;
|
||||
}
|
||||
|
||||
console.log(chalk.bold('\nTask Details:'));
|
||||
console.log(` ID: ${task.id}`);
|
||||
console.log(` Status: ${task.status}`);
|
||||
console.log(` Priority: ${task.priority}`);
|
||||
console.log(` Prompt: ${task.prompt}`);
|
||||
console.log(` Working Dir: ${task.workingDir}`);
|
||||
if (task.assignedSessionId) {
|
||||
console.log(` Session: ${task.assignedSessionId}`);
|
||||
}
|
||||
if (task.error) {
|
||||
console.log(` Error: ${chalk.red(task.error)}`);
|
||||
}
|
||||
if (task.output) {
|
||||
console.log(chalk.bold('\nOutput:'));
|
||||
console.log(task.output.slice(0, 500) + (task.output.length > 500 ? '...' : ''));
|
||||
}
|
||||
console.log('');
|
||||
});
|
||||
|
||||
taskCmd
|
||||
.command('remove <id>')
|
||||
.alias('rm')
|
||||
.description('Remove a task')
|
||||
.action((id) => {
|
||||
const queue = getTaskQueue();
|
||||
if (queue.removeTask(id)) {
|
||||
console.log(chalk.green(`✓ Task removed: ${id}`));
|
||||
} else {
|
||||
console.log(chalk.red(`Task ${id} not found`));
|
||||
}
|
||||
});
|
||||
|
||||
taskCmd
|
||||
.command('clear')
|
||||
.description('Clear completed/failed tasks')
|
||||
.option('-a, --all', 'Clear all tasks')
|
||||
.option('-f, --failed', 'Clear only failed tasks')
|
||||
.action((options) => {
|
||||
const queue = getTaskQueue();
|
||||
let count: number;
|
||||
|
||||
if (options.all) {
|
||||
count = queue.clearAll();
|
||||
console.log(chalk.green(`✓ Cleared ${count} tasks`));
|
||||
} else if (options.failed) {
|
||||
count = queue.clearFailed();
|
||||
console.log(chalk.green(`✓ Cleared ${count} failed tasks`));
|
||||
} else {
|
||||
count = queue.clearCompleted();
|
||||
console.log(chalk.green(`✓ Cleared ${count} completed tasks`));
|
||||
}
|
||||
});
|
||||
|
||||
// ============ Ralph Loop Commands ============
|
||||
|
||||
const ralphCmd = program
|
||||
.command('ralph')
|
||||
.alias('r')
|
||||
.description('Control the Ralph autonomous loop');
|
||||
|
||||
ralphCmd
|
||||
.command('start')
|
||||
.description('Start the Ralph loop')
|
||||
.option('-m, --min-hours <hours>', 'Minimum duration in hours')
|
||||
.option('--no-auto-generate', 'Disable auto-generating follow-up tasks')
|
||||
.action(async (options) => {
|
||||
const loop = getRalphLoop({
|
||||
autoGenerateTasks: options.autoGenerate,
|
||||
});
|
||||
|
||||
if (options.minHours) {
|
||||
loop.setMinDuration(parseFloat(options.minHours));
|
||||
}
|
||||
|
||||
if (loop.isRunning()) {
|
||||
console.log(chalk.yellow('Ralph loop is already running'));
|
||||
return;
|
||||
}
|
||||
|
||||
loop.on('taskAssigned', (taskId, sessionId) => {
|
||||
console.log(chalk.cyan(`→ Task ${taskId.slice(0, 8)} assigned to session ${sessionId.slice(0, 8)}`));
|
||||
});
|
||||
|
||||
loop.on('taskCompleted', (taskId) => {
|
||||
console.log(chalk.green(`✓ Task ${taskId.slice(0, 8)} completed`));
|
||||
});
|
||||
|
||||
loop.on('taskFailed', (taskId, error) => {
|
||||
console.log(chalk.red(`✗ Task ${taskId.slice(0, 8)} failed: ${error}`));
|
||||
});
|
||||
|
||||
loop.on('stopped', () => {
|
||||
console.log(chalk.yellow('\nRalph loop stopped'));
|
||||
printStats(loop.getStats());
|
||||
process.exit(0);
|
||||
});
|
||||
|
||||
await loop.start();
|
||||
console.log(chalk.green('✓ Ralph loop started'));
|
||||
if (options.minHours) {
|
||||
console.log(` Minimum duration: ${options.minHours} hours`);
|
||||
}
|
||||
console.log(chalk.gray(' Press Ctrl+C to stop\n'));
|
||||
|
||||
// Keep process running
|
||||
process.on('SIGINT', () => {
|
||||
console.log(chalk.yellow('\nStopping Ralph loop...'));
|
||||
loop.stop();
|
||||
});
|
||||
});
|
||||
|
||||
ralphCmd
|
||||
.command('stop')
|
||||
.description('Stop the Ralph loop')
|
||||
.action(() => {
|
||||
const loop = getRalphLoop();
|
||||
if (!loop.isRunning()) {
|
||||
console.log(chalk.yellow('Ralph loop is not running'));
|
||||
return;
|
||||
}
|
||||
loop.stop();
|
||||
console.log(chalk.green('✓ Ralph loop stopped'));
|
||||
});
|
||||
|
||||
ralphCmd
|
||||
.command('status')
|
||||
.description('Show Ralph loop status')
|
||||
.action(() => {
|
||||
const loop = getRalphLoop();
|
||||
const stats = loop.getStats();
|
||||
printStats(stats);
|
||||
});
|
||||
|
||||
function printStats(stats: ReturnType<ReturnType<typeof getRalphLoop>['getStats']>) {
|
||||
const statusColor =
|
||||
stats.status === 'running' ? chalk.green :
|
||||
stats.status === 'paused' ? chalk.yellow :
|
||||
chalk.gray;
|
||||
|
||||
console.log(chalk.bold('\nRalph Loop Status:'));
|
||||
console.log(` Status: ${statusColor(stats.status)}`);
|
||||
console.log(` Elapsed: ${stats.elapsedHours.toFixed(2)} hours`);
|
||||
if (stats.minDurationMs) {
|
||||
const minHours = stats.minDurationMs / (1000 * 60 * 60);
|
||||
console.log(` Min Duration: ${minHours.toFixed(2)} hours (${stats.minDurationReached ? 'reached' : 'not reached'})`);
|
||||
}
|
||||
|
||||
console.log(chalk.bold('\nTasks:'));
|
||||
console.log(` Pending: ${stats.pending}`);
|
||||
console.log(` Running: ${stats.running}`);
|
||||
console.log(` Completed: ${stats.completed} (${stats.tasksCompleted} this session)`);
|
||||
console.log(` Failed: ${stats.failed}`);
|
||||
console.log(` Generated: ${stats.tasksGenerated}`);
|
||||
|
||||
console.log(chalk.bold('\nSessions:'));
|
||||
console.log(` Active: ${stats.activeSessions}`);
|
||||
console.log(` Idle: ${stats.idleSessions}`);
|
||||
console.log(` Busy: ${stats.busySessions}`);
|
||||
console.log('');
|
||||
}
|
||||
|
||||
// ============ Utility Commands ============
|
||||
|
||||
program
|
||||
.command('status')
|
||||
.description('Show overall status')
|
||||
.action(() => {
|
||||
const manager = getSessionManager();
|
||||
const queue = getTaskQueue();
|
||||
const loop = getRalphLoop();
|
||||
|
||||
const sessions = manager.getAllSessions();
|
||||
const taskCounts = queue.getCount();
|
||||
const loopStatus = loop.status;
|
||||
|
||||
console.log(chalk.bold('\nClaudeman Status'));
|
||||
console.log('─'.repeat(40));
|
||||
|
||||
console.log(chalk.bold('\nSessions:'));
|
||||
console.log(` Active: ${sessions.length}`);
|
||||
console.log(` Idle: ${sessions.filter((s) => s.isIdle()).length}`);
|
||||
console.log(` Busy: ${sessions.filter((s) => s.isBusy()).length}`);
|
||||
|
||||
console.log(chalk.bold('\nTasks:'));
|
||||
console.log(` Total: ${taskCounts.total}`);
|
||||
console.log(` Pending: ${taskCounts.pending}`);
|
||||
console.log(` Running: ${taskCounts.running}`);
|
||||
console.log(` Completed: ${taskCounts.completed}`);
|
||||
console.log(` Failed: ${taskCounts.failed}`);
|
||||
|
||||
const statusColor =
|
||||
loopStatus === 'running' ? chalk.green :
|
||||
loopStatus === 'paused' ? chalk.yellow :
|
||||
chalk.gray;
|
||||
console.log(chalk.bold('\nRalph Loop:'));
|
||||
console.log(` Status: ${statusColor(loopStatus)}`);
|
||||
console.log('');
|
||||
});
|
||||
|
||||
program
|
||||
.command('reset')
|
||||
.description('Reset all state')
|
||||
.option('-f, --force', 'Skip confirmation')
|
||||
.action(async (options) => {
|
||||
if (!options.force) {
|
||||
console.log(chalk.yellow('This will stop all sessions and clear all state.'));
|
||||
console.log(chalk.yellow('Use --force to confirm.'));
|
||||
return;
|
||||
}
|
||||
|
||||
const manager = getSessionManager();
|
||||
const store = getStore();
|
||||
|
||||
await manager.stopAllSessions();
|
||||
store.reset();
|
||||
|
||||
console.log(chalk.green('✓ All state reset'));
|
||||
});
|
||||
|
||||
// Shorthand commands at root level
|
||||
program
|
||||
.command('start')
|
||||
.description('Start a new session (shorthand)')
|
||||
.option('-d, --dir <path>', 'Working directory', process.cwd())
|
||||
.action(async (options) => {
|
||||
const manager = getSessionManager();
|
||||
const session = await manager.createSession(options.dir);
|
||||
console.log(chalk.green(`✓ Session started: ${session.id}`));
|
||||
});
|
||||
|
||||
program
|
||||
.command('list')
|
||||
.alias('ls')
|
||||
.description('List all sessions (shorthand)')
|
||||
.action(() => {
|
||||
const manager = getSessionManager();
|
||||
const sessions = manager.getAllSessions();
|
||||
const stored = manager.getStoredSessions();
|
||||
|
||||
if (sessions.length === 0 && Object.keys(stored).length === 0) {
|
||||
console.log(chalk.yellow('No sessions found'));
|
||||
return;
|
||||
}
|
||||
|
||||
console.log(chalk.bold('\nActive Sessions:'));
|
||||
if (sessions.length === 0) {
|
||||
console.log(' (none)');
|
||||
} else {
|
||||
for (const session of sessions) {
|
||||
const status = session.status === 'idle'
|
||||
? chalk.green('idle')
|
||||
: session.status === 'busy'
|
||||
? chalk.yellow('busy')
|
||||
: chalk.red(session.status);
|
||||
console.log(` ${chalk.cyan(session.id.slice(0, 8))} ${status} ${session.workingDir}`);
|
||||
}
|
||||
}
|
||||
console.log('');
|
||||
});
|
||||
|
||||
export { program };
|
||||
@@ -0,0 +1,17 @@
|
||||
#!/usr/bin/env node
|
||||
|
||||
import { program } from './cli.js';
|
||||
|
||||
// Handle uncaught errors
|
||||
process.on('uncaughtException', (err) => {
|
||||
console.error('Uncaught exception:', err.message);
|
||||
process.exit(1);
|
||||
});
|
||||
|
||||
process.on('unhandledRejection', (reason) => {
|
||||
console.error('Unhandled rejection:', reason);
|
||||
process.exit(1);
|
||||
});
|
||||
|
||||
// Run CLI
|
||||
program.parse();
|
||||
@@ -0,0 +1,397 @@
|
||||
import { EventEmitter } from 'node:events';
|
||||
import { getSessionManager, SessionManager } from './session-manager.js';
|
||||
import { getTaskQueue, TaskQueue } from './task-queue.js';
|
||||
import { getStore, StateStore } from './state-store.js';
|
||||
import { Session } from './session.js';
|
||||
import { Task } from './task.js';
|
||||
import { RalphLoopStatus } from './types.js';
|
||||
|
||||
export interface RalphLoopEvents {
|
||||
started: () => void;
|
||||
stopped: () => void;
|
||||
taskAssigned: (taskId: string, sessionId: string) => void;
|
||||
taskCompleted: (taskId: string) => void;
|
||||
taskFailed: (taskId: string, error: string) => void;
|
||||
error: (error: Error) => void;
|
||||
}
|
||||
|
||||
export interface RalphLoopOptions {
|
||||
pollIntervalMs?: number;
|
||||
minDurationMs?: number;
|
||||
autoGenerateTasks?: boolean;
|
||||
}
|
||||
|
||||
export class RalphLoop extends EventEmitter {
|
||||
private sessionManager: SessionManager;
|
||||
private taskQueue: TaskQueue;
|
||||
private store: StateStore;
|
||||
private pollIntervalMs: number;
|
||||
private minDurationMs: number | null;
|
||||
private autoGenerateTasks: boolean;
|
||||
private loopTimer: NodeJS.Timeout | null = null;
|
||||
private _status: RalphLoopStatus = 'stopped';
|
||||
private startedAt: number | null = null;
|
||||
private tasksCompleted: number = 0;
|
||||
private tasksGenerated: number = 0;
|
||||
|
||||
constructor(options: RalphLoopOptions = {}) {
|
||||
super();
|
||||
this.sessionManager = getSessionManager();
|
||||
this.taskQueue = getTaskQueue();
|
||||
this.store = getStore();
|
||||
|
||||
const config = this.store.getConfig();
|
||||
this.pollIntervalMs = options.pollIntervalMs ?? config.pollIntervalMs;
|
||||
this.minDurationMs = options.minDurationMs ?? null;
|
||||
this.autoGenerateTasks = options.autoGenerateTasks ?? true;
|
||||
|
||||
// Load state from store
|
||||
const savedState = this.store.getRalphLoopState();
|
||||
if (savedState.status === 'running') {
|
||||
// If we crashed while running, reset to stopped
|
||||
this._status = 'stopped';
|
||||
this.store.setRalphLoopState({ status: 'stopped' });
|
||||
}
|
||||
|
||||
this.setupEventHandlers();
|
||||
}
|
||||
|
||||
private setupEventHandlers(): void {
|
||||
this.sessionManager.on('sessionCompletion', (sessionId: string, phrase: string) => {
|
||||
this.handleSessionCompletion(sessionId, phrase);
|
||||
});
|
||||
|
||||
this.sessionManager.on('sessionError', (sessionId: string, error: string) => {
|
||||
this.handleSessionError(sessionId, error);
|
||||
});
|
||||
|
||||
this.sessionManager.on('sessionStopped', (sessionId: string) => {
|
||||
this.handleSessionStopped(sessionId);
|
||||
});
|
||||
}
|
||||
|
||||
get status(): RalphLoopStatus {
|
||||
return this._status;
|
||||
}
|
||||
|
||||
isRunning(): boolean {
|
||||
return this._status === 'running';
|
||||
}
|
||||
|
||||
getElapsedMs(): number {
|
||||
if (!this.startedAt) {
|
||||
return 0;
|
||||
}
|
||||
return Date.now() - this.startedAt;
|
||||
}
|
||||
|
||||
getElapsedHours(): number {
|
||||
return this.getElapsedMs() / (1000 * 60 * 60);
|
||||
}
|
||||
|
||||
isMinDurationReached(): boolean {
|
||||
if (!this.minDurationMs) {
|
||||
return true;
|
||||
}
|
||||
return this.getElapsedMs() >= this.minDurationMs;
|
||||
}
|
||||
|
||||
getStats() {
|
||||
const taskCounts = this.taskQueue.getCount();
|
||||
return {
|
||||
status: this._status,
|
||||
elapsedMs: this.getElapsedMs(),
|
||||
elapsedHours: this.getElapsedHours(),
|
||||
minDurationMs: this.minDurationMs,
|
||||
minDurationReached: this.isMinDurationReached(),
|
||||
tasksCompleted: this.tasksCompleted,
|
||||
tasksGenerated: this.tasksGenerated,
|
||||
...taskCounts,
|
||||
activeSessions: this.sessionManager.getSessionCount(),
|
||||
idleSessions: this.sessionManager.getIdleSessions().length,
|
||||
busySessions: this.sessionManager.getBusySessions().length,
|
||||
};
|
||||
}
|
||||
|
||||
async start(): Promise<void> {
|
||||
if (this._status === 'running') {
|
||||
return;
|
||||
}
|
||||
|
||||
this._status = 'running';
|
||||
this.startedAt = Date.now();
|
||||
this.tasksCompleted = 0;
|
||||
this.tasksGenerated = 0;
|
||||
|
||||
this.store.setRalphLoopState({
|
||||
status: 'running',
|
||||
startedAt: this.startedAt,
|
||||
minDurationMs: this.minDurationMs,
|
||||
tasksCompleted: 0,
|
||||
tasksGenerated: 0,
|
||||
});
|
||||
|
||||
this.emit('started');
|
||||
this.runLoop();
|
||||
}
|
||||
|
||||
stop(): void {
|
||||
if (this._status === 'stopped') {
|
||||
return;
|
||||
}
|
||||
|
||||
this._status = 'stopped';
|
||||
|
||||
if (this.loopTimer) {
|
||||
clearTimeout(this.loopTimer);
|
||||
this.loopTimer = null;
|
||||
}
|
||||
|
||||
this.store.setRalphLoopState({
|
||||
status: 'stopped',
|
||||
lastCheckAt: Date.now(),
|
||||
});
|
||||
|
||||
this.emit('stopped');
|
||||
}
|
||||
|
||||
pause(): void {
|
||||
if (this._status !== 'running') {
|
||||
return;
|
||||
}
|
||||
|
||||
this._status = 'paused';
|
||||
|
||||
if (this.loopTimer) {
|
||||
clearTimeout(this.loopTimer);
|
||||
this.loopTimer = null;
|
||||
}
|
||||
|
||||
this.store.setRalphLoopState({ status: 'paused' });
|
||||
}
|
||||
|
||||
resume(): void {
|
||||
if (this._status !== 'paused') {
|
||||
return;
|
||||
}
|
||||
|
||||
this._status = 'running';
|
||||
this.store.setRalphLoopState({ status: 'running' });
|
||||
this.runLoop();
|
||||
}
|
||||
|
||||
private runLoop(): void {
|
||||
if (this._status !== 'running') {
|
||||
return;
|
||||
}
|
||||
|
||||
this.tick()
|
||||
.catch((err) => {
|
||||
this.emit('error', err);
|
||||
})
|
||||
.finally(() => {
|
||||
if (this._status === 'running') {
|
||||
this.loopTimer = setTimeout(() => this.runLoop(), this.pollIntervalMs);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
private async tick(): Promise<void> {
|
||||
this.store.setRalphLoopState({ lastCheckAt: Date.now() });
|
||||
|
||||
// Check for timed out tasks
|
||||
await this.checkTimeouts();
|
||||
|
||||
// Assign tasks to idle sessions
|
||||
await this.assignTasks();
|
||||
|
||||
// Check if we should auto-generate tasks
|
||||
if (this.autoGenerateTasks && this.shouldGenerateTasks()) {
|
||||
await this.generateFollowUpTasks();
|
||||
}
|
||||
|
||||
// Check if we're done
|
||||
if (this.shouldStop()) {
|
||||
this.stop();
|
||||
}
|
||||
}
|
||||
|
||||
private async assignTasks(): Promise<void> {
|
||||
const idleSessions = this.sessionManager.getIdleSessions();
|
||||
|
||||
for (const session of idleSessions) {
|
||||
const task = this.taskQueue.next();
|
||||
if (!task) {
|
||||
break;
|
||||
}
|
||||
|
||||
await this.assignTaskToSession(task, session);
|
||||
}
|
||||
}
|
||||
|
||||
private async assignTaskToSession(task: Task, session: Session): Promise<void> {
|
||||
try {
|
||||
task.assign(session.id);
|
||||
session.assignTask(task.id);
|
||||
this.taskQueue.updateTask(task);
|
||||
|
||||
// Send the prompt to the session
|
||||
await session.sendInput(task.prompt);
|
||||
|
||||
this.emit('taskAssigned', task.id, session.id);
|
||||
} catch (err) {
|
||||
task.fail((err as Error).message);
|
||||
session.clearTask();
|
||||
this.taskQueue.updateTask(task);
|
||||
this.emit('taskFailed', task.id, (err as Error).message);
|
||||
}
|
||||
}
|
||||
|
||||
private async checkTimeouts(): Promise<void> {
|
||||
for (const task of this.taskQueue.getRunningTasks()) {
|
||||
if (task.isTimedOut()) {
|
||||
task.fail('Task timed out');
|
||||
this.taskQueue.updateTask(task);
|
||||
|
||||
const session = this.sessionManager.getSession(task.assignedSessionId!);
|
||||
if (session) {
|
||||
session.clearTask();
|
||||
}
|
||||
|
||||
this.emit('taskFailed', task.id, 'Task timed out');
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private handleSessionCompletion(sessionId: string, phrase: string): void {
|
||||
const session = this.sessionManager.getSession(sessionId);
|
||||
if (!session) {
|
||||
return;
|
||||
}
|
||||
|
||||
const taskId = session.currentTaskId;
|
||||
if (!taskId) {
|
||||
return;
|
||||
}
|
||||
|
||||
const task = this.taskQueue.getTask(taskId);
|
||||
if (!task) {
|
||||
return;
|
||||
}
|
||||
|
||||
// Append output and check for completion
|
||||
task.appendOutput(session.getOutput());
|
||||
|
||||
if (task.checkCompletion(session.getOutput()) || phrase) {
|
||||
task.complete();
|
||||
this.taskQueue.updateTask(task);
|
||||
session.clearTask();
|
||||
this.tasksCompleted++;
|
||||
this.store.setRalphLoopState({ tasksCompleted: this.tasksCompleted });
|
||||
this.emit('taskCompleted', task.id);
|
||||
}
|
||||
}
|
||||
|
||||
private handleSessionError(sessionId: string, error: string): void {
|
||||
const session = this.sessionManager.getSession(sessionId);
|
||||
if (!session) {
|
||||
return;
|
||||
}
|
||||
|
||||
const taskId = session.currentTaskId;
|
||||
if (!taskId) {
|
||||
return;
|
||||
}
|
||||
|
||||
const task = this.taskQueue.getTask(taskId);
|
||||
if (!task) {
|
||||
return;
|
||||
}
|
||||
|
||||
task.setError(error);
|
||||
// Don't fail the task immediately on stderr - some tools write to stderr normally
|
||||
}
|
||||
|
||||
private handleSessionStopped(sessionId: string): void {
|
||||
const task = this.taskQueue.getRunningTaskForSession(sessionId);
|
||||
if (task) {
|
||||
task.fail('Session stopped unexpectedly');
|
||||
this.taskQueue.updateTask(task);
|
||||
this.emit('taskFailed', task.id, 'Session stopped unexpectedly');
|
||||
}
|
||||
}
|
||||
|
||||
private shouldGenerateTasks(): boolean {
|
||||
// Generate tasks if:
|
||||
// 1. No pending tasks
|
||||
// 2. Min duration not reached
|
||||
// 3. We have idle sessions
|
||||
const counts = this.taskQueue.getCount();
|
||||
return (
|
||||
counts.pending === 0 &&
|
||||
!this.isMinDurationReached() &&
|
||||
this.sessionManager.getIdleSessions().length > 0
|
||||
);
|
||||
}
|
||||
|
||||
private async generateFollowUpTasks(): Promise<void> {
|
||||
// This is a placeholder for auto-generating follow-up tasks
|
||||
// In a real implementation, this could:
|
||||
// - Analyze completed tasks to find optimization opportunities
|
||||
// - Generate tasks for code cleanup, tests, documentation
|
||||
// - Use Claude to suggest improvements
|
||||
|
||||
const suggestions = [
|
||||
'Review and optimize recently changed code',
|
||||
'Add tests for uncovered code paths',
|
||||
'Update documentation for changed APIs',
|
||||
'Check for security vulnerabilities',
|
||||
'Run linting and fix any issues',
|
||||
];
|
||||
|
||||
// Only generate one task at a time
|
||||
const suggestion = suggestions[this.tasksGenerated % suggestions.length];
|
||||
const defaultDir = process.cwd();
|
||||
|
||||
this.taskQueue.addTask({
|
||||
prompt: suggestion,
|
||||
workingDir: defaultDir,
|
||||
priority: -1, // Lower priority than user-added tasks
|
||||
});
|
||||
|
||||
this.tasksGenerated++;
|
||||
this.store.setRalphLoopState({ tasksGenerated: this.tasksGenerated });
|
||||
}
|
||||
|
||||
private shouldStop(): boolean {
|
||||
const counts = this.taskQueue.getCount();
|
||||
|
||||
// Don't stop if there are pending or running tasks
|
||||
if (counts.pending > 0 || counts.running > 0) {
|
||||
return false;
|
||||
}
|
||||
|
||||
// Don't stop if min duration not reached and auto-generate is on
|
||||
if (!this.isMinDurationReached() && this.autoGenerateTasks) {
|
||||
return false;
|
||||
}
|
||||
|
||||
// All tasks done and conditions met
|
||||
return true;
|
||||
}
|
||||
|
||||
setMinDuration(hours: number): void {
|
||||
this.minDurationMs = hours * 60 * 60 * 1000;
|
||||
this.store.setRalphLoopState({ minDurationMs: this.minDurationMs });
|
||||
}
|
||||
}
|
||||
|
||||
// Singleton instance
|
||||
let loopInstance: RalphLoop | null = null;
|
||||
|
||||
export function getRalphLoop(options?: RalphLoopOptions): RalphLoop {
|
||||
if (!loopInstance) {
|
||||
loopInstance = new RalphLoop(options);
|
||||
}
|
||||
return loopInstance;
|
||||
}
|
||||
@@ -0,0 +1,158 @@
|
||||
import { EventEmitter } from 'node:events';
|
||||
import { Session } from './session.js';
|
||||
import { getStore } from './state-store.js';
|
||||
import { SessionState } from './types.js';
|
||||
|
||||
export interface SessionManagerEvents {
|
||||
sessionStarted: (session: Session) => void;
|
||||
sessionStopped: (sessionId: string) => void;
|
||||
sessionError: (sessionId: string, error: string) => void;
|
||||
sessionOutput: (sessionId: string, output: string) => void;
|
||||
sessionCompletion: (sessionId: string, phrase: string) => void;
|
||||
}
|
||||
|
||||
export class SessionManager extends EventEmitter {
|
||||
private sessions: Map<string, Session> = new Map();
|
||||
private store = getStore();
|
||||
|
||||
constructor() {
|
||||
super();
|
||||
this.loadFromStore();
|
||||
}
|
||||
|
||||
private loadFromStore(): void {
|
||||
const storedSessions = this.store.getSessions();
|
||||
// Note: We don't restore actual processes, just the state
|
||||
// Dead sessions are marked as stopped
|
||||
for (const [id, state] of Object.entries(storedSessions)) {
|
||||
if (state.status !== 'stopped') {
|
||||
state.status = 'stopped';
|
||||
state.pid = null;
|
||||
this.store.setSession(id, state);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async createSession(workingDir: string): Promise<Session> {
|
||||
const config = this.store.getConfig();
|
||||
|
||||
if (this.sessions.size >= config.maxConcurrentSessions) {
|
||||
throw new Error(`Maximum concurrent sessions (${config.maxConcurrentSessions}) reached`);
|
||||
}
|
||||
|
||||
const session = new Session({ workingDir });
|
||||
|
||||
// Set up event forwarding
|
||||
session.on('output', (data) => {
|
||||
this.emit('sessionOutput', session.id, data);
|
||||
this.updateSessionState(session);
|
||||
});
|
||||
|
||||
session.on('error', (data) => {
|
||||
this.emit('sessionError', session.id, data);
|
||||
this.updateSessionState(session);
|
||||
});
|
||||
|
||||
session.on('completion', (phrase) => {
|
||||
this.emit('sessionCompletion', session.id, phrase);
|
||||
});
|
||||
|
||||
session.on('exit', () => {
|
||||
this.emit('sessionStopped', session.id);
|
||||
this.updateSessionState(session);
|
||||
});
|
||||
|
||||
await session.start();
|
||||
|
||||
this.sessions.set(session.id, session);
|
||||
this.store.setSession(session.id, session.toState());
|
||||
|
||||
this.emit('sessionStarted', session);
|
||||
return session;
|
||||
}
|
||||
|
||||
async stopSession(id: string): Promise<void> {
|
||||
const session = this.sessions.get(id);
|
||||
if (!session) {
|
||||
// Update store to mark as stopped if it exists there
|
||||
const storedSession = this.store.getSession(id);
|
||||
if (storedSession) {
|
||||
storedSession.status = 'stopped';
|
||||
storedSession.pid = null;
|
||||
this.store.setSession(id, storedSession);
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
await session.stop();
|
||||
this.sessions.delete(id);
|
||||
this.updateSessionState(session);
|
||||
}
|
||||
|
||||
async stopAllSessions(): Promise<void> {
|
||||
const stopPromises = Array.from(this.sessions.keys()).map((id) =>
|
||||
this.stopSession(id)
|
||||
);
|
||||
await Promise.all(stopPromises);
|
||||
}
|
||||
|
||||
getSession(id: string): Session | undefined {
|
||||
return this.sessions.get(id);
|
||||
}
|
||||
|
||||
getAllSessions(): Session[] {
|
||||
return Array.from(this.sessions.values());
|
||||
}
|
||||
|
||||
getIdleSessions(): Session[] {
|
||||
return this.getAllSessions().filter((s) => s.isIdle());
|
||||
}
|
||||
|
||||
getBusySessions(): Session[] {
|
||||
return this.getAllSessions().filter((s) => s.isBusy());
|
||||
}
|
||||
|
||||
getSessionCount(): number {
|
||||
return this.sessions.size;
|
||||
}
|
||||
|
||||
hasSession(id: string): boolean {
|
||||
return this.sessions.has(id);
|
||||
}
|
||||
|
||||
private updateSessionState(session: Session): void {
|
||||
this.store.setSession(session.id, session.toState());
|
||||
}
|
||||
|
||||
getStoredSessions(): Record<string, SessionState> {
|
||||
return this.store.getSessions();
|
||||
}
|
||||
|
||||
async sendToSession(sessionId: string, input: string): Promise<void> {
|
||||
const session = this.sessions.get(sessionId);
|
||||
if (!session) {
|
||||
throw new Error(`Session ${sessionId} not found`);
|
||||
}
|
||||
await session.sendInput(input);
|
||||
}
|
||||
|
||||
getSessionOutput(sessionId: string): string | null {
|
||||
const session = this.sessions.get(sessionId);
|
||||
return session?.getOutput() ?? null;
|
||||
}
|
||||
|
||||
getSessionError(sessionId: string): string | null {
|
||||
const session = this.sessions.get(sessionId);
|
||||
return session?.getError() ?? null;
|
||||
}
|
||||
}
|
||||
|
||||
// Singleton instance
|
||||
let managerInstance: SessionManager | null = null;
|
||||
|
||||
export function getSessionManager(): SessionManager {
|
||||
if (!managerInstance) {
|
||||
managerInstance = new SessionManager();
|
||||
}
|
||||
return managerInstance;
|
||||
}
|
||||
+235
@@ -0,0 +1,235 @@
|
||||
import { spawn, ChildProcess } from 'node:child_process';
|
||||
import { EventEmitter } from 'node:events';
|
||||
import { v4 as uuidv4 } from 'uuid';
|
||||
import { SessionState, SessionStatus, SessionConfig } from './types.js';
|
||||
|
||||
export interface SessionEvents {
|
||||
output: (data: string) => void;
|
||||
error: (data: string) => void;
|
||||
exit: (code: number | null) => void;
|
||||
completion: (phrase: string) => void;
|
||||
}
|
||||
|
||||
export class Session extends EventEmitter {
|
||||
readonly id: string;
|
||||
readonly workingDir: string;
|
||||
readonly createdAt: number;
|
||||
|
||||
private process: ChildProcess | null = null;
|
||||
private _status: SessionStatus = 'idle';
|
||||
private _currentTaskId: string | null = null;
|
||||
private _outputBuffer: string = '';
|
||||
private _errorBuffer: string = '';
|
||||
private _lastActivityAt: number;
|
||||
|
||||
constructor(config: Partial<SessionConfig> & { workingDir: string }) {
|
||||
super();
|
||||
this.id = config.id || uuidv4();
|
||||
this.workingDir = config.workingDir;
|
||||
this.createdAt = config.createdAt || Date.now();
|
||||
this._lastActivityAt = this.createdAt;
|
||||
}
|
||||
|
||||
get status(): SessionStatus {
|
||||
return this._status;
|
||||
}
|
||||
|
||||
get currentTaskId(): string | null {
|
||||
return this._currentTaskId;
|
||||
}
|
||||
|
||||
get pid(): number | null {
|
||||
return this.process?.pid ?? null;
|
||||
}
|
||||
|
||||
get outputBuffer(): string {
|
||||
return this._outputBuffer;
|
||||
}
|
||||
|
||||
get errorBuffer(): string {
|
||||
return this._errorBuffer;
|
||||
}
|
||||
|
||||
get lastActivityAt(): number {
|
||||
return this._lastActivityAt;
|
||||
}
|
||||
|
||||
isIdle(): boolean {
|
||||
return this._status === 'idle';
|
||||
}
|
||||
|
||||
isBusy(): boolean {
|
||||
return this._status === 'busy';
|
||||
}
|
||||
|
||||
isRunning(): boolean {
|
||||
return this._status === 'idle' || this._status === 'busy';
|
||||
}
|
||||
|
||||
toState(): SessionState {
|
||||
return {
|
||||
id: this.id,
|
||||
pid: this.pid,
|
||||
status: this._status,
|
||||
workingDir: this.workingDir,
|
||||
currentTaskId: this._currentTaskId,
|
||||
createdAt: this.createdAt,
|
||||
lastActivityAt: this._lastActivityAt,
|
||||
};
|
||||
}
|
||||
|
||||
start(): Promise<void> {
|
||||
return new Promise((resolve, reject) => {
|
||||
if (this.process) {
|
||||
reject(new Error('Session already started'));
|
||||
return;
|
||||
}
|
||||
|
||||
try {
|
||||
// Spawn claude CLI with --print flag for non-interactive mode
|
||||
this.process = spawn('claude', ['--print'], {
|
||||
cwd: this.workingDir,
|
||||
stdio: ['pipe', 'pipe', 'pipe'],
|
||||
env: { ...process.env },
|
||||
});
|
||||
|
||||
this._status = 'idle';
|
||||
this._lastActivityAt = Date.now();
|
||||
|
||||
this.process.stdout?.on('data', (data: Buffer) => {
|
||||
const text = data.toString();
|
||||
this._outputBuffer += text;
|
||||
this._lastActivityAt = Date.now();
|
||||
this.emit('output', text);
|
||||
this.checkForCompletion(text);
|
||||
});
|
||||
|
||||
this.process.stderr?.on('data', (data: Buffer) => {
|
||||
const text = data.toString();
|
||||
this._errorBuffer += text;
|
||||
this._lastActivityAt = Date.now();
|
||||
this.emit('error', text);
|
||||
});
|
||||
|
||||
this.process.on('error', (err) => {
|
||||
this._status = 'error';
|
||||
this.emit('error', err.message);
|
||||
reject(err);
|
||||
});
|
||||
|
||||
this.process.on('exit', (code) => {
|
||||
this._status = 'stopped';
|
||||
this.process = null;
|
||||
this.emit('exit', code);
|
||||
});
|
||||
|
||||
// Give process time to start
|
||||
setTimeout(() => {
|
||||
if (this.process && !this.process.killed) {
|
||||
resolve();
|
||||
}
|
||||
}, 100);
|
||||
} catch (err) {
|
||||
this._status = 'error';
|
||||
reject(err);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
stop(): Promise<void> {
|
||||
return new Promise((resolve) => {
|
||||
if (!this.process) {
|
||||
this._status = 'stopped';
|
||||
resolve();
|
||||
return;
|
||||
}
|
||||
|
||||
const cleanup = () => {
|
||||
this.process = null;
|
||||
this._status = 'stopped';
|
||||
this._currentTaskId = null;
|
||||
resolve();
|
||||
};
|
||||
|
||||
this.process.once('exit', cleanup);
|
||||
|
||||
// Try graceful shutdown first
|
||||
this.process.kill('SIGTERM');
|
||||
|
||||
// Force kill after timeout
|
||||
setTimeout(() => {
|
||||
if (this.process && !this.process.killed) {
|
||||
this.process.kill('SIGKILL');
|
||||
}
|
||||
}, 5000);
|
||||
});
|
||||
}
|
||||
|
||||
async sendInput(input: string): Promise<void> {
|
||||
if (!this.process || !this.process.stdin) {
|
||||
throw new Error('Session not started or stdin not available');
|
||||
}
|
||||
|
||||
this._status = 'busy';
|
||||
this._lastActivityAt = Date.now();
|
||||
|
||||
return new Promise((resolve, reject) => {
|
||||
this.process!.stdin!.write(input + '\n', (err) => {
|
||||
if (err) {
|
||||
reject(err);
|
||||
} else {
|
||||
resolve();
|
||||
}
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
assignTask(taskId: string): void {
|
||||
this._currentTaskId = taskId;
|
||||
this._status = 'busy';
|
||||
this._outputBuffer = '';
|
||||
this._errorBuffer = '';
|
||||
this._lastActivityAt = Date.now();
|
||||
}
|
||||
|
||||
clearTask(): void {
|
||||
this._currentTaskId = null;
|
||||
this._status = 'idle';
|
||||
this._lastActivityAt = Date.now();
|
||||
}
|
||||
|
||||
getOutput(): string {
|
||||
return this._outputBuffer;
|
||||
}
|
||||
|
||||
getError(): string {
|
||||
return this._errorBuffer;
|
||||
}
|
||||
|
||||
clearBuffers(): void {
|
||||
this._outputBuffer = '';
|
||||
this._errorBuffer = '';
|
||||
}
|
||||
|
||||
private checkForCompletion(text: string): void {
|
||||
// Check for completion phrase pattern: <promise>...</promise>
|
||||
const promiseMatch = text.match(/<promise>([^<]+)<\/promise>/);
|
||||
if (promiseMatch) {
|
||||
this.emit('completion', promiseMatch[1]);
|
||||
}
|
||||
|
||||
// Also check for common completion indicators
|
||||
const completionIndicators = [
|
||||
/Task completed successfully/i,
|
||||
/All tasks done/i,
|
||||
/✓ Complete/i,
|
||||
];
|
||||
|
||||
for (const pattern of completionIndicators) {
|
||||
if (pattern.test(text)) {
|
||||
this.emit('completion', 'auto-detected');
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,123 @@
|
||||
import { readFileSync, writeFileSync, existsSync, mkdirSync } from 'node:fs';
|
||||
import { homedir } from 'node:os';
|
||||
import { dirname, join } from 'node:path';
|
||||
import { AppState, createInitialState } from './types.js';
|
||||
|
||||
export class StateStore {
|
||||
private state: AppState;
|
||||
private filePath: string;
|
||||
|
||||
constructor(filePath?: string) {
|
||||
this.filePath = filePath || join(homedir(), '.claudeman', 'state.json');
|
||||
this.state = this.load();
|
||||
this.state.config.stateFilePath = this.filePath;
|
||||
}
|
||||
|
||||
private ensureDir(): void {
|
||||
const dir = dirname(this.filePath);
|
||||
if (!existsSync(dir)) {
|
||||
mkdirSync(dir, { recursive: true });
|
||||
}
|
||||
}
|
||||
|
||||
private load(): AppState {
|
||||
try {
|
||||
if (existsSync(this.filePath)) {
|
||||
const data = readFileSync(this.filePath, 'utf-8');
|
||||
const parsed = JSON.parse(data) as Partial<AppState>;
|
||||
// Merge with initial state to ensure all fields exist
|
||||
const initial = createInitialState();
|
||||
return {
|
||||
...initial,
|
||||
...parsed,
|
||||
sessions: { ...parsed.sessions },
|
||||
tasks: { ...parsed.tasks },
|
||||
ralphLoop: { ...initial.ralphLoop, ...parsed.ralphLoop },
|
||||
config: { ...initial.config, ...parsed.config },
|
||||
};
|
||||
}
|
||||
} catch (err) {
|
||||
console.error('Failed to load state, using initial state:', err);
|
||||
}
|
||||
return createInitialState();
|
||||
}
|
||||
|
||||
save(): void {
|
||||
this.ensureDir();
|
||||
writeFileSync(this.filePath, JSON.stringify(this.state, null, 2), 'utf-8');
|
||||
}
|
||||
|
||||
getState(): AppState {
|
||||
return this.state;
|
||||
}
|
||||
|
||||
getSessions() {
|
||||
return this.state.sessions;
|
||||
}
|
||||
|
||||
getSession(id: string) {
|
||||
return this.state.sessions[id] || null;
|
||||
}
|
||||
|
||||
setSession(id: string, session: AppState['sessions'][string]) {
|
||||
this.state.sessions[id] = session;
|
||||
this.save();
|
||||
}
|
||||
|
||||
removeSession(id: string) {
|
||||
delete this.state.sessions[id];
|
||||
this.save();
|
||||
}
|
||||
|
||||
getTasks() {
|
||||
return this.state.tasks;
|
||||
}
|
||||
|
||||
getTask(id: string) {
|
||||
return this.state.tasks[id] || null;
|
||||
}
|
||||
|
||||
setTask(id: string, task: AppState['tasks'][string]) {
|
||||
this.state.tasks[id] = task;
|
||||
this.save();
|
||||
}
|
||||
|
||||
removeTask(id: string) {
|
||||
delete this.state.tasks[id];
|
||||
this.save();
|
||||
}
|
||||
|
||||
getRalphLoopState() {
|
||||
return this.state.ralphLoop;
|
||||
}
|
||||
|
||||
setRalphLoopState(ralphLoop: Partial<AppState['ralphLoop']>) {
|
||||
this.state.ralphLoop = { ...this.state.ralphLoop, ...ralphLoop };
|
||||
this.save();
|
||||
}
|
||||
|
||||
getConfig() {
|
||||
return this.state.config;
|
||||
}
|
||||
|
||||
setConfig(config: Partial<AppState['config']>) {
|
||||
this.state.config = { ...this.state.config, ...config };
|
||||
this.save();
|
||||
}
|
||||
|
||||
reset(): void {
|
||||
this.state = createInitialState();
|
||||
this.state.config.stateFilePath = this.filePath;
|
||||
this.save();
|
||||
}
|
||||
}
|
||||
|
||||
// Singleton instance
|
||||
let storeInstance: StateStore | null = null;
|
||||
|
||||
export function getStore(filePath?: string): StateStore {
|
||||
if (!storeInstance) {
|
||||
storeInstance = new StateStore(filePath);
|
||||
}
|
||||
return storeInstance;
|
||||
}
|
||||
@@ -0,0 +1,179 @@
|
||||
import { EventEmitter } from 'node:events';
|
||||
import { Task, CreateTaskOptions } from './task.js';
|
||||
import { getStore } from './state-store.js';
|
||||
import { TaskState } from './types.js';
|
||||
|
||||
export interface TaskQueueEvents {
|
||||
taskAdded: (task: Task) => void;
|
||||
taskRemoved: (taskId: string) => void;
|
||||
taskUpdated: (task: Task) => void;
|
||||
}
|
||||
|
||||
export class TaskQueue extends EventEmitter {
|
||||
private tasks: Map<string, Task> = new Map();
|
||||
private store = getStore();
|
||||
|
||||
constructor() {
|
||||
super();
|
||||
this.loadFromStore();
|
||||
}
|
||||
|
||||
private loadFromStore(): void {
|
||||
const storedTasks = this.store.getTasks();
|
||||
for (const [id, state] of Object.entries(storedTasks)) {
|
||||
const task = Task.fromState(state);
|
||||
this.tasks.set(id, task);
|
||||
}
|
||||
}
|
||||
|
||||
addTask(options: CreateTaskOptions): Task {
|
||||
const task = new Task(options);
|
||||
this.tasks.set(task.id, task);
|
||||
this.store.setTask(task.id, task.toState());
|
||||
this.emit('taskAdded', task);
|
||||
return task;
|
||||
}
|
||||
|
||||
getTask(id: string): Task | undefined {
|
||||
return this.tasks.get(id);
|
||||
}
|
||||
|
||||
removeTask(id: string): boolean {
|
||||
const removed = this.tasks.delete(id);
|
||||
if (removed) {
|
||||
this.store.removeTask(id);
|
||||
this.emit('taskRemoved', id);
|
||||
}
|
||||
return removed;
|
||||
}
|
||||
|
||||
updateTask(task: Task): void {
|
||||
this.tasks.set(task.id, task);
|
||||
this.store.setTask(task.id, task.toState());
|
||||
this.emit('taskUpdated', task);
|
||||
}
|
||||
|
||||
getAllTasks(): Task[] {
|
||||
return Array.from(this.tasks.values());
|
||||
}
|
||||
|
||||
getPendingTasks(): Task[] {
|
||||
return this.getAllTasks()
|
||||
.filter((t) => t.isPending())
|
||||
.sort((a, b) => {
|
||||
// Sort by priority (higher first), then by creation time (older first)
|
||||
if (a.priority !== b.priority) {
|
||||
return b.priority - a.priority;
|
||||
}
|
||||
return a.createdAt - b.createdAt;
|
||||
});
|
||||
}
|
||||
|
||||
getRunningTasks(): Task[] {
|
||||
return this.getAllTasks().filter((t) => t.isRunning());
|
||||
}
|
||||
|
||||
getCompletedTasks(): Task[] {
|
||||
return this.getAllTasks().filter((t) => t.isCompleted());
|
||||
}
|
||||
|
||||
getFailedTasks(): Task[] {
|
||||
return this.getAllTasks().filter((t) => t.isFailed());
|
||||
}
|
||||
|
||||
hasNext(): boolean {
|
||||
return this.getNextAvailable() !== null;
|
||||
}
|
||||
|
||||
getNextAvailable(): Task | null {
|
||||
const pending = this.getPendingTasks();
|
||||
|
||||
for (const task of pending) {
|
||||
if (this.areDependenciesSatisfied(task)) {
|
||||
return task;
|
||||
}
|
||||
}
|
||||
|
||||
return null;
|
||||
}
|
||||
|
||||
next(): Task | null {
|
||||
return this.getNextAvailable();
|
||||
}
|
||||
|
||||
private areDependenciesSatisfied(task: Task): boolean {
|
||||
for (const depId of task.dependencies) {
|
||||
const dep = this.tasks.get(depId);
|
||||
if (!dep || !dep.isCompleted()) {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
getTasksBySession(sessionId: string): Task[] {
|
||||
return this.getAllTasks().filter((t) => t.assignedSessionId === sessionId);
|
||||
}
|
||||
|
||||
getRunningTaskForSession(sessionId: string): Task | null {
|
||||
return this.getAllTasks().find(
|
||||
(t) => t.isRunning() && t.assignedSessionId === sessionId
|
||||
) || null;
|
||||
}
|
||||
|
||||
getCount(): { total: number; pending: number; running: number; completed: number; failed: number } {
|
||||
const tasks = this.getAllTasks();
|
||||
return {
|
||||
total: tasks.length,
|
||||
pending: tasks.filter((t) => t.isPending()).length,
|
||||
running: tasks.filter((t) => t.isRunning()).length,
|
||||
completed: tasks.filter((t) => t.isCompleted()).length,
|
||||
failed: tasks.filter((t) => t.isFailed()).length,
|
||||
};
|
||||
}
|
||||
|
||||
clearCompleted(): number {
|
||||
let count = 0;
|
||||
for (const task of this.getAllTasks()) {
|
||||
if (task.isCompleted()) {
|
||||
this.removeTask(task.id);
|
||||
count++;
|
||||
}
|
||||
}
|
||||
return count;
|
||||
}
|
||||
|
||||
clearFailed(): number {
|
||||
let count = 0;
|
||||
for (const task of this.getAllTasks()) {
|
||||
if (task.isFailed()) {
|
||||
this.removeTask(task.id);
|
||||
count++;
|
||||
}
|
||||
}
|
||||
return count;
|
||||
}
|
||||
|
||||
clearAll(): number {
|
||||
const count = this.tasks.size;
|
||||
for (const id of this.tasks.keys()) {
|
||||
this.store.removeTask(id);
|
||||
}
|
||||
this.tasks.clear();
|
||||
return count;
|
||||
}
|
||||
|
||||
getStoredTasks(): Record<string, TaskState> {
|
||||
return this.store.getTasks();
|
||||
}
|
||||
}
|
||||
|
||||
// Singleton instance
|
||||
let queueInstance: TaskQueue | null = null;
|
||||
|
||||
export function getTaskQueue(): TaskQueue {
|
||||
if (!queueInstance) {
|
||||
queueInstance = new TaskQueue();
|
||||
}
|
||||
return queueInstance;
|
||||
}
|
||||
+186
@@ -0,0 +1,186 @@
|
||||
import { v4 as uuidv4 } from 'uuid';
|
||||
import { TaskDefinition, TaskState, TaskStatus } from './types.js';
|
||||
|
||||
export interface CreateTaskOptions {
|
||||
prompt: string;
|
||||
workingDir?: string;
|
||||
priority?: number;
|
||||
dependencies?: string[];
|
||||
completionPhrase?: string;
|
||||
timeoutMs?: number;
|
||||
}
|
||||
|
||||
export class Task {
|
||||
readonly id: string;
|
||||
readonly prompt: string;
|
||||
readonly workingDir: string;
|
||||
readonly priority: number;
|
||||
readonly dependencies: string[];
|
||||
readonly completionPhrase: string | undefined;
|
||||
readonly timeoutMs: number | undefined;
|
||||
readonly createdAt: number;
|
||||
|
||||
private _status: TaskStatus = 'pending';
|
||||
private _assignedSessionId: string | null = null;
|
||||
private _startedAt: number | null = null;
|
||||
private _completedAt: number | null = null;
|
||||
private _output: string = '';
|
||||
private _error: string | null = null;
|
||||
|
||||
constructor(options: CreateTaskOptions, id?: string) {
|
||||
this.id = id || uuidv4();
|
||||
this.prompt = options.prompt;
|
||||
this.workingDir = options.workingDir || process.cwd();
|
||||
this.priority = options.priority ?? 0;
|
||||
this.dependencies = options.dependencies || [];
|
||||
this.completionPhrase = options.completionPhrase;
|
||||
this.timeoutMs = options.timeoutMs;
|
||||
this.createdAt = Date.now();
|
||||
}
|
||||
|
||||
get status(): TaskStatus {
|
||||
return this._status;
|
||||
}
|
||||
|
||||
get assignedSessionId(): string | null {
|
||||
return this._assignedSessionId;
|
||||
}
|
||||
|
||||
get startedAt(): number | null {
|
||||
return this._startedAt;
|
||||
}
|
||||
|
||||
get completedAt(): number | null {
|
||||
return this._completedAt;
|
||||
}
|
||||
|
||||
get output(): string {
|
||||
return this._output;
|
||||
}
|
||||
|
||||
get error(): string | null {
|
||||
return this._error;
|
||||
}
|
||||
|
||||
isPending(): boolean {
|
||||
return this._status === 'pending';
|
||||
}
|
||||
|
||||
isRunning(): boolean {
|
||||
return this._status === 'running';
|
||||
}
|
||||
|
||||
isCompleted(): boolean {
|
||||
return this._status === 'completed';
|
||||
}
|
||||
|
||||
isFailed(): boolean {
|
||||
return this._status === 'failed';
|
||||
}
|
||||
|
||||
isDone(): boolean {
|
||||
return this._status === 'completed' || this._status === 'failed';
|
||||
}
|
||||
|
||||
toDefinition(): TaskDefinition {
|
||||
return {
|
||||
id: this.id,
|
||||
prompt: this.prompt,
|
||||
workingDir: this.workingDir,
|
||||
priority: this.priority,
|
||||
dependencies: this.dependencies,
|
||||
completionPhrase: this.completionPhrase,
|
||||
timeoutMs: this.timeoutMs,
|
||||
};
|
||||
}
|
||||
|
||||
toState(): TaskState {
|
||||
return {
|
||||
...this.toDefinition(),
|
||||
status: this._status,
|
||||
assignedSessionId: this._assignedSessionId,
|
||||
createdAt: this.createdAt,
|
||||
startedAt: this._startedAt,
|
||||
completedAt: this._completedAt,
|
||||
output: this._output,
|
||||
error: this._error,
|
||||
};
|
||||
}
|
||||
|
||||
static fromState(state: TaskState): Task {
|
||||
const task = new Task(
|
||||
{
|
||||
prompt: state.prompt,
|
||||
workingDir: state.workingDir,
|
||||
priority: state.priority,
|
||||
dependencies: state.dependencies,
|
||||
completionPhrase: state.completionPhrase,
|
||||
timeoutMs: state.timeoutMs,
|
||||
},
|
||||
state.id
|
||||
);
|
||||
task._status = state.status;
|
||||
task._assignedSessionId = state.assignedSessionId;
|
||||
task._startedAt = state.startedAt;
|
||||
task._completedAt = state.completedAt;
|
||||
task._output = state.output;
|
||||
task._error = state.error;
|
||||
return task;
|
||||
}
|
||||
|
||||
assign(sessionId: string): void {
|
||||
if (this._status !== 'pending') {
|
||||
throw new Error(`Cannot assign task ${this.id}: status is ${this._status}`);
|
||||
}
|
||||
this._status = 'running';
|
||||
this._assignedSessionId = sessionId;
|
||||
this._startedAt = Date.now();
|
||||
}
|
||||
|
||||
appendOutput(output: string): void {
|
||||
this._output += output;
|
||||
}
|
||||
|
||||
setError(error: string): void {
|
||||
this._error = error;
|
||||
}
|
||||
|
||||
complete(): void {
|
||||
this._status = 'completed';
|
||||
this._completedAt = Date.now();
|
||||
}
|
||||
|
||||
fail(error?: string): void {
|
||||
this._status = 'failed';
|
||||
this._completedAt = Date.now();
|
||||
if (error) {
|
||||
this._error = error;
|
||||
}
|
||||
}
|
||||
|
||||
reset(): void {
|
||||
this._status = 'pending';
|
||||
this._assignedSessionId = null;
|
||||
this._startedAt = null;
|
||||
this._completedAt = null;
|
||||
this._output = '';
|
||||
this._error = null;
|
||||
}
|
||||
|
||||
checkCompletion(output: string): boolean {
|
||||
if (this.completionPhrase) {
|
||||
// Check for exact completion phrase in promise tags
|
||||
const pattern = new RegExp(`<promise>${this.completionPhrase}</promise>`);
|
||||
return pattern.test(output);
|
||||
}
|
||||
// Default: check for any promise tag completion
|
||||
return /<promise>[^<]+<\/promise>/.test(output);
|
||||
}
|
||||
|
||||
isTimedOut(): boolean {
|
||||
if (!this.timeoutMs || !this._startedAt) {
|
||||
return false;
|
||||
}
|
||||
return Date.now() - this._startedAt > this.timeoutMs;
|
||||
}
|
||||
}
|
||||
+104
@@ -0,0 +1,104 @@
|
||||
export type SessionStatus = 'idle' | 'busy' | 'stopped' | 'error';
|
||||
export type TaskStatus = 'pending' | 'running' | 'completed' | 'failed';
|
||||
export type RalphLoopStatus = 'stopped' | 'running' | 'paused';
|
||||
|
||||
export interface SessionConfig {
|
||||
id: string;
|
||||
workingDir: string;
|
||||
createdAt: number;
|
||||
}
|
||||
|
||||
export interface SessionState {
|
||||
id: string;
|
||||
pid: number | null;
|
||||
status: SessionStatus;
|
||||
workingDir: string;
|
||||
currentTaskId: string | null;
|
||||
createdAt: number;
|
||||
lastActivityAt: number;
|
||||
}
|
||||
|
||||
export interface TaskDefinition {
|
||||
id: string;
|
||||
prompt: string;
|
||||
workingDir: string;
|
||||
priority: number;
|
||||
dependencies: string[];
|
||||
completionPhrase?: string;
|
||||
timeoutMs?: number;
|
||||
}
|
||||
|
||||
export interface TaskState {
|
||||
id: string;
|
||||
prompt: string;
|
||||
workingDir: string;
|
||||
priority: number;
|
||||
dependencies: string[];
|
||||
completionPhrase?: string;
|
||||
timeoutMs?: number;
|
||||
status: TaskStatus;
|
||||
assignedSessionId: string | null;
|
||||
createdAt: number;
|
||||
startedAt: number | null;
|
||||
completedAt: number | null;
|
||||
output: string;
|
||||
error: string | null;
|
||||
}
|
||||
|
||||
export interface RalphLoopState {
|
||||
status: RalphLoopStatus;
|
||||
startedAt: number | null;
|
||||
minDurationMs: number | null;
|
||||
tasksCompleted: number;
|
||||
tasksGenerated: number;
|
||||
lastCheckAt: number | null;
|
||||
}
|
||||
|
||||
export interface AppState {
|
||||
sessions: Record<string, SessionState>;
|
||||
tasks: Record<string, TaskState>;
|
||||
ralphLoop: RalphLoopState;
|
||||
config: AppConfig;
|
||||
}
|
||||
|
||||
export interface AppConfig {
|
||||
pollIntervalMs: number;
|
||||
defaultTimeoutMs: number;
|
||||
maxConcurrentSessions: number;
|
||||
stateFilePath: string;
|
||||
}
|
||||
|
||||
export interface SessionOutput {
|
||||
stdout: string;
|
||||
stderr: string;
|
||||
exitCode: number | null;
|
||||
}
|
||||
|
||||
export interface TaskAssignment {
|
||||
sessionId: string;
|
||||
taskId: string;
|
||||
assignedAt: number;
|
||||
}
|
||||
|
||||
export const DEFAULT_CONFIG: AppConfig = {
|
||||
pollIntervalMs: 1000,
|
||||
defaultTimeoutMs: 300000, // 5 minutes
|
||||
maxConcurrentSessions: 5,
|
||||
stateFilePath: '',
|
||||
};
|
||||
|
||||
export function createInitialState(): AppState {
|
||||
return {
|
||||
sessions: {},
|
||||
tasks: {},
|
||||
ralphLoop: {
|
||||
status: 'stopped',
|
||||
startedAt: null,
|
||||
minDurationMs: null,
|
||||
tasksCompleted: 0,
|
||||
tasksGenerated: 0,
|
||||
lastCheckAt: null,
|
||||
},
|
||||
config: { ...DEFAULT_CONFIG },
|
||||
};
|
||||
}
|
||||
Reference in New Issue
Block a user