mirror of
https://github.com/Ark0N/Codeman.git
synced 2026-10-02 13:39:41 +02:00
chore: version packages
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
This commit is contained in:
@@ -1,5 +1,11 @@
|
||||
# aicodeman
|
||||
|
||||
## 0.5.4
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Fix terminal flicker regression — re-add server-side DEC 2026 synchronized output wrapping around batched terminal data. Ink spinner frames (cursor-up + redraw cycles) do not emit their own DEC 2026 markers, so without the server wrapper each partial cursor update rendered individually causing visible flicker. Also: extract SSE stream management, session listener wiring, and respawn event wiring from server.ts into dedicated modules; deduplicate error message extraction across 7 files with shared getErrorMessage() helper; update SSE event count in CLAUDE.md (106 → 117).
|
||||
|
||||
## 0.5.3
|
||||
|
||||
### Patch Changes
|
||||
|
||||
@@ -52,7 +52,7 @@ When user says "COM":
|
||||
4. **Sync CLAUDE.md version**: Update the `**Version**` line below to match the new version from `package.json`
|
||||
5. **Commit and deploy**: `git add -A && git commit -m "chore: version packages" && git push && npm run build && systemctl --user restart codeman-web`
|
||||
|
||||
**Version**: 0.5.3 (must match `package.json`)
|
||||
**Version**: 0.5.4 (must match `package.json`)
|
||||
|
||||
## Project Overview
|
||||
|
||||
@@ -167,7 +167,7 @@ Frontend JS modules have `@fileoverview` with `@dependency`/`@loadorder` tags. L
|
||||
|
||||
### SSE Event Registry
|
||||
|
||||
~106 event types in `src/web/sse-events.ts` (backend) and `SSE_EVENTS` in `constants.js` (frontend). Both must be kept in sync.
|
||||
~117 event types in `src/web/sse-events.ts` (backend) and `SSE_EVENTS` in `constants.js` (frontend). Both must be kept in sync.
|
||||
|
||||
### API Routes
|
||||
|
||||
|
||||
+1
-1
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "aicodeman",
|
||||
"version": "0.5.3",
|
||||
"version": "0.5.4",
|
||||
"description": "The missing control plane for AI coding agents - run 20 autonomous agents with real-time monitoring and session persistence",
|
||||
"type": "module",
|
||||
"main": "dist/index.js",
|
||||
|
||||
@@ -30,6 +30,7 @@ import { join } from 'node:path';
|
||||
import { EventEmitter } from 'node:events';
|
||||
import { getAugmentedPath, ANSI_ESCAPE_PATTERN_SIMPLE } from './utils/index.js';
|
||||
import { AI_CHECK_MAX_BACKOFF_MS } from './config/ai-defaults.js';
|
||||
import { getErrorMessage } from './types.js';
|
||||
|
||||
// ========== Security Validation ==========
|
||||
|
||||
@@ -293,7 +294,7 @@ export abstract class AiCheckerBase<
|
||||
this.emit('checkCompleted', result);
|
||||
return result;
|
||||
} catch (err) {
|
||||
const errorMsg = err instanceof Error ? err.message : String(err);
|
||||
const errorMsg = getErrorMessage(err);
|
||||
this.handleError(errorMsg);
|
||||
const result = this.createErrorResult(errorMsg, Date.now() - this.checkStartTime);
|
||||
this.emit('checkFailed', errorMsg);
|
||||
@@ -412,9 +413,7 @@ export abstract class AiCheckerBase<
|
||||
});
|
||||
muxProcess.unref();
|
||||
} catch (err) {
|
||||
throw new Error(
|
||||
`Failed to spawn ${this.checkDescription} tmux session: ${err instanceof Error ? err.message : String(err)}`
|
||||
);
|
||||
throw new Error(`Failed to spawn ${this.checkDescription} tmux session: ${getErrorMessage(err)}`);
|
||||
}
|
||||
|
||||
// Poll the temp file for completion
|
||||
|
||||
@@ -16,6 +16,7 @@ import { existsSync, statSync, realpathSync } from 'node:fs';
|
||||
import { resolve, relative, isAbsolute } from 'node:path';
|
||||
import { homedir } from 'node:os';
|
||||
import { EventEmitter } from 'node:events';
|
||||
import { getErrorMessage } from './types.js';
|
||||
import { CLEANUP_CHECK_INTERVAL_MS, INACTIVITY_TIMEOUT_MS } from './config/server-timing.js';
|
||||
|
||||
// ========== Configuration Constants ==========
|
||||
@@ -172,10 +173,7 @@ export class FileStreamManager extends EventEmitter {
|
||||
}
|
||||
} catch (err) {
|
||||
const errorCode = err instanceof Error && 'code' in err ? (err as NodeJS.ErrnoException).code : 'UNKNOWN';
|
||||
console.warn(
|
||||
`[FileStreamManager] Failed to stat file "${absolutePath}" (${errorCode}):`,
|
||||
err instanceof Error ? err.message : String(err)
|
||||
);
|
||||
console.warn(`[FileStreamManager] Failed to stat file "${absolutePath}" (${errorCode}):`, getErrorMessage(err));
|
||||
return { success: false, error: 'File not found or not accessible' };
|
||||
}
|
||||
|
||||
|
||||
@@ -20,7 +20,13 @@
|
||||
*/
|
||||
|
||||
import type { Session } from './session.js';
|
||||
import type { OrchestratorPhase, OrchestratorConfig, VerificationResult, VerificationCheck } from './types.js';
|
||||
import {
|
||||
getErrorMessage,
|
||||
type OrchestratorPhase,
|
||||
type OrchestratorConfig,
|
||||
type VerificationResult,
|
||||
type VerificationCheck,
|
||||
} from './types.js';
|
||||
|
||||
// ═══════════════════════════════════════════════════════════════
|
||||
// Constants
|
||||
@@ -111,7 +117,7 @@ export class OrchestratorVerifier {
|
||||
type: 'test_command',
|
||||
description: `Run: ${command}`,
|
||||
passed: false,
|
||||
output: err instanceof Error ? err.message : String(err),
|
||||
output: getErrorMessage(err),
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -164,7 +170,7 @@ export class OrchestratorVerifier {
|
||||
type: 'ai_review',
|
||||
description: `AI review of "${phase.name}"`,
|
||||
passed: false,
|
||||
output: `AI review timed out or failed: ${err instanceof Error ? err.message : String(err)}`,
|
||||
output: `AI review timed out or failed: ${getErrorMessage(err)}`,
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
@@ -20,7 +20,7 @@ import type { TerminalMultiplexer } from './mux-interface.js';
|
||||
import { existsSync, mkdirSync, writeFileSync } from 'node:fs';
|
||||
import { join } from 'node:path';
|
||||
import { RESEARCH_AGENT_PROMPT, PLANNER_PROMPT } from './prompts/index.js';
|
||||
import type { PlanItem } from './types.js';
|
||||
import { getErrorMessage, type PlanItem } from './types.js';
|
||||
|
||||
// Re-export for backward compatibility
|
||||
export type { PlanItem };
|
||||
@@ -355,7 +355,7 @@ export class PlanOrchestrator {
|
||||
} catch (err) {
|
||||
return {
|
||||
success: false,
|
||||
error: err instanceof Error ? err.message : String(err),
|
||||
error: getErrorMessage(err),
|
||||
};
|
||||
}
|
||||
}
|
||||
@@ -523,7 +523,7 @@ export class PlanOrchestrator {
|
||||
return result;
|
||||
} catch (err) {
|
||||
const durationMs = Date.now() - startTime;
|
||||
const error = err instanceof Error ? err.message : String(err);
|
||||
const error = getErrorMessage(err);
|
||||
this._emitAgentFailure(onSubagent, agentId, 'research', this.researchModel, error, durationMs);
|
||||
return {
|
||||
success: false,
|
||||
@@ -651,7 +651,7 @@ export class PlanOrchestrator {
|
||||
return { success: true, items, gaps, warnings };
|
||||
} catch (err) {
|
||||
const durationMs = Date.now() - startTime;
|
||||
const error = err instanceof Error ? err.message : String(err);
|
||||
const error = getErrorMessage(err);
|
||||
this._emitAgentFailure(onSubagent, agentId, 'planner', this.plannerModel, error, durationMs);
|
||||
return { success: false, error };
|
||||
} finally {
|
||||
|
||||
@@ -70,12 +70,13 @@ import {
|
||||
AI_PLAN_CHECK_TIMEOUT_MS,
|
||||
AI_PLAN_CHECK_COOLDOWN_MS,
|
||||
} from './config/ai-defaults.js';
|
||||
import type {
|
||||
RespawnCycleMetrics,
|
||||
RespawnAggregateMetrics,
|
||||
RalphLoopHealthScore,
|
||||
TimingHistory,
|
||||
CycleOutcome,
|
||||
import {
|
||||
getErrorMessage,
|
||||
type RespawnCycleMetrics,
|
||||
type RespawnAggregateMetrics,
|
||||
type RalphLoopHealthScore,
|
||||
type TimingHistory,
|
||||
type CycleOutcome,
|
||||
} from './types.js';
|
||||
|
||||
// ========== Constants ==========
|
||||
@@ -2172,7 +2173,7 @@ export class RespawnController extends EventEmitter {
|
||||
}
|
||||
if (this._state === 'stopped') return; // Guard against stopped state
|
||||
if (this._state === 'ai_checking') {
|
||||
const errorMsg = err instanceof Error ? err.message : String(err);
|
||||
const errorMsg = getErrorMessage(err);
|
||||
this.logAction('ai-check', `Failed: ${errorMsg.substring(0, 50)}`);
|
||||
this.emit('aiCheckFailed', errorMsg);
|
||||
this.setState('watching');
|
||||
@@ -2371,7 +2372,7 @@ export class RespawnController extends EventEmitter {
|
||||
}
|
||||
})
|
||||
.catch((err) => {
|
||||
const errorMsg = err instanceof Error ? err.message : String(err);
|
||||
const errorMsg = getErrorMessage(err);
|
||||
this.emit('planCheckFailed', errorMsg);
|
||||
this.logAction('plan-check', `Failed: ${errorMsg.substring(0, 50)}`);
|
||||
});
|
||||
|
||||
+2
-1
@@ -40,6 +40,7 @@ import {
|
||||
ActiveBashTool,
|
||||
NiceConfig,
|
||||
DEFAULT_NICE_CONFIG,
|
||||
getErrorMessage,
|
||||
type ClaudeMode,
|
||||
type SessionMode,
|
||||
type OpenCodeConfig,
|
||||
@@ -1955,7 +1956,7 @@ export class Session extends EventEmitter {
|
||||
this._status = 'busy';
|
||||
this._lastActivityAt = Date.now();
|
||||
this.runPrompt(input).catch((err) => {
|
||||
const errorMsg = err instanceof Error ? err.message : String(err);
|
||||
const errorMsg = getErrorMessage(err);
|
||||
// Clean up task state so the task queue doesn't get stuck
|
||||
if (this._currentTaskId) {
|
||||
const taskId = this._currentTaskId;
|
||||
|
||||
@@ -29,6 +29,7 @@ import {
|
||||
RESTART_DELAY_MS,
|
||||
FORCE_KILL_MS,
|
||||
} from './config/tunnel-config.js';
|
||||
import { getErrorMessage } from './types.js';
|
||||
|
||||
// ========== Types ==========
|
||||
|
||||
@@ -164,7 +165,7 @@ export class TunnelManager extends EventEmitter {
|
||||
detached: false,
|
||||
});
|
||||
} catch (err) {
|
||||
this.emit('error', `Failed to spawn cloudflared: ${err instanceof Error ? err.message : String(err)}`);
|
||||
this.emit('error', `Failed to spawn cloudflared: ${getErrorMessage(err)}`);
|
||||
return;
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,360 @@
|
||||
/**
|
||||
* @fileoverview Respawn event wiring — pure functions that connect RespawnController
|
||||
* events to SSE broadcasts, push notifications, and run summary tracking.
|
||||
*
|
||||
* Extracted from WebServer to keep respawn-specific event plumbing separate from
|
||||
* HTTP/session concerns. Follows the same DI pattern as session-listener-wiring.ts.
|
||||
*/
|
||||
|
||||
import { Session } from '../session.js';
|
||||
import { RespawnController, RespawnConfig, RespawnState } from '../respawn-controller.js';
|
||||
import type { PersistedRespawnConfig } from '../types.js';
|
||||
import type { RunSummaryTracker } from '../run-summary.js';
|
||||
import type { TerminalMultiplexer } from '../mux-interface.js';
|
||||
import type { TeamWatcher } from '../team-watcher.js';
|
||||
import { SseEvent } from './sse-events.js';
|
||||
|
||||
// ============================================================================
|
||||
// Dependency Interface
|
||||
// ============================================================================
|
||||
|
||||
export interface RespawnWiringDeps {
|
||||
broadcast(event: string, data: unknown): void;
|
||||
sendPushNotifications(event: string, data: Record<string, unknown>): void;
|
||||
persistSessionState(session: Session): void;
|
||||
getSession(sessionId: string): Session | undefined;
|
||||
sessionExists(sessionId: string): boolean;
|
||||
getRunSummaryTracker(sessionId: string): RunSummaryTracker | undefined;
|
||||
getRespawnControllers(): Map<string, RespawnController>;
|
||||
getRespawnTimers(): Map<string, { timer: NodeJS.Timeout; endAt: number; startedAt: number }>;
|
||||
getPendingRespawnStarts(): Map<string, NodeJS.Timeout>;
|
||||
teamWatcher: TeamWatcher;
|
||||
serverStartTime: number;
|
||||
respawnRestoreGracePeriodMs: number;
|
||||
mux: TerminalMultiplexer;
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// Respawn Listener Wiring
|
||||
// ============================================================================
|
||||
|
||||
/**
|
||||
* Wire a RespawnController's events to SSE broadcasts, push notifications,
|
||||
* and run summary tracking.
|
||||
*/
|
||||
export function wireRespawnListeners(sessionId: string, controller: RespawnController, deps: RespawnWiringDeps): void {
|
||||
// Wire team watcher for team-aware idle detection
|
||||
controller.setTeamWatcher(deps.teamWatcher);
|
||||
|
||||
// Helper to get tracker lazily (may not exist at setup time for restored sessions)
|
||||
const getTracker = () => deps.getRunSummaryTracker(sessionId);
|
||||
|
||||
// ─── Respawn State Machine ──────────────────────────────
|
||||
|
||||
/** Broadcasts `respawn:stateChanged` — state machine transition (e.g., IDLE → DETECTING → RESPAWNING) */
|
||||
controller.on('stateChanged', (state: RespawnState, prevState: RespawnState) => {
|
||||
deps.broadcast(SseEvent.RespawnStateChanged, { sessionId, state, prevState });
|
||||
const tracker = getTracker();
|
||||
if (tracker) tracker.recordStateChange(state, `${prevState} → ${state}`);
|
||||
});
|
||||
|
||||
// ─── Respawn Cycle Lifecycle ────────────────────────────
|
||||
|
||||
/** Broadcasts `respawn:cycleStarted` — new respawn cycle begins */
|
||||
controller.on('respawnCycleStarted', (cycleNumber: number) => {
|
||||
deps.broadcast(SseEvent.RespawnCycleStarted, { sessionId, cycleNumber });
|
||||
});
|
||||
|
||||
/** Broadcasts `respawn:cycleCompleted` — respawn cycle finished */
|
||||
controller.on('respawnCycleCompleted', (cycleNumber: number) => {
|
||||
deps.broadcast(SseEvent.RespawnCycleCompleted, { sessionId, cycleNumber });
|
||||
});
|
||||
|
||||
/** Broadcasts `respawn:blocked` + push notification — respawn blocked by error/circuit breaker */
|
||||
controller.on('respawnBlocked', (data: { reason: string; details: string }) => {
|
||||
deps.broadcast(SseEvent.RespawnBlocked, { sessionId, reason: data.reason, details: data.details });
|
||||
const sessionForPush = deps.getSession(sessionId);
|
||||
deps.sendPushNotifications(SseEvent.RespawnBlocked, {
|
||||
sessionId,
|
||||
sessionName: sessionForPush?.name ?? sessionId.slice(0, 8),
|
||||
reason: data.reason,
|
||||
});
|
||||
const tracker = getTracker();
|
||||
if (tracker) tracker.recordWarning(`Respawn blocked: ${data.reason}`, data.details);
|
||||
});
|
||||
|
||||
// ─── Respawn Step Progress ──────────────────────────────
|
||||
|
||||
/** Broadcasts `respawn:stepSent` — respawn step input sent (e.g., /clear, kickstart prompt) */
|
||||
controller.on('stepSent', (step: string, input: string) => {
|
||||
deps.broadcast(SseEvent.RespawnStepSent, { sessionId, step, input });
|
||||
});
|
||||
|
||||
/** Broadcasts `respawn:stepCompleted` — respawn step finished */
|
||||
controller.on('stepCompleted', (step: string) => {
|
||||
deps.broadcast(SseEvent.RespawnStepCompleted, { sessionId, step });
|
||||
});
|
||||
|
||||
/** Broadcasts `respawn:detectionUpdate` — idle/completion detection state changed */
|
||||
controller.on('detectionUpdate', (detection: unknown) => {
|
||||
deps.broadcast(SseEvent.RespawnDetectionUpdate, { sessionId, detection });
|
||||
});
|
||||
|
||||
/** Broadcasts `respawn:autoAcceptSent` — auto-accepted a permission prompt */
|
||||
controller.on('autoAcceptSent', () => {
|
||||
deps.broadcast(SseEvent.RespawnAutoAcceptSent, { sessionId });
|
||||
});
|
||||
|
||||
// ─── AI Checker Events ──────────────────────────────────
|
||||
|
||||
/** Broadcasts `respawn:aiCheckStarted` — AI idle checker invoked */
|
||||
controller.on('aiCheckStarted', () => {
|
||||
deps.broadcast(SseEvent.RespawnAiCheckStarted, { sessionId });
|
||||
});
|
||||
|
||||
/** Broadcasts `respawn:aiCheckCompleted` — AI idle check returned verdict (idle/working/stuck) */
|
||||
controller.on('aiCheckCompleted', (result: { verdict: string; reasoning: string; durationMs: number }) => {
|
||||
deps.broadcast(SseEvent.RespawnAiCheckCompleted, {
|
||||
sessionId,
|
||||
verdict: result.verdict,
|
||||
reasoning: result.reasoning,
|
||||
durationMs: result.durationMs,
|
||||
});
|
||||
const tracker = getTracker();
|
||||
if (tracker) tracker.recordAiCheckResult(result.verdict);
|
||||
});
|
||||
|
||||
/** Broadcasts `respawn:aiCheckFailed` — AI idle check errored */
|
||||
controller.on('aiCheckFailed', (error: string) => {
|
||||
deps.broadcast(SseEvent.RespawnAiCheckFailed, { sessionId, error });
|
||||
const tracker = getTracker();
|
||||
if (tracker) tracker.recordError('AI check failed', error);
|
||||
});
|
||||
|
||||
/** Broadcasts `respawn:aiCheckCooldown` — AI check on cooldown after failure */
|
||||
controller.on('aiCheckCooldown', (active: boolean, endsAt: number | null) => {
|
||||
deps.broadcast(SseEvent.RespawnAiCheckCooldown, { sessionId, active, endsAt });
|
||||
});
|
||||
|
||||
// ─── Plan Checker Events ────────────────────────────────
|
||||
|
||||
/** Broadcasts `respawn:planCheckStarted` — AI plan completion checker invoked */
|
||||
controller.on('planCheckStarted', () => {
|
||||
deps.broadcast(SseEvent.RespawnPlanCheckStarted, { sessionId });
|
||||
});
|
||||
|
||||
/** Broadcasts `respawn:planCheckCompleted` — plan check returned verdict */
|
||||
controller.on('planCheckCompleted', (result: { verdict: string; reasoning: string; durationMs: number }) => {
|
||||
deps.broadcast(SseEvent.RespawnPlanCheckCompleted, {
|
||||
sessionId,
|
||||
verdict: result.verdict,
|
||||
reasoning: result.reasoning,
|
||||
durationMs: result.durationMs,
|
||||
});
|
||||
});
|
||||
|
||||
/** Broadcasts `respawn:planCheckFailed` — plan check errored */
|
||||
controller.on('planCheckFailed', (error: string) => {
|
||||
deps.broadcast(SseEvent.RespawnPlanCheckFailed, { sessionId, error });
|
||||
});
|
||||
|
||||
// ─── Timer Events (UI countdown display) ────────────────
|
||||
|
||||
/** Broadcasts `respawn:timerStarted` — countdown timer started (idle, cooldown, etc.) */
|
||||
controller.on('timerStarted', (timer) => {
|
||||
deps.broadcast(SseEvent.RespawnTimerStarted, { sessionId, timer });
|
||||
});
|
||||
|
||||
/** Broadcasts `respawn:timerCancelled` — timer cancelled before expiry */
|
||||
controller.on('timerCancelled', (timerName, reason) => {
|
||||
deps.broadcast(SseEvent.RespawnTimerCancelled, { sessionId, timerName, reason });
|
||||
});
|
||||
|
||||
/** Broadcasts `respawn:timerCompleted` — timer expired */
|
||||
controller.on('timerCompleted', (timerName) => {
|
||||
deps.broadcast(SseEvent.RespawnTimerCompleted, { sessionId, timerName });
|
||||
});
|
||||
|
||||
// ─── Logging & Errors ───────────────────────────────────
|
||||
|
||||
/** Broadcasts `respawn:actionLog` — respawn action logged for audit/debugging */
|
||||
controller.on('actionLog', (action) => {
|
||||
deps.broadcast(SseEvent.RespawnActionLog, { sessionId, action });
|
||||
});
|
||||
|
||||
/** Broadcasts `respawn:log` — general respawn log message */
|
||||
controller.on('log', (message: string) => {
|
||||
deps.broadcast(SseEvent.RespawnLog, { sessionId, message });
|
||||
});
|
||||
|
||||
/** Broadcasts `respawn:error` — respawn controller error */
|
||||
controller.on('error', (error: Error) => {
|
||||
deps.broadcast(SseEvent.RespawnError, { sessionId, error: error.message });
|
||||
const tracker = getTracker();
|
||||
if (tracker) tracker.recordError('Respawn error', error.message);
|
||||
});
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// Timed Respawn
|
||||
// ============================================================================
|
||||
|
||||
/**
|
||||
* Set up a duration-limited respawn timer that stops respawn after N minutes.
|
||||
*/
|
||||
export function setupTimedRespawn(sessionId: string, durationMinutes: number, deps: RespawnWiringDeps): void {
|
||||
const timers = deps.getRespawnTimers();
|
||||
|
||||
// Clear existing timer if any
|
||||
const existing = timers.get(sessionId);
|
||||
if (existing) {
|
||||
clearTimeout(existing.timer);
|
||||
}
|
||||
|
||||
const now = Date.now();
|
||||
const endAt = now + durationMinutes * 60 * 1000;
|
||||
|
||||
const timer = setTimeout(
|
||||
() => {
|
||||
// Stop respawn when time is up
|
||||
const controllers = deps.getRespawnControllers();
|
||||
const controller = controllers.get(sessionId);
|
||||
if (controller) {
|
||||
controller.stop();
|
||||
controller.removeAllListeners();
|
||||
controllers.delete(sessionId);
|
||||
deps.broadcast(SseEvent.RespawnStopped, { sessionId, reason: 'duration_expired' });
|
||||
}
|
||||
timers.delete(sessionId);
|
||||
// Update persisted state (respawn no longer active)
|
||||
const session = deps.getSession(sessionId);
|
||||
if (session) {
|
||||
deps.persistSessionState(session);
|
||||
}
|
||||
},
|
||||
durationMinutes * 60 * 1000
|
||||
);
|
||||
|
||||
timers.set(sessionId, { timer, endAt, startedAt: now });
|
||||
deps.broadcast(SseEvent.RespawnTimerStarted, { sessionId, durationMinutes, endAt, startedAt: now });
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// Respawn Controller Restore
|
||||
// ============================================================================
|
||||
|
||||
/**
|
||||
* Restore a RespawnController from persisted configuration.
|
||||
* Creates the controller, wires listeners, and starts after a grace period.
|
||||
*/
|
||||
export function restoreRespawnController(
|
||||
session: Session,
|
||||
config: PersistedRespawnConfig,
|
||||
source: string,
|
||||
deps: RespawnWiringDeps
|
||||
): void {
|
||||
const controller = new RespawnController(session, {
|
||||
idleTimeoutMs: config.idleTimeoutMs,
|
||||
updatePrompt: config.updatePrompt,
|
||||
interStepDelayMs: config.interStepDelayMs,
|
||||
enabled: true,
|
||||
sendClear: config.sendClear,
|
||||
sendInit: config.sendInit,
|
||||
kickstartPrompt: config.kickstartPrompt,
|
||||
completionConfirmMs: config.completionConfirmMs,
|
||||
noOutputTimeoutMs: config.noOutputTimeoutMs,
|
||||
autoAcceptPrompts: config.autoAcceptPrompts,
|
||||
autoAcceptDelayMs: config.autoAcceptDelayMs,
|
||||
aiIdleCheckEnabled: config.aiIdleCheckEnabled,
|
||||
aiIdleCheckModel: config.aiIdleCheckModel,
|
||||
aiIdleCheckMaxContext: config.aiIdleCheckMaxContext,
|
||||
aiIdleCheckTimeoutMs: config.aiIdleCheckTimeoutMs,
|
||||
aiIdleCheckCooldownMs: config.aiIdleCheckCooldownMs,
|
||||
aiPlanCheckEnabled: config.aiPlanCheckEnabled,
|
||||
aiPlanCheckModel: config.aiPlanCheckModel,
|
||||
aiPlanCheckMaxContext: config.aiPlanCheckMaxContext,
|
||||
aiPlanCheckTimeoutMs: config.aiPlanCheckTimeoutMs,
|
||||
aiPlanCheckCooldownMs: config.aiPlanCheckCooldownMs,
|
||||
});
|
||||
|
||||
const controllers = deps.getRespawnControllers();
|
||||
controllers.set(session.id, controller);
|
||||
wireRespawnListeners(session.id, controller, deps);
|
||||
|
||||
// Calculate delay: wait until grace period after server start before starting respawn
|
||||
// This prevents false idle detection immediately after a server restart/rebuild
|
||||
const timeSinceStart = Date.now() - deps.serverStartTime;
|
||||
const delayMs = Math.max(0, deps.respawnRestoreGracePeriodMs - timeSinceStart);
|
||||
|
||||
const pendingStarts = deps.getPendingRespawnStarts();
|
||||
|
||||
if (delayMs > 0) {
|
||||
console.log(
|
||||
`[Server] Restored respawn controller for session ${session.id} from ${source} (will start in ${Math.ceil(delayMs / 1000)}s)`
|
||||
);
|
||||
const delayTimer = setTimeout(() => {
|
||||
pendingStarts.delete(session.id);
|
||||
// Verify session still exists (may have been deleted during grace period)
|
||||
if (!deps.sessionExists(session.id)) {
|
||||
console.log(`[Server] Skipping restored respawn start - session ${session.id} no longer exists`);
|
||||
return;
|
||||
}
|
||||
// Double-check controller still exists and is stopped
|
||||
const ctrl = controllers.get(session.id);
|
||||
if (ctrl && ctrl.state === 'stopped') {
|
||||
ctrl.start();
|
||||
deps.broadcast(SseEvent.RespawnStarted, { sessionId: session.id });
|
||||
console.log(`[Server] Restored respawn controller started for session ${session.id}`);
|
||||
}
|
||||
}, delayMs);
|
||||
pendingStarts.set(session.id, delayTimer);
|
||||
} else {
|
||||
// Grace period has passed, start immediately
|
||||
controller.start();
|
||||
console.log(`[Server] Restored respawn controller for session ${session.id} from ${source} (started immediately)`);
|
||||
}
|
||||
|
||||
if (config.durationMinutes && config.durationMinutes > 0) {
|
||||
setupTimedRespawn(session.id, config.durationMinutes, deps);
|
||||
}
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// Respawn Config Persistence
|
||||
// ============================================================================
|
||||
|
||||
/**
|
||||
* Save respawn config to mux for restart recovery.
|
||||
*/
|
||||
export function saveRespawnConfig(
|
||||
sessionId: string,
|
||||
config: RespawnConfig,
|
||||
mux: TerminalMultiplexer,
|
||||
durationMinutes?: number
|
||||
): void {
|
||||
const persistedConfig: PersistedRespawnConfig = {
|
||||
enabled: config.enabled,
|
||||
idleTimeoutMs: config.idleTimeoutMs,
|
||||
updatePrompt: config.updatePrompt,
|
||||
interStepDelayMs: config.interStepDelayMs,
|
||||
sendClear: config.sendClear,
|
||||
sendInit: config.sendInit,
|
||||
kickstartPrompt: config.kickstartPrompt,
|
||||
autoAcceptPrompts: config.autoAcceptPrompts,
|
||||
autoAcceptDelayMs: config.autoAcceptDelayMs,
|
||||
completionConfirmMs: config.completionConfirmMs,
|
||||
noOutputTimeoutMs: config.noOutputTimeoutMs,
|
||||
aiIdleCheckEnabled: config.aiIdleCheckEnabled,
|
||||
aiIdleCheckModel: config.aiIdleCheckModel,
|
||||
aiIdleCheckMaxContext: config.aiIdleCheckMaxContext,
|
||||
aiIdleCheckTimeoutMs: config.aiIdleCheckTimeoutMs,
|
||||
aiIdleCheckCooldownMs: config.aiIdleCheckCooldownMs,
|
||||
aiPlanCheckEnabled: config.aiPlanCheckEnabled,
|
||||
aiPlanCheckModel: config.aiPlanCheckModel,
|
||||
aiPlanCheckMaxContext: config.aiPlanCheckMaxContext,
|
||||
aiPlanCheckTimeoutMs: config.aiPlanCheckTimeoutMs,
|
||||
aiPlanCheckCooldownMs: config.aiPlanCheckCooldownMs,
|
||||
durationMinutes,
|
||||
};
|
||||
mux.updateRespawnConfig(sessionId, persistedConfig);
|
||||
}
|
||||
+118
-1083
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,392 @@
|
||||
/**
|
||||
* @fileoverview Session event listener wiring — creates, attaches, and detaches session listeners.
|
||||
*
|
||||
* Extracted from server.ts for modularity. Provides:
|
||||
* - `SessionListenerRefs` interface (named listener references for leak-free cleanup)
|
||||
* - `createSessionListeners()` — builds all 25 listener handlers via dependency injection
|
||||
* - `attachSessionListeners()` / `detachSessionListeners()` — symmetric attach/detach
|
||||
*
|
||||
* The detach function deduplicates a pattern that was previously copy-pasted 3 times
|
||||
* in server.ts (_doCleanupSession, exit handler, stop()).
|
||||
*
|
||||
* @dependencies session.ts (Session, event types), sse-events.ts, types.ts
|
||||
* @consumedby web/server.ts (WebServer delegates listener lifecycle here)
|
||||
*
|
||||
* @module web/session-listener-wiring
|
||||
*/
|
||||
|
||||
import type {
|
||||
Session,
|
||||
ClaudeMessage,
|
||||
BackgroundTask,
|
||||
RalphTrackerState,
|
||||
RalphTodoItem,
|
||||
ActiveBashTool,
|
||||
} from '../session.js';
|
||||
import type { RalphStatusBlock, CircuitBreakerStatus } from '../types.js';
|
||||
import { SseEvent } from './sse-events.js';
|
||||
import { getLifecycleLog } from '../session-lifecycle-log.js';
|
||||
import { fileStreamManager } from '../file-stream-manager.js';
|
||||
|
||||
/** Stored listener references for session cleanup (prevents memory leaks) */
|
||||
export interface SessionListenerRefs {
|
||||
terminal: (data: string) => void;
|
||||
clearTerminal: () => void;
|
||||
needsRefresh: () => void;
|
||||
message: (msg: ClaudeMessage) => void;
|
||||
error: (error: string) => void;
|
||||
completion: (result: string, cost: number) => void;
|
||||
exit: (code: number | null) => void;
|
||||
working: () => void;
|
||||
idle: () => void;
|
||||
taskCreated: (task: BackgroundTask) => void;
|
||||
taskUpdated: (task: BackgroundTask) => void;
|
||||
taskCompleted: (task: BackgroundTask) => void;
|
||||
taskFailed: (task: BackgroundTask, error: string) => void;
|
||||
autoClear: (data: { tokens: number; threshold: number }) => void;
|
||||
autoCompact: (data: { tokens: number; threshold: number; prompt?: string }) => void;
|
||||
cliInfoUpdated: (data: { version?: string; model?: string; accountType?: string; latestVersion?: string }) => void;
|
||||
ralphLoopUpdate: (state: RalphTrackerState) => void;
|
||||
ralphTodoUpdate: (todos: RalphTodoItem[]) => void;
|
||||
ralphCompletionDetected: (phrase: string) => void;
|
||||
ralphStatusBlockDetected: (block: RalphStatusBlock) => void;
|
||||
ralphCircuitBreakerUpdate: (status: CircuitBreakerStatus) => void;
|
||||
ralphExitGateMet: (data: { completionIndicators: number; exitSignal: boolean }) => void;
|
||||
bashToolStart: (tool: ActiveBashTool) => void;
|
||||
bashToolEnd: (tool: ActiveBashTool) => void;
|
||||
bashToolsUpdate: (tools: ActiveBashTool[]) => void;
|
||||
}
|
||||
|
||||
/** Dependencies injected by WebServer — keeps listener creation decoupled from server internals. */
|
||||
export interface SessionListenerDeps {
|
||||
broadcast(event: string, data: unknown): void;
|
||||
batchTerminalData(sessionId: string, data: string): void;
|
||||
batchTaskUpdate(sessionId: string, task: BackgroundTask): void;
|
||||
broadcastSessionStateDebounced(sessionId: string): void;
|
||||
sendPushNotifications(event: string, data: Record<string, unknown>): void;
|
||||
persistSessionState(session: Session): void;
|
||||
getSessionStateWithRespawn(session: Session): unknown;
|
||||
getRunSummaryTracker(sessionId: string): import('../run-summary.js').RunSummaryTracker | undefined;
|
||||
stopTranscriptWatcher(sessionId: string): void;
|
||||
cleanupSessionBatches(sessionId: string): void;
|
||||
cancelPersistDebounce(sessionId: string): void;
|
||||
removeRunSummaryTracker(sessionId: string): void;
|
||||
removeSessionListenerRefs(sessionId: string): void;
|
||||
cleanupRespawnOnExit(sessionId: string): void;
|
||||
getStore(): import('../state-store.js').StateStore;
|
||||
}
|
||||
|
||||
/**
|
||||
* Creates all 25 session listener handlers, capturing dependencies via closure.
|
||||
* Call `attachSessionListeners()` after to wire them to the session.
|
||||
*/
|
||||
export function createSessionListeners(session: Session, deps: SessionListenerDeps): SessionListenerRefs {
|
||||
return {
|
||||
// ─── Terminal Output ─────────────────────────────────────
|
||||
|
||||
/** Batches PTY output → broadcasts `session:terminal` at 16-50ms intervals */
|
||||
terminal: (data) => {
|
||||
deps.batchTerminalData(session.id, data);
|
||||
},
|
||||
|
||||
/** Broadcasts `session:clearTerminal` — tells clients to wipe their xterm buffer (after mux attach) */
|
||||
clearTerminal: () => {
|
||||
deps.broadcast(SseEvent.SessionClearTerminal, { id: session.id });
|
||||
},
|
||||
|
||||
/** Broadcasts `session:needsRefresh` — tells clients to reload buffer */
|
||||
needsRefresh: () => {
|
||||
deps.broadcast(SseEvent.SessionNeedsRefresh, { id: session.id });
|
||||
},
|
||||
|
||||
// ─── Session Messages & Errors ──────────────────────────
|
||||
|
||||
/** Broadcasts `session:message` — structured Claude JSON messages (assistant, tool_use, etc.) */
|
||||
message: (msg: ClaudeMessage) => {
|
||||
deps.broadcast(SseEvent.SessionMessage, { id: session.id, message: msg });
|
||||
},
|
||||
|
||||
/** Broadcasts `session:error` + sends push notification */
|
||||
error: (error) => {
|
||||
deps.broadcast(SseEvent.SessionError, { id: session.id, error });
|
||||
deps.sendPushNotifications(SseEvent.SessionError, {
|
||||
sessionId: session.id,
|
||||
sessionName: session.name,
|
||||
error: String(error),
|
||||
});
|
||||
const tracker = deps.getRunSummaryTracker(session.id);
|
||||
if (tracker) tracker.recordError('Session error', String(error));
|
||||
},
|
||||
|
||||
/** Broadcasts `session:completion` + `session:updated` — prompt finished, persists state */
|
||||
completion: (result, cost) => {
|
||||
deps.broadcast(SseEvent.SessionCompletion, { id: session.id, result, cost });
|
||||
deps.broadcast(SseEvent.SessionUpdated, deps.getSessionStateWithRespawn(session));
|
||||
deps.persistSessionState(session);
|
||||
const tracker = deps.getRunSummaryTracker(session.id);
|
||||
if (tracker) tracker.recordTokens(session.inputTokens, session.outputTokens);
|
||||
},
|
||||
|
||||
// ─── Session Lifecycle ──────────────────────────────────
|
||||
|
||||
/** Broadcasts `session:exit` + `session:updated` — PTY process exited; cleans up respawn, timers, listeners */
|
||||
exit: (code) => {
|
||||
getLifecycleLog().log({
|
||||
event: 'exit',
|
||||
sessionId: session.id,
|
||||
name: session.name,
|
||||
exitCode: code,
|
||||
});
|
||||
// Wrap in try/catch to ensure cleanup always happens
|
||||
try {
|
||||
deps.broadcast(SseEvent.SessionExit, { id: session.id, code });
|
||||
deps.broadcast(SseEvent.SessionUpdated, deps.getSessionStateWithRespawn(session));
|
||||
deps.persistSessionState(session);
|
||||
} catch (err) {
|
||||
console.error(`[Server] Error broadcasting session exit for ${session.id}:`, err);
|
||||
}
|
||||
|
||||
// Always clean up respawn controller, even if broadcast failed
|
||||
try {
|
||||
deps.cleanupRespawnOnExit(session.id);
|
||||
} catch (err) {
|
||||
console.error(`[Server] Error cleaning up respawn controller for ${session.id}:`, err);
|
||||
}
|
||||
|
||||
// Clean up per-session resources that are stale after PTY exit.
|
||||
try {
|
||||
// Transcript watcher is tied to the specific PTY run
|
||||
deps.stopTranscriptWatcher(session.id);
|
||||
|
||||
// Finalize run summary tracker
|
||||
deps.removeRunSummaryTracker(session.id);
|
||||
|
||||
// Flush/clear terminal batching state (no more output coming)
|
||||
deps.cleanupSessionBatches(session.id);
|
||||
|
||||
// Clear pending persist-debounce timer
|
||||
deps.cancelPersistDebounce(session.id);
|
||||
|
||||
// Close any active file streams
|
||||
fileStreamManager.closeSessionStreams(session.id);
|
||||
|
||||
// Remove stored listener refs to break closure references (prevents memory leak).
|
||||
deps.removeSessionListenerRefs(session.id);
|
||||
} catch (err) {
|
||||
console.error(`[Server] Error cleaning up session resources on exit for ${session.id}:`, err);
|
||||
}
|
||||
},
|
||||
|
||||
// ─── Activity State ─────────────────────────────────────
|
||||
|
||||
/** Broadcasts `session:working` — Claude started processing */
|
||||
working: () => {
|
||||
deps.broadcast(SseEvent.SessionWorking, { id: session.id });
|
||||
const tracker = deps.getRunSummaryTracker(session.id);
|
||||
if (tracker) {
|
||||
tracker.recordWorking();
|
||||
tracker.recordTokens(session.inputTokens, session.outputTokens);
|
||||
}
|
||||
},
|
||||
|
||||
/** Broadcasts `session:idle` — Claude finished processing, waiting for input */
|
||||
idle: () => {
|
||||
deps.broadcast(SseEvent.SessionIdle, { id: session.id });
|
||||
deps.broadcastSessionStateDebounced(session.id);
|
||||
const tracker = deps.getRunSummaryTracker(session.id);
|
||||
if (tracker) {
|
||||
tracker.recordIdle();
|
||||
tracker.recordTokens(session.inputTokens, session.outputTokens);
|
||||
}
|
||||
},
|
||||
|
||||
// ─── Background Task Events ──────────────────────────────
|
||||
|
||||
/** Broadcasts `task:created` — new background task discovered */
|
||||
taskCreated: (task: BackgroundTask) => {
|
||||
deps.broadcast(SseEvent.TaskCreated, { sessionId: session.id, task });
|
||||
deps.broadcastSessionStateDebounced(session.id);
|
||||
},
|
||||
|
||||
/** Batched broadcast of `task:updated` — high-frequency progress updates */
|
||||
taskUpdated: (task: BackgroundTask) => {
|
||||
deps.batchTaskUpdate(session.id, task);
|
||||
},
|
||||
|
||||
/** Broadcasts `task:completed` — background task finished successfully */
|
||||
taskCompleted: (task: BackgroundTask) => {
|
||||
deps.broadcast(SseEvent.TaskCompleted, { sessionId: session.id, task });
|
||||
deps.broadcastSessionStateDebounced(session.id);
|
||||
},
|
||||
|
||||
/** Broadcasts `task:failed` — background task errored */
|
||||
taskFailed: (task: BackgroundTask, error: string) => {
|
||||
deps.broadcast(SseEvent.TaskFailed, { sessionId: session.id, task, error });
|
||||
deps.broadcastSessionStateDebounced(session.id);
|
||||
},
|
||||
|
||||
// ─── Auto-Operations ────────────────────────────────────
|
||||
|
||||
/** Broadcasts `session:autoClear` — context window auto-cleared at token threshold */
|
||||
autoClear: (data: { tokens: number; threshold: number }) => {
|
||||
deps.broadcast(SseEvent.SessionAutoClear, { sessionId: session.id, ...data });
|
||||
deps.broadcastSessionStateDebounced(session.id);
|
||||
const tracker = deps.getRunSummaryTracker(session.id);
|
||||
if (tracker) tracker.recordAutoClear(data.tokens, data.threshold);
|
||||
},
|
||||
|
||||
/** Broadcasts `session:autoCompact` — context window auto-compacted at token threshold */
|
||||
autoCompact: (data: { tokens: number; threshold: number; prompt?: string }) => {
|
||||
deps.broadcast(SseEvent.SessionAutoCompact, { sessionId: session.id, ...data });
|
||||
deps.broadcastSessionStateDebounced(session.id);
|
||||
const tracker = deps.getRunSummaryTracker(session.id);
|
||||
if (tracker) tracker.recordAutoCompact(data.tokens, data.threshold);
|
||||
},
|
||||
|
||||
// ─── CLI Info ────────────────────────────────────────────
|
||||
|
||||
/** Broadcasts `session:cliInfo` — Claude Code version, model, account type parsed from terminal */
|
||||
cliInfoUpdated: (data: { version?: string; model?: string; accountType?: string; latestVersion?: string }) => {
|
||||
deps.broadcast(SseEvent.SessionCliInfo, { sessionId: session.id, ...data });
|
||||
deps.broadcastSessionStateDebounced(session.id);
|
||||
},
|
||||
|
||||
// ─── Ralph Tracking Events ──────────────────────────────
|
||||
|
||||
/** Broadcasts `session:ralphLoopUpdate` — Ralph tracker loop state changed (iteration, phase) */
|
||||
ralphLoopUpdate: (state: RalphTrackerState) => {
|
||||
deps.broadcast(SseEvent.SessionRalphLoopUpdate, { sessionId: session.id, state });
|
||||
deps.getStore().updateRalphState(session.id, { loop: state });
|
||||
},
|
||||
|
||||
/** Broadcasts `session:ralphTodoUpdate` — todo items added, completed, or modified */
|
||||
ralphTodoUpdate: (todos: RalphTodoItem[]) => {
|
||||
deps.broadcast(SseEvent.SessionRalphTodoUpdate, { sessionId: session.id, todos });
|
||||
deps.getStore().updateRalphState(session.id, { todos });
|
||||
},
|
||||
|
||||
/** Broadcasts `session:ralphCompletionDetected` + push notification — completion phrase matched */
|
||||
ralphCompletionDetected: (phrase: string) => {
|
||||
deps.broadcast(SseEvent.SessionRalphCompletionDetected, { sessionId: session.id, phrase });
|
||||
deps.sendPushNotifications(SseEvent.SessionRalphCompletionDetected, {
|
||||
sessionId: session.id,
|
||||
sessionName: session.name,
|
||||
phrase,
|
||||
});
|
||||
const tracker = deps.getRunSummaryTracker(session.id);
|
||||
if (tracker) tracker.recordRalphCompletion(phrase);
|
||||
},
|
||||
|
||||
/** Broadcasts `session:ralphStatusUpdate` — RALPH_STATUS block parsed from output */
|
||||
ralphStatusBlockDetected: (block: RalphStatusBlock) => {
|
||||
deps.broadcast(SseEvent.SessionRalphStatusUpdate, { sessionId: session.id, block });
|
||||
const tracker = deps.getRunSummaryTracker(session.id);
|
||||
if (tracker) {
|
||||
tracker.addEvent(
|
||||
block.status === 'BLOCKED' ? 'warning' : 'idle_detected',
|
||||
block.status === 'BLOCKED' ? 'warning' : 'info',
|
||||
`Ralph Status: ${block.status}`,
|
||||
`Tasks: ${block.tasksCompletedThisLoop}, Files: ${block.filesModified}, Tests: ${block.testsStatus}`
|
||||
);
|
||||
}
|
||||
},
|
||||
|
||||
/** Broadcasts `session:circuitBreakerUpdate` — circuit breaker state changed (CLOSED/HALF_OPEN/OPEN) */
|
||||
ralphCircuitBreakerUpdate: (status: CircuitBreakerStatus) => {
|
||||
deps.broadcast(SseEvent.SessionCircuitBreakerUpdate, { sessionId: session.id, status });
|
||||
const tracker = deps.getRunSummaryTracker(session.id);
|
||||
if (tracker && status.state === 'OPEN') {
|
||||
tracker.addEvent('warning', 'warning', 'Circuit Breaker Opened', status.reason);
|
||||
}
|
||||
},
|
||||
|
||||
/** Broadcasts `session:exitGateMet` — all completion indicators met, ready to exit */
|
||||
ralphExitGateMet: (data: { completionIndicators: number; exitSignal: boolean }) => {
|
||||
deps.broadcast(SseEvent.SessionExitGateMet, { sessionId: session.id, ...data });
|
||||
const tracker = deps.getRunSummaryTracker(session.id);
|
||||
if (tracker) {
|
||||
tracker.addEvent(
|
||||
'ralph_completion',
|
||||
'success',
|
||||
'Exit Gate Met',
|
||||
`Indicators: ${data.completionIndicators}, EXIT_SIGNAL: ${data.exitSignal}`
|
||||
);
|
||||
}
|
||||
},
|
||||
|
||||
// ─── Bash Tool Tracking ────────────────────────────────
|
||||
|
||||
/** Broadcasts `session:bashToolStart` — bash tool invocation started */
|
||||
bashToolStart: (tool: ActiveBashTool) => {
|
||||
deps.broadcast(SseEvent.SessionBashToolStart, { sessionId: session.id, tool });
|
||||
},
|
||||
|
||||
/** Broadcasts `session:bashToolEnd` — bash tool invocation completed */
|
||||
bashToolEnd: (tool: ActiveBashTool) => {
|
||||
deps.broadcast(SseEvent.SessionBashToolEnd, { sessionId: session.id, tool });
|
||||
},
|
||||
|
||||
/** Broadcasts `session:bashToolsUpdate` — full active bash tools list refreshed */
|
||||
bashToolsUpdate: (tools: ActiveBashTool[]) => {
|
||||
deps.broadcast(SseEvent.SessionBashToolsUpdate, { sessionId: session.id, tools });
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
/** Attach all listeners to a session. */
|
||||
export function attachSessionListeners(session: Session, refs: SessionListenerRefs): void {
|
||||
session.on('terminal', refs.terminal);
|
||||
session.on('clearTerminal', refs.clearTerminal);
|
||||
session.on('needsRefresh', refs.needsRefresh);
|
||||
session.on('message', refs.message);
|
||||
session.on('error', refs.error);
|
||||
session.on('completion', refs.completion);
|
||||
session.on('exit', refs.exit);
|
||||
session.on('working', refs.working);
|
||||
session.on('idle', refs.idle);
|
||||
session.on('taskCreated', refs.taskCreated);
|
||||
session.on('taskUpdated', refs.taskUpdated);
|
||||
session.on('taskCompleted', refs.taskCompleted);
|
||||
session.on('taskFailed', refs.taskFailed);
|
||||
session.on('autoClear', refs.autoClear);
|
||||
session.on('autoCompact', refs.autoCompact);
|
||||
session.on('cliInfoUpdated', refs.cliInfoUpdated);
|
||||
session.on('ralphLoopUpdate', refs.ralphLoopUpdate);
|
||||
session.on('ralphTodoUpdate', refs.ralphTodoUpdate);
|
||||
session.on('ralphCompletionDetected', refs.ralphCompletionDetected);
|
||||
session.on('ralphStatusBlockDetected', refs.ralphStatusBlockDetected);
|
||||
session.on('ralphCircuitBreakerUpdate', refs.ralphCircuitBreakerUpdate);
|
||||
session.on('ralphExitGateMet', refs.ralphExitGateMet);
|
||||
session.on('bashToolStart', refs.bashToolStart);
|
||||
session.on('bashToolEnd', refs.bashToolEnd);
|
||||
session.on('bashToolsUpdate', refs.bashToolsUpdate);
|
||||
}
|
||||
|
||||
/** Detach all listeners from a session (prevents memory leaks from closure references). */
|
||||
export function detachSessionListeners(session: Session, refs: SessionListenerRefs): void {
|
||||
session.off('terminal', refs.terminal);
|
||||
session.off('clearTerminal', refs.clearTerminal);
|
||||
session.off('needsRefresh', refs.needsRefresh);
|
||||
session.off('message', refs.message);
|
||||
session.off('error', refs.error);
|
||||
session.off('completion', refs.completion);
|
||||
session.off('exit', refs.exit);
|
||||
session.off('working', refs.working);
|
||||
session.off('idle', refs.idle);
|
||||
session.off('taskCreated', refs.taskCreated);
|
||||
session.off('taskUpdated', refs.taskUpdated);
|
||||
session.off('taskCompleted', refs.taskCompleted);
|
||||
session.off('taskFailed', refs.taskFailed);
|
||||
session.off('autoClear', refs.autoClear);
|
||||
session.off('autoCompact', refs.autoCompact);
|
||||
session.off('cliInfoUpdated', refs.cliInfoUpdated);
|
||||
session.off('ralphLoopUpdate', refs.ralphLoopUpdate);
|
||||
session.off('ralphTodoUpdate', refs.ralphTodoUpdate);
|
||||
session.off('ralphCompletionDetected', refs.ralphCompletionDetected);
|
||||
session.off('ralphStatusBlockDetected', refs.ralphStatusBlockDetected);
|
||||
session.off('ralphCircuitBreakerUpdate', refs.ralphCircuitBreakerUpdate);
|
||||
session.off('ralphExitGateMet', refs.ralphExitGateMet);
|
||||
session.off('bashToolStart', refs.bashToolStart);
|
||||
session.off('bashToolEnd', refs.bashToolEnd);
|
||||
session.off('bashToolsUpdate', refs.bashToolsUpdate);
|
||||
}
|
||||
@@ -0,0 +1,488 @@
|
||||
/**
|
||||
* @fileoverview SSE stream manager — owns all SSE client state, broadcasting, and event batching.
|
||||
*
|
||||
* Extracted from server.ts for modularity. Handles:
|
||||
* - SSE client connection tracking with subscription filtering
|
||||
* - Backpressure-aware message delivery
|
||||
* - Terminal data batching with adaptive intervals (16-50ms for 60fps)
|
||||
* - Task update and session state batching
|
||||
* - Dead client cleanup and keepalive
|
||||
* - Cloudflare tunnel padding for proxy buffer flushing
|
||||
*
|
||||
* @dependencies CleanupManager (managed timers), config/server-timing (constants)
|
||||
* @consumedby web/server.ts (WebServer delegates all SSE operations here)
|
||||
*
|
||||
* @module web/sse-stream-manager
|
||||
*/
|
||||
|
||||
import type { FastifyReply } from 'fastify';
|
||||
import type { BackgroundTask } from '../session.js';
|
||||
import { CleanupManager, StaleExpirationMap } from '../utils/index.js';
|
||||
import { SseEvent } from './sse-events.js';
|
||||
import {
|
||||
TERMINAL_BATCH_INTERVAL,
|
||||
TASK_UPDATE_BATCH_INTERVAL,
|
||||
STATE_UPDATE_DEBOUNCE_INTERVAL,
|
||||
BATCH_FLUSH_THRESHOLD,
|
||||
SSE_PADDING_SIZE,
|
||||
INACTIVITY_TIMEOUT_MS,
|
||||
} from '../config/server-timing.js';
|
||||
|
||||
// SSE padding for Cloudflare tunnel buffer flushing.
|
||||
// Cloudflare quick tunnels buffer small SSE responses, causing lag for real-time events.
|
||||
// Appending SSE comment padding (ignored by EventSource) forces the proxy to flush.
|
||||
// Pre-computed once at startup to avoid repeated string allocation.
|
||||
const SSE_PADDING = ':' + 'p'.repeat(SSE_PADDING_SIZE) + '\n';
|
||||
|
||||
/** Dependencies injected by WebServer — keeps SseStreamManager decoupled from session/respawn state. */
|
||||
export interface SseStreamManagerDeps {
|
||||
/** Get session state with respawn info for session:updated broadcasts */
|
||||
getSessionStateWithRespawn(sessionId: string): unknown;
|
||||
}
|
||||
|
||||
export class SseStreamManager {
|
||||
// ─── SSE Client Tracking ────────────────────────────────
|
||||
/**
|
||||
* SSE clients mapped to their session subscription filter.
|
||||
* Value is a Set of session IDs the client wants events for,
|
||||
* or `null` meaning "receive all events" (backwards-compatible default).
|
||||
*/
|
||||
private sseClients: Map<FastifyReply, Set<string> | null> = new Map();
|
||||
/** SSE clients connecting from non-localhost (i.e. through tunnel) */
|
||||
private remoteSseClients: Set<FastifyReply> = new Set();
|
||||
/** Clients with backpressure — skip writes until 'drain' fires */
|
||||
private backpressuredClients: Set<FastifyReply> = new Set();
|
||||
|
||||
// ─── Tunnel State ───────────────────────────────────────
|
||||
/** Cached tunnel active state — updated on TunnelStarted/TunnelStopped to avoid getUrl() on every broadcast */
|
||||
private _isTunnelActive: boolean = false;
|
||||
|
||||
// ─── Terminal Batching ──────────────────────────────────
|
||||
private terminalBatches: Map<string, string[]> = new Map();
|
||||
private terminalBatchSizes: Map<string, number> = new Map(); // Running total avoids O(n) reduce per push
|
||||
private terminalBatchTimers: Map<string, NodeJS.Timeout> = new Map(); // Per-session timers (staggered flushes)
|
||||
// Adaptive batching: track rapid events to extend batch window (per-session)
|
||||
// StaleExpirationMap auto-cleans entries for sessions that stop generating output
|
||||
private lastTerminalEventTime: StaleExpirationMap<string, number>;
|
||||
|
||||
// ─── Event Batching ─────────────────────────────────────
|
||||
private taskUpdateBatches: Map<string, { sessionId: string; task: BackgroundTask }> = new Map();
|
||||
private taskUpdateBatchTimerId: string | null = null;
|
||||
// State update batching (reduce expensive toDetailedState() serialization)
|
||||
private stateUpdatePending: Set<string> = new Set();
|
||||
private stateUpdateTimerId: string | null = null;
|
||||
|
||||
// ─── Lifecycle ──────────────────────────────────────────
|
||||
private _isStopping: boolean = false;
|
||||
|
||||
constructor(
|
||||
private deps: SseStreamManagerDeps,
|
||||
private cleanup: CleanupManager
|
||||
) {
|
||||
this.lastTerminalEventTime = new StaleExpirationMap({
|
||||
ttlMs: INACTIVITY_TIMEOUT_MS, // 5 minutes - auto-expire stale session timing data
|
||||
refreshOnGet: false, // Don't refresh on reads, only on explicit sets
|
||||
});
|
||||
}
|
||||
|
||||
// ========== SSE Connection Management ==========
|
||||
|
||||
get clientCount(): number {
|
||||
return this.sseClients.size;
|
||||
}
|
||||
|
||||
get remoteClientCount(): number {
|
||||
return this.remoteSseClients.size;
|
||||
}
|
||||
|
||||
get isTunnelActive(): boolean {
|
||||
return this._isTunnelActive;
|
||||
}
|
||||
|
||||
setTunnelActive(active: boolean): void {
|
||||
this._isTunnelActive = active;
|
||||
}
|
||||
|
||||
addClient(reply: FastifyReply, sessionFilter: Set<string> | null, isRemote: boolean): void {
|
||||
this.sseClients.set(reply, sessionFilter);
|
||||
if (isRemote) {
|
||||
this.remoteSseClients.add(reply);
|
||||
}
|
||||
}
|
||||
|
||||
removeClient(reply: FastifyReply): void {
|
||||
this.sseClients.delete(reply);
|
||||
this.remoteSseClients.delete(reply);
|
||||
this.backpressuredClients.delete(reply);
|
||||
}
|
||||
|
||||
/** Send a single SSE event to a specific client. */
|
||||
sendSSE(reply: FastifyReply, event: string, data: unknown): void {
|
||||
try {
|
||||
reply.raw.write(`event: ${event}\ndata: ${JSON.stringify(data)}\n\n`);
|
||||
} catch {
|
||||
this.sseClients.delete(reply);
|
||||
this.remoteSseClients.delete(reply);
|
||||
}
|
||||
}
|
||||
|
||||
/** Send pre-formatted tunnel padding to a specific client. */
|
||||
sendPadding(reply: FastifyReply): void {
|
||||
if (!this._isTunnelActive) return;
|
||||
try {
|
||||
reply.raw.write(SSE_PADDING);
|
||||
} catch {
|
||||
/* client gone */
|
||||
}
|
||||
}
|
||||
|
||||
// Optimized: send pre-formatted SSE message to a client
|
||||
// Returns false if client is backpressured or dead
|
||||
private sendSSEPreformatted(reply: FastifyReply, message: string): void {
|
||||
// Skip backpressured clients to prevent unbounded memory growth.
|
||||
// Terminal data dropped here is recovered via session:needsRefresh on drain.
|
||||
if (this.backpressuredClients.has(reply)) return;
|
||||
|
||||
try {
|
||||
const ok = reply.raw.write(message);
|
||||
if (!ok) {
|
||||
// Buffer is full — mark as backpressured, resume on drain
|
||||
this.backpressuredClients.add(reply);
|
||||
reply.raw.once('drain', () => {
|
||||
this.backpressuredClients.delete(reply);
|
||||
// Client may have missed terminal data during backpressure.
|
||||
// Tell it to reload the active session's buffer to recover.
|
||||
try {
|
||||
const drainPadding = this._isTunnelActive ? SSE_PADDING : '';
|
||||
reply.raw.write(`event: ${SseEvent.SessionNeedsRefresh}\ndata: {}\n\n${drainPadding}`);
|
||||
} catch {
|
||||
/* client gone */
|
||||
}
|
||||
});
|
||||
}
|
||||
} catch {
|
||||
this.sseClients.delete(reply);
|
||||
this.remoteSseClients.delete(reply);
|
||||
this.backpressuredClients.delete(reply);
|
||||
}
|
||||
}
|
||||
|
||||
// ========== Broadcasting ==========
|
||||
|
||||
broadcast(event: string, data: unknown): void {
|
||||
// Skip serialization entirely when no clients are listening
|
||||
if (this.sseClients.size === 0) return;
|
||||
|
||||
// Performance optimization: serialize JSON once for all clients.
|
||||
// Only append Cloudflare tunnel padding for latency-sensitive events —
|
||||
// Recovery events need immediate proxy flush; low-frequency metadata events
|
||||
// (session:created, ralph:*, respawn:*, etc.) don't need padding.
|
||||
// Note: session:terminal has its own padding in flushSessionTerminalBatch().
|
||||
const needsPadding = this._isTunnelActive && event === SseEvent.SessionNeedsRefresh;
|
||||
const padding = needsPadding ? SSE_PADDING : '';
|
||||
let message: string;
|
||||
try {
|
||||
message = `event: ${event}\ndata: ${JSON.stringify(data)}\n\n` + padding;
|
||||
} catch (err) {
|
||||
// Handle circular references or non-serializable values
|
||||
console.error(`[Server] Failed to serialize SSE event "${event}":`, err);
|
||||
return;
|
||||
}
|
||||
// Extract sessionId from event data for subscription filtering.
|
||||
const eventSessionId = this.extractSessionId(event, data);
|
||||
|
||||
for (const [client, filter] of this.sseClients) {
|
||||
// No filter (null) = receive everything. Otherwise, skip if event is
|
||||
// session-scoped and the session isn't in the client's subscription set.
|
||||
if (filter && eventSessionId && !filter.has(eventSessionId)) continue;
|
||||
this.sendSSEPreformatted(client, message);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Extract the session ID from an event's data payload for subscription filtering.
|
||||
* Returns the sessionId string if the event is session-scoped, or null for global events.
|
||||
*/
|
||||
private extractSessionId(event: string, data: unknown): string | null {
|
||||
if (data == null || typeof data !== 'object') return null;
|
||||
const record = data as Record<string, unknown>;
|
||||
|
||||
// Most session-scoped events use `sessionId`
|
||||
if (typeof record.sessionId === 'string') return record.sessionId;
|
||||
|
||||
// Session lifecycle events (session:*) use `id` from the session state object
|
||||
if (typeof record.id === 'string' && event.startsWith('session:')) return record.id;
|
||||
|
||||
// No session ID found — treat as global event (sent to all clients)
|
||||
return null;
|
||||
}
|
||||
|
||||
// ========== Terminal Data Batching ==========
|
||||
|
||||
// Batch terminal data for better performance (60fps)
|
||||
// Uses per-session timers with adaptive intervals to prevent thundering herd:
|
||||
// each session flushes independently rather than all sessions flushing in one burst.
|
||||
batchTerminalData(sessionId: string, data: string): void {
|
||||
// Skip if server is stopping
|
||||
if (this._isStopping) return;
|
||||
|
||||
let chunks = this.terminalBatches.get(sessionId);
|
||||
if (!chunks) {
|
||||
chunks = [];
|
||||
this.terminalBatches.set(sessionId, chunks);
|
||||
}
|
||||
chunks.push(data);
|
||||
const prevSize = this.terminalBatchSizes.get(sessionId) ?? 0;
|
||||
const totalLength = prevSize + data.length;
|
||||
this.terminalBatchSizes.set(sessionId, totalLength);
|
||||
|
||||
// Adaptive batching: detect rapid events and extend batch window (per-session)
|
||||
const now = Date.now();
|
||||
const lastEvent = this.lastTerminalEventTime.get(sessionId) ?? 0;
|
||||
const eventGap = now - lastEvent;
|
||||
this.lastTerminalEventTime.set(sessionId, now);
|
||||
|
||||
// Adjust batch interval based on event frequency (per-session)
|
||||
// Rapid events (<10ms gap) = 50ms batch, moderate (<20ms) = 32ms, else 16ms
|
||||
let sessionInterval: number;
|
||||
if (eventGap > 0 && eventGap < 10) {
|
||||
sessionInterval = 50;
|
||||
} else if (eventGap > 0 && eventGap < 20) {
|
||||
sessionInterval = 32;
|
||||
} else {
|
||||
sessionInterval = TERMINAL_BATCH_INTERVAL;
|
||||
}
|
||||
|
||||
// Flush immediately if batch is large for responsiveness
|
||||
if (totalLength > BATCH_FLUSH_THRESHOLD) {
|
||||
const existingTimer = this.terminalBatchTimers.get(sessionId);
|
||||
if (existingTimer) {
|
||||
clearTimeout(existingTimer);
|
||||
this.terminalBatchTimers.delete(sessionId);
|
||||
}
|
||||
this.flushSessionTerminalBatch(sessionId);
|
||||
return;
|
||||
}
|
||||
|
||||
// Start per-session batch timer if not already running
|
||||
// Each session flushes independently — prevents one busy session from
|
||||
// forcing all sessions to flush at its rate (thundering herd)
|
||||
if (!this.terminalBatchTimers.has(sessionId)) {
|
||||
this.terminalBatchTimers.set(
|
||||
sessionId,
|
||||
setTimeout(() => {
|
||||
this.terminalBatchTimers.delete(sessionId);
|
||||
this.flushSessionTerminalBatch(sessionId);
|
||||
}, sessionInterval)
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
/** Flush a single session's batched terminal data */
|
||||
private flushSessionTerminalBatch(sessionId: string): void {
|
||||
if (this._isStopping) {
|
||||
this.terminalBatches.delete(sessionId);
|
||||
this.terminalBatchSizes.delete(sessionId);
|
||||
return;
|
||||
}
|
||||
const chunks = this.terminalBatches.get(sessionId);
|
||||
if (chunks && chunks.length > 0) {
|
||||
// Join chunks only at flush time (avoids O(n^2) string concatenation in batchTerminalData)
|
||||
const data = chunks.join('');
|
||||
// Wrap batched output in DEC 2026 synchronized output markers so xterm.js
|
||||
// renders the entire batch atomically. Ink spinner frames (cursor-up + redraw)
|
||||
// do NOT emit their own 2026 markers, so without this wrapper each partial
|
||||
// cursor update renders individually, causing visible flicker.
|
||||
// xterm.js 6.0+ handles DEC 2026 natively: it buffers everything between
|
||||
// 2026h/2026l and renders in one pass.
|
||||
const syncData = '\x1b[?2026h' + data + '\x1b[?2026l';
|
||||
// Fast path: build SSE message directly without JSON.stringify on wrapper object.
|
||||
// Only the terminal data string needs escaping; sessionId is a UUID (safe to template).
|
||||
const escapedData = JSON.stringify(syncData);
|
||||
// Append tunnel padding for immediate Cloudflare proxy flush —
|
||||
// terminal data is high-frequency and latency-sensitive.
|
||||
const padding = this._isTunnelActive ? SSE_PADDING : '';
|
||||
const message = `event: session:terminal\ndata: {"id":"${sessionId}","data":${escapedData}}\n\n` + padding;
|
||||
for (const [client, filter] of this.sseClients) {
|
||||
// Skip clients that have a session filter and aren't subscribed to this session
|
||||
if (filter && !filter.has(sessionId)) continue;
|
||||
this.sendSSEPreformatted(client, message);
|
||||
}
|
||||
}
|
||||
this.terminalBatches.delete(sessionId);
|
||||
this.terminalBatchSizes.delete(sessionId);
|
||||
}
|
||||
|
||||
// ========== Task Update Batching ==========
|
||||
|
||||
// Batch task:updated events at 100ms - only send latest update per task
|
||||
// Key is sessionId:taskId to avoid collisions when multiple tasks update concurrently
|
||||
batchTaskUpdate(sessionId: string, task: BackgroundTask): void {
|
||||
// Skip if server is stopping
|
||||
if (this._isStopping) return;
|
||||
|
||||
// Use composite key to avoid losing updates when multiple tasks update in same batch window
|
||||
const key = `${sessionId}:${task.id}`;
|
||||
this.taskUpdateBatches.set(key, { sessionId, task });
|
||||
|
||||
if (!this.taskUpdateBatchTimerId) {
|
||||
this.taskUpdateBatchTimerId = this.cleanup.setTimeout(
|
||||
() => {
|
||||
this.taskUpdateBatchTimerId = null;
|
||||
this.flushTaskUpdateBatches();
|
||||
},
|
||||
TASK_UPDATE_BATCH_INTERVAL,
|
||||
{ description: 'task update batch flush' }
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
private flushTaskUpdateBatches(): void {
|
||||
// Skip if server is stopping (timer may have been queued before stop() was called)
|
||||
if (this._isStopping) {
|
||||
this.taskUpdateBatches.clear();
|
||||
return;
|
||||
}
|
||||
for (const [, { sessionId, task }] of this.taskUpdateBatches) {
|
||||
this.broadcast(SseEvent.TaskUpdated, { sessionId, task });
|
||||
}
|
||||
this.taskUpdateBatches.clear();
|
||||
}
|
||||
|
||||
// ========== Session State Batching ==========
|
||||
|
||||
/**
|
||||
* Debounce expensive session:updated broadcasts.
|
||||
* Instead of calling toDetailedState() on every event, batch requests
|
||||
* and only serialize once per STATE_UPDATE_DEBOUNCE_INTERVAL.
|
||||
*/
|
||||
broadcastSessionStateDebounced(sessionId: string): void {
|
||||
// Skip if server is stopping
|
||||
if (this._isStopping) return;
|
||||
|
||||
this.stateUpdatePending.add(sessionId);
|
||||
|
||||
if (!this.stateUpdateTimerId) {
|
||||
this.stateUpdateTimerId = this.cleanup.setTimeout(
|
||||
() => {
|
||||
this.stateUpdateTimerId = null;
|
||||
this.flushStateUpdates();
|
||||
},
|
||||
STATE_UPDATE_DEBOUNCE_INTERVAL,
|
||||
{ description: 'state update debounce flush' }
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
private flushStateUpdates(): void {
|
||||
// Skip if server is stopping (timer may have been queued before stop() was called)
|
||||
if (this._isStopping) {
|
||||
this.stateUpdatePending.clear();
|
||||
return;
|
||||
}
|
||||
for (const sessionId of this.stateUpdatePending) {
|
||||
// Single expensive serialization per batch interval
|
||||
const state = this.deps.getSessionStateWithRespawn(sessionId);
|
||||
if (state) {
|
||||
this.broadcast(SseEvent.SessionUpdated, state);
|
||||
}
|
||||
}
|
||||
this.stateUpdatePending.clear();
|
||||
}
|
||||
|
||||
// ========== Client Health ==========
|
||||
|
||||
/**
|
||||
* Clean up dead SSE clients and send keep-alive comments.
|
||||
* Keep-alive prevents proxy/load-balancer timeouts on idle connections.
|
||||
* Dead client cleanup prevents memory leaks from abruptly terminated connections.
|
||||
*/
|
||||
cleanupDeadClients(): void {
|
||||
const deadClients: FastifyReply[] = [];
|
||||
|
||||
for (const [client] of this.sseClients) {
|
||||
try {
|
||||
// Check if the underlying socket is still writable
|
||||
const socket = client.raw.socket;
|
||||
if (!socket || socket.destroyed || !socket.writable) {
|
||||
deadClients.push(client);
|
||||
} else {
|
||||
// Send SSE comment as keep-alive. Only add padding when tunnel is
|
||||
// active — it flushes Cloudflare proxy buffers but wastes bandwidth
|
||||
// for direct/Tailscale connections.
|
||||
const ka = this._isTunnelActive ? ':keepalive\n' + SSE_PADDING : ':keepalive\n\n';
|
||||
client.raw.write(ka);
|
||||
}
|
||||
} catch {
|
||||
// Error accessing socket means client is dead
|
||||
deadClients.push(client);
|
||||
}
|
||||
}
|
||||
|
||||
// Remove dead clients
|
||||
for (const client of deadClients) {
|
||||
this.sseClients.delete(client);
|
||||
this.remoteSseClients.delete(client);
|
||||
this.backpressuredClients.delete(client);
|
||||
}
|
||||
|
||||
if (deadClients.length > 0) {
|
||||
console.log(`[Server] Cleaned up ${deadClients.length} dead SSE client(s)`);
|
||||
}
|
||||
}
|
||||
|
||||
// ========== Session Cleanup ==========
|
||||
|
||||
/** Clean up all batching state for a session (call on session exit or deletion). */
|
||||
cleanupSessionBatches(sessionId: string): void {
|
||||
this.terminalBatches.delete(sessionId);
|
||||
this.terminalBatchSizes.delete(sessionId);
|
||||
const batchTimer = this.terminalBatchTimers.get(sessionId);
|
||||
if (batchTimer) {
|
||||
clearTimeout(batchTimer);
|
||||
this.terminalBatchTimers.delete(sessionId);
|
||||
}
|
||||
this.taskUpdateBatches.delete(sessionId);
|
||||
this.stateUpdatePending.delete(sessionId);
|
||||
this.lastTerminalEventTime.delete(sessionId);
|
||||
}
|
||||
|
||||
// ========== Lifecycle ==========
|
||||
|
||||
setStopping(): void {
|
||||
this._isStopping = true;
|
||||
}
|
||||
|
||||
/** Graceful shutdown: notify clients, close connections, clear all state. */
|
||||
stop(): void {
|
||||
this._isStopping = true;
|
||||
|
||||
// Gracefully close all SSE connections before clearing
|
||||
for (const [client] of this.sseClients) {
|
||||
try {
|
||||
// Send a final event to notify clients of shutdown
|
||||
this.sendSSE(client, 'server:shutdown', { reason: 'Server stopping' });
|
||||
client.raw.end();
|
||||
} catch {
|
||||
// Client may already be disconnected
|
||||
}
|
||||
}
|
||||
this.sseClients.clear();
|
||||
this.remoteSseClients.clear();
|
||||
this.backpressuredClients.clear();
|
||||
|
||||
// Clear per-session batch timers
|
||||
for (const timer of this.terminalBatchTimers.values()) {
|
||||
clearTimeout(timer);
|
||||
}
|
||||
this.terminalBatchTimers.clear();
|
||||
this.terminalBatches.clear();
|
||||
this.terminalBatchSizes.clear();
|
||||
|
||||
this.taskUpdateBatches.clear();
|
||||
this.stateUpdatePending.clear();
|
||||
|
||||
// Dispose StaleExpirationMap (stops internal cleanup timer)
|
||||
this.lastTerminalEventTime.dispose();
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user