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
1 change: 1 addition & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@ jobs:
uses: actions/checkout@v4
with:
repository: puffball1567/koutendb
ref: v0.12.0
path: koutendb-core

- name: Set up Python
Expand Down
6 changes: 4 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@ Applications pass a human-readable ring name, and KoutenDB returns a typed ID.

## Status

- package: PyPI [`koutendb`](https://pypi.org/project/koutendb/) v0.2.0
- package: PyPI [`koutendb`](https://pypi.org/project/koutendb/) v0.2.1
- current mode: native TCP wire driver
- Python: 3.10+
- runtime dependencies: none
Expand All @@ -23,13 +23,15 @@ Implemented:
- `query` / `query_encoded` / `query_text` / `query_json`
- codec metadata negotiation with `CODECMETA ON`
- `batch_get`
- direct owner redirects from extended `FWD ... owner` responses
- routed multi-node `batch_get` fallback with stable input ordering
- typed `KoutenId`
- one reconnect retry
- context manager support
- username/password, shared-secret transport, and TLS authentication

Planned:

- authentication / secret-key handshake support
- retrieve / atlas wire APIs once the public wire contract is finalized for drivers
- ring-read filters/projection once the public wire contract is finalized for drivers
- connection pooling
Expand Down
51 changes: 46 additions & 5 deletions koutendb/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -284,8 +284,35 @@ def query_json(self, doc_id: KoutenId, selection: str, node: Optional[int] = Non
value = self.query_text(doc_id, selection, node=node)
return None if value is None else json.loads(value)

def batch_get(self, ids: Iterable[KoutenId], node: int = 0) -> list[Optional[bytes]]:
def batch_get(
self, ids: Iterable[KoutenId], node: Optional[int] = None
) -> list[Optional[bytes]]:
id_list = list(ids)
if not id_list:
return []
if node is not None:
return self._batch_get_node(id_list, node)

result: list[Optional[bytes]] = [None] * len(id_list)
missing = list(range(len(id_list)))
for peer_node in range(len(self.peers)):
if not missing:
break
values = self._batch_get_node(
[id_list[index] for index in missing], peer_node
)
still_missing: list[int] = []
for index, value in zip(missing, values):
if value is None:
still_missing.append(index)
else:
result[index] = value
missing = still_missing
return result

def _batch_get_node(
self, id_list: list[KoutenId], node: int
) -> list[Optional[bytes]]:
body = "".join(
f"{doc_id.parent} {doc_id.seq} {doc_id.period} {doc_id.head} {doc_id.t_write}\n"
for doc_id in id_list
Expand All @@ -311,7 +338,12 @@ def batch_get(self, ids: Iterable[KoutenId], node: int = 0) -> list[Optional[byt
return out

def _read_id_encoded(
self, op: str, doc_id: KoutenId, selection: bytes, node: int
self,
op: str,
doc_id: KoutenId,
selection: bytes,
node: int,
redirects_left: int = 2,
) -> Optional[EncodedPayload]:
header = (
f"{op} {doc_id.parent} {doc_id.epoch} {doc_id.seq} "
Expand All @@ -322,13 +354,15 @@ def _read_id_encoded(
parts = self._rpc(node, header, selection)
if not parts:
raise KoutenError(f"{op} returned an empty response")
if parts[0] == "MISS":
if parts[0] in ("MISS", "GONE"):
return None
if parts[0] == "ERR":
raise KoutenError(" ".join(parts[1:]))
if parts[0] == "FWD":
if len(parts) != 7:
if len(parts) not in (7, 8):
raise KoutenError("invalid FWD response: " + " ".join(parts))
if redirects_left <= 0:
raise KoutenError("too many FWD redirects")
fwd = KoutenId(
parent=int(parts[1]),
epoch=int(parts[2]),
Expand All @@ -337,7 +371,14 @@ def _read_id_encoded(
period=float(parts[5]),
head=float(parts[6]),
)
return self._read_id_encoded(op, fwd, selection, node=node)
target_node = int(parts[7]) if len(parts) == 8 else node
return self._read_id_encoded(
op,
fwd,
selection,
node=target_node,
redirects_left=redirects_left - 1,
)
if parts[0] != "VAL" or len(parts) not in (3, 4):
raise KoutenError(f"{op} failed: " + " ".join(parts))
codec = _codec(parts[3]) if len(parts) == 4 else ("json" if op == "QRYID" else "raw")
Expand Down
2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"

[project]
name = "koutendb"
version = "0.2.0"
version = "0.2.1"
description = "Pure Python TCP driver for KoutenDB"
readme = "README.md"
requires-python = ">=3.10"
Expand Down
Loading