Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
@@ -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"
Expand Down
23 changes: 23 additions & 0 deletions serverish/base/exceptions.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
28 changes: 28 additions & 0 deletions serverish/connection/connection_jets.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
import re

import nats
import nats.js.errors
from nats.js import JetStreamContext
import param

Expand Down Expand Up @@ -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
Expand Down
4 changes: 4 additions & 0 deletions serverish/messenger/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
52 changes: 52 additions & 0 deletions serverish/messenger/messenger.py
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down
111 changes: 111 additions & 0 deletions serverish/messenger/msg_kv.py
Original file line number Diff line number Diff line change
@@ -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.<bucket>`` 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}]'
78 changes: 78 additions & 0 deletions serverish/messenger/msg_kv_read.py
Original file line number Diff line number Diff line change
@@ -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)
Loading
Loading