spacr.flowview.collector¶
Thread-safe bounded collection and event folding for FlowView.
Classes¶
Own the event queue and the only mutable copy of a run graph. |
Module Contents¶
- class spacr.flowview.collector.Collector(graph: spacr.flowview.model.RunGraph, *, max_queue_size: int = 2000)[source]¶
Own the event queue and the only mutable copy of a run graph.
Producers call
emit(), which never waits for queue capacity but may briefly contend for the queue lock. A saturated queue drops its oldest event and records that the resulting display is sampled. Renderers consume recursively detached graph snapshots rather than sharing the collector’s dictionaries.- Parameters:
graph – run graph to copy recursively; later changes to the caller’s node parameter dictionaries cannot alter collector state.
max_queue_size – positive integer count of events that may wait before the oldest is dropped. It is a bound on memory, not on throughput: a producer never waits for capacity, so this is the number of events a slow renderer may fall behind by before the display becomes sampled rather than complete.
- Raises:
ValueError – when
max_queue_sizeis not a positive integer.
Copy one graph and initialize its bounded event stream.
- Parameters:
graph – run graph whose nodes, metrics, parameters, and edge container are detached from caller-owned state.
max_queue_size – positive integer maximum of pending events retained before the oldest is discarded and
sampledbecomes true.
- Raises:
ValueError – if
max_queue_sizeis not a positive integer.
Recursive-copy errors from caller-provided node payloads propagate.
- clear_sampled() None[source]¶
Clear the sticky saturation indicator without draining the queue.
- Returns:
None.
- drain(limit: int | None = None) int[source]¶
Consume queued events in order and fold one batch.
- Parameters:
limit – maximum events to consume, all pending events when
None; a negative integer behaves as zero.- Returns:
number consumed, including unrecognized objects that entered the queue through a dynamically typed caller.
- Raises:
ValueError – if a non-
Nonelimit is not an integer or is a boolean.
- emit(event: spacr.flowview.events.FlowEvent) bool[source]¶
Queue an event without waiting for capacity.
- Parameters:
event – FlowView event to retain for the next drain; node-bearing events are recursively detached before queueing.
- Returns:
true when no event was discarded; false when the oldest queued event was replaced and
sampledwas set.
- fold(event: object) bool[source]¶
Fold one event immediately.
- Parameters:
event – candidate event to apply under the state lock; node-bearing events are recursively detached first.
- Returns:
true for a recognized event whose update was accepted; duplicate node and edge declarations remain recognized even when graph content does not change.
- snapshot() spacr.flowview.model.RunGraph[source]¶
Return a recursively detached, renderer-safe graph snapshot.
- Returns:
graph whose containers, node metrics, and node parameters can be mutated without changing collector state.
- property revision: int[source]¶
Monotonic graph revision for renderers that can skip idle work.
It advances once per recognized immediate event and once per drained batch containing at least one recognized event.
- property sampled: bool[source]¶
Whether saturation occurred since
clear_sampled().
Nested helpers¶
- Collector._fold_unlocked.add_metric(node: Node) Node¶
Return a detached node carrying the captured metric event.
- Parameters:
node – existing run-graph node to update.
- Returns:
a dataclass copy whose detached metrics mapping sets the captured event’s name to its value. The input node and its original mapping are not mutated.
spacr/flowview/collector.py:279