mirror of
https://github.com/Ark0N/Codeman.git
synced 2026-10-01 13:09:42 +02:00
Compare commits
6
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
7752325c90 | ||
|
|
6b284598cf | ||
|
|
94bcf524a2 | ||
|
|
98966def03 | ||
|
|
e87b03b6c2 | ||
|
|
edd494ec5f |
@@ -1,5 +1,16 @@
|
||||
# aicodeman
|
||||
|
||||
## 0.6.9
|
||||
|
||||
### Patch Changes
|
||||
|
||||
- Terminal renderer hardening, SSE bandwidth cut, image paste, and a security tightening on the new live filter:
|
||||
- **Multi-primitive yield for write pacing** (#85): replaces six raw `requestAnimationFrame` callsites in the xterm.js write pipeline with a yielding helper that races `requestAnimationFrame`, `setTimeout(50)`, and a tick Worker. Keeps the terminal responsive when the tab is backgrounded or occluded — Chrome's intensive-throttling no longer stalls long writes.
|
||||
- **WebGL longtask auto-fallback** (#83): a `PerformanceObserver` watches for ≥200ms WebGL frames; three within a 30s window disposes the WebGL addon and falls back to the canvas renderer. Decision is persisted in localStorage for 7 days, and `?webgl=force` clears it.
|
||||
- **Per-client live SSE subscription filter** (#86): each connected client gets a stable UUID and can narrow its terminal stream to one session via `POST /api/events/subscribe` — no EventSource reconnect on tab switches. Cuts SSE bandwidth roughly N× when N sessions are open. Lifecycle/metadata events (`session:*`, `case:*`, `ralph:*`, `hook:*`) now broadcast to every client so sidebars stay in sync.
|
||||
- **Image paste and drag-and-drop into the terminal** (#84): `Ctrl+V` and dropped images upload to `POST /api/sessions/:id/paste-image`, save under `${workingDir}/.claude-images/paste-${ts}.${ext}` and type the path into the terminal. Hard 10MB cap, server-generated filename (no traversal), `.svg` deliberately excluded from the allowlist to avoid a same-origin XSS path through `file-raw`.
|
||||
- **SSE clientId validation**: the per-client identifier introduced in #86 is now constrained to `[A-Za-z0-9_-]{8,64}` at both ingress points. Without this, an authenticated attacker could send another tab's clientId to silently evict it from broadcasts, mutate any clientId's session filter to blackhole the victim's terminal stream, or grow `sseClientsById` unboundedly via long IDs. The subscribe payload is also capped at 64 session entries of ≤128 chars each.
|
||||
|
||||
## 0.6.8
|
||||
|
||||
### Patch Changes
|
||||
|
||||
@@ -56,7 +56,7 @@ When user says "COM":
|
||||
|
||||
CI runs `npm run check:lockfile` on every push/PR, so lockfile drift fails the build even if the `version-packages` script is bypassed.
|
||||
|
||||
**Version**: 0.6.8 (must match `package.json`)
|
||||
**Version**: 0.6.9 (must match `package.json`)
|
||||
|
||||
## Project Overview
|
||||
|
||||
|
||||
Generated
+2
-2
@@ -1,12 +1,12 @@
|
||||
{
|
||||
"name": "aicodeman",
|
||||
"version": "0.6.8",
|
||||
"version": "0.6.9",
|
||||
"lockfileVersion": 3,
|
||||
"requires": true,
|
||||
"packages": {
|
||||
"": {
|
||||
"name": "aicodeman",
|
||||
"version": "0.6.8",
|
||||
"version": "0.6.9",
|
||||
"hasInstallScript": true,
|
||||
"license": "MIT",
|
||||
"workspaces": [
|
||||
|
||||
+1
-1
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "aicodeman",
|
||||
"version": "0.6.8",
|
||||
"version": "0.6.9",
|
||||
"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",
|
||||
|
||||
@@ -93,6 +93,7 @@ console.log('\n[build] content-hash cache busting');
|
||||
'ralph-wizard.js',
|
||||
'api-client.js',
|
||||
'subagent-windows.js',
|
||||
'image-input.js',
|
||||
'vendor/xterm-zerolag-input.js',
|
||||
];
|
||||
const manifest = {};
|
||||
|
||||
+85
-2
@@ -286,6 +286,12 @@ class CodemanApp {
|
||||
this.totalTokens = 0;
|
||||
this.globalStats = null; // Global token/cost stats across all sessions
|
||||
this.eventSource = null;
|
||||
// Stable per-page client ID — lets the server target this connection
|
||||
// for live filter updates (POST /api/events/subscribe) without forcing
|
||||
// an SSE reconnect on session switches.
|
||||
this._clientId = (typeof crypto !== 'undefined' && crypto.randomUUID)
|
||||
? crypto.randomUUID()
|
||||
: 'c-' + Math.random().toString(36).slice(2) + Date.now().toString(36);
|
||||
this.terminal = null;
|
||||
this.fitAddon = null;
|
||||
this.activeSessionId = null;
|
||||
@@ -617,14 +623,57 @@ class CodemanApp {
|
||||
this._webglAddon = new WebglAddon.WebglAddon();
|
||||
this._webglAddon.onContextLoss(() => {
|
||||
console.error('[CRASH-DIAG] WebGL context LOST — falling back to canvas renderer');
|
||||
this._webglAddon.dispose();
|
||||
_crashDiag.log('WEBGL_LOST');
|
||||
this._disableWebGLSticky('context-lost');
|
||||
this._webglAddon?.dispose();
|
||||
this._webglAddon = null;
|
||||
});
|
||||
this.terminal.loadAddon(this._webglAddon);
|
||||
console.log('[CRASH-DIAG] WebGL renderer enabled');
|
||||
this._installWebGLLongTaskGuard();
|
||||
} catch (_e) { /* WebGL2 unavailable — canvas renderer used */ }
|
||||
}
|
||||
|
||||
/**
|
||||
* Watch for sustained main-thread stalls that indicate WebGL/GPU trouble.
|
||||
* After 3 long tasks (>=200ms each) within 30s, dispose the WebGL addon and
|
||||
* persist a sticky disable so subsequent reloads also use the DOM renderer.
|
||||
* 5s grace period skips initial-load stalls. Force-re-enable: ?webgl=force.
|
||||
*/
|
||||
_installWebGLLongTaskGuard() {
|
||||
if (typeof PerformanceObserver === 'undefined' || this._webglLongTaskObserver) return;
|
||||
const installedAt = performance.now();
|
||||
const recent = [];
|
||||
try {
|
||||
this._webglLongTaskObserver = new PerformanceObserver((list) => {
|
||||
if (!this._webglAddon) return;
|
||||
const now = performance.now();
|
||||
if (now - installedAt < 5000) return;
|
||||
for (const entry of list.getEntries()) {
|
||||
if (entry.duration >= 200) recent.push(entry.startTime);
|
||||
}
|
||||
while (recent.length && now - recent[0] > 30000) recent.shift();
|
||||
if (recent.length >= 3) {
|
||||
console.warn(`[CRASH-DIAG] WebGL long-task threshold (${recent.length} stalls/30s) — falling back to canvas renderer`);
|
||||
_crashDiag.log(`WEBGL_FALLBACK: ${recent.length}`);
|
||||
this._disableWebGLSticky('long-tasks');
|
||||
this._webglAddon?.dispose();
|
||||
this._webglAddon = null;
|
||||
try { this._webglLongTaskObserver.disconnect(); } catch {}
|
||||
this._webglLongTaskObserver = null;
|
||||
try { this.terminal.refresh(0, this.terminal.rows - 1); } catch {}
|
||||
}
|
||||
});
|
||||
this._webglLongTaskObserver.observe({ type: 'longtask', buffered: false });
|
||||
} catch { /* longtask not supported */ }
|
||||
}
|
||||
|
||||
_disableWebGLSticky(reason) {
|
||||
try {
|
||||
localStorage.setItem('codeman-webgl-disabled', JSON.stringify({ reason, at: Date.now() }));
|
||||
} catch {}
|
||||
}
|
||||
|
||||
// ═══════════════════════════════════════════════════════════════
|
||||
// Event Listeners (Keyboard Shortcuts, Resize, Beforeunload)
|
||||
// ═══════════════════════════════════════════════════════════════
|
||||
@@ -696,6 +745,28 @@ class CodemanApp {
|
||||
// SSE Connection
|
||||
// ═══════════════════════════════════════════════════════════════
|
||||
|
||||
/**
|
||||
* POST a live subscription update so the server filters terminal events
|
||||
* to the given session(s) for this client. Fire-and-forget — failures
|
||||
* are non-fatal because we'll still get every event we don't want
|
||||
* (just at higher cost), and the next reconnect carries the filter via
|
||||
* the SSE query string.
|
||||
*/
|
||||
_updateSseSubscription(sessionId) {
|
||||
try {
|
||||
const body = JSON.stringify({
|
||||
clientId: this._clientId,
|
||||
sessions: sessionId ? [sessionId] : null,
|
||||
});
|
||||
fetch('/api/events/subscribe', {
|
||||
method: 'POST',
|
||||
headers: { 'Content-Type': 'application/json' },
|
||||
body,
|
||||
keepalive: true,
|
||||
}).catch(() => { /* non-fatal */ });
|
||||
} catch { /* non-fatal */ }
|
||||
}
|
||||
|
||||
connectSSE() {
|
||||
// Check if browser is offline
|
||||
if (!navigator.onLine) {
|
||||
@@ -725,7 +796,13 @@ class CodemanApp {
|
||||
this.setConnectionStatus('reconnecting');
|
||||
}
|
||||
|
||||
this.eventSource = new EventSource('/api/events');
|
||||
// Build URL with stable client ID and (if known) the active-session
|
||||
// filter so the server only streams session:terminal events for the
|
||||
// session we're rendering. Lifecycle/metadata events are sent globally
|
||||
// regardless of filter (server side).
|
||||
const _sseParams = new URLSearchParams({ clientId: this._clientId });
|
||||
if (this.activeSessionId) _sseParams.set('sessions', this.activeSessionId);
|
||||
this.eventSource = new EventSource(`/api/events?${_sseParams.toString()}`);
|
||||
|
||||
// Store all event listeners for cleanup on reconnect
|
||||
const listeners = [];
|
||||
@@ -2424,6 +2501,12 @@ class CodemanApp {
|
||||
this._cleanupPreviousSession(sessionId);
|
||||
this.activeSessionId = sessionId;
|
||||
try { localStorage.setItem('codeman-active-session', sessionId); } catch {}
|
||||
// Narrow SSE filter to the active session — server stops streaming
|
||||
// session:terminal events for other sessions to this client. Cuts
|
||||
// SSE traffic ~Nx for N concurrent sessions. Fire-and-forget; on the
|
||||
// rare race where server doesn't know our clientId yet, the next
|
||||
// selectSession or reconnect catches up.
|
||||
this._updateSseSubscription(sessionId);
|
||||
this.hideWelcome();
|
||||
// Clear idle hooks on view, but keep action hooks until user interacts
|
||||
this.clearPendingHooks(sessionId, 'idle_prompt');
|
||||
|
||||
@@ -0,0 +1,143 @@
|
||||
/**
|
||||
* Image Input Mixin - Clipboard paste and drag-and-drop image support
|
||||
*
|
||||
* For paste: intercepts Ctrl+V at the xterm keyboard level, creates a temporary
|
||||
* hidden contenteditable div ("paste trap"), lets the browser's native paste fill
|
||||
* it, then checks for image data. This works on HTTP (no secure context needed).
|
||||
*
|
||||
* For drag-and-drop: listens on the terminal container for file drops.
|
||||
*
|
||||
* @dependency app.js (uses global `app` for sendInput, activeSessionId, showToast)
|
||||
* @dependency panels-ui.js (provides showToast)
|
||||
*/
|
||||
|
||||
Object.assign(CodemanApp.prototype, {
|
||||
|
||||
initImageInput() {
|
||||
// Drag-and-drop handlers on terminal container
|
||||
const container = document.getElementById('terminalContainer');
|
||||
if (!container) return;
|
||||
|
||||
container.addEventListener('dragover', (e) => {
|
||||
e.preventDefault();
|
||||
if (e.dataTransfer && e.dataTransfer.types.includes('Files')) {
|
||||
container.classList.add('drag-active');
|
||||
}
|
||||
});
|
||||
|
||||
container.addEventListener('dragleave', (e) => {
|
||||
if (!container.contains(e.relatedTarget)) {
|
||||
container.classList.remove('drag-active');
|
||||
}
|
||||
});
|
||||
|
||||
container.addEventListener('drop', (e) => {
|
||||
e.preventDefault();
|
||||
container.classList.remove('drag-active');
|
||||
|
||||
if (!this.activeSessionId) return;
|
||||
if (!e.dataTransfer || !e.dataTransfer.files.length) return;
|
||||
|
||||
const imageFiles = Array.from(e.dataTransfer.files).filter((f) => f.type.startsWith('image/'));
|
||||
if (imageFiles.length === 0) {
|
||||
this.showToast('Only image files are supported', 'error');
|
||||
return;
|
||||
}
|
||||
this._uploadAndInsertImages(imageFiles);
|
||||
});
|
||||
},
|
||||
|
||||
// Called from customKeyEventHandler in terminal-ui.js on Ctrl+V keydown.
|
||||
// Creates a hidden paste trap, lets the browser paste into it, then inspects
|
||||
// the result for images. Works on plain HTTP (no Clipboard API needed).
|
||||
_handleImagePaste() {
|
||||
const self = this;
|
||||
|
||||
// Create a hidden contenteditable div to receive the paste
|
||||
const trap = document.createElement('div');
|
||||
trap.contentEditable = 'true';
|
||||
trap.style.cssText = 'position:fixed;left:-9999px;top:0;width:1px;height:1px;opacity:0;overflow:hidden';
|
||||
document.body.appendChild(trap);
|
||||
trap.focus();
|
||||
|
||||
// Listen for the paste event on our trap
|
||||
trap.addEventListener('paste', function(e) {
|
||||
e.stopPropagation();
|
||||
|
||||
// Check for images in clipboard items
|
||||
var imageFiles = [];
|
||||
var items = e.clipboardData && e.clipboardData.items;
|
||||
if (items) {
|
||||
for (var i = 0; i < items.length; i++) {
|
||||
if (items[i].type.startsWith('image/')) {
|
||||
var blob = items[i].getAsFile();
|
||||
if (blob) imageFiles.push(blob);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Clean up the trap
|
||||
setTimeout(function() {
|
||||
if (trap.parentNode) trap.parentNode.removeChild(trap);
|
||||
// Refocus the terminal
|
||||
if (self.terminal) self.terminal.focus();
|
||||
}, 0);
|
||||
|
||||
if (imageFiles.length > 0) {
|
||||
e.preventDefault();
|
||||
self._uploadAndInsertImages(imageFiles);
|
||||
} else {
|
||||
// No image -- extract text and send to terminal
|
||||
var text = e.clipboardData ? e.clipboardData.getData('text/plain') : '';
|
||||
e.preventDefault();
|
||||
if (text) self.sendInput(text);
|
||||
}
|
||||
});
|
||||
|
||||
// Trigger the browser's native paste via execCommand
|
||||
// (this fires the paste event on our focused trap element)
|
||||
document.execCommand('paste');
|
||||
},
|
||||
|
||||
async _uploadAndInsertImages(files) {
|
||||
const sessionId = this.activeSessionId;
|
||||
if (!sessionId) return;
|
||||
|
||||
this.showToast('Uploading ' + files.length + ' image' + (files.length > 1 ? 's' : '') + '...', 'info');
|
||||
|
||||
const paths = [];
|
||||
for (const file of files) {
|
||||
try {
|
||||
const path = await this._uploadPasteImage(sessionId, file);
|
||||
paths.push(path);
|
||||
} catch (err) {
|
||||
this.showToast('Upload failed: ' + (err.message || 'unknown error'), 'error');
|
||||
}
|
||||
}
|
||||
|
||||
if (paths.length > 0) {
|
||||
const pathStr = paths.join(' ');
|
||||
await this.sendInput(pathStr);
|
||||
this.showToast(paths.length + ' image' + (paths.length > 1 ? 's' : '') + ' ready', 'success');
|
||||
}
|
||||
},
|
||||
|
||||
async _uploadPasteImage(sessionId, file) {
|
||||
const form = new FormData();
|
||||
form.append('image', file);
|
||||
|
||||
const resp = await fetch('/api/sessions/' + sessionId + '/paste-image', {
|
||||
method: 'POST',
|
||||
body: form,
|
||||
});
|
||||
|
||||
if (!resp.ok) {
|
||||
const data = await resp.json().catch(() => ({}));
|
||||
throw new Error(data.error || 'HTTP ' + resp.status);
|
||||
}
|
||||
|
||||
const data = await resp.json();
|
||||
return data.path;
|
||||
},
|
||||
|
||||
});
|
||||
@@ -1804,5 +1804,6 @@
|
||||
<script defer src="ralph-wizard.js"></script>
|
||||
<script defer src="api-client.js"></script>
|
||||
<script defer src="subagent-windows.js"></script>
|
||||
<script defer src="image-input.js"></script>
|
||||
</body>
|
||||
</html>
|
||||
|
||||
@@ -8591,3 +8591,23 @@ kbd {
|
||||
margin-top: 4px;
|
||||
font-size: 0.7rem;
|
||||
}
|
||||
|
||||
/* Image drag-and-drop overlay */
|
||||
#terminalContainer.drag-active {
|
||||
outline: 2px dashed #4a9eff;
|
||||
outline-offset: -2px;
|
||||
position: relative;
|
||||
}
|
||||
#terminalContainer.drag-active::after {
|
||||
content: 'Drop image here';
|
||||
position: absolute;
|
||||
inset: 0;
|
||||
display: flex;
|
||||
align-items: center;
|
||||
justify-content: center;
|
||||
background: rgba(74, 158, 255, 0.08);
|
||||
color: #4a9eff;
|
||||
font-size: 1.2rem;
|
||||
pointer-events: none;
|
||||
z-index: 100;
|
||||
}
|
||||
|
||||
+110
-10
@@ -83,6 +83,15 @@ Object.assign(CodemanApp.prototype, {
|
||||
// Let Alt+digit pass through to browser (tab switching)
|
||||
if (ev.altKey && ev.key >= '0' && ev.key <= '9') return false;
|
||||
|
||||
// Ctrl+V / Cmd+V: intercept before xterm sends ^V to PTY.
|
||||
// Route through our paste trap which handles both images and text.
|
||||
if ((ev.ctrlKey || ev.metaKey) && ev.key === 'v' && ev.type === 'keydown') {
|
||||
if (this.activeSessionId && this._handleImagePaste) {
|
||||
this._handleImagePaste();
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
// Shift+Enter / Ctrl+Enter: insert newline for multi-line input.
|
||||
// xterm.js sends plain \r for all Enter variants, so Claude Code (Ink) can't
|
||||
// distinguish them. We use tmux send-keys -H to send a line feed byte (0x0a)
|
||||
@@ -174,10 +183,36 @@ Object.assign(CodemanApp.prototype, {
|
||||
// but the 48KB/frame flush cap in flushPendingWrites() now prevents
|
||||
// oversized terminal.write() calls that triggered the stalls.
|
||||
// Disable with ?nowebgl URL param if GPU issues return.
|
||||
// Auto-fallback: _initWebGL installs a long-task watchdog that disables
|
||||
// WebGL sticky in localStorage after repeated GPU stalls (see app.js).
|
||||
// Force re-enable after sticky disable with ?webgl=force.
|
||||
// Lazy-loaded: script downloaded only on desktop (saves 244KB on mobile).
|
||||
this._webglAddon = null;
|
||||
const skipWebGL = MobileDetection.getDeviceType() !== 'desktop';
|
||||
if (!skipWebGL && !new URLSearchParams(location.search).has('nowebgl')) {
|
||||
const _params = new URLSearchParams(location.search);
|
||||
if (_params.get('webgl') === 'force') {
|
||||
try { localStorage.removeItem('codeman-webgl-disabled'); } catch {}
|
||||
}
|
||||
const _stickyDisabled = (() => {
|
||||
try {
|
||||
const raw = localStorage.getItem('codeman-webgl-disabled');
|
||||
if (!raw) return false;
|
||||
const { at } = JSON.parse(raw);
|
||||
// Auto-expire after 7 days so we retry (driver may have been fixed)
|
||||
if (Date.now() - at > 7 * 24 * 60 * 60 * 1000) {
|
||||
localStorage.removeItem('codeman-webgl-disabled');
|
||||
return false;
|
||||
}
|
||||
return true;
|
||||
} catch { return false; }
|
||||
})();
|
||||
const skipWebGL =
|
||||
MobileDetection.getDeviceType() !== 'desktop' ||
|
||||
_params.has('nowebgl') ||
|
||||
_stickyDisabled;
|
||||
if (_stickyDisabled) {
|
||||
console.log('[CRASH-DIAG] WebGL sticky-disabled from prior stalls — DOM renderer in use. Re-enable: ?webgl=force');
|
||||
}
|
||||
if (!skipWebGL) {
|
||||
if (typeof WebglAddon !== 'undefined') {
|
||||
this._initWebGL();
|
||||
} else {
|
||||
@@ -341,6 +376,9 @@ Object.assign(CodemanApp.prototype, {
|
||||
// Welcome message
|
||||
this.showWelcome();
|
||||
|
||||
// Image paste and drag-and-drop support
|
||||
this.initImageInput();
|
||||
|
||||
// Generation counter for chunkedTerminalWrite — aborts stale writes on tab switch
|
||||
this._chunkedWriteGen = 0;
|
||||
|
||||
@@ -1151,7 +1189,7 @@ Object.assign(CodemanApp.prototype, {
|
||||
|
||||
if (!this.writeFrameScheduled) {
|
||||
this.writeFrameScheduled = true;
|
||||
requestAnimationFrame(() => {
|
||||
this._safeYield(() => {
|
||||
// xterm.js 6.0 handles DEC 2026 sync markers natively — it buffers
|
||||
// content between 2026h/2026l and renders atomically. No need for
|
||||
// client-side incomplete-block detection; just flush every frame.
|
||||
@@ -1176,7 +1214,7 @@ Object.assign(CodemanApp.prototype, {
|
||||
// Trigger a normal flush
|
||||
if (!this.writeFrameScheduled) {
|
||||
this.writeFrameScheduled = true;
|
||||
requestAnimationFrame(() => {
|
||||
this._safeYield(() => {
|
||||
this.flushPendingWrites();
|
||||
this.writeFrameScheduled = false;
|
||||
});
|
||||
@@ -1264,7 +1302,7 @@ Object.assign(CodemanApp.prototype, {
|
||||
deferred = true;
|
||||
if (!this.writeFrameScheduled) {
|
||||
this.writeFrameScheduled = true;
|
||||
requestAnimationFrame(() => {
|
||||
this._safeYield(() => {
|
||||
this.flushPendingWrites();
|
||||
this.writeFrameScheduled = false;
|
||||
});
|
||||
@@ -1336,9 +1374,70 @@ Object.assign(CodemanApp.prototype, {
|
||||
}
|
||||
},
|
||||
|
||||
/**
|
||||
* Schedule cb via THREE racing primitives so data-pacing makes progress
|
||||
* regardless of which scheduling primitive Chrome is throttling:
|
||||
* 1. requestAnimationFrame — primary, fires at compositor rate
|
||||
* (may be 0Hz when window is occluded / on backgrounded monitor).
|
||||
* 2. setTimeout(50) — fallback for occluded-but-visible windows
|
||||
* (clamped to 1Hz by Chrome's intensive wake-up throttling
|
||||
* after ~5 min of no user interaction).
|
||||
* 3. Worker postMessage — bypasses intensive throttling entirely;
|
||||
* Workers are not subject to background-tab / idle-tab throttling
|
||||
* (the React Scheduler trick).
|
||||
* Whichever fires first wins; the others are no-ops thanks to the
|
||||
* `done` guard. Without all three, chunkedTerminalWrite and the deferred
|
||||
* path of flushPendingWrites stall indefinitely when the substrate is
|
||||
* degraded (visible-but-occluded window, OR idle-throttled tab, OR
|
||||
* background tab on a different monitor).
|
||||
*/
|
||||
_safeYield(cb) {
|
||||
let done = false;
|
||||
const wrapped = () => {
|
||||
if (done) return;
|
||||
done = true;
|
||||
cb();
|
||||
};
|
||||
requestAnimationFrame(wrapped);
|
||||
setTimeout(wrapped, 50);
|
||||
this._workerYield(wrapped);
|
||||
},
|
||||
|
||||
/**
|
||||
* Lazy-init a tiny "tick" worker whose only job is to postMessage back to
|
||||
* us as fast as possible, escaping main-thread throttling. The worker's
|
||||
* setTimeout(0) is not subject to Chrome's intensive wake-up throttling
|
||||
* even when the parent tab is idle.
|
||||
*/
|
||||
_workerYield(cb) {
|
||||
try {
|
||||
if (this._yieldWorker === undefined) {
|
||||
// First call: build the worker (or mark unavailable). Each
|
||||
// postMessage in produces exactly one postMessage out — we count on
|
||||
// FIFO 1:1 to drain queue entries.
|
||||
const src = "onmessage=()=>setTimeout(()=>postMessage(0),0);";
|
||||
const blob = new Blob([src], { type: 'application/javascript' });
|
||||
const url = URL.createObjectURL(blob);
|
||||
this._yieldWorker = new Worker(url);
|
||||
URL.revokeObjectURL(url);
|
||||
this._yieldQueue = [];
|
||||
this._yieldWorker.onmessage = () => {
|
||||
const fn = this._yieldQueue.shift();
|
||||
if (fn) fn();
|
||||
};
|
||||
}
|
||||
if (!this._yieldWorker) return;
|
||||
this._yieldQueue.push(cb);
|
||||
this._yieldWorker.postMessage(0);
|
||||
} catch {
|
||||
this._yieldWorker = null; // mark unavailable, future calls skip
|
||||
}
|
||||
},
|
||||
|
||||
/**
|
||||
* Write large buffer to terminal in chunks to avoid UI jank.
|
||||
* Uses requestAnimationFrame to spread work across frames.
|
||||
* Uses _safeYield to spread work across frames; falls back to setTimeout
|
||||
* and a tick-Worker so progress continues on occluded / idle-throttled tabs.
|
||||
* @param {string} buffer - The full terminal buffer to write
|
||||
* @param {number} chunkSize - Size of each chunk (default 128KB for smooth 60fps)
|
||||
* @returns {Promise<void>} - Resolves when all chunks written
|
||||
@@ -1397,7 +1496,7 @@ Object.assign(CodemanApp.prototype, {
|
||||
`[CRASH-DIAG] chunkedTerminalWrite complete: ${cleanBuffer.length} bytes in ${_chunkCount} chunks, ${_totalMs.toFixed(0)}ms total`
|
||||
);
|
||||
// Wait one more frame for xterm to finish rendering before resolving
|
||||
requestAnimationFrame(finish);
|
||||
this._safeYield(finish);
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -1412,12 +1511,13 @@ Object.assign(CodemanApp.prototype, {
|
||||
);
|
||||
offset += chunkSize;
|
||||
|
||||
// Schedule next chunk on next frame
|
||||
requestAnimationFrame(writeChunk);
|
||||
// Schedule next chunk; rAF if possible, else setTimeout/Worker
|
||||
// fallback so progress doesn't stall on occluded/unfocused windows.
|
||||
this._safeYield(writeChunk);
|
||||
};
|
||||
|
||||
// Start writing
|
||||
requestAnimationFrame(writeChunk);
|
||||
this._safeYield(writeChunk);
|
||||
});
|
||||
},
|
||||
|
||||
|
||||
@@ -5,7 +5,7 @@
|
||||
*/
|
||||
|
||||
import { FastifyInstance } from 'fastify';
|
||||
import { join, dirname } from 'node:path';
|
||||
import { join, dirname, extname } from 'node:path';
|
||||
import { homedir } from 'node:os';
|
||||
import { existsSync, statSync, mkdirSync, writeFileSync } from 'node:fs';
|
||||
import { execFile } from 'node:child_process';
|
||||
@@ -1437,4 +1437,109 @@ export function registerSessionRoutes(
|
||||
|
||||
return { sessions: results.slice(0, 50) };
|
||||
});
|
||||
|
||||
// ═══════════════════════════════════════════════════════════════
|
||||
// Paste Image (clipboard / drag-drop upload)
|
||||
// ═══════════════════════════════════════════════════════════════
|
||||
|
||||
const MAX_PASTE_IMAGE_SIZE = 10 * 1024 * 1024; // 10 MB
|
||||
const ALLOWED_IMAGE_EXTS = new Set(['.png', '.jpg', '.jpeg', '.gif', '.webp', '.bmp']);
|
||||
|
||||
app.post('/api/sessions/:id/paste-image', async (req, reply) => {
|
||||
const { id } = req.params as { id: string };
|
||||
const session = findSessionOrFail(ctx, id);
|
||||
|
||||
const contentType = req.headers['content-type'] ?? '';
|
||||
if (!contentType.includes('multipart/form-data')) {
|
||||
reply.code(400);
|
||||
return createErrorResponse(ApiErrorCode.INVALID_INPUT, 'Expected multipart/form-data');
|
||||
}
|
||||
|
||||
// Parse multipart boundary
|
||||
const boundaryMatch = contentType.match(/boundary=(.+?)(?:;|$)/);
|
||||
if (!boundaryMatch) {
|
||||
reply.code(400);
|
||||
return createErrorResponse(ApiErrorCode.INVALID_INPUT, 'Missing boundary');
|
||||
}
|
||||
|
||||
// Collect raw body with size limit
|
||||
const chunks: Buffer[] = [];
|
||||
let totalSize = 0;
|
||||
for await (const chunk of req.raw) {
|
||||
totalSize += chunk.length;
|
||||
if (totalSize > MAX_PASTE_IMAGE_SIZE) {
|
||||
reply.code(413);
|
||||
return createErrorResponse(ApiErrorCode.INVALID_INPUT, 'File too large (max 10MB)');
|
||||
}
|
||||
chunks.push(chunk as Buffer);
|
||||
}
|
||||
const body = Buffer.concat(chunks);
|
||||
|
||||
// Extract image from multipart body
|
||||
const boundary = '--' + boundaryMatch[1];
|
||||
const boundaryBuf = Buffer.from(boundary);
|
||||
const parts: { headers: string; data: Buffer }[] = [];
|
||||
let pos = 0;
|
||||
|
||||
while (pos < body.length) {
|
||||
const start = body.indexOf(boundaryBuf, pos);
|
||||
if (start === -1) break;
|
||||
const afterBoundary = start + boundaryBuf.length;
|
||||
if (body[afterBoundary] === 0x2d && body[afterBoundary + 1] === 0x2d) break;
|
||||
const headerStart = afterBoundary + 2;
|
||||
const headerEnd = body.indexOf(Buffer.from('\r\n\r\n'), headerStart);
|
||||
if (headerEnd === -1) break;
|
||||
const headers = body.subarray(headerStart, headerEnd).toString();
|
||||
const dataStart = headerEnd + 4;
|
||||
const nextBoundary = body.indexOf(boundaryBuf, dataStart);
|
||||
const dataEnd = nextBoundary === -1 ? body.length : nextBoundary - 2;
|
||||
parts.push({ headers, data: body.subarray(dataStart, dataEnd) });
|
||||
pos = nextBoundary === -1 ? body.length : nextBoundary;
|
||||
}
|
||||
|
||||
const imagePart = parts.find((p) => p.headers.includes('name="image"'));
|
||||
if (!imagePart || imagePart.data.length === 0) {
|
||||
reply.code(400);
|
||||
return createErrorResponse(ApiErrorCode.INVALID_INPUT, 'No image uploaded');
|
||||
}
|
||||
|
||||
// Determine extension from filename or Content-Type
|
||||
let ext = '.png';
|
||||
const filenameMatch = imagePart.headers.match(/filename="(.+?)"/);
|
||||
if (filenameMatch) {
|
||||
const origExt = extname(filenameMatch[1]).toLowerCase();
|
||||
if (ALLOWED_IMAGE_EXTS.has(origExt)) ext = origExt;
|
||||
}
|
||||
const ctMatch = imagePart.headers.match(/Content-Type:\s*image\/(png|jpeg|jpg|webp|gif|bmp)/i);
|
||||
if (ctMatch) {
|
||||
const map: Record<string, string> = {
|
||||
png: '.png',
|
||||
jpeg: '.jpg',
|
||||
jpg: '.jpg',
|
||||
webp: '.webp',
|
||||
gif: '.gif',
|
||||
bmp: '.bmp',
|
||||
};
|
||||
ext = map[ctMatch[1].toLowerCase()] ?? ext;
|
||||
}
|
||||
|
||||
if (!ALLOWED_IMAGE_EXTS.has(ext)) {
|
||||
reply.code(400);
|
||||
return createErrorResponse(
|
||||
ApiErrorCode.INVALID_INPUT,
|
||||
`Unsupported image type: ${ext}. Allowed: ${[...ALLOWED_IMAGE_EXTS].join(', ')}`
|
||||
);
|
||||
}
|
||||
|
||||
// Save to {workingDir}/.claude-images/
|
||||
const imageDir = join(session.workingDir, '.claude-images');
|
||||
if (!existsSync(imageDir)) {
|
||||
mkdirSync(imageDir, { recursive: true });
|
||||
}
|
||||
const filename = `paste-${Date.now()}${ext}`;
|
||||
const filepath = join(imageDir, filename);
|
||||
await fs.writeFile(filepath, imagePart.data);
|
||||
|
||||
return { success: true, path: filepath, filename };
|
||||
});
|
||||
}
|
||||
|
||||
+39
-5
@@ -34,7 +34,7 @@ import fastifyStatic from '@fastify/static';
|
||||
import fastifyWebsocket from '@fastify/websocket';
|
||||
import { join, dirname } from 'node:path';
|
||||
import { fileURLToPath } from 'node:url';
|
||||
import { existsSync, mkdirSync, readFileSync, chmodSync } from 'node:fs';
|
||||
import { existsSync, mkdirSync, readFileSync, chmodSync, rmSync } from 'node:fs';
|
||||
import fs from 'node:fs/promises';
|
||||
import { execSync } from 'node:child_process';
|
||||
import { homedir, hostname as getHostname } from 'node:os';
|
||||
@@ -119,6 +119,11 @@ import {
|
||||
|
||||
const __dirname = dirname(fileURLToPath(import.meta.url));
|
||||
|
||||
// Bounded, predictable shape for SSE client identifiers: alphanumerics, `_`, `-`.
|
||||
// Length range covers crypto.randomUUID() (36 chars) plus any short stable IDs,
|
||||
// while capping growth of `sseClientsById` and blocking pathological inputs.
|
||||
const SSE_CLIENT_ID_RE = /^[A-Za-z0-9_-]{8,64}$/;
|
||||
|
||||
function escapeHtmlText(value: string): string {
|
||||
return value.replaceAll('&', '&').replaceAll('<', '<').replaceAll('>', '>');
|
||||
}
|
||||
@@ -579,9 +584,11 @@ export class WebServer extends EventEmitter {
|
||||
}
|
||||
|
||||
// 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 };
|
||||
// /api/events?sessions=id1,id2 — client only receives session:terminal
|
||||
// events for those sessions (other events broadcast to all clients).
|
||||
// /api/events?clientId=<uuid> — enables live filter updates via
|
||||
// POST /api/events/subscribe without reconnecting.
|
||||
const query = req.query as { sessions?: string; clientId?: string };
|
||||
let sessionFilter: Set<string> | null = null;
|
||||
if (query.sessions) {
|
||||
const ids = query.sessions
|
||||
@@ -592,6 +599,8 @@ export class WebServer extends EventEmitter {
|
||||
sessionFilter = new Set(ids);
|
||||
}
|
||||
}
|
||||
const clientId =
|
||||
typeof query.clientId === 'string' && SSE_CLIENT_ID_RE.test(query.clientId) ? query.clientId : undefined;
|
||||
|
||||
reply.raw.writeHead(200, {
|
||||
'Content-Type': 'text/event-stream',
|
||||
@@ -603,7 +612,7 @@ export class WebServer extends EventEmitter {
|
||||
// Track tunnel clients — cloudflared proxies locally so req.ip is always
|
||||
// 127.0.0.1; detect tunnel traffic via Cf-Connecting-Ip header instead.
|
||||
const isRemote = !!req.headers['cf-connecting-ip'];
|
||||
this.sse.addClient(reply, sessionFilter, isRemote);
|
||||
this.sse.addClient(reply, sessionFilter, isRemote, clientId);
|
||||
|
||||
// Send initial state
|
||||
// Use light state for SSE init to avoid sending 2MB+ terminal buffers
|
||||
@@ -618,6 +627,22 @@ export class WebServer extends EventEmitter {
|
||||
});
|
||||
});
|
||||
|
||||
// Live subscription update — change a connected client's session filter
|
||||
// without forcing an SSE reconnect. Body: { clientId, sessions: string[] | null }
|
||||
// Empty/null sessions array = remove filter (receive all session:terminal events).
|
||||
this.app.post('/api/events/subscribe', (req, reply) => {
|
||||
const body = (req.body || {}) as { clientId?: string; sessions?: string[] | null };
|
||||
if (typeof body.clientId !== 'string' || !SSE_CLIENT_ID_RE.test(body.clientId)) {
|
||||
reply.code(400).send({ error: 'clientId required' });
|
||||
return;
|
||||
}
|
||||
const sessions = Array.isArray(body.sessions)
|
||||
? body.sessions.filter((s) => typeof s === 'string' && s.length > 0 && s.length <= 128).slice(0, 64)
|
||||
: null;
|
||||
const updated = this.sse.updateClientFilter(body.clientId, sessions);
|
||||
reply.code(updated ? 204 : 404).send();
|
||||
});
|
||||
|
||||
// Global error handler for structured errors thrown by findSessionOrFail
|
||||
this.app.setErrorHandler((error, _req, reply) => {
|
||||
const statusCode = (error as { statusCode?: number }).statusCode ?? 500;
|
||||
@@ -926,6 +951,15 @@ export class WebServer extends EventEmitter {
|
||||
fileStreamManager.closeSessionStreams(sessionId);
|
||||
// Stop watching for images in this session's directory
|
||||
imageWatcher.unwatchSession(sessionId);
|
||||
// Clean up pasted images directory for this session
|
||||
if (killMux && session.workingDir) {
|
||||
const pasteImageDir = join(session.workingDir, '.claude-images');
|
||||
try {
|
||||
rmSync(pasteImageDir, { recursive: true, force: true });
|
||||
} catch {
|
||||
// Best-effort cleanup
|
||||
}
|
||||
}
|
||||
await session.stop(killMux);
|
||||
this.sessions.delete(sessionId);
|
||||
// Only remove from state.json if we're also killing the mux session.
|
||||
|
||||
@@ -48,6 +48,8 @@ export class SseStreamManager {
|
||||
* or `null` meaning "receive all events" (backwards-compatible default).
|
||||
*/
|
||||
private sseClients: Map<FastifyReply, Set<string> | null> = new Map();
|
||||
/** Optional client-supplied IDs → reply, for live filter updates without reconnecting */
|
||||
private sseClientsById: Map<string, FastifyReply> = 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 */
|
||||
@@ -103,17 +105,43 @@ export class SseStreamManager {
|
||||
this._isTunnelActive = active;
|
||||
}
|
||||
|
||||
addClient(reply: FastifyReply, sessionFilter: Set<string> | null, isRemote: boolean): void {
|
||||
addClient(reply: FastifyReply, sessionFilter: Set<string> | null, isRemote: boolean, clientId?: string): void {
|
||||
this.sseClients.set(reply, sessionFilter);
|
||||
if (isRemote) {
|
||||
this.remoteSseClients.add(reply);
|
||||
}
|
||||
if (clientId) {
|
||||
// If a previous reply registered the same id (reconnect), drop the old one.
|
||||
const prev = this.sseClientsById.get(clientId);
|
||||
if (prev && prev !== reply) {
|
||||
this.sseClients.delete(prev);
|
||||
this.remoteSseClients.delete(prev);
|
||||
this.backpressuredClients.delete(prev);
|
||||
}
|
||||
this.sseClientsById.set(clientId, reply);
|
||||
}
|
||||
}
|
||||
|
||||
removeClient(reply: FastifyReply): void {
|
||||
this.sseClients.delete(reply);
|
||||
this.remoteSseClients.delete(reply);
|
||||
this.backpressuredClients.delete(reply);
|
||||
// Clear any clientId mappings pointing at this reply
|
||||
for (const [id, r] of this.sseClientsById) {
|
||||
if (r === reply) this.sseClientsById.delete(id);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Update an existing client's session subscription filter without forcing
|
||||
* an SSE reconnect. Returns true if the client was found and updated.
|
||||
*/
|
||||
updateClientFilter(clientId: string, sessions: string[] | null): boolean {
|
||||
const reply = this.sseClientsById.get(clientId);
|
||||
if (!reply || !this.sseClients.has(reply)) return false;
|
||||
const filter = sessions && sessions.length > 0 ? new Set(sessions) : null;
|
||||
this.sseClients.set(reply, filter);
|
||||
return true;
|
||||
}
|
||||
|
||||
/** Send a single SSE event to a specific client. */
|
||||
@@ -188,35 +216,18 @@ export class SseStreamManager {
|
||||
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;
|
||||
// Subscription filtering is intentionally NOT applied here. The
|
||||
// `?sessions=` filter is intended to suppress only the high-volume
|
||||
// terminal stream — lifecycle/metadata events (session:created,
|
||||
// session:updated, ralph:*, hook:*, etc.) are needed for correct UI
|
||||
// state across all sessions even when the client subscribes to a single
|
||||
// active session's terminal output. Terminal events bypass this method
|
||||
// entirely (see flushSessionTerminalBatch — it applies the filter).
|
||||
for (const [client] of this.sseClients) {
|
||||
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)
|
||||
|
||||
@@ -3,7 +3,8 @@
|
||||
*
|
||||
* Covers:
|
||||
* - SSE subscription filter edge cases (empty params, whitespace, duplicates)
|
||||
* - extractSessionId logic (sessionId vs id field, global events)
|
||||
* - Lifecycle-event broadcast contract (session:*, case:* fan out to all clients;
|
||||
* only session:terminal is filtered by subscription)
|
||||
* - Tab switching: terminal buffer loading, session creation + switch
|
||||
* - Terminal data cap / backpressure recovery
|
||||
* - Lazy teammate terminal lifecycle
|
||||
@@ -254,11 +255,11 @@ describe('Operation Lightspeed', () => {
|
||||
});
|
||||
|
||||
// ═══════════════════════════════════════════════════════════════
|
||||
// extractSessionId — Event Classification
|
||||
// Lifecycle Event Broadcast — Event Classification
|
||||
// ═══════════════════════════════════════════════════════════════
|
||||
|
||||
describe('extractSessionId via SSE Filtering', () => {
|
||||
it('should route session:updated events by id field', async () => {
|
||||
describe('Lifecycle Event Broadcast Contract', () => {
|
||||
it('should deliver session:updated events to all clients regardless of filter', async () => {
|
||||
// Create two sessions
|
||||
const session1 = await createSession(baseUrl);
|
||||
const session2 = await createSession(baseUrl);
|
||||
@@ -310,21 +311,23 @@ describe('Operation Lightspeed', () => {
|
||||
|
||||
const events = parseSSEEvents(receivedData);
|
||||
|
||||
// Should receive session:updated for session1 only
|
||||
// New contract: session:updated is a lifecycle event that broadcasts to ALL clients.
|
||||
// The subscription filter only applies to session:terminal.
|
||||
const updatedEvents = events.filter((e) => e.event === 'session:updated');
|
||||
const session1Updated = updatedEvents.find((e) => (e.data as any).id === session1);
|
||||
const session2Updated = updatedEvents.find((e) => (e.data as any).id === session2);
|
||||
|
||||
expect(session1Updated).toBeDefined();
|
||||
expect(session2Updated).toBeUndefined();
|
||||
expect(session2Updated).toBeDefined();
|
||||
|
||||
// Cleanup
|
||||
await deleteSession(baseUrl, session1);
|
||||
await deleteSession(baseUrl, session2);
|
||||
});
|
||||
|
||||
it('should filter session:deleted by session ID (sessionId extraction from id field)', async () => {
|
||||
// Tests extractSessionId's fallback path: session:* events use `id` not `sessionId`
|
||||
it('should deliver session:deleted events to all clients regardless of filter', async () => {
|
||||
// New contract: lifecycle events (session:*) broadcast to every connected client;
|
||||
// the per-client filter no longer gates them. Only session:terminal is filtered.
|
||||
const target = await createSession(baseUrl);
|
||||
const other = await createSession(baseUrl);
|
||||
|
||||
@@ -367,13 +370,12 @@ describe('Operation Lightspeed', () => {
|
||||
|
||||
const events = parseSSEEvents(receivedData);
|
||||
|
||||
// Target deletion should arrive (extractSessionId matches `id` field for session:* events)
|
||||
// Both deletions arrive regardless of the per-client filter
|
||||
const targetDeleted = events.find((e) => e.event === 'session:deleted' && (e.data as any).id === target);
|
||||
expect(targetDeleted).toBeDefined();
|
||||
|
||||
// Other deletion should NOT arrive
|
||||
const otherDeleted = events.find((e) => e.event === 'session:deleted' && (e.data as any).id === other);
|
||||
expect(otherDeleted).toBeUndefined();
|
||||
expect(otherDeleted).toBeDefined();
|
||||
});
|
||||
});
|
||||
|
||||
@@ -488,13 +490,13 @@ describe('Operation Lightspeed', () => {
|
||||
expect(events.find((e) => e.event === 'init')).toBeDefined();
|
||||
});
|
||||
|
||||
it('should handle multiple SSE clients with different filters', async () => {
|
||||
it('should fan lifecycle events out to all SSE clients regardless of filter', async () => {
|
||||
const session1 = await createSession(baseUrl);
|
||||
const session2 = await createSession(baseUrl);
|
||||
|
||||
// Client A: subscribes to session1
|
||||
// Client B: subscribes to session2
|
||||
// Client C: no filter (all events)
|
||||
// Client A: subscribes to session1, Client B: subscribes to session2, Client C: no filter.
|
||||
// Under the broadcast contract, all three see every session:deleted event — the filter
|
||||
// only narrows session:terminal traffic.
|
||||
const controllerA = new AbortController();
|
||||
const controllerB = new AbortController();
|
||||
const controllerC = new AbortController();
|
||||
@@ -585,15 +587,13 @@ describe('Operation Lightspeed', () => {
|
||||
const eventsB = parseSSEEvents(dataB);
|
||||
const eventsC = parseSSEEvents(dataC);
|
||||
|
||||
// Client A: sees session1 deleted, not session2
|
||||
// Every client sees both deletions — lifecycle events are not filter-gated.
|
||||
expect(eventsA.find((e) => e.event === 'session:deleted' && (e.data as any).id === session1)).toBeDefined();
|
||||
expect(eventsA.find((e) => e.event === 'session:deleted' && (e.data as any).id === session2)).toBeUndefined();
|
||||
expect(eventsA.find((e) => e.event === 'session:deleted' && (e.data as any).id === session2)).toBeDefined();
|
||||
|
||||
// Client B: sees session2 deleted, not session1
|
||||
expect(eventsB.find((e) => e.event === 'session:deleted' && (e.data as any).id === session1)).toBeDefined();
|
||||
expect(eventsB.find((e) => e.event === 'session:deleted' && (e.data as any).id === session2)).toBeDefined();
|
||||
expect(eventsB.find((e) => e.event === 'session:deleted' && (e.data as any).id === session1)).toBeUndefined();
|
||||
|
||||
// Client C: sees both
|
||||
expect(eventsC.find((e) => e.event === 'session:deleted' && (e.data as any).id === session1)).toBeDefined();
|
||||
expect(eventsC.find((e) => e.event === 'session:deleted' && (e.data as any).id === session2)).toBeDefined();
|
||||
});
|
||||
@@ -991,13 +991,13 @@ describe('Operation Lightspeed', () => {
|
||||
});
|
||||
|
||||
// ═══════════════════════════════════════════════════════════════
|
||||
// extractSessionId — Additional Edge Cases
|
||||
// Lifecycle Event Broadcast — Additional Edge Cases
|
||||
// ═══════════════════════════════════════════════════════════════
|
||||
|
||||
describe('extractSessionId — Edge Cases via SSE', () => {
|
||||
describe('Lifecycle Event Broadcast — Edge Cases via SSE', () => {
|
||||
it('should treat non-session: events with id field as global (not filtered)', async () => {
|
||||
// Events like case:created have an `id` field but aren't session:* events.
|
||||
// extractSessionId should NOT use the `id` field for non-session:* events.
|
||||
// Under the broadcast contract they reach every connected client.
|
||||
const controller = new AbortController();
|
||||
let receivedData = '';
|
||||
|
||||
@@ -1052,8 +1052,9 @@ describe('Operation Lightspeed', () => {
|
||||
}
|
||||
});
|
||||
|
||||
it('should deliver session:created for a newly created session to unfiltered client but not mismatched filter', async () => {
|
||||
// session:created uses `id` field and starts with `session:` — extractSessionId should match it
|
||||
it('should deliver session:created to every client, even those with a mismatched filter', async () => {
|
||||
// Under the broadcast contract, lifecycle events ignore the per-client filter.
|
||||
// A client subscribed only to `existing` still receives `session:created` for `newSession`.
|
||||
const existing = await createSession(baseUrl);
|
||||
|
||||
// Subscribe to existing session only
|
||||
@@ -1091,9 +1092,9 @@ describe('Operation Lightspeed', () => {
|
||||
}
|
||||
|
||||
const events = parseSSEEvents(receivedData);
|
||||
// session:created for newSession should be filtered OUT (id doesn't match our filter)
|
||||
// session:created reaches the filtered client even though its id doesn't match the filter.
|
||||
const createdEvent = events.find((e) => e.event === 'session:created' && (e.data as any).id === newSession);
|
||||
expect(createdEvent).toBeUndefined();
|
||||
expect(createdEvent).toBeDefined();
|
||||
|
||||
await Promise.all([deleteSession(baseUrl, existing), deleteSession(baseUrl, newSession)]);
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user