Activity#

Tracking what the middleware is doing while it does it.

Live activity feed for the middleware machinery.

This module turns the middleware internal logging stream into a live feed for browsers. It is library-level and mode-agnostic. A controller and a monitor are the same library with different configuration. They inherit this feed from the same code with no role-specific branching. That is the architectural claim the demo exists to make.

What it shows: wiring decisions, message traffic, registrations, heartbeats.

What it does not show: resource state. That is the monitor job. Keeping the two apart is a decided boundary. Mixing them would conflate the middleware health with the device status. Debugging would become ambiguous.

Why the stream is polled rather than pushed: Log records arrive from two places. The event loop sends records. Worker threads send records (this codebase calls anyio.to_thread.run_sync for every graph write). Pushing from the logging handler into an asyncio.Queue is not thread-safe. It would require capture of the running loop. This breaks when no loop runs (construction time, tests). Instead, the handler appends to a thread-safe collections.deque. The SSE endpoint polls that deque on a short sleep. This decouples the logging path (synchronous, multi-threaded) from the serving path (asynchronous, single-threaded). No loop introspection or call_soon_threadsafe is required.

class kapps_semantic_middleware.activity.ActivityRecord(seq: int, timestamp: float, level: str, logger: str, message: str)[source]#

Bases: object

One entry in the activity feed.

Immutable so that records handed to the SSE generator cannot receive mutation by later logging activity. The seq field allows clients to request only what is new.

class kapps_semantic_middleware.activity.ActivityFeed(capacity: int = 200)[source]#

Bases: object

Thread-safe buffer for activity records.

Holds a rolling window of records. Appends receive protection by a lock because the logging handler runs in arbitrary threads. Reads receive protection for the same reason. The SSE generator holds the lock only long enough to copy data out.

append(*, timestamp: float, level: str, logger: str, message: str) → ActivityRecord[source]#

Buffer one record. Assign its sequence number. Drop the oldest when full.

Numbering and insertion happen under one lock, together. That is the whole point of this method that takes fields rather than a finished record. Split them. Take a number in one critical section. Insert in another. Two threads interleave. The deque becomes ordered [.., 6, 5]. last_seq then reports 5. The stream re-sends record 6 on every poll, forever. The race is reachable here rather than theoretical. This codebase logs from worker threads (anyio.to_thread.run_sync wraps every graph write) and from the event loop.

since(seq: int) → List[ActivityRecord][source]#

Return records with .seq > seq, oldest first.

Called by the SSE generator. The lock is held only for the duration of the slice operation, not during the yield. The event loop is not blocked.

snapshot() → List[ActivityRecord][source]#

Return everything currently buffered.

Used to seed a client that connects mid-run. Same locking discipline as since.

property last_seq: int#

The highest sequence number currently buffered. Zero if empty.

Not a resume token. Eviction begins. This is still the newest record. The oldest moves forward under it. A consumer stores this value. The consumer goes away. The consumer comes back. It asks for everything after that value. It has no way to learn what fell out of the window meanwhile. Pair it with oldest_seq to detect that. The stream does.

property oldest_seq: int#

The lowest sequence number still buffered. Zero if empty.

A consumer expected n. It finds this value greater than n. The feed outran it. The records between are gone. Exposed so that loss is reportable rather than invisible. A feed that quietly drops lines is worse than one that says it dropped them. The viewer cannot tell “nothing happened” from “I missed it”.

kapps_semantic_middleware.activity.enable_activity_feed(middleware: Any, *, capacity: int = 200, level: int = 20, logger_name: str = 'kapps_semantic_middleware') → ActivityFeed[source]#

Enable the activity feed on a middleware instance.

Builds an ActivityFeed. Attaches a logging handler to capture records. Mounts the HTTP routes on middleware.app. A second call on the same middleware returns the existing feed without attach of a second handler.

The feed is stored on the middleware as middleware.activity_feed. Other components (tests, diagnostics) can inspect it directly.