diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 9dcf23c..8ca0869 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -27,6 +27,7 @@ jobs: uses: actions/checkout@v4 with: repository: puffball1567/koutendb + ref: v0.12.0 path: koutendb-core - name: Set up Python diff --git a/README.md b/README.md index ed69d1d..c1f168f 100644 --- a/README.md +++ b/README.md @@ -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 @@ -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 diff --git a/koutendb/client.py b/koutendb/client.py index 4cfc213..cc9611a 100644 --- a/koutendb/client.py +++ b/koutendb/client.py @@ -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 @@ -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} " @@ -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]), @@ -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") diff --git a/pyproject.toml b/pyproject.toml index 92e810c..cebc296 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -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"