perf: Operation Lightspeed — 5 parallel performance optimizations

1. Session-scoped SSE subscriptions: server filters events by session ID,
   clients can subscribe via ?sessions=id1,id2 (backwards-compatible)
2. Lazy xterm.js for subagent windows: terminals created on restore,
   disposed on minimize — saves ~3.75MB DOM at 50 agents
3. Targeted badge updates: badge count changes update the <span> directly
   instead of rebuilding the entire session tab sidebar (O(1) vs O(n))
4. Conditional SSE padding: 8KB Cloudflare padding only on session:terminal
   and session:needsRefresh, not every event (~70% bandwidth reduction)
5. Canvas renderer on mobile: skip WebGL addon on mobile devices to reduce
   GPU pressure and prevent context loss on weaker mobile GPUs

All 5 implemented in parallel via isolated git worktrees, merged conflict-free.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
This commit is contained in:
arkon
2026-03-07 06:14:16 +01:00
co-authored by Claude Opus 4.6
parent 415f02e680
commit 1f30ed445c
5 changed files with 1138 additions and 37 deletions
+423
View File
@@ -0,0 +1,423 @@
# Performance Analysis & Optimization Opportunities
**Date**: 2026-03-07
**Scope**: Full-stack performance analysis — backend PTY handling, SSE broadcasting, frontend terminal rendering, local echo overlay, DOM updates, config/scaling limits.
**Constraint**: All recommendations preserve existing functionality including local echo, backpressure, anti-flicker pipeline, and mobile support.
---
## Executive Summary
The codebase is already well-optimized in critical paths. The multi-layer backpressure system, adaptive terminal batching, DEC 2026 sync markers, and incremental state serialization are strong. The main opportunities are in **reducing unnecessary work** (SSE filtering, DOM rebuilds, lazy terminal init) rather than algorithmic changes.
**Top 5 high-impact opportunities:**
| # | Optimization | Impact | Risk | Effort |
|---|-------------|--------|------|--------|
| 1 | Session-scoped SSE subscriptions | Bandwidth -60-80%, CPU -40% | Medium | Medium |
| 2 | Lazy xterm.js for minimized subagent windows | Memory -3.5MB at 50 agents | Low | Low |
| 3 | Targeted badge update (skip full tab rebuild) | Eliminates O(n) reflow on badge change | Low | Low |
| 4 | Conditional SSE padding (tunnel-only, terminal-only) | Bandwidth -70% when tunneled | Low | Low |
| 5 | Canvas renderer on mobile | GPU pressure reduction, battery savings | Low | Low |
---
## 1. SSE Broadcasting
### Current State
- **92 event types** broadcast to all connected clients (max 100)
- Single `JSON.stringify()` per event, shared across all clients (efficient)
- **No per-client filtering** — every client receives every event regardless of which session they're viewing
- 8KB padding appended to **every** event when tunnel is active (forces Cloudflare proxy flush)
- Backpressure: clients marked as backpressured if `reply.raw.write()` returns false; recovery via `session:needsRefresh`
### Bottlenecks
**B1: No session-scoped SSE subscriptions** (`server.ts:1986`)
- Client viewing session A still receives all events for sessions B through T
- With 20 active sessions, ~95% of terminal events are irrelevant to any given client
- Cost: wasted bandwidth, CPU for JSON parsing, and event handler dispatch on client
**B2: Unconditional 8KB padding** (`server.ts:1977`)
- Every event gets 8KB comment padding when tunnel is active
- A `task:updated` event (~200 bytes payload) becomes ~8.2KB
- High-frequency events like `session:terminal` need the padding; low-frequency events like `session:created` don't
### Recommendations
**R1: Session-scoped SSE subscriptions** (High impact)
- Add `?sessions=id1,id2` query param to `/api/events` SSE endpoint
- Server filters events by session ID before broadcasting
- Client subscribes to active session + "global" events (session lifecycle, system)
- Re-subscribes on tab switch (or subscribe to all with client-side filter as fallback)
- **Savings**: ~80% bandwidth reduction for single-session viewers; ~60% for multi-session dashboards
**R2: Tiered SSE padding** (Medium impact)
- Only pad `session:terminal` events and SSE heartbeats (the two that need proxy flush)
- Skip padding for low-frequency structural events (`session:created`, `task:updated`, etc.)
- **Savings**: ~70% padding overhead reduction; terminal events already large enough to flush
---
## 2. Terminal Rendering
### Current State (Well-Optimized)
- **6-layer anti-flicker pipeline**: Server batching (adaptive 16-50ms) → DEC 2026 sync wrap → single JSON serialize → client rAF batching → sync segment parser → chunked buffer loading (32KB/frame)
- **64KB/frame write budget** with DEC 2026 sync-segment awareness (prevents 141KB single-frame freezes)
- **3-layer backpressure**: SSE cap (128KB queued → drop + refresh), frame budget (64KB/frame), chunked restore (32KB/frame)
- WebGL renderer enabled by default with canvas fallback on context loss
- Typical latency: 16-32ms; worst case: ~115ms (50ms server batch + 50ms sync wait + 16ms rAF)
### Bottlenecks
**B3: WebGL on mobile** (`app.js:627-637`)
- Mobile GPUs are weaker; WebGL context loss more likely on low-end devices
- Canvas renderer is sufficient for mobile (typically 1 session, smaller viewport)
**B4: Static scrollback for all sessions** (`app.js:572`)
- Default 5000 lines scrollback for all sessions regardless of activity level
- Heavy output sessions (build logs, test runners) accumulate large scroll buffers
**B5: No addon lazy loading**
- FitAddon, Unicode11Addon, and WebGLAddon all loaded at terminal init
- Unicode11Addon only needed for CJK content; WebGLAddon is large
### Recommendations
**R3: Force canvas renderer on mobile** (Low risk)
- Detect `MobileDetection.isMobile()` and skip WebGL addon loading
- Reduces GPU memory pressure, prevents context loss crashes
- Mobile typically has 1-2 sessions — canvas performance is more than adequate
**R4: Dynamic scrollback based on session activity** (Low risk)
- Active sessions (working state): 5000 lines (current default)
- Inactive/idle sessions: reduce to 2000 lines
- Restore on session select (fetch from server buffer)
- **Savings**: ~60% scrollback memory for idle sessions
**R5: Lazy-load Unicode11Addon** (Low risk)
- Only load when CJK content is detected in terminal output
- Detection: check for characters in CJK Unicode ranges during ANSI stripping (already iterating)
- Most sessions never need it
---
## 3. DOM & Session Tab Rendering
### Current State
- Session tabs use **intelligent incremental updates** with debounced 100ms rendering
- Incremental path: only updates changed properties (classes, textContent, badges) when session list is stable
- Full rebuild path: triggered when sessions added/removed **or badge count changes**
- Subagent windows: per-window xterm.js instances, even when minimized
### Bottlenecks
**B6: Badge count change triggers full tab rebuild** (`app.js:3207-3209`)
- A single subagent badge increment on one tab triggers `_fullRenderSessionTabs()` — rebuilds entire sidebar HTML via `innerHTML =`
- With 20 sessions, this is an O(n) reflow for a single badge number change
- Badge changes are frequent during active subagent work
**B7: Minimized subagent windows retain xterm.js instances** (`subagent-windows.js`)
- 50 subagent windows × ~75KB per xterm.js instance = ~3.75MB DOM memory
- Minimized windows are invisible but their terminals remain in DOM
- xterm.js instances continue processing resize events even when hidden
**B8: `backdrop-filter: blur()` on overlays** (`styles.css:2246-2247, 3098`)
- Forces new stacking context, disables browser compositing optimizations
- 50-100ms layout thrashing on modal open/close
- Only 2 uses, but they're on frequently toggled overlays
### Recommendations
**R6: Targeted badge update without full rebuild** (Low risk)
- When badge count changes but session list is stable, update only the badge `<span>` textContent
- Keep incremental path for badge changes; only use full rebuild for structural changes (add/remove sessions)
- **Savings**: Eliminates O(n) reflow per badge change; reduces to O(1) targeted update
**R7: Lazy xterm.js initialization for subagent windows** (Medium impact)
- Only create xterm.js Terminal instance when window is restored/maximized
- On minimize: serialize terminal buffer, dispose Terminal instance, keep buffer in memory
- On restore: create new Terminal, write buffer back
- **Savings**: ~3.5MB DOM reduction at 50 minimized agents; eliminates hidden resize processing
- **Trade-off**: ~200-500ms restore delay (buffer write), mitigated by chunked loading
**R8: Replace `backdrop-filter: blur()` with `background: rgba()`** (Low risk)
- Use semi-transparent background instead of blur effect
- Or use `will-change: transform` hint if blur is kept
- **Savings**: Eliminates forced recomposition layer; 50-100ms faster overlay open
---
## 4. Backend PTY & State Management
### Current State (Excellent)
- **BufferAccumulator**: Array-based chunking with lazy join on read — avoids O(n) string concatenation
- **ANSI stripping**: Throttled at 150ms intervals with lazy evaluation (not per-chunk)
- **State persistence**: 500ms debounce + incremental JSON caching per session (only dirty sessions re-serialized)
- **Expensive parsers**: Throttled to 150ms window, accumulated data capped at 64KB
- **Memory**: All buffers have hard limits (2MB terminal, 1MB text, 1000 messages, 64KB line buffer)
### Bottlenecks
**B9: Pending clean data cap at 64KB** (`session.ts:1097-1133`)
- Between 150ms processing windows, raw PTY data accumulates in `_pendingCleanData`
- Capped at 64KB — excess data rolls off (old data discarded)
- During heavy output (large build logs), this means parsers may miss content
- Acceptable trade-off for performance, but worth documenting
**B10: `LRUMap.delete()` is O(n) worst case** (`utils/lru-map.ts:137-138`)
- When deleting the newest entry, iterates all keys to find new newest
- Rare in practice (delete is uncommon; set/get are hot paths)
- Could matter during mass cleanup of 500 agents
### Recommendations
**R9: Consider adaptive pending data cap** (Low priority)
- During idle detection (critical to get right), increase cap to 128KB
- During active working state, keep at 64KB (parsers less critical)
- **Benefit**: More accurate idle detection during heavy output
**R10: Track second-newest in LRUMap** (Low priority)
- Maintain a `_secondNewestKey` alongside `_newestKey`
- On delete of newest, promote second-newest without iteration
- Only matters at scale (500+ agents with frequent eviction)
---
## 5. Local Echo & Input Path
### Current State (Well-Designed)
- **DOM overlay approach** — `<span>` elements in `.xterm-screen` at z-index 7, completely independent of `terminal.write()`
- **Render caching**: `_lastRenderKey` includes text, position, column offsets — skips redundant re-renders
- **Input flow**: Char accumulation → Enter triggers flush → 80ms delay before `\r` (ensures text reaches PTY first)
- **Tab completion**: Baseline snapshot → detect buffer change → 300ms fallback timer
- **CJK support**: Per-character width detection with `terminal.unicode.getStringCellWidth()` preferred, manual fallback
- **Prompt detection**: Bottom-up line scan, O(rows) — cached position, column-lock prevents jitter
### Bottlenecks
**B11: tmux send-keys latency** (~50-100ms per input)
- Each `writeViaMux()` spawns a child process (`tmux send-keys`)
- Text and Enter sent separately with 50ms delay between
- For rapid typing: characters batch before Enter, so overhead is per-command not per-keystroke
- **Acceptable trade-off** for session persistence (tmux survives server restarts)
**B12: 80ms delay between text flush and Enter** (`app.js:872-875`)
- Intentional: ensures text reaches PTY before Enter, preventing Ink from processing empty input
- Adds 80ms to perceived Enter-to-response latency
- Could potentially be reduced with acknowledgment-based approach
**B13: Scroll listener on terminal viewport** (`zerolag-input-addon.ts:139`)
- 50ms debounced re-render on scroll — acceptable but fires frequently during heavy output
- Overlay hidden when scrolled up (correct behavior), shown when at bottom
### Recommendations
**R11: Reduce Enter delay from 80ms to 50ms** (Low risk, test carefully)
- The tmux `send-keys` already has 50ms internal delay
- Combined with network latency, 80ms client-side may be excessive
- Test with Ink-heavy sessions (Claude Code's status bar) — if text arrives before Enter at 50ms, reduce
- **Savings**: 30ms perceived latency reduction per command
**R12: Batch tmux send-keys via stdin pipe** (Medium effort, high impact for rapid input)
- Instead of spawning `tmux send-keys` per input, maintain a persistent connection
- Use `tmux -C` (control mode) for programmatic interaction without child process spawning
- **Savings**: Eliminate ~50-100ms process spawn overhead per input
- **Risk**: Control mode has different semantics; needs careful testing with session persistence
**R13: Skip overlay re-render during heavy output scroll** (Low risk)
- When terminal is receiving >10KB/s output, hide overlay entirely (user isn't typing during heavy output)
- Re-show overlay after 500ms of output silence
- **Savings**: Eliminates unnecessary DOM overlay re-renders during build logs / test output
---
## 6. Polling & File Watchers
### Current State
- **SubagentWatcher**: 1s base poll, full scan throttled to every 5s, fs.watch() on known directories
- **TranscriptWatcher**: 1 per session, fs.watch() primary with 1s poll fallback
- **ImageWatcher**: chokidar per session with 100ms stability poll, burst limit 20/10s
- **TeamWatcher**: chokidar primary with 30s poll fallback, LRU caches (50 teams, 200 tasks)
- **RalphTracker**: Todo cleanup every 5 minutes
### Scaling Profile (20 sessions)
| Component | Instances | Frequency | Total ops/sec |
|-----------|-----------|-----------|---------------|
| SubagentWatcher | 1 (global) | Full scan every 5s | 0.2/s |
| TranscriptWatcher | 20 | 1s poll (fallback) | 20/s max |
| ImageWatcher | 20 | 100ms poll (during writes only) | 200/s burst |
| TeamWatcher | 1 (global) | 30s poll (fallback) | 0.03/s |
| SSE heartbeat | 1 (global) | 15s | 0.07/s |
| SSE dead client check | 1 (global) | 30s | 0.03/s |
| Mux stats collection | 1 (global) | 2s | 0.5/s |
| **Total steady-state** | | | **~21/s** |
### Recommendations
**R14: Increase TranscriptWatcher poll interval to 2s** (Low risk)
- Transcript changes are infrequent (new messages every few seconds at most)
- fs.watch() is the primary mechanism; polling is fallback
- **Savings**: Halves fallback filesystem checks (20/s → 10/s for 20 sessions)
**R15: Share chokidar instances for co-located session directories** (Medium effort)
- Sessions in the same parent directory could share a single chokidar watcher with depth:3
- Common case: multiple sessions in `~/projects/foo/` — one watcher covers all
- **Savings**: Reduce chokidar instances from 20 to ~5-10 for typical workloads
---
## 7. Frontend Asset Delivery
### Current State
- **app.js**: 12,027 lines (source) → esbuild minified → gzip/brotli compressed (~30-40KB gzipped)
- **Static caching**: `maxAge: '1y'` via `@fastify/static`
- **Service worker**: Push notification handler only — no asset caching
- **No code splitting**: Single monolithic app.js bundle
### Bottlenecks
**B14: No cache-busting mechanism**
- `maxAge: '1y'` means browsers cache aggressively
- After deployment, users need `Ctrl+Shift+R` to see updates
- No content hash in filenames or ETags for automatic invalidation
**B15: Monolithic app.js**
- All 12K lines loaded on initial page load regardless of which features are used
- Ralph wizard, plan orchestrator UI, team management — all loaded upfront
- Mobile loads the same bundle as desktop
### Recommendations
**R16: Add content hash to asset filenames** (Medium impact)
- Build step: rename `app.js` → `app.[hash].js`
- Generate a manifest or inject hash into HTML template
- Keep `maxAge: '1y'` — cache invalidation happens via filename change
- **Savings**: Eliminates stale cache issues after deployment; removes need for manual hard refresh
**R17: Code-split app.js into core + feature modules** (High effort, medium impact)
- Core (~4K lines): terminal, SSE, session management, tabs, input handling
- Deferred (~8K lines): Ralph wizard, plan UI, team management, subagent windows, image viewer
- Load deferred modules on first use via dynamic `import()` or lazy `<script>` injection
- **Savings**: ~60% reduction in initial load size; faster time-to-interactive
- **Risk**: Complexity increase; need to handle loading states for deferred features
- **Note**: May not be worth the effort given the app is already gzipped to ~30-40KB
---
## 8. CSS Performance
### Current State
- **styles.css**: 7,153 lines with ~45 box-shadow uses, 2 backdrop-filter uses
- Animations: GPU-accelerated keyframes for pulsing alerts, loading spinners
- Z-index layering: well-organized (subagent 1000, plan 1100, log 2000, image 3000, overlay 7)
### Recommendations
**R18: Replace backdrop-filter with opaque overlay** (Low risk, covered in R8)
**R19: Use `contain: content` on subagent windows** (Low risk)
- Add CSS containment to subagent window containers
- Prevents layout changes inside windows from triggering reflow on parent
- Especially valuable with 50 windows: changes in one window won't invalidate others
- ```css
.subagent-window { contain: content; }
```
- **Savings**: Reduces layout recalculation scope from global to per-window
**R20: Use `content-visibility: auto` on off-screen subagent windows** (Low risk)
- Browser skips rendering of off-screen windows entirely
- Combined with `contain-intrinsic-size` to prevent layout shift
- ```css
.subagent-window.minimized { content-visibility: hidden; }
```
- **Savings**: Browser skips paint/layout for minimized windows; complements R7
---
## 9. Memory & Scaling Limits
### Current Budget (20 sessions)
| Component | Per Session | Total | Status |
|-----------|-----------|-------|--------|
| Terminal buffer | 2MB | 40MB | Hard-limited, auto-trim |
| Text output | 1MB | 20MB | Hard-limited, auto-trim |
| Messages | ~1MB | 20MB | Capped at 1000, trims to 800 |
| Respawn buffer | 1MB | 20MB | Hard-limited |
| **Buffers total** | | **100MB** | Acceptable |
| TranscriptWatcher | ~100KB | 2MB | |
| ImageWatcher | ~50KB | 1MB | |
| SubagentWatcher | ~500KB | 500KB | Global |
| Frontend terminal cache | ~256KB | 5MB | LRU, max 20 entries |
| **Total estimated** | | **~110MB** | Comfortable |
### At Max Scale (50 sessions)
- Buffers: ~250MB
- Watchers: ~5MB
- **Total: ~255MB** + Node.js overhead — acceptable on modern hardware
### Potential Leak Vectors (All Mitigated)
- `_shortIdCache` in server — unbounded Map, but entries are tiny (string→string); grows at O(sessions created), not O(events)
- All CleanupManager-registered resources tracked and disposed on session stop
- `isStopped` guard prevents new timers after session cleanup
---
## 10. Implementation Priority Matrix
### Phase 1 — Quick Wins (1-2 hours each, low risk)
| # | Optimization | Files to Change |
|---|-------------|-----------------|
| R6 | Targeted badge update | `app.js` (3207-3209) |
| R3 | Canvas renderer on mobile | `app.js` (627-637) |
| R8 | Replace backdrop-filter blur | `styles.css` (2246, 3098) |
| R19 | CSS containment on subagent windows | `styles.css` |
| R20 | `content-visibility: hidden` on minimized windows | `styles.css` |
### Phase 2 — Medium Effort (half-day each)
| # | Optimization | Files to Change |
|---|-------------|-----------------|
| R2 | Tiered SSE padding | `server.ts` (broadcast function) |
| R7 | Lazy xterm.js for minimized subagents | `subagent-windows.js` |
| R11 | Reduce Enter delay to 50ms | `app.js` (872-875), test with Ink |
| R14 | TranscriptWatcher 2s poll | `transcript-watcher.ts` |
| R16 | Content-hash asset filenames | `build.mjs`, `server.ts` |
### Phase 3 — Larger Initiatives (1-2 days each)
| # | Optimization | Files to Change |
|---|-------------|-----------------|
| R1 | Session-scoped SSE subscriptions | `server.ts`, `app.js` (SSE connect) |
| R5 | Lazy Unicode11Addon loading | `app.js`, build pipeline |
| R12 | Persistent tmux control mode | `tmux-manager.ts` |
| R17 | Code-split app.js | `app.js`, `build.mjs`, HTML template |
### Not Recommended (Low ROI or High Risk)
| # | Why Not |
|---|---------|
| R4 | Dynamic scrollback adds complexity; memory savings marginal vs total budget |
| R9 | Adaptive pending data cap adds state; current 64KB cap rarely matters |
| R10 | LRUMap.delete() O(n) is theoretical; never triggered at current scale |
| R15 | Shared chokidar instances add directory-matching complexity for minimal gain |
---
## Appendix: Key File Locations
| Area | File | Key Lines |
|------|------|-----------|
| SSE broadcast | `src/web/server.ts` | 1961-1989 (broadcast), 1934-1959 (backpressure) |
| Terminal batching | `src/web/server.ts` | 1994-2048 (per-session adaptive batching) |
| Frame budget | `src/web/public/app.js` | 1370-1478 (flushPendingWrites, 64KB cap) |
| Flicker filter | `src/web/public/app.js` | 1176-1255 (50ms sync wait, 256KB safety) |
| Tab rendering | `src/web/public/app.js` | 3108-3357 (incremental + full rebuild) |
| Tab switching | `src/web/public/app.js` | 3560-3760 (cache + chunked load + deferred UI) |
| Local echo | `packages/xterm-zerolag-input/src/` | All files (overlay, prompt, CJK) |
| Local echo integration | `src/web/public/app.js` | 640, 815-988 (input flow) |
| Subagent windows | `src/web/public/subagent-windows.js` | Full file (window mgmt, drag, minimize) |
| State persistence | `src/state-store.ts` | 161-250 (debounced save, incremental JSON) |
| Buffer accumulator | `src/utils/buffer-accumulator.ts` | Full file (array chunks, lazy join) |
| PTY handling | `src/session.ts` | 1046-1133 (data flow), 1173-1230 (parsing) |
| Config limits | `src/config/` | 9 files (buffer, map, timing, auth, etc.) |
| Anti-flicker docs | `docs/terminal-anti-flicker.md` | Architecture reference |
| CSS | `src/web/public/styles.css` | 2246 (backdrop-filter), full file |
| Build pipeline | `scripts/build.mjs` | 59-68 (minify + compress) |
+59 -9
View File
@@ -624,7 +624,8 @@ class CodemanApp {
// oversized terminal.write() calls that triggered the stalls.
// Disable with ?nowebgl URL param if GPU issues return.
this._webglAddon = null;
if (!new URLSearchParams(location.search).has('nowebgl') && typeof WebglAddon !== 'undefined') {
const isMobile = MobileDetection.getDeviceType() === 'mobile';
if (!isMobile && !new URLSearchParams(location.search).has('nowebgl') && typeof WebglAddon !== 'undefined') {
try {
this._webglAddon = new WebglAddon.WebglAddon();
this._webglAddon.onContextLoss(() => {
@@ -3198,15 +3199,38 @@ class CodemanApp {
return;
}
// Update subagent badge - check if count changed
// Update subagent badge - targeted update without full rebuild
const subagentBadgeEl = tab.querySelector('.tab-subagent-badge');
const minimizedAgents = this.minimizedSubagents.get(id);
const minimizedCount = minimizedAgents?.size || 0;
const currentCount = subagentBadgeEl ? parseInt(subagentBadgeEl.querySelector('.subagent-count')?.textContent || '0') : 0;
if (minimizedCount !== currentCount) {
// Count changed - need full rebuild for dropdown update
this._fullRenderSessionTabs();
return;
if (minimizedCount > 0 && subagentBadgeEl) {
// Badge exists and still has agents - update label and dropdown in-place
const labelEl = subagentBadgeEl.querySelector('.subagent-label');
const newLabel = minimizedCount === 1 ? 'AGENT' : `AGENTS (${minimizedCount})`;
if (labelEl && labelEl.textContent !== newLabel) {
labelEl.textContent = newLabel;
}
// Rebuild dropdown items (agent list may have changed)
const dropdownEl = subagentBadgeEl.querySelector('.subagent-dropdown');
if (dropdownEl) {
const newBadgeHtml = this.renderSubagentTabBadge(id, minimizedAgents);
const temp = document.createElement('div');
temp.innerHTML = newBadgeHtml;
const newDropdown = temp.querySelector('.subagent-dropdown');
if (newDropdown) {
dropdownEl.innerHTML = newDropdown.innerHTML;
}
}
} else if (minimizedCount > 0 && !subagentBadgeEl) {
// Need to add badge - insert before gear icon
const badgeHtml = this.renderSubagentTabBadge(id, minimizedAgents);
const gearEl = tab.querySelector('.tab-gear');
if (gearEl) {
gearEl.insertAdjacentHTML('beforebegin', badgeHtml);
}
} else if (minimizedCount === 0 && subagentBadgeEl) {
// Count went to 0 - remove badge
subagentBadgeEl.remove();
}
}
} else {
@@ -9599,10 +9623,16 @@ class CodemanApp {
// Show window (unless it was minimized by user)
if (!windowInfo.minimized) {
windowInfo.element.style.display = 'flex';
// Lazily re-create teammate terminal if it was disposed when hidden
if (windowInfo._lazyTerminal) {
this._restoreTeammateTerminalFromLazy(agentId);
}
}
windowInfo.hidden = false;
} else {
// Hide window (but don't close it)
// Dispose teammate terminal to free memory while hidden on inactive tab
this._disposeTeammateTerminalForMinimize(agentId);
windowInfo.element.style.display = 'none';
windowInfo.hidden = true;
}
@@ -9685,6 +9715,8 @@ class CodemanApp {
minimizeSubagentWindow(agentId) {
const windowData = this.subagentWindows.get(agentId);
if (windowData) {
// Dispose teammate terminal on minimize to free DOM/memory (lazy re-creation on restore)
this._disposeTeammateTerminalForMinimize(agentId);
windowData.element.style.display = 'none';
windowData.minimized = true;
this.updateConnectionLines();
@@ -9695,6 +9727,10 @@ class CodemanApp {
// Debounced wrapper — coalesces rapid subagent events (tool_call, progress,
// message) into a single DOM update per 100ms per agent window.
scheduleSubagentWindowRender(agentId) {
// Skip DOM updates for windows with lazy (disposed) terminals — they're minimized
const windowData = this.subagentWindows.get(agentId);
if (windowData?.minimized) return;
if (!this._subagentWindowRenderTimeouts) this._subagentWindowRenderTimeouts = new Map();
if (this._subagentWindowRenderTimeouts.has(agentId)) {
clearTimeout(this._subagentWindowRenderTimeouts.get(agentId));
@@ -9708,6 +9744,9 @@ class CodemanApp {
renderSubagentWindowContent(agentId) {
// Skip if this window has a live terminal (don't overwrite xterm with activity HTML)
if (this.teammateTerminals.has(agentId)) return;
// Skip if this window has a lazy (disposed) terminal — it will be re-created on restore
const windowData = this.subagentWindows.get(agentId);
if (windowData?._lazyTerminal) return;
const body = document.getElementById(`subagent-window-body-${agentId}`);
if (!body) return;
@@ -10110,8 +10149,19 @@ class CodemanApp {
resizeObserver.observe(win);
this.subagentWindows.get(windowId).resizeObserver = resizeObserver;
// Init the xterm.js terminal
this.initTeammateTerminal(windowId, paneData, win);
// Init the xterm.js terminal (lazy if hidden)
if (shouldHide) {
// Window starts hidden — defer terminal creation until visible (lazy init)
const windowEntry = this.subagentWindows.get(windowId);
if (windowEntry) {
windowEntry._lazyTerminal = true;
windowEntry._lazyPaneTarget = paneData.paneTarget;
windowEntry._lazySessionId = paneData.sessionId;
windowEntry._lazyBuffer = '';
}
} else {
this.initTeammateTerminal(windowId, paneData, win);
}
// Animate in
requestAnimationFrame(() => {
+149 -18
View File
@@ -121,7 +121,7 @@ Object.assign(CodemanApp.prototype, {
if (!windowData.minimized) {
openWindows.push({
agentId,
position: windowData.position || null
position: windowData.position || null,
});
}
}
@@ -163,8 +163,7 @@ Object.assign(CodemanApp.prototype, {
// Use the PERSISTENT parent map (THE source of truth)
// Fall back to saved sessionId only if it exists in current sessions
const parentFromMap = this.subagentParentMap.get(agentId);
const correctSessionId = parentFromMap ||
(this.sessions.has(savedSessionId) ? savedSessionId : null);
const correctSessionId = parentFromMap || (this.sessions.has(savedSessionId) ? savedSessionId : null);
if (correctSessionId) {
// Ensure the parent map has this association
@@ -184,7 +183,7 @@ Object.assign(CodemanApp.prototype, {
// Restore open windows (for recent, non-completed agents only)
const now = Date.now();
const maxAgeMs = 10 * 60 * 1000; // 10 minutes - don't restore windows for old agents
for (const { agentId, position } of (states.open || [])) {
for (const { agentId, position } of states.open || []) {
const agent = this.subagents.get(agentId);
// Only restore window if agent exists, is recent, and is still active/idle
const agentAge = agent?.startedAt ? now - agent.startedAt : Infinity;
@@ -458,6 +457,113 @@ Object.assign(CodemanApp.prototype, {
}
},
// ═══════════════════════════════════════════════════════════════
// Lazy Terminal Lifecycle
// ═══════════════════════════════════════════════════════════════
//
// Teammate terminal windows use xterm.js Terminal instances that consume
// ~75KB of DOM memory each. With 50 agents minimized, that's ~3.75MB of
// invisible terminals. To avoid this, we dispose the Terminal when a
// window is minimized and lazily re-create it when restored.
//
// Flow:
// minimize → _disposeTeammateTerminalForMinimize() → sets _lazyTerminal flag
// restore → _restoreTeammateTerminalFromLazy() → re-creates Terminal
// create (hidden/minimized) → skip initTeammateTerminal, set _lazyTerminal
//
// The pane buffer is always re-fetched from the API on restore, so no
// client-side buffer accumulation is needed (the tmux pane is the source
// of truth). Regular (non-teammate) subagent windows use activity HTML
// and are unaffected by this optimization.
/** Max bytes to buffer for a minimized teammate terminal (256KB). */
_LAZY_TERMINAL_BUFFER_CAP: 256 * 1024,
/**
* Dispose a teammate terminal when its window is minimized.
* Saves pane metadata so the terminal can be re-created on restore.
* No-op if the window has no teammate terminal.
*/
_disposeTeammateTerminalForMinimize(agentId) {
const termData = this.teammateTerminals.get(agentId);
if (!termData) return; // Not a teammate terminal window
const windowData = this.subagentWindows.get(agentId);
// Save pane metadata needed to re-create the terminal on restore
if (windowData) {
windowData._lazyTerminal = true;
windowData._lazyPaneTarget = termData.paneTarget;
windowData._lazySessionId = termData.sessionId;
// Buffer for any data that arrives while minimized (from pendingData or future writes)
windowData._lazyBuffer = '';
}
// Dispose the resize observer
if (termData.resizeObserver) {
termData.resizeObserver.disconnect();
}
// Dispose the xterm.js Terminal instance (frees DOM nodes and internal buffers)
if (termData.terminal) {
try {
termData.terminal.dispose();
} catch {}
}
// Remove from teammateTerminals map so renderSubagentWindowContent won't skip this window
// (the activity HTML can serve as a lightweight placeholder while minimized)
this.teammateTerminals.delete(agentId);
},
/**
* Re-create a teammate terminal when its window is restored from minimized state.
* Fetches the current pane buffer from the API (tmux is the source of truth).
* No-op if the window doesn't have the _lazyTerminal flag.
*/
_restoreTeammateTerminalFromLazy(agentId) {
const windowData = this.subagentWindows.get(agentId);
if (!windowData || !windowData._lazyTerminal) return;
const paneTarget = windowData._lazyPaneTarget;
const sessionId = windowData._lazySessionId;
// Clear lazy state
windowData._lazyTerminal = false;
windowData._lazyPaneTarget = null;
windowData._lazySessionId = null;
windowData._lazyBuffer = null;
if (!paneTarget || !sessionId) return;
// Re-create the terminal using the same initTeammateTerminal flow
const paneInfo = { paneTarget, sessionId };
this.initTeammateTerminal(agentId, paneInfo, windowData.element);
},
/**
* Append terminal data to a minimized teammate terminal's lazy buffer.
* Called when SSE data arrives for a minimized window. Caps at _LAZY_TERMINAL_BUFFER_CAP.
* Returns true if the data was buffered, false if the window is not in lazy mode.
*/
_bufferLazyTerminalData(agentId, data) {
const windowData = this.subagentWindows.get(agentId);
if (!windowData || !windowData._lazyTerminal) return false;
if (windowData._lazyBuffer === null || windowData._lazyBuffer === undefined) {
windowData._lazyBuffer = '';
}
// Append data, capping total size
windowData._lazyBuffer += data;
if (windowData._lazyBuffer.length > this._LAZY_TERMINAL_BUFFER_CAP) {
// Keep only the tail to stay under cap
windowData._lazyBuffer = windowData._lazyBuffer.slice(-this._LAZY_TERMINAL_BUFFER_CAP);
}
return true;
},
// ═══════════════════════════════════════════════════════════════
// Subagent Floating Windows
// ═══════════════════════════════════════════════════════════════
@@ -495,7 +601,7 @@ Object.assign(CodemanApp.prototype, {
// Only open windows for agents that belong to a Codeman-managed session tab.
// Agents from external Claude sessions (not tracked by Codeman) should not pop up.
if (agent.sessionId) {
const hasMatchingTab = Array.from(this.sessions.values()).some(s => s.claudeSessionId === agent.sessionId);
const hasMatchingTab = Array.from(this.sessions.values()).some((s) => s.claudeSessionId === agent.sessionId);
if (!hasMatchingTab) return;
}
@@ -600,9 +706,7 @@ Object.assign(CodemanApp.prototype, {
}
// Get parent TAB element for spawn animation
const parentTab = parentSessionId
? document.querySelector(`.session-tab[data-id="${parentSessionId}"]`)
: null;
const parentTab = parentSessionId ? document.querySelector(`.session-tab[data-id="${parentSessionId}"]`) : null;
// Create window element
const win = document.createElement('div');
@@ -611,17 +715,19 @@ Object.assign(CodemanApp.prototype, {
win.style.zIndex = ++this.subagentWindowZIndex;
// Build parent header if we have parent info
const parentHeader = parentSessionId && parentSessionName
? `<div class="subagent-window-parent" data-parent-session="${parentSessionId}">
const parentHeader =
parentSessionId && parentSessionName
? `<div class="subagent-window-parent" data-parent-session="${parentSessionId}">
<span class="parent-label">from</span>
<span class="parent-name" onclick="app.selectSession('${escapeHtml(parentSessionId)}')">${escapeHtml(parentSessionName)}</span>
</div>`
: '';
: '';
const teammateInfo = this.getTeammateInfo(agent);
const windowTitle = teammateInfo ? teammateInfo.name : (agent.description || agentId.substring(0, 7));
const windowTitle = teammateInfo ? teammateInfo.name : agent.description || agentId.substring(0, 7);
const maxTitleLen = isMobile ? 30 : 50;
const truncatedTitle = windowTitle.length > maxTitleLen ? windowTitle.substring(0, maxTitleLen) + '...' : windowTitle;
const truncatedTitle =
windowTitle.length > maxTitleLen ? windowTitle.substring(0, maxTitleLen) + '...' : windowTitle;
const modelBadge = agent.modelShort
? `<span class="subagent-model-badge ${agent.modelShort}">${agent.modelShort}</span>`
: '';
@@ -695,7 +801,19 @@ Object.assign(CodemanApp.prototype, {
// Render content — check if this teammate has a tmux pane
const paneInfo = teammateInfo ? this.teammatePanesByName.get(teammateInfo.name) : null;
if (paneInfo) {
this.initTeammateTerminal(agentId, paneInfo, win);
if (shouldHide) {
// Window starts hidden — defer terminal creation until visible (lazy init).
// Saves ~75KB of DOM memory per hidden teammate terminal window.
const windowEntry = this.subagentWindows.get(agentId);
if (windowEntry) {
windowEntry._lazyTerminal = true;
windowEntry._lazyPaneTarget = paneInfo.paneTarget;
windowEntry._lazySessionId = paneInfo.sessionId;
windowEntry._lazyBuffer = '';
}
} else {
this.initTeammateTerminal(agentId, paneInfo, win);
}
} else {
this.renderSubagentWindowContent(agentId);
}
@@ -755,6 +873,10 @@ Object.assign(CodemanApp.prototype, {
this.setAgentParentSessionId(agentId, parentSessionId);
}
// Dispose teammate terminal on minimize to free DOM/memory (~75KB per instance).
// The terminal will be lazily re-created on restore via initTeammateTerminal().
this._disposeTeammateTerminalForMinimize(agentId);
// Always minimize to tab
windowData.element.style.display = 'none';
windowData.minimized = true;
@@ -829,7 +951,9 @@ Object.assign(CodemanApp.prototype, {
for (const [, termData] of this.teammateTerminals) {
if (termData.resizeObserver) termData.resizeObserver.disconnect();
if (termData.terminal) {
try { termData.terminal.dispose(); } catch {}
try {
termData.terminal.dispose();
} catch {}
}
}
this.teammateTerminals.clear();
@@ -959,6 +1083,13 @@ Object.assign(CodemanApp.prototype, {
windowData.hidden = false;
}
windowData.minimized = false;
// Lazily re-create teammate terminal if it was disposed on minimize.
// Only re-create when the window is actually becoming visible.
if (shouldShow && windowData._lazyTerminal) {
this._restoreTeammateTerminalFromLazy(agentId);
}
this.updateConnectionLines();
// Restack all visible mobile windows so restored ones don't overlap
this.relayoutMobileSubagentWindows();
@@ -1067,7 +1198,7 @@ Object.assign(CodemanApp.prototype, {
if (!dropdown || dropdown.classList.contains('open')) return;
// Close other dropdowns first
document.querySelectorAll('.subagent-dropdown.open').forEach(d => {
document.querySelectorAll('.subagent-dropdown.open').forEach((d) => {
d.classList.remove('open', 'pinned');
if (d.parentElement === document.body && d._originalParent) {
d._originalParent.appendChild(d);
@@ -1089,8 +1220,8 @@ Object.assign(CodemanApp.prototype, {
// Schedule hide after delay (allows moving mouse to dropdown)
scheduleHideSubagentDropdown(badgeEl) {
this._subagentHideTimeout = setTimeout(() => {
const dropdown = badgeEl?.querySelector?.('.subagent-dropdown') ||
document.querySelector('.subagent-dropdown.open');
const dropdown =
badgeEl?.querySelector?.('.subagent-dropdown') || document.querySelector('.subagent-dropdown.open');
if (dropdown && !dropdown.classList.contains('pinned')) {
dropdown.classList.remove('open');
if (dropdown._originalParent) {
+62 -10
View File
@@ -210,7 +210,12 @@ export class WebServer extends EventEmitter {
// Store session listener references for explicit cleanup (prevents memory leaks)
private sessionListenerRefs: Map<string, SessionListenerRefs> = new Map();
private scheduledRuns: Map<string, ScheduledRun> = new Map();
private sseClients: Set<FastifyReply> = new Set();
/**
* 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();
/** Clients with backpressure — skip writes until 'drain' fires */
private backpressuredClients: Set<FastifyReply> = new Set();
private store = getStore();
@@ -590,6 +595,18 @@ export class WebServer extends EventEmitter {
return;
}
// Parse optional session subscription filter from query parameter.
// /api/events?sessions=id1,id2 — client only receives events for those sessions.
// /api/events (no param) — client receives all events (backwards-compatible).
const query = req.query as { sessions?: string };
let sessionFilter: Set<string> | null = null;
if (query.sessions) {
const ids = query.sessions.split(',').filter(Boolean);
if (ids.length > 0) {
sessionFilter = new Set(ids);
}
}
reply.raw.writeHead(200, {
'Content-Type': 'text/event-stream',
'Cache-Control': 'no-cache',
@@ -597,7 +614,7 @@ export class WebServer extends EventEmitter {
'X-Accel-Buffering': 'no', // Disable nginx buffering
});
this.sseClients.add(reply);
this.sseClients.set(reply, sessionFilter);
// Send initial state
// Use light state for SSE init to avoid sending 2MB+ terminal buffers
@@ -1972,9 +1989,15 @@ export class WebServer extends EventEmitter {
this.cachedSessionsList = null;
}
// Performance optimization: serialize JSON once for all clients.
// Only append Cloudflare tunnel padding when tunnel is actually active —
// direct/Tailscale clients don't need 8KB padding on every event.
const padding = this._isTunnelActive ? SSE_PADDING : '';
// Only append Cloudflare tunnel padding for latency-sensitive events —
// high-frequency terminal data and recovery events need immediate proxy flush,
// but low-frequency metadata events (session:created, ralph:*, respawn:*, etc.)
// are small and infrequent enough that proxy buffering doesn't matter.
// Note: session:terminal bypasses broadcast() via flushSessionTerminalBatch(),
// but is included here for completeness in case the path changes.
const needsPadding =
this._isTunnelActive && (event === SseEvent.SessionTerminal || event === SseEvent.SessionNeedsRefresh);
const padding = needsPadding ? SSE_PADDING : '';
let message: string;
try {
message = `event: ${event}\ndata: ${JSON.stringify(data)}\n\n` + padding;
@@ -1983,11 +2006,35 @@ export class WebServer extends EventEmitter {
console.error(`[Server] Failed to serialize SSE event "${event}":`, err);
return;
}
for (const client of this.sseClients) {
// 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;
}
// 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.
@@ -2066,8 +2113,13 @@ export class WebServer extends EventEmitter {
// 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);
const message = `event: session:terminal\ndata: {"id":"${sessionId}","data":${escapedData}}\n\n`;
for (const client of this.sseClients) {
// 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);
}
}
@@ -2240,7 +2292,7 @@ export class WebServer extends EventEmitter {
private cleanupDeadSSEClients(): void {
const deadClients: FastifyReply[] = [];
for (const client of this.sseClients) {
for (const [client] of this.sseClients) {
try {
// Check if the underlying socket is still writable
const socket = client.raw.socket;
@@ -2640,7 +2692,7 @@ export class WebServer extends EventEmitter {
this.cleanup.dispose();
// Gracefully close all SSE connections before clearing
for (const client of this.sseClients) {
for (const [client] of this.sseClients) {
try {
// Send a final event to notify clients of shutdown
this.sendSSE(client, 'server:shutdown', { reason: 'Server stopping' });
+445
View File
@@ -0,0 +1,445 @@
import { describe, it, expect, beforeAll, afterAll } from 'vitest';
import { WebServer } from '../src/web/server.js';
const TEST_PORT = 3212;
// Helper to parse SSE events from raw text
function parseSSEEvents(text: string): Array<{ event: string; data: unknown }> {
const events: Array<{ event: string; data: unknown }> = [];
const lines = text.split('\n');
let currentEvent = '';
let currentData = '';
for (const line of lines) {
if (line.startsWith('event: ')) {
currentEvent = line.substring(7);
} else if (line.startsWith('data: ')) {
currentData = line.substring(6);
} else if (line === '') {
if (currentEvent && currentData) {
try {
events.push({ event: currentEvent, data: JSON.parse(currentData) });
} catch {
events.push({ event: currentEvent, data: currentData });
}
}
currentEvent = '';
currentData = '';
}
}
return events;
}
// Helper to collect SSE events for a given duration
async function collectSSEEvents(
baseUrl: string,
queryParams: string,
durationMs: number
): Promise<Array<{ event: string; data: unknown }>> {
const controller = new AbortController();
let receivedData = '';
const fetchPromise = fetch(`${baseUrl}/api/events${queryParams}`, {
signal: controller.signal,
}).then(async (response) => {
const reader = response.body?.getReader();
if (reader) {
try {
while (true) {
const { done, value } = await reader.read();
if (done) break;
receivedData += new TextDecoder().decode(value);
}
} catch {
/* AbortError expected */
}
}
});
await new Promise((resolve) => setTimeout(resolve, durationMs));
controller.abort();
try {
await fetchPromise;
} catch {
/* AbortError expected */
}
return parseSSEEvents(receivedData);
}
describe('SSE Subscription Filtering', () => {
let server: WebServer;
let baseUrl: string;
beforeAll(async () => {
server = new WebServer(TEST_PORT, false, true);
await server.start();
baseUrl = `http://localhost:${TEST_PORT}`;
});
afterAll(async () => {
await server.stop();
}, 60000);
it('should accept sessions query parameter without error', async () => {
const controller = new AbortController();
const timeout = setTimeout(() => controller.abort(), 500);
try {
const response = await fetch(`${baseUrl}/api/events?sessions=abc,def`, {
signal: controller.signal,
});
expect(response.headers.get('content-type')).toBe('text/event-stream');
} catch (err: any) {
if (err.name !== 'AbortError') throw err;
} finally {
clearTimeout(timeout);
}
});
it('should send init event regardless of session filter', async () => {
const events = await collectSSEEvents(baseUrl, '?sessions=nonexistent-id', 500);
const initEvent = events.find((e) => e.event === 'init');
expect(initEvent).toBeDefined();
expect((initEvent?.data as any).sessions).toBeDefined();
});
it('should receive all events when no sessions param is provided (backwards-compatible)', async () => {
// Start two SSE listeners: one with no filter, one with a filter for a nonexistent session
const controller1 = new AbortController();
const controller2 = new AbortController();
let unfilteredData = '';
let filteredData = '';
const fetch1 = fetch(`${baseUrl}/api/events`, {
signal: controller1.signal,
}).then(async (response) => {
const reader = response.body?.getReader();
if (reader) {
try {
while (true) {
const { done, value } = await reader.read();
if (done) break;
unfilteredData += new TextDecoder().decode(value);
}
} catch {
/* expected */
}
}
});
const fetch2 = fetch(`${baseUrl}/api/events?sessions=nonexistent-session`, {
signal: controller2.signal,
}).then(async (response) => {
const reader = response.body?.getReader();
if (reader) {
try {
while (true) {
const { done, value } = await reader.read();
if (done) break;
filteredData += new TextDecoder().decode(value);
}
} catch {
/* expected */
}
}
});
// Wait for connections to establish
await new Promise((resolve) => setTimeout(resolve, 200));
// Create a session — this emits session:created (a global-ish event that has id but
// the nonexistent filter won't match it)
const createRes = await fetch(`${baseUrl}/api/sessions`, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ workingDir: '/tmp' }),
});
const createData = await createRes.json();
const sessionId = createData.session.id;
// Wait for events to arrive
await new Promise((resolve) => setTimeout(resolve, 300));
// Stop listening
controller1.abort();
controller2.abort();
try {
await fetch1;
} catch {
/* expected */
}
try {
await fetch2;
} catch {
/* expected */
}
const unfilteredEvents = parseSSEEvents(unfilteredData);
const filteredEvents = parseSSEEvents(filteredData);
// Unfiltered client should receive the session:created event
const unfilteredCreated = unfilteredEvents.find((e) => e.event === 'session:created');
expect(unfilteredCreated).toBeDefined();
expect((unfilteredCreated?.data as any).id).toBe(sessionId);
// Filtered client (subscribed to nonexistent-session) should NOT receive session:created
// because session:created has an `id` field that doesn't match the filter
const filteredCreated = filteredEvents.find((e) => e.event === 'session:created');
expect(filteredCreated).toBeUndefined();
// Both should have received the init event (it has no sessionId)
expect(unfilteredEvents.find((e) => e.event === 'init')).toBeDefined();
expect(filteredEvents.find((e) => e.event === 'init')).toBeDefined();
// Cleanup
await fetch(`${baseUrl}/api/sessions/${sessionId}`, { method: 'DELETE' });
});
it('should deliver session events to a client subscribed to that session', async () => {
// First create a session so we know the ID
const createRes = await fetch(`${baseUrl}/api/sessions`, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ workingDir: '/tmp' }),
});
const createData = await createRes.json();
const sessionId = createData.session.id;
// Now connect SSE with the session filter
const controller = new AbortController();
let receivedData = '';
const fetchPromise = fetch(`${baseUrl}/api/events?sessions=${sessionId}`, {
signal: controller.signal,
}).then(async (response) => {
const reader = response.body?.getReader();
if (reader) {
try {
while (true) {
const { done, value } = await reader.read();
if (done) break;
receivedData += new TextDecoder().decode(value);
}
} catch {
/* expected */
}
}
});
// Wait for connection
await new Promise((resolve) => setTimeout(resolve, 200));
// Delete the session — this emits session:deleted with { id: sessionId }
await fetch(`${baseUrl}/api/sessions/${sessionId}`, { method: 'DELETE' });
// Wait for events
await new Promise((resolve) => setTimeout(resolve, 300));
controller.abort();
try {
await fetchPromise;
} catch {
/* expected */
}
const events = parseSSEEvents(receivedData);
// Should receive session:deleted because we're subscribed to this session
const deletedEvent = events.find((e) => e.event === 'session:deleted');
expect(deletedEvent).toBeDefined();
expect((deletedEvent?.data as any).id).toBe(sessionId);
});
it('should filter out events for sessions not in the subscription', async () => {
// Create two sessions
const createRes1 = await fetch(`${baseUrl}/api/sessions`, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ workingDir: '/tmp' }),
});
const session1 = (await createRes1.json()).session;
const createRes2 = await fetch(`${baseUrl}/api/sessions`, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ workingDir: '/tmp' }),
});
const session2 = (await createRes2.json()).session;
// Connect SSE subscribed ONLY to session1
const controller = new AbortController();
let receivedData = '';
const fetchPromise = fetch(`${baseUrl}/api/events?sessions=${session1.id}`, {
signal: controller.signal,
}).then(async (response) => {
const reader = response.body?.getReader();
if (reader) {
try {
while (true) {
const { done, value } = await reader.read();
if (done) break;
receivedData += new TextDecoder().decode(value);
}
} catch {
/* expected */
}
}
});
// Wait for connection
await new Promise((resolve) => setTimeout(resolve, 200));
// Delete both sessions — emits session:deleted for each
await fetch(`${baseUrl}/api/sessions/${session1.id}`, { method: 'DELETE' });
await fetch(`${baseUrl}/api/sessions/${session2.id}`, { method: 'DELETE' });
// Wait for events
await new Promise((resolve) => setTimeout(resolve, 300));
controller.abort();
try {
await fetchPromise;
} catch {
/* expected */
}
const events = parseSSEEvents(receivedData);
// Should receive session:deleted for session1
const deleted1 = events.find((e) => e.event === 'session:deleted' && (e.data as any).id === session1.id);
expect(deleted1).toBeDefined();
// Should NOT receive session:deleted for session2
const deleted2 = events.find((e) => e.event === 'session:deleted' && (e.data as any).id === session2.id);
expect(deleted2).toBeUndefined();
});
it('should support subscribing to multiple sessions', async () => {
// Create two sessions
const createRes1 = await fetch(`${baseUrl}/api/sessions`, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ workingDir: '/tmp' }),
});
const session1 = (await createRes1.json()).session;
const createRes2 = await fetch(`${baseUrl}/api/sessions`, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ workingDir: '/tmp' }),
});
const session2 = (await createRes2.json()).session;
// Connect SSE subscribed to both sessions
const controller = new AbortController();
let receivedData = '';
const fetchPromise = fetch(`${baseUrl}/api/events?sessions=${session1.id},${session2.id}`, {
signal: controller.signal,
}).then(async (response) => {
const reader = response.body?.getReader();
if (reader) {
try {
while (true) {
const { done, value } = await reader.read();
if (done) break;
receivedData += new TextDecoder().decode(value);
}
} catch {
/* expected */
}
}
});
// Wait for connection
await new Promise((resolve) => setTimeout(resolve, 200));
// Delete both sessions
await fetch(`${baseUrl}/api/sessions/${session1.id}`, { method: 'DELETE' });
await fetch(`${baseUrl}/api/sessions/${session2.id}`, { method: 'DELETE' });
// Wait for events
await new Promise((resolve) => setTimeout(resolve, 300));
controller.abort();
try {
await fetchPromise;
} catch {
/* expected */
}
const events = parseSSEEvents(receivedData);
// Should receive session:deleted for both sessions
const deleted1 = events.find((e) => e.event === 'session:deleted' && (e.data as any).id === session1.id);
const deleted2 = events.find((e) => e.event === 'session:deleted' && (e.data as any).id === session2.id);
expect(deleted1).toBeDefined();
expect(deleted2).toBeDefined();
});
it('should deliver global events (no sessionId) to filtered clients', async () => {
// Connect with a filter — global events like case:created should still arrive
const controller = new AbortController();
let receivedData = '';
const fetchPromise = fetch(`${baseUrl}/api/events?sessions=some-session-id`, {
signal: controller.signal,
}).then(async (response) => {
const reader = response.body?.getReader();
if (reader) {
try {
while (true) {
const { done, value } = await reader.read();
if (done) break;
receivedData += new TextDecoder().decode(value);
}
} catch {
/* expected */
}
}
});
// Wait for connection
await new Promise((resolve) => setTimeout(resolve, 200));
// Create a case — case:created is a global event (no sessionId or id matching a session)
const caseName = `test-filter-case-${Date.now()}`;
await fetch(`${baseUrl}/api/cases`, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ name: caseName }),
});
// Wait for events
await new Promise((resolve) => setTimeout(resolve, 300));
controller.abort();
try {
await fetchPromise;
} catch {
/* expected */
}
const events = parseSSEEvents(receivedData);
// case:created should arrive because it has no sessionId — it's a global event
const caseCreated = events.find((e) => e.event === 'case:created');
expect(caseCreated).toBeDefined();
expect((caseCreated?.data as any).name).toBe(caseName);
// Cleanup
const { rmSync } = await import('node:fs');
const { join } = await import('node:path');
const { homedir } = await import('node:os');
try {
rmSync(join(homedir(), 'codeman-cases', caseName), { recursive: true });
} catch {
/* may not exist */
}
});
});