spacr.flowview.feeder

Failure-isolated multiprocessing bridge for FlowView events.

Classes

MultiprocessingFeeder

Feed a multiprocessing-compatible queue into an in-process collector.

Functions

is_transport_event(→ bool)

Return whether a candidate is a bounded public transport event.

put_event_nowait(→ bool)

Validate and offer one event without waiting for queue capacity.

Module Contents

class spacr.flowview.feeder.MultiprocessingFeeder(source: _QueueSource, collector: spacr.flowview.collector.Collector, *, poll_interval: float = 0.05, max_event_bytes: int = MAX_EVENT_BYTES)[source]

Feed a multiprocessing-compatible queue into an in-process collector.

The source remains owned by the caller. Stopping does not close or intentionally empty it, but one value returned by an in-flight read after shutdown can be consumed and discarded rather than emitted.

Initialise a stopped queue-to-collector bridge.

Parameters:
  • source – caller-owned queue-like source whose get(timeout=...) returns candidate events.

  • collector – in-process collector that receives validated public events.

  • poll_interval – positive finite seconds for each source read and fault backoff.

  • max_event_bytes – positive integer byte ceiling applied when an event is re-pickled before forwarding. This limits forwarded data; the producer helper must be used to enforce it before enqueue.

Raises:

ValueError – if either limit is invalid.

start() → MultiprocessingFeeder[source]

Start one daemon feeder thread if none is alive and return this feeder.

stop(timeout: float = 1.0) → bool[source]

Request shutdown and report whether no feeder thread remains.

Parameters:

timeout – non-negative finite seconds to wait for the thread.

Returns:

True if no thread remains, or False when the timeout expires or the feeder thread itself requests shutdown.

Raises:

ValueError – if timeout is negative or non-finite.

property running: bool[source]

Return whether this feeder currently has a live daemon thread.

spacr.flowview.feeder.is_transport_event(value: object, *, max_event_bytes: int = MAX_EVENT_BYTES) → bool[source]

Return whether a candidate is a bounded public transport event.

Parameters:
  • value – candidate value to inspect.

  • max_event_bytes – positive integer ceiling for its highest-protocol pickle.

Returns:

False for the wrong type, a pickle failure, or an oversized pickle; otherwise True.

Raises:

ValueError – if max_event_bytes is not a positive integer.

spacr.flowview.feeder.put_event_nowait(destination: _QueueSink, event: object, *, max_event_bytes: int = MAX_EVENT_BYTES) → bool[source]

Validate and offer one event without waiting for queue capacity.

Parameters:
  • destination – queue-like sink providing put_nowait.

  • event – candidate public FlowView event.

  • max_event_bytes – positive integer pickle-size ceiling.

Returns:

True only when the event validates and put_nowait returns without raising; otherwise False.

Raises:

ValueError – if max_event_bytes is not a positive integer.