File size: 3,569 Bytes
1605cbb
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
"""
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)})