streaming-proof
streaming-proof is a black-box harness that checks the delivery and ordering guarantees of a running streaming system.
The Lakestream specs describe contracts. This page describes a harness for checking whether a running system delivers and orders messages as expected: streaming-proof, an open source correctness-verification framework for distributed streaming systems that uses no Lakestream APIs.
It is black-box. It drives a deployment through its ordinary client protocol, observes what comes back, and reports violations. It does not read internal state, inspect storage, or depend on Lakestream APIs, which is what lets the same harness run against different systems and against the same system under different configurations.
What it checks
The harness produces uniquely sequenced messages across a set of keys and verifies their reception and order on the consumer side. A run is called a proof and names the guarantees it is checking through a features list:
at_least_once— every produced message is eventually consumed.ordering— messages for a key arrive in the sequence they were produced in.exactly_once— for Kafka, produce through an embedded transactional processor and verify that the output topic contains each message once.
A separate mode covers Pulsar Shared subscriptions, where messages are distributed round-robin across consumers and per-key ordering is not expected; verification there compares the highest contiguous sequence per key across all consumers rather than a last-sequence value.
Drivers in the repository cover Kafka and Kafka-API-compatible systems, Pulsar, and MQTT. A driver is an implementation of a ProofDriver interface that knows how to create producers and consumers and move messages for one target system.
How a run works
A coordinator manages the lifecycle and one or more workers do the producing and consuming, so a run can be spread across nodes or availability zones.
- Register workers and driver settings through
PUT /configs. - Create a proof with
POST /proofs, specifying the driver, features, topic, partitions, producer and consumer counts, message rate, key count, and durations. - Producers emit
(key, sequence)pairs and track the last sequence sent per key. Consumers track the last valid sequence received per key and record duplicates, out-of-order arrivals, and sequence gaps. - The coordinator takes periodic checkpoints from both sides, aggregates them, and compares the expected producer state against the observed consumer state. A
timeoutbounds how long a produced sequence may remain unverified before it is flagged;finalWaitSecondsgives consumers extra time to catch up after producers stop. - Query progress with
GET /proofs/{id}orGET /proofs/{id}/report, and stop the run withPUT /proofs/{id}/stop.
Verification is per key, which is what makes gap detection meaningful: a consumer that has seen sequence ranges [1-5] and [10-15] for a key is missing 6-9, and the harness reports those four as missed rather than waiting for a total-order comparison.
Proof status fields
A proof's summary carries the counters below. The first six are the correctness result; the last two describe whether verification was making progress.
| Field | Meaning |
|---|---|
verified | Total messages successfully verified |
errors | Failures during message publication |
outOfOrders | Messages received out of sequence |
missed | Messages not received for verification |
duplicates | Duplicate messages detected |
timeouts | Verification attempts that exceeded the configured timeout |
verifiedStallSeconds | Wall-clock seconds since verified last increased; resets when it moves again, and freezes when the run completes |
maxVerifiedStallSeconds | Peak stall observed during the run; retained across recoveries and frozen at completion |
The two stall fields are always reported. To make a long stall fail a run rather than merely appear in the output, set maxStallSeconds when creating the proof: if maxVerifiedStallSeconds exceeds it, the result is marked failed with the observed and configured values. The check is opt-in and defaults to 0, which disables it.
A passing run is evidence, not a guarantee
A proof reports what happened during one run, against one deployment, under one workload and one set of conditions. A run with no violations does not establish that the guarantee holds in general, under different failure modes, or at a different scale. Treat results as bounded evidence, and record the configuration alongside them.
Performance benchmarking is an explicit non-goal of the framework: it measures correctness, not throughput or latency. So are automated failure injection — the framework is designed to run alongside chaos tooling rather than supply it — and verification of application-level stream processing logic.
Running it
The repository builds with Maven on Java 21 or later, with a Docker profile for building the image. It ships a Docker Compose stack that brings up the coordinator, a worker, and a Kafka broker with ZooKeeper for local runs, and a Helm chart for Kubernetes. The chart deliberately excludes worker and broker connection details; supply those through PUT /configs or a private values file rather than committing credentials. The coordinator also serves a small UI prototype over the same port as the API.
Relationship to the formal models
These are two different kinds of checking, and neither substitutes for the other.
- The protocol specifications are model-checked: a bounded exhaustive search over the state space of an abstract protocol, which can show that no reachable state violates a stated property. It says nothing about whether a given codebase implements that protocol.
- streaming-proof is empirical: it exercises a real deployment and reports the violations it observed. It can find implementation and integration defects that a protocol model cannot express, and it can miss anything the run did not happen to trigger.
Together they cover different failure surfaces. A protocol can be sound while its implementation is not; an implementation can pass a run while retaining a latent protocol-level flaw the workload never reached.
No certification
There is no certification program for Lakestream, and running this harness does not produce one. As the Verification index states, this documentation defines no cross-implementation certification suite and reports no passing conformance result; see Implementation status for what each implementation actually supports. streaming-proof is a tool an operator or implementer can run against their own deployment; the results belong to that deployment and that run, and a claim about them should carry the configuration that produced them.