spacr.flowview.collector

Thread-safe bounded collection and event folding for FlowView.

Classes

Collector

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_size is 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 sampled becomes true.

Raises:

ValueError – if max_queue_size is 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-None limit 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 sampled was 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 pending: int[source]

Point-in-time count of events waiting under concurrent producers.

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