slap_cluster

slap_cluster manages a fixed number of SlateDB databases ("shards") in one store, using slap_slatedb. It assigns each shard to a node, opens its database there, routes calls to its owner, and reports when writes become durable. If SlateDB fences a writer, the shard stops and the placement strategy can assign it elsewhere. The application provides the processes that use each shard and decides what they store.

Use slap_cluster when you are building a sharded service with your own data model and need to route operations to shard owners or move shards between nodes after failures. It manages one SlateDB writer per shard; your shard children implement the service's reads, writes, and responses. For one database without shard placement, use slap_slatedb directly. If your data fits the Streams or partitioned KV APIs, use slap_streams or slap_kv, which already build on this package.

Installation

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

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

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

Example

In an application that depends on slap_cluster, define a cluster module:

defmodule MyApp.Cluster do
  use Slap.Cluster, otp_app: :my_app
end

Configure it in config/config.exs. This local configuration sketch assumes the application provides MyApp.ShardChildren.child_specs/1 to start its processes for each shard:

config :my_app, MyApp.Cluster,
  store: {:local, "/tmp/my-app-cluster"},
  shards: 8,
  shard_children: {MyApp.ShardChildren, :child_specs, []}

Add MyApp.Cluster to the application's supervision tree after defining its shard children.

In a shard child, ctx is its context, ops is a SlateDB write batch, and ref is a token chosen by the child to identify the notification. The child can request a durability notification for a write:

{:ok, seq} = Slap.SlateDB.write(ctx.db, ops)
Slap.Cluster.notify_when_durable(ctx, seq, {:ack, ref})

The child later receives {:slap_cluster_durable, {:ack, ref}}. A caller can route an operation to the owner of a key's shard. Here MyApp.Streams.append/2 is an application function, and path and body are its arguments:

shard = MyApp.Cluster.shard_for("some/key")
{:ok, result} = MyApp.Cluster.call(shard, {MyApp.Streams, :append, [path, body]})

The outer {:ok, result} means routing succeeded. result is the application function's return value, including any {:error, reason} it returns.

Slap.Cluster's docs cover the options, shard children, durability and telemetry.

How it works

The cluster probes the store's create-if-absent behavior before opening shards. ObjectLease also requires If-Match support. Each shard has one database handle and one durability subscription. Application shard children stop before the handle closes. A fenced or crashed database stops its shard; the placement strategy decides where to reopen it. Reopening without that decision could fence a new owner. Children that need the stop reason use shard_status/1; their supervisor exit reason is :shutdown even when the shard was fenced.

When a child calls notify_when_durable/3, the shard compares the write's sequence number with SlateDB's current durable sequence number, not with the last durability message it handled. A write that is already durable is reported at once, even if that message is still in the shard's mailbox. shard_for/1 uses xxh64(key) mod shards; changing the shard count or hash moves existing keys.

Several nodes

To run a cluster on several nodes, use a shared object store and a placement strategy. For an S3-compatible store such as RustFS, aws_endpoint sets its endpoint; the store must support conditional writes:

config :my_app, MyApp.Cluster,
  store:
    {:url, "s3://bucket/prefix",
     aws_endpoint: "http://rustfs:9000", aws_allow_http: true},
  shards: 64,
  settings: %{flush_interval: "10ms"},
  cache: [capacity_bytes: 2 * 1024 * 1024 * 1024],
  shard_children: {MyApp.ShardChildren, :child_specs, []},
  strategy: {Slap.Cluster.Strategy.ObjectLease, lease_ttl: 15_000}

Nodes connect through distributed Erlang. ObjectLease stores conditional leases and node heartbeats in the bucket. An owner renews its leases and stops a shard if renewal fails past its local deadline, before another node can claim it. A paused owner that misses the deadline is fenced by SlateDB when a new owner opens the database. A clean stop closes shards before freeing their leases. Self-fencing runs separately from the renewal loop so a blocked store call cannot delay it.

call/4 routes to the owner over :erpc. If the reply is :not_owner or :unassigned, or the owner cannot be reached before the call is sent, call/4 refreshes the placement and retries. A call that was already sent is not retried, because it may have run; the caller receives {:error, {:erpc, :noconnection}}. Lease data does not create node-name atoms, so peers must be discovered or configured before routing to them.

Placement strategies

ObjectLease requires clocks within :max_clock_skew_ms (1 s by default); the local file store cannot provide If-Match. Lower net_ticktime with Distributed if failure detection must be faster.

Tests

mix test                                                     # local directories
SLAP_TEST_S3_ENDPOINT=http://127.0.0.1:9000 mix test      # also 64 shards on RustFS

The tests exercise durability, fencing, placement and routing with a counter server on each shard. Multi-node tests use separate OS processes and a shared store; SLAP_TEST_S3_ENDPOINT enables the RustFS cases. See test/support/failover.ex for the fault test harness.

License

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

Copyright (c) 2026, Michael Russo.