refactor: rename state_sync namespace to webui and simplify handler event registration
- Rename /state_sync namespace to /webui throughout codebase - Remove get_event_types() from WebSocketHandler - handlers now process all events for their namespace - Replace per-event handler registration with namespace-wide registration - Add validate_event_type() class method for runtime event name validation - Update UserMessage instantiation to use keyword arguments (message=, attachments=) - Move send_data
frdel committed
Mar 20, 2026 at 15:34 UTC
a5620506d51161824194abfcfe373780221f00ca
25 files changed
+256
-259
api/api_message.py
+1
-1
@@ -142,7 +142,7 @@ class ApiMessage(ApiHandler):
142
)
143
144
# Send message to agent
145
- task = context.communicate(UserMessage(message, attachment_paths))
145
+ task = context.communicate(UserMessage(message=message, attachments=attachment_paths))
146
result = await task.result()
147
148
# Clean up expired chats
api/message.py
+1
-1
@@ -68,4 +68,4 @@ class Message(ApiHandler):
68
# Log to console and UI using helper function
69
mq.log_user_message(context, message, attachment_paths, message_id)
70
71
- return context.communicate(UserMessage(message, attachment_paths)), context
71
+ return context.communicate(UserMessage(message=message, attachments=attachment_paths)), context
docs/developer/websockets.md
+5
-5
@@ -139,7 +139,7 @@ These expanded flows complement the operation matrix later in the guide, ensurin
139
140
Handlers are discovered deterministically from `python/websocket_handlers/`:
141
142
-- **File entry**: `python/websocket_handlers/state_sync_handler.py` → namespace `/state_sync`
142
+- **File entry**: `python/websocket_handlers/webui_handler.py` → namespace `/webui`
143
- **Folder entry**: `python/websocket_handlers/orders/` or `python/websocket_handlers/orders_handler/` → namespace `/orders` (loads `*.py` one level deep; ignores `__init__.py` and deeper nesting)
144
- **Reserved root**: `python/websocket_handlers/_default.py` → namespace `/` (diagnostics-only by default)
145
@@ -273,7 +273,7 @@ console.log(window.runtimeInfo.id, window.runtimeInfo.isDevelopment);
273
274
### Namespaces (end-state)
275
276
-- The root namespace (`/`) is reserved and intentionally unhandled by default for application events. Feature code should connect to an explicit namespace (for example `/state_sync`).
276
+- The root namespace (`/`) is reserved and intentionally unhandled by default for application events. Feature code should connect to an explicit namespace (for example `/webui`).
277
- The frontend exposes `createNamespacedClient(namespace)` and `getNamespacedClient(namespace)` (one client instance per namespace per tab). Namespaced clients expose the same minimal API: `emit`, `request`, `on`, `off`.
278
- Unknown namespaces are rejected deterministically during the Socket.IO connect handshake with a `connect_error` payload:
279
- `err.message === "UNKNOWN_NAMESPACE"`
@@ -342,7 +342,7 @@ Example:
342
```javascript
343
import { getNamespacedClient, createCorrelationId, validateServerEnvelope } from '/js/websocket.js';
344
345
-const websocket = getNamespacedClient('/state_sync');
345
+const websocket = getNamespacedClient('/webui');
346
347
const { results } = await websocket.request(
348
'hello_request',
@@ -381,7 +381,7 @@ Example – request()
381
```javascript
382
import { getNamespacedClient } from '/js/websocket.js'
383
384
-const websocket = getNamespacedClient('/state_sync')
384
+const websocket = getNamespacedClient('/webui')
385
386
function renderError(code, message) {
387
// Map codes to UI copy; keep messages concise
@@ -409,7 +409,7 @@ Subscriptions – envelope handler
409
```javascript
410
import { getNamespacedClient } from '/js/websocket.js'
411
412
-const websocket = getNamespacedClient('/state_sync')
412
+const websocket = getNamespacedClient('/webui')
413
414
websocket.on('example_broadcast', ({ data, handlerId, eventId, correlationId }) => {
415
// handle data; errors should not typically arrive via broadcast
extensions/python/webui_ws_connect/_10_state_sync.py
new
+15
@@ -0,0 +1,15 @@
1
+from helpers.extension import Extension
2
+from helpers.print_style import PrintStyle
3
+from helpers.state_monitor import get_state_monitor, _ws_debug_enabled
4
+
5
+
6
+class StateSync(Extension):
7
+ async def execute(self, instance=None, sid: str = "", **kwargs):
8
+ if instance is None:
9
+ return
10
+
11
+ monitor = get_state_monitor()
12
+ monitor.bind_manager(instance.manager, handler_id=instance.identifier)
13
+ monitor.register_sid(instance.namespace, sid)
14
+ if _ws_debug_enabled():
15
+ PrintStyle.debug(f"[WebuiHandler] connect sid={sid}")
extensions/python/webui_ws_disconnect/_10_state_sync.py
new
+13
@@ -0,0 +1,13 @@
1
+from helpers.extension import Extension
2
+from helpers.print_style import PrintStyle
3
+from helpers.state_monitor import get_state_monitor, _ws_debug_enabled
4
+
5
+
6
+class StateSync(Extension):
7
+ async def execute(self, instance=None, sid: str = "", **kwargs):
8
+ if instance is None:
9
+ return
10
+
11
+ get_state_monitor().unregister_sid(instance.namespace, sid)
12
+ if _ws_debug_enabled():
13
+ PrintStyle.debug(f"[WebuiHandler] disconnect sid={sid}")
extensions/python/webui_ws_event/_10_state_sync.py
new
+66
@@ -0,0 +1,66 @@
1
+from helpers import runtime
2
+from helpers.extension import Extension
3
+from helpers.print_style import PrintStyle
4
+from helpers.state_monitor import get_state_monitor, _ws_debug_enabled
5
+from helpers.state_snapshot import (
6
+ StateRequestValidationError,
7
+ parse_state_request_payload,
8
+)
9
+
10
+
11
+class StateSync(Extension):
12
+ async def execute(
13
+ self,
14
+ instance=None,
15
+ sid: str = "",
16
+ event_type: str = "",
17
+ data: dict | None = None,
18
+ response_data: dict | None = None,
19
+ **kwargs,
20
+ ):
21
+ if instance is None or data is None:
22
+ return
23
+
24
+ if event_type != "state_request":
25
+ return
26
+
27
+ correlation_id = data.get("correlationId")
28
+ try:
29
+ request = parse_state_request_payload(data)
30
+ except StateRequestValidationError as exc:
31
+ PrintStyle.warning(
32
+ f"[WebuiHandler] INVALID_REQUEST sid={sid} reason={exc.reason} details={exc.details!r}"
33
+ )
34
+ if response_data is not None:
35
+ response_data["code"] = "INVALID_REQUEST"
36
+ response_data["message"] = str(exc)
37
+ return
38
+
39
+ if _ws_debug_enabled():
40
+ PrintStyle.debug(
41
+ f"[WebuiHandler] state_request sid={sid} context={request.context!r} "
42
+ f"log_from={request.log_from} notifications_from={request.notifications_from} timezone={request.timezone!r} "
43
+ f"correlation_id={correlation_id}"
44
+ )
45
+
46
+ seq_base = 1
47
+ monitor = get_state_monitor()
48
+ monitor.update_projection(
49
+ instance.namespace,
50
+ sid,
51
+ request=request,
52
+ seq_base=seq_base,
53
+ )
54
+ monitor.mark_dirty(
55
+ instance.namespace,
56
+ sid,
57
+ reason="webui_ws_event.StateSync.state_request",
58
+ )
59
+ if _ws_debug_enabled():
60
+ PrintStyle.debug(
61
+ f"[WebuiHandler] state_request accepted sid={sid} seq_base={seq_base}"
62
+ )
63
+
64
+ if response_data is not None:
65
+ response_data["runtime_epoch"] = runtime.get_runtime_id()
66
+ response_data["seq_base"] = seq_base
extensions/webui/webui_ws_push/clear_cache.js
renamed
helpers/plugins.py
+1
-1
@@ -186,7 +186,7 @@ def clear_plugin_cache():
186
187
DeferredTask().start_task(
188
send_data,
189
- endpoint_name="/state_sync",
189
+ endpoint_name="/webui",
190
event_name="clear_cache",
191
data={"areas": areas},
192
)
helpers/websocket.py
+17
-35
@@ -273,8 +273,9 @@ class WebSocketHandler(ABC):
273
"""Base class for WebSocket event handlers.
274
275
The interface mirrors :class:`helpers.api.ApiHandler` with declarative
276
- security configuration and lifecycle hooks while enforcing event-naming
277
- conventions.
276
+ security configuration and lifecycle hooks. Handlers are namespace-wide:
277
+ every inbound event for the bound namespace is dispatched to
278
+ :meth:`process_event`, which decides whether and how to respond.
279
"""
280
281
_instances: dict[type["WebSocketHandler"], "WebSocketHandler"] = {}
@@ -343,39 +344,20 @@ class WebSocketHandler(ABC):
344
WebSocketHandler._construction_tokens.pop(cls, None)
345
346
@classmethod
346
- @abstractmethod
347
- def get_event_types(cls) -> list[str]:
348
- """Return the list of event types this handler subscribes to."""
349
-
350
- @classmethod
351
- def validate_event_types(cls, event_types: Iterable[str]) -> list[str]:
352
- """Validate event type declarations.
353
-
354
- Ensures that every event name follows ``lowercase_snake_case`` naming,
355
- does not collide with Socket.IO reserved events, and that the handler
356
- does not declare duplicates.
357
- """
358
-
359
- validated: list[str] = []
360
- seen: set[str] = set()
361
- for event in event_types:
362
- if not isinstance(event, str):
363
- raise TypeError("Event type declarations must be strings")
364
- if not _EVENT_NAME_PATTERN.fullmatch(event):
365
- raise ValueError(
366
- f"Invalid event type '{event}' – must match lowercase_snake_case"
367
- )
368
- if event in _RESERVED_EVENT_NAMES:
369
- raise ValueError(
370
- f"Event type '{event}' is reserved by Socket.IO and cannot be used"
371
- )
372
- if event in seen:
373
- raise ValueError(f"Duplicate event type '{event}' declared in handler")
374
- seen.add(event)
375
- validated.append(event)
376
- if not validated:
377
- raise ValueError("Handlers must declare at least one event type")
378
- return validated
347
+ def validate_event_type(cls, event_type: str) -> str:
348
+ """Validate a runtime event name before dispatch."""
349
+
350
+ if not isinstance(event_type, str):
351
+ raise TypeError("Event type must be a string")
352
+ if not _EVENT_NAME_PATTERN.fullmatch(event_type):
353
+ raise ValueError(
354
+ f"Invalid event type '{event_type}' – must match lowercase_snake_case"
355
+ )
356
+ if event_type in _RESERVED_EVENT_NAMES:
357
+ raise ValueError(
358
+ f"Event type '{event_type}' is reserved by Socket.IO and cannot be used"
359
+ )
360
+ return event_type
361
362
@classmethod
363
def requires_auth(cls) -> bool:
helpers/websocket_manager.py
+63
-62
@@ -22,6 +22,17 @@ BUFFER_TTL = timedelta(hours=1)
22
_shared_websocket_manager: WebSocketManager | None = None
23
24
25
+async def send_data(
26
+ event_name: str,
27
+ data: dict[str, Any],
28
+ endpoint_name: str = "/webui",
29
+ connection_id: str | None = None,
30
+) -> None:
31
+ manager = get_shared_websocket_manager()
32
+ print(f"Sending data to {endpoint_name}/{event_name} with data {data}")
33
+ await manager.send_data(endpoint_name, event_name, data, connection_id)
34
+
35
+
36
def _utcnow() -> datetime:
37
return datetime.now(timezone.utc)
38
@@ -38,17 +49,6 @@ def get_shared_websocket_manager() -> "WebSocketManager":
49
return manager
50
51
41
-async def send_data(
42
- endpoint_name: str,
43
- event_name: str,
44
- data: dict[str, Any],
45
- connection_id: str | None = None,
46
-) -> None:
47
- manager = get_shared_websocket_manager()
48
- print(f"Sending data to {endpoint_name}/{event_name} with data {data}")
49
- await manager.send_data(endpoint_name, event_name, data, connection_id)
50
-
51
-
52
@dataclass
53
class BufferedEvent:
54
event_type: str
@@ -85,11 +85,11 @@ class WebSocketManager:
85
def __init__(self, socketio: socketio.AsyncServer, lock) -> None:
86
self.socketio = socketio
87
self.lock = lock
88
- self.handlers: defaultdict[str, defaultdict[str, List[WebSocketHandler]]] = defaultdict(
89
- lambda: defaultdict(list)
90
- )
88
+ self.handlers: defaultdict[str, List[WebSocketHandler]] = defaultdict(list)
89
self.connections: Dict[ConnectionIdentity, ConnectionInfo] = {}
92
- self.buffers: defaultdict[ConnectionIdentity, Deque[BufferedEvent]] = defaultdict(deque)
90
+ self.buffers: defaultdict[ConnectionIdentity, Deque[BufferedEvent]] = (
91
+ defaultdict(deque)
92
+ )
93
self._known_sids: Set[ConnectionIdentity] = set()
94
self._identifier: str = f"{self.__class__.__module__}.{self.__class__.__name__}"
95
# Session tracking (single-user default)
@@ -261,9 +261,7 @@ class WebSocketManager:
261
262
asyncio.create_task(_broadcast())
263
264
- def _normalize_handler_filter(
265
- self, value: Any, field_name: str
266
- ) -> Set[str] | None:
264
+ def _normalize_handler_filter(self, value: Any, field_name: str) -> Set[str] | None:
265
if value is None:
266
return None
267
if isinstance(value, str):
@@ -271,7 +269,9 @@ class WebSocketManager:
269
try:
270
iterator = iter(value)
271
except TypeError as exc: # pragma: no cover - defensive
274
- raise ValueError(f"{field_name} must be an array of handler identifiers") from exc
272
+ raise ValueError(
273
+ f"{field_name} must be an array of handler identifiers"
274
+ ) from exc
275
276
normalized: Set[str] = set()
277
for item in iterator:
@@ -282,9 +282,7 @@ class WebSocketManager:
282
normalized.add(item)
283
return normalized
284
285
- def _normalize_sid_filter(
286
- self, value: str | Iterable[str] | None
287
- ) -> Set[str]:
285
+ def _normalize_sid_filter(self, value: str | Iterable[str] | None) -> Set[str]:
286
if value is None:
287
return set()
288
if isinstance(value, str):
@@ -297,12 +295,11 @@ class WebSocketManager:
295
def _select_handlers(
296
self,
297
namespace: str,
300
- event_type: str,
298
*,
299
include: Set[str] | None,
300
exclude: Set[str] | None,
301
) -> tuple[list[WebSocketHandler], Set[str]]:
305
- registered = self.handlers.get(namespace, {}).get(event_type, [])
302
+ registered = self.handlers.get(namespace, [])
303
available_ids = {handler.identifier for handler in registered}
304
305
if include is not None:
@@ -346,33 +343,23 @@ class WebSocketManager:
343
for namespace, handlers in handlers_by_namespace.items():
344
for handler in handlers:
345
handler.bind_manager(self, namespace=namespace)
349
- declared = handler.get_event_types()
350
- try:
351
- validated_events = handler.validate_event_types(declared)
352
- except Exception as exc:
353
- PrintStyle.error(
354
- f"Failed to register handler {handler.identifier}: {exc}"
355
- )
356
- raise
357
-
346
if _ws_debug_enabled():
347
PrintStyle.info(
360
- "Registered WebSocket handler %s namespace=%s for events: %s"
361
- % (handler.identifier, namespace, ", ".join(validated_events))
348
+ "Registered WebSocket handler %s namespace=%s"
349
+ % (handler.identifier, namespace)
350
)
363
- for event_type in validated_events:
364
- existing = self.handlers[namespace].get(event_type)
365
- if existing:
366
- PrintStyle.warning(
367
- f"Duplicate handler registration for namespace '{namespace}' event '{event_type}'"
368
- )
369
- self.handlers[namespace][event_type].append(handler)
370
- self._debug(
371
- f"Registered handler {handler.identifier} namespace={namespace} event='{event_type}'"
351
+ existing = self.handlers.get(namespace, [])
352
+ if handler in existing:
353
+ PrintStyle.warning(
354
+ f"Duplicate handler registration for namespace '{namespace}'"
355
)
356
+ self.handlers[namespace].append(handler)
357
+ self._debug(
358
+ f"Registered handler {handler.identifier} namespace={namespace}"
359
+ )
360
361
def iter_event_types(self, namespace: str) -> Iterable[str]:
375
- return list(self.handlers.get(namespace, {}).keys())
362
+ return []
363
364
def iter_namespaces(self) -> list[str]:
365
return list(self.handlers.keys())
@@ -592,13 +579,26 @@ class WebSocketManager:
579
include = include_handlers or include_meta
580
exclude = exclude_handlers or (exclude_meta if allow_exclude else None)
581
595
- registered = self.handlers.get(namespace, {}).get(event_type, [])
582
+ try:
583
+ WebSocketHandler.validate_event_type(event_type)
584
+ except (TypeError, ValueError) as exc:
585
+ error = self._build_error_result(
586
+ handler_id=handler_id or self._identifier,
587
+ code="INVALID_EVENT",
588
+ message=str(exc),
589
+ correlation_id=correlation_id,
590
+ )
591
+ if ack:
592
+ ack({"correlationId": correlation_id, "results": [error]})
593
+ return {"correlationId": correlation_id, "results": [error]}
594
+
595
+ registered = self.handlers.get(namespace, [])
596
if not registered:
597
- PrintStyle.warning(f"No handlers registered for event '{event_type}'")
597
+ PrintStyle.warning(f"No handlers registered for namespace '{namespace}'")
598
error = self._build_error_result(
599
handler_id=handler_id or self._identifier,
600
code="NO_HANDLERS",
601
- message=f"No handler for namespace '{namespace}' event '{event_type}'",
601
+ message=f"No handler for namespace '{namespace}'",
602
correlation_id=correlation_id,
603
)
604
if ack:
@@ -607,7 +607,7 @@ class WebSocketManager:
607
608
try:
609
selected_handlers, _ = self._select_handlers(
610
- namespace, event_type, include=include, exclude=exclude
610
+ namespace, include=include, exclude=exclude
611
)
612
except ValueError as exc:
613
error = self._build_error_result(
@@ -862,7 +862,9 @@ class WebSocketManager:
862
863
try:
864
task = asyncio.create_task(_dispatch())
865
- return await asyncio.wait_for(asyncio.shield(task), timeout=timeout_seconds)
865
+ return await asyncio.wait_for(
866
+ asyncio.shield(task), timeout=timeout_seconds
867
+ )
868
except asyncio.TimeoutError:
869
PrintStyle.warning(
870
f"requestAll timeout for sid {target_sid} correlation={correlation_id}"
@@ -885,9 +887,7 @@ class WebSocketManager:
887
],
888
}
889
888
- tasks = {
889
- sid: asyncio.create_task(_invoke_for_sid(sid)) for sid in active_sids
890
- }
890
+ tasks = {sid: asyncio.create_task(_invoke_for_sid(sid)) for sid in active_sids}
891
892
aggregated: list[dict[str, Any]] = []
893
for sid, task in tasks.items():
@@ -1064,15 +1064,16 @@ class WebSocketManager:
1064
}
1065
)
1066
1067
- async def _run_lifecycle(self, namespace: str, fn: Callable[[WebSocketHandler], Any]) -> None:
1067
+ async def _run_lifecycle(
1068
+ self, namespace: str, fn: Callable[[WebSocketHandler], Any]
1069
+ ) -> None:
1070
seen: Set[WebSocketHandler] = set()
1071
coros: list[Any] = []
1070
- for handler_list in self.handlers.get(namespace, {}).values():
1071
- for handler in handler_list:
1072
- if handler in seen:
1073
- continue
1074
- seen.add(handler)
1075
- coros.append(self._get_handler_worker().execute_inside(fn, handler))
1072
+ for handler in self.handlers.get(namespace, []):
1073
+ if handler in seen:
1074
+ continue
1075
+ seen.add(handler)
1076
+ coros.append(self._get_handler_worker().execute_inside(fn, handler))
1077
if coros:
1078
await asyncio.gather(*coros, return_exceptions=True)
1079
@@ -1175,12 +1176,12 @@ class WebSocketManager:
1176
"""Return SIDs for a user; single-user default returns all active SIDs."""
1177
with self.lock:
1178
bucket = self._ALL_USERS_BUCKET if user is None else user
1178
- return list(self.user_to_sids.get(bucket, set())) # type: ignore
1179
+ return list(self.user_to_sids.get(bucket, set())) # type: ignore
1180
1181
def get_user_for_sid(self, sid: str) -> str | None:
1182
"""Return user identifier for a SID or None."""
1183
with self.lock:
1183
- return self.sid_to_user.get(sid) # type: ignore
1184
+ return self.sid_to_user.get(sid) # type: ignore
1185
1186
def set_server_restart_broadcast(self, enabled: bool) -> None:
1187
"""Enable or disable automatic server restart broadcasts."""
python/websocket_handlers/_default.py
-5
@@ -20,11 +20,6 @@ class RootDefaultHandler(WebSocketHandler):
20
def requires_csrf(cls) -> bool:
21
return False
22
23
- @classmethod
24
- def get_event_types(cls) -> list[str]:
25
- # Diagnostics-only noop endpoint.
26
- return ["ws_root_echo"]
27
-
23
async def process_event(
24
self, event_type: str, data: dict[str, Any], sid: str
25
) -> dict[str, Any] | WebSocketResult | None:
python/websocket_handlers/dev_websocket_test_handler.py
-13
@@ -11,19 +11,6 @@ from helpers.websocket import WebSocketHandler, WebSocketResult
11
class DevWebsocketTestHandler(WebSocketHandler):
12
"""Test harness handler powering the developer WebSocket validation component."""
13
14
- @classmethod
15
- def get_event_types(cls) -> list[str]:
16
- return [
17
- "ws_tester_emit",
18
- "ws_tester_request",
19
- "ws_tester_request_delayed",
20
- "ws_tester_trigger_persistence",
21
- "ws_tester_request_all",
22
- "ws_tester_broadcast_demo_trigger",
23
- "ws_event_console_subscribe",
24
- "ws_event_console_unsubscribe",
25
- ]
26
-
14
async def process_event(
15
self, event_type: str, data: Dict[str, Any], sid: str
16
) -> dict[str, Any] | WebSocketResult | None:
python/websocket_handlers/hello_handler.py
-4
@@ -7,10 +7,6 @@ from helpers.websocket import WebSocketHandler
7
class HelloHandler(WebSocketHandler):
8
"""Sample handler used for foundational testing."""
9
10
- @classmethod
11
- def get_event_types(cls) -> list[str]:
12
- return ["hello_request"]
13
-
10
async def process_event(self, event_type: str, data: dict, sid: str):
11
name = data.get("name") or "stranger"
12
PrintStyle.info(f"hello_request from {sid} ({name})")
python/websocket_handlers/state_sync_handler.py
deleted
-76
@@ -1,76 +0,0 @@
1
-from __future__ import annotations
2
-
3
-from helpers import runtime
4
-from helpers.print_style import PrintStyle
5
-from helpers.websocket import WebSocketHandler, WebSocketResult
6
-from helpers.state_monitor import get_state_monitor, _ws_debug_enabled
7
-from helpers.state_snapshot import (
8
- StateRequestValidationError,
9
- parse_state_request_payload,
10
-)
11
-
12
-
13
-class StateSyncHandler(WebSocketHandler):
14
- @classmethod
15
- def get_event_types(cls) -> list[str]:
16
- return ["state_request"]
17
-
18
- async def on_connect(self, sid: str) -> None:
19
- monitor = get_state_monitor()
20
- monitor.bind_manager(self.manager, handler_id=self.identifier)
21
- monitor.register_sid(self.namespace, sid)
22
- if _ws_debug_enabled():
23
- PrintStyle.debug(f"[StateSyncHandler] connect sid={sid}")
24
-
25
- async def on_disconnect(self, sid: str) -> None:
26
- get_state_monitor().unregister_sid(self.namespace, sid)
27
- if _ws_debug_enabled():
28
- PrintStyle.debug(f"[StateSyncHandler] disconnect sid={sid}")
29
-
30
- async def process_event(self, event_type: str, data: dict, sid: str) -> dict | WebSocketResult | None:
31
- correlation_id = data.get("correlationId")
32
- try:
33
- request = parse_state_request_payload(data)
34
- except StateRequestValidationError as exc:
35
- PrintStyle.warning(
36
- f"[StateSyncHandler] INVALID_REQUEST sid={sid} reason={exc.reason} details={exc.details!r}"
37
- )
38
- return self.result_error(
39
- code="INVALID_REQUEST",
40
- message=str(exc),
41
- correlation_id=correlation_id,
42
- )
43
-
44
- if _ws_debug_enabled():
45
- PrintStyle.debug(
46
- f"[StateSyncHandler] state_request sid={sid} context={request.context!r} "
47
- f"log_from={request.log_from} notifications_from={request.notifications_from} timezone={request.timezone!r} "
48
- f"correlation_id={correlation_id}"
49
- )
50
-
51
- # Baseline sequence must be reset on every state_request (new sync period).
52
- # V1 policy: seq_base starts >0 to allow simple gating checks.
53
- seq_base = 1
54
- monitor = get_state_monitor()
55
- monitor.update_projection(
56
- self.namespace,
57
- sid,
58
- request=request,
59
- seq_base=seq_base,
60
- )
61
- # INVARIANT.STATE.INITIAL_SNAPSHOT: schedule a full snapshot quickly after handshake.
62
- monitor.mark_dirty(
63
- self.namespace,
64
- sid,
65
- reason="state_sync_handler.StateSyncHandler.state_request",
66
- )
67
- if _ws_debug_enabled():
68
- PrintStyle.debug(f"[StateSyncHandler] state_request accepted sid={sid} seq_base={seq_base}")
69
-
70
- return self.result_ok(
71
- {
72
- "runtime_epoch": runtime.get_runtime_id(),
73
- "seq_base": seq_base,
74
- },
75
- correlation_id=correlation_id,
76
- )
python/websocket_handlers/webui_handler.py
new
+34
@@ -0,0 +1,34 @@
1
+from helpers.websocket import WebSocketHandler, WebSocketResult
2
+from helpers import extension
3
+
4
+
5
+class WebuiHandler(WebSocketHandler):
6
+ async def on_connect(self, sid: str) -> None:
7
+ await extension.call_extensions_async(
8
+ "webui_ws_connect", agent=None, instance=self, sid=sid
9
+ )
10
+
11
+ async def on_disconnect(self, sid: str) -> None:
12
+ await extension.call_extensions_async(
13
+ "webui_ws_disconnect", agent=None, instance=self, sid=sid
14
+ )
15
+
16
+ async def process_event(
17
+ self, event_type: str, data: dict, sid: str
18
+ ) -> dict | WebSocketResult | None:
19
+ response_data: dict = {}
20
+
21
+ await extension.call_extensions_async(
22
+ "webui_ws_event",
23
+ agent=None,
24
+ instance=self,
25
+ sid=sid,
26
+ event_type=event_type,
27
+ data=data,
28
+ response_data=response_data,
29
+ )
30
+
31
+ return self.result_ok(
32
+ response_data,
33
+ correlation_id=data.get("correlationId"),
34
+ )
run_ui.py
-16
@@ -336,22 +336,6 @@ def configure_websocket_namespaces(
336
async def _disconnect(sid, _namespace: str = namespace): # type: ignore[override]
337
await websocket_manager.handle_disconnect(_namespace, sid)
338
339
- def _register_socketio_event(event_type: str) -> None:
340
- @socketio_server.on(event_type, namespace=namespace)
341
- async def _event_handler(
342
- sid,
343
- data,
344
- _event_type: str = event_type,
345
- _namespace: str = namespace,
346
- ):
347
- payload = data or {}
348
- return await websocket_manager.route_event(
349
- _namespace, _event_type, payload, sid
350
- )
351
-
352
- for _event_type in websocket_manager.iter_event_types(namespace):
353
- _register_socketio_event(_event_type)
354
-
339
@socketio_server.on("*", namespace=namespace)
340
async def _catch_all(event, sid, data, _namespace: str = namespace):
341
payload = data or {}
tests/test_multi_tab_isolation.py
+2
-2
@@ -18,7 +18,7 @@ async def test_state_monitor_per_sid_isolation_independent_snapshots_seq_and_cur
18
snapshot_calls: list[dict[str, object]] = []
19
emitted: list[dict[str, object]] = []
20
21
- namespace = "/state_sync"
21
+ namespace = "/webui"
22
23
async def fake_build_snapshot_from_request(*, request):
24
context = request.context
@@ -133,7 +133,7 @@ async def test_state_monitor_mark_dirty_for_context_scopes_to_active_context():
133
from helpers.state_snapshot import StateRequestV1
134
135
monitor = StateMonitor(debounce_seconds=60.0)
136
- namespace = "/state_sync"
136
+ namespace = "/webui"
137
monitor.register_sid(namespace, "sid-a")
138
monitor.register_sid(namespace, "sid-b")
139
tests/test_state_monitor.py
+1
-1
@@ -13,7 +13,7 @@ async def test_state_monitor_debounce_coalesces_without_postponing_and_cleanup_c
13
from helpers.state_monitor import StateMonitor
14
from helpers.state_snapshot import StateRequestV1
15
16
- namespace = "/state_sync"
16
+ namespace = "/webui"
17
monitor = StateMonitor(debounce_seconds=10.0)
18
monitor.register_sid(namespace, "sid-1")
19
monitor.bind_manager(type("FakeManager", (), {"_dispatcher_loop": None})())
tests/test_state_sync_handler.py
+7
-7
@@ -12,7 +12,7 @@ if str(PROJECT_ROOT) not in sys.path:
12
13
from helpers.websocket_manager import WebSocketManager
14
15
-NAMESPACE = "/state_sync"
15
+NAMESPACE = "/webui"
16
17
18
class FakeSocketIOServer:
@@ -27,12 +27,12 @@ async def _create_manager() -> WebSocketManager:
27
socketio = FakeSocketIOServer()
28
manager = WebSocketManager(socketio, threading.RLock())
29
30
- from python.websocket_handlers.state_sync_handler import StateSyncHandler
30
+ from python.websocket_handlers.webui_handler import WebuiHandler
31
from helpers.state_monitor import _reset_state_monitor_for_testing
32
33
_reset_state_monitor_for_testing()
34
- StateSyncHandler._reset_instance_for_testing()
35
- handler = StateSyncHandler.get_instance(socketio, threading.RLock())
34
+ WebuiHandler._reset_instance_for_testing()
35
+ handler = WebuiHandler.get_instance(socketio, threading.RLock())
36
manager.register_handlers({NAMESPACE: [handler]})
37
await manager.handle_connect(NAMESPACE, "sid-1")
38
return manager
@@ -42,12 +42,12 @@ async def _create_manager_with_socketio() -> tuple[WebSocketManager, FakeSocketI
42
socketio = FakeSocketIOServer()
43
manager = WebSocketManager(socketio, threading.RLock())
44
45
- from python.websocket_handlers.state_sync_handler import StateSyncHandler
45
+ from python.websocket_handlers.webui_handler import WebuiHandler
46
from helpers.state_monitor import _reset_state_monitor_for_testing
47
48
_reset_state_monitor_for_testing()
49
- StateSyncHandler._reset_instance_for_testing()
50
- handler = StateSyncHandler.get_instance(socketio, threading.RLock())
49
+ WebuiHandler._reset_instance_for_testing()
50
+ handler = WebuiHandler.get_instance(socketio, threading.RLock())
51
manager.register_handlers({NAMESPACE: [handler]})
52
await manager.handle_connect(NAMESPACE, "sid-1")
53
return manager, socketio
tests/test_state_sync_welcome_screen.py
+4
-4
@@ -12,7 +12,7 @@ if str(PROJECT_ROOT) not in sys.path:
12
13
from helpers.websocket_manager import WebSocketManager
14
15
-NAMESPACE = "/state_sync"
15
+NAMESPACE = "/webui"
16
17
18
class FakeSocketIOServer:
@@ -32,14 +32,14 @@ async def test_state_sync_handshake_and_initial_snapshot_work_with_no_selected_c
32
33
from helpers.state_snapshot import validate_snapshot_schema_v1
34
from helpers.state_monitor import _reset_state_monitor_for_testing
35
- from python.websocket_handlers.state_sync_handler import StateSyncHandler
35
+ from python.websocket_handlers.webui_handler import WebuiHandler
36
37
socketio = FakeSocketIOServer()
38
manager = WebSocketManager(socketio, threading.RLock())
39
40
_reset_state_monitor_for_testing()
41
- StateSyncHandler._reset_instance_for_testing()
42
- handler = StateSyncHandler.get_instance(socketio, threading.RLock())
41
+ WebuiHandler._reset_instance_for_testing()
42
+ handler = WebuiHandler.get_instance(socketio, threading.RLock())
43
manager.register_handlers({NAMESPACE: [handler]})
44
await manager.handle_connect(NAMESPACE, "sid-1")
45
tests/test_websocket_handlers.py
+4
-4
@@ -141,17 +141,17 @@ def test_get_instance_returns_singleton():
141
@pytest.mark.asyncio
142
async def test_state_sync_handler_registers_and_routes_state_request():
143
from helpers.websocket_manager import WebSocketManager
144
- from python.websocket_handlers.state_sync_handler import StateSyncHandler
144
+ from python.websocket_handlers.webui_handler import WebuiHandler
145
from helpers.state_monitor import _reset_state_monitor_for_testing
146
147
_reset_state_monitor_for_testing()
148
- StateSyncHandler._reset_instance_for_testing()
148
+ WebuiHandler._reset_instance_for_testing()
149
150
socketio = _FakeSocketIO()
151
lock = threading.RLock()
152
manager = WebSocketManager(socketio, lock)
153
- handler = StateSyncHandler.get_instance(socketio, lock)
154
- namespace = "/state_sync"
153
+ handler = WebuiHandler.get_instance(socketio, lock)
154
+ namespace = "/webui"
155
manager.register_handlers({namespace: [handler]})
156
await manager.handle_connect(namespace, "sid-1")
157
tests/test_websocket_namespaces.py
+16
-16
@@ -102,7 +102,7 @@ async def test_namespace_isolation_state_sync_vs_dev_websocket_test() -> None:
102
"""
103
CONTRACT.INVARIANT.NS.ISOLATION: no cross-namespace delivery for application events.
104
105
- Acceptance proof for `/state_sync` vs `/dev_websocket_test` namespaces.
105
+ Acceptance proof for `/webui` vs `/dev_websocket_test` namespaces.
106
"""
107
108
from flask import Flask
@@ -161,7 +161,7 @@ async def test_namespace_isolation_state_sync_vs_dev_websocket_test() -> None:
161
socketio_server=sio,
162
websocket_manager=manager,
163
handlers_by_namespace={
164
- "/state_sync": [StateHandler.get_instance(sio, lock)],
164
+ "/webui": [StateHandler.get_instance(sio, lock)],
165
"/dev_websocket_test": [DevHandler.get_instance(sio, lock)],
166
},
167
)
@@ -188,19 +188,19 @@ async def test_namespace_isolation_state_sync_vs_dev_websocket_test() -> None:
188
async def _on_tester_broadcast_state(_payload: Any) -> None:
189
tester_broadcast_state.set()
190
191
- client.on("state_push", _on_state_push_state, namespace="/state_sync")
191
+ client.on("state_push", _on_state_push_state, namespace="/webui")
192
client.on("state_push", _on_state_push_dev, namespace="/dev_websocket_test")
193
client.on("ws_tester_broadcast", _on_tester_broadcast_dev, namespace="/dev_websocket_test")
194
- client.on("ws_tester_broadcast", _on_tester_broadcast_state, namespace="/state_sync")
194
+ client.on("ws_tester_broadcast", _on_tester_broadcast_state, namespace="/webui")
195
196
await client.connect(
197
base_url,
198
- namespaces=["/state_sync", "/dev_websocket_test"],
198
+ namespaces=["/webui", "/dev_websocket_test"],
199
headers={"Origin": base_url},
200
wait_timeout=2,
201
)
202
try:
203
- await client.call("state_request", {"context": None}, namespace="/state_sync", timeout=2)
203
+ await client.call("state_request", {"context": None}, namespace="/webui", timeout=2)
204
await asyncio.wait_for(state_push_state.wait(), timeout=2)
205
await asyncio.sleep(0.05)
206
assert state_push_dev.is_set() is False
@@ -237,7 +237,7 @@ async def test_diagnostics_include_source_namespace_and_deliver_on_dev_namespace
237
manager = WebSocketManager(socketio, threading.RLock())
238
manager._schedule_lifecycle_broadcast = lambda *_args, **_kwargs: None # type: ignore[assignment]
239
240
- ns_state = "/state_sync"
240
+ ns_state = "/webui"
241
ns_dev = "/dev_websocket_test"
242
243
handler = DummyHandler.get_instance(socketio, threading.RLock())
@@ -275,13 +275,13 @@ def test_namespace_discovery_maps_core_handlers_to_expected_namespaces() -> None
275
)
276
by_namespace = {entry.namespace: entry for entry in discoveries}
277
278
- assert "/state_sync" in by_namespace
278
+ assert "/webui" in by_namespace
279
assert "/dev_websocket_test" in by_namespace
280
281
- state_cls_names = [cls.__name__ for cls in by_namespace["/state_sync"].handler_classes]
281
+ state_cls_names = [cls.__name__ for cls in by_namespace["/webui"].handler_classes]
282
dev_cls_names = [cls.__name__ for cls in by_namespace["/dev_websocket_test"].handler_classes]
283
284
- assert state_cls_names == ["StateSyncHandler"]
284
+ assert state_cls_names == ["WebuiHandler"]
285
assert dev_cls_names == ["DevWebsocketTestHandler"]
286
287
@@ -290,15 +290,15 @@ def test_run_ui_builds_namespace_handler_map_without_cross_registration() -> Non
290
291
handlers_by_namespace = _build_websocket_handlers_by_namespace(object(), threading.RLock())
292
293
- assert "/state_sync" in handlers_by_namespace
293
+ assert "/webui" in handlers_by_namespace
294
assert "/dev_websocket_test" in handlers_by_namespace
295
296
assert all(
297
handler.__class__.__name__ != "DevWebsocketTestHandler"
298
- for handler in handlers_by_namespace["/state_sync"]
298
+ for handler in handlers_by_namespace["/webui"]
299
)
300
assert all(
301
- handler.__class__.__name__ != "StateSyncHandler"
301
+ handler.__class__.__name__ != "WebuiHandler"
302
for handler in handlers_by_namespace["/dev_websocket_test"]
303
)
304
@@ -315,7 +315,7 @@ async def test_route_event_dispatches_only_within_connected_namespace_and_result
315
manager = WebSocketManager(socketio, threading.RLock())
316
manager._schedule_lifecycle_broadcast = lambda *_args, **_kwargs: None # type: ignore[assignment]
317
318
- ns_state = "/state_sync"
318
+ ns_state = "/webui"
319
ns_dev = "/dev_websocket_test"
320
321
calls: list[str] = []
@@ -373,7 +373,7 @@ async def test_lifecycle_broadcasts_deliver_only_within_the_namespace() -> None:
373
socketio = FakeSocketIOServer()
374
manager = WebSocketManager(socketio, threading.RLock())
375
376
- ns_state = "/state_sync"
376
+ ns_state = "/webui"
377
ns_dev = "/dev_websocket_test"
378
379
# Connect events should broadcast only within their namespace.
@@ -428,7 +428,7 @@ async def test_request_semantics_no_handlers_and_timeouts_are_namespace_scoped_a
428
manager = WebSocketManager(socketio, threading.RLock())
429
manager._schedule_lifecycle_broadcast = lambda *_args, **_kwargs: None # type: ignore[assignment]
430
431
- ns_state = "/state_sync"
431
+ ns_state = "/webui"
432
ns_dev = "/dev_websocket_test"
433
434
class Alpha(WebSocketHandler):
tests/test_websocket_root_namespace.py
+2
-2
@@ -88,7 +88,7 @@ async def test_root_namespace_request_style_calls_resolve_with_no_handlers() ->
88
socketio_server=sio,
89
websocket_manager=manager,
90
handlers_by_namespace={
91
- "/state_sync": [HelloHandler.get_instance(sio, lock)],
91
+ "/webui": [HelloHandler.get_instance(sio, lock)],
92
},
93
)
94
@@ -161,7 +161,7 @@ async def test_root_namespace_fire_and_forget_does_not_invoke_application_handle
161
socketio_server=sio,
162
websocket_manager=manager,
163
handlers_by_namespace={
164
- "/state_sync": [SideEffectHandler.get_instance(sio, lock)],
164
+ "/webui": [SideEffectHandler.get_instance(sio, lock)],
165
},
166
)
167
webui/components/settings/developer/websocket-test-store.js
+1
-1
@@ -12,7 +12,7 @@ const MAX_PAYLOAD_BYTES = 50 * 1024 * 1024;
12
const TOAST_DURATION = 5;
13
14
const websocket = getNamespacedClient("/dev_websocket_test");
15
-const stateSocket = getNamespacedClient("/state_sync");
15
+const stateSocket = getNamespacedClient("/webui");
16
17
function now() {
18
return new Date().toISOString();
webui/components/sync/sync-store.js
+3
-3
@@ -6,7 +6,7 @@ import { store as chatTopStore } from "/components/chat/top-section/chat-top-sto
6
import { store as notificationStore } from "/components/notifications/notification-store.js";
7
import * as Extensions from "/js/extensions.js"
8
9
-const stateSocket = getNamespacedClient("/state_sync");
9
+const stateSocket = getNamespacedClient("/webui");
10
11
const SYNC_MODES = {
12
DISCONNECTED: "DISCONNECTED",
@@ -317,7 +317,7 @@ const model = {
317
318
// handle all requests with extensions
319
await stateSocket.on("*", (eventType, envelope) => {
320
- console.log(`[syncStore] *${eventType} received`);
320
+ // console.log(`[syncStore] *${eventType} received`);
321
this.handleEvent(eventType, envelope)
322
});
323
@@ -330,7 +330,7 @@ const model = {
330
},
331
332
async handleEvent(eventType, envelope){
333
- await Extensions.callJsExtensions("ws_sync_push", eventType, envelope);
333
+ await Extensions.callJsExtensions("webui_ws_push", eventType, envelope);
334
},
335
336
async sendStateRequest(options = {}) {