Skip to main content

Python client

CI PyPI Python license

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)

Source distribution for diavasi-client 0.1.0
File Size Uploaded
diavasi_client-0.1.0.tar.gz 13.7 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for diavasi-client 0.1.0
File Interpreter ABI Platform
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

Release history Release notifications | RSS feed

This release

0.1.0 This release

2 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page