slipstream.core

Core module.

Attributes

Classes

PausableStream

Can signal source stream to pause.

Module Contents

slipstream.core.READ_FROM_START = -2[source]
slipstream.core.READ_FROM_END = -1[source]
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.PAUSE pauses consumption; Signal.RESUME resumes it when no other pause reasons remain.

Pause and resume are tracked per reason so independent callers compose (checkpoint downtime vs produce retry). Use counted=True when overlapping pauses share a reason, such as concurrent produces.

property iterable: collections.abc.AsyncIterable[Any][source]

Get iterable.

signal: slipstream.utils.Signal | Any = None[source]
running: asyncio.Event[source]
send_signal(signal: slipstream.utils.Signal | Any, reason: str = '', *, counted: bool = False) → None[source]

Send signal to stream.

PAUSE/RESUME compose by reason. A counted reason increments on pause and decrements on resume; an uncounted reason latches until resumed. Resume is applied only when no reasons remain.