slipstream¶
Top level objects.
Submodules¶
Classes¶
Track one stream against the streams it depends on. |
|
The application configuration singleton. |
Functions¶
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
heartbeatandcheck_pulseautomatically:>>> 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
Downtimemap of name → check (orTrueif 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
signalto every registered source stream.
- 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