A multi-tenant streaming service in Rust, built on SlateDB — an object-store-native LSM — so that the object store is the only stateful tier. Appends are acknowledged only after their bytes are durable in object storage; servers are stateless and hold only caches and in-flight buffers; every stream is encrypted with a per-stream key the service never persists.
The service exposes two HTTP surfaces:
/v1/streams/{name}— the Prisma product API (the primary product). Typed collections with routing-key records, producer sessions, consumer groups, watches, forks, and a seal lifecycle. UsesPrisma-*headers and product cursors. This is what the@prisma/streamsSDK speaks and what applications should build against./v1/stream/{name}— the pinned Durable Streams standards surface. The append-only default-key sequence of the Durable Streams protocol, preserved byte-for-byte and verified against the upstream conformance suite. It is the raw, standards-compliant view of a collection's default routing key — use it for protocol interoperability, not as the product API. Under shared-cell enforcement this surface is INTERNAL-ONLY (docs/MULTITENANCY.md §14.3): it takes fleet identity, never customer tokens, and the product surface is the only public customer API.
Pre-launch posture. This is a clean, single-format cutover: there are no legacy decoders, translators, aliases, or dual layouts. Removed experimental names (
Stream-Encryption-Key,Stream-Key, stream profiles) are rejected on the wire, never translated. Descriptors are written at oneLAYOUT_VERSION. Historical designs live in docs/history/.
- Two coherent surfaces, one engine. The product route's default routing key is the raw route's sequence — the same records, the same durability, seen two ways.
- Committed-before-ack durability. Concurrent appends bundle into shared WAL PUTs; the ack races nothing — a 2xx means the bytes are in object storage. Every response whose truth depends on state (duplicates, idempotent closes, producer/sequence conflicts, seal fences) waits behind that same durability barrier.
- Tenant isolation by cryptography. Each collection's data is AES-GCM
encrypted under a caller-supplied key (
Prisma-Encryption-Key) the service never stores. Backups are ciphertext. - Product lifecycle. Typed creation documents, producer sessions with exactly-once semantics, consumer groups (pull/settle, leases, per-key FIFO, DLQ), watches with signed observation URLs, forks with resumable lineage, and a generation-fenced seal state machine (Open → Sealing → Sealed).
- Coordination-free horizontal scale. Shard placement derives from fleet heartbeats via rendezvous hashing; correctness never depends on routing being right — object-store CAS fencing makes a stale owner's writes fail, not corrupt. Hot routing keys split into real physical child segments.
| metric | value |
|---|---|
| fleet of 4 sustained | ~1,250 req/s avg, peaks 2,700+ req/s |
| durable-ack p50 under load | 50–65 ms (25–50 ms WAL flush + Tigris PUT) |
| single instance, direct path | ~1,180 req/s max observed |
| 2 h soak | flat p50 at saturation, zero deaths |
| chaos (kill N−2 under load) | survivors absorb, zero data loss |
Full history: BENCHMARKS.md, EXPERIMENT-PILOT.md, REPORT.md.
The @prisma/streams SDK is the canonical getting-started path;
its README is the tutorial. In brief:
import { StreamsClient } from "@prisma/streams";
const client = new StreamsClient({ url, token });
const orders = await client.createStream("orders", {
encryptionKey, // 32 bytes, base64url — never stored
format: { kind: "json" },
watches: [{ name: "by-customer", fields: ["/customerId"] }],
});
await orders.append({ customerId: "c1", total: 42 }, { routingKey: "c1" });
for await (const record of orders.subscribe()) { /* ... */ }cargo build --release
./target/release/s3lite --listen 127.0.0.1:9500 --latency-ms 5 &
./target/release/streams-slate \
--listen 127.0.0.1:8090 \
--s3-endpoint http://127.0.0.1:9500 --bucket streams --region auto \
--access-key-id test --secret-access-key test \
--initial-shards 4 --auth-token devtoken &
KEY=$(./target/release/streams-keys generate) # 32 bytes, base64url
# Product route: create a typed collection, append under a routing key, read.
curl -X PUT http://127.0.0.1:8090/v1/streams/orders \
-H "authorization: Bearer devtoken" -H "Prisma-Encryption-Key: $KEY" \
-H 'content-type: application/json' -d '{"format":{"kind":"json"}}'
curl -X POST http://127.0.0.1:8090/v1/streams/orders/records \
-H "authorization: Bearer devtoken" -H "Prisma-Encryption-Key: $KEY" \
-H 'Prisma-Routing-Key: c1' -H 'content-type: application/json' -d '{"n":1}'
curl "http://127.0.0.1:8090/v1/streams/orders/records?routingKey=c1" \
-H "authorization: Bearer devtoken" -H "Prisma-Encryption-Key: $KEY"
# The raw Durable Streams surface is the DEFAULT routing key's view —
# records appended under other routing keys (like c1 above) do not
# appear there. Append one default-key record (no Prisma-Routing-Key),
# then read it back through the standards route:
curl -X POST http://127.0.0.1:8090/v1/streams/orders/records \
-H "authorization: Bearer devtoken" -H "Prisma-Encryption-Key: $KEY" \
-H 'content-type: application/json' -d '{"n":2}'
curl http://127.0.0.1:8090/v1/stream/orders \
-H "authorization: Bearer devtoken" -H "Prisma-Encryption-Key: $KEY"The SDK is dependency-free and derives watch keys via WebCrypto.
- Node 18 and 22 — gated in CI (built, packed, installed from the tarball, smoke-tested end to end).
- Bun and Deno — gated in CI.
- Browsers — expected to work (WebCrypto +
fetchonly), not yet verified; a browser integration gate is required before browser support is claimed.
RUNBOOK.md is the operator manual — building (including the mandatory x86_64-musl cross-compile for Prisma Compute), the configuration reference, fleet mode and autoscaling, admission control, debug endpoints, deployment, monitoring baselines, and a symptom→cause→fix matrix.
OPERATIONS.md covers the durability/security posture: object-store requirements, backup/PITR, tenant identity and key custody, SLOs.
scripts/release-provenance.sh binds a release report to the exact artifact
(server commit, SlateDB pin, SDK tarball SHA, layout version, conformance
pin, DST scenario count).
- Rust quality policy — adopted coding standard, pinned tools, local/CI gates and invariant verification.
| document | what it answers |
|---|---|
| sdk/README.md | getting started with the product API |
| docs/RELEASE-PRODUCT-SURFACE.md | the product surface: gates, audit rounds, contracts |
| SPEC.md | architecture, decision log, guarantees |
| RUNBOOK.md | build, run, deploy, scale, monitor, debug |
| CONFORMANCE.md | the pinned Durable Streams suite; how to run it |
| docs/ROUTING-V3.md | routing keys, physical scaling, postings, cost |
| OPERATIONS.md | provider requirements, backup/PITR, identity, SLOs |
| docs/OBSERVABILITY-BILLING.md | billing/usage/telemetry design (normative) |
| docs/OBSERVABILITY-BILLING-STATUS.md | billing implementation matrix + open gates |
| SECURITY.md | key custody, watch capabilities, tenant isolation |
| docs/dst/ | the deterministic-simulation program + scenario catalogue |
| docs/history/ | removed designs (profiles, …) — provenance only |
| repro-edge-404/ | the Compute edge-publication platform ticket |
Agents and contributors start with AGENTS.md (how to work and verify here) and docs/README.md (a map of every document).
src/ server crate (streams-slate) + bins
http.rs both HTTP surfaces + admission + debug endpoints
product.rs the /v1/streams product surface (lifecycle, consumers, watches)
shard.rs shard engine: commit pipeline, group PUTs, watermarks, fences
registry.rs the descriptor control plane (incarnations, seal state)
scaler3.rs routing-key distribution sketches + physical split/merge
history.rs history tier + absorber
crypto.rs stream-key envelope (AES-GCM)
fleet.rs heartbeats, load vector, desired-count computation
dst/ deterministic simulation tests
(all 51 modules: python3 scripts/dev/impact.py --codemap)
sdk/ @prisma/streams TypeScript client SDK (canonical entry point)
conformance/ the pinned Durable Streams suite runner
scripts/ field gate, release provenance, analysis
docs/ routing/cost/soak campaigns, RELEASE-PRODUCT-SURFACE, dst/, history/
repro-edge-404/ the Compute edge-publication reproduction package
Apache-2.0 (see LICENSE).