slipstream package

Top level objects.

class slipstream.Cache(*args: Any, **kwargs: Any)[source]

Bases: ICache

Create a RocksDB database in the specified folder.

>>> cache = Cache('db/mycache')  

The cache instance acts as a callable to store data:

>>> cache('key', {'msg': 'Hello World!'})  
>>> cache['key']  
{'msg': 'Hello World!'}
cancel_all_background(wait: bool = True) → None[source]

Request stopping background work.

close() → None[source]

Flush memory to disk, and drop the current column family.

columns(backwards: bool = False, from_key: str | int | float | bytes | bool | None = None, read_opt: ReadOptions | None = None, prefix: str | int | float | bytes | bool | None = None) → Iterator[list[tuple[Any, Any]]][source]

Get values as widecolumns.

compact_range(begin: str | int | float | bytes | bool | None, end: str | int | float | bytes | bool | None, compact_opt: CompactOptions | None = None) → None[source]

Run manual compaction on range for the current column family.

create_column_family(name: str, options: Options | None = None) → Rdict[source]

Create column family.

delete(key: str | int | float | bytes | bool, write_opt: WriteOptions | None = None) → None[source]

Delete item from database.

delete_range(begin: str | int | float | bytes | bool, end: str | int | float | bytes | bool, write_opt: WriteOptions | None = None) → None[source]

Delete database items, excluding end.

destroy(options: Options | None = None) → None[source]

Delete the database.

drop_column_family(name: str) → None[source]

Drop column family by name.

entities(backwards: bool = False, from_key: str | int | float | bytes | bool | None = None, read_opt: ReadOptions | None = None, prefix: str | int | float | bytes | bool | None = None) → Iterator[tuple[str | int | float | bytes | bool, list[tuple[Any, Any]]]][source]

Get keys and entities.

flush(wait: bool = True) → None[source]

Manually flush the current column family.

flush_wal(sync: bool = True) → None[source]

Manually flush the WAL buffer.

get(key: str | int | float | bytes | bool | list[str | int | float | bytes | bool], default: T | None = None, read_opt: ReadOptions | None = None) → Any | T[source]

Get item from database by key.

get_column_family(name: str) → Rdict[source]

Get column family by name.

get_column_family_handle(name: str) → ColumnFamily[source]

Get column family handle by name.

get_entity(key: str | int | float | bytes | bool | list[str | int | float | bytes | bool], default: Any | None = None, read_opt: ReadOptions | None = None) → list[tuple[Any, Any]] | None[source]

Get wide-column from database by key.

>>> cache.get_entity('key')  
[('a', 1), ('b', 2)]
ingest_external_file(paths: list[str], opts: IngestExternalFileOptions | None = None) → None[source]

Load list of SST files into current column family.

items(backwards: bool = False, from_key: str | int | float | bytes | bool | None = None, read_opt: ReadOptions | None = None, prefix: str | int | float | bytes | bool | None = None) → Iterator[tuple[str | int | float | bytes | bool, Any]][source]

Get tuples of key-value pairs.

iter(read_opt: ReadOptions | None = None) → RdictIter[source]

Get iterable.

key_may_exist(key: str | int | float | bytes | bool, fetch: bool = False, read_opt: ReadOptions | None = None) → bool | tuple[bool, Any][source]

Check if a key exist without performing IO operations.

keys(backwards: bool = False, from_key: str | int | float | bytes | bool | None = None, read_opt: ReadOptions | None = None, prefix: str | int | float | bytes | bool | None = None) → Iterator[str | int | float | bytes | bool][source]

Get keys.

latest_sequence_number() → int[source]

Get sequence number of the most recent transaction.

list_cf(path: str, options: Options | None = None) → list[str][source]

List column families.

live_files() → list[dict[str, Any]][source]

Get list of all table files with their level, start/end key.

path() → str[source]

Get current database path.

property_int_value(name: str) → int | None[source]

Get property as int by name from current column family.

property_value(name: str) → str | None[source]

Get property by name from current column family.

put(key: str | int | float | bytes | bool, value: Any, write_opt: WriteOptions | None = None) → None[source]

Put item in database using key.

put_entity(key: str | int | float | bytes | bool, names: list[Any], values: list[Any], write_opt: WriteOptions | None = None) → None[source]

Put wide-column in database using key.

>>> cache.put_entity('key', ['a', 'b'], [1, 2])  
repair(path: str, options: Options | None = None) → None[source]

Repair the database.

set_dumps(dumps: Callable[[Any], bytes]) → None[source]

Set custom dumps function.

set_loads(loads: Callable[[bytes], Any]) → None[source]

Set custom loads function.

set_options(options: dict[str, str]) → None[source]

Set options for current column family.

set_read_options(read_opt: ReadOptions) → None[source]

Set custom read options.

set_write_options(write_opt: WriteOptions) → None[source]

Set custom write options.

snapshot() → Snapshot[source]

Create snapshot of current column family.

transaction(key: str | int | float | bytes | bool) → AsyncGenerator[Cache, None][source]

Lock the db entry while using the context manager.

>>> async with cache.transaction('fish'):  
...     cache['fish'] = '🐟'
  • This works for asynchronous code (not multi-threading/processing)

  • While locked, other transactions on the same key will block

  • Actions outside of transaction blocks ignore ongoing transactions

  • Reads aren’t limited by ongoing transactions

try_catch_up_with_primary() → None[source]

Try to catch up with the primary by reading log files.

values(backwards: bool = False, from_key: str | int | float | bytes | bool | None = None, read_opt: ReadOptions | None = None, prefix: str | int | float | bytes | bool | None = None) → Iterator[Any][source]

Get values.

write(write_batch: WriteBatch, write_opt: WriteOptions | None = None) → None[source]

Write a batch.

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

Bases: object

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.

async check_pulse(marker: datetime | Any, checkpoint_state: Any | None = 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.

async heartbeat(marker: datetime | Any, dependency_name: str | None = None, checkpoint_state: Any | None = 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.

class slipstream.Conf(*args: Any, **kwargs: Any)[source]

Bases: object

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)
exit_hooks: ClassVar[set[Callable[[P], T | Awaitable[T]]]] = {}
iterables: ClassVar[dict[str, PausableStream]] = {}
pipes: ClassVar[dict[Callable[[P], T | Awaitable[T]], tuple[str, tuple[Callable[[AsyncIterable[Any]], AsyncIterable[Any]], ...]]]] = {}
pubsub = <slipstream.utils.PubSub object>
register_exit_hook(exit_hook: Callable[[P], T | Awaitable[T]]) → None[source]

Add exit hook that’s called on shutdown.

register_handler(key: str, handler: Callable[[P], T | Awaitable[T]], *pipe: Callable[[AsyncIterable[Any]], AsyncIterable[Any]]) → None[source]

Add handler to global Conf.

register_iterable(key: str, it: AsyncIterable[Any]) → None[source]

Add iterable to global Conf.

signal_iterables(signal: 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.

class slipstream.Topic(name: str, conf: dict[str, Any] | None = None, offset: int | dict[int, int] | None = None, codec: ICodec | None = None, dry: bool = False)

Bases: object

Act as a consumer and producer.

>>> topic = Topic(
...     'emoji',
...     {
...         'bootstrap_servers': 'localhost:29091',
...         'auto_offset_reset': 'earliest',
...         'group_id': 'demo',
...     },
... )

Loop over topic (iterable) to consume from it:

>>> async for msg in topic:  
...     print(msg.value)

Call topic (callable) with data to produce to it:

>>> await topic({'msg': 'Hello World!'})  

Produce retries default to fail-fast (produce_retries=0). Set produce_retries (and optional produce_retry_backoff seconds) on Conf or the topic conf to retry retriable broker errors. Sources are paused via Conf.signal_iterables while retrying.

property admin: AIOKafkaClient

Get started instance of Kafka admin client.

async asend(value: Any) → ConsumerRecord[Any, Any]

Send data to generator.

async exit_hook() → None

Cleanup and finalization.

async get_consumer() → AIOKafkaConsumer

Get started instance of Kafka consumer.

async get_producer() → AIOKafkaProducer

Get started instance of Kafka producer.

async init_generator() → AsyncGenerator[Literal[Signal.SENTINEL] | ConsumerRecord[Any, Any], bool | None]

Initialize generator.

async seek(offset: int | dict[int, int], consumer: AIOKafkaConsumer | None = None, timeout: float = 30.0) → None

Seek to offset.

slipstream.handle(*iterable: AsyncIterable[Any], pipe: Iterable[Callable[[AsyncIterable[Any]], AsyncIterable[Any]]] = [], sink: Iterable[Callable | Callable[[P], T | Awaitable[T]]] = []) → Callable[[Callable[[P], T | Awaitable[T]]], 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) → Coroutine[None, None, None][source]

Start processing iterables bound by handle function.

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

Submodules

slipstream.caching module

Slipstream caching.

class slipstream.caching.Cache(*args: Any, **kwargs: Any)[source]

Bases: ICache

Create a RocksDB database in the specified folder.

>>> cache = Cache('db/mycache')  

The cache instance acts as a callable to store data:

>>> cache('key', {'msg': 'Hello World!'})  
>>> cache['key']  
{'msg': 'Hello World!'}
cancel_all_background(wait: bool = True) → None[source]

Request stopping background work.

close() → None[source]

Flush memory to disk, and drop the current column family.

columns(backwards: bool = False, from_key: str | int | float | bytes | bool | None = None, read_opt: ReadOptions | None = None, prefix: str | int | float | bytes | bool | None = None) → Iterator[list[tuple[Any, Any]]][source]

Get values as widecolumns.

compact_range(begin: str | int | float | bytes | bool | None, end: str | int | float | bytes | bool | None, compact_opt: CompactOptions | None = None) → None[source]

Run manual compaction on range for the current column family.

create_column_family(name: str, options: Options | None = None) → Rdict[source]

Create column family.

delete(key: str | int | float | bytes | bool, write_opt: WriteOptions | None = None) → None[source]

Delete item from database.

delete_range(begin: str | int | float | bytes | bool, end: str | int | float | bytes | bool, write_opt: WriteOptions | None = None) → None[source]

Delete database items, excluding end.

destroy(options: Options | None = None) → None[source]

Delete the database.

drop_column_family(name: str) → None[source]

Drop column family by name.

entities(backwards: bool = False, from_key: str | int | float | bytes | bool | None = None, read_opt: ReadOptions | None = None, prefix: str | int | float | bytes | bool | None = None) → Iterator[tuple[str | int | float | bytes | bool, list[tuple[Any, Any]]]][source]

Get keys and entities.

flush(wait: bool = True) → None[source]

Manually flush the current column family.

flush_wal(sync: bool = True) → None[source]

Manually flush the WAL buffer.

get(key: str | int | float | bytes | bool | list[str | int | float | bytes | bool], default: T | None = None, read_opt: ReadOptions | None = None) → Any | T[source]

Get item from database by key.

get_column_family(name: str) → Rdict[source]

Get column family by name.

get_column_family_handle(name: str) → ColumnFamily[source]

Get column family handle by name.

get_entity(key: str | int | float | bytes | bool | list[str | int | float | bytes | bool], default: Any | None = None, read_opt: ReadOptions | None = None) → list[tuple[Any, Any]] | None[source]

Get wide-column from database by key.

>>> cache.get_entity('key')  
[('a', 1), ('b', 2)]
ingest_external_file(paths: list[str], opts: IngestExternalFileOptions | None = None) → None[source]

Load list of SST files into current column family.

items(backwards: bool = False, from_key: str | int | float | bytes | bool | None = None, read_opt: ReadOptions | None = None, prefix: str | int | float | bytes | bool | None = None) → Iterator[tuple[str | int | float | bytes | bool, Any]][source]

Get tuples of key-value pairs.

iter(read_opt: ReadOptions | None = None) → RdictIter[source]

Get iterable.

key_may_exist(key: str | int | float | bytes | bool, fetch: bool = False, read_opt: ReadOptions | None = None) → bool | tuple[bool, Any][source]

Check if a key exist without performing IO operations.

keys(backwards: bool = False, from_key: str | int | float | bytes | bool | None = None, read_opt: ReadOptions | None = None, prefix: str | int | float | bytes | bool | None = None) → Iterator[str | int | float | bytes | bool][source]

Get keys.

latest_sequence_number() → int[source]

Get sequence number of the most recent transaction.

list_cf(path: str, options: Options | None = None) → list[str][source]

List column families.

live_files() → list[dict[str, Any]][source]

Get list of all table files with their level, start/end key.

path() → str[source]

Get current database path.

property_int_value(name: str) → int | None[source]

Get property as int by name from current column family.

property_value(name: str) → str | None[source]

Get property by name from current column family.

put(key: str | int | float | bytes | bool, value: Any, write_opt: WriteOptions | None = None) → None[source]

Put item in database using key.

put_entity(key: str | int | float | bytes | bool, names: list[Any], values: list[Any], write_opt: WriteOptions | None = None) → None[source]

Put wide-column in database using key.

>>> cache.put_entity('key', ['a', 'b'], [1, 2])  
repair(path: str, options: Options | None = None) → None[source]

Repair the database.

set_dumps(dumps: Callable[[Any], bytes]) → None[source]

Set custom dumps function.

set_loads(loads: Callable[[bytes], Any]) → None[source]

Set custom loads function.

set_options(options: dict[str, str]) → None[source]

Set options for current column family.

set_read_options(read_opt: ReadOptions) → None[source]

Set custom read options.

set_write_options(write_opt: WriteOptions) → None[source]

Set custom write options.

snapshot() → Snapshot[source]

Create snapshot of current column family.

transaction(key: str | int | float | bytes | bool) → AsyncGenerator[Cache, None][source]

Lock the db entry while using the context manager.

>>> async with cache.transaction('fish'):  
...     cache['fish'] = '🐟'
  • This works for asynchronous code (not multi-threading/processing)

  • While locked, other transactions on the same key will block

  • Actions outside of transaction blocks ignore ongoing transactions

  • Reads aren’t limited by ongoing transactions

try_catch_up_with_primary() → None[source]

Try to catch up with the primary by reading log files.

values(backwards: bool = False, from_key: str | int | float | bytes | bool | None = None, read_opt: ReadOptions | None = None, prefix: str | int | float | bytes | bool | None = None) → Iterator[Any][source]

Get values.

write(write_batch: WriteBatch, write_opt: WriteOptions | None = None) → None[source]

Write a batch.

class slipstream.caching.Proxy(*args: Any, **kwargs: Any)[source]

Bases: object

Proxy class to publish/subscribe messages.

slipstream.checkpointing module

Slipstream checkpointing.

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

Bases: object

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.

async check_pulse(marker: datetime | Any, checkpoint_state: Any | None = 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.

async heartbeat(marker: datetime | Any, dependency_name: str | None = None, checkpoint_state: Any | None = 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.

class slipstream.checkpointing.Dependency(name: str, dependency: AsyncIterable[Any], downtime_threshold: Any = datetime.timedelta(seconds=600), downtime_check: Callable[[Checkpoint, Dependency], Any | Awaitable[Any]] | None = None, recovery_check: Callable[[Checkpoint, Dependency], bool | Awaitable[bool]] | None = None, marker: Callable[[...], Any] | str | None = None, state: Callable[[Any, dict[str, Any]], dict[str, Any]] | None = None)[source]

Bases: object

Track the dependent stream state to recover from downtime.

The dependency name should not be changed once created, it is used to persist the dependency in the cache.

>>> async def emoji():
...     for emoji in '🏆📞🐟👌':
...         yield emoji
>>> Dependency('emoji', emoji())
{'checkpoint_state': None, 'checkpoint_marker': None}
load(cache: ICache, cache_key_prefix: str) → None[source]

Load checkpoint state from cache.

save(cache: ICache, cache_key_prefix: str, checkpoint_state: Any, checkpoint_marker: datetime) → None[source]

Save checkpoint state to cache.

uses_default_downtime_check() → bool[source]

Return whether first-pulse event-time seeding applies.

class slipstream.checkpointing.Downtime[source]

Bases: dict

Per-dependency pulse result.

Falsy when empty. A single entry compares equal to that value, so downtime == timedelta(...) still holds for one leader.

slipstream.checkpointing.bind_checkpoint(f: Callable[[...], Any], handler: Callable[[...], Awaitable[Any]], checkpoint: Checkpoint) → Callable[[...], Awaitable[Any]][source]

Heartbeat dependencies and pulse the dependent around a handler.

slipstream.codecs module

Slipstream codecs.

class slipstream.codecs.JsonCodec[source]

Bases: ICodec

Serialize/deserialize json messages.

decode(s: bytes) → object[source]

Deserialize message.

>>> c = JsonCodec()
>>> c.decode(b'{"key": 1}')
{'key': 1}
encode(obj: Any) → bytes[source]

Serialize message.

>>> c = JsonCodec()
>>> c.encode({'key': 1})
b'{"key": 1}'

slipstream.core module

Core module.

class slipstream.core.PausableStream(it: AsyncIterable[Any])[source]

Bases: object

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: AsyncIterable[Any][source]

Get iterable.

send_signal(signal: 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.

class slipstream.core.Topic(name: str, conf: dict[str, Any] | None = None, offset: int | dict[int, int] | None = None, codec: ICodec | None = None, dry: bool = False)

Bases: object

Act as a consumer and producer.

>>> topic = Topic(
...     'emoji',
...     {
...         'bootstrap_servers': 'localhost:29091',
...         'auto_offset_reset': 'earliest',
...         'group_id': 'demo',
...     },
... )

Loop over topic (iterable) to consume from it:

>>> async for msg in topic:  
...     print(msg.value)

Call topic (callable) with data to produce to it:

>>> await topic({'msg': 'Hello World!'})  

Produce retries default to fail-fast (produce_retries=0). Set produce_retries (and optional produce_retry_backoff seconds) on Conf or the topic conf to retry retriable broker errors. Sources are paused via Conf.signal_iterables while retrying.

property admin: AIOKafkaClient

Get started instance of Kafka admin client.

async asend(value: Any) → ConsumerRecord[Any, Any]

Send data to generator.

async exit_hook() → None

Cleanup and finalization.

async get_consumer() → AIOKafkaConsumer

Get started instance of Kafka consumer.

async get_producer() → AIOKafkaProducer

Get started instance of Kafka producer.

async init_generator() → AsyncGenerator[Literal[Signal.SENTINEL] | ConsumerRecord[Any, Any], bool | None]

Initialize generator.

async seek(offset: int | dict[int, int], consumer: AIOKafkaConsumer | None = None, timeout: float = 30.0) → None

Seek to offset.

slipstream.interfaces module

Slipstream interfaces.

class slipstream.interfaces.ICache(*args: Any, **kwargs: Any)[source]

Bases: object

Base class for cache implementations.

>>> class MyCache(ICache):
...     def __init__(self):
...         self.db = {}
...
...     def __contains__(self, key: Key) -> bool:
...         return key in self.db
...
...     def __delitem__(self, key: Key) -> None:
...         del self.db[key]
...
...     def __getitem__(self, key: Key | list[Key]) -> Any:
...         return self.db.get(key, None)
...
...     def __setitem__(self, key: Key, val: Any) -> None:
...         self.db[key] = val
>>> cache = MyCache()
>>> cache['prize'] = '🏆'
>>> cache['prize']
'🏆'
>>> del cache['prize']
>>> 'prize' in cache
False
class slipstream.interfaces.ICodec[source]

Bases: object

Base class for codecs.

abstract decode(s: bytes) → object[source]

Deserialize object.

abstract encode(obj: Any) → bytes[source]

Serialize object.

class slipstream.interfaces.SourceSinkMeta(name, bases, namespace, **kwargs)[source]

Bases: ABCMeta

Metaclass adds default source/sink functionalities.

slipstream.utils module

Slipstream utilities.

class slipstream.utils.AsyncSynchronizedGenerator(gen: AsyncIterable[Any])[source]

Bases: object

Async generator that synchronizes values across copies.

copy() → _GeneratorCopy[source]

Create a synchronized copy of this generator.

property value: Any[source]

Get current value the generator is holding.

class slipstream.utils.PubSub(*args: Any, **kwargs: Any)[source]

Bases: object

Singleton publish subscribe pattern class.

async apublish(topic: str, *args: Any, **kwargs: Any) → None[source]

Publish message to subscribers of topic.

async iter_topic(topic: str) → AsyncIterator[Any][source]

Asynchronously iterate over messages published to a topic.

publish(topic: str, *args: Any, **kwargs: Any) → None[source]

Publish message to subscribers of topic.

subscribe(topic: str, listener: Callable[[P], T | Awaitable[T]]) → None[source]

Subscribe callable to topic.

unsubscribe(topic: str, listener: Callable[[P], T | Awaitable[T]]) → None[source]

Unsubscribe callable from topic.

class slipstream.utils.Signal(value)[source]

Bases: Enum

Signals can be exchanged with streams.

SENTINEL represents an absent yield value PAUSE represents the signal to pause stream RESUME represents the signal to resume stream STOP represents an exhausted stream

PAUSE = 1[source]
RESUME = 2[source]
SENTINEL = 0[source]
STOP = 3[source]
class slipstream.utils.Singleton[source]

Bases: type

Maintain a single instance of a class.

async slipstream.utils.awaitable(x: Any) → Any[source]

Convert into awaitable.

slipstream.utils.get_param_names(o: Any) → tuple[str, ...][source]

Return function parameter names.