Falsify / falsify /events.py
Aaryan Kumar
deploy to hugging face
1605cbb
Raw
History Blame Contribute Delete
3.57 kB
"""
FALSIFY real-time event bus — the seam between the belief pipeline and the live UI.
The whole point of the web demo is that a judge *watches* a belief die: the refute
flash, the cascade sweep, the forget-dissolve. Those animations are driven by events
that this module fans out to every connected browser over Server-Sent Events.
Design constraints that shaped this module:
* **Leaf module.** It imports nothing from ``falsify`` so that ``graph_ops`` can import
it without a cycle (``graph_ops -> events`` only, never the reverse).
* **No-op when nobody's watching.** ``emit`` iterates an empty subscriber set under the
CLI (``python main.py``), so the two hooks in ``graph_ops`` cost effectively nothing
and the offline demo behaves exactly as before.
* **Never block the pipeline.** A slow browser must not stall belief revision, so we
``put_nowait`` and drop on a full queue rather than awaiting backpressure. The client
re-syncs via ``GET /api/graph`` on reconnect, so a dropped frame is cosmetic.
Each open ``GET /api/events`` connection owns one queue (``subscribe`` / ``unsubscribe``);
``server.py`` drains it and serializes each event as an SSE ``data:`` frame.
"""
from __future__ import annotations
import logging
from asyncio import Queue, QueueFull
from typing import Any, Dict, Set
logger = logging.getLogger("falsify.events")
# One queue per live SSE connection. Empty set == CLI / offline == emit is a no-op.
_subscribers: Set[Queue] = set()
# Per-connection buffer. Large enough that a whole cascade never overflows a healthy
# client; a client slow enough to fill 1000 frames is already gone.
_QUEUE_MAXSIZE = 1000
def subscribe() -> Queue:
"""Register a new SSE connection and return its private event queue."""
q: Queue = Queue(maxsize=_QUEUE_MAXSIZE)
_subscribers.add(q)
logger.debug("SSE subscriber added (now %d)", len(_subscribers))
return q
def unsubscribe(q: Queue) -> None:
"""Drop a disconnected SSE connection's queue (idempotent)."""
_subscribers.discard(q)
logger.debug("SSE subscriber removed (now %d)", len(_subscribers))
def has_subscribers() -> bool:
"""True if at least one browser is listening — lets the server pace animations
(small inter-step sleeps) only when someone is actually watching."""
return bool(_subscribers)
async def emit(event: Dict[str, Any]) -> None:
"""Fan one event out to every live connection. No-op when none. Never blocks."""
for q in list(_subscribers):
try:
q.put_nowait(event)
except QueueFull: # slow client: drop the frame, don't stall the pipeline
logger.debug("dropping event for a full subscriber queue")
async def emit_state_change(node_id: str, state: str, epoch: int) -> None:
"""A node's truth-state changed (refuted / invalidated / superseded / alive)."""
await emit({"type": "node_state_changed", "id": str(node_id),
"state": str(state), "epoch": int(epoch)})
async def emit_forgotten(node_id: str) -> None:
"""A node was hard-deleted from the graph (the forget-dissolve animation)."""
await emit({"type": "node_forgotten", "id": str(node_id)})
async def emit_graph_reset() -> None:
"""The graph was rebuilt/reseeded — clients should refetch the full snapshot."""
await emit({"type": "graph_reset"})
async def emit_step(step: str, detail: str = "") -> None:
"""A human-readable pipeline milestone, for the revision-log feed."""
await emit({"type": "pipeline_step", "step": str(step), "detail": str(detail)})