diff --git a/examples/desktop/README.md b/examples/desktop/README.md index 42250e4796..712645211d 100644 --- a/examples/desktop/README.md +++ b/examples/desktop/README.md @@ -104,6 +104,17 @@ a toggle), and `context/message` / `steering/message` events surface as and why. A `Cancel` button appears mid-turn and cuts the stream via `session/cancel`. +The wire accepts only one in-flight prompt per session, so pressing Enter +while a turn is running **queues** your message instead of erroring: a strip +above the composer shows one chip per queued message with a `queued N` +counter, and each chip lets you edit the text inline, delete it, or bump it to +the front (**send next**). When the turn ends, the head of the queue is sent +automatically — one message per turn completion — so a burst of follow-ups +plays out in order without you babysitting each turn. The queue is +per-session (switching sessions shows that session's queue), survives a +cancelled turn, and is cleared with a notice if the runtime restarts or you +switch profiles. + ![](docs/readme-shots/03-chat.png) #### Session Tree diff --git a/examples/desktop/docs/qa-msg-queue/01-inflight-two-queued.png b/examples/desktop/docs/qa-msg-queue/01-inflight-two-queued.png new file mode 100644 index 0000000000..605b4d8480 Binary files /dev/null and b/examples/desktop/docs/qa-msg-queue/01-inflight-two-queued.png differ diff --git a/examples/desktop/docs/qa-msg-queue/02-after-turn-end-one-left.png b/examples/desktop/docs/qa-msg-queue/02-after-turn-end-one-left.png new file mode 100644 index 0000000000..5b22de61db Binary files /dev/null and b/examples/desktop/docs/qa-msg-queue/02-after-turn-end-one-left.png differ diff --git a/examples/desktop/scripts/qa-cdp-shoot-msg-queue.mjs b/examples/desktop/scripts/qa-cdp-shoot-msg-queue.mjs new file mode 100644 index 0000000000..3080036d66 --- /dev/null +++ b/examples/desktop/scripts/qa-cdp-shoot-msg-queue.mjs @@ -0,0 +1,243 @@ +// QA verification script for lane-msg-queue. Boots an isolated Electron on +// its own CDP port, drives the composer message queue entirely through the +// renderer's `window.__dshRenderer` seam (no live model needed), and captures +// two screenshots into docs/qa-msg-queue/: +// +// 01-inflight-two-queued.png — a turn in flight with TWO messages queued; +// the strip shows two chips + a "queued 2" counter. +// 02-after-turn-end-one-left.png — after turn/end auto-drains the head, the +// first queued message appears as a sent bubble and the strip shows one +// remaining chip ("queued 1"). +// +// Isolation follows the 2026-07-18 postmortem baked into +// scripts/qa-cdp-shoot-nav-optional.mjs: +// 1. --user-data-dir= isolates Chromium userdata. +// 2. DSH_DESKTOP_HOME= isolates the main-process config root so we +// never write into ~/.dsh-desktop. +// 3. own CDP port (≥9290) so a lingering Electron helper from another +// lane's shoot can't hijack our DevTools endpoint. +// +// Why drive via the seam and not a real prompt: the message-queue behaviour +// is a pure renderer concern (enqueue-on-inflight + auto-drain on turn/end). +// `window.__dshRenderer` exposes send() + onSessionEvent() + listMsgQueue() +// ungated, so we can create an in-flight turn, enqueue two messages, and step +// the turn to completion deterministically — no daemon round-trip, no model +// key. The keyless daemon-echo include still boots the runtime so the shell +// is in its normal chat state. + +import { spawn } from 'node:child_process' +import { existsSync, mkdirSync, writeFileSync, rmSync, statSync } from 'node:fs' +import { resolve, join } from 'node:path' +import { setTimeout as sleep } from 'node:timers/promises' +import { tmpdir } from 'node:os' + +const WORKTREE = resolve(process.env.DSH_WORKTREE || process.cwd()) +const PARENT = resolve(process.env.DSH_REPO || '/Users/ziya/harness/dsh-desktop-demo') +const ELECTRON = join(PARENT, 'node_modules/.bin/electron') +const CDP_PORT = Number(process.env.DSH_MSG_QUEUE_PORT || 9290) +const OUTDIR = join(WORKTREE, 'docs/qa-msg-queue') + +if (!existsSync(ELECTRON)) { + console.error(`electron binary not found at ${ELECTRON}`) + process.exit(2) +} +mkdirSync(OUTDIR, { recursive: true }) + +function seedHome(dshHome) { + // Minimal keyless overlay + config + onboarded sentinel so the first-run + // modal doesn't fire and steal the window / rewrite our seed. daemon-echo + // is the keyless demo runtime; we never send a real prompt through it. + const seedOverlay = [ + '# QA msg-queue shoot seed overlay (tmp, per-run).', + 'plugins:', + ` - "@cordisjs/plugin-include":`, + ` path: ${join(WORKTREE, 'config/daemon-echo.yml')}`, + '', + ].join('\n') + writeFileSync(join(dshHome, 'user-overlay.cordis.yml'), seedOverlay) + writeFileSync(join(dshHome, 'config.json'), JSON.stringify({ role: 'coding', approvalMode: 'never' }, null, 2)) + writeFileSync(join(dshHome, '.onboarded'), new Date().toISOString()) +} + +async function bootElectron(dshHome, userData, port) { + const child = spawn(ELECTRON, [ + `--remote-debugging-port=${port}`, + `--user-data-dir=${userData}`, + '--disable-gpu', + '--no-sandbox', + '.', + ], { + cwd: WORKTREE, + env: { ...process.env, DSH_DESKTOP_HOME: dshHome, DSH_MAXIMIZE: '1' }, + stdio: ['ignore', 'pipe', 'pipe'], + }) + const logs = [] + child.stdout.on('data', d => logs.push(String(d))) + child.stderr.on('data', d => logs.push(String(d))) + for (let i = 0; i < 40; i++) { + await sleep(500) + try { + const r = await fetch(`http://localhost:${port}/json/list`) + if (r.ok) return { child, logs } + } catch {} + } + child.kill('SIGKILL') + console.error('electron CDP did not come up in 20s. logs:\n' + logs.join('')) + process.exit(3) +} + +async function newCdp(port) { + const targets = await (await fetch(`http://localhost:${port}/json/list`)).json() + const target = targets.find(t => t.type === 'page') + if (!target) throw new Error('no page target on port ' + port) + const ws = new WebSocket(target.webSocketDebuggerUrl) + await new Promise((ok, err) => { ws.onopen = ok; ws.onerror = e => err(e) }) + let id = 1 + const pending = new Map() + ws.onmessage = ev => { + const data = typeof ev.data === 'string' ? ev.data : String(ev.data) + let msg + try { msg = JSON.parse(data) } catch { return } + if (msg.id != null && pending.has(msg.id)) { + const [ok, err] = pending.get(msg.id); pending.delete(msg.id) + if (msg.error) err(new Error(msg.error.message)); else ok(msg.result) + } + } + const call = (m, p = {}, ms = 15000) => new Promise((ok, err) => { + const _id = id++ + const t = setTimeout(() => { pending.delete(_id); err(new Error('cdp timeout: ' + m)) }, ms) + pending.set(_id, [v => { clearTimeout(t); ok(v) }, e => { clearTimeout(t); err(e) }]) + ws.send(JSON.stringify({ id: _id, method: m, params: p })) + }) + const evj = async expr => { + const r = await call('Runtime.evaluate', { + expression: `(async()=>{try{return (${expr})}catch(e){return {__err:String(e)}}})()`, + returnByValue: true, awaitPromise: true, + }) + if (r.exceptionDetails) throw new Error(r.exceptionDetails.exception?.description || r.exceptionDetails.text) + return r.result?.value + } + return { ws, call, evj } +} + +// Clip to the composer footer so the strip is unmistakably in frame (not a +// giant empty chat pane). Falls back to a full-page shot if the selector +// isn't found. +async function shoot(call, evj, name) { + const clip = await evj(` + (() => { + const f = document.querySelector('.composer.composer-shell') + if (!f) return null + const r = f.getBoundingClientRect() + // pad upward so a freshly-sent bubble just above the composer is caught. + const top = Math.max(0, r.y - 220) + return { x: r.x, y: top, width: r.width, height: (r.bottom - top) + 12, scale: 1 } + })() + `) + const shotArgs = { format: 'png', captureBeyondViewport: true } + if (clip) shotArgs.clip = clip + const shot = await call('Page.captureScreenshot', shotArgs, 30000) + if (!shot || !shot.data) throw new Error('captureScreenshot returned no data') + const outPath = join(OUTDIR, name) + writeFileSync(outPath, Buffer.from(shot.data, 'base64')) + const bytes = statSync(outPath).size + console.log(` wrote ${outPath} (${bytes} bytes)`) + if (bytes < 20000) throw new Error(`screenshot ${name} is only ${bytes} bytes (<20KB) — likely blank`) + return { path: outPath, bytes } +} + +async function main() { + const dshHome = join(tmpdir(), 'dsh-msg-queue-home') + const userData = join(tmpdir(), 'dsh-msg-queue-userdata') + for (const dir of [dshHome, userData]) { + try { rmSync(dir, { recursive: true, force: true }) } catch {} + mkdirSync(dir, { recursive: true }) + } + seedHome(dshHome) + console.log(`booting on CDP :${CDP_PORT}`) + const { child } = await bootElectron(dshHome, userData, CDP_PORT) + try { + await sleep(1500) + const { call, evj } = await newCdp(CDP_PORT) + await call('Page.enable') + + // Step 1: create a session, make it active, and open an in-flight turn. + // We drive the renderer directly through the always-exposed seam. + const setup = await evj(` + (async () => { + const R = window.__dshRenderer + if (!R) return { __err: 'no __dshRenderer seam' } + R.ensureSession('qa-mq', { title: 'Message queue demo', header: {}, hasUserMessage: true }) + await R.selectSession('qa-mq') + R.onSessionEvent('qa-mq', { type: 'turn/start', seq: 1 }) + // Enqueue two follow-ups via the real composer send() path. + const input = document.getElementById('input') + input.value = 'Summarise the three failing tests' + await R.send() + input.value = 'Then open a PR with the fix' + await R.send() + return { active: R.getActiveSessionId(), queued: R.listMsgQueue('qa-mq').map(x => x.text) } + })() + `) + console.log(' setup:', JSON.stringify(setup)) + if (setup && setup.__err) throw new Error(setup.__err) + if (!setup || setup.queued.length !== 2) throw new Error('expected 2 queued messages, got ' + JSON.stringify(setup)) + + // Assert the strip is actually visible with two chips before shooting. + const stripA = await evj(` + (() => { + const s = document.getElementById('msg-queue-strip') + return { + hidden: s.hidden, + chips: s.querySelectorAll('.msg-queue-chip').length, + counter: (s.querySelector('.msg-queue-count') || {}).textContent || '', + } + })() + `) + console.log(' strip (inflight):', JSON.stringify(stripA)) + if (stripA.hidden || stripA.chips !== 2) throw new Error('strip not showing 2 chips: ' + JSON.stringify(stripA)) + await sleep(300) + const shotA = await shoot(call, evj, '01-inflight-two-queued.png') + + // Step 2: end the turn. turn/end auto-drains exactly one item — the head + // renders as an optimistic user bubble and the strip drops to one chip. + const drained = await evj(` + (async () => { + const R = window.__dshRenderer + R.onSessionEvent('qa-mq', { type: 'turn/end', seq: 2 }) + await new Promise(r => setTimeout(r, 60)) + const bubbles = Array.from(document.querySelectorAll('#stream .msg.user .role-label')).length + return { remaining: R.listMsgQueue('qa-mq').map(x => x.text), userBubbles: bubbles } + })() + `) + console.log(' after turn/end:', JSON.stringify(drained)) + if (!drained || drained.remaining.length !== 1) throw new Error('expected 1 remaining, got ' + JSON.stringify(drained)) + + const stripB = await evj(` + (() => { + const s = document.getElementById('msg-queue-strip') + return { + hidden: s.hidden, + chips: s.querySelectorAll('.msg-queue-chip').length, + counter: (s.querySelector('.msg-queue-count') || {}).textContent || '', + } + })() + `) + console.log(' strip (after drain):', JSON.stringify(stripB)) + if (stripB.hidden || stripB.chips !== 1) throw new Error('strip not showing 1 chip after drain: ' + JSON.stringify(stripB)) + await sleep(300) + const shotB = await shoot(call, evj, '02-after-turn-end-one-left.png') + + console.log('\n--- SUMMARY ---') + console.log(`inflight (2 queued): ${shotA.path} (${shotA.bytes} bytes)`) + console.log(`after drain (1 left): ${shotB.path} (${shotB.bytes} bytes)`) + } finally { + try { child.kill('SIGKILL') } catch {} + for (let i = 0; i < 20; i++) { + await sleep(500) + try { await fetch(`http://localhost:${CDP_PORT}/json/list`) } catch { break } + } + } +} + +main().catch(err => { console.error(err); process.exit(1) }) diff --git a/examples/desktop/src/renderer/index.html b/examples/desktop/src/renderer/index.html index 993ca81f5a..cadcddd1e3 100644 --- a/examples/desktop/src/renderer/index.html +++ b/examples/desktop/src/renderer/index.html @@ -517,6 +517,13 @@ matching the widget action catalog (see next-actions.js). Empty by default; hidden entirely when no suggestions land. --> + +
@@ -1487,6 +1494,7 @@ drill-down view models. Loaded before subagent-view.js so buildInlineSubagentTrace can pick up the drill-down helpers. --> + diff --git a/examples/desktop/src/renderer/msg-queue-model.js b/examples/desktop/src/renderer/msg-queue-model.js new file mode 100644 index 0000000000..1dc44640da --- /dev/null +++ b/examples/desktop/src/renderer/msg-queue-model.js @@ -0,0 +1,142 @@ +// Pure model for the composer message queue (lane-msg-queue). +// +// The DSH wire accepts exactly ONE in-flight prompt per session — a second +// session/prompt while a turn is running fails on the wire. So when the user +// hits Enter mid-turn we don't drop the text and we don't error: we park it +// in a per-session FIFO and auto-send the head when the turn completes. +// +// This module is the data structure only — no DOM, no wire. The renderer +// owns the "when to enqueue vs send" decision and the "drain on turn end" +// timing; this file just holds the queues and enforces their invariants so +// the ordering/isolation semantics can be locked in `node --test` without an +// Electron harness. +// +// Shape: +// queues: Map> +// Item ids are monotonic across the whole module (not per-session) so a +// chip's id is globally unique — the UI keys DOM nodes on it and never has +// to disambiguate by session. +// +// Invariants the renderer relies on: +// - FIFO: enqueue appends to the tail; drain pops the head. +// - promote(id) moves that item to the head (the "send next" affordance), +// so the very next drain sends it regardless of arrival order. +// - Per-session isolation: no operation on session A can read, mutate, or +// drain session B. A drain of an empty/unknown session returns null. +// - Empty / whitespace-only text never enqueues (returns null) — a blank +// chip is a trap the user can't tell apart from a real queued message. + +'use strict' + +function createMsgQueue() { + /** @type {Map>} */ + const queues = new Map() + let _seq = 0 + + function nextId() { _seq += 1; return _seq } + + function ensure(sessionId) { + let q = queues.get(sessionId) + if (!q) { q = []; queues.set(sessionId, q) } + return q + } + + // Append `text` to sessionId's queue. Returns the new item id, or null if + // the text is empty/whitespace-only (never enqueue a blank). + function enqueue(sessionId, text) { + if (sessionId == null) return null + const trimmed = typeof text === 'string' ? text.trim() : '' + if (!trimmed) return null + const id = nextId() + ensure(sessionId).push({ id, text: trimmed }) + return id + } + + // Drop the item with `id` from sessionId's queue. Returns true if removed. + function remove(sessionId, id) { + const q = queues.get(sessionId) + if (!q) return false + const i = q.findIndex((it) => it.id === id) + if (i < 0) return false + q.splice(i, 1) + return true + } + + // Rewrite the text of an existing queued item. Empty/whitespace text is + // rejected (returns false) — the caller should treat that as "delete" if + // it wants blank edits to remove the chip. Returns true on success. + function update(sessionId, id, text) { + const q = queues.get(sessionId) + if (!q) return false + const trimmed = typeof text === 'string' ? text.trim() : '' + if (!trimmed) return false + const it = q.find((x) => x.id === id) + if (!it) return false + it.text = trimmed + return true + } + + // Move `id` to the head of its session queue (the "send next" affordance). + // No-op returning false if the item isn't found; a single-item or + // already-head move succeeds (idempotent) and returns true. + function promote(sessionId, id) { + const q = queues.get(sessionId) + if (!q) return false + const i = q.findIndex((it) => it.id === id) + if (i < 0) return false + if (i === 0) return true + const [it] = q.splice(i, 1) + q.unshift(it) + return true + } + + // Pop and return the head item ({id,text}) of sessionId's queue, or null + // when the queue is empty/unknown. This is the single "send one on turn + // end" primitive — the renderer calls it exactly once per turn completion. + function drain(sessionId) { + const q = queues.get(sessionId) + if (!q || q.length === 0) return null + return q.shift() + } + + // Return a shallow copy of sessionId's queue (safe to render from; callers + // can't mutate the live array). Empty array for an unknown session. + function list(sessionId) { + const q = queues.get(sessionId) + return q ? q.map((it) => ({ id: it.id, text: it.text })) : [] + } + + function size(sessionId) { + const q = queues.get(sessionId) + return q ? q.length : 0 + } + + // Empty one session's queue. Returns the number of items dropped. + function clear(sessionId) { + const q = queues.get(sessionId) + if (!q) return 0 + const n = q.length + queues.delete(sessionId) + return n + } + + // Empty every session's queue (runtime crash / profile switch — the old + // session ids are a different namespace against the restarted daemon, so + // holding their queued text would send it into the void). Returns total + // items dropped across all sessions. + function clearAll() { + let n = 0 + for (const q of queues.values()) n += q.length + queues.clear() + return n + } + + return { enqueue, remove, update, promote, drain, list, size, clear, clearAll } +} + +if (typeof module !== 'undefined' && module.exports) { + module.exports = { createMsgQueue } +} +if (typeof window !== 'undefined') { + window.__dshMsgQueueModel = { createMsgQueue } +} diff --git a/examples/desktop/src/renderer/renderer.js b/examples/desktop/src/renderer/renderer.js index 08fa96bccc..2d22edba58 100644 --- a/examples/desktop/src/renderer/renderer.js +++ b/examples/desktop/src/renderer/renderer.js @@ -99,6 +99,19 @@ const statusText = document.getElementById('status-text') const modelBadge = document.getElementById('model-badge') const profileSelect = document.getElementById('profile') const newSessionBtn = document.getElementById('new-session') +const msgQueueStripEl = document.getElementById('msg-queue-strip') + +// Composer message queue (lane-msg-queue). The wire accepts one in-flight +// prompt per session; a mid-turn Enter would otherwise fire a second +// session/prompt that the daemon rejects. Instead we park the text in a +// per-session FIFO (msg-queue-model.js) and auto-drain the head when the +// turn completes. The model is a pure data structure; all DOM + wire timing +// lives in this file (see §"message queue"). Guarded so a stripped build +// without the module script simply disables queueing (send stays direct). +const msgQueue = (typeof window !== 'undefined' && window.__dshMsgQueueModel + && typeof window.__dshMsgQueueModel.createMsgQueue === 'function') + ? window.__dshMsgQueueModel.createMsgQueue() + : null const state = { activeSessionId: null, @@ -586,6 +599,9 @@ async function selectSession(id) { // their seed boundary. During live streaming we grow the same map when // subagent.started fires. installKnownForkMarkers(id) + // Per-session queue: repaint the strip for the session we just switched to. + // Strict isolation — the strip only ever shows the active session's queue. + renderMsgQueueStrip() } async function replayHistory(id) { @@ -4856,6 +4872,11 @@ function onSessionEvent(sessionId, event) { // without a badge instead of one from the previous turn. const startData = event.data || event meta.currentTurnTrigger = (startData && startData.trigger) || null + // Arm the once-per-turn queue drain. turn/end and session.finished both + // signal completion (a clean turn emits turn/end; an errored/cancelled + // turn may arrive only as session.finished). Whichever fires first + // consumes this flag so the queue drains exactly one item per turn. + meta._turnDrainPending = true if (sessionId === state.activeSessionId) { state.inflightTurn = true; updateCancelButton(); updateCompactButton(); updateForkButtons() } renderSessionList() } @@ -4864,6 +4885,11 @@ function onSessionEvent(sessionId, event) { meta.currentTurnTrigger = null if (sessionId === state.activeSessionId) { state.inflightTurn = false; updateCancelButton(); updateCompactButton(); updateForkButtons() } renderSessionList() + // Auto-drain one queued message for this session now that its turn is + // done. Fires for the exact turn-completion condition that re-enables + // Cancel; a cancelled turn ends with turn/end too, so the queue drains + // there just like a clean completion (queue survives Cancel by design). + void drainMsgQueueOnce(sessionId) } if (event.time && event.time > meta.lastEventTime) meta.lastEventTime = event.time // §1.1 trace bucket: non-boundary events get bucketed into @@ -6071,9 +6097,133 @@ function onInterruptInvalidate({ interruptId, reason }) { // -- send / cancel ----------------------------------------------------------- +// -- message queue ---------------------------------------------------------- +// +// Renders the strip above the composer for the ACTIVE session's queue. The +// strip is hidden when the queue is empty; each queued message is one chip +// with truncated text (click to inline-edit), a "send next" (promote) button, +// and a delete (×). A "queued N" counter badge leads the row. Per-session: +// selectSession + every mutation calls this, so switching sessions shows that +// session's queue and nothing bleeds across sessions. + +const MSG_QUEUE_CHIP_MAX = 48 + +// Truncate for the chip label. Kept in sync with MSG_QUEUE_CHIP_MAX; the full +// text lives on the chip's title attribute so hover reveals the whole thing. +function truncateQueueText(text) { + const t = String(text || '') + return t.length > MSG_QUEUE_CHIP_MAX ? t.slice(0, MSG_QUEUE_CHIP_MAX - 1) + '…' : t +} + +// Swap a chip's static label for an inline seeded with the full text. +// Enter / blur commits via msgQueue.update (blank ⇒ delete, matching the +// model's "never keep a blank" rule); Escape re-renders to cancel. +function beginQueueChipEdit(sessionId, id, chipEl, fullText) { + if (!chipEl) return + const input = document.createElement('input') + input.type = 'text' + input.className = 'msg-queue-chip-edit' + input.value = fullText + let done = false + const commit = (save) => { + if (done) return + done = true + if (save) { + const v = input.value.trim() + if (!v) msgQueue.remove(sessionId, id) + else msgQueue.update(sessionId, id, v) + } + renderMsgQueueStrip() + } + input.addEventListener('keydown', (e) => { + if (e.key === 'Enter') { e.preventDefault(); commit(true) } + else if (e.key === 'Escape') { e.preventDefault(); commit(false) } + }) + input.addEventListener('blur', () => commit(true)) + chipEl.textContent = '' + chipEl.appendChild(input) + input.focus() +} + +function makeQueueChip(sessionId, item) { + const chip = document.createElement('span') + chip.className = 'msg-queue-chip' + chip.dataset.queueId = String(item.id) + + const label = document.createElement('button') + label.type = 'button' + label.className = 'msg-queue-chip-text' + label.textContent = truncateQueueText(item.text) + label.title = item.text + '\n(click to edit)' + label.addEventListener('click', () => beginQueueChipEdit(sessionId, item.id, chip, item.text)) + + // "send next" — promote to head so the next drain sends this item first. + const promoteBtn = document.createElement('button') + promoteBtn.type = 'button' + promoteBtn.className = 'msg-queue-chip-btn msg-queue-chip-promote' + promoteBtn.title = 'Send this one next' + promoteBtn.setAttribute('aria-label', 'Send this message next') + promoteBtn.textContent = '↑' + promoteBtn.addEventListener('click', () => { msgQueue.promote(sessionId, item.id); renderMsgQueueStrip() }) + + const delBtn = document.createElement('button') + delBtn.type = 'button' + delBtn.className = 'msg-queue-chip-btn msg-queue-chip-del' + delBtn.title = 'Remove from queue' + delBtn.setAttribute('aria-label', 'Remove queued message') + delBtn.textContent = '×' + delBtn.addEventListener('click', () => { msgQueue.remove(sessionId, item.id); renderMsgQueueStrip() }) + + chip.append(label, promoteBtn, delBtn) + return chip +} + +function renderMsgQueueStrip() { + if (!msgQueueStripEl || !msgQueue) return + const sid = state.activeSessionId + const items = sid ? msgQueue.list(sid) : [] + msgQueueStripEl.textContent = '' + if (items.length === 0) { + msgQueueStripEl.hidden = true + return + } + msgQueueStripEl.hidden = false + + const badge = document.createElement('span') + badge.className = 'msg-queue-count' + badge.textContent = `queued ${items.length}` + msgQueueStripEl.appendChild(badge) + + // Only the head is a promote target worth showing an active arrow on, but + // we render the affordance on every non-head chip; the head chip's arrow is + // greyed (already next) via the .is-head class. + for (let i = 0; i < items.length; i++) { + const chip = makeQueueChip(sid, items[i]) + if (i === 0) chip.classList.add('is-head') + msgQueueStripEl.appendChild(chip) + } +} + async function send() { const text = inputEl.value.trim() if (!text) return + // Mid-turn Enter: the active session already has a prompt in flight (the + // wire accepts only one at a time). Park this text in the session's queue + // and clear the composer as if sent — it auto-drains on turn/end. Only + // applies when there's an active session AND a turn is running; a fresh + // "+" session or an idle session falls through to the normal send path. + if (msgQueue && state.activeSessionId && state.inflightTurn) { + const sid = state.activeSessionId + const meta = state.sessions.get(sid) + if (meta && !meta.hasUserMessage) { + meta.hasUserMessage = true + if (!meta.title) meta.title = text.slice(0, 40) + } + msgQueue.enqueue(sid, text) + inputEl.value = '' + renderMsgQueueStrip() + return + } if (!state.activeSessionId) { const { id } = await window.dsh.newSession() ensureSession(id, { title: text.slice(0, 40), hasUserMessage: true }) @@ -6091,22 +6241,66 @@ async function send() { } } const sid = state.activeSessionId - // Optimistic bubble — the server echoes user/message shortly after, and the - // event handler adopts this element in place instead of re-rendering. - appendMessage({ role: 'user', text, optimistic: true }) - updateEmptyStateVisibility() - inputEl.value = '' - sendBtn.disabled = true - state.inflightTurn = true; updateCancelButton(); updateCompactButton(); updateForkButtons() + await dispatchPrompt(sid, text, { clearComposer: true }) +} + +// The shared "actually send a prompt" tail: optimistic bubble + inflight +// flag + sendPrompt, with the existing error surface. Extracted from send() +// so the queue's auto-drain sends a parked message through the exact same +// path a live Enter takes. `opts.clearComposer` empties the textarea (live +// send); the drain path leaves it alone (the user may be mid-typing the next +// message). The optimistic bubble + inflight flag only apply when `sid` is +// the active session — a drain firing for a background session's turn/end +// must not paint into the foreground stream (the echoed user/message renders +// when that session is next viewed). +async function dispatchPrompt(sid, text, opts = {}) { + const isActive = sid === state.activeSessionId + if (isActive) { + // Optimistic bubble — the server echoes user/message shortly after, and + // the event handler adopts this element in place instead of re-rendering. + appendMessage({ role: 'user', text, optimistic: true }) + updateEmptyStateVisibility() + if (opts.clearComposer) inputEl.value = '' + sendBtn.disabled = true + state.inflightTurn = true; updateCancelButton(); updateCompactButton(); updateForkButtons() + } + const meta = state.sessions.get(sid) + if (meta) meta.running = true try { await window.dsh.sendPrompt(sid, text) } catch (err) { + // Surface the error on the stream (matches the live-send path). We do + // NOT clear the rest of the queue — a failed drain leaves the remaining + // items parked so the user can retry / edit rather than losing them. appendSystem(`error: ${err.message}`) } finally { - sendBtn.disabled = false + if (isActive) sendBtn.disabled = false } } +// Drain exactly one queued message for `sessionId` and send it — guarded so +// it fires at most once per turn even when both turn/end and +// session.finished arrive for the same turn. The `_turnDrainPending` flag is +// armed on turn/start and consumed here. ONE item per turn completion, so +// each queued follow-up waits for its own turn to finish (the head's +// turn/end drains the next head, and so on). Runs after the completion +// handler has flipped inflightTurn off, so dispatchPrompt re-arms it for the +// drained send. +async function drainMsgQueueOnce(sessionId) { + if (!msgQueue) return + const meta = state.sessions.get(sessionId) + // Guard: only the first completion signal for this turn drains. Absent a + // meta (defensive) we still allow one drain. + if (meta) { + if (!meta._turnDrainPending) return + meta._turnDrainPending = false + } + const next = msgQueue.drain(sessionId) + renderMsgQueueStrip() + if (!next) return + await dispatchPrompt(sessionId, next.text, { clearComposer: false }) +} + async function cancel() { if (!state.activeSessionId) return try { @@ -6171,6 +6365,11 @@ window.dsh.onNotify(({ method, params }) => { if (params.sessionId === state.activeSessionId) { state.inflightTurn = false; updateCancelButton(); updateForkButtons() } // Refresh listing so title/lastEventTime catch up. void refreshSessionList() + // Drain-once: if this session's turn ended via session.finished without a + // clean turn/end (error/cancel paths), the guard fires the single drain + // here. If turn/end already drained, `_turnDrainPending` is false and this + // is a no-op. + void drainMsgQueueOnce(params.sessionId) } else if (method === 'subagent.started') { // Grow the tree eagerly: register the child session with a placeholder // parent link, refresh the sidebar, and drop a fork marker on the parent @@ -6615,6 +6814,19 @@ window.dsh.onInitialized((info) => { titleEl.textContent = 'New chat' updateEmptyStateVisibility() updateCancelButton() + // Runtime crash / profile switch = a fresh session namespace (the old + // session ids can never come back). Any messages still queued against the + // dead runtime's sessions would send into the void, so wipe every queue + // and — if we actually dropped anything — drop a one-line notice so the + // user knows their parked follow-ups didn't survive the restart. Mirrors + // the interrupt-invalidate posture (stale-against-restart state is a trap). + if (msgQueue) { + const dropped = msgQueue.clearAll() + if (dropped > 0) { + appendSystem(`runtime restarted — cleared ${dropped} queued message${dropped === 1 ? '' : 's'}`) + } + renderMsgQueueStrip() + } // New runtime = unknown compact support; button falls back to greyed // until the first successful call proves the daemon accepts the method // (or a MethodNotFound flips it to `false`). @@ -8211,6 +8423,14 @@ window.__dshRenderer = { sessionIds: Array.from(state.sessions.keys()), replayingId: state.replayingId, }), + // Message queue (lane-msg-queue) seam: expose the send entry point, the + // pure queue handle, and the strip renderer so renderer-level tests and QA + // scripts can walk the enqueue-on-inflight + auto-drain wiring without a + // real wire round-trip. `getMsgQueue` returns the live model instance. + send, + getMsgQueue: () => msgQueue, + listMsgQueue: (sid) => (msgQueue ? msgQueue.list(sid) : []), + renderMsgQueueStrip, } // `⌘.` (Ctrl-. on non-Mac) toggles every trace-event-row diff --git a/examples/desktop/src/renderer/style.css b/examples/desktop/src/renderer/style.css index 2e4dd202bb..ea7dc5ec3c 100644 --- a/examples/desktop/src/renderer/style.css +++ b/examples/desktop/src/renderer/style.css @@ -3673,6 +3673,78 @@ textarea:focus-visible { cursor: pointer; } .composer-model-warn a:hover { color: var(--accent); } + +/* Message queue strip (lane-msg-queue). Sits above the composer frame, shown + only when the active session has queued mid-turn messages. One chip per + queued message + a leading "queued N" counter. Chip grammar mirrors the + composer chip pills (999px radius, muted surface, 11.5px label) so it reads + as part of the composer, not a separate widget. */ +.msg-queue-strip { + display: flex; align-items: center; flex-wrap: wrap; + gap: 6px; + margin: 0 4px 6px 4px; + padding: 4px 2px; +} +.msg-queue-count { + display: inline-flex; align-items: center; + padding: 2px 8px; + border-radius: 999px; + background: var(--accent-soft); + color: var(--accent); + font-size: 10.5px; font-weight: 600; + letter-spacing: 0.03em; + text-transform: uppercase; + flex: 0 0 auto; +} +.msg-queue-chip { + display: inline-flex; align-items: center; gap: 2px; + padding: 2px 4px 2px 8px; + border-radius: 999px; + background: var(--surface-2); + border: 1px solid var(--border); + max-width: 280px; + transition: border-color 120ms ease, background 120ms ease; +} +.msg-queue-chip:hover { border-color: var(--accent); } +.msg-queue-chip-text { + border: none; background: transparent; + padding: 0; margin: 0; + color: var(--text-secondary); + font: inherit; font-size: 11.5px; line-height: 1.4; + max-width: 200px; + overflow: hidden; text-overflow: ellipsis; white-space: nowrap; + cursor: text; +} +.msg-queue-chip-text:hover { color: var(--text); } +.msg-queue-chip-edit { + border: 1px solid var(--accent); + border-radius: var(--radius-sm); + background: var(--surface); + color: var(--text); + font: inherit; font-size: 11.5px; line-height: 1.4; + padding: 1px 6px; + min-width: 180px; + outline: none; +} +.msg-queue-chip-btn { + display: inline-flex; align-items: center; justify-content: center; + width: 18px; height: 18px; + padding: 0; border: none; + border-radius: 999px; + background: transparent; + color: var(--text-tertiary); + font-size: 13px; line-height: 1; + cursor: pointer; + transition: background 120ms ease, color 120ms ease; +} +.msg-queue-chip-btn:hover { background: var(--surface-hover); color: var(--text); } +.msg-queue-chip-del:hover { color: var(--err); } +/* The head chip is already "next" — dim its promote arrow so the affordance + reads as inert without removing it (keeps chip layout stable). */ +.msg-queue-chip.is-head .msg-queue-chip-promote { + opacity: 0.35; + pointer-events: none; +} .composer-chip--mode { color: var(--text-secondary); } diff --git a/examples/desktop/test/msg-queue-model.test.js b/examples/desktop/test/msg-queue-model.test.js new file mode 100644 index 0000000000..a35021a711 --- /dev/null +++ b/examples/desktop/test/msg-queue-model.test.js @@ -0,0 +1,145 @@ +// Unit tests for msg-queue-model.js — the per-session FIFO backing the +// composer's mid-turn message queue. Pure data structure, no DOM. + +'use strict' + +const test = require('node:test') +const assert = require('node:assert/strict') + +const { createMsgQueue } = require('../src/renderer/msg-queue-model.js') + +test('enqueue returns a monotonic id and list preserves FIFO order', () => { + const q = createMsgQueue() + const a = q.enqueue('s1', 'first') + const b = q.enqueue('s1', 'second') + const c = q.enqueue('s1', 'third') + assert.ok(a < b && b < c, 'ids are monotonic') + assert.deepEqual(q.list('s1').map((x) => x.text), ['first', 'second', 'third']) +}) + +test('enqueue trims text and rejects empty / whitespace-only', () => { + const q = createMsgQueue() + assert.equal(q.enqueue('s1', ''), null) + assert.equal(q.enqueue('s1', ' '), null) + assert.equal(q.enqueue('s1', '\n\t '), null) + const id = q.enqueue('s1', ' hi ') + assert.ok(id) + assert.equal(q.list('s1')[0].text, 'hi', 'stored text is trimmed') + assert.equal(q.size('s1'), 1, 'rejected blanks never entered the queue') +}) + +test('enqueue with null sessionId is a no-op', () => { + const q = createMsgQueue() + assert.equal(q.enqueue(null, 'x'), null) + assert.equal(q.enqueue(undefined, 'x'), null) +}) + +test('drain pops the head and returns null when empty/unknown', () => { + const q = createMsgQueue() + assert.equal(q.drain('s1'), null, 'unknown session drains to null') + q.enqueue('s1', 'one') + q.enqueue('s1', 'two') + assert.equal(q.drain('s1').text, 'one', 'head first') + assert.equal(q.drain('s1').text, 'two') + assert.equal(q.drain('s1'), null, 'emptied queue drains to null') +}) + +test('drain-once semantics: one call removes exactly one item', () => { + const q = createMsgQueue() + q.enqueue('s1', 'a') + q.enqueue('s1', 'b') + q.enqueue('s1', 'c') + const first = q.drain('s1') + assert.equal(first.text, 'a') + assert.equal(q.size('s1'), 2, 'the other two wait for their own turn ends') + assert.deepEqual(q.list('s1').map((x) => x.text), ['b', 'c']) +}) + +test('remove drops a specific item by id; returns false when absent', () => { + const q = createMsgQueue() + const a = q.enqueue('s1', 'a') + const b = q.enqueue('s1', 'b') + assert.equal(q.remove('s1', a), true) + assert.deepEqual(q.list('s1').map((x) => x.text), ['b']) + assert.equal(q.remove('s1', a), false, 'already gone') + assert.equal(q.remove('nope', b), false, 'unknown session') +}) + +test('update rewrites text, trims, and rejects blank', () => { + const q = createMsgQueue() + const id = q.enqueue('s1', 'old') + assert.equal(q.update('s1', id, ' new '), true) + assert.equal(q.list('s1')[0].text, 'new') + assert.equal(q.update('s1', id, ' '), false, 'blank edit rejected') + assert.equal(q.list('s1')[0].text, 'new', 'text unchanged after rejected edit') + assert.equal(q.update('s1', 9999, 'x'), false, 'unknown id') +}) + +test('promote moves an item to the head so the next drain sends it', () => { + const q = createMsgQueue() + q.enqueue('s1', 'a') + const b = q.enqueue('s1', 'b') + q.enqueue('s1', 'c') + assert.equal(q.promote('s1', b), true) + assert.deepEqual(q.list('s1').map((x) => x.text), ['b', 'a', 'c']) + assert.equal(q.drain('s1').text, 'b', 'promoted item drains first') +}) + +test('promote is idempotent on the head and single-item queues', () => { + const q = createMsgQueue() + const a = q.enqueue('s1', 'a') + assert.equal(q.promote('s1', a), true, 'single item head promote succeeds') + assert.deepEqual(q.list('s1').map((x) => x.text), ['a']) + const b = q.enqueue('s1', 'b') + assert.equal(q.promote('s1', a), true, 'already-head promote succeeds') + assert.deepEqual(q.list('s1').map((x) => x.text), ['a', 'b']) + assert.equal(q.promote('s1', 4242), false, 'unknown id fails') +}) + +test('per-session isolation: operations never cross sessions', () => { + const q = createMsgQueue() + const a1 = q.enqueue('s1', 's1-a') + q.enqueue('s2', 's2-a') + q.enqueue('s2', 's2-b') + assert.equal(q.size('s1'), 1) + assert.equal(q.size('s2'), 2) + // drain s1 leaves s2 untouched + assert.equal(q.drain('s1').text, 's1-a') + assert.equal(q.size('s2'), 2, 's2 unaffected by s1 drain') + // remove/promote keyed on the wrong session is a no-op + assert.equal(q.remove('s2', a1), false, 's1 id not found in s2') + assert.equal(q.promote('s2', a1), false) + assert.deepEqual(q.list('s2').map((x) => x.text), ['s2-a', 's2-b']) +}) + +test('clear empties one session and returns the count dropped', () => { + const q = createMsgQueue() + q.enqueue('s1', 'a') + q.enqueue('s1', 'b') + q.enqueue('s2', 'c') + assert.equal(q.clear('s1'), 2) + assert.equal(q.size('s1'), 0) + assert.equal(q.size('s2'), 1, 'other session survives a targeted clear') + assert.equal(q.clear('unknown'), 0) +}) + +test('clearAll wipes every session (crash / profile switch) and totals the drop', () => { + const q = createMsgQueue() + q.enqueue('s1', 'a') + q.enqueue('s1', 'b') + q.enqueue('s2', 'c') + assert.equal(q.clearAll(), 3) + assert.equal(q.size('s1'), 0) + assert.equal(q.size('s2'), 0) + assert.deepEqual(q.list('s1'), []) +}) + +test('list returns copies — mutating the result cannot corrupt the queue', () => { + const q = createMsgQueue() + q.enqueue('s1', 'a') + const snap = q.list('s1') + snap[0].text = 'HACKED' + snap.push({ id: 999, text: 'injected' }) + assert.equal(q.list('s1')[0].text, 'a', 'internal item untouched') + assert.equal(q.size('s1'), 1, 'internal length untouched') +}) diff --git a/examples/desktop/test/renderer-collisions.test.js b/examples/desktop/test/renderer-collisions.test.js index 52f3bd3331..3923fe3edb 100644 --- a/examples/desktop/test/renderer-collisions.test.js +++ b/examples/desktop/test/renderer-collisions.test.js @@ -143,6 +143,12 @@ const NON_IIFE_ALLOWLIST = new Set([ 'intervention-timeline.js', 'compact-config-model.js', 'subagent-drilldown.js', + // lane-msg-queue (2026-07-20) composer message queue: dual-exported pure + // FIFO model — module.exports for node --test, window.__dshMsgQueueModel + // for the renderer. Same shape as compact-config-model.js; preloadPure + // require()s it so it must not be IIFE-wrapped. Sole top-level binding is + // `function createMsgQueue`, unique across the shared scope. + 'msg-queue-model.js', ]) function listRendererScripts() { diff --git a/examples/desktop/test/renderer-harness.js b/examples/desktop/test/renderer-harness.js index f7a45c0991..eb1998a7ae 100644 --- a/examples/desktop/test/renderer-harness.js +++ b/examples/desktop/test/renderer-harness.js @@ -419,6 +419,7 @@ async function loadRenderer(customStubs = {}, options = {}) { ['tool-cards.js', '__dshToolCards'], ['widgets.js', '__dshWidgets'], ['capabilities.js', '__dshCapabilities'], + ['msg-queue-model.js', '__dshMsgQueueModel'], ] for (const [file, key] of preloadPure) { const p = path.join(__dirname, '..', 'src', 'renderer', file) diff --git a/examples/desktop/test/renderer-msg-queue.test.js b/examples/desktop/test/renderer-msg-queue.test.js new file mode 100644 index 0000000000..76fcc23983 --- /dev/null +++ b/examples/desktop/test/renderer-msg-queue.test.js @@ -0,0 +1,196 @@ +// Renderer-level tests for the composer message queue (lane-msg-queue). +// +// Covers the wiring the pure msg-queue-model.js tests can't reach: +// - send() enqueues instead of sendPrompt when a turn is in flight +// - the queue strip renders/hides per the active session's queue +// - turn/end auto-drains exactly ONE queued item through sendPrompt +// - session.finished drains once too, and doesn't double-drain with turn/end +// - per-session isolation: a background session's turn/end never touches +// the foreground session's composer +// - runtime restart (onInitialized) clears every queue +// +// The harness loads the whole renderer against a DOM stub and exposes the +// send entry point + queue handle via window.__dshRenderer. + +'use strict' + +const test = require('node:test') +const assert = require('node:assert/strict') +const { loadRenderer } = require('./renderer-harness.js') + +// Drive an in-flight turn on `sid` and make it active, returning refs. +async function activeInflight(renderer, document, sid = 's1') { + renderer.ensureSession(sid, { title: 'sess', header: {}, hasUserMessage: true }) + await renderer.selectSession(sid) + renderer.onSessionEvent(sid, { type: 'turn/start', seq: 1 }) + return { input: document.getElementById('input') } +} + +function sendCalls(dsh) { + return dsh.__calls.filter((c) => c[0] === 'sendPrompt') +} + +test('send() enqueues instead of sendPrompt while a turn is in flight', async () => { + const { renderer, document, dsh } = await loadRenderer() + const { input } = await activeInflight(renderer, document) + const before = sendCalls(dsh).length + input.value = 'queued follow-up' + await renderer.send() + assert.equal(sendCalls(dsh).length, before, 'no sendPrompt fired mid-turn') + assert.deepEqual(renderer.listMsgQueue('s1').map((x) => x.text), ['queued follow-up']) + assert.equal(input.value, '', 'composer cleared as if sent') +}) + +test('whitespace-only mid-turn input never enqueues', async () => { + const { renderer, document } = await loadRenderer() + const { input } = await activeInflight(renderer, document) + input.value = ' ' + await renderer.send() + assert.equal(renderer.listMsgQueue('s1').length, 0) +}) + +test('queue strip renders one chip per queued message + counter, hides when empty', async () => { + const { renderer, document } = await loadRenderer() + const { input } = await activeInflight(renderer, document) + const strip = document.getElementById('msg-queue-strip') + assert.equal(strip.hidden, true, 'hidden with an empty queue') + input.value = 'first'; await renderer.send() + input.value = 'second'; await renderer.send() + assert.equal(strip.hidden, false, 'shown once queued') + assert.equal(strip.querySelectorAll('.msg-queue-chip').length, 2) + const badge = strip.querySelector('.msg-queue-count') + assert.match(badge.textContent, /queued 2/) +}) + +test('turn/end auto-drains exactly one queued item through sendPrompt', async () => { + const { renderer, document, dsh } = await loadRenderer() + const { input } = await activeInflight(renderer, document) + input.value = 'one'; await renderer.send() + input.value = 'two'; await renderer.send() + const before = sendCalls(dsh).length + renderer.onSessionEvent('s1', { type: 'turn/end', seq: 2 }) + // drain is async (dispatchPrompt awaits sendPrompt) — let microtasks settle. + await new Promise((r) => setTimeout(r, 5)) + const after = sendCalls(dsh) + assert.equal(after.length - before, 1, 'exactly one drained send') + assert.equal(after[after.length - 1][2], 'one', 'FIFO: head sent first') + assert.deepEqual(renderer.listMsgQueue('s1').map((x) => x.text), ['two'], + 'the second item waits for its own turn to finish') +}) + +test('a second turn/end drains the next item (one per completion)', async () => { + const { renderer, document, dsh } = await loadRenderer() + const { input } = await activeInflight(renderer, document) + input.value = 'one'; await renderer.send() + input.value = 'two'; await renderer.send() + renderer.onSessionEvent('s1', { type: 'turn/end', seq: 2 }) + await new Promise((r) => setTimeout(r, 5)) + // the drained 'one' send flipped inflightTurn back on; simulate its turn. + renderer.onSessionEvent('s1', { type: 'turn/start', seq: 3 }) + renderer.onSessionEvent('s1', { type: 'turn/end', seq: 4 }) + await new Promise((r) => setTimeout(r, 5)) + const sent = sendCalls(dsh).map((c) => c[2]) + assert.deepEqual(sent, ['one', 'two'], 'both drained in FIFO order across two turns') + assert.equal(renderer.listMsgQueue('s1').length, 0) +}) + +test('turn/end + session.finished for one turn drains only once', async () => { + const { renderer, document, dsh, listeners } = await loadRenderer() + const { input } = await activeInflight(renderer, document) + input.value = 'one'; await renderer.send() + input.value = 'two'; await renderer.send() + const before = sendCalls(dsh).length + // Some daemons emit both boundaries for the same turn. + renderer.onSessionEvent('s1', { type: 'turn/end', seq: 2 }) + await new Promise((r) => setTimeout(r, 5)) + listeners.onNotify({ method: 'session.finished', params: { sessionId: 's1', status: 'ok' } }) + await new Promise((r) => setTimeout(r, 5)) + assert.equal(sendCalls(dsh).length - before, 1, 'drain-once guard held across both signals') + assert.deepEqual(renderer.listMsgQueue('s1').map((x) => x.text), ['two']) +}) + +test('session.finished alone (no turn/end) still drains once', async () => { + const { renderer, document, dsh, listeners } = await loadRenderer() + const { input } = await activeInflight(renderer, document) + input.value = 'only'; await renderer.send() + const before = sendCalls(dsh).length + listeners.onNotify({ method: 'session.finished', params: { sessionId: 's1', status: 'error' } }) + await new Promise((r) => setTimeout(r, 5)) + assert.equal(sendCalls(dsh).length - before, 1, 'error-path completion drained the head') +}) + +test('cancelled turn (turn/end after cancel) still drains — queue survives Cancel', async () => { + const { renderer, document, dsh } = await loadRenderer() + const { input } = await activeInflight(renderer, document) + input.value = 'after cancel'; await renderer.send() + // User cancels; the wire still closes the turn with turn/end. + const before = sendCalls(dsh).length + renderer.onSessionEvent('s1', { type: 'turn/end', seq: 2 }) + await new Promise((r) => setTimeout(r, 5)) + assert.equal(sendCalls(dsh).length - before, 1, 'queued follow-up sent on the cancelled turn end') +}) + +test('per-session isolation: a background turn/end never touches the active composer', async () => { + const { renderer, document, dsh } = await loadRenderer() + // s1 active + in flight, one queued item. + const { input } = await activeInflight(renderer, document, 's1') + input.value = 's1-queued'; await renderer.send() + // s2 exists, has its own in-flight turn + queue, but is NOT active. + renderer.ensureSession('s2', { title: 's2', header: {}, hasUserMessage: true, running: true }) + renderer.getMsgQueue().enqueue('s2', 's2-queued') + // Arm s2's drain flag via a background turn/start (not the active session, + // so it doesn't touch the composer), then end it. + renderer.onSessionEvent('s2', { type: 'turn/start', seq: 8 }) + const before = sendCalls(dsh).length + renderer.onSessionEvent('s2', { type: 'turn/end', seq: 9 }) + await new Promise((r) => setTimeout(r, 5)) + const sent = sendCalls(dsh).slice(before).map((c) => [c[1], c[2]]) + assert.deepEqual(sent, [['s2', 's2-queued']], 'only s2 drained, sent against s2') + // s1's queue + composer untouched. + assert.deepEqual(renderer.listMsgQueue('s1').map((x) => x.text), ['s1-queued']) + const strip = document.getElementById('msg-queue-strip') + assert.equal(strip.querySelectorAll('.msg-queue-chip').length, 1, 'active strip still shows s1') +}) + +test('switching sessions shows the target session\'s queue (strict isolation)', async () => { + const { renderer, document } = await loadRenderer() + const { input } = await activeInflight(renderer, document, 's1') + input.value = 's1-only'; await renderer.send() + const strip = document.getElementById('msg-queue-strip') + assert.equal(strip.querySelectorAll('.msg-queue-chip').length, 1) + // Switch to a fresh idle session — its queue is empty, strip hides. + renderer.ensureSession('s2', { title: 's2', header: {}, hasUserMessage: true }) + await renderer.selectSession('s2') + assert.equal(strip.hidden, true, 'empty queue on s2 hides the strip') + // Switch back — s1's queue is intact. + await renderer.selectSession('s1') + assert.equal(strip.querySelectorAll('.msg-queue-chip').length, 1) +}) + +test('runtime restart (onInitialized) clears every queue and posts a notice', async () => { + const { renderer, document, listeners } = await loadRenderer() + const { input } = await activeInflight(renderer, document, 's1') + input.value = 'doomed'; await renderer.send() + renderer.getMsgQueue().enqueue('s2', 'also doomed') + // Fire the initialize handshake (new daemon / profile switch). + listeners.onInitialized({ serverInfo: { name: 'echo', version: '1' }, protocolVersion: 1 }) + await new Promise((r) => setTimeout(r, 5)) + assert.equal(renderer.getMsgQueue().size('s1'), 0) + assert.equal(renderer.getMsgQueue().size('s2'), 0) + assert.match(renderer.getStreamText(), /cleared 2 queued messages/) +}) + +test('drained send re-arms inflightTurn so a follow-up Enter queues again', async () => { + const { renderer, document, dsh } = await loadRenderer() + const { input } = await activeInflight(renderer, document) + input.value = 'one'; await renderer.send() + input.value = 'two'; await renderer.send() + renderer.onSessionEvent('s1', { type: 'turn/end', seq: 2 }) + await new Promise((r) => setTimeout(r, 5)) + // 'one' is now in flight (dispatchPrompt set inflightTurn). A new Enter + // should queue behind 'two', not fire a second concurrent sendPrompt. + const before = sendCalls(dsh).length + input.value = 'three'; await renderer.send() + assert.equal(sendCalls(dsh).length, before, 'follow-up queued, not sent concurrently') + assert.deepEqual(renderer.listMsgQueue('s1').map((x) => x.text), ['two', 'three']) +})