slap_streams

slap_streams is a Durable Streams server that persists data using SlateDB, built on slap_cluster and slap_slatedb. It is meant to be embedded: your application starts Slap.Streams.Cluster, calls Slap.Streams in-process, and can mount Slap.Streams.HTTP.Router in its Plug pipeline.

Use slap_streams for named, append-only streams of messages that clients read from a saved offset, so that they can catch up after a disconnect or restart. For example, an application can append a document's changes to a stream, and each client continues reading from its last offset. The package handles durable appends, reads by offset, waits for new messages, producer deduplication, stream expiry, and forks. It does not define your message format or apply messages to your application's state.

Also see the slap package, which provides a Mix task for running slap_streams as a standalone server (mix slap.server --streams --store memory). It contains the code for running the official Durable Streams conformance suite and benchmarks, the crash and cluster tests, and the HTTP load benchmark.

Section numbers (§) in the API documentation refer to the Durable Streams specification.

Installation

Add slap_streams to the dependencies in your application's mix.exs:

defp deps do
[
{:slap_streams, "~> 0.1.0"}
]
end

Run mix deps.get. API documentation is on HexDocs.

Example

In an application that depends on slap_streams, run this in iex -S mix with a fresh local directory. The in-process API is Slap.Streams:

iex> {:ok, _} = Slap.Streams.Cluster.start_link(store: {:local, "/tmp/slap-streams"}, shards: 8); :ok
:ok
iex> {:ok, :created, %{next_offset: 0}} = Slap.Streams.create("/chat/1", content_type: "application/json"); :created
:created
iex> {:ok, %{result: :appended, next_offset: 12}} =
...> Slap.Streams.append("/chat/1", ~s([{"m": 1}]), producer: {"client-a", 0, 0}); {:appended, 12}
{:appended, 12}
iex> {:ok, %{messages: messages, up_to_date: true}} = Slap.Streams.read("/chat/1", 0); messages
[{0, "{\"m\": 1}"}]

The functions return protocol results, not HTTP responses; Slap.Streams documents the HTTP status for each result. The next offset is 12 because the JSON element occupies eight bytes and each message adds four bytes of overhead.

Embedding

Add Bandit to the host application's dependencies. Then add these children to its supervision tree, with the cluster first:

children = [
{Slap.Streams.Cluster, store: {:local, "/tmp/slap-streams"}, shards: 8},
{Bandit, plug: Slap.Streams.HTTP.Router, port: 4437}
]

This serves /v1/stream/ directly. An application can instead mount the router with Plug.Router.forward/2. A stream's name is its full HTTP request path, so the stream at /v1/stream/chat/1 is Slap.Streams.head("/v1/stream/chat/1") in-process. With a forwarded mount, the mount path is part of the name.

To use an application-owned cluster module:

defmodule MyApp.StreamsCluster do
use Slap.Streams.Cluster, otp_app: :my_app
end
children = [{MyApp.StreamsCluster, store: {:local, "/tmp/slap-streams"}, shards: 8}]
{:ok, :created, _} = Slap.Streams.create("/notes", cluster: MyApp.StreamsCluster)

Pass cluster: MyApp.StreamsCluster to the router and to any Slap.Yjs.Docs instance using it. Run one Streams cluster per VM: metrics and the HTTP request-body budget are shared across the VM. A KV cluster can run on the same VM.

Using S3

Create the bucket first and configure AWS_REGION and credentials in the environment used to start the application, for example with AWS_ACCESS_KEY_ID and AWS_SECRET_ACCESS_KEY. Then, in iex -S mix in an application that depends on slap_streams, start a cluster on a new prefix in that bucket:

iex> {:ok, _} = Slap.Streams.Cluster.start_link(
...> store: {:url, "s3://my-bucket/my-app/streams"},
...> shards: 8
...> ); :ok
:ok
iex> {:ok, :created, _} = Slap.Streams.create("/notes", content_type: "text/plain"); :created
:created
iex> {:ok, %{result: :appended}} = Slap.Streams.append("/notes", "hello"); :appended
:appended
iex> {:ok, %{messages: messages}} = Slap.Streams.read("/notes", 0); messages
[{0, "hello"}]

Use the same store URL and shard count when restarting the application. In a supervised application, put Slap.Streams.Cluster in the supervision tree as shown above. The example uses the default local placement on one node; see Several nodes to run on several nodes.

How it works

HTTP

Slap.Streams.HTTP.Router is a Plug over Slap.Streams. Streams live under /v1/stream/, and a stream's name is the request path. By default the router matches the official Go server: the same headers, status codes, cursors, CORS headers, and security headers. For HTTP request examples, see slap.

Router options

Option Default Effect
:long_poll_timeout 30,000 ms How long a long-poll waits for data.
:sse_timeout 60,000 ms How long an SSE response stays open.
:max_read 1 MiB About how many bytes one read returns.
:max_body 64 MiB The largest request body. A larger Content-Length gets 413 before the body is read.
:max_path 512 bytes The longest stream path (414 above it).
:max_buffered 512 MiB Request-body bytes the VM holds at once (503 above it).
:trust_forwarded false Build a created stream's Location from X-Forwarded-Proto and X-Forwarded-Host.
:private_cache false Send Cache-Control: private instead of public.
:cors "*" The allowed origin, or false to omit CORS headers.

Authentication

The router does not authenticate requests or limit their rate. Put both in your application's pipeline or in a proxy in front of it. Two things to cover:

Backpressure

An append gets 503 with Retry-After: 1 if it would take a stream's bytes in flight (written but not yet durable) over 64 MiB, or its shard's over 256 MiB. A stream with nothing in flight always accepts one request. Set the limits with the cluster child_options keys :max_inflight_bytes_per_stream and :max_inflight_bytes_per_shard.

Metrics

Slap.Streams.Telemetry emits telemetry events for:

Slap.Streams.Metrics keeps these, together with slap_cluster's durability lag and shard events, as Prometheus metrics; its documentation lists them. Slap.Streams.HTTP.Metrics is a Plug that serves them in the Prometheus text format. Mount it where Prometheus can reach it and the public cannot.

Lifecycle

Storage after deletes

Deleting rows writes tombstones. SlateDB reclaims their space only after compaction and garbage collection; old sorted runs and checkpoints can delay reclamation. The local file store does not collect old manifests, so use an object store when reclaiming deleted data matters. bench/storage_soak.exs measures this behavior.

Several nodes

With a placement strategy, slap_cluster spreads the shards over several nodes and moves them when nodes join, leave, or fail. Connect the nodes with distributed Erlang (for example with libcluster). Any node can serve any request:

{Slap.Streams.Cluster, store: {:url, "s3://streams/ds"}, shards: 64,
strategy: {Slap.Cluster.Strategy.ObjectLease, lease_ttl: 15_000}}

Placement strategies and consistency

Any of slap_cluster's strategies works (--placement for mix slap.server):

Without leases, reads are not linearizable. A node that was paused, or that disagrees with others about which nodes are connected, may serve reads from its old state until it finds that SlateDB has fenced it (about a second). Its writes fail with 503, because SlateDB rejects writes from a fenced writer. Before it reports a producer sequence gap, it confirms its ownership with a write, so that reply also gets 503. The Jepsen suite checks both kinds of placement.

Storage

Offsets use the official wire format (%016d_%016d) and advance by four bytes plus the body length per message. An append stores its messages, tail, producer state, and metadata changes in one SlateDB batch. JSON mode keeps each array element's original bytes.

Tests

mix test # SLAP_STREAMS_PROP_RUNS=1000 for more model sequences

The model test covers stream operations, crashes, and clock changes against a pure protocol model. Set SLAP_STREAMS_PROP_RUNS to run more sequences; SLAP_STREAMS_PROP_STATS=1 prints outcome counts.

The tests that need the server as an OS process live in slap: the official Durable Streams conformance suite and benchmarks, the crash tests and the cluster tests.

Benchmarks

mix run bench/append_latency.exs
mix run bench/read_scaling.exs

These measure the library in-process. slap runs the official Durable Streams HTTP benchmarks and an HTTP append load test against a running server.

bench/append_latency.exs measures durable append latency on the in-memory store. Set SLAP_BENCH_FLUSH_INTERVAL to change SlateDB's flush interval.

bench/read_scaling.exs times first server reads, warm random reads, and narrow and broad lists as stream count grows. It also times reads at the start, middle, and tail of a growing stream. It uses one shard, flushes data before timing, and reports p50, p99, and maximum latency. The first server reads include process startup and state loading; the other operations use warmed paths. The default store is a temporary local directory, removed after the run. Use a fresh object-store prefix to measure that backend:

mix run bench/read_scaling.exs -- --streams 1000 --messages 1024 --samples 200 \
--cache disabled --store s3:s3://bucket/read-bench

Increase --streams, --messages, or --message-bytes to test larger stores. --cache disabled turns off SlateDB's in-memory cache; operating-system and object-store caches may still affect results. --page-bytes and --flush-interval control read size and SlateDB's flush interval. --store memory avoids local and object-store I/O. Object-store data remains after the run and must be removed separately.

License

Slap is released under the terms of the Apache License 2.0.

Copyright (c) 2026, Michael Russo.