diff --git a/docs/console.md b/docs/console.md index 874b6488..32917bf6 100644 --- a/docs/console.md +++ b/docs/console.md @@ -61,6 +61,8 @@ The collector (`turnstone/console/collector.py`) maintains an in-memory snapshot 3. **Poll loop** — fetches `GET /v1/api/dashboard` and `GET /health` from each known node every 10 seconds. Uses `ThreadPoolExecutor(max_workers=50)` for parallelism. Each poll replaces the node's workstream list with the authoritative server data. +A `get_snapshot()` method builds the full cluster state under a single lock acquisition — overview aggregates and per-node workstream lists in one atomic read. This is served both as a REST endpoint and as the initial SSE event on client connect. + ### Thread Safety All reads and writes to the node/workstream map are protected by a single `threading.Lock`. Query methods acquire the lock, copy data, and release before returning. @@ -146,6 +148,38 @@ Single node detail with all its workstreams. } ``` +### `GET /v1/api/cluster/snapshot` + +Full cluster state in a single response — all nodes with their workstreams plus overview aggregates. Built under a single lock for internal consistency. Used by the browser on initial load and SSE reconnect. + +```json +{ + "nodes": [ + { + "node_id": "db-west-04", + "server_url": "http://10.0.3.4:8080", + "max_ws": 10, + "reachable": true, + "version": "0.3.0", + "health": {"status": "ok", "version": "0.3.0"}, + "aggregate": {"total_tokens": 48200, "total_tool_calls": 156}, + "workstreams": [ + {"id": "a1b2c3d4", "name": "perf-db-west", "state": "running", ...} + ] + } + ], + "overview": { + "nodes": 847, + "workstreams": 4219, + "states": {"running": 1847, "thinking": 312, "attention": 89, "idle": 1940, "error": 31}, + "aggregate": {"total_tokens": 12400000, "total_tool_calls": 34200}, + "version_drift": false, + "versions": ["0.3.0"] + }, + "timestamp": 1709294400.0 +} +``` + ### `POST /v1/api/cluster/workstreams/new` Create a new workstream on a target node. Dispatches a `CreateWorkstreamMessage` through the Redis MQ pipeline — the bridge on the target node picks it up and creates the workstream on the server. Requires `write` scope. @@ -182,7 +216,7 @@ Creation is asynchronous — the response confirms the MQ message was dispatched ### `GET /v1/api/cluster/events` -Server-Sent Events stream for real-time cluster updates. +Server-Sent Events stream for real-time cluster updates. The first event is always a `snapshot` containing the full cluster state (same shape as `GET /v1/api/cluster/snapshot` with an added `type: "snapshot"` field), followed by incremental events: ``` data: {"type":"cluster_state","ws_id":"a1b2","node_id":"db-west-04","state":"running"} @@ -368,6 +402,8 @@ On submit, `POST /v1/api/cluster/workstreams/new` dispatches the creation reques All five views receive live updates via SSE — state cards update counts, node rows update metrics, workstream rows update state indicators. +The browser maintains a local `clusterState` object that mirrors the cluster snapshot. It is initialized from the SSE `snapshot` event on connect (or via `GET /v1/api/cluster/snapshot` on initial page load) and updated incrementally by SSE events. View navigation reads from local state — no API round-trips needed after the initial snapshot. + ### 5. Admin Panel Accessed via the "admin" button in the header (visible when authenticated diff --git a/docs/diagrams/11-console-data-flow.puml b/docs/diagrams/11-console-data-flow.puml index 58f8e4e9..2e056570 100644 --- a/docs/diagrams/11-console-data-flow.puml +++ b/docs/diagrams/11-console-data-flow.puml @@ -89,10 +89,15 @@ deactivate CC Browser -> Server : GET /v1/api/cluster/events activate Server +Server -> CC : get_snapshot() +CC --> Server : ClusterSnapshot\n(full current state) + Server -> CC : register_listener(queue) note right : Per-client queue.Queue(maxsize=500)\nSSE via EventSourceResponse + run_in_executor() -loop continuous +Server -> Browser : data: {"type":"snapshot",...}\n(full state as first SSE event) + +loop continuous (incremental updates) CC -> Server : event via listener queue\n(from any of the 3 threads) Server -> Browser : data: {"type":"cluster_state",...}\n\n end @@ -105,6 +110,13 @@ Browser -> Server : connection closed Server -> CC : unregister_listener(queue) deactivate Server +== Browser REST: Snapshot == + +Browser -> Server : GET /v1/api/cluster/snapshot +Server -> CC : get_snapshot() +CC --> Server : ClusterSnapshot\n(full current state) +Server --> Browser : JSON response + == Browser REST Requests == Browser -> Server : GET /v1/api/cluster/overview diff --git a/docs/diagrams/13-sdk-architecture.puml b/docs/diagrams/13-sdk-architecture.puml index eff1a3ff..a6984374 100644 --- a/docs/diagrams/13-sdk-architecture.puml +++ b/docs/diagrams/13-sdk-architecture.puml @@ -45,6 +45,7 @@ package "turnstone/sdk/ (Python)" { + nodes() + workstreams() + node_detail() + + snapshot() + create_workstream() + stream_cluster_events() + login() / logout() @@ -129,6 +130,7 @@ package "sdk/typescript/ (TypeScript)" { class "TurnstoneConsole" as TSConsole <> { + overview() + nodes() + + snapshot() + clusterEvents() ... } diff --git a/docs/sdk.md b/docs/sdk.md index 608dea90..3eb2c147 100644 --- a/docs/sdk.md +++ b/docs/sdk.md @@ -95,6 +95,7 @@ Both `TurnstoneConsole` (sync) and `AsyncTurnstoneConsole` (async) expose: | | `nodes(*, sort, limit, offset)` | `ClusterNodesResponse` | | | `workstreams(*, state, node, search, sort, page, per_page)` | `ClusterWorkstreamsResponse` | | | `node_detail(node_id)` | `NodeDetailResponse` | +| | `snapshot()` | `ClusterSnapshotResponse` | | | `create_workstream(*, node_id, name, model, initial_message)` | `ConsoleCreateWsResponse` | | **Schedules** | `list_schedules()` | `ListSchedulesResponse` | | | `create_schedule(*, name, schedule_type, initial_message, ...)` | `ScheduleInfo` | @@ -146,6 +147,9 @@ SSE events are deserialized into typed dataclasses. Use `event.type` to discrimi | `node_lost` | `NodeLostEvent` | `node_id` | | `cluster_state` | `ClusterStateEvent` | `ws_id`, `node_id`, `state`, `tokens` | | `ws_created` | `ClusterWsCreatedEvent` | `ws_id`, `node_id`, `name` | +| `ws_closed` | `ClusterWsClosedEvent` | `ws_id` | +| `ws_rename` | `ClusterWsRenameEvent` | `ws_id`, `name` | +| `snapshot` | `ClusterSnapshotEvent` | `nodes`, `overview`, `timestamp` | ### TurnResult diff --git a/sdk/typescript/src/console.ts b/sdk/typescript/src/console.ts index bd455241..2e517492 100644 --- a/sdk/typescript/src/console.ts +++ b/sdk/typescript/src/console.ts @@ -6,6 +6,7 @@ import type { AuthStatusResponse, ClusterNodesResponse, ClusterOverviewResponse, + ClusterSnapshotResponse, ClusterWorkstreamsResponse, ConsoleCreateWsRequest, ConsoleCreateWsResponse, @@ -33,6 +34,10 @@ export class TurnstoneConsole extends BaseClient { return this.request("GET", "/v1/api/cluster/overview"); } + async snapshot(): Promise { + return this.request("GET", "/v1/api/cluster/snapshot"); + } + async nodes(opts?: NodesOptions): Promise { return this.request("GET", "/v1/api/cluster/nodes", { params: { diff --git a/sdk/typescript/src/events.ts b/sdk/typescript/src/events.ts index 54967279..e335d3d1 100644 --- a/sdk/typescript/src/events.ts +++ b/sdk/typescript/src/events.ts @@ -1,3 +1,5 @@ +import type { ClusterOverviewResponse, ClusterSnapshotNode } from "./types.js"; + // --------------------------------------------------------------------------- // Server SSE events // --------------------------------------------------------------------------- @@ -191,6 +193,13 @@ export interface ClusterWsRenameEvent { name: string; } +export interface ClusterSnapshotEvent { + type: "snapshot"; + nodes: ClusterSnapshotNode[]; + overview: ClusterOverviewResponse; + timestamp: number; +} + /** Discriminated union of all console cluster SSE event types. */ export type ClusterEvent = | NodeJoinedEvent @@ -198,7 +207,8 @@ export type ClusterEvent = | ClusterStateEvent | ClusterWsCreatedEvent | ClusterWsClosedEvent - | ClusterWsRenameEvent; + | ClusterWsRenameEvent + | ClusterSnapshotEvent; // --------------------------------------------------------------------------- // Type guards diff --git a/sdk/typescript/src/index.ts b/sdk/typescript/src/index.ts index 3fdc8455..eadb4c9b 100644 --- a/sdk/typescript/src/index.ts +++ b/sdk/typescript/src/index.ts @@ -55,6 +55,7 @@ export type { ClusterWsCreatedEvent, ClusterWsClosedEvent, ClusterWsRenameEvent, + ClusterSnapshotEvent, } from "./events.js"; export { @@ -97,6 +98,8 @@ export type { ClusterOverviewResponse, ClusterNodeInfo, ClusterNodesResponse, + ClusterSnapshotNode, + ClusterSnapshotResponse, ClusterWorkstreamInfo, ClusterWorkstreamsResponse, NodeDetailResponse, diff --git a/sdk/typescript/src/types.ts b/sdk/typescript/src/types.ts index 6a070800..b637668d 100644 --- a/sdk/typescript/src/types.ts +++ b/sdk/typescript/src/types.ts @@ -244,6 +244,23 @@ export interface NodeDetailResponse { aggregate: ClusterAggregate; } +export interface ClusterSnapshotNode { + node_id: string; + server_url: string; + max_ws: number; + reachable: boolean; + version: string; + health: Record; + aggregate: Record; + workstreams: ClusterWorkstreamInfo[]; +} + +export interface ClusterSnapshotResponse { + nodes: ClusterSnapshotNode[]; + overview: ClusterOverviewResponse; + timestamp: number; +} + export interface ConsoleCreateWsRequest { node_id?: string; name?: string; diff --git a/tests/test_console.py b/tests/test_console.py index c2941f64..81cd6b46 100644 --- a/tests/test_console.py +++ b/tests/test_console.py @@ -446,6 +446,44 @@ class TestCollectorQueries: def test_get_node_detail_not_found(self, populated_collector): assert populated_collector.get_node_detail("nonexistent") is None + def test_get_snapshot_empty(self): + c = _make_collector() + snap = c.get_snapshot() + assert snap["nodes"] == [] + assert snap["overview"]["nodes"] == 0 + assert snap["overview"]["workstreams"] == 0 + assert snap["overview"]["states"]["running"] == 0 + assert "timestamp" in snap + + def test_get_snapshot_with_nodes(self, populated_collector): + snap = populated_collector.get_snapshot() + assert len(snap["nodes"]) == 2 + assert snap["overview"]["nodes"] == 2 + assert snap["overview"]["workstreams"] == 3 + assert snap["overview"]["states"]["running"] == 1 + assert snap["overview"]["states"]["attention"] == 1 + assert snap["overview"]["states"]["idle"] == 1 + assert snap["overview"]["aggregate"]["total_tokens"] == 17000 + assert snap["timestamp"] > 0 + # Each node should embed its workstreams + node_ids = {n["node_id"] for n in snap["nodes"]} + assert node_ids == {"node-a", "node-b"} + for n in snap["nodes"]: + if n["node_id"] == "node-a": + assert len(n["workstreams"]) == 2 + elif n["node_id"] == "node-b": + assert len(n["workstreams"]) == 1 + + def test_get_snapshot_consistency(self, populated_collector): + """Snapshot overview should match get_overview().""" + snap = populated_collector.get_snapshot() + overview = populated_collector.get_overview() + assert snap["overview"]["nodes"] == overview["nodes"] + assert snap["overview"]["workstreams"] == overview["workstreams"] + assert snap["overview"]["states"] == overview["states"] + assert snap["overview"]["aggregate"] == overview["aggregate"] + assert snap["overview"]["version_drift"] == overview["version_drift"] + # --------------------------------------------------------------------------- # ClusterStateEvent protocol tests @@ -536,6 +574,31 @@ class TestConsoleHTTPEndpoints: "workstreams": [], "aggregate": {}, } + collector.get_snapshot.return_value = { + "nodes": [ + { + "node_id": "node-a", + "server_url": "http://a:8080", + "max_ws": 10, + "reachable": True, + "version": "0.5.0", + "health": {}, + "aggregate": {"total_tokens": 50000, "total_tool_calls": 200}, + "workstreams": [ + {"id": "ws1", "name": "test", "state": "running", "node": "node-a"}, + ], + }, + ], + "overview": { + "nodes": 3, + "workstreams": 15, + "states": {"running": 5, "thinking": 2, "attention": 1, "idle": 6, "error": 1}, + "aggregate": {"total_tokens": 50000, "total_tool_calls": 200}, + "version_drift": False, + "versions": ["0.5.0"], + }, + "timestamp": 1234567890.0, + } return collector @pytest.fixture() @@ -615,6 +678,16 @@ class TestConsoleHTTPEndpoints: assert status == 404 assert "error" in data + def test_get_snapshot(self, client, mock_collector): + status, data = self._get(client, "/v1/api/cluster/snapshot") + assert status == 200 + assert len(data["nodes"]) == 1 + assert data["nodes"][0]["node_id"] == "node-a" + assert data["overview"]["nodes"] == 3 + assert data["overview"]["workstreams"] == 15 + assert data["timestamp"] == 1234567890.0 + mock_collector.get_snapshot.assert_called_once() + def test_health_endpoint(self, client, mock_collector): status, data = self._get(client, "/health") assert status == 200 diff --git a/turnstone/api/console_schemas.py b/turnstone/api/console_schemas.py index 25c153cd..50b80b65 100644 --- a/turnstone/api/console_schemas.py +++ b/turnstone/api/console_schemas.py @@ -2,6 +2,8 @@ from __future__ import annotations +from typing import Any + from pydantic import BaseModel, Field # --------------------------------------------------------------------------- @@ -48,7 +50,7 @@ class ClusterNodeInfo(BaseModel): total_tokens: int = 0 started: float = 0.0 reachable: bool = True - health: dict[str, str] = Field(default_factory=dict) + health: dict[str, Any] = Field(default_factory=dict) version: str = "" @@ -91,12 +93,34 @@ class ClusterWorkstreamsResponse(BaseModel): class NodeDetailResponse(BaseModel): node_id: str server_url: str = "" - health: dict[str, str] = Field(default_factory=dict) + health: dict[str, Any] = Field(default_factory=dict) workstreams: list[ClusterWorkstreamInfo] = [] aggregate: dict[str, int] = Field(default_factory=dict) reachable: bool = True +# --------------------------------------------------------------------------- +# Cluster snapshot +# --------------------------------------------------------------------------- + + +class ClusterSnapshotNode(BaseModel): + node_id: str + server_url: str = "" + max_ws: int = 10 + reachable: bool = True + version: str = "" + health: dict[str, Any] = Field(default_factory=dict) + aggregate: dict[str, int] = Field(default_factory=dict) + workstreams: list[ClusterWorkstreamInfo] = [] + + +class ClusterSnapshotResponse(BaseModel): + nodes: list[ClusterSnapshotNode] + overview: ClusterOverviewResponse + timestamp: float = 0.0 + + # --------------------------------------------------------------------------- # Workstream creation # --------------------------------------------------------------------------- diff --git a/turnstone/api/console_spec.py b/turnstone/api/console_spec.py index 7b62d03b..2618014f 100644 --- a/turnstone/api/console_spec.py +++ b/turnstone/api/console_spec.py @@ -10,6 +10,7 @@ if TYPE_CHECKING: from turnstone.api.console_schemas import ( ClusterNodesResponse, ClusterOverviewResponse, + ClusterSnapshotResponse, ClusterWorkstreamsResponse, ConsoleCreateWsRequest, ConsoleCreateWsResponse, @@ -97,14 +98,23 @@ CONSOLE_ENDPOINTS: list[EndpointSpec] = [ error_codes=[400, 404, 503], tags=["Cluster"], ), + EndpointSpec( + "/v1/api/cluster/snapshot", + "GET", + "Full cluster state snapshot", + description="Returns the complete cluster state: all nodes with their workstreams " + "and overview aggregates. Used for initial load and reconnection.", + response_model=ClusterSnapshotResponse, + tags=["Cluster"], + ), # --- Streaming --- EndpointSpec( "/v1/api/cluster/events", "GET", "Cluster SSE event stream", description="Server-Sent Events stream for real-time cluster updates. " - "Returns text/event-stream with node_joined, node_lost, cluster_state, " - "ws_created, ws_closed, ws_rename events.", + "First event is a 'snapshot' with full cluster state, followed by " + "node_joined, node_lost, cluster_state, ws_created, ws_closed, ws_rename events.", tags=["Streaming"], ), # --- Auth --- @@ -270,6 +280,7 @@ _ALL_MODELS: list[type[BaseModel]] = [ ClusterNodesResponse, ClusterWorkstreamsResponse, NodeDetailResponse, + ClusterSnapshotResponse, ConsoleCreateWsRequest, ConsoleCreateWsResponse, ConsoleHealthResponse, diff --git a/turnstone/console/collector.py b/turnstone/console/collector.py index 3bd047c7..914f00c6 100644 --- a/turnstone/console/collector.py +++ b/turnstone/console/collector.py @@ -379,11 +379,11 @@ class ClusterCollector: ) total = len(items) - # Sort + # Sort (secondary key: node_id for stable ordering) if sort_by == "activity": - items.sort(key=lambda n: n["ws_running"] + n["ws_attention"], reverse=True) + items.sort(key=lambda n: (-(n["ws_running"] + n["ws_attention"]), n["node_id"])) elif sort_by == "tokens": - items.sort(key=lambda n: n["total_tokens"], reverse=True) + items.sort(key=lambda n: (-n["total_tokens"], n["node_id"])) elif sort_by == "name": items.sort(key=lambda n: n["node_id"]) @@ -455,6 +455,89 @@ class ClusterCollector: "reachable": node.reachable, } + def get_snapshot(self) -> dict[str, Any]: + """Build a complete cluster snapshot under a single lock. + + Returns everything the UI needs to render the full dashboard: + all nodes with their workstreams plus pre-computed overview aggregates. + """ + with self._lock: + return self._build_snapshot_locked() + + def get_snapshot_and_register(self, q: queue.Queue[dict[str, Any]]) -> dict[str, Any]: + """Build snapshot and register listener atomically. + + Acquiring both locks ensures no event can be published between + the snapshot read and the listener registration — the client + receives the snapshot followed by every subsequent event with + no gap. + """ + with self._lock: + snap = self._build_snapshot_locked() + with self._listeners_lock: + self._listeners.append(q) + return snap + + def _build_snapshot_locked(self) -> dict[str, Any]: + """Build snapshot data — caller must hold ``_lock``.""" + nodes_out = [] + states: dict[str, int] = { + "running": 0, + "thinking": 0, + "attention": 0, + "idle": 0, + "error": 0, + } + total_tokens = 0 + total_tool_calls = 0 + total_ws = 0 + versions: set[str] = set() + + for node in self._nodes.values(): + ws_list = [] + for ws in node.workstreams.values(): + ws_list.append(dict(ws)) + s = ws.get("state", "idle") + states[s] = states.get(s, 0) + 1 + total_ws += 1 + + total_tokens += node.aggregate.get("total_tokens", 0) + total_tool_calls += node.aggregate.get("total_tool_calls", 0) + ver = node.health.get("version", "") + if ver: + versions.add(ver) + + nodes_out.append( + { + "node_id": node.node_id, + "server_url": node.server_url, + "max_ws": node.max_ws, + "reachable": node.reachable, + "version": ver, + "health": dict(node.health), + "aggregate": dict(node.aggregate), + "workstreams": ws_list, + } + ) + + node_count = len(self._nodes) + + return { + "nodes": nodes_out, + "overview": { + "nodes": node_count, + "workstreams": total_ws, + "states": states, + "aggregate": { + "total_tokens": total_tokens, + "total_tool_calls": total_tool_calls, + }, + "version_drift": len(versions) > 1, + "versions": sorted(versions), + }, + "timestamp": time.time(), + } + # -- SSE listener management --------------------------------------------- def register_listener(self, q: queue.Queue[dict[str, Any]]) -> None: diff --git a/turnstone/console/server.py b/turnstone/console/server.py index 969be692..2cee10de 100644 --- a/turnstone/console/server.py +++ b/turnstone/console/server.py @@ -40,7 +40,7 @@ from turnstone.console.collector import ClusterCollector from turnstone.core.auth import JWT_AUD_CONSOLE, AuthMiddleware if TYPE_CHECKING: - from collections.abc import AsyncGenerator + from collections.abc import AsyncGenerator, AsyncIterator from starlette.requests import Request @@ -245,14 +245,25 @@ async def cluster_node_detail(request: Request) -> JSONResponse: return JSONResponse({"error": "Node not found"}, status_code=404) +async def cluster_snapshot(request: Request) -> JSONResponse: + collector: ClusterCollector = request.app.state.collector + return JSONResponse(collector.get_snapshot()) + + async def cluster_events_sse(request: Request) -> Response: collector: ClusterCollector = request.app.state.collector client_queue: queue.Queue[dict[str, Any]] = queue.Queue(maxsize=500) - collector.register_listener(client_queue) async def event_generator() -> AsyncGenerator[dict[str, str], None]: loop = asyncio.get_running_loop() try: + # Atomic snapshot+register — no event gap possible. + snap = await loop.run_in_executor( + None, collector.get_snapshot_and_register, client_queue + ) + snap["type"] = "snapshot" + yield {"data": json.dumps(snap)} + while True: try: event = await loop.run_in_executor( @@ -596,10 +607,36 @@ async def _proxy_sse( ) yield f"event: error\ndata: Upstream returned status {response.status_code}\n\n".encode() return - async for chunk in response.aiter_bytes(): - if await request.is_disconnected(): - return - yield chunk + byte_iter = response.aiter_bytes().__aiter__() + + async def _read_next(it: AsyncIterator[bytes]) -> bytes: + return await it.__anext__() + + read_task: asyncio.Task[bytes] | None = None + try: + while True: + if await request.is_disconnected(): + return + if read_task is None: + read_task = asyncio.create_task(_read_next(byte_iter)) + ping_wait = asyncio.create_task(asyncio.sleep(3)) + done, _ = await asyncio.wait( + {read_task, ping_wait}, + return_when=asyncio.FIRST_COMPLETED, + ) + if read_task in done: + ping_wait.cancel() + try: + yield read_task.result() + except StopAsyncIteration: + return + read_task = None + else: + # No upstream data in 3s — keepalive comment + yield b": proxy-ping\n\n" + finally: + if read_task is not None: + read_task.cancel() except httpx.HTTPError: log.debug("SSE proxy stream ended for %s", target) @@ -1195,6 +1232,7 @@ def create_app( Route("/api/cluster/workstreams", cluster_workstreams), Route("/api/cluster/workstreams/new", create_workstream, methods=["POST"]), Route("/api/cluster/node/{node_id}", cluster_node_detail), + Route("/api/cluster/snapshot", cluster_snapshot), Route("/api/cluster/events", cluster_events_sse), Route("/api/auth/login", auth_login, methods=["POST"]), Route("/api/auth/logout", auth_logout, methods=["POST"]), diff --git a/turnstone/console/static/app.js b/turnstone/console/static/app.js index fd15f3ef..1e5bfb50 100644 --- a/turnstone/console/static/app.js +++ b/turnstone/console/static/app.js @@ -1,9 +1,6 @@ // --- Shared hooks --- window.onLoginSuccess = function () { connectSSE(); - if (currentView === "overview") loadOverview(); - else if (currentView === "node") drillDownToNode(currentNodeId); - else if (currentView === "filtered") loadFilteredWorkstreams(); }; window.onLogout = function () { if (evtSource) { @@ -33,6 +30,8 @@ var _lastOverviewJson = ""; var _lastNodesJson = ""; var evtSource = null; var retryDelay = 1000; +var clusterState = null; +var _navigatingFromPopstate = false; // --- Constants --- var STATE_DISPLAY = { @@ -44,6 +43,236 @@ var STATE_DISPLAY = { }; var STATE_ORDER = ["running", "thinking", "attention", "error", "idle"]; +// --- Cluster State Model --- +function applySnapshot(data) { + clusterState = { + nodes: {}, + overview: data.overview || {}, + timestamp: data.timestamp || 0, + }; + (data.nodes || []).forEach(function (n) { + clusterState.nodes[n.node_id] = n; + }); + renderFromState(); +} + +function patchClusterState(data) { + if (!clusterState) return; + var t = data.type; + if (t === "cluster_state") { + var node = clusterState.nodes[data.node_id]; + if (node) { + (node.workstreams || []).forEach(function (ws) { + if (ws.id === data.ws_id) { + if ("state" in data) ws.state = data.state; + if ("tokens" in data) ws.tokens = data.tokens; + if ("context_ratio" in data) ws.context_ratio = data.context_ratio; + if ("activity" in data) ws.activity = data.activity; + if ("activity_state" in data) ws.activity_state = data.activity_state; + } + }); + } + } else if (t === "ws_created") { + var targetNode = clusterState.nodes[data.node_id]; + if (targetNode) { + targetNode.workstreams = targetNode.workstreams || []; + targetNode.workstreams.push({ + id: data.ws_id, + name: data.name || "", + state: "idle", + node: data.node_id, + server_url: targetNode.server_url || "", + title: data.title || "", + tokens: 0, + context_ratio: 0.0, + activity: "", + activity_state: "", + tool_calls: 0, + }); + } + } else if (t === "ws_closed") { + Object.keys(clusterState.nodes).forEach(function (nid) { + var n = clusterState.nodes[nid]; + n.workstreams = (n.workstreams || []).filter(function (ws) { + return ws.id !== data.ws_id; + }); + }); + } else if (t === "ws_rename") { + Object.keys(clusterState.nodes).forEach(function (nid) { + (clusterState.nodes[nid].workstreams || []).forEach(function (ws) { + if (ws.id === data.ws_id) ws.name = data.name || ""; + }); + }); + } else if (t === "node_joined") { + if (!clusterState.nodes[data.node_id]) { + clusterState.nodes[data.node_id] = { + node_id: data.node_id, + server_url: "", + max_ws: 10, + reachable: true, + version: "", + health: {}, + aggregate: {}, + workstreams: [], + }; + } + } else if (t === "node_lost") { + delete clusterState.nodes[data.node_id]; + } else { + return; + } + scheduleRender(); +} + +var _renderTimer = null; +function scheduleRender() { + if (_renderTimer) return; + _renderTimer = requestAnimationFrame(function () { + _renderTimer = null; + recomputeOverview(); + renderFromState(); + }); +} + +function recomputeOverview() { + if (!clusterState) return; + var states = { running: 0, thinking: 0, attention: 0, idle: 0, error: 0 }; + var totalTokens = 0, + totalToolCalls = 0, + totalWs = 0; + var versions = {}; + Object.keys(clusterState.nodes).forEach(function (nid) { + var node = clusterState.nodes[nid]; + var nodeWsTokens = 0; + (node.workstreams || []).forEach(function (ws) { + var s = ws.state || "idle"; + states[s] = (states[s] || 0) + 1; + totalWs++; + nodeWsTokens += ws.tokens || 0; + }); + var aggTokens = (node.aggregate || {}).total_tokens || 0; + totalTokens += aggTokens || nodeWsTokens; + totalToolCalls += (node.aggregate || {}).total_tool_calls || 0; + if (node.version) versions[node.version] = true; + }); + var versionList = Object.keys(versions).sort(); + clusterState.overview = { + nodes: Object.keys(clusterState.nodes).length, + workstreams: totalWs, + states: states, + aggregate: { + total_tokens: totalTokens, + total_tool_calls: totalToolCalls, + }, + version_drift: versionList.length > 1, + versions: versionList, + }; +} + +function buildNodeInfoFromSnapshot(node) { + var states = { running: 0, thinking: 0, attention: 0, idle: 0, error: 0 }; + var ws = node.workstreams || []; + ws.forEach(function (w) { + var s = w.state || "idle"; + states[s] = (states[s] || 0) + 1; + }); + var aggTokens = (node.aggregate || {}).total_tokens || 0; + if (!aggTokens) { + ws.forEach(function (w) { + aggTokens += w.tokens || 0; + }); + } + return { + node_id: node.node_id, + server_url: node.server_url || "", + ws_total: ws.length, + ws_running: states.running, + ws_thinking: states.thinking, + ws_attention: states.attention, + ws_idle: states.idle, + ws_error: states.error, + total_tokens: aggTokens, + ws_tokens: aggTokens, + max_ws: node.max_ws || 10, + started: node.started || 0, + reachable: node.reachable !== false, + health: node.health || {}, + version: node.version || "", + }; +} + +function renderFromState() { + if (!clusterState) return; + renderStatusBar(clusterState.overview); + if (currentView === "overview") { + var nodesList = Object.keys(clusterState.nodes).map(function (nid) { + return buildNodeInfoFromSnapshot(clusterState.nodes[nid]); + }); + nodesList.sort(function (a, b) { + var d = b.ws_running + b.ws_attention - (a.ws_running + a.ws_attention); + return d !== 0 ? d : a.node_id.localeCompare(b.node_id); + }); + renderNodeGroups(nodesList, nodesList.length); + document.getElementById("cluster-summary").textContent = + clusterState.overview.nodes + + " nodes \u00b7 " + + formatCount(clusterState.overview.workstreams) + + " workstreams"; + } else if (currentView === "node" && currentNodeId) { + var snapNode = clusterState.nodes[currentNodeId]; + if (snapNode) { + var wsList = snapNode.workstreams || []; + var active = wsList.filter(function (w) { + return w.state !== "idle"; + }).length; + document.getElementById("node-ws-summary").textContent = + active + " active \u00b7 " + wsList.length + " total"; + renderWsTable(document.getElementById("node-ws-table"), wsList); + } + } else if (currentView === "filtered") { + var allWs = []; + Object.keys(clusterState.nodes).forEach(function (nid) { + (clusterState.nodes[nid].workstreams || []).forEach(function (ws) { + allWs.push(ws); + }); + }); + if (currentFilter.state) { + allWs = allWs.filter(function (ws) { + return ws.state === currentFilter.state; + }); + } + if (currentFilter.node) { + allWs = allWs.filter(function (ws) { + return ws.node === currentFilter.node; + }); + } + var stateOrder = { + running: 0, + thinking: 1, + attention: 2, + error: 3, + idle: 4, + }; + allWs.sort(function (a, b) { + return (stateOrder[a.state] || 9) - (stateOrder[b.state] || 9); + }); + var total = allWs.length; + var perPage = currentFilter.per_page || 50; + var pages = Math.max(1, Math.ceil(total / perPage)); + var page = Math.min(currentFilter.page || 1, pages); + var start = (page - 1) * perPage; + var pageWs = allWs.slice(start, start + perPage); + document.getElementById("filtered-summary").textContent = + "Page " + page + " of " + pages + " (" + total + " total)"; + renderWsTable(document.getElementById("filtered-ws-table"), pageWs); + renderPagination( + document.getElementById("filtered-pagination"), + page, + pages, + ); + } +} + // --- SSE Connection --- function connectSSE() { if (evtSource) { @@ -94,19 +323,11 @@ function connectSSE() { }; } -var _refreshTimer = null; -function scheduleRefresh() { - if (_refreshTimer) return; - _refreshTimer = setTimeout(function () { - _refreshTimer = null; - if (currentView === "overview") loadOverview(); - else if (currentView === "node" && currentNodeId) - loadNodeDetail(currentNodeId); - else if (currentView === "filtered") loadFilteredWorkstreams(); - }, 250); -} - function handleClusterEvent(data) { + if (data.type === "snapshot") { + applySnapshot(data); + return; + } if ( data.type === "cluster_state" || data.type === "ws_created" || @@ -115,7 +336,7 @@ function handleClusterEvent(data) { data.type === "node_joined" || data.type === "node_lost" ) { - scheduleRefresh(); + patchClusterState(data); } if (data.type === "ws_closed" && data.reason === "evicted") { showToast("Evicted" + (data.name ? ": " + data.name : "") + " (capacity)"); @@ -135,28 +356,18 @@ function showOverview() { if (adminView) adminView.style.display = "none"; document.getElementById("breadcrumb").style.display = "none"; document.getElementById("main").scrollTop = 0; - loadOverview(); - history.pushState({ view: "overview" }, ""); + if (clusterState) renderFromState(); + else loadOverview(); + if (!_navigatingFromPopstate) history.pushState({ view: "overview" }, ""); } function loadOverview() { - var overviewP = authFetch("/v1/api/cluster/overview").then(function (r) { - return r.json(); - }); - var nodesP = authFetch("/v1/api/cluster/nodes?sort=activity&limit=1000").then( - function (r) { + authFetch("/v1/api/cluster/snapshot") + .then(function (r) { return r.json(); - }, - ); - Promise.all([overviewP, nodesP]) - .then(function (res) { - renderStatusBar(res[0]); - renderNodeGroups(res[1].nodes, res[1].total); - document.getElementById("cluster-summary").textContent = - res[0].nodes + - " nodes \u00b7 " + - formatCount(res[0].workstreams) + - " workstreams"; + }) + .then(function (data) { + applySnapshot(data); }) .catch(function () { document.getElementById("node-table").innerHTML = @@ -310,7 +521,8 @@ function groupNodes(nodes) { }); groupOrder.forEach(function (prefix) { groupMap[prefix].nodes.sort(function (a, b) { - return b.ws_running + b.ws_attention - (a.ws_running + a.ws_attention); + var d = b.ws_running + b.ws_attention - (a.ws_running + a.ws_attention); + return d !== 0 ? d : a.node_id.localeCompare(b.node_id); }); }); var groups = groupOrder.map(function (p) { @@ -653,38 +865,37 @@ function drillDownToNode(nodeId, serverUrl) { link.href = "/node/" + encodeURIComponent(nodeId) + "/"; link.style.display = ""; document.getElementById("main").scrollTop = 0; - document.getElementById("node-ws-table").innerHTML = - '
Loading workstreams...
'; - loadNodeDetail(nodeId); + if (clusterState && clusterState.nodes[nodeId]) { + renderFromState(); + } else { + document.getElementById("node-ws-table").innerHTML = + '
Loading workstreams...
'; + loadNodeDetail(nodeId); + } document.getElementById("breadcrumb-home").focus(); - history.pushState({ view: "node", nodeId: nodeId, serverUrl: serverUrl }, ""); + if (!_navigatingFromPopstate) + history.pushState( + { view: "node", nodeId: nodeId, serverUrl: serverUrl }, + "", + ); } function loadNodeDetail(nodeId) { - var detailP = authFetch( - "/v1/api/cluster/node/" + encodeURIComponent(nodeId), - ).then(function (r) { - return r.json(); - }); - var overviewP = authFetch("/v1/api/cluster/overview").then(function (r) { - return r.json(); - }); - Promise.all([detailP, overviewP]).then(function (res) { - var data = res[0]; - renderStatusBar(res[1]); - if (data.error) { + authFetch("/v1/api/cluster/snapshot") + .then(function (r) { + return r.json(); + }) + .then(function (data) { + applySnapshot(data); + if (!clusterState || !clusterState.nodes[nodeId]) { + document.getElementById("node-ws-table").innerHTML = + '
Node not found
'; + } + }) + .catch(function () { document.getElementById("node-ws-table").innerHTML = - '
' + escapeHtml(data.error) + "
"; - return; - } - var ws = data.workstreams || []; - var active = ws.filter(function (w) { - return w.state !== "idle"; - }).length; - document.getElementById("node-ws-summary").textContent = - active + " active \u00b7 " + ws.length + " total"; - renderWsTable(document.getElementById("node-ws-table"), ws); - }); + '
Failed to load
'; + }); } // --- Drill-down: Filtered --- @@ -703,9 +914,11 @@ function drillDownByState(state) { document.getElementById("filtered-title").textContent = "WORKSTREAMS — " + sd.label.toUpperCase(); document.getElementById("main").scrollTop = 0; - loadFilteredWorkstreams(); + if (clusterState) renderFromState(); + else loadFilteredWorkstreams(); document.getElementById("breadcrumb-home").focus(); - history.pushState({ view: "filtered", filter: currentFilter }, ""); + if (!_navigatingFromPopstate) + history.pushState({ view: "filtered", filter: currentFilter }, ""); } function drillDownByNode(nodeId) { @@ -721,48 +934,20 @@ function drillDownByNode(nodeId) { document.getElementById("filtered-title").textContent = "WORKSTREAMS — " + nodeId; document.getElementById("main").scrollTop = 0; - loadFilteredWorkstreams(); + if (clusterState) renderFromState(); + else loadFilteredWorkstreams(); document.getElementById("breadcrumb-home").focus(); - history.pushState({ view: "filtered", filter: currentFilter }, ""); + if (!_navigatingFromPopstate) + history.pushState({ view: "filtered", filter: currentFilter }, ""); } function loadFilteredWorkstreams() { - var params = - "page=" + currentFilter.page + "&per_page=" + currentFilter.per_page; - if (currentFilter.state) - params += "&state=" + encodeURIComponent(currentFilter.state); - if (currentFilter.node) - params += "&node=" + encodeURIComponent(currentFilter.node); - var wsP = authFetch("/v1/api/cluster/workstreams?" + params).then( - function (r) { + authFetch("/v1/api/cluster/snapshot") + .then(function (r) { return r.json(); - }, - ); - var overviewP = authFetch("/v1/api/cluster/overview").then(function (r) { - return r.json(); - }); - Promise.all([wsP, overviewP]) - .then(function (res) { - var data = res[0]; - renderStatusBar(res[1]); - document.getElementById("main").scrollTop = 0; - document.getElementById("filtered-summary").textContent = - "Page " + - data.page + - " of " + - data.pages + - " (" + - data.total + - " total)"; - renderWsTable( - document.getElementById("filtered-ws-table"), - data.workstreams, - ); - renderPagination( - document.getElementById("filtered-pagination"), - data.page, - data.pages, - ); + }) + .then(function (data) { + applySnapshot(data); }) .catch(function () { document.getElementById("filtered-ws-table").innerHTML = @@ -778,7 +963,8 @@ function renderPagination(container, page, pages) { prev.disabled = page <= 1; prev.onclick = function () { currentFilter.page--; - loadFilteredWorkstreams(); + if (clusterState) renderFromState(); + else loadFilteredWorkstreams(); }; container.appendChild(prev); var info = document.createElement("span"); @@ -789,7 +975,8 @@ function renderPagination(container, page, pages) { next.disabled = page >= pages; next.onclick = function () { currentFilter.page++; - loadFilteredWorkstreams(); + if (clusterState) renderFromState(); + else loadFilteredWorkstreams(); }; container.appendChild(next); } @@ -923,19 +1110,24 @@ function renderWsTable(container, wsList) { window.addEventListener("popstate", function (e) { var overlay = document.getElementById("login-overlay"); if (overlay && overlay.style.display !== "none") return; - if (!e.state) { - showOverview(); - return; - } - if (e.state.view === "overview") showOverview(); - else if (e.state.view === "admin" && typeof showAdmin === "function") - showAdmin(); - else if (e.state.view === "node" && e.state.nodeId) - drillDownToNode(e.state.nodeId, e.state.serverUrl); - else if (e.state.view === "filtered" && e.state.filter) { - currentFilter = e.state.filter; - if (currentFilter.state) drillDownByState(currentFilter.state); - else if (currentFilter.node) drillDownByNode(currentFilter.node); + _navigatingFromPopstate = true; + try { + if (!e.state) { + showOverview(); + return; + } + if (e.state.view === "overview") showOverview(); + else if (e.state.view === "admin" && typeof showAdmin === "function") + showAdmin(); + else if (e.state.view === "node" && e.state.nodeId) + drillDownToNode(e.state.nodeId, e.state.serverUrl); + else if (e.state.view === "filtered" && e.state.filter) { + currentFilter = e.state.filter; + if (currentFilter.state) drillDownByState(currentFilter.state); + else if (currentFilter.node) drillDownByNode(currentFilter.node); + } + } finally { + _navigatingFromPopstate = false; } }); diff --git a/turnstone/sdk/console.py b/turnstone/sdk/console.py index 224fac00..db7ea1db 100644 --- a/turnstone/sdk/console.py +++ b/turnstone/sdk/console.py @@ -16,6 +16,7 @@ from typing import TYPE_CHECKING, Any from turnstone.api.console_schemas import ( ClusterNodesResponse, ClusterOverviewResponse, + ClusterSnapshotResponse, ClusterWorkstreamsResponse, ConsoleCreateWsResponse, ConsoleHealthResponse, @@ -102,6 +103,11 @@ class AsyncTurnstoneConsole(_BaseClient): "GET", f"/v1/api/cluster/node/{node_id}", response_model=NodeDetailResponse ) + async def snapshot(self) -> ClusterSnapshotResponse: + return await self._request( + "GET", "/v1/api/cluster/snapshot", response_model=ClusterSnapshotResponse + ) + async def create_workstream( self, *, @@ -342,6 +348,9 @@ class TurnstoneConsole: def node_detail(self, node_id: str) -> NodeDetailResponse: return self._runner.run(self._async.node_detail(node_id)) + def snapshot(self) -> ClusterSnapshotResponse: + return self._runner.run(self._async.snapshot()) + def create_workstream( self, *, diff --git a/turnstone/sdk/events.py b/turnstone/sdk/events.py index ad9a476e..17e43ec9 100644 --- a/turnstone/sdk/events.py +++ b/turnstone/sdk/events.py @@ -247,6 +247,14 @@ class ClusterWsRenameEvent(ClusterEvent): name: str = "" +@dataclass +class ClusterSnapshotEvent(ClusterEvent): + type: str = "snapshot" + nodes: list[dict[str, Any]] = field(default_factory=list) + overview: dict[str, Any] = field(default_factory=dict) + timestamp: float = 0.0 + + # --------------------------------------------------------------------------- # Type registries (built after all classes are defined) # --------------------------------------------------------------------------- @@ -296,5 +304,6 @@ _CLUSTER_REGISTRY: dict[str, type[ClusterEvent]] = { ClusterWsCreatedEvent, ClusterWsClosedEvent, ClusterWsRenameEvent, + ClusterSnapshotEvent, ] }