slipstream

Top level objects.

Submodules

Classes

Checkpoint

Track one stream against the streams it depends on.

Conf

The application configuration singleton.

Functions

handle(...)

Bind sources and sinks to the handler function.

stream(→ collections.abc.Coroutine[None, None, None])

Start processing iterables bound by handle function.

Package Contents

class slipstream.Checkpoint(dependent: collections.abc.AsyncIterable[Any], *leaders: collections.abc.AsyncIterable[Any] | Dependency, dependencies: collections.abc.AsyncIterable[Any] | Dependency | list[collections.abc.AsyncIterable[Any] | Dependency] | None = None, name: str | None = None, on_downtime: collections.abc.Callable[[Checkpoint, Dependency], Any] | None = None, on_recovery: collections.abc.Callable[[Checkpoint, Dependency], Any] | None = None, cache: slipstream.interfaces.ICache | None = None, cache_key_prefix: str = '_', pause_dependent: bool = True, downtime_threshold: datetime.timedelta | None = None, marker: collections.abc.Callable[..., Any] | str | None = None, state: collections.abc.Callable[[Any, dict[str, Any]], dict[str, Any]] | None = None)[source]

Track one stream against the streams it depends on.

>>> async def emoji():
...     for emoji in '🏆📞🐟👌':
...         yield emoji
>>> dependent, dependency = emoji(), emoji()
>>> checkpoint = Checkpoint(
...     dependent,
...     Dependency('dependency', dependency),
...     name='dependent',
... )

Pass a marker and bind the checkpoint to call heartbeat and check_pulse automatically:

>>> checkpoint = Checkpoint(
...     dependent,
...     Dependency('dependency', dependency, marker='timestamp'),
...     name='dependent',
...     marker='timestamp',
... )
>>> from slipstream import handle
>>> @handle(checkpoint)
... async def dependent_handler(msg, checkpoint=None):
...     yield msg

If no cache is provided, the checkpoint lasts only for this process.

name
dependent
dependencies
pause_dependent = True
marker = None
state_extractor = None
downtime: Any | None = None
state
state_marker: Any = None
async heartbeat(marker: datetime.datetime | Any, dependency_name: str | None = None, checkpoint_state: Any = None) → dict[source]

Update checkpoint to latest state.

Args:
marker (datetime | Any): Typically the event timestamp that is

compared to the event timestamp of a dependent stream.

dependency_name (str, optional): Required when there are multiple

dependencies to specify which one the heartbeat is for.

checkpoint_state: Complete caller-selected dependency state.

async check_pulse(marker: datetime.datetime | Any, checkpoint_state: Any = None, **kwargs: Any) → Any | None[source]

Update state that can be used as checkpoint.

Args:
marker (datetime | Any): Typically the event timestamp that is

compared to the event timestamp of a dependency stream.

checkpoint_state: Complete caller-selected dependent state. kwargs (Any): Any information that can be used for reprocessing any

incorrect data that was sent out during downtime of a dependency stream, stored in state.

Returns:

None when every dependency is healthy. Otherwise a Downtime map of name → check (or True if still down but not over threshold). One leader still compares equal to its timedelta / True.

class slipstream.Conf(conf: dict[str, Any] | None = None)[source]

The application configuration singleton.

Register iterables (sources) and handlers (sinks): >>> from slipstream import handle

>>> async def messages():
...     for emoji in '🏆📞🐟👌':
...         yield emoji
>>> @handle(messages(), sink=[print])
... def handle_message(msg):
...     yield f'Hello {msg}!'

Set application kafka configuration (optional):

>>> Conf({'bootstrap_servers': 'localhost:29091'})
{'bootstrap_servers': 'localhost:29091'}

Provide exit hooks:

>>> async def exit_hook():
...     print('Shutting down application.')
>>> c = Conf()
>>> c.register_exit_hook(exit_hook)
pubsub
iterables: ClassVar[dict[str, PausableStream]]
pipes: ClassVar[dict[slipstream.utils.AsyncCallable, tuple[str, tuple[slipstream.utils.Pipe, ...]]]]
exit_hooks: ClassVar[set[slipstream.utils.AsyncCallable]]
conf: dict[str, Any]
register_iterable(key: str, it: collections.abc.AsyncIterable[Any]) → None[source]

Add iterable to global Conf.

register_handler(key: str, handler: slipstream.utils.AsyncCallable, *pipe: slipstream.utils.Pipe) → None[source]

Add handler to global Conf.

register_exit_hook(exit_hook: slipstream.utils.AsyncCallable) → None[source]

Add exit hook that’s called on shutdown.

signal_iterables(signal: slipstream.utils.Signal, reason: str = '', *, counted: bool = False) → None[source]

Send signal to every registered source stream.

async start(**kwargs: Any) → None[source]

Start processing registered iterables.

slipstream.handle(*iterable: collections.abc.AsyncIterable[Any], pipe: collections.abc.Iterable[slipstream.utils.Pipe] = [], sink: collections.abc.Iterable[collections.abc.Callable | slipstream.utils.AsyncCallable] = []) → collections.abc.Callable[[slipstream.utils.AsyncCallable], Handler][source]

Bind sources and sinks to the handler function.

Ex:
>>> topic = Topic('demo')  
>>> cache = Cache('state/demo')  
>>> @handle(topic, sink=[print, cache])  
... def handler(msg, **kwargs):
...     return msg.key, msg.value
slipstream.stream(**kwargs: Any) → collections.abc.Coroutine[None, None, None][source]

Start processing iterables bound by handle function.

Ex:
>>> from asyncio import run
>>> kwargs = {
...     'env': 'DEV',
... }
>>> run(stream(**kwargs))