distkit
A toolkit of distributed systems primitives for Rust, backed by Redis.
What is distkit?
distkit provides building blocks for distributed applications. It ships
distributed counters (strict and lax), instance-aware counters, distributed
locks (Mutex, RwLock), and rate limiting, all backed by Redis.
Documentation and guides: https://distkit.davidoyinbo.com
Features
- StrictCounter -- every operation executes a Redis Lua script atomically. Reads always reflect the latest write. Best for billing, inventory, or anything where accuracy is critical.
- LaxCounter -- buffers increments in memory and flushes to Redis every ~20 ms. Sub-microsecond latency on the hot path. Best for analytics and high-throughput metrics.
- Instance-aware counters -- each running instance owns a named slice of the total, with automatic cleanup of contributions from instances that stop heartbeating.
- Mutex / RwLock (opt-in
lockfeature) -- Redis-backed distributed locks mirroringtokio::sync::Mutex/tokio::sync::RwLock. RAII guards, background lease refresh, writer-preferring reader-writer locking. - Rate limiting (opt-in
trypemafeature) -- sliding-window rate limiting with local, Redis-backed, and hybrid providers. Supports absolute and probabilistic suppression strategies. - Safe by default --
#![forbid(unsafe_code)], no panics in library code.
Feature flags
| Feature | Default | Description |
|---|---|---|
counter |
yes | Distributed counters (StrictCounter, LaxCounter) |
instance-aware-counter |
no | Per-instance counters (StrictInstanceAwareCounter, LaxInstanceAwareCounter) |
lock |
no | Distributed locks (Mutex, RwLock) |
trypema |
no | Rate limiting via the trypema crate |
Installation
Or add to Cargo.toml:
[]
= "0.7"
To enable instance-aware counters or rate limiting:
[]
= { = "0.7", = ["instance-aware-counter", "trypema"] }
Counters and locks require Redis 5.0+. Trypema's local provider does not contact Redis; its Redis and hybrid providers require Redis 7.2+.
Quick start
use ;
async
Counter types
StrictCounter
Every call is a single Redis round-trip executing an atomic Lua script. The counter value is always authoritative.
let key = try_from?;
strict.inc.await?; // HINCRBY via Lua
strict.set.await?; // HSET via Lua
strict.del.await?; // HDEL, returns old value
strict.clear.await?; // DEL on the hash
Conditional writes use CounterComparator and return (new, old). When the
comparison fails, new == old.
use CounterComparator;
strict.set.await?;
assert_eq!;
assert_eq!;
Batch increments follow the same rules and preserve input order.
let results = strict
.inc_all_if
.await?;
assert_eq!;
LaxCounter
Writes are buffered in a local DashMap and flushed to Redis in batched
pipelines every allowed_lag (default 20 ms). Reads return the local view
(remote_total + pending_delta), which is always consistent within the same
process.
let key = try_from?;
lax.inc.await?; // local atomic add, sub-microsecond
let val = lax.get.await?; // reads local state, no Redis hit
A background Tokio task handles flushing. It holds a Weak reference to the
counter, so it stops automatically when the counter is dropped.
Choosing a counter
StrictCounter |
LaxCounter |
StrictInstanceAwareCounter |
LaxInstanceAwareCounter |
|
|---|---|---|---|---|
| Consistency | Immediate | Eventual (default: ~20 ms lag) | Immediate | Eventual (flush_interval lag) |
inc latency |
Redis round-trip | Sub-microsecond (warm path) | Redis round-trip | Sub-microsecond (warm path) |
| Redis I/O | Every operation | Batched on interval | Every inc |
Batched on interval |
set / del |
Immediate | Immediate | Immediate (bumps epoch) | Flushes pending delta, then immediate |
| Per-instance tracking | No | No | Yes | Yes |
| Dead-instance cleanup | No | No | Yes | Yes |
| Feature flag | counter (default) |
counter (default) |
instance-aware-counter |
instance-aware-counter |
| Use case | Billing, inventory, exact global count | Analytics, high-throughput metrics | Connection counts, exact live metrics | High-frequency per-node throughput metrics |
Instance-aware counters
Enable the instance-aware-counter feature:
[]
= { = "0.7", = ["instance-aware-counter"] }
Instance-aware counters track each running instance's contribution separately.
The cumulative total is the sum of all live instances. When an instance stops
heartbeating for longer than dead_instance_threshold_ms (default 30 s), its
contribution is automatically subtracted from the cumulative on the next
operation by any surviving instance.
This makes them well-suited for:
- Connection pool sizing -- each server reports its active connection count; the cumulative is the cluster-wide total.
- Live session counting -- contributions disappear naturally when a node restarts or crashes.
- Per-node metrics -- see both the global total and each instance's slice.
Conditional instance-aware writes follow the same rule set:
inc_ifandset_ifcompare against the cumulative total.set_on_instance_ifcompares against the calling instance's slice.- Failed comparisons return the current
(cumulative, instance_count)unchanged.
StrictInstanceAwareCounter
Every call is immediately consistent with Redis. set and del bump a
per-key epoch that causes stale instances to reset their stored count on
their next operation, preventing double-counting.
use ;
use DistkitRedisKey;
let client = open?;
let conn = client.get_connection_manager.await?;
let prefix = try_from?;
let counter = new;
let key = try_from?;
// Increment this instance's contribution; returns (cumulative, instance_count).
let = counter.inc.await?;
// Decrement this instance's contribution.
let = counter.dec.await?;
// Read without modifying.
let = counter.get.await?;
// Set this instance's slice to an exact value without bumping the epoch.
let = counter.set_on_instance.await?;
// Set the global total to an exact value and bump the epoch.
let = counter.set.await?;
// Remove only this instance's contribution.
let = counter.del_on_instance.await?;
// Delete the key globally and bump the epoch.
let = counter.del.await?;
Dead-instance cleanup
Each instance sends a heartbeat on every operation. If a process silently dies, surviving instances automatically remove its contribution the next time any of them touches the same key.
use ;
use DistkitRedisKey;
let client = open?;
let conn1 = client.get_connection_manager.await?;
let conn2 = client.get_connection_manager.await?;
let prefix = try_from?;
let key = try_from?;
let opts = ;
let server_a = new;
let server_b = new;
server_a.inc.await?; // cumulative = 10
server_b.inc.await?; // cumulative = 15
// server_a goes offline. After 30 s, server_b's next call removes its
// contribution automatically.
let = server_b.get.await?; // total = 5 once cleaned up
LaxInstanceAwareCounter
A buffered wrapper around StrictInstanceAwareCounter. inc calls accumulate
locally and are flushed to the strict counter in bulk every flush_interval
(default 20 ms). Global operations (set, del, clear) flush any pending
delta first, then delegate immediately.
Use this when you have many inc/dec calls per second and can tolerate a
small consistency lag.
use ;
use DistkitRedisKey;
use Duration;
let client = open?;
let conn = client.get_connection_manager.await?;
let prefix = try_from?;
let counter = new;
let key = try_from?;
// Returns the local estimate immediately — no Redis round-trip on warm path.
let = counter.inc.await?;
// Decrement also stays local until flushed.
let = counter.dec.await?;
// get() also returns the local estimate (cumulative + pending delta).
let = counter.get.await?;
Distributed locks
Enable the lock feature for Redis-backed Mutex and RwLock:
[]
= { = "0.7", = ["lock"] }
Both mirror the surface of tokio::sync::Mutex / tokio::sync::RwLock. The
guards hold no inner data — they are pure access tokens. A held lock renews its
lease in the background (every ttl/3) and releases on drop, with an explicit
awaitable release() for callers who want to observe the final state. Each guard
also reports get_on_attempt — the zero-based acquire poll that won the lock
(0 on the first try, higher under contention). RwLock is writer-preferring (a
waiting writer blocks new readers).
use ;
async
Acquire forms per lock: waiting (lock / read / write), non-blocking
(try_lock / try_read / try_write), time-bounded
(try_lock_with_timeout / try_read_with_timeout / try_write_with_timeout), and
retry-bounded (try_lock_with_retries / try_read_with_retries /
try_write_with_retries). The bounded forms poll at the lock's configured
retry_interval; the retry-bounded forms give up with LockError::RetriesExhausted
after max_retries retries. (The older try_*_for(timeout, retry_interval) forms
are deprecated.) Tune ttl, max_wait, retry_interval, owner_id, and
namespace via LockOptions or LockOptions::builder.
Rate limiting (trypema)
Enable the trypema feature to access sliding-window rate limiting.
Trypema documentation website: https://trypema.davidoyinbo.com
[]
= { = "0.7", = ["trypema"] }
All public types from trypema 2 are re-exported
under distkit::trypema. Trypema 2 constructs each provider independently:
- Sliding-window rate limiting with configurable window size and rate.
- Three providers -- local (in-process), Redis-backed (distributed), and hybrid (local fast-path with periodic Redis sync).
- Two strategies -- absolute (binary allow/reject) and suppressed (probabilistic degradation that smoothly ramps rejection probability).
Local rate limiting
use ;
let provider = builder
.window_size
.bucket_size
.build
.unwrap;
let rate = per_second_or_panic;
match provider.absolute.inc
build() returns an Arc and starts stale-state cleanup. Use
.disable_cleanup() while building to opt out, or the provider's idempotent
start_cleanup_loop() and stop_cleanup_loop() methods after construction.
Redis-backed and hybrid rate limiting
For distributed enforcement across multiple processes or servers, construct a Redis or hybrid provider with a connection manager. Redis-backed providers require Redis 7.2 or newer.
use ;
let client = open?;
let conn = client.get_connection_manager.await?;
let window = minutes?;
let bucket = milliseconds?;
let redis_provider = builder
.window_size
.bucket_size
.build?;
let hybrid_provider = builder
.window_size
.bucket_size
.sync_interval
.build?;
let key = try_from?;
let rate = per_second?;
// Distributed absolute enforcement
let decision = redis_provider.absolute.inc.await?;
// Local fast-path with periodic Redis synchronization
let decision = hybrid_provider.absolute.inc.await?;
Trypema 2 removed RateLimiter, RateLimiterOptions, and the provider option
structs. It also renamed WindowSizeSeconds to WindowSize, RateGroupSizeMs
to BucketSize, SuppressionFactorCacheMs to
SuppressionFactorCachePeriod, and SyncIntervalMs to SyncInterval.
See the trypema documentation for full API details and advanced configuration.
Development
Prerequisites
- Rust (latest stable)
- Docker (for the test Redis instance)
Commands
Tests and benchmarks require the REDIS_URL environment variable.
The make targets set this automatically.
License
MIT