From 52bfb64cf5dc3ad1c2b0c95b2efb512c143cbd85 Mon Sep 17 00:00:00 2001 From: TadMSTR <69825253+TadMSTR@users.noreply.github.com> Date: Wed, 30 Sep 2026 14:25:09 -0400 Subject: [PATCH] perf: one list read per refresh, parallel loads, no duplicate watcher refresh (v0.12.0) - /tasks (unfiltered) and /headless-runs join one in-flight upstream GET /tasks (SharedRead). Only the read is shared; nothing is reused once it settles. Every mutation goes through mutate() -> invalidatingWrite, which invalidates before the route responds; the watcher invalidates on the raw fs event. - The UI's three reads, and a button's list + detail refresh, run in parallel. Results apply last-started-wins (Latest). - A watcher event already covered by a successful refresh that started after the change is skipped (coveredByRefresh), removing the second full refresh after every button press. - Per-task in-flight guard: a second click (Start included) is dropped before confirm(); mutations are never parallel or duplicated (task-queue-mcp F-01). - Includes 0.11.1's redirect fix. Measured vs live API: tab load 0.372 s / 3 upstream reads -> 0.250 s / 2; Park cycle 0.382 s / 4 -> 0.268 s / 3. Programme task-queue-read-perf-2026-09 part 3; vikunja#1003. Co-Authored-By: Claude Opus 5.5 --- AGENTS.md | 12 +- CHANGELOG.md | 42 ++++++ README.md | 2 +- manifest.json | 2 +- package-lock.json | 4 +- package.json | 2 +- src/index.ts | 244 ++++++++++++++++++-------------- src/refresh-rules.ts | 46 ++++++ src/server.ts | 50 ++++++- src/shared-read.ts | 65 +++++++++ src/tests/refresh-rules.test.ts | 74 ++++++++++ src/tests/shared-read.test.ts | 189 +++++++++++++++++++++++++ 12 files changed, 612 insertions(+), 120 deletions(-) create mode 100644 src/refresh-rules.ts create mode 100644 src/shared-read.ts create mode 100644 src/tests/refresh-rules.test.ts create mode 100644 src/tests/shared-read.test.ts diff --git a/AGENTS.md b/AGENTS.md index 4bff6e5..991e403 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -21,6 +21,11 @@ src/ (queueGet) and mutations (callControlApi). Extracted from server.ts so the auth and transport guards are unit-testable without booting the server. + shared-read.ts One in-flight upstream read shared by concurrent callers, and + the invalidate-after-write wrapper every mutation goes through. + refresh-rules.ts The UI's ordering rules as pure functions: last-started-wins, + and when a watcher event is already covered by a refresh. + Also owns WATCH_DEBOUNCE_MS, imported by server.ts and index.ts. queue-token.ts Loads the plugin's client token from the fixed file under $HOME, failing closed on mode or content. ws-guard.ts The WebSocket upgrade decision, as a pure function of @@ -73,6 +78,11 @@ importable. esbuild and `tsc` (`allowImportingTsExtensions`) both accept `.ts`. ## Invariants - **The plugin never reads or writes queue YAML directly.** Every mutation goes through `control-api.ts` to the MCP control API, inheriting its transition validation, `fcntl` locking, and atomic writes. Since v0.11.0 every read does too (`queueGet` → `GET /tasks`, `GET /tasks/{id}`), inheriting the queue's TTL, dead-letter and status rules; before that this file globbed the YAML and applied none of them. A new mutation means a new control-API action, never an `fs.writeFile`; a new read means an API call, never a `yamlLoad` of a queue file. The `fs.watch` on the queue root is a change trigger and reads nothing. +- **Reads in parallel, mutations never.** A refresh issues its three reads concurrently, and a button's refresh reads the list and the detail concurrently. Each button's mutation is a single awaited call, and `pendingActions` in `index.ts` drops a second click on the same task (Start included) before any `confirm()`. task-queue-mcp's audit accepted an unpark race (part 1, F-01) on the premise that clients do not send duplicate mutations. Do not relax that guard, and do not fire a mutation alongside anything else for the same task. +- **The unfiltered list is shared as a READ, never as a cached RESULT.** `/tasks` (no filters) and `queueIndexByPrefix` join one in-flight `GET /tasks` through `SharedRead`. Nothing is reused once that read settles, so there is no time window. Every mutation goes through `mutate()` → `invalidatingWrite`, which invalidates before the route responds, whatever the outcome. The watcher invalidates on the raw fs event. Tests pin both, and pin that `server.ts` calls `callControlApi` only there. A new mutation path that calls `callControlApi` directly would serve a pre-mutation list after its own click. Filtered lists and the dead-letters read are separate queries and stay unshared. +- **Refresh results apply last-started-wins.** Parallel loads can resolve out of order, so `loadTasks` and `loadTaskDetail` apply a result only if no newer load has begun (`Latest`). The detail load also drops a result for a task no longer selected. +- **A watcher event already covered by a refresh is skipped.** `coveredByRefresh`: if a refresh whose list read **succeeded** started at or after `eventAt - WATCH_DEBOUNCE_MS` (the latest the reported change can have happened), it read after the write. The check runs when the 2 s UI debounce fires, not when the event arrives, so a refresh still in flight at arrival has settled by then. A failed refresh covers nothing. That removes the second full refresh a second after every button press. Changes made elsewhere still refresh. Lateness only errs toward refreshing. If you change the backend debounce, change the constant, which both sides import. +- **Parallel requests do not make the API parallel.** task-queue-mcp serves HTTP reads on worker threads, but the parse is CPU-bound Python, so two concurrent list reads take as long as two sequential ones (measured 2026-09-30: 0.215 s sequential vs 0.229 s parallel). The saving is in making fewer reads, not in overlapping them. - **`truncated` is rendered, never dropped.** `GET /tasks` returns at most 1000 records and says when it cut some off. The header and the dead-letters badge show it. If a view ever needs more than one page, report the number and the use rather than paging silently or asking for the cap to be raised. - **`ControlAction` must match the MCP's route set.** The union type in `control-api.ts`, the route regex in `server.ts`, and the MCP's custom routes are three copies of one contract. Change one, change all three. The two copies that live in *this* repo are now pinned to each other by a source-level test in `control-api.test.ts` — nothing detected the drift before, because adding an action to the union alone compiles and the failure mode is a button that 404s against the plugin's own backend. The third copy is in another repo and still needs a human. - **The task-queue vocabulary lives in `vocabulary.ts`, once, and is gated against its owner.** This plugin does not own the queue's statuses, task types, or workflow modes — `task-queue-mcp`'s `src/tools/queue.py` does. It used to carry four partial hand-written copies (`STATUS_ORDER`, `NON_TERMINAL_STATUSES`, `DETAIL_NON_TERMINAL_STATUSES`, the `statusColor` switch); none of them learned about `routing-failed`, so for months the status most in need of an operator sorted *below* `cancelled`, rendered the same grey as `parked`, and was not offered by the status filter. `manual-then-auto` was the same omission one field over. Two mechanisms hold the line and they are different in kind: `npm run gate:vocabulary` fetches the MCP's `main` and fails on any difference, which catches "upstream changed and we did not"; and the UI maps are `Record` keyed by the vocabulary itself, so adding a status without giving it a sort position and a colour is a `tsc` error rather than a silent fallthrough to `?? 9` and `muted`. Do not weaken either — a `Record` accepts anything and covers nothing, which is precisely how this happened. @@ -114,7 +124,7 @@ npm test npm run gate:vocabulary # parity with task-queue-mcp main — needs network, by design ``` -Tests cover `queue-token.ts` (real files in a tmpdir: missing, empty, directory, every group/other mode bit refused, `0600` and `0400` accepted, and no error carrying the file's content), `control-api.ts` (the token gate, the `X-Task-Queue-Token` header and the absence of `Authorization` and the retired secret header, the read helper and its query builder, the manifest's `env:` grants, task-id validation, header and body shape per action, transport-failure mapping, pass-through of the MCP's authorization rejections, and the union/route-regex drift gate), `dead-letters.ts` (the real dispatcher record shape including the `Date`-valued `created` js-yaml hands back, the id-less and reason-less fallbacks, and the grouping rule against the live seventeen-identical-reasons case), `ws-guard.ts` (all three upgrade cases, including the loopback-with-no-Origin one that v0.4.0 broke), `launch-policy.ts` (every closed-set rejection, whole-document rejection, and both argv shapes), `path-guard.ts` (a **real** symlink escape, traversal, the `/comms-other` sibling case — with real files in a tmpdir, because a mocked `fs` cannot demonstrate that realpath runs first), `launch-log.ts` (round-trip against the real `launchLogName`, the live bare-UUID orphans, path-shaped route ids, fence extraction including the dropped unterminated case, and the birthtime-after-mtime fallback), `vocabulary.ts` (every status has a sort position and a colour, `routing-failed` sorts above `in-progress` and is not muted, the derived non-terminal set, and the `manual-then-auto` pass-through through `buildLaunchArgv`), `gates/python-sets.ts` (the real `queue.py` shape including interleaved comments, and the two constructs that must NOT parse as string sets — a derived set and an annotated dict), and the reconnect schedule. The UI panels are not otherwise unit-tested — verify them in CloudCLI after `./deploy.sh && pm2 restart cloudcli`. +Tests cover `queue-token.ts` (real files in a tmpdir: missing, empty, directory, every group/other mode bit refused, `0600` and `0400` accepted, and no error carrying the file's content), `control-api.ts` (the token gate, the `X-Task-Queue-Token` header and the absence of `Authorization` and the retired secret header, the read helper and its query builder, the manifest's `env:` grants, task-id validation, header and body shape per action, transport-failure mapping, pass-through of the MCP's authorization rejections, and the union/route-regex drift gate), `dead-letters.ts` (the real dispatcher record shape including the `Date`-valued `created` js-yaml hands back, the id-less and reason-less fallbacks, and the grouping rule against the live seventeen-identical-reasons case), `ws-guard.ts` (all three upgrade cases, including the loopback-with-no-Origin one that v0.4.0 broke), `launch-policy.ts` (every closed-set rejection, whole-document rejection, and both argv shapes), `path-guard.ts` (a **real** symlink escape, traversal, the `/comms-other` sibling case — with real files in a tmpdir, because a mocked `fs` cannot demonstrate that realpath runs first), `launch-log.ts` (round-trip against the real `launchLogName`, the live bare-UUID orphans, path-shaped route ids, fence extraction including the dropped unterminated case, and the birthtime-after-mtime fallback), `vocabulary.ts` (every status has a sort position and a colour, `routing-failed` sorts above `in-progress` and is not muted, the derived non-terminal set, and the `manual-then-auto` pass-through through `buildLaunchArgv`), `gates/python-sets.ts` (the real `queue.py` shape including interleaved comments, and the two constructs that must NOT parse as string sets — a derived set and an annotated dict), `shared-read.ts` (one upstream GET for concurrent callers against an injected fetch, no reuse after settle, a post-mutation list that is fresh while a pre-mutation read is still in flight, invalidation on refused and failed writes, and failures never cached), source pins that `server.ts` routes every mutation through `mutate()` and invalidates in the watcher, `refresh-rules.ts` (a stale out-of-order response losing to a newer one, and the watcher-skip rule's own-action, elsewhere, mid-refresh and late-delivery cases), and the reconnect schedule. The UI panels are not otherwise unit-tested — verify them in CloudCLI after `./deploy.sh && pm2 restart cloudcli`. Two build/test gotchas worth knowing before you touch either script: diff --git a/CHANGELOG.md b/CHANGELOG.md index dde996b..5edd608 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,48 @@ All notable changes to this project will be documented in this file. +## [0.12.0] - 2026-09-30 + +Fewer, parallel reads per refresh. Programme `task-queue-read-perf-2026-09` part 3; +vikunja#1003. **Includes 0.11.1's security fix** (redirects never followed on +token-bearing requests), so deploying 0.12.0 supersedes the pending 0.11.1 deploy. +No change to authentication, the token file or the manifest's grants. + +### Changed + +- **One upstream `GET /tasks` per refresh, down from two.** `/tasks` (unfiltered) and + `/headless-runs` used to make the same list read each. They now join one in-flight read + (`shared-read.ts`). Only the read is shared: nothing is served once it settles. Every + mutation route and a Start's history write invalidate it before responding, and so + does the queue watcher, so a list requested after a mutation is always read after it. + Filtered lists and the dead-letters read are not shared. +- **The tab's three reads run in parallel**, and after a button press the list and the + task detail are re-read in parallel. A failed runs or dead-letters read still never + blanks the task list. Results apply last-started-wins, so a slow older refresh cannot + overwrite a newer one. +- **No second refresh after a button press.** The watcher's `tasks` event for the tab's + own write arrives about a second later and used to trigger another full refresh. It is + now skipped when a refresh already started after the change it reports. Changes from + agents and the dispatcher still refresh live. +- **A second click on the same task is ignored while its action is on the wire**, Start + included (a double click used to be able to launch two sessions). Mutations are never + sent in parallel or twice. task-queue-mcp's accepted unpark race (F-01) depends on this. + +### Measured + +Built backend run locally against forge's live task-queue-mcp v0.13.0, via a counting +proxy (park stubbed at the proxy, so the queue was not written). Median of 7: + +| | 0.11.1 | 0.12.0 | +|---|---|---| +| Tab load | 0.372 s, 3 upstream reads | 0.250 s, 2 upstream reads | +| Park + refresh | 0.382 s, 4 upstream reads | 0.268 s, 3 upstream reads | +| Watcher refresh after that Park | 3 more reads | skipped (by rule and unit tests; seen in CloudCLI only after deploy) | + +The saving comes from the dropped read, not from overlap. task-queue-mcp's list parse is +CPU-bound, so two concurrent reads take as long as two sequential ones (0.229 s vs +0.215 s). + ## [0.11.1] - 2026-09-30 Security fix from the `operator-panel-2026-09-p2-queue-read-api` audit (F-01, Medium; diff --git a/README.md b/README.md index 33cb23e..d2060c0 100644 --- a/README.md +++ b/README.md @@ -32,7 +32,7 @@ Routing writes through `task-queue-mcp` means mutations inherit its transition v - **UI** (`dist/index.js`) — renders the tab panel: a filterable task list and a detail view with history timeline, amendments, and context-ref previews. - **Backend** (`dist/server.js`) — HTTP + WebSocket server launched by CloudCLI. Picks a free ephemeral port at startup and reports it to CloudCLI as JSON on stdout. The UI reaches it through CloudCLI's plugin RPC API (`api.rpc()`). -Live updates arrive over WebSocket: the backend watches the queue directory and pushes a `tasks` event when files change; the UI debounces, then re-reads through the API. +Live updates arrive over WebSocket: the backend watches the queue directory and pushes a `tasks` event when files change; the UI debounces, then re-reads through the API. A refresh makes its three reads (tasks, headless runs, dead letters) in parallel, and the backend answers the first two from one upstream `GET /tasks`. After a button press the UI refreshes once, and skips the watcher's event for the same write (since v0.12.0). If the API truncated a read (it returns at most 1000 records per call), the header says **truncated: showing N of M** and the dead-letters badge says the count may be low. It is never hidden. diff --git a/manifest.json b/manifest.json index 783ca6f..0b4dbd3 100644 --- a/manifest.json +++ b/manifest.json @@ -1,7 +1,7 @@ { "name": "task-queue", "displayName": "Task Queue", - "version": "0.11.1", + "version": "0.12.0", "description": "Task queue dashboard \u2014 view, filter, and launch agent tasks.", "author": "TadMSTR", "icon": "icon.svg", diff --git a/package-lock.json b/package-lock.json index 3450310..74b0ced 100644 --- a/package-lock.json +++ b/package-lock.json @@ -1,12 +1,12 @@ { "name": "cloudcli-plugin-task-queue", - "version": "0.11.1", + "version": "0.12.0", "lockfileVersion": 3, "requires": true, "packages": { "": { "name": "cloudcli-plugin-task-queue", - "version": "0.11.1", + "version": "0.12.0", "dependencies": { "js-yaml": "^4.3.2", "ws": "^8.20.0" diff --git a/package.json b/package.json index 3dee279..bdd98e7 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "cloudcli-plugin-task-queue", - "version": "0.11.1", + "version": "0.12.0", "private": true, "type": "module", "scripts": { diff --git a/src/index.ts b/src/index.ts index 735cbc6..8389dfc 100644 --- a/src/index.ts +++ b/src/index.ts @@ -5,6 +5,7 @@ import { renderTaskDetail } from './panels/task-detail.ts'; import { renderHeadlessRuns } from './panels/headless-runs.ts'; import { renderDeadLetters } from './panels/dead-letters.ts'; import { createWsClient, WsClient } from './panels/ws-client.ts'; +import { Latest, coveredByRefresh } from './refresh-rules.ts'; // ── State ────────────────────────────────────────────────────────────── @@ -82,64 +83,74 @@ export function mount(container: HTMLElement, api: PluginAPI): void { }); container.appendChild(root); - // Debounce WS-triggered refreshes + // Debounce WS-triggered refreshes. `eventAt` is when the latest watcher event arrived. + // The coverage check runs when the timer fires rather than on arrival, so a refresh + // still in flight when the event came in has had time to succeed or fail. let refreshTimer: ReturnType | null = null; - function debouncedRefresh(delayMs = 2000): void { + function debouncedRefresh(eventAt: number, delayMs = 2000): void { if (refreshTimer) clearTimeout(refreshTimer); refreshTimer = setTimeout(() => { refreshTimer = null; + // Skipped when a refresh already read the list after the change this event reports, + // usually this tab's own button a second earlier. See coveredByRefresh. + if (coveredByRefresh(lastGoodLoadStartedAt, eventAt)) return; loadTasks(); }, delayMs); } // ── Data loading ────────────────────────────────────────────────── + // A refresh applies its results only if no newer one has started since: with the three + // reads in parallel, an older refresh can resolve after a newer one, and must not put the + // older list back on screen. + const listLoads = new Latest(); + // When the latest loadTasks() whose list read SUCCEEDED started, on the monotonic clock. + // A failed refresh covers nothing: its error stays on screen until something re-reads. + let lastGoodLoadStartedAt = -Infinity; + async function loadTasks(): Promise { - try { - const res = await api.rpc('GET', 'tasks') as { tasks: Task[]; count?: number; truncated?: boolean }; - state.tasks = res.tasks ?? []; - state.taskCount = res.count ?? state.tasks.length; - state.tasksTruncated = res.truncated === true; + const isLatest = listLoads.begin(); + const startedAt = performance.now(); + // Reads only, in parallel. The backend answers /tasks and /headless-runs from ONE + // upstream GET /tasks when they arrive together (see shared-read.ts). + const [tasks, runs, dead] = await Promise.allSettled([ + api.rpc('GET', 'tasks') as Promise<{ tasks: Task[]; count?: number; truncated?: boolean }>, + // Loaded UNFILTERED. The agent filter is applied for display in the panel, so it + // costs no round trip and — more importantly — runForTask() can still find a task's + // run while the section is filtered to a different agent. The route's ?agent= + // parameter still exists for direct API use. + api.rpc('GET', 'headless-runs') as Promise<{ runs: HeadlessRun[] }>, + // Loaded on every refresh, not only when the section is expanded: the collapsed + // heading shows the count, and a count that only appears after the operator opens + // the section is a count nobody sees. That is the failure mode this whole surface + // exists to end. + api.rpc('GET', 'dead-letters') as Promise<{ deadLetters: DeadLetter[]; truncated?: boolean }>, + ]); + if (!isLatest()) return; + + if (tasks.status === 'fulfilled') { + state.tasks = tasks.value.tasks ?? []; + state.taskCount = tasks.value.count ?? state.tasks.length; + state.tasksTruncated = tasks.value.truncated === true; state.error = null; - } catch (err) { - state.error = (err as Error).message; - } - await loadHeadlessRuns(); - await loadDeadLetters(); - state.loading = false; - render(api.context); - } - - // Runs refresh on the same cadence as the task list rather than on their own timer. - // There is nothing to stream: `claude -p` writes its final message on exit, so a run - // only ever changes once, and this surface is opened deliberately. - async function loadHeadlessRuns(): Promise { - // Loaded UNFILTERED. The agent filter is applied for display in the panel, so it - // costs no round trip and — more importantly — runForTask() can still find a task's - // run while the section is filtered to a different agent. The route's ?agent= - // parameter still exists for direct API use. - try { - const res = await api.rpc('GET', 'headless-runs') as { runs: HeadlessRun[] }; - state.headlessRuns = res.runs ?? []; - } catch { - // A failure here must not blank the task list — the runs section is secondary. - state.headlessRuns = []; + lastGoodLoadStartedAt = Math.max(lastGoodLoadStartedAt, startedAt); + } else { + state.error = (tasks.reason as Error).message; } - } - - // Loaded on every refresh, not only when the section is expanded: the collapsed heading - // shows the count, and a count that only appears after the operator opens the section is - // a count nobody sees. That is the failure mode this whole surface exists to end. - async function loadDeadLetters(): Promise { - try { - const res = await api.rpc('GET', 'dead-letters') as { deadLetters: DeadLetter[]; truncated?: boolean }; - state.deadLetters = res.deadLetters ?? []; - state.deadLettersTruncated = res.truncated === true; - } catch { - // A failure here must not blank the task list — this section is secondary. + // Runs refresh on the same cadence as the task list rather than on their own timer. + // There is nothing to stream: `claude -p` writes its final message on exit, so a run + // only ever changes once, and this surface is opened deliberately. + // A failure in either secondary read must not blank the task list. + state.headlessRuns = runs.status === 'fulfilled' ? runs.value.runs ?? [] : []; + if (dead.status === 'fulfilled') { + state.deadLetters = dead.value.deadLetters ?? []; + state.deadLettersTruncated = dead.value.truncated === true; + } else { state.deadLetters = []; state.deadLettersTruncated = false; } + state.loading = false; + render(api.context); } async function loadRunDetail(id: string): Promise { @@ -174,9 +185,16 @@ export function mount(container: HTMLElement, api: PluginAPI): void { } } + // Same rule as listLoads, for the detail view. A detail read also loses to a change of + // selection: the operator has moved on, and the answer is for a task not on screen. + const detailLoads = new Latest(); + async function loadTaskDetail(taskId: string): Promise { + const isLatest = detailLoads.begin(); + const current = (): boolean => isLatest() && state.selectedTaskId === taskId; try { const res = await api.rpc('GET', `tasks/${taskId}`) as { task: Task; previews: Record }; + if (!current()) return; state.selectedTask = res.task; state.contextPreviews = new Map(); if (res.previews) { @@ -185,23 +203,56 @@ export function mount(container: HTMLElement, api: PluginAPI): void { } } } catch (err) { + if (!current()) return; state.error = (err as Error).message; } render(api.context); } - async function handleApprove(taskId: string): Promise { + // Task ids with a mutation or Start on the wire. A second click on the same task before + // the first returns is dropped, before any confirm() or prompt(). This is what keeps the + // plugin from ever sending a duplicate mutation: part 1's audit accepted an unpark race + // (F-01) on the premise that clients do not, since two concurrent unparks can both write + // a history entry. Reads run in parallel; mutations never do, per task. + const pendingActions = new Set(); + + /** + * One button press: the mutation, awaited alone, then the list and (if this task is + * open) its detail, read in parallel. The backend invalidates its shared list before the + * mutation's response, so the refresh reads the new state. + */ + async function mutateThenRefresh( + taskId: string, + request: () => Promise, + toast: string | null, + ): Promise { + if (pendingActions.has(taskId)) return; + pendingActions.add(taskId); try { - await api.rpc('POST', `tasks/${taskId}/approve`); - await loadTasks(); - if (state.selectedTaskId === taskId) await loadTaskDetail(taskId); + await request(); } catch (err) { state.error = (err as Error).message; render(api.context); + return; + } finally { + pendingActions.delete(taskId); } + state.error = null; + if (toast) showToast(toast); + await Promise.all([ + loadTasks(), + state.selectedTaskId === taskId ? loadTaskDetail(taskId) : undefined, + ]); + } + + async function handleApprove(taskId: string): Promise { + await mutateThenRefresh(taskId, () => api.rpc('POST', `tasks/${taskId}/approve`), null); } async function handleStart(taskId: string, mode: 'review' | 'auto'): Promise { + // Guarded like a mutation: a double click here would launch two sessions. + if (pendingActions.has(taskId)) return; + pendingActions.add(taskId); try { const res = await api.rpc('POST', `tasks/${taskId}/start`, { mode }) as { note?: string }; state.error = null; @@ -213,70 +264,56 @@ export function mount(container: HTMLElement, api: PluginAPI): void { } catch (err) { state.error = (err as Error).message; render(api.context); + } finally { + pendingActions.delete(taskId); } } async function handleCancel(taskId: string): Promise { + if (pendingActions.has(taskId)) return; if (!confirm('Cancel this task? It becomes a terminal record — recoverable as a record, never deleted.')) return; - try { - await api.rpc('POST', `tasks/${taskId}/cancel`, { note: 'Cancelled via CloudCLI' }); - state.error = null; - showToast('Task cancelled'); - await loadTasks(); - if (state.selectedTaskId === taskId) await loadTaskDetail(taskId); - } catch (err) { - state.error = (err as Error).message; - render(api.context); - } + await mutateThenRefresh( + taskId, + () => api.rpc('POST', `tasks/${taskId}/cancel`, { note: 'Cancelled via CloudCLI' }), + 'Task cancelled', + ); } async function handlePark(taskId: string): Promise { + if (pendingActions.has(taskId)) return; if (!confirm("Park this task? It stays in the list, marked parked, and won't be picked up until you unpark it.")) return; - try { - await api.rpc('POST', `tasks/${taskId}/park`, { note: 'Parked via CloudCLI' }); - state.error = null; - showToast('Task parked'); - // The task stays visible — no need to leave the detail view. - await loadTasks(); - if (state.selectedTaskId === taskId) await loadTaskDetail(taskId); - } catch (err) { - state.error = (err as Error).message; - render(api.context); - } + // The task stays visible — no need to leave the detail view. + await mutateThenRefresh( + taskId, + () => api.rpc('POST', `tasks/${taskId}/park`, { note: 'Parked via CloudCLI' }), + 'Task parked', + ); } async function handleUnpark(taskId: string): Promise { - try { - await api.rpc('POST', `tasks/${taskId}/unpark`, { note: 'Unparked via CloudCLI' }); - state.error = null; - showToast('Task unparked'); - await loadTasks(); - if (state.selectedTaskId === taskId) await loadTaskDetail(taskId); - } catch (err) { - state.error = (err as Error).message; - render(api.context); - } + await mutateThenRefresh( + taskId, + () => api.rpc('POST', `tasks/${taskId}/unpark`, { note: 'Unparked via CloudCLI' }), + 'Task unparked', + ); } async function handleAmend(taskId: string): Promise { + if (pendingActions.has(taskId)) return; const amendment = prompt( 'Append an amendment. The original description is never rewritten — this is added ' + 'below it.\n\nIf a task needs more than one or two amendments, cancel and re-queue instead.', ); if (!amendment || !amendment.trim()) return; - try { - await api.rpc('POST', `tasks/${taskId}/amend`, { amendment, reason: 'Amended via CloudCLI' }); - state.error = null; - showToast('Amendment appended'); - await loadTasks(); - if (state.selectedTaskId === taskId) await loadTaskDetail(taskId); - } catch (err) { - state.error = (err as Error).message; - render(api.context); - } + await mutateThenRefresh( + taskId, + () => api.rpc('POST', `tasks/${taskId}/amend`, { amendment, reason: 'Amended via CloudCLI' }), + 'Amendment appended', + ); } async function handleRequeue(taskId: string, summary: string): Promise { + if (pendingActions.has(taskId)) return; // Confirmed, because it puts work back in front of an agent. The caveat is in the // prompt rather than only in the docs: requeueing does not fix why the task was // dropped, and all seventeen of the records this shipped against would dead-letter @@ -286,32 +323,23 @@ export function mount(container: HTMLElement, api: PluginAPI): void { + 'reset. This does NOT fix why it was dropped — if the cause is still live it will ' + 'be dead-lettered again.', )) return; - try { - await api.rpc('POST', `tasks/${taskId}/requeue`, { note: 'Requeued via CloudCLI' }); - state.error = null; - showToast('Task requeued'); - await loadTasks(); - } catch (err) { - state.error = (err as Error).message; - render(api.context); - } + await mutateThenRefresh( + taskId, + () => api.rpc('POST', `tasks/${taskId}/requeue`, { note: 'Requeued via CloudCLI' }), + 'Task requeued', + ); } async function handleSetStatus(taskId: string, status: string): Promise { - try { - await api.rpc('POST', `tasks/${taskId}/status`, { + await mutateThenRefresh( + taskId, + () => api.rpc('POST', `tasks/${taskId}/status`, { status, note: 'Status changed via CloudCLI', allow_override: true, - }); - state.error = null; - showToast(`Status set to ${status}`); - await loadTasks(); - if (state.selectedTaskId === taskId) await loadTaskDetail(taskId); - } catch (err) { - state.error = (err as Error).message; - render(api.context); - } + }), + `Status set to ${status}`, + ); } function showToast(message: string): void { @@ -496,7 +524,7 @@ export function mount(container: HTMLElement, api: PluginAPI): void { state.wsConnected = next; updateConnectionBadge(); } else if (event.type === 'tasks') { - debouncedRefresh(); + debouncedRefresh(performance.now()); } }); diff --git a/src/refresh-rules.ts b/src/refresh-rules.ts new file mode 100644 index 0000000..42ae228 --- /dev/null +++ b/src/refresh-rules.ts @@ -0,0 +1,46 @@ +// The two ordering rules the Task Queue tab's refreshes depend on, as pure code so they +// have tests. index.ts calls listTasks, headless runs and dead letters in parallel, and +// parallel reads can resolve in any order; sequential ones could not. + +/** + * How long the backend's queue watcher waits after the last `.yml` change before it + * broadcasts `tasks`. server.ts's startWatcher and the UI's refresh-skip rule both import + * this, so the two cannot drift apart. + */ +export const WATCH_DEBOUNCE_MS = 1000; + +/** + * Last-started-wins. `begin()` marks a new load and returns a check that stays true only + * until the next `begin()`. A load applies its result only if its check is still true, so + * an older load that resolves late never overwrites what a newer one showed. + */ +export class Latest { + private generation = 0; + + begin(): () => boolean { + const mine = ++this.generation; + return () => mine === this.generation; + } +} + +/** + * Whether a watcher `tasks` event is already covered by a refresh that started at + * `lastLoadStartedAt`, so a second one would re-read what the first already read. + * + * The backend broadcasts `watchDebounceMs` after the LAST change in a burst, so that change + * happened no later than `eventAt - watchDebounceMs`. A refresh that started at or after + * that point read the queue after the write. The case this exists for is the tab's own + * button: the mutation returns, its refresh starts, and a second later the watcher reports + * the same write. A change made elsewhere while no refresh was running still refreshes. + * + * Timer lateness and WebSocket latency only make `eventAt` later, which makes the estimated + * change time later than the real one, so both err toward refreshing, never toward + * skipping. Both times are on the same clock (`performance.now()` in the browser). + */ +export function coveredByRefresh( + lastLoadStartedAt: number, + eventAt: number, + watchDebounceMs: number = WATCH_DEBOUNCE_MS, +): boolean { + return lastLoadStartedAt >= eventAt - watchDebounceMs; +} diff --git a/src/server.ts b/src/server.ts index ec23c96..6140405 100644 --- a/src/server.ts +++ b/src/server.ts @@ -18,6 +18,8 @@ import { import { loadToken, tokenPath, type TokenResult } from './queue-token.ts'; import { evaluateUpgrade, allowedOrigins } from './ws-guard.ts'; import { resolveAllowedPath } from './path-guard.ts'; +import { SharedRead, invalidatingWrite } from './shared-read.ts'; +import { WATCH_DEBOUNCE_MS } from './refresh-rules.ts'; import type { DeadLetter, HeadlessRun, HeadlessRunDetail } from './types.ts'; import { toDeadLetter } from './dead-letters.ts'; import { isTerminal } from './vocabulary.ts'; @@ -191,6 +193,33 @@ async function listTasks( return page; } +/** + * The unfiltered list, read once for every caller that asks while a read is in flight. + * + * `/tasks` with no filters and `/headless-runs` (through queueIndexByPrefix) both need + * exactly this read, and the UI requests them together. Filtered lists and the + * dead-letters read are different queries and are never shared. See shared-read.ts for + * the rules: nothing is reused once a read settles, and mutate() and the watcher + * invalidate it. + */ +const unfilteredList = new SharedRead(() => listTasks()); + +/** + * Every queue mutation goes through here, never through callControlApi directly. + * + * The shared list is invalidated once the mutation returns and before the route responds, + * so a list the UI asks for after its click never comes from a read that started before + * the write, whatever the write's outcome (see invalidatingWrite). + * A test pins that server.ts has no other callControlApi call site. + */ +async function mutate( + taskId: string, + action: ControlAction, + body: Record, +): ReturnType { + return invalidatingWrite(unfilteredList, () => callControlApi(taskId, action, body, apiOpts())); +} + /** GET /tasks/{id}. Null when no record has the id; throws when the API fails. */ async function getTask(taskId: string): Promise { if (!VALID_ID.test(taskId)) return null; @@ -358,11 +387,11 @@ async function recordStartInQueue( // Terminal and unknown statuses are refused by the handler anyway; not asking is // quieter than asking and logging a rejection on every Start of a closed task. if (!current || isTerminal(current)) return; - const { status, data } = await callControlApi(taskId, 'status', { + const { status, data } = await mutate(taskId, 'status', { status: current, allow_override: true, note: `Session launched from CloudCLI in ${mode} mode`, - }, apiOpts()); + }); if (status !== 200) { console.error( `[task-queue] launched ${taskId} but could not record it in the task's history: ` @@ -550,7 +579,7 @@ async function queueIndexByPrefix(): Promise void): void { // would stream a surface nobody is watching in real time. fs.watch(TASK_QUEUE_DIR, { persistent: false }, (_eventType, filename) => { if (!filename?.endsWith('.yml')) return; + // Invalidate on the raw event, not after the debounce: an agent's write has landed, + // so a read already in flight may predate it. Plugin-initiated writes are covered + // earlier, by mutate(). + unfilteredList.invalidate(); if (watchDebounce) clearTimeout(watchDebounce); watchDebounce = setTimeout(() => { watchDebounce = null; @@ -749,7 +782,7 @@ function startWatcher(broadcast: (msg: object) => void): void { // It used to carry a file count, which the UI never read, and which was one more // reader of the queue directory with its own idea of what counts as a task. broadcast({ type: 'tasks', changed: filename }); - }, 1000); + }, WATCH_DEBOUNCE_MS); }); } @@ -800,7 +833,12 @@ const server = http.createServer(async (req, res) => { if (status) filters.status = status; if (taskType) filters.task_type = taskType; - res.end(JSON.stringify(await listTasks(filters))); + // The UI's own list read has no filters, and shares its upstream read with the + // concurrent /headless-runs request. A filtered call is its own query. + const page = Object.keys(filters).length === 0 + ? await unfilteredList.get() + : await listTasks(filters); + res.end(JSON.stringify(page)); return; } @@ -916,7 +954,7 @@ const server = http.createServer(async (req, res) => { } } - const { status: apiStatus, data } = await callControlApi(mTaskId, action, body, apiOpts()); + const { status: apiStatus, data } = await mutate(mTaskId, action, body); res.statusCode = apiStatus; res.end(JSON.stringify(data)); return; diff --git a/src/shared-read.ts b/src/shared-read.ts new file mode 100644 index 0000000..83c8c4b --- /dev/null +++ b/src/shared-read.ts @@ -0,0 +1,65 @@ +// One in-flight read, shared by every caller that asks while it is running. +// +// The Task Queue tab loads `/tasks` and `/headless-runs` together, and both need the same +// unfiltered `GET /tasks` from task-queue-mcp: the first to render the list, the second to +// match launch logs to task status (queueIndexByPrefix). Before v0.12.0 each route made its +// own read, so every refresh fetched the whole active queue twice. +// +// This shares the READ, never a RESULT. A settled promise is dropped at once, so nothing is +// served from memory after the read that produced it has finished, and there is no time +// window to reason about. Two calls share a read only if the second arrives while the first +// is still on the wire, which is exactly the case the UI's parallel loads create. +// +// `invalidate()` detaches the in-flight read, so a caller arriving after it starts a fresh +// one. The backend calls it when a mutation returns and when the queue watcher fires. A +// read that began before a mutation may reflect the old state, and a request made after the +// mutation returned must not be handed it. Callers that already joined the old read keep +// it: they asked before the mutation finished, so an answer from before it is not stale to +// them. +// +// A failed read is shared with the callers already waiting on it, and with no one after. + +export class SharedRead { + private inflight: Promise | null = null; + // Not a parameter property: `npm test` runs this file under Node's type stripping, + // which rejects them. + private readonly load: () => Promise; + + constructor(load: () => Promise) { + this.load = load; + } + + get(): Promise { + if (this.inflight) return this.inflight; + const p = this.load(); + this.inflight = p; + // Clear only if this read is still the current one. After an invalidate() and a newer + // read, the older one settling must not detach the newer. + const clear = (): void => { + if (this.inflight === p) this.inflight = null; + }; + p.then(clear, clear); + return p; + } + + invalidate(): void { + this.inflight = null; + } +} + +/** + * Run a write, then invalidate `shared` before returning, whether the write succeeded, + * was refused, or failed in transport (a lost response does not prove the write did not + * land). The caller responds only after this returns, so any read requested after the + * response starts fresh. + */ +export async function invalidatingWrite( + shared: SharedRead, + write: () => Promise, +): Promise { + try { + return await write(); + } finally { + shared.invalidate(); + } +} diff --git a/src/tests/refresh-rules.test.ts b/src/tests/refresh-rules.test.ts new file mode 100644 index 0000000..c5a8d3c --- /dev/null +++ b/src/tests/refresh-rules.test.ts @@ -0,0 +1,74 @@ +import assert from 'node:assert/strict'; +import test from 'node:test'; +import fs from 'node:fs'; +import url from 'node:url'; + +import { Latest, coveredByRefresh, WATCH_DEBOUNCE_MS } from '../refresh-rules.ts'; + +test('a stale out-of-order response does not overwrite newer state', async () => { + // The shape of index.ts's loadTasks: begin, await the reads, apply only if still latest. + const loads = new Latest(); + const state = { tasks: [] as string[] }; + let releaseOld!: (v: string[]) => void; + let releaseNew!: (v: string[]) => void; + + const load = async (read: Promise) => { + const isLatest = loads.begin(); + const tasks = await read; + if (!isLatest()) return; + state.tasks = tasks; + }; + + const older = load(new Promise(r => { releaseOld = r; })); + const newer = load(new Promise(r => { releaseNew = r; })); + // The newer refresh answers first, then the older one arrives late. + releaseNew(['parked']); + await newer; + releaseOld(['approved']); + await older; + assert.deepEqual(state.tasks, ['parked']); +}); + +test('in order, each load applies', async () => { + const loads = new Latest(); + const first = loads.begin(); + assert.equal(first(), true); + const second = loads.begin(); + assert.equal(first(), false); + assert.equal(second(), true); +}); + +test("the tab's own action: a refresh started after the write covers the watcher event", () => { + // t=0 park written; t=40 mutation returns, refresh starts; the watcher broadcasts at + // t=1000 (last change + debounce) and the WS event lands at t=1010. + assert.equal(coveredByRefresh(40, 1010), true); +}); + +test('a change made elsewhere, with no refresh since, is not covered', () => { + // Last refresh 30 s ago; an agent writes at t=30000; the event lands at t=31005. + assert.equal(coveredByRefresh(0, 31005), false); +}); + +test('a write that lands DURING a refresh is not covered by it', () => { + // Refresh started at t=100; an agent wrote at t=150, after the refresh read; the event + // lands at t=1150. The refresh may predate the write, so it must not suppress this. + assert.equal(coveredByRefresh(100, 1150), false); +}); + +test('late delivery only errs toward refreshing', () => { + // Same own-action case as above, but the event is 2 s late. The estimated change time + // moves later, past the refresh start, so the rule refreshes rather than skipping. + assert.equal(coveredByRefresh(40, 3010), false); +}); + +test('the UI rule and the server watcher use the same debounce constant', () => { + const read = (rel: string) => + fs.readFileSync(url.fileURLToPath(new URL(rel, import.meta.url)), 'utf-8'); + assert.equal(WATCH_DEBOUNCE_MS, 1000); + assert.match(read('../server.ts'), /\}, WATCH_DEBOUNCE_MS\);/); + const index = read('../index.ts'); + assert.match(index, /if \(coveredByRefresh\(lastGoodLoadStartedAt, eventAt\)\) return;/); + assert.match(index, /debouncedRefresh\(performance\.now\(\)\)/); + // Only a refresh whose list read succeeded counts as covering an event. + assert.match(index, /state\.error = null;\s*lastGoodLoadStartedAt = Math\.max\(lastGoodLoadStartedAt, startedAt\);/); +}); diff --git a/src/tests/shared-read.test.ts b/src/tests/shared-read.test.ts new file mode 100644 index 0000000..ac9d440 --- /dev/null +++ b/src/tests/shared-read.test.ts @@ -0,0 +1,189 @@ +import assert from 'node:assert/strict'; +import test from 'node:test'; +import fs from 'node:fs'; +import url from 'node:url'; + +import { SharedRead, invalidatingWrite } from '../shared-read.ts'; +import { callControlApi, queueGet, tasksQuery, LIST_PAGE_MAX, type ControlApiOptions } from '../control-api.ts'; + +/** A promise you resolve or reject from outside, to hold a read on the wire. */ +function deferred(): { promise: Promise; resolve: (v: T) => void; reject: (e: Error) => void } { + let resolve!: (v: T) => void; + let reject!: (e: Error) => void; + const promise = new Promise((res, rej) => { resolve = res; reject = rej; }); + return { promise, resolve, reject }; +} + +/** + * A stand-in task-queue API: one task whose status a POST .../park changes, and a count of + * GET /tasks requests. Each GET is held until the test releases it, so "in flight" is a + * state the test controls rather than a timing it hopes for. + */ +function fakeApi() { + let status = 'approved'; + const gets: Array<{ release: () => void }> = []; + const fetchImpl = (async (input: string | URL | Request, init?: RequestInit) => { + const u = String(input); + if (init?.method === 'POST' && u.endsWith('/park')) { + status = 'parked'; + return new Response(JSON.stringify({ ok: true }), { status: 200 }); + } + // Snapshot at the moment the request arrives: that is what a real read returns. + const body = JSON.stringify({ tasks: [{ id: 'abc12345-0000', status }], count: 1, truncated: false }); + const held = deferred(); + gets.push({ release: () => held.resolve() }); + await held.promise; + return new Response(body, { status: 200 }); + }) as typeof fetch; + const opts: ControlApiOptions = { + apiBase: 'http://127.0.0.1:8485', + token: { ok: true, token: 't' }, + fetchImpl, + }; + const read = async () => { + const r = await queueGet(tasksQuery({ limit: LIST_PAGE_MAX }), opts); + return (r.data as { tasks: Array<{ status: string }> }).tasks[0].status; + }; + return { opts, gets, read }; +} + +/** Let queued microtasks (the fetch mock reaching its hold) run. */ +const tick = () => new Promise(r => setImmediate(r)); + +test('concurrent callers share ONE upstream GET /tasks', async () => { + const api = fakeApi(); + const shared = new SharedRead(api.read); + // What the UI's parallel /tasks + /headless-runs do to the backend. + const a = shared.get(); + const b = shared.get(); + await tick(); + assert.equal(api.gets.length, 1); + api.gets[0].release(); + assert.deepEqual(await Promise.all([a, b]), ['approved', 'approved']); +}); + +test('nothing is reused once a read settles: the next caller reads again', async () => { + const api = fakeApi(); + const shared = new SharedRead(api.read); + const a = shared.get(); + await tick(); + api.gets[0].release(); + await a; + const b = shared.get(); + await tick(); + assert.equal(api.gets.length, 2, 'a settled read must not be served from memory'); + api.gets[1].release(); + await b; +}); + +test('a mutation followed by a list read returns the post-mutation state', async () => { + // The hard constraint: a read begun BEFORE the park is still on the wire when the park + // returns. A list requested after the park must not join it. + const api = fakeApi(); + const shared = new SharedRead(api.read); + const before = shared.get(); + await tick(); + + const res = await invalidatingWrite(shared, () => + callControlApi('abc12345-0000', 'park', {}, api.opts)); + assert.equal(res.status, 200); + + const after = shared.get(); + await tick(); + assert.equal(api.gets.length, 2, 'the post-mutation read must be a new upstream GET'); + api.gets[0].release(); + api.gets[1].release(); + // The early caller asked before the park returned; an answer from before it is its answer. + assert.equal(await before, 'approved'); + assert.equal(await after, 'parked'); +}); + +test('a refused or failed write still invalidates', async () => { + for (const failure of ['refused', 'transport'] as const) { + const api = fakeApi(); + const shared = new SharedRead(api.read); + shared.get(); + await tick(); + const write = failure === 'refused' + ? () => Promise.resolve({ status: 409, data: { ok: false } }) + : () => Promise.reject(new Error('socket hang up')); + await invalidatingWrite(shared, write).catch(() => undefined); + shared.get(); + await tick(); + assert.equal(api.gets.length, 2, `${failure}: a lost response does not prove the write did not land`); + for (const g of api.gets) g.release(); + } +}); + +test('an older read settling after invalidate() does not detach the newer one', async () => { + const api = fakeApi(); + const shared = new SharedRead(api.read); + const old = shared.get(); + await tick(); + shared.invalidate(); + const fresh = shared.get(); + await tick(); + api.gets[0].release(); + await old; + // The old read's cleanup ran. A caller now must still join `fresh`, not start a third. + const joiner = shared.get(); + await tick(); + assert.equal(api.gets.length, 2); + api.gets[1].release(); + assert.equal(await joiner, await fresh); +}); + +test('a failed read is shared only with callers already waiting, never cached', async () => { + let calls = 0; + const pending: Array>> = []; + const shared = new SharedRead(() => { + calls++; + const d = deferred(); + pending.push(d); + return d.promise; + }); + const a = shared.get(); + const b = shared.get(); + pending[0].reject(new Error('502')); + await assert.rejects(a, /502/); + await assert.rejects(b, /502/); + assert.equal(calls, 1); + const c = shared.get(); + assert.equal(calls, 2, 'the next caller after a failure retries'); + pending[1].resolve('ok'); + assert.equal(await c, 'ok'); +}); + +// ── Source-level pins on server.ts ────────────────────────────────────── +// server.ts calls listen() at import time, so its routes cannot be imported into a test. +// These pin the wiring the behaviour tests above rely on. The failure they catch is a new +// mutation path added with a direct callControlApi call, which compiles and quietly serves +// a pre-mutation list after it. + +const serverSrc = fs.readFileSync(url.fileURLToPath(new URL('../server.ts', import.meta.url)), 'utf-8'); +// Comments name these functions too; only code counts. +const serverCode = serverSrc.replace(/\/\*[\s\S]*?\*\//g, '').replace(/\/\/.*$/gm, ''); + +test('server.ts has exactly one callControlApi call, inside invalidatingWrite', () => { + const calls = serverCode.match(/callControlApi\(/g) ?? []; + assert.equal(calls.length, 1, 'every mutation must go through mutate()'); + assert.match(serverCode, /invalidatingWrite\(unfilteredList, \(\) => callControlApi\(/); +}); + +test('both mutation paths (the action routes and a Start) go through mutate()', () => { + assert.match(serverCode, /await mutate\(mTaskId, action, body\)/); + assert.match(serverCode, /await mutate\(taskId, 'status',/); +}); + +test('the watcher invalidates the shared list on the raw event, before its debounce', () => { + const watcher = serverCode.slice(serverCode.indexOf('function startWatcher')); + const inv = watcher.indexOf('unfilteredList.invalidate()'); + const debounce = watcher.indexOf('setTimeout('); + assert.ok(inv > 0 && inv < debounce); +}); + +test('only the unfiltered list is shared; filtered and dead-letter reads are their own', () => { + assert.match(serverCode, /new SharedRead\(\(\) => listTasks\(\)\)/); + assert.match(serverCode, /Object\.keys\(filters\)\.length === 0\s*\?\s*await unfilteredList\.get\(\)\s*:\s*await listTasks\(filters\)/); + assert.match(serverCode, /listTasks\(\{ status: 'failed', include_dead_letters: true \}\)/); +});