From 5390c9d60b299f36405dfd80942f1d91bc5bf128 Mon Sep 17 00:00:00 2001 From: majkelx Date: Wed, 5 Aug 2026 14:14:26 +0200 Subject: [PATCH 1/2] Add KV bucket support: MsgKvStore, MsgKvReader, MsgKvSubscriber --- serverish/base/exceptions.py | 23 ++++ serverish/connection/connection_jets.py | 28 +++++ serverish/messenger/__init__.py | 4 + serverish/messenger/messenger.py | 52 ++++++++ serverish/messenger/msg_kv.py | 111 +++++++++++++++++ serverish/messenger/msg_kv_read.py | 78 ++++++++++++ serverish/messenger/msg_kv_store.py | 141 +++++++++++++++++++++ serverish/messenger/msg_kv_sub.py | 111 +++++++++++++++++ tests/test_messenger_kv.py | 158 ++++++++++++++++++++++++ 9 files changed, 706 insertions(+) create mode 100644 serverish/messenger/msg_kv.py create mode 100644 serverish/messenger/msg_kv_read.py create mode 100644 serverish/messenger/msg_kv_store.py create mode 100644 serverish/messenger/msg_kv_sub.py create mode 100644 tests/test_messenger_kv.py diff --git a/serverish/base/exceptions.py b/serverish/base/exceptions.py index ced4cb5..851fe63 100644 --- a/serverish/base/exceptions.py +++ b/serverish/base/exceptions.py @@ -38,4 +38,27 @@ class MessengerRequestTimeout(MessengerRequestNoResponse, TimeoutError): pass class MessengerRequestNoResultYet(MessengerRequestNoResponse): + pass + + +class MessengerKvError(Exception): + pass + + +class MessengerKvBucketNotFound(MessengerKvError): + """Raised when a KV bucket does not exist and the driver was not allowed to create it.""" + pass + + +class MessengerKvKeyNotFound(MessengerKvError, KeyError): + pass + + +class MessengerKvMalformed(MessengerKvError): + """Raised when a KV entry does not carry a serverish envelope. + + Serverish stores full ``{"data": ..., "meta": ...}`` envelopes in KV buckets. + An entry written by a non-serverish client fails loudly instead of being + silently returned as empty. + """ pass \ No newline at end of file diff --git a/serverish/connection/connection_jets.py b/serverish/connection/connection_jets.py index 1677a30..f74c17e 100644 --- a/serverish/connection/connection_jets.py +++ b/serverish/connection/connection_jets.py @@ -5,6 +5,7 @@ import re import nats +import nats.js.errors from nats.js import JetStreamContext import param @@ -147,6 +148,33 @@ async def ensure_subject_in_stream(self, stream: str, subject: str, cfg.subjects.append(subject) await js.update_stream(config=cfg) + async def ensure_kv_bucket(self, bucket: str, create_if_needed: bool = False, **config): + """Binds to a JetStream KV bucket, optionally creating it. + + Like streams (see `ensure_subject_in_stream`), buckets are not created on the fly + by default, because usually one has to control the bucket parameters. + + Args: + bucket (str): KV bucket name + create_if_needed (bool): create the bucket if it does not exist + config: `nats.js.api.KeyValueConfig` parameters (e.g. history, ttl, max_bytes) + used only when the bucket is being created + + Returns: + nats.js.kv.KeyValue: bound KV bucket handle + + Raises: + nats.js.errors.BucketNotFoundError: bucket does not exist and create_if_needed is False + """ + js: JetStreamContext = self.js + try: + return await js.key_value(bucket) + except nats.js.errors.BucketNotFoundError: + if not create_if_needed: + raise + logger.info(f"Creating KV bucket {bucket}") + return await js.create_key_value(bucket=bucket, **config) + async def diagnose_stream_config(self) -> StatusReport: """Diagnoses stream configuration diff --git a/serverish/messenger/__init__.py b/serverish/messenger/__init__.py index bb3edf1..75c7db6 100644 --- a/serverish/messenger/__init__.py +++ b/serverish/messenger/__init__.py @@ -19,3 +19,7 @@ from .msg_core_sub import MsgCoreSub, get_coresubscriber from .msg_cmd_pub import MsgCommandPublisher, get_commandpublisher from .msg_cmd_sub import MsgCommandSubscriber, get_commandsubscriber +from .msg_kv import MsgKvDriver +from .msg_kv_store import MsgKvStore, get_kvstore, kv_get, kv_put +from .msg_kv_read import MsgKvReader, get_kvreader +from .msg_kv_sub import MsgKvSubscriber, get_kvsubscriber diff --git a/serverish/messenger/messenger.py b/serverish/messenger/messenger.py index 268179b..81870d9 100644 --- a/serverish/messenger/messenger.py +++ b/serverish/messenger/messenger.py @@ -627,6 +627,58 @@ def get_commandsubscriber(subject: str) -> 'MsgCommandSubscriber': from serverish.messenger.msg_cmd_sub import MsgCommandSubscriber return MsgCommandSubscriber(subject=subject, parent=Messenger()) + @staticmethod + def get_kvstore(bucket: str, **kwargs) -> 'MsgKvStore': + """Returns a key-value store for a given JetStream KV bucket + + Values are full serverish envelopes, validated like published messages. + + Args: + bucket (str): KV bucket name + kwargs: additional driver arguments (e.g. create_bucket, bucket_config) + + Returns: + MsgKvStore: a key-value store for the given bucket + """ + from serverish.messenger.msg_kv_store import MsgKvStore + return MsgKvStore(bucket=bucket, parent=Messenger(), **kwargs) + + @staticmethod + def get_kvreader(bucket: str, key: str = '>', **kwargs) -> 'MsgKvReader': + """Returns a KV reader (async iterator over key changes) for a given bucket + + Args: + bucket (str): KV bucket name + key (str): key or wildcard pattern to watch + kwargs: additional driver arguments (e.g. include_history, ignore_deletes) + + Returns: + MsgKvReader: a KV reader for the given bucket + + Usage: + reader = Messenger.get_kvreader('bucket', key='telescope.>') + await reader.open() + async for data, meta in reader: + print(meta['kv']['key'], data) + """ + from serverish.messenger.msg_kv_read import MsgKvReader + return MsgKvReader(bucket=bucket, key=key, parent=Messenger(), **kwargs) + + @staticmethod + def get_kvsubscriber(bucket: str, key: str = '>', **kwargs) -> 'MsgKvSubscriber': + """Returns a callback-based KV subscriber for a given bucket + + Args: + bucket (str): KV bucket name + key (str): key or wildcard pattern to watch + kwargs: additional driver arguments (e.g. include_history, ignore_deletes) + + Returns: + MsgKvSubscriber: a callback-based KV subscriber for the given bucket + """ + from serverish.messenger.msg_kv_sub import MsgKvSubscriber + return MsgKvSubscriber(bucket=bucket, key=key, parent=Messenger(), **kwargs) + class MsgDriver(Manageable): subject: str = param.String(default=None, allow_None=True, doc="User subject to publish to, prefix may be added") diff --git a/serverish/messenger/msg_kv.py b/serverish/messenger/msg_kv.py new file mode 100644 index 0000000..cc52723 --- /dev/null +++ b/serverish/messenger/msg_kv.py @@ -0,0 +1,111 @@ +"""Key-value bucket support for Messenger. + +Provides `MsgKvDriver`, the common base for all KV drivers (`MsgKvStore`, +`MsgKvReader`, `MsgKvSubscriber`). NATS JetStream KV buckets are the backend; +values are always full serverish envelopes ``{"data": ..., "meta": ...}``, +so schema validation and message metadata work exactly as for stream messages. + +Delete and purge markers are the only entries without an envelope - the KV +protocol stores them as empty tombstones. Drivers surface them as +``(None, meta)`` with ``meta['kv']['operation']`` set to ``'DEL'`` or ``'PURGE'``. +""" +from __future__ import annotations + +import logging + +import jsonschema +import nats.js.errors +import param +from nats.js.kv import KeyValue + +from serverish.base import MessengerKvBucketNotFound, MessengerKvMalformed, dt_ensure_array +from serverish.messenger.messenger import MsgDriver + +log = logging.getLogger(__name__.rsplit('.')[-1]) + + +class MsgKvDriver(MsgDriver): + """Base for KV bucket operators + + KV drivers operate on a bucket (and keys within it) instead of a subject. + The inherited `subject` is set to the underlying ``$KV.`` prefix + for diagnostics only. + + Args: + bucket (str): KV bucket name + create_bucket (bool): create the bucket on open() if it does not exist + bucket_config (dict): `nats.js.api.KeyValueConfig` parameters (e.g. history, ttl) + used only when the bucket is being created + """ + bucket = param.String(default=None, allow_None=True, doc="KV bucket name") + create_bucket = param.Boolean(default=False, doc="Create the bucket on open() if it does not exist") + bucket_config = param.Dict(default={}, doc="KeyValueConfig parameters used when creating the bucket") + + def __init__(self, **kwargs) -> None: + self.kv: KeyValue | None = None + super().__init__(**kwargs) + if self.subject is None and self.bucket: + self.subject = f'$KV.{self.bucket}' # diagnostic only, KV entries are addressed by bucket/key + + async def open(self) -> None: + try: + self.kv = await self.connection.ensure_kv_bucket(self.bucket, + create_if_needed=self.create_bucket, + **self.bucket_config) + except nats.js.errors.BucketNotFoundError as e: + raise MessengerKvBucketNotFound( + f"KV bucket '{self.bucket}' does not exist, " + f"create it externally or open the driver with create_bucket=True") from e + await super().open() + + async def close(self) -> None: + self.kv = None + await super().close() + + def encode_envelope(self, data: dict | None, meta: dict | None) -> tuple[dict, bytes]: + """Creates and validates a serverish envelope, returns (message, encoded bytes)""" + msg = self.messenger.create_msg(data, meta) + try: + self.messenger.msg_validate(msg) + except jsonschema.ValidationError as e: + log.error(f"Message {msg['meta']['id']} validation error: {e}") + raise e + return msg, self.messenger.encode(msg) + + def decode_entry(self, entry: KeyValue.Entry) -> tuple[dict | None, dict]: + """Decodes a KV entry into (data, meta) + + Delete/purge markers carry no envelope: data is None and meta contains + only the 'kv' section. + + Raises: + MessengerKvMalformed: entry value is not a serverish envelope + """ + operation = entry.operation or 'PUT' + kv_meta = { + 'bucket': entry.bucket, + 'key': entry.key, + 'revision': entry.revision, + 'operation': operation, + } + if entry.created is not None: + kv_meta['created'] = dt_ensure_array(entry.created) + if operation != 'PUT': + return None, {'kv': kv_meta} + try: + msg = self.messenger.decode(entry.value) + except (ValueError, TypeError) as e: + raise MessengerKvMalformed( + f"Entry '{entry.key}'@{entry.revision} in KV bucket '{entry.bucket}' is not " + f"a serverish envelope, was it written by a non-serverish client?") from e + if not isinstance(msg, dict) or 'meta' not in msg: + raise MessengerKvMalformed( + f"Entry '{entry.key}'@{entry.revision} in KV bucket '{entry.bucket}' carries " + f"no meta, was it written by a non-serverish client?") + data = msg.get('data', {}) + meta = msg['meta'] + meta['kv'] = kv_meta + return data, meta + + def __str__(self): + return f'{self.name} [{self.bucket}]' diff --git a/serverish/messenger/msg_kv_read.py b/serverish/messenger/msg_kv_read.py new file mode 100644 index 0000000..e3bdc62 --- /dev/null +++ b/serverish/messenger/msg_kv_read.py @@ -0,0 +1,78 @@ +from __future__ import annotations + +import param +from nats.js.kv import KeyValue + +from serverish.base import MessengerReaderStopped +from serverish.messenger.messenger import Messenger +from serverish.messenger.msg_kv import MsgKvDriver, log + + +class MsgKvReader(MsgKvDriver): + """An async-iterator over changes of KV keys + + Iteration yields (data, meta) pairs like `MsgReader`; the key, revision and + operation are available in `meta['kv']`. On open, current values of matching + keys are delivered first, then live updates as they happen. + + Delete/purge markers are yielded as (None, meta) with + meta['kv']['operation'] set to 'DEL'/'PURGE' (unless ignore_deletes is set). + + Usage: + reader = get_kvreader('bucket', key='telescope.>') + await reader.open() + async for data, meta in reader: + print(meta['kv']['key'], data) + """ + key = param.String(default='>', doc="Key or wildcard pattern to watch") + include_history = param.Boolean(default=False, doc="Deliver historical revisions of the keys on open") + ignore_deletes = param.Boolean(default=False, doc="Skip delete/purge markers") + + def __init__(self, **kwargs) -> None: + self._watcher: KeyValue.KeyWatcher | None = None + super().__init__(**kwargs) + + async def open(self) -> None: + await super().open() + self._watcher = await self.kv.watch(self.key, + include_history=self.include_history, + ignore_deletes=self.ignore_deletes) + + async def close(self) -> None: + if self._watcher is not None: + try: + await self._watcher.stop() + except Exception as e: + log.debug(f"Error stopping KV watcher for {self}: {e}") + self._watcher = None + await super().close() + + def __aiter__(self): + return self + + async def __anext__(self) -> tuple[dict | None, dict]: + if self._watcher is None: + raise MessengerReaderStopped(f"KV reader {self} is not open") + while True: + entry = await self._watcher.__anext__() # raises StopAsyncIteration when watcher is stopped + if entry is None: + # nats-py sends a None marker once the initial replay is done, it is not a value + continue + return self.decode_entry(entry) + + def __str__(self): + return f'{self.name} [{self.bucket}/{self.key}]' + + +def get_kvreader(bucket: str, key: str = '>', **kwargs) -> MsgKvReader: + """Returns a KV reader for a given bucket and key pattern + + Args: + bucket (str): KV bucket name + key (str): key or wildcard pattern to watch + kwargs: additional driver arguments (e.g. include_history, ignore_deletes) + + Returns: + MsgKvReader: a KV reader for the given bucket + """ + return Messenger.get_kvreader(bucket, key=key, **kwargs) diff --git a/serverish/messenger/msg_kv_store.py b/serverish/messenger/msg_kv_store.py new file mode 100644 index 0000000..3dc584e --- /dev/null +++ b/serverish/messenger/msg_kv_store.py @@ -0,0 +1,141 @@ +from __future__ import annotations + +import nats.js.errors +from nats.js.kv import KeyValue + +from serverish.base import MessengerKvKeyNotFound +from serverish.messenger.messenger import Messenger, MsgDriver +from serverish.messenger.msg_kv import MsgKvDriver, log + + +class MsgKvStore(MsgKvDriver): + """A key-value store over a NATS JetStream KV bucket + + Values are full serverish envelopes, so `get` returns the usual + (data, meta) pair and `put` validates the message against schemas. + + All operations are decorated with `ensure_open`, so the store can be used + one-shot without explicit open/close, like other messenger drivers. + """ + + @MsgDriver.ensure_open + async def put(self, key: str, data: dict | None = None, meta: dict | None = None) -> dict: + """Stores a value for the key, returns the published message + + Args: + key (str): key to store the value under + data (dict): message data + meta (dict): message metadata + + Returns: + dict: published message, with meta['kv'] carrying bucket/key/revision + """ + msg, bdata = self.encode_envelope(data, meta) + revision = await self.kv.put(key, bdata) + msg['meta']['kv'] = {'bucket': self.bucket, 'key': key, 'revision': revision, 'operation': 'PUT'} + self.messenger.log_msg_trace(msg.get('data', {}), msg['meta'], f"KV PUT {self.bucket}/{key}") + return msg + + @MsgDriver.ensure_open + async def get(self, key: str, revision: int | None = None) -> tuple[dict, dict]: + """Returns (data, meta) of the latest (or given revision) value for the key + + Args: + key (str): key to read + revision (int): specific revision to read, latest if None + + Raises: + MessengerKvKeyNotFound: key does not exist (or is deleted) + MessengerKvMalformed: entry is not a serverish envelope + """ + try: + entry = await self.kv.get(key, revision=revision) + except nats.js.errors.KeyNotFoundError as e: + raise MessengerKvKeyNotFound(f"Key '{key}' not found in KV bucket '{self.bucket}'") from e + data, meta = self.decode_entry(entry) + self.messenger.log_msg_trace(data, meta, f"KV GET {self.bucket}/{key}") + return data, meta + + @MsgDriver.ensure_open + async def delete(self, key: str) -> None: + """Places a delete marker for the key and removes previous revisions""" + await self.kv.delete(key) + log.debug(f"KV DEL {self.bucket}/{key}") + + @MsgDriver.ensure_open + async def purge(self, key: str) -> None: + """Removes the key including all its revisions""" + await self.kv.purge(key) + log.debug(f"KV PURGE {self.bucket}/{key}") + + @MsgDriver.ensure_open + async def keys(self) -> list[str]: + """Returns list of keys in the bucket, empty list for an empty bucket""" + try: + return await self.kv.keys() + except nats.js.errors.NoKeysError: + return [] + + @MsgDriver.ensure_open + async def history(self, key: str) -> list[tuple[dict | None, dict]]: + """Returns the revision history of the key, oldest first + + Delete/purge markers are returned as (None, meta) entries. + Note: the bucket must be created with history > 1 to keep past revisions. + + Raises: + MessengerKvKeyNotFound: key has no history in the bucket + """ + try: + entries = await self.kv.history(key) + except nats.js.errors.NoKeysError as e: + raise MessengerKvKeyNotFound(f"Key '{key}' has no history in KV bucket '{self.bucket}'") from e + return [self.decode_entry(entry) for entry in entries] + + @MsgDriver.ensure_open + async def status(self) -> KeyValue.BucketStatus: + """Returns the status of the underlying KV bucket""" + return await self.kv.status() + + +def get_kvstore(bucket: str, **kwargs) -> MsgKvStore: + """Returns a key-value store for a given KV bucket + + Args: + bucket (str): KV bucket name + kwargs: additional driver arguments (e.g. create_bucket, bucket_config) + + Returns: + MsgKvStore: a key-value store for the given bucket + """ + return Messenger.get_kvstore(bucket, **kwargs) + + +async def kv_put(bucket: str, key: str, data: dict | None = None, meta: dict | None = None, **kwargs) -> dict: + """Stores a single value in a KV bucket (one-shot) + + Args: + bucket (str): KV bucket name + key (str): key to store the value under + data (dict): message data + meta (dict): message metadata + kwargs: additional driver arguments (e.g. create_bucket) + + Returns: + dict: published message + """ + return await get_kvstore(bucket, **kwargs).put(key, data, meta) + + +async def kv_get(bucket: str, key: str, **kwargs) -> tuple[dict, dict]: + """Reads a single value from a KV bucket (one-shot) + + Args: + bucket (str): KV bucket name + key (str): key to read + kwargs: additional driver arguments + + Returns: + tuple[dict, dict]: (data, meta) of the stored message + """ + return await get_kvstore(bucket, **kwargs).get(key) diff --git a/serverish/messenger/msg_kv_sub.py b/serverish/messenger/msg_kv_sub.py new file mode 100644 index 0000000..c3ac7a0 --- /dev/null +++ b/serverish/messenger/msg_kv_sub.py @@ -0,0 +1,111 @@ +from __future__ import annotations + +import asyncio +from asyncio import CancelledError, Event +from typing import Callable + +import param + +from serverish.base import Task, create_task +from serverish.messenger.messenger import Messenger +from serverish.messenger.msg_kv_read import MsgKvReader +from serverish.messenger.msg_kv import log + + +class MsgKvSubscriber(MsgKvReader): + """A class for watching KV keys and calling a callback function on each change + + This class works like `MsgKvReader`, but allows to specify a callback function + for each change instead of iterating. + + The callback is called with (data, meta); data is None for delete/purge markers + (check meta['kv']['operation']). Note that current values of matching keys are + delivered on subscribe, before live updates. + + Usage: + def callback(data, meta): + print(meta['kv']['key'], data) + + sub = get_kvsubscriber('bucket', key='telescope.>') + await sub.open() + await sub.subscribe(callback) + """ + callback = param.Callable(default=None, doc="Callback function to call on each change") + task = param.ClassSelector(default=None, class_=Task, doc="Task for watching changes") + + def __init__(self, **kwargs) -> None: + self._stop_event = Event() + super().__init__(**kwargs) + + async def close(self) -> None: + await self.stop() + if self.task is not None: + self.task.cancel() + return await super().close() + + async def stop(self) -> None: + """Stops watching changes""" + self._stop_event.set() + + async def subscribe(self, callback: Callable[[dict | None, dict], bool] | + Callable[[dict | None, dict], asyncio.Future]) -> Task: + """Sets a callback function for each change of watched keys + + Args: + callback: a callback function to call on each change, may be asynchronous + callback is called with two arguments: data dict (None for delete/purge + markers) and metadata dict, and may return False to stop watching. + Any other return value (including None — i.e. callbacks with no explicit + return — and True) keeps the subscription running. + """ + self.callback = callback + if asyncio.iscoroutinefunction(callback): + self.task = await create_task(self._task_body(acb=callback), f'KVASUB.{self.bucket}.{self.key}') + else: + self.task = await create_task(self._task_body(scb=callback), f'KVSSUB.{self.bucket}.{self.key}') + return self.task + + async def _task_body(self, + scb: Callable[[dict | None, dict], bool] | None = None, + acb: Callable[[dict | None, dict], asyncio.Future] | None = None + ) -> None: + + assert scb is not None or acb is not None + assert not (scb is not None and acb is not None) + cb = scb or acb + cont: object = True + log.debug(f"Entering KV watch iteration {self}") + async for data, meta in self: + try: + if scb is not None: + log.debug(f"Calling sync callback {cb} for KV change {meta}{str(data):20}") + cont = scb(data, meta) + else: + log.debug(f"Calling async callback {cb} for KV change {meta}{str(data):20}") + cont = await acb(data, meta) + except CancelledError: + log.debug(f'Cancelled {self}') + break + except Exception as e: + log.exception(f'Error in callback {cb} for KV change {meta}{str(data):20}: {e}') + # Stop only on explicit False — None (the implicit "no return" value + # of a Python function) and any truthy value keep the subscription + # alive. Treating None as "stop" would silently kill any callback + # that doesn't bother returning anything, which is the common case. + if cont is False or self._stop_event.is_set(): + break + log.debug(f"Exiting KV watch iteration {self}") + + +def get_kvsubscriber(bucket: str, key: str = '>', **kwargs) -> MsgKvSubscriber: + """Returns a callback-based KV subscriber for a given bucket and key pattern + + Args: + bucket (str): KV bucket name + key (str): key or wildcard pattern to watch + kwargs: additional driver arguments (e.g. include_history, ignore_deletes) + + Returns: + MsgKvSubscriber: a callback-based KV subscriber for the given bucket + """ + return Messenger.get_kvsubscriber(bucket, key=key, **kwargs) diff --git a/tests/test_messenger_kv.py b/tests/test_messenger_kv.py new file mode 100644 index 0000000..2bf56e9 --- /dev/null +++ b/tests/test_messenger_kv.py @@ -0,0 +1,158 @@ +"""Tests for KV bucket support: MsgKvStore, MsgKvReader, MsgKvSubscriber.""" +from __future__ import annotations + +import asyncio +import uuid + +import pytest +import pytest_asyncio + +from serverish.base import MessengerKvBucketNotFound, MessengerKvKeyNotFound, MessengerKvMalformed +from serverish.messenger import (Messenger, get_kvstore, get_kvreader, get_kvsubscriber, + kv_get, kv_put) + + +@pytest_asyncio.fixture(loop_scope='session') +async def kv_bucket(messenger): + """Provide a unique KV bucket name, delete the bucket on teardown.""" + bucket = f"test-kv-{uuid.uuid4().hex[:8]}" + yield bucket + try: + await Messenger().connection.js.delete_key_value(bucket) + except Exception: + pass # bucket may have never been created + + +@pytest.mark.nats_js +async def test_kv_put_get(messenger, kv_bucket): + store = get_kvstore(kv_bucket, create_bucket=True) + async with store: + msg = await store.put('telescope.focus', {'position': 1250}) + assert msg['meta']['kv']['bucket'] == kv_bucket + assert msg['meta']['kv']['key'] == 'telescope.focus' + assert msg['meta']['kv']['revision'] >= 1 + + data, meta = await store.get('telescope.focus') + assert data == {'position': 1250} + # full envelope round-trip: standard meta fields survive + assert meta['id'] == msg['meta']['id'] + assert 'ts' in meta + assert meta['kv']['operation'] == 'PUT' + assert meta['kv']['revision'] == msg['meta']['kv']['revision'] + + +@pytest.mark.nats_js +async def test_kv_get_missing_key(messenger, kv_bucket): + async with get_kvstore(kv_bucket, create_bucket=True) as store: + with pytest.raises(MessengerKvKeyNotFound): + await store.get('no.such.key') + + +@pytest.mark.nats_js +async def test_kv_bucket_not_found(messenger): + store = get_kvstore(f"test-kv-missing-{uuid.uuid4().hex[:8]}") + with pytest.raises(MessengerKvBucketNotFound): + await store.open() + + +@pytest.mark.nats_js +async def test_kv_delete(messenger, kv_bucket): + async with get_kvstore(kv_bucket, create_bucket=True) as store: + await store.put('ephemeral', {'v': 1}) + await store.delete('ephemeral') + with pytest.raises(MessengerKvKeyNotFound): + await store.get('ephemeral') + + +@pytest.mark.nats_js +async def test_kv_keys_and_history(messenger, kv_bucket): + async with get_kvstore(kv_bucket, create_bucket=True, bucket_config={'history': 5}) as store: + assert await store.keys() == [] + + await store.put('alpha', {'v': 1}) + await store.put('alpha', {'v': 2}) + await store.put('beta', {'v': 3}) + + assert sorted(await store.keys()) == ['alpha', 'beta'] + + history = await store.history('alpha') + assert [data for data, meta in history] == [{'v': 1}, {'v': 2}] + revisions = [meta['kv']['revision'] for data, meta in history] + assert revisions == sorted(revisions) + + +@pytest.mark.nats_js +async def test_kv_malformed_entry(messenger, kv_bucket): + async with get_kvstore(kv_bucket, create_bucket=True) as store: + # bypass serverish and write a raw value, as a foreign client would + await store.kv.put('foreign', b'not a serverish envelope') + with pytest.raises(MessengerKvMalformed): + await store.get('foreign') + + +@pytest.mark.nats_js +async def test_kv_oneshot(messenger, kv_bucket): + await kv_put(kv_bucket, 'dome.status', {'open': True}, create_bucket=True) + data, meta = await kv_get(kv_bucket, 'dome.status') + assert data == {'open': True} + assert meta['kv']['key'] == 'dome.status' + + +@pytest.mark.nats_js +async def test_kv_reader_watch(messenger, kv_bucket): + async with get_kvstore(kv_bucket, create_bucket=True) as store: + await store.put('alpha', {'v': 1}) + + reader = get_kvreader(kv_bucket) + await reader.open() + try: + got = [] + + async def consume(): + async for data, meta in reader: + got.append((data, meta)) + if len(got) >= 3: + break + + consumer = asyncio.create_task(consume()) + await asyncio.sleep(0.3) # let the watcher deliver the initial value + await store.put('beta', {'v': 2}) + await store.delete('alpha') + await asyncio.wait_for(consumer, timeout=5) + finally: + await reader.close() + + # initial replay of the existing key, then live updates in order + assert got[0][0] == {'v': 1} + assert got[0][1]['kv']['key'] == 'alpha' + assert got[1][0] == {'v': 2} + assert got[1][1]['kv']['key'] == 'beta' + # delete marker carries no envelope + assert got[2][0] is None + assert got[2][1]['kv'] == {**got[2][1]['kv'], 'key': 'alpha', 'operation': 'DEL'} + + +@pytest.mark.nats_js +async def test_kv_subscriber_callback(messenger, kv_bucket): + async with get_kvstore(kv_bucket, create_bucket=True) as store: + events = [] + done = asyncio.Event() + + def callback(data, meta): + events.append((data, meta)) + if len(events) >= 2: + done.set() + return False + + sub = get_kvsubscriber(kv_bucket) + await sub.open() + try: + await sub.subscribe(callback) + await store.put('k1', {'n': 1}) + await store.put('k2', {'n': 2}) + await asyncio.wait_for(done.wait(), timeout=5) + finally: + await sub.close() + + assert [data for data, meta in events] == [{'n': 1}, {'n': 2}] + assert [meta['kv']['key'] for data, meta in events] == ['k1', 'k2'] From 875d4f1bc89d806f60c13658e353b5405c3fd3f3 Mon Sep 17 00:00:00 2001 From: majkelx Date: Wed, 5 Aug 2026 14:15:10 +0200 Subject: [PATCH 2/2] ver bump --- pyproject.toml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pyproject.toml b/pyproject.toml index 7670b0b..f174330 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -1,6 +1,6 @@ [tool.poetry] name = "serverish" -version = "2.0.5" +version = "2.1.0" description = "helpers for server alike projects" authors = ["Mikołaj Kałuszyński", "MMME team"] readme = "README.md"