Python client
diavasi_client.consume is a thin client of diavasi.data.v1. It opens a TLS stream, sends the bearer token, Hello version 1, then JoinGroup. The iterator yields each batch. Continuing the iterator acks that batch_id. The client stores no cursor and does not dedupe on record_id. A dropped stream is how unacked batches return. Reconnect with the same consumer_id and the server replays them.
proto/data.proto in this repository is the copy of diavasi.data.v1 from github.com/diavasis/diavasi tag v0.13.0. Package diavasi-client is version 0.1.0.
Install
pip install diavasi-client==0.1.0
From a checkout of this repository:
python -m venv .venv
.venv/bin/pip install -r requirements.txt
export PYTHONPATH=.
Library
from diavasi_client import CallError, ProtocolError, consume
try:
for batch in consume(
addr="127.0.0.1:7710",
ca="/tmp/diavasi-sdk/dataplane-ca.crt",
token="sdk-demo",
group_id="demo",
consumer_id="python",
expect_records=8,
):
for record in batch.records:
print(f"batch {batch.batch_id} record {record.record_id} ({len(record.payload)} bytes)")
except ProtocolError as err:
print(f"protocol {err.code}: {err.message}")
except CallError as err:
print(f"grpc {err.status}: {err.message}")
The same program is examples/process.py. From the repo root, after the server and the demo group are up:
PYTHONPATH=. .venv/bin/python examples/process.py
consume() acks in a finally around the yield, so breaking out of the loop still acks the current batch. expect_records sends Leave once that many records are acked. halt_after_acks closes after that many acks and does not send Leave. Session is the same stream when the caller wants to call ack itself.
ProtocolError carries the protocol code. CallError carries a gRPC status. A bad token raises CallError with status UNAUTHENTICATED and message unauthorized. A group that is not running raises ProtocolError with code 5.
| Code | Meaning |
|---|---|
| 1 | Bad version |
| 2 | Bad state |
| 3 | Unknown ack |
| 4 | Duplicate ack |
| 5 | Group is not running |
| 6 | Unsupported |
| 7 | Internal |
| 8 | Heartbeat timeout |
Example
PYTHONPATH=. python -m diavasi_client \
--addr 127.0.0.1:7710 --ca /tmp/diavasi-sdk/dataplane-ca.crt \
--token sdk-demo --group demo --consumer python --total 8
Flags: --addr, --ca, --token, --group, --consumer, --total, --max-in-flight (default 1), --halt-after. The last occurrence of a flag wins. The example prints record_ids and batch_ids.
docker compose -f clients/docker-compose.yml --profile python up --abort-on-container-exit
Test
python -m unittest test_consume.py returns immediately until DIAVASI_DATA_ADDR, DIAVASI_CA, and DIAVASI_API_TOKEN are set. With those set, it consumes DIAVASI_TOTAL records (default 8) from DIAVASI_GROUP.
Stage 0 bench
The TCP bench client in this tree speaks a different protocol from data.proto.
mise install
mise exec -- python -m venv .venv
mise exec -- .venv/bin/pip install -r requirements.txt
mise exec -- .venv/bin/python -m diavasi_bench.gen_proto
# terminal 1
cargo run -p diavasi --bin diavasi-transport-bench --features transport-bench -- \
--transport tcp --role server --listen 127.0.0.1:9800 --smoke
# terminal 2
PYTHONPATH=. .venv/bin/python -m diavasi_bench \
--connect 127.0.0.1:9800 --total-records 200 \
--output docs/bench/results.jsonl
Metadata
Release files for diavasi-client 0.1.0
For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.
Source distribution (sdist)
| File | Size | Uploaded | |
|---|---|---|---|
| diavasi_client-0.1.0.tar.gz | 13.7 kB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| diavasi_client-0.1.0-py3-none-any.whl | Python 3 | none | any | Details |
Total release size: 27.2 kB
Release files / diavasi_client-0.1.0.tar.gz
| Download URL | diavasi_client-0.1.0.tar.gz |
|---|---|
| Size | 13.7 kB |
| Tags | Source |
|
SHA-256 checksum How to use checksums |
eea617fbd6db256d533475c90a312b3037449cea2141d14a027f97e70a05b7eb
|
|
BLAKE2b-256 checksum How to use checksums |
5918ac29787737eb378bf9738a6c663344e20a70c698643a0ef090ab928834df
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
Yes |
| Uploaded via |
twine/7.0.0 CPython/3.12.14
|
Release files / diavasi_client-0.1.0-py3-none-any.whl
| Download URL | diavasi_client-0.1.0-py3-none-any.whl |
|---|---|
| Size | 13.5 kB |
| Tags | Python 3 |
|
SHA-256 checksum How to use checksums |
663d0d379d6db63e08c9aedfb72b5093bf81df47bf13663fac14d00163120612
|
|
BLAKE2b-256 checksum How to use checksums |
5a96652d19b81ff97991b9ecbddf1a419458c603a14af2d81124e9049986a7ae
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
Yes |
| Uploaded via |
twine/7.0.0 CPython/3.12.14
|