slipstream.core¶
Core module.
Attributes¶
Classes¶
Can signal source stream to pause. |
Module Contents¶
- class slipstream.core.PausableStream(it: collections.abc.AsyncIterable[Any])[source]¶
Can signal source stream to pause.
If it is of type AsyncGenerator, it will receive the signal through the yield send syntax in order to handle the state change appropriately. Alternatively, the signal property can be used directly.
For example, the Topic class uses the signal to pause the Consumer.
Only
Signal.PAUSEpauses consumption;Signal.RESUMEresumes it when no other pause reasons remain.Pause and resume are tracked per
reasonso independent callers compose (checkpoint downtime vs produce retry). Usecounted=Truewhen overlapping pauses share a reason, such as concurrent produces.- signal: slipstream.utils.Signal | Any = None[source]¶
- send_signal(signal: slipstream.utils.Signal | Any, reason: str = '', *, counted: bool = False) None[source]¶
Send signal to stream.
PAUSE/RESUMEcompose byreason. A counted reason increments on pause and decrements on resume; an uncounted reason latches until resumed. Resume is applied only when no reasons remain.