Compare commits

..

3 Commits

Author SHA1 Message Date
Patrick Buckley 7ea150fa71 Bump version to 0.5.2 2026-03-09 13:42:28 -07:00
Patrick Buckley 4d665a5f62 fix: SSE reconnect loop — remove _sse_generation single-consumer lock
The _sse_generation mechanism assumed one SSE consumer per workstream,
but the bridge also maintains an SSE connection to each workstream.
When a new client connected (browser, proxy, or test), it incremented
the generation counter, killing the bridge's connection. The bridge
reconnected, killing the new client's connection — creating a
mutual-kill cascade that closed every SSE connection after one ping
cycle (5s).

Fix: remove _sse_generation entirely. sse-starlette handles disconnect
detection via its own ASGI task. Also remove the redundant
request.is_disconnected() check which raced with sse-starlette's
disconnect listener in Starlette 0.52.

Root cause confirmed via raw socket test: the server was sending
a zero-length chunked terminator (0\r\n\r\n) at exactly 5s,
cleanly ending the HTTP response body.
2026-03-09 13:40:36 -07:00
Patrick Buckley 3bc3250869 fix: recovered workstreams invisible in console UI (#35)
* fix: recovered workstreams invisible in console UI

Bridge startup recovery (_recover_workstreams) re-registered workstream
ownership but never published WorkstreamCreatedEvent to the cluster
channel. The collector's poll loop would pick up the workstream in its
internal state, but _apply_poll never fanned out SSE events to connected
browsers. Combined, this made channel-resumed workstreams invisible in
the console while remaining accessible through the proxied node UI.

- Bridge: emit WorkstreamCreatedEvent for each recovered workstream
- Collector: diff poll results and fan out synthetic ws_created/ws_closed
  events for workstream additions and removals
- Skip workstreams with empty IDs in poll processing
- Add 4 tests for poll-diff fanout behavior
- Update console data-flow diagram and architecture docs

* fix: address PR review — filter empty ws IDs, stable event ordering

- Filter empty-string keys from old_ids to avoid phantom ws_closed
  events if a previous poll inserted a workstream under key "".
- Sort set diffs before iterating so ws_created/ws_closed fanout
  order is deterministic across poll cycles.
2026-03-09 13:39:47 -07:00
8 changed files with 123 additions and 17 deletions
+7 -1
View File
@@ -1186,6 +1186,9 @@ for existing workstreams are auto-routed via `turnstone:ws:{ws_id}` ownership ke
If a bridge picks up a shared-queue message for a workstream owned by another node, it
re-routes to that node's queue (1 extra hop). Bridges publish heartbeats to
`turnstone:node:{node_id}` with configurable TTL for node discovery.
On startup, `_recover_workstreams` re-registers ownership of existing
workstreams and publishes `WorkstreamCreatedEvent` to the cluster channel
so the console collector picks them up immediately.
### Cluster Console
@@ -1211,7 +1214,10 @@ The console HTTP layer is a Starlette/ASGI app served by uvicorn. The SSE
endpoint uses `EventSourceResponse` with the same listener queue pattern as
the main server. `ClusterCollector`'s background threads (event subscriber,
node discovery, poll loop) use sync Redis clients and `ThreadPoolExecutor`
for parallel HTTP polling.
for parallel HTTP polling. The poll loop diffs workstream IDs between poll
cycles and fans out synthetic `ws_created`/`ws_closed` SSE events for any
changes, ensuring browser clients stay in sync even when real-time cluster
events are missed (e.g. bridge startup recovery).
The console has two write-path capabilities:
+12
View File
@@ -78,7 +78,19 @@ activate NodeA
NodeA --> CC : {status:"ok", version:"0.3.0",\nmodel:"...", workstreams:{...}}
deactivate NodeA
CC -> CC : Diff old vs new workstream IDs
CC -> CC : Replace NodeSnapshot["nodeA"]\n.workstreams, .health, .aggregate
CC -> CC : _fanout(ws_created) for\nnewly appeared workstreams
CC -> CC : _fanout(ws_closed) for\nremoved workstreams
note right of CC
Poll-diff fanout ensures
browser SSE clients learn
about workstreams that
appeared without a real-time
cluster event (e.g. bridge
startup recovery).
end note
CC -x NodeB : (SKIPPED: sim:// URL)
+1 -1
View File
@@ -4,7 +4,7 @@ build-backend = "hatchling.build"
[project]
name = "turnstone"
version = "0.5.1"
version = "0.5.2"
description = "Multi-node AI orchestration platform with tool use, agent routing, and cluster simulation."
readme = "README.md"
license = "BUSL-1.1"
+62
View File
@@ -202,6 +202,68 @@ class TestCollectorPolling:
# Should not raise
c._apply_poll("unknown", _dashboard_response(), {})
def test_apply_poll_emits_ws_created_for_new_workstream(self):
c = _make_collector()
c._nodes["node-a"] = NodeSnapshot(node_id="node-a", server_url="http://a:8080")
q: queue.Queue[dict] = queue.Queue()
c.register_listener(q)
dashboard = _dashboard_response(
workstreams=[{"id": "ws1", "name": "new-task", "state": "idle"}]
)
c._apply_poll("node-a", dashboard, {})
event = q.get_nowait()
assert event["type"] == "ws_created"
assert event["ws_id"] == "ws1"
assert event["name"] == "new-task"
assert event["node_id"] == "node-a"
def test_apply_poll_emits_ws_closed_for_removed_workstream(self):
c = _make_collector()
c._nodes["node-a"] = NodeSnapshot(
node_id="node-a",
server_url="http://a:8080",
workstreams={"ws1": {"id": "ws1", "name": "old", "state": "idle"}},
)
q: queue.Queue[dict] = queue.Queue()
c.register_listener(q)
c._apply_poll("node-a", _dashboard_response(), {})
event = q.get_nowait()
assert event["type"] == "ws_closed"
assert event["ws_id"] == "ws1"
def test_apply_poll_no_events_when_unchanged(self):
c = _make_collector()
c._nodes["node-a"] = NodeSnapshot(
node_id="node-a",
server_url="http://a:8080",
workstreams={"ws1": {"id": "ws1", "name": "same", "state": "idle"}},
)
q: queue.Queue[dict] = queue.Queue()
c.register_listener(q)
dashboard = _dashboard_response(
workstreams=[{"id": "ws1", "name": "same", "state": "running"}]
)
c._apply_poll("node-a", dashboard, {})
assert q.empty()
def test_apply_poll_skips_empty_id_workstream(self):
c = _make_collector()
c._nodes["node-a"] = NodeSnapshot(node_id="node-a", server_url="http://a:8080")
q: queue.Queue[dict] = queue.Queue()
c.register_listener(q)
dashboard = _dashboard_response(workstreams=[{"name": "no-id", "state": "idle"}])
c._apply_poll("node-a", dashboard, {})
assert q.empty()
assert len(c._nodes["node-a"].workstreams) == 0
class TestCollectorEvents:
"""Real-time event handling from cluster channel."""
+1 -1
View File
@@ -1,3 +1,3 @@
"""turnstone - Multi-node AI orchestration platform with tool use, agent routing, and cluster simulation."""
__version__ = "0.5.1"
__version__ = "0.5.2"
+27 -3
View File
@@ -273,6 +273,7 @@ class ClusterCollector:
"""Apply polled data to the in-memory node snapshot."""
ws_list = dashboard.get("workstreams", [])
aggregate = dashboard.get("aggregate", {})
pending_events: list[dict[str, Any]] = []
with self._lock:
node = self._nodes.get(node_id)
if not node:
@@ -281,12 +282,35 @@ class ClusterCollector:
node.reachable = True
node.health = health
node.aggregate = aggregate
# Replace workstreams entirely from the authoritative poll
node.workstreams = {}
# Build new workstream map
old_ids = {k for k in node.workstreams if k}
new_ws: dict[str, dict[str, Any]] = {}
for ws in ws_list:
ws_id = ws.get("id", "")
if not ws_id:
continue
ws["node"] = node_id
ws["server_url"] = node.server_url
node.workstreams[ws.get("id", "")] = ws
new_ws[ws_id] = ws
new_ids = set(new_ws.keys())
# Detect additions not yet known to SSE clients
for ws_id in sorted(new_ids - old_ids):
ws = new_ws[ws_id]
pending_events.append(
{
"type": "ws_created",
"ws_id": ws_id,
"name": ws.get("name", ""),
"node_id": node_id,
}
)
# Detect removals
for ws_id in sorted(old_ids - new_ids):
pending_events.append({"type": "ws_closed", "ws_id": ws_id})
node.workstreams = new_ws
# Fan out diffs to SSE listeners outside the lock
for event in pending_events:
self._fanout(event)
# -- query methods (thread-safe) -----------------------------------------
+9 -1
View File
@@ -206,9 +206,17 @@ class Bridge:
data = resp.json()
for ws in data.get("workstreams", []):
ws_id = ws["id"]
log.info("Recovered workstream %s (%s)", ws_id, ws.get("name", ""))
ws_name = ws.get("name", "")
log.info("Recovered workstream %s (%s)", ws_id, ws_name)
self._broker.set_ws_owner(ws_id, self._node_id)
self._start_ws_sse(ws_id)
self._publish_cluster(
WorkstreamCreatedEvent(
ws_id=ws_id,
name=ws_name,
node_id=self._node_id,
)
)
except Exception as exc:
log.warning("Could not recover workstreams: %s", exc)
+4 -10
View File
@@ -473,11 +473,9 @@ async def events_sse(request: Request) -> Response:
if not ws or not ui:
return JSONResponse({"error": "Unknown workstream"}, status_code=404)
ui._sse_generation += 1
my_gen = ui._sse_generation
# Drain stale events. A race with the worker thread is acceptable:
# worst case we discard one fresh event, and the client catches up
# via the history replay above.
# Drain stale events so this client starts fresh. A race with the
# worker thread is acceptable: worst case we discard one fresh event,
# and the client catches up via the history replay below.
while not ui._event_queue.empty():
try:
ui._event_queue.get_nowait()
@@ -509,7 +507,7 @@ async def events_sse(request: Request) -> Response:
_metrics.record_sse_connect()
try:
loop = asyncio.get_running_loop()
while my_gen == ui._sse_generation:
while True:
try:
event = await loop.run_in_executor(
None, functools.partial(ui._event_queue.get, timeout=5)
@@ -517,8 +515,6 @@ async def events_sse(request: Request) -> Response:
yield {"data": json.dumps(event)}
except queue.Empty:
pass
if await request.is_disconnected():
break
finally:
_metrics.record_sse_disconnect()
@@ -545,8 +541,6 @@ async def global_events_sse(request: Request) -> Response:
yield {"data": json.dumps(event)}
except queue.Empty:
pass
if await request.is_disconnected():
break
finally:
_metrics.record_sse_disconnect()
with listeners_lock: