Files
turnstone/tests/test_interactive_pane_js.py
T
Patrick Buckley 7daf3b782d fix(#900): stream generation closes the reconnect-inside-the-await render; jitter both retries
Round-2 review. The render-time cursor-safety gate was point-in-time: a
transport that dropped AND finished re-establishing inside the /history
await reads back OPEN and is indistinguishable from one that never moved.
It is not — the redial re-presented the frozen cursor, the server answered
replay_ok, and the quiesce buffered that slice, so the render commits rows
the flush then repaints on top. Object identity cannot see it either,
since a native reconnect reuses the same EventSource; only a counter can.

_connectEpoch is bumped in onopen and nowhere else. Native auto-reconnect
calls neither connectSSE nor disconnectSSE, so those two are blind to the
exact case this exists for; connectSSE would also false-bump on its
document.hidden early return, which establishes no stream; and a closed
source can never fire a late open. Captured at dispatch, required
unchanged before a seedless render commits.

This is original-strata residual, not a regression this branch introduced:
before #900 the render was ungated entirely. The branch closed the
fire-time half and the still-down cases; these are the drop-and-recover
ones that were always open.

Two rulings written in at the gate, since neither is closed: a
fresh/truncated reconnect inside the await declines a render that would
have been safe (one wasted /history, self-healing via the flushed
synthetic state_change), and a refetch dispatched between onopen and the
replay slice arriving still renders past the frozen cursor — replay_ok
emits no end-of-replay marker, so no client-side signal exists (#903).
Coord's half of the same gate is #904; its exposure is a race rather than
this determinism, so it is not ported blind.

Also corrected: the claim that the idle-edge backstop's stream is live by
construction. It isn't — handleEvent also runs from the quiesce flush, so
a queued idle edge reaches the backstop with the transport down. The
render-time gate is what covers it. The clear_ui retry gains additive
jitter in BOTH clients from one shared constant: a declined render now
leaves the latch set, so a successful fetch can arm the retry, and the
decline trigger is herd-shaped. Kept small deliberately — the spread works
against #884's single-flight, which coalesces a lockstep herd.

test_coordinator_page.py anchored the fire guard on a literal `}, 2000);`
and on exact indentation; both would have ERRORED rather than failed once
the delay became an expression.
2026-07-24 18:15:03 -07:00

1298 lines
68 KiB
Python

"""Static smoke guards for the shared interactive pane module.
``turnstone/shared_static/interactive.js`` is the per-workstream conversational
``Pane`` lifted out of ``ui/static/app.js`` (L-shell step 5a) so BOTH the
standalone ``turnstone-server`` UI and the console L-shell can mount it. The
load-bearing invariants of that extraction are pinned here — like the rest of
the WebUI, the module has no JS test framework, so these are Python-side
string-presence assertions that catch the silent one-line regression.
"""
from __future__ import annotations
import re
from pathlib import Path
import pytest
_ROOT = Path(__file__).resolve().parent.parent
_INTERACTIVE = _ROOT / "turnstone/shared_static/interactive.js"
_COMPOSER = _ROOT / "turnstone/shared_static/composer.js"
_AUTH = _ROOT / "turnstone/shared_static/auth.js"
_APP = _ROOT / "turnstone/ui/static/app.js"
_UI_INDEX = _ROOT / "turnstone/ui/static/index.html"
def _strip_comments(js: str) -> str:
js = re.sub(r"/\*.*?\*/", "", js, flags=re.S)
js = re.sub(r"//[^\n]*", "", js)
return js
def test_interactive_is_esm_imported_by_the_shell() -> None:
"""Real ES module: it ``export``s the factory the shell imports in BOTH
deployments. Step 6 retired the window bridge (no window.InteractivePane)
and the standalone HTML no longer script-tags interactive.js — shell.js
pulls it via ``import``."""
body = _INTERACTIVE.read_text(encoding="utf-8")
assert "export { Pane as InteractivePane, createInteractivePane };" in body
assert "window.InteractivePane = Pane" not in body, (
"the window bridge is retired — the shell imports the factory (ESM)."
)
html = _UI_INDEX.read_text(encoding="utf-8")
assert "/shared/interactive.js" not in html, (
"the standalone HTML must NOT script-tag interactive.js — shell.js imports it."
)
def test_pane_constructor_takes_transport_and_host_seam() -> None:
"""The constructor takes the ``(wsId, opts)`` seam: a transport ``base``
(the node-proxy prefix) and a ``host`` adapter for the few things only the
surrounding shell knows. The old ``embedded`` flag is gone — every pane is
L-shell-hosted since the step-6 fork collapse."""
body = _INTERACTIVE.read_text(encoding="utf-8")
assert "constructor(wsId, opts) {" in body
for field in (
"this._base = opts.base",
"this._host = opts.host || INTERACTIVE_DEFAULT_HOST",
):
assert field in body, f"missing constructor seam: {field!r}"
assert "opts.embedded" not in body, (
"the embedded flag is retired — every pane is L-shell-hosted."
)
def test_transport_urls_are_base_prefixed() -> None:
"""Every per-ws request is prefixed with ``this._base`` so a console pane
proxies through ``/node/{id}`` (the LOCALITY invariant: an interactive
session lives on a cluster node). A bare ``/v1/api/workstreams/`` URL would
hit the console instead of the node and silently 404 / cross-talk."""
body = _strip_comments(_INTERACTIVE.read_text(encoding="utf-8"))
# Collapse whitespace so a prettier line-wrap (``this._base +`` on the line
# ABOVE the URL string) doesn't read as a bare URL: every workstream URL
# must be preceded by ``this._base +``.
collapsed = re.sub(r"\s+", " ", body)
bad = re.findall(r'(?<!this\._base \+ )"/v1/api/workstreams/"', collapsed)
assert not bad, (
"found a bare '/v1/api/workstreams/' URL not prefixed by this._base — "
"a console pane would route it to the console, not the owning node."
)
# The EventSource + history + send all go through the base.
assert 'this._base + "/v1/api/workstreams/"' in collapsed
assert "new EventSource(evtUrl" in body
def test_split_pane_chrome_is_retired() -> None:
"""The standalone split-pane chrome is GONE, not gated: no pane header
(name / persona / state live in the tab + rail; the conversation owns the
full pane height), no focus tracking, no split/close buttons. The dead
``!this._embedded`` branches referenced shell globals (setFocusedPane,
splitPane, splitRoot…) that no longer exist anywhere — reaching them was a
guaranteed ReferenceError, so their removal is a bugfix too."""
body = _INTERACTIVE.read_text(encoding="utf-8")
assert "_embedded" not in body, "the embedded gate is retired (always-on)"
assert 'className = "pane pane--embedded"' in body, (
"the pane root must carry pane--embedded unconditionally — "
"interactive.css scopes the slim-chrome layout to it"
)
for gone in (
"setFocusedPane",
"showPaneContextMenu",
"splitPane(",
"splitRoot",
"this.headerEl",
'"pane-header"',
'"pane-action-btn"',
"updateWsName",
):
assert gone not in body, f"retired split-pane symbol {gone!r} resurfaced"
# The persona tag stays gone (the rail's INT/COORD vocabulary shows it).
assert '"pane-persona-tag"' not in body
assert '"INTERACTIVE"' not in body
def test_factory_returns_lifecycle_over_node_proxy() -> None:
"""``createInteractivePane`` is the console factory (mirrors
``createCoordinatorPane``): it derives the node-proxy base from ``nodeId``
and returns the lifecycle controller the shell drives."""
body = _INTERACTIVE.read_text(encoding="utf-8")
assert "function createInteractivePane(root, wsId, opts) {" in body
assert '"/node/" + encodeURIComponent(opts.nodeId)' in body
for hook in ("connect()", "deactivate()", "onLogin()", "destroy()"):
assert hook in body, f"factory controller missing lifecycle hook {hook!r}"
# Teardown must close the stream so a backgrounded pane can't leak an
# upstream node connection.
assert "pane.disconnectSSE();" in body
def test_host_seam_routes_shell_couplings() -> None:
"""Every coupling to the surrounding shell goes through ``this._host`` — so
the same Pane works standalone (real adapter) and console-embedded (no-op /
Tier-1 adapter). No direct ``focusedPaneId`` / ``workstreams`` / consent
badge reference survives in the module."""
body = _INTERACTIVE.read_text(encoding="utf-8")
for call in (
"this._host.isFocused(this)",
"this._host.onStreamError(this)",
"this._host.warningTarget(this)",
"this._host.onConsentDetected(",
):
assert call in body, f"missing host seam call {call!r}"
assert "getWsName" not in body, (
"getWsName left the host seam with the pane header — the tab + rail "
"own the workstream name now."
)
code = _strip_comments(body)
# The classic split-pane shell globals must not leak into the module as
# bare code references (URL path strings excepted, handled above).
assert not re.search(r"(?<![\w$./\"])focusedPaneId(?![\w$])", code), (
"focusedPaneId leaked into the shared module — route it through "
"host.isFocused so the console (which has no such global) still works."
)
assert "_pendingConsentServers" not in code, (
"the consent-badge state must stay in the standalone shell; the pane "
"only notifies via host.onConsentDetected."
)
def test_standalone_opens_sessions_via_the_shell_pane_manager() -> None:
"""Step 6 retired the standalone's local split-pane construction: app.js no
longer builds panes via window.InteractivePane / STANDALONE_HOST. Sessions
open through the shared shell's PaneManager — openSessionPane delegates to
openPane('interactive', wsId)."""
app = _APP.read_text(encoding="utf-8")
assert "STANDALONE_HOST" not in app, "the standalone host adapter is retired."
assert "new window.InteractivePane(" not in app, (
"the standalone no longer constructs panes locally."
)
start = app.index("function openSessionPane(wsId)")
fn = app[start : start + 300]
assert 'openPane("interactive", wsId)' in fn, (
"openSessionPane must open the session as a pane via the shell PaneManager."
)
def test_approval_keyboard_shortcuts_wired() -> None:
"""The converged card advertises y/n/a (+Enter/Esc) kbd hints, so the pane
must route those keys to resolveApproval when a pending approval is up —
pane-owned on this.el (the fork collapse retired the old app.js global
handler + getFocusedPane). Guards against the chips over-promising."""
body = _INTERACTIVE.read_text(encoding="utf-8")
assert "if (!this.pendingApproval || !this.approvalBlockEl) return;" in body, (
"approval keydown must early-return unless a pending approval is up"
)
assert "e.key.toLowerCase()" in body, "the y/n/a shortcut branch"
assert ".conv-feedback" in body, (
"the feedback field uses the converged .conv-feedback, not the retired "
".ts-approval-feedback"
)
def test_media_playback_lifted_and_pane_owned() -> None:
"""The media Play affordance is rendered by the pane (buildPlayButton /
buildMediaEmbed), so its activation must live in the pane too — the old
standalone wired a DOCUMENT-level click/keydown listener in app.js, which
the console host never loaded (so the button was dead in console-hosted
panes). The fix mirrors the approval-keydown pattern: a pane-owned listener
on this.el, root-scoped via closest(".media-play-btn"). Pin both the
lifted helpers and the pane wiring so the document-level regression can't
silently come back."""
body = _INTERACTIVE.read_text(encoding="utf-8")
# The lifted activation machinery now lives in the shared module.
for fn in (
"function _loadHls(",
"function _isHlsUrl(",
"function _activatePlayer(",
"function activateMediaPlayButton(",
):
assert fn in body, f"media player helper must be lifted into the pane: {fn}"
# The HLS vendor is fetched by absolute /shared/ URL (resolves in BOTH the
# standalone server and the console, where /shared is mounted at the root).
assert 'script.src = "/shared/hls-1.6.16/hls.min.js";' in body
# Pane-owned + root-scoped — NOT a document-level delegated listener.
assert 'this.el.addEventListener("click"' in body, (
"media play must be wired on this.el (pane-owned), not document"
)
assert 'e.target.closest(".media-play-btn")' in body, (
"the play handler must be root-scoped via closest, not a document-wide id"
)
assert "activateMediaPlayButton(btn)" in body
collapsed = _strip_comments(body)
assert 'document.addEventListener("click"' not in collapsed, (
"the pane must not register a document-level click delegate — that is "
"the standalone regression that left console panes dead"
)
def test_controller_terminal_dead_state() -> None:
"""Lifecycle round 2: the console controller must STOP reconnect-polling a
session that is gone (closed / evicted / node restarted) — three consecutive
CLOSED recovery beats → give up: stream closed, status bar terminal,
``opts.onDead()`` fired once. A successful stream open resets the counter
(the new host.onStreamOpen seam). ``isDead()`` / ``markDead()`` / ``base``
are the shell's revive surface; a dead controller also ignores the login
re-arm (recovery may need a DIFFERENT node — the shell's revive owns it)."""
body = _INTERACTIVE.read_text(encoding="utf-8")
# The give-up ladder.
assert "let dead = false;" in body and "let failCount = 0;" in body
assert "const giveUp = function () {" in body
assert "failCount += 1;" in body and "if (failCount >= 3) giveUp();" in body
assert 'pane._sbTokens.textContent = "Disconnected"' in body, (
"the terminal state must be worded distinctly from the transient Reconnecting…"
)
assert "opts.onDead" in body, "the shell must hear about the give-up"
# The reset seam: Pane.connectSSE onopen → host.onStreamOpen → failCount = 0.
assert "this._host.onStreamOpen(this)" in body
assert "onStreamOpen() {}" in body, "the default host must carry the no-op"
# The shell-facing surface.
assert "isDead()" in body and "markDead: giveUp," in body
assert "base: base," in body, "the controller must expose its transport base"
# Dead controllers don't reconnect on re-auth.
assert "if (connected && !dead) pane._loadHistoryThenConnect(wsId);" in body
def test_stream_pipeline_is_wedge_proof() -> None:
"""Long-session hardening (perf audit P0): the SSE pipeline must not be
able to permanently wedge the pane. ``onmessage`` guards BOTH the
``JSON.parse`` and the ``handleEvent`` dispatch (an exception escaping it
doesn't close the EventSource, so an unguarded throw left the streaming
refs poisoned for the rest of the session), and ``stream_end`` resets the
segment refs BEFORE the finalize render, with a plain-text fallback —
with the old order a finalize throw skipped the clears and every later
delta painted into the dead segment."""
body = _INTERACTIVE.read_text(encoding="utf-8")
assert "dropping malformed SSE frame" in body
assert "handleEvent failed for" in body
case = body.index('case "stream_end"')
seg = body[case : body.index("break;", case)]
clears = seg.index("this.currentAssistantBodyEl = null;")
finalize = seg.index("streamingRenderFinalize(")
assert clears < finalize, (
"stream_end must clear segment refs BEFORE finalize — the old "
"finalize-first order wedged all later assistant output on a throw."
)
assert "doneBodyEl.textContent = doneBuffer;" in seg
def test_rebuild_quiesces_live_events_and_releases_agent_tracking() -> None:
"""clear_ui re-render race (perf audit P0): live SSE events painted
between the history snapshot and ``replaceChildren()`` were wiped with
no redelivery, and streaming refs kept pointing at detached nodes.
Pinned: the quiesce queue sits on the handleEvent hot path, the
live-stream re-render trigger (clear_ui) arms it, ``replayHistory``
resets the streaming refs and clears the agent-card/orphan maps (the
detached-DOM retention leak), and the mid-stream guard covers the
reasoning bubble. (replay_truncated no longer arms the quiesce — it
tears the stream down and runs the full fresh-connect flow instead;
see test_truncated_resync_is_full_fresh_connect.)"""
body = _INTERACTIVE.read_text(encoding="utf-8")
assert "this._replayQueue.events.push(evt);" in body
assert body.count("this._beginReplayQuiesce(") >= 1, (
"clear_ui must arm the quiesce (it re-renders on a LIVE stream)"
)
assert "!this.currentAssistantEl && !this.currentReasoningEl" in body
replay = body.index("replayHistory(messages) {")
seg = body[replay : replay + 2600]
for line in (
"this._resetStreamingRefs();",
"this._clearAgentTracking();",
):
assert line in seg, f"replayHistory must reset: {line!r}"
assert "this._agentCards.clear();" in body
# Review-hardened lifecycle: the card entry SURVIVES the terminal
# tool_result (a late child event finding no Map entry would rebuild a
# duplicate empty card beside the finished one), and transport-only
# reconnects preserve the maps + any armed quiesce queue — clearing them
# in disconnectSSE duplicated cards and dropped buffered orphan steps on
# every transient stream blip. Full-reload cleanup lives in
# _loadHistoryThenConnect; terminal cleanup in the factory's destroy().
assert "this._agentCards.delete(callId);" not in body
disc = body.index("disconnectSSE() {")
disc_seg = body[disc : body.index("_loadHistoryThenConnect(wsId) {", disc)]
assert "this._clearAgentTracking();" not in disc_seg
assert "this._replayQueue = null;" not in disc_seg
load = body.index("_loadHistoryThenConnect(wsId) {")
load_seg = body[load : body.index("async _refetchHistory(", load)]
assert "this._clearAgentTracking();" in load_seg
assert "this._replayQueue = null;" in load_seg
# A mid-stream replay_truncated DEFERS the re-sync (flag consumed on the
# idle edge) instead of dropping it — skipping left the lost-event gap
# unrepaired for the rest of the session.
assert "this._pendingTruncatedResync = true;" in body
# The refetch FAILURE branch is a DOM/ref no-op (#890): the full
# guard-before-wipe contract — failure branch, clear_ui, the
# resumability-gated reload reset, and the affordance gates — is
# pinned in test_interactive_refetch_failure_preserves_the_pane.
def test_interactive_refetch_failure_preserves_the_pane() -> None:
"""A FAILED /history fetch must leave the transcript, streaming refs,
and repair-intent state untouched on EVERY re-render route (#890 —
the interactive mirror of coord's #882 G3 guard-before-wipe, pinned
there by test_coordinator_refetch_failure_preserves_the_pane). The
wipe + resets live in replayHistory, reached only on success:
- the clear_ui case must NOT pre-wipe (pre-#890 a failed fetch
during a rewind blanked the highest-traffic pane on a live
stream);
- the _refetchHistory failure branch must be empty except for the
quiesce release (no showEmptyState — the stray hint below stale
content was the resync-route wart; no ref reset);
- the failed-first-paint placeholder survives ONLY via the factory
connect() pre-seed, which must stay ahead of the load call;
- _loadHistoryThenConnect must reset streaming refs UNLESS a
truncation resync is armed (the armed route's cursor replay
resumes the mid-jitter bubble; every non-resumable flavor —
ws switch, first paint, idle edge, unarmed re-auth reload —
resets, or a stale ref would concatenate the next turn into the
old bubble);
- the row-level mutating affordances (rewind / edit /
edit-and-resend) must gate on the _historyStale latch alongside
busy — the latch spans clear_ui arrival through the next
SUCCESSFUL render, covering the fetch window AND the failed-fetch
aftermath (a quiesce-based gate reopened on the failure exit and
let a second rewind over-rewind off the stale DOM);
- the latch heals TRANSPORT-FREE: one turn-free bounded retry from
the clear_ui .then, and a quiesced same-token refetch at organic
idle edges as the double-failure backstop. Neither may touch
the stream — a reload's fresh reconnect draws the server's
synthetic state_change:idle back into the backstop's own trigger
(the round-5 zero-backoff reconnect storm).
"""
body = _INTERACTIVE.read_text(encoding="utf-8")
# clear_ui: no pre-wipe before the fetch, and the quiesce must be
# armed BEFORE the fetch (queued live events land in the rebuilt —
# or, on failure, the stale-but-real — pane, never the void).
cl = body.index('case "clear_ui":')
cl_seg = body[cl : body.index("break;", cl)]
assert "this._refetchHistory(this.wsId, token)" in cl_seg
assert "this.messagesEl.replaceChildren();" not in cl_seg, (
"clear_ui must not pre-wipe the transcript (#890)"
)
assert "this._resetStreamingRefs();" not in cl_seg, (
"clear_ui must not pre-reset streaming refs (#890)"
)
assert cl_seg.index("this._beginReplayQuiesce(token);") < cl_seg.index(
"this._refetchHistory("
), "clear_ui must arm the quiesce BEFORE the fetch"
# Render-time cursor-safety gate (#900): a SEEDLESS render commits
# /history's as-of-now truth without advancing _lastEventId, so
# committing it while the transport is down strands the cursor below
# the rows just painted and the next connect's replay paints them
# again. The gate must be seedless-SCOPED — _loadHistoryThenConnect
# disconnects first, so its evtSource is null for its whole fetch and
# gating it would break every first paint, ws switch and resync.
ref = body.index("async _refetchHistory(wsId, token, seedCursor = false) {")
ref_seg = body[ref : body.index("\n _beginReplayQuiesce(token) {", ref)]
gate = ref_seg.index("const cursorSafe =")
gate_seg = ref_seg[gate : ref_seg.index(";", gate)]
assert "seedCursor ||" in gate_seg, (
"the render-time gate must exempt seeded (disconnect-first) loads (#900)"
)
assert "this.evtSource.readyState === EventSource.OPEN" in gate_seg
# Connection-generation term (#900 r2): readyState alone cannot see a
# transport that DROPPED and finished RE-ESTABLISHING inside the await —
# it reads OPEN either way, while the redial re-presented the frozen
# cursor and the quiesce buffered the replay the render would duplicate.
# Object identity is not a substitute: a NATIVE reconnect reuses the same
# EventSource, so only a counter can see it.
assert "this._connectEpoch === epoch" in gate_seg, (
"the seedless render must require an unchanged stream generation (#900)"
)
assert "if (data && cursorSafe) {" in ref_seg, (
"the seedless render must be gated on a live cursor (#900)"
)
# Captured at DISPATCH, not read live at render time — a live read would
# compare the epoch against itself and the term would be vacuous.
assert "const epoch = this._connectEpoch;" in ref_seg, (
"the gate must compare against the generation captured at dispatch (#900)"
)
assert ref_seg.index("const epoch = this._connectEpoch;") < ref_seg.index("await authFetch("), (
"the generation must be captured BEFORE the fetch await (#900)"
)
# The bump belongs to onopen and NOWHERE else. A native auto-reconnect
# calls neither connectSSE nor disconnectSSE, so those two are blind to
# the very case the term exists for; connectSSE would additionally
# FALSE-bump on its document.hidden early return, which establishes no
# stream. Absence from both is therefore load-bearing, not incidental.
assert "this._connectEpoch = 0;" in body, (
"the generation must be initialised — undefined makes every compare "
"NaN-false and silently declines every seedless render (#900)"
)
onopen = body.index("this.evtSource.onopen = () => {")
onopen_seg = body[onopen : body.index("\n };", onopen)]
assert "this._connectEpoch += 1;" in onopen_seg, (
"the stream generation must be bumped in onopen — the only site that "
"fires for a NATIVE auto-reconnect (#900)"
)
conn = body.index("connectSSE(wsId) {")
assert "_connectEpoch" not in _strip_comments(body[conn:onopen]), (
"connectSSE must NOT bump the generation — it would false-bump on the "
"document.hidden early return, which establishes no stream (#900)"
)
dis = re.search(r"\n disconnectSSE\(\) \{(.*?)\n \}\n", body, re.S)
assert dis is not None, "disconnectSSE not found"
# Comment-stripped: the method's inventory comment NAMES both fields as
# deliberately-not-cleared, which is the ruling — assert on code only.
dis_code = _strip_comments(dis.group(1))
assert "_connectEpoch" not in dis_code, (
"disconnectSSE must NOT bump the generation — a teardown is decided "
"by the presence term, and a re-establish by the next onopen (#900)"
)
# THE load-bearing negative: disconnectSSE deliberately does NOT cancel
# the clear_ui retry (transport-only redials keep the pending heal
# intent). That is exactly why the retry can fire against a dead
# transport, which is what its OPEN fire-guard term exists to handle —
# so a "symmetry" cleanup adding the cancel here would silently make
# that guard unreachable and E5's detector vacuous, with nothing else
# failing. Coord pins the same invariant (test_coordinator_page.py).
assert "_staleRetryTimer" not in dis_code, (
"disconnectSSE must NOT cancel the clear_ui failure retry — "
"transport-only reconnects keep the pending heal intent (#890/#900)"
)
# Failure branch: quiesce release only.
fail = body.index("Failed fetch = DOM + ref + repair-intent no-op")
fail_seg = body[fail : fail + 1100]
assert "this._endReplayQuiesce(token);" in fail_seg
assert "this.showEmptyState();" not in fail_seg, (
"a failed fetch must not append an empty-state hint (#890)"
)
assert "this._resetStreamingRefs();" not in fail_seg
# The failure branch must RESOLVE, never throw/reject: the clear_ui
# .then dispatches the queued edit-and-resend after a failed fetch
# too (the rewind already committed server-side) — a throw here
# would route to .catch and strand the resend.
assert "throw " not in fail_seg
assert "Promise.reject" not in fail_seg
# Success path still owns the wipe + resets.
replay = body.index("replayHistory(messages) {")
replay_seg = body[replay : replay + 2600]
assert "this.messagesEl.replaceChildren();" in replay_seg
assert "this._resetStreamingRefs();" in replay_seg
# Factory pre-seed ahead of the load — the failed-first-paint
# placeholder's only remaining producer.
conn = body.index("Load-bearing pre-seed (#890)")
conn_seg = body[conn : conn + 600]
seed = conn_seg.index("pane.showEmptyState();")
load = conn_seg.index("pane._loadHistoryThenConnect(wsId);")
assert seed < load, "connect() must seed the empty-state BEFORE the fetch"
# Resumability-gated ref reset in _loadHistoryThenConnect: refs
# survive a reload only when an armed truncation cursor lets the
# reconnect resume into them; every other flavor (ws switch,
# unarmed re-auth reload) resets.
lh = body.index("_loadHistoryThenConnect(wsId) {")
lh_seg = body[lh : body.index("async _refetchHistory(", lh)]
assert "if (this._truncatedFromCursor == null) this._resetStreamingRefs();" in lh_seg, (
"reload must reset streaming refs unless a truncation resync is armed (#890)"
)
# The mutating row affordances gate on the staleness latch alongside
# busy — scoped per method so a gate migrating off one affordance
# cannot hide inside a file-wide occurrence count.
for sig in (
"_rewindToMessage(msgEl) {",
"_startEdit(msgEl, originalText) {",
"_editAndResend(msgEl, newText) {",
):
m = body.index(sig)
assert "if (this.busy || this._historyStale) return;" in body[m : m + 400], (
f"missing busy || _historyStale gate in {sig!r} (#890)"
)
# Latch lifecycle: set at clear_ui arrival BEFORE the quiesce is
# armed (so no event can interleave between restructure-signal and
# gate-close); cleared ONLY in replayHistory's render, which also
# cancels the pending failure retry; the retry is scheduled from
# the clear_ui .then (bounded by construction — a retry's own
# failure cannot re-schedule) and fire-time-gated turn-free; the
# idle edge backstops a double failure via else-if BEHIND the
# truncated branch (whose own render heals the latch too).
assert cl_seg.index("this._historyStale = true;") < cl_seg.index(
"this._beginReplayQuiesce(token);"
), "the staleness latch must be set before the quiesce is armed"
assert "this._historyStale = false;" in replay_seg, (
"replayHistory must clear the staleness latch (its only clear site)"
)
assert "clearTimeout(this._staleRetryTimer);" in replay_seg
assert "this._staleRetryTimer = setTimeout(" in cl_seg
assert "!this.currentAssistantEl &&" in cl_seg, (
"the clear_ui failure retry must be turn-free-gated"
)
assert "} else if (this._historyStale && !this._replayQueue) {" in body, (
"the idle edge must backstop the latch behind the truncated branch, "
"skipping edges with a quiesced fetch already in flight"
)
# The backstop must be TRANSPORT-FREE: a quiesced same-token refetch,
# never _loadHistoryThenConnect — the reload's fresh reconnect draws
# the server's synthetic state_change:idle back into this branch's
# own trigger (the round-5 storm).
backstop = body.index("} else if (this._historyStale && !this._replayQueue) {")
backstop_seg = body[backstop : backstop + 2200]
assert "this._refetchHistory(this.wsId, staleToken);" in backstop_seg, (
"the staleness backstop must heal via a quiesced REST refetch"
)
assert "this._loadHistoryThenConnect(" not in backstop_seg, (
"the staleness backstop must never touch the transport (#890 r5)"
)
# The retry yields to an in-flight quiesce (no same-token stomp).
retry = cl_seg.index("this._staleRetryTimer = setTimeout(")
assert "!this._replayQueue &&" in cl_seg[retry : retry + 700], (
"the clear_ui retry must yield to an in-flight quiesced fetch"
)
# ...and must not fire against a DOWN transport (#900): disconnectSSE
# deliberately keeps the timer armed, so a hidden tab / degraded
# cooldown / native redial can hold the fire. A seedless refetch then
# paints rows the frozen _lastEventId still sits below, and the next
# connect replays that slice on top. OPEN, not merely present — a
# CONNECTING source has a frozen cursor with a replay pending.
assert "this.evtSource.readyState === EventSource.OPEN" in cl_seg[retry:], (
"the clear_ui retry must require a live stream at fire time (#900)"
)
# Terminal teardown must invalidate in-flight loads AND cancel the
# failure retry. The token bump (#900) is the chokepoint: without it
# destroy() left _loadHistoryThenConnect's .finally free to reopen an
# EventSource on the detached pane (re-registering the document-level
# visibilitychange listener destroy just removed), a settling
# _refetchHistory free to replayHistory into detached DOM, and the
# clear_ui .then free to RE-ARM the timer destroy had cancelled. The
# clearTimeout stays because the timer may already be armed and a
# timer into a destroyed pane must be dead, not merely inert.
# (disconnectSSE deliberately does NOT cancel it — transport-only
# reconnects keep the pending heal intent.)
dest = body.index("destroy() {")
dest_seg = body[dest : body.index("\n },", dest)]
assert "clearTimeout(pane._staleRetryTimer);" in dest_seg, (
"destroy() must cancel the clear_ui failure retry (#890)"
)
assert "pane._historyLoadToken = (pane._historyLoadToken || 0) + 1;" in dest_seg, (
"destroy() must bump the load token to invalidate in-flight loads (#900)"
)
# Same rule on the other terminal path: giveUp() already bumped the
# token, so the timer is inert there — but inert is not dead.
give = body.index("const giveUp = function () {")
give_seg = body[give : body.index("\n };", give)]
assert "pane._historyLoadToken = (pane._historyLoadToken || 0) + 1;" in give_seg
assert "clearTimeout(pane._staleRetryTimer);" in give_seg, (
"giveUp() must cancel the clear_ui failure retry too (#900 symmetry)"
)
def test_per_token_hot_path_avoids_container_scans() -> None:
"""P1 (perf audit): per-token work must stay O(1) in transcript length.
The thinking indicator is an instance ref (the class-selector miss walked
the whole transcript on EVERY content/reasoning delta); near-bottom state
comes from the passive scroll listener instead of a forced-layout
geometry read per event; the scroll pin is rAF-coalesced; per-tool
row/stream lookups resolve through the self-healing caches."""
body = _INTERACTIVE.read_text(encoding="utf-8")
stripped = _strip_comments(body)
assert 'querySelector(".thinking-indicator")' not in stripped, (
"thinking indicator must use the instance ref, not a container scan"
)
assert "this._thinkingEl" in body
near = body.index("isNearBottom() {")
assert "return this._nearBottom;" in body[near : near + 700]
assert "passive: true" in body
# The rAF pin re-checks the flag AT FIRE TIME (a user scroll landing in
# the schedule→rAF window must win over a stale pin), with force
# requests latched across the coalescing window; resizes re-derive the
# flag via ResizeObserver since they move the bottom without a scroll.
assert "this._scrollPinForce = false;" in body
assert "ResizeObserver" in body
for helper in ("_toolRow(callId) {", "_streamEl(callId) {"):
assert helper in body, f"missing lookup-cache helper: {helper!r}"
# -- Shared-workstream cross-user send gate -----------------------------------
#
# The UX complement to the server-side CrossUserInterjectionError (a 409): while
# another participant's turn is in flight, this viewer's send button is disabled
# so they can't interject under the initiator's credentials / be misattributed.
# The wiring spans three modules; these string-presence guards catch the silent
# one-line regression the way the rest of this file does (no JS test framework).
def test_composer_exposes_hard_send_block() -> None:
"""The composer has an independent hard-block axis, reconciled with busy,
so a caller can disable send even in queueWhileBusy (queue) mode."""
body = _COMPOSER.read_text(encoding="utf-8")
assert "Composer.prototype.setSendBlocked = function" in body
assert "Composer.prototype._reconcileDisabled = function" in body
assert "this._sendBlocked = false;" in body
# setBusy must route the disabled write through the reconciler (not clobber
# the block with a direct sendBtn.disabled assignment).
stripped = _strip_comments(body)
setbusy = stripped.index("Composer.prototype.setBusy = function")
setbusy_end = stripped.index("Composer.prototype._reconcileDisabled")
assert "this._reconcileDisabled();" in stripped[setbusy:setbusy_end]
assert "this.sendBtn.disabled =" not in stripped[setbusy:setbusy_end], (
"setBusy must not write sendBtn.disabled directly — reconcile owns it"
)
def test_auth_retains_user_id_for_gate() -> None:
"""whoami's opaque user_id is retained (separately from the display
username) so the pane can compare it against the acting-user id."""
body = _AUTH.read_text(encoding="utf-8")
assert 'sessionStorage.setItem("ts.user_id", data.user_id);' in body
assert 'sessionStorage.removeItem("ts.user_id");' in body
def test_pane_gates_send_on_cross_user_busy() -> None:
"""The pane tracks the acting user from state_change, compares it against
the viewer's own id, and blocks send while another participant is busy."""
body = _INTERACTIVE.read_text(encoding="utf-8")
assert "_reconcileSendBlock() {" in body
# tracks the acting user from the state_change event...
assert "this._actingUserId = evt.acting_user_id;" in body
assert "this._actingUserId = null;" in body # cleared when the turn settles
# ...compares against the viewer's own id from /whoami...
assert 'sessionStorage.getItem("ts.user_id")' in body
assert "this._actingUserId !== me" in body
# ...and drives the composer's hard block, re-run on every busy edge.
assert "this.composer.setSendBlocked(" in body
stripped = _strip_comments(body)
setbusy = stripped.index("setBusy(b, source) {")
assert "this._reconcileSendBlock();" in stripped[setbusy : setbusy + 800]
def test_pane_handles_cross_user_409() -> None:
"""The reactive fallback: a 409 (button not yet disabled) surfaces a clean
message, not the generic 'Connection error' catch. The pane converts
the 409 body at the fetch stage; the status ARM itself lives in the
shared settle helper (composer_queue.settleSendResponse) with the rest
of the response matrix."""
body = _INTERACTIVE.read_text(encoding="utf-8")
assert "r.status === 409" in body
assert 'status: "cross_user_interjection"' in body
helper = (_ROOT / "turnstone/shared_static/composer_queue.js").read_text(encoding="utf-8")
assert 'status === "cross_user_interjection"' in helper
def test_sync_approval_state_prunes_orphan_cycles() -> None:
"""``_syncApprovalState`` prunes cycles whose block elements are no longer
in the living DOM (``.isConnected === false``). This covers the rare case
where an ``approve_request`` event is processed between a DOM wipe
(``clear_ui`` / ``replay_truncated`` / ``replaceChildren``) and the
refetch-restore — the cycle card lives in a detached subtree, the matching
``approval_resolved`` never arrives, and the send button stays disabled
forever without this guard. The pin guards against a future refactor that
drops the orphan prune but doesn't otherwise break ``_syncApprovalState``."""
body = _INTERACTIVE.read_text(encoding="utf-8")
fn_start = body.index("_syncApprovalState() {")
assert "entry.blockEls && !entry.blockEls.some((el) => el.isConnected)" in body, (
"orphan pruning must check .isConnected on block elements"
)
tail = body[fn_start : body.index("_oldestCycleId()", fn_start)]
assert "this.approvalCycles.delete(cid);" in tail, (
"orphan pruning must delete the cycle from the Map"
)
# ---------------------------------------------------------------------------
# SSE overflow recovery + close-on-hide (fast-stream corruption fixes)
# ---------------------------------------------------------------------------
def test_stream_overflow_case_counts_and_rate_limits() -> None:
"""The server closes an overflowed stream after an id-less
``stream_overflow`` frame; the pane must count it (field
instrumentation for the drop-vs-render-wedge diagnosis) and route it
through the reconnect limiter so a persistently slow consumer trips
the degraded catch-up instead of churning reconnect/replay cycles."""
body = _INTERACTIVE.read_text(encoding="utf-8")
assert 'case "stream_overflow":' in body
assert "this._noteStreamOverflow();" in body
# The health object carries all four field-forensics counters; the
# truncated-resync counter distinguishes the replay-window class from
# the dropped-events class (overflows).
health = re.search(r"_streamHealth = \{(.*?)\};", body, re.S)
assert health is not None, "_streamHealth initializer not found"
for field in ("overflows: 0", "renderThrows: 0", "malformedFrames: 0", "truncatedGaps: 0"):
assert field in health.group(1), f"_streamHealth must init {field!r}"
# Both wedge-class catch sites increment the render-throw counter,
# and the malformed-frame drop counts too — the C-OVERDETERMINED
# instrumentation that tells drops apart from wedges in the field.
assert body.count("this._streamHealth.renderThrows += 1;") == 2
assert "this._streamHealth.malformedFrames += 1;" in body
assert "this._streamHealth.overflows += 1;" in body
def test_truncated_resync_is_full_fresh_connect_with_churn_limit() -> None:
"""replay_truncated = the stream admitted losing events past recovery.
The pane must treat the connection as DEAD: run the full fresh-connect
flow (``_loadHistoryThenConnect`` — disconnect first, REST /history,
adopt the resume cursor, reconnect) rather than an in-place refetch.
The in-place shape discarded the /history cursor while /history TRIMS
the trailing in-flight turn whenever it returns one — a mid-run
truncation wiped the executing turn (task cards included) with no
redelivery, and sub-agent children then escaped to top-level rows
after the orphan grace (the 2026-07 field reports).
Pinned: (1) the immediate branch routes through
``_noteTruncatedResync()`` and SKIPS the resync when the limiter just
tripped — the cooldown disconnected the stream, and the flow's
``.finally`` reconnect would defeat it; on no-trip it SCHEDULES the
fresh connect behind the herd-spreading jitter rather than starting
it inline (a node restart makes every stale-cursor tab resync inside
the EventSource retry window, and the per-tab limiter cannot see a
cross-tab herd); (2) the mid-stream guard still DEFERS (detachable
bubble); (3) the idle-edge consumption runs the same fresh-connect
flow, unjittered (staggered by turn-settle timing); (4) the limiter
feeds the SAME churn window as overflow closes and enters degraded
catch-up on trip; (5) neither truncated branch arms the quiesce or
calls the in-place ``_refetchHistory`` — the old shape must not come
back; (6) every ``truncatedGaps`` bump routes through the one
increment+log step (``_recordTruncatedGap``), whose class-of-event
wording never asserts a resync a trip may skip; (7) the scheduler
nulls its handle before loading and ``disconnectSSE`` cancels a
pending one, so a torn-down stream's resync can't fire against the
pane's next workstream; (8) the gap-repair guarantee is
interleaving-proof via the connect chokepoint; (9) a same-ws load
supersession must not strand a queued edit-and-resend."""
body = _INTERACTIVE.read_text(encoding="utf-8")
trunc = re.search(r'case "replay_truncated":(.*?)break;', body, re.S)
assert trunc is not None, "replay_truncated case not found"
t = trunc.group(1)
# (1) immediate branch: limiter check gates the JITTERED fresh connect.
assert "if (!this._noteTruncatedResync())" in t
assert "this._scheduleTruncatedResync();" in t
# (2) mid-stream defer unchanged.
assert "!this.currentAssistantEl && !this.currentReasoningEl" in t
assert "this._pendingTruncatedResync = true;" in t
# (5) the old in-place shape must not come back.
assert "_beginReplayQuiesce" not in t
assert "_refetchHistory" not in t
# (3) idle-edge consumption: consume the latch, record, fresh-connect.
idle = re.search(
r"if \(this\._pendingTruncatedResync\) \{(.*?)\n \}",
body,
re.S,
)
assert idle is not None, "idle-edge truncated consumption not found"
i = idle.group(1)
assert "this._pendingTruncatedResync = false;" in i
assert "this._recordTruncatedGap();" in i
assert "this._loadHistoryThenConnect(this.wsId);" in i
assert "_refetchHistory" not in i
# (4) the limiter method: records the gap, then delegates churn
# accounting to the ONE shared trip step (also used by
# _noteStreamOverflow) so the trip parameters cannot silently diverge
# between the overflow and truncated classes.
note = re.search(r"_noteTruncatedResync\(\)\s*\{(.*?)\n \}", body, re.S)
assert note is not None, "_noteTruncatedResync method not found"
n = note.group(1)
assert "this._recordTruncatedGap();" in n
assert "return this._recordChurnAndMaybeTrip();" in n
churn = re.search(r"_recordChurnAndMaybeTrip\(\)\s*\{(.*?)\n \}", body, re.S)
assert churn is not None, "_recordChurnAndMaybeTrip method not found"
c = churn.group(1)
assert "this._overflowTimes.push(now);" in c
assert "overflowWindowTripped(" in c
assert "this._enterDegradedCatchup();" in c
assert "return true;" in c
assert "return false;" in c
over = re.search(r"_noteStreamOverflow\(\)\s*\{(.*?)\n \}", body, re.S)
assert over is not None, "_noteStreamOverflow method not found"
assert "this._recordChurnAndMaybeTrip();" in over.group(1), (
"overflow closes must feed the same shared churn step"
)
# (6) the one increment+log step: every bump carries the running count,
# class-of-event wording (no action a trip may skip); the counter is
# named for what it counts — gap detections, not resyncs performed.
rec = re.search(r"_recordTruncatedGap\(\)\s*\{(.*?)\n \}", body, re.S)
assert rec is not None, "_recordTruncatedGap method not found"
r = rec.group(1)
assert "this._streamHealth.truncatedGaps += 1;" in r
assert "console.warn(" in r
assert "resyncing history from REST" not in r, (
"the shared log line must not assert a resync the trip branch skips"
)
assert body.count("this._streamHealth.truncatedGaps += 1;") == 1, (
"all truncatedGaps bumps must route through _recordTruncatedGap "
"so the running-count console invariant holds"
)
# (7) jittered scheduler: dedup guard, null-before-load, cancel on
# disconnect.
sched = re.search(r"_scheduleTruncatedResync\(\)\s*\{(.*?)\n \}", body, re.S)
assert sched is not None, "_scheduleTruncatedResync method not found"
s = sched.group(1)
assert "if (this._resyncTimer != null) return;" in s
assert "Math.random() * TRUNCATED_RESYNC_JITTER_MS" in s
null_then_load = s.index("this._resyncTimer = null;")
load = s.index("this._loadHistoryThenConnect(this.wsId);")
assert null_then_load < load, (
"the firing path must null the handle BEFORE loading, or "
"disconnectSSE (called inside the load) would cancel the work it "
"is part of"
)
dis = re.search(r"disconnectSSE\(\)\s*\{(.*?)\n \}", body, re.S)
assert dis is not None, "disconnectSSE not found"
assert "clearTimeout(this._resyncTimer)" in dis.group(1)
# (8) interleaving-proof gap repair: the truncation-time cursor is a
# single record — captured keep-oldest at the envelope, cleared by
# any full committed-history render (replayHistory) or a ws switch —
# and the CONNECT CHOKEPOINT consumes it: while set, every manual
# (re)connect presents it instead of the since-advanced live cursor,
# so the server re-answers replay_truncated and the resync re-arms
# no matter which teardown (hide/show, degraded cooldown, recover
# beat, failed /history fetch) cancelled the pending jittered timer.
# A cancellable timer alone silently lost the repair when a
# hide/show cycle landed inside the jitter window.
assert re.search(
r"if \(this\._truncatedFromCursor == null\) \{\s*"
r"this\._truncatedFromCursor = this\._lastEventId;",
t,
), "the truncated case must record the truncation-time cursor keep-oldest"
load = body.index("_loadHistoryThenConnect(wsId) {")
load_seg = body[load : body.index("async _refetchHistory(", load)]
assert "if (this.wsId !== wsId) this._truncatedFromCursor = null;" in load_seg, (
"a ws switch must drop the old ws's truncation record"
)
replay_fn = body.index("replayHistory(messages) {")
replay_head = body[replay_fn : replay_fn + 1200]
assert "this._truncatedFromCursor = null;" in replay_head, (
"a successful full-history render must clear the truncation record"
)
# (8b) ...and supersede ALL pending repair intent in the same breath —
# the deferred mid-turn latch and any pending jittered timer. Without
# these, a clear_ui heal left them armed and the next idle edge fired a
# phantom _loadHistoryThenConnect against the repaired gap (false
# truncatedGaps bump; on its failed-fetch leg, a cursorless reconnect
# with nothing armed). Every _loadHistoryThenConnect flavor
# clears both BEFORE its fetch, so these are no-ops on the load paths —
# the clear_ui heal is the path they exist for. Mirrors the
# coordinator's refetchHistory supersession.
assert "this._pendingTruncatedResync = false;" in replay_head, (
"replayHistory must clear the deferred-resync latch — a latch "
"surviving a heal fires a phantom resync at the next idle edge"
)
assert "clearTimeout(this._resyncTimer)" in replay_head, (
"replayHistory must cancel a pending jittered resync — the render "
"just repaired the gap it was scheduled for"
)
conn = body.index("connectSSE(wsId) {")
conn_seg = body[conn : body.index("this.evtSource = new EventSource", conn)]
assert "this._truncatedFromCursor != null" in conn_seg, (
"connectSSE must consult the truncation record"
)
assert "encodeURIComponent(connectCursor)" in conn_seg
# (9) same-ws supersession must not strand the queued edit-and-resend
# (the jittered resync bumps _historyLoadToken through
# _loadHistoryThenConnect while a clear_ui edit flow is in flight);
# only a ws SWITCH discards the resend, and then it recovers the
# composer instead of leaving the latch armed for the next ws.
clear = re.search(r'case "clear_ui": \{(.*?)break;', body, re.S)
assert clear is not None, "clear_ui case not found"
cl = clear.group(1)
assert "const editWs = this.wsId;" in cl
assert "if (token !== this._historyLoadToken && this.wsId !== editWs)" in cl, (
"only a cross-ws supersession may discard the edit-and-resend"
)
assert "this._pendingEditSend = null;" in cl
assert "this.setBusy(false);" in cl
# Supersession is centralized at the render (replayHistory) — clear_ui
# must not grow a path-local resync cancel of its own.
assert "clearTimeout(this._resyncTimer)" not in cl
def test_degraded_catchup_stops_live_stream_and_retries() -> None:
"""Degraded catch-up contract: close the stream FIRST (which also
clears any earlier degraded timer — disconnectSSE owns that), show a
plain-language status, then arm the retry timer with a doubling
cooldown. The retry must defer to the show edge when the tab is
hidden (reopening into a throttled tab would overflow again)."""
body = _INTERACTIVE.read_text(encoding="utf-8")
m = re.search(r"_enterDegradedCatchup\(\)\s*\{(.*?)\n \}", body, re.S)
assert m is not None, "_enterDegradedCatchup method not found"
method = m.group(1)
# Order matters: disconnect before arming the timer, or the fresh
# timer would be cancelled by its own disconnect.
assert method.index("this.disconnectSSE()") < method.index("this._degradedTimer = setTimeout")
assert "Connection is slow" in method, "degraded state must use plain language"
assert "DEGRADED_COOLDOWN_MAX_MS" in method
assert "document.hidden" in method
# disconnectSSE owns the timer teardown (ws-switch / giveUp / destroy
# all supersede a pending degraded retry through it).
dis = re.search(r"disconnectSSE\(\)\s*\{(.*?)\n \}", body, re.S)
assert dis is not None
assert "clearTimeout(this._degradedTimer)" in dis.group(1)
def test_visibilitychange_closes_on_hide_reconnects_on_show() -> None:
"""Close-on-hide / replay-on-show: a hidden tab's throttled drain is
the likeliest slow consumer behind server-side overflow (the old
"PR-G closes those connections on hide" comment described a handler
that never existed). The pane installs one visibilitychange
listener, marks ITS OWN hide-closes via ``_hiddenDisconnect`` so a
show edge never resurrects a deliberately-closed stream, and the
factory's destroy removes the listener (it strongly references the
pane)."""
body = _INTERACTIVE.read_text(encoding="utf-8")
assert 'document.addEventListener("visibilitychange", this._visHandler);' in body
assert 'document.removeEventListener("visibilitychange", this._visHandler);' in body
vis = re.search(r"_onVisibilityChange\(\)\s*\{(.*?)\n \}", body, re.S)
assert vis is not None, "_onVisibilityChange method not found"
method = vis.group(1)
assert "this.disconnectSSE();" in method
assert "this._hiddenDisconnect = true;" in method
assert "this.connectSSE(this.wsId);" in method
# Reconnect only consumes OUR hide-close marker.
assert "else if (this._hiddenDisconnect)" in method
# Teardown: the factory controller removes the listener on destroy.
assert "pane._removeVisibilityHandler();" in body
# The streaming buffers survive a hide-close: disconnectSSE stays
# transport-only (no contentBuffer wipe) so the visible tail is
# intact when the tab returns.
dis = re.search(r"disconnectSSE\(\)\s*\{(.*?)\n \}", body, re.S)
assert dis is not None
assert "contentBuffer" not in dis.group(1)
def test_no_global_sse_gap_detector() -> None:
"""Live event ids are NOT strictly monotonic across concurrent
tool+content emit (the fan-out runs outside the listeners lock), so
a naive ``id !== lastEventId + 1`` gap check would false-positive.
Recovery is server-signalled (``stream_overflow``) + reconnect
replay instead. This tripwire pins the absence of the naive
arithmetic — if gap detection is ever added, it must be scoped to
the content stream only (content-vs-content never reorders)."""
code = _strip_comments(_INTERACTIVE.read_text(encoding="utf-8"))
assert not re.search(r"_lastEventId\s*[+\-]\s*1", code), (
"found lastEventId +/- 1 arithmetic — a global gap detector "
"false-positives on legal concurrent tool/content id inversion"
)
def test_overflow_helpers_extracted_to_shared_module() -> None:
"""The storm-guard constants + the two pure helpers were extracted to the
shared ``sse_overflow.js`` module (its own runtime probes live in
``test_sse_overflow_js.py``) so the interactive and coordinator panes can't
drift. Pin that the pane IMPORTS them rather than re-declaring a local
copy: a stray local ``function overflowWindowTripped`` / ``const
OVERFLOW_TRIP_COUNT`` would silently fork the trip math again."""
body = _INTERACTIVE.read_text(encoding="utf-8")
m = re.search(
r"import \{([^}]*)\} from \"\./sse_overflow\.js\";",
body,
re.S,
)
assert m is not None, "interactive pane must import the shared overflow helpers"
imported = m.group(1)
for name in (
"OVERFLOW_TRIP_COUNT",
"OVERFLOW_TRIP_WINDOW_MS",
"DEGRADED_COOLDOWN_BASE_MS",
"DEGRADED_COOLDOWN_MAX_MS",
"DEGRADED_COOLDOWN_RESET_MS",
"overflowWindowTripped",
"degradedCooldownStep",
):
assert name in imported, f"{name} must be imported from sse_overflow.js"
# No local fork of the extracted definitions.
assert not re.search(r"^function overflowWindowTripped\(", body, re.M), (
"overflowWindowTripped must be imported, not re-declared locally"
)
assert not re.search(r"^function degradedCooldownStep\(", body, re.M), (
"degradedCooldownStep must be imported, not re-declared locally"
)
assert not re.search(r"^const OVERFLOW_TRIP_COUNT\s*=", body, re.M), (
"the trip constants must be imported, not re-declared locally"
)
def test_note_stream_overflow_does_not_reset_cooldown() -> None:
"""The exact finding [0] bug shape must not regress: _noteStreamOverflow
only counts + trips; it must NOT touch _degradedCooldownMs (the reset
that defeated the ladder lived here). The ladder decision lives solely
in _enterDegradedCatchup, keyed off _lastDegradedAt via
degradedCooldownStep."""
body = _INTERACTIVE.read_text(encoding="utf-8")
note = re.search(r"_noteStreamOverflow\(\)\s*\{(.*?)\n \}", body, re.S)
assert note is not None, "_noteStreamOverflow not found"
assert "_degradedCooldownMs" not in note.group(1), (
"_noteStreamOverflow must not write _degradedCooldownMs — that reset "
"was the bug that stopped the ladder escalating"
)
enter = re.search(r"_enterDegradedCatchup\(\)\s*\{(.*?)\n \}", body, re.S)
assert enter is not None
assert "degradedCooldownStep(" in enter.group(1)
assert "this._lastDegradedAt = now" in enter.group(1)
def test_recover_beat_defers_reconnect_when_tab_hidden() -> None:
"""Review round-2 finding [1]: the factory's transient-error recovery
beat (recoverTimer) must NOT reopen an EventSource into a hidden tab —
that re-creates the throttled slow-consumer overflow that close-on-hide
exists to prevent. It guards on document.hidden and defers to the
visibilitychange show edge (marking _hiddenDisconnect)."""
body = _INTERACTIVE.read_text(encoding="utf-8")
beat = re.search(r"recoverTimer = setTimeout\(\(\) => \{(.*?)\n \}, 5000\);", body, re.S)
assert beat is not None, "recoverTimer setTimeout body not found"
b = beat.group(1)
assert "document.hidden" in b, "recovery beat must guard on document.hidden"
assert "pane._hiddenDisconnect = true" in b, (
"recovery beat must defer to the show edge when hidden"
)
# The hidden guard must precede the reconnect (connectSSE) so it can't fall
# through to reopening the stream.
assert b.index("document.hidden") < b.index("pane.connectSSE(pane.wsId)")
def test_giveup_removes_visibility_handler() -> None:
"""Review round-2 finding [3]: giveUp() (markDead) must detach the
visibility handler and clear _hiddenDisconnect, or a tab hidden before
the give-up resurrects the dead controller's stream on return (the show
edge would connectSSE the closed ws and 404-reconnect it forever)."""
body = _INTERACTIVE.read_text(encoding="utf-8")
give = re.search(r"const giveUp = function \(\) \{(.*?)\n \};", body, re.S)
assert give is not None, "giveUp function body not found"
g = give.group(1)
assert "pane._removeVisibilityHandler();" in g, (
"giveUp must remove the visibility handler so a show edge can't resurrect a dead controller"
)
# _removeVisibilityHandler also clears _hiddenDisconnect (pinned in its body).
rvh = re.search(r"_removeVisibilityHandler\(\)\s*\{(.*?)\n \}", body, re.S)
assert rvh is not None
assert "this._hiddenDisconnect = false" in rvh.group(1)
def test_connectsse_defers_open_when_tab_hidden() -> None:
"""PR #805 review (Copilot + R3): connectSSE is the single connect
chokepoint and must not open an EventSource into a hidden tab. The
fresh-connect path (_loadHistoryThenConnect) has no timer guard, so a
first load in a background tab would otherwise open a throttled stream —
the slow-consumer overflow this PR exists to prevent. The guard sits
AFTER the visibilitychange-handler install (so the show edge can
reconnect) and AFTER the wsId assignment (so it targets the right ws),
and BEFORE `new EventSource` (so nothing opens)."""
body = _INTERACTIVE.read_text(encoding="utf-8")
start = body.index("connectSSE(wsId) {")
open_at = body.index("new EventSource(evtUrl)", start)
head = body[start:open_at] # connectSSE up to the EventSource open
assert "if (document.hidden) {" in head, (
"connectSSE must guard on document.hidden BEFORE opening the stream"
)
assert "this._hiddenDisconnect = true;" in head, (
"the deferred connect must mark _hiddenDisconnect so the show edge reconnects"
)
assert head.index("this.wsId = wsId;") < head.index("if (document.hidden) {")
assert head.index('addEventListener("visibilitychange"') < head.index("if (document.hidden) {")
def test_send_post_abort_machinery_is_gone() -> None:
"""The parked-POST era's client abort machinery must stay deleted in
BOTH panes: sends during a command window are answered "queued"
immediately (server-side defer-and-drain), so there is no long-lived
POST for a compaction-aware bound (``sendAbortMs``) to protect, and
dismissal is bind() → server-confirmed DELETE — never a POST abort
(``_sendAbort``), which fired on the interjection path too and
dispatched "dismissed" messages anyway. Reintroducing either hook
means re-parking the POST; that design deterministically dropped
messages from every timeout-bounded caller (coordinator client and
console proxy at 30s, SDKs, stock proxies)."""
interactive = _INTERACTIVE.read_text(encoding="utf-8")
coordinator = (_ROOT / "turnstone/console/static/coordinator/coordinator.js").read_text(
encoding="utf-8"
)
conversation = (_ROOT / "turnstone/shared_static/conversation.js").read_text(encoding="utf-8")
composer_queue = (_ROOT / "turnstone/shared_static/composer_queue.js").read_text(
encoding="utf-8"
)
for name, src in (
("interactive.js", interactive),
("coordinator.js", coordinator),
("conversation.js", conversation),
("composer_queue.js", composer_queue),
):
assert "sendAbortMs" not in src, f"{name}: the compaction-aware abort bound is dead"
assert "_sendAbort" not in src, f"{name}: dismiss must be bind() → DELETE, not a POST abort"
# The flat wedged-node bound stands in both panes: every /send answers
# within RTT now (dispatched / queued / deferred-with-msg_id).
assert "sendCtrl.abort(), 15000" in interactive
assert "sendCtrl.abort(), 15000" in coordinator
# The deferred-attachment count rides bind()'s documented options seam
# (controller dataset) — the per-pane element expando is dead.
for name, src in (
("interactive.js", interactive),
("coordinator.js", coordinator),
("composer_queue.js", composer_queue),
):
assert "_deferredAttachments" not in src, (
f"{name}: deferred state must ride bind(el, msgId, opts), not an expando"
)
def test_deferred_send_settle_protocol_pins() -> None:
"""The deferred-chip settle protocol (round 7, C4): a deferred send's
queued chip keeps its retraction affordance exactly until the message
truly leaves the parked list. Pins the controller's contract and both
panes' wiring — losing any of these silently re-promotes parked
messages to "sent" while the server still honors DELETE (loss
disguised as delivery on a node restart)."""
interactive = _INTERACTIVE.read_text(encoding="utf-8")
coordinator = (_ROOT / "turnstone/console/static/coordinator/coordinator.js").read_text(
encoding="utf-8"
)
composer_queue = (_ROOT / "turnstone/shared_static/composer_queue.js").read_text(
encoding="utf-8"
)
# Controller: bind() stores the options on its own dataset state...
assert "function bind(el, msgId, opts)" in composer_queue
assert 'el.dataset.deferred = "1"' in composer_queue
assert "el.dataset.attachedCount = String(opts.attachedCount)" in composer_queue
# ...the idle sweep skips deferred AND unbound chips (the "idle ⇒
# drained" invariant is untrue for both)...
assert "if (el.dataset.deferred) return;" in composer_queue
assert "if (!el.dataset.msgId) return;" in composer_queue
# ...and settleDeferred branches on the fold-in arm: clear the flag
# only (DELETE still genuinely removes a folded message until the seam
# drains), promote only on the fresh-spawn arm.
assert "function settleDeferred(msgId, folded)" in composer_queue
assert "delete target.dataset.deferred;" in composer_queue
assert "settleDeferred: settleDeferred" in composer_queue
# A barrier-deferred entry can dispatch within milliseconds of its ack,
# so the SSE settle can beat the POST response's bind(): the controller
# parks chip-absent settles and bind() reconciles them — without this a
# raced chip stays flagged deferred and the idle sweep skips it forever.
# Expiry is TTL-based: a size cap evicted exactly this tab's raced
# settle when a window closed with a burst of deferred sends (ours
# parks FIRST, the foreign settles behind it overflow the cap).
assert "_preBindSettles" in composer_queue
assert "_preBindSettles.has(msgId)" in composer_queue
assert "PRE_BIND_SETTLE_TTL_MS" in composer_queue
assert "_preBindSettles.size" not in composer_queue, "size-cap eviction must stay dead"
# The full send-response settle matrix lives ONCE, in the shared
# helper — retro-convert (a parked, still-retractable message must not
# render as a sent bubble), the deferred busy-undo, and the queue_full
# idle-pane cleanup (bubble removed + busy restored: the refusal can
# now fire with no worker and no drain alive, so no state event would
# ever unstick the composer).
assert "export function settleSendResponse(queue, data, ctx)" in composer_queue
assert "!queuedEl && data.deferred" in composer_queue
assert "deferred: !!data.deferred" in composer_queue
assert "attachedCount: (data.attached_ids || []).length" in composer_queue
assert "ctx.busyIsOptimistic()" in composer_queue
assert composer_queue.count("ctx.optimisticEl.remove()") >= 2, (
"both the retro-convert and queue_full arms must clear the optimistic bubble"
)
# The missed-edge settle: a non-deferred chip binding onto an
# already-idle pane missed its only sweep — the post-bind promote
# (keyed on POST-bind chip state, honoring the aria-busy
# dismiss-in-flight discipline) is what settles it.
assert "!ctx.paneIsBusy()" in composer_queue
assert 'queuedEl.hasAttribute("aria-busy")' in composer_queue
# Both panes route their parsed /send response through the helper and
# consume the pane-tier settle event; the busy stamp is centralized in
# each pane's setBusy (source defaults to "server" — only the send
# flow's optimistic flip may ever be undone).
for name, src in (("interactive.js", interactive), ("coordinator.js", coordinator)):
assert "settleSendResponse(" in src, f"{name}: settle matrix must be the shared helper"
assert "busyIsOptimistic" in src, name
assert "paneIsBusy" in src, f"{name}: the missed-edge settle needs the live flag"
assert 'setBusy(true, "optimistic")' in src, f"{name}: optimistic flip must stamp"
assert "parsePriority(" in src, f"{name}: shared !!! parse"
assert 'case "message_dispatched"' in src, f"{name}: settle event not consumed"
assert "settleDeferred(" in src, name
# /command's degraded outcomes are ALL surfaced: busy, running (the
# backstop answer), error (503 — the worker never spawned), the
# status-less non-2xx arm (404 / proxy 502), and the transport catch —
# silence at any of them reads as success.
assert 'body.status === "running"' in interactive
assert 'body.status === "error"' in interactive
assert "Command failed (HTTP " in interactive
assert '"Command failed: " + err.message' in interactive
def test_settle_send_response_missed_edge_behavior(tmp_path) -> None:
"""Execute the shared settle helper under node and pin the
missed-edge matrix behaviorally (not just textually): a non-deferred
chip binding onto an idle pane promotes; a busy pane, a deferred
chip, and a dismiss-in-flight chip do not."""
import shutil
import subprocess
if shutil.which("node") is None:
pytest.skip("node binary not available on PATH")
helper = _ROOT / "turnstone/shared_static/composer_queue.js"
script = tmp_path / "settle_harness.mjs"
script.write_text(
f'const {{ settleSendResponse }} = await import("file://{helper}");\n'
+ """
function makeEl(over) {
const el = {
isConnected: true,
classList: { contains: (c) => c === "msg-queued" },
dataset: {},
hasAttribute: () => false,
};
return Object.assign(el, over || {});
}
function run(queuedEl, paneBusy, data) {
const calls = [];
const queue = {
bind: (el, id, opts) => {
calls.push("bind");
// Mirror the real bind: stamp the deferred flag from opts.
if (opts && opts.deferred) el.dataset.deferred = "1";
},
promote: () => calls.push("promote"),
remove: () => calls.push("remove"),
addQueuedMessage: () => makeEl(),
};
settleSendResponse(queue, data, {
queuedEl,
optimisticEl: null,
isBusy: true,
displayText: "t",
priority: "notice",
setBusy: () => {},
busyIsOptimistic: () => false,
paneIsBusy: () => paneBusy,
renderError: () => {},
consumeAttachments: () => {},
});
return calls;
}
const queued = { status: "queued", msg_id: "m1" };
let c = run(makeEl(), false, queued);
if (!(c.includes("bind") && c.includes("promote")))
throw new Error("missed-edge chip must promote: " + c);
c = run(makeEl(), true, queued);
if (c.includes("promote")) throw new Error("busy pane must not promote: " + c);
c = run(makeEl(), false, { status: "queued", msg_id: "m1", deferred: true });
if (c.includes("promote"))
throw new Error("deferred chip is message_dispatched's: " + c);
c = run(makeEl({ hasAttribute: (a) => a === "aria-busy" }), false, queued);
if (c.includes("promote"))
throw new Error("dismiss-in-flight chip must be left to its DELETE verdict: " + c);
// Null / non-object 2xx body (a misbehaving proxy answering `200 null`): the
// helper normalizes it to {} so neither call site guards — it must fall through
// to the unknown/"ok" arm and SETTLE the optimistic chip (promote), never throw
// and strand a delivered message as a connection error. (The no-op
// consumeAttachments stub cannot prevent this: data.attached_ids is evaluated to
// build the :639 call args, so against unfixed code this line throws and crashes
// the harness.)
c = run(makeEl(), false, null);
if (!c.includes("promote"))
throw new Error("null body must settle via unknown-ok, not throw: " + c);
console.log("settle matrix OK");
""",
encoding="utf-8",
)
proc = subprocess.run(
["node", str(script)],
capture_output=True,
text=True,
timeout=15,
)
assert proc.returncode == 0, f"settle harness failed:\n{proc.stderr}\n{proc.stdout}"