slap_streams

slap_streams is a server for Durable Streams, an HTTP protocol for append-only streams of messages. Clients read a stream from an offset, and resume from their last offset after a disconnect or restart.

slap_streams stores streams in SlateDB, and acknowledges an append only once it is durably stored in object storage. It passes the official Durable Streams conformance suite, and is built on top of slap_cluster and slap_slatedb.

Use slap_streams to deliver a sequence of messages to clients. For example, an application can append a document's changes to a stream, and each client keeps 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.

slap_streams runs inside your application, on one node or several. You start Slap.Streams.Cluster, call Slap.Streams in-process, and can mount Slap.Streams.HTTP.Router in your Plug pipeline. Streams are divided among SlateDB databases called shards, and a call for a stream goes to the node that owns its shard.

The slap package runs slap_streams as a standalone server (mix slap.server --streams --store memory).

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

About Slap {: .info}

This package is part of Slap.

Slap provides Elixir bindings for SlateDB (an embedded key-value engine optimized for storing data in object storage), as well as services built on top of these bindings. The services run inside your application, on one node or several.

Services:

Building blocks:

Standalone server:

Installation

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

defp deps do
  [
    {:slap_streams, "~> 0.2.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. Offsets in Slap.Streams are these integers; over HTTP they use the official wire format, and Slap.Streams.Offset converts between the two.

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}
]

Bandit serves the router at /v1/stream/. A {:local, dir} store is for development, because it never deletes old manifests (see Storage after deletes); Using S3 describes the alternative.

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.

Mounting the router in an application

A stream's name is its full request path, including any path the router is mounted under. Served at the root, the stream at /v1/stream/chat/1 is Slap.Streams.head("/v1/stream/chat/1") in-process. Mounted under /streams, the same stream is at /streams/v1/stream/chat/1, and that is its name.

The router reads the request body itself, so it must run before Plug.Parsers. After it, a request whose content type Plug.Parsers parses reaches the router with its body already read. With Phoenix's default parsers (JSON, URL-encoded and multipart), an append to a JSON stream then gets 400 ("empty body not allowed"). A Phoenix router runs after the endpoint's Plug.Parsers, so mount the Streams router in the endpoint, with a function plug placed before plug Plug.Parsers:

# In MyAppWeb.Endpoint, before `plug Plug.Parsers`:
plug :streams

@streams Slap.Streams.HTTP.Router.init([])

defp streams(%Plug.Conn{path_info: ["streams" | rest]} = conn, _opts),
  do: conn |> Plug.forward(rest, Slap.Streams.HTTP.Router, @streams) |> halt()

defp streams(conn, _opts), do: conn

In a Plug.Router that has no Plug.Parsers before it, forward "/streams", to: Slap.Streams.HTTP.Router does the same.

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.