Delivery Semantics and Consistency Model
The behavioral contract: exactly what Felix promises about delivery, ordering, durability, and consistency — and, just as deliberately, what it does not. Applications should rely on what is written here and nothing stronger.
Philosophy: Explicit Over Implicit
Section titled “Philosophy: Explicit Over Implicit”Felix makes trade-offs explicit rather than hiding them behind ambiguous guarantees. Every semantic choice has observable behavior that can be tested and reasoned about.
Pub/Sub Delivery Semantics
Section titled “Pub/Sub Delivery Semantics”Delivery Guarantees
Section titled “Delivery Guarantees”A stream’s guarantee follows from how it is registered.
An ephemeral stream is at-most-once. This is the default:
- Messages are delivered to subscribers zero or one time
- No retries or redelivery
- No acknowledgements from subscribers
- Slow subscribers may drop messages without notification
A durable stream is at-least-once. Every record is written to disk before
the publish is acknowledged, and a subscriber replays from any retained offset,
so a record survives a broker restart and can be read again. A stream declared
Quorum waits for a majority of its replicas before acknowledging, so the
record also survives losing the broker that accepted it.
A consumer group is at-least-once, and redelivers. A record handed to a consumer that does not answer is handed to another once the visibility timeout lapses. See Projections.
At-most-once is appropriate for:
- Real-time signals where latest value matters most
- High-frequency metrics and telemetry
- Workloads where occasional loss is acceptable
- Applications that implement their own deduplication
Example at-most-once workload:
// Real-time sensor data where latest reading matters mostlet mut subscription = client.subscribe("acme", "sensors", "temperature").await?;
while let Some(event) = subscription.next_event().await? { // Process latest temperature reading // If we miss a reading, the next one will arrive soon update_dashboard(event.payload);}At-least-once, and why there is no third guarantee
Section titled “At-least-once, and why there is no third guarantee”At-least-once is implemented, two ways, and both are described above: replay a durable stream from a checkpointed offset, or consume through a consumer group, which requires an acknowledgement per record, redelivers anything unanswered once its visibility timeout lapses, and dead-letters a record that has been attempted too many times.
Idempotent producers are implemented; exactly-once delivery is not. A
producer takes an id from the broker and numbers its batches, and a batch
re-sent after a lost acknowledgement lands once: the shard’s leader answers a
sequence it already holds from memory rather than appending it again. That
closes the ambiguous-outcome gap on the publish side (ClusterClient::idempotent_producer,
negotiated as FEATURE_IDEMPOTENT_PRODUCER). It does not make delivery
exactly-once: a consumer can still see a record twice on redelivery, and
end-to-end exactly-once would also need transactional coordination across the
log and the consumer’s own state, and deduplication on receive — which has to
live in the application regardless, because the application is the only thing
that knows what makes two records the same. Deduplicate there, keyed on
something the record carries.
Consistency: how many brokers must hold it
Section titled “Consistency: how many brokers must hold it”A durable stream is replicated to a set of brokers — one leader and its
replicas. consistency on the stream decides how many of them must hold a
record before the publisher is told it is safe.
Leader — the default. The leader writes the record to its own log,
durably, and answers. Replication still happens; the acknowledgement simply
does not wait for it. One round trip.
Quorum — the leader writes durably, ships the record to its replicas
concurrently, and answers once a majority of the replica set, counting
itself, holds it. On a set of three that is two, so one unreachable replica
costs nothing — the leader is not waiting for all of them, only for enough.
Note what Quorum does not change. The record is written the same way, to the
same log, with the same fsync policy. What changes is what the acknowledgement
means:
A
Quorumacknowledgement survives losing the leader. ALeaderacknowledgement is a promise only that one broker can keep.
What each one costs
Section titled “What each one costs”Quorum costs latency, and it costs availability at the other end: a stream
that cannot reach a majority stops accepting writes rather than accepting
ones it might not keep. A publish with no reachable majority is refused, and a
refusal means “this cannot be vouched for” rather than “this did not
happen” — the record may well have landed on the leader. Retry through an
idempotent producer, which re-sends under the same sequence and cannot land
it twice.
Leader is one round trip instead of two, and it moves the moment you find out.
If the leader dies holding a record nothing else has, the control plane will not
promote a replica, because promoting one would open the shard without that
record and no reader could tell. The shard is left unavailable until the old
leader returns with its disk.
So the trade is not really safety against latency. Both refuse to lose an
acknowledged record; they differ in when you learn there is a problem —
Quorum at publish time, while you still hold the record, or Leader at
failover time, when the only copy is on a broker that is gone.
task cluster:consistency runs exactly
that: the same fault put to both, on a real three-node cluster.
Message Ordering
Section titled “Message Ordering”Within a stream: Ordering is preserved per publisher-broker-subscriber path.
graph LR
P[Publisher] -->|msg 1, 2, 3| B[Broker]
B -->|msg 1, 2, 3| S[Subscriber]
style P fill:#e3f2fd,stroke:#334155,color:#111827
style B fill:#fff9c4,stroke:#334155,color:#111827
style S fill:#c8e6c9,stroke:#334155,color:#111827
Guarantees:
- Messages from a single publisher to a stream arrive in send order
- A single subscriber sees messages in the order they were enqueued
- Order is preserved through batching and fanout
Across streams: No ordering guarantees.
graph LR
P[Publisher]
P -->|msg A| S1[Stream 1]
P -->|msg B| S2[Stream 2]
Sub[Subscriber]
S1 --> Sub
S2 --> Sub
Note[msg A and B may arrive in any order]
style P fill:#e3f2fd,stroke:#334155,color:#111827
style Sub fill:#c8e6c9,stroke:#334155,color:#111827
Example:
// Publish to two streamsuse felix_wire::AckMode;let publisher = client.publisher().await?;publisher .publish("acme", "prod", "user-login", login_event, AckMode::None) .await?;publisher .publish("acme", "prod", "audit-log", audit_event, AckMode::None) .await?;
// Subscribers to user-login and audit-log may see events in any relative orderOrdering within batches:
// Batch publish preserves order within the batchuse felix_wire::AckMode;let publisher = client.publisher().await?;let messages = vec![msg1, msg2, msg3];publisher .publish_batch("acme", "prod", "orders", messages, AckMode::PerBatch) .await?;
// Subscribers will see msg1, msg2, msg3 in that orderFanout Fairness
Section titled “Fanout Fairness”Felix enforces subscriber isolation: slow subscribers never block fast subscribers.
graph TB
P[Publisher] --> B[Broker]
B --> S1[Fast Subscriber]
B --> S2[Slow Subscriber]
B --> S3[Fast Subscriber]
S1 -->|Processing msgs 1-100| D1[Dashboard]
S2 -->|Still on msg 23, dropping msgs| D2[Slow System]
S3 -->|Processing msgs 1-100| D3[Analytics]
style S1 fill:#c8e6c9,stroke:#334155,color:#111827
style S2 fill:#ffccbc,stroke:#334155,color:#111827
style S3 fill:#c8e6c9,stroke:#334155,color:#111827
Isolation mechanism:
Each subscription maintains an independent buffer:
pub struct Subscription { buffer: BoundedQueue<Event>, // Per-subscription buffer event_stream: UnidirectionalStream, // Independent QUIC stream}Buffer behavior:
- Each subscriber has
subscriber_queue_capacitybuffer slots (default: 512) - When buffer fills, new events are dropped for that subscriber only
- Other subscribers continue receiving events normally
- A drop is not announced. For a durable stream a subscriber can detect one, because delivered records carry log offsets and a jump between consecutive offsets is exactly a drop.
Configuration:
# Broker configsubscriber_queue_capacity: 512 # Per-subscriber buffer sizesubscriber_writer_lanes: 4subscriber_lane_shard: autoA larger buffer tolerates more bursty subscribers, trading memory for burst tolerance:
subscriber_queue_capacity: 4096Publisher Backpressure
Section titled “Publisher Backpressure”Publisher behavior: Publishing never blocks on subscriber speed.
sequenceDiagram
participant P as Publisher
participant B as Broker Queue
participant F as Fanout Workers
participant S1 as Fast Subscriber
participant S2 as Slow Subscriber
P->>B: publish_batch
B-->>P: ack (immediate)
B->>F: dequeue for fanout
par Independent fanout
F->>S1: deliver events
and
F->>S2: deliver events (buffer fills, drops)
end
Note over P: Publisher never waits for subscribers
Publisher queue:
Publishers write to a bounded queue with configurable depth:
pub_queue_depth: 64 # Bounded publish queuepublish_queue_wait_timeout_ms: 2000 # Timeout if queue fullWhen the publish queue is full:
- New publishes block up to
publish_queue_wait_timeout_ms - After timeout, publish fails with error
- This indicates broker overload (too many publishes, insufficient workers)
Tuning publish pipeline:
# Increase parallelismpub_workers_per_conn: 4
# Increase buffer (trades latency for burst tolerance)pub_queue_depth: 256
# Faster timeout for fail-fast behaviorpublish_queue_wait_timeout_ms: 1000Disconnection Behavior
Section titled “Disconnection Behavior”Subscriber disconnects:
- Subscription is immediately removed from registry
- Buffered events for that subscriber are discarded
- No redelivery: a plain subscription has no record of what was handled. A consumer group does, and redelivers anything claimed but never acknowledged once its visibility timeout lapses
- Subscriber must re-subscribe. On a durable stream it can resume at a checkpointed offset rather than restarting at the tail; on an ephemeral one the tail is all there is
Publisher disconnects:
- In-flight publishes may be lost if not acknowledged
- No automatic retry or persistence of unacked publishes
- Application must handle reconnection and retry logic
Broker restarts:
- In-memory state is lost. Durable streams, the log-backed cache, and consumer-group positions are on disk and survive.
- Active subscriptions are terminated
- Clients detect connection loss and must reconnect
- No historical replay available
Cache Semantics
Section titled “Cache Semantics”Consistency Model
Section titled “Consistency Model”Felix cache provides eventual consistency with read-your-writes for single clients:
sequenceDiagram
participant C1 as Client 1
participant B as Broker Cache
participant C2 as Client 2
C1->>B: put(key=X, value=1)
B-->>C1: ok
C1->>B: get(key=X)
B-->>C1: value=1
Note over C2: Concurrent get may see old value briefly
C2->>B: get(key=X)
B-->>C2: value=1 (or old value)
Guarantees:
- Read-your-writes: Client sees its own writes immediately
- Monotonic reads: Client never sees older values after newer ones (single session)
- Eventual consistency: All clients eventually see the latest value
- No dirty reads: Clients never see partial or uncommitted writes
Not guaranteed:
- Linearizability across clients
- Causal consistency across keys
- Multi-key transactions
TTL and Expiration
Section titled “TTL and Expiration”TTL semantics:
// Store with 60-second TTLuse bytes::Bytes;client .cache_put( "acme", "prod", "session", session_id, Bytes::from(session_data), Some(60_000), ) .await?;
// Store without expirationclient .cache_put("acme", "prod", "config", config_key, Bytes::from(config_value), None) .await?;Expiration behavior:
- TTL countdown starts when
cache_putreturnsok - Expiration is lazy: checked on access, not proactively
- Expired entries return
nulloncache_get - Expired entries may occupy memory until accessed or evicted
Cache Scoping
Section titled “Cache Scoping”Cache entries are scoped to (tenant_id, namespace, cache_name, key):
// These are independent cache entries:client.cache_put_scoped("tenant1", "prod", "sessions", "user123", data).await?;client.cache_put_scoped("tenant1", "staging", "sessions", "user123", data).await?;client.cache_put_scoped("tenant2", "prod", "sessions", "user123", data).await?;Isolation guarantees:
- Different tenants cannot access each other’s cache entries
- Different namespaces within a tenant are isolated
- Keys are unique only within their (tenant, namespace, cache) scope
Eviction Policy
Section titled “Eviction Policy”Today: the in-memory cache evicts best-effort under pressure. The log-backed cache does not evict — it compacts, reclaiming superseded and expired records.
- No guaranteed LRU or LFU policy
- Eviction is opportunistic
- Applications should not rely on specific eviction order
Concurrency and Race Conditions
Section titled “Concurrency and Race Conditions”Concurrent writes to same key:
sequenceDiagram
participant C1 as Client 1
participant C2 as Client 2
participant B as Broker
par Concurrent puts
C1->>B: put(key=X, value=A)
and
C2->>B: put(key=X, value=B)
end
Note over B: Last write wins (undefined order)
C1->>B: get(key=X)
B-->>C1: value=A or value=B
Behavior: Last write wins, but order is undefined for concurrent writes.
No atomic operations:
- No compare-and-swap
- No atomic increment
- No multi-key transactions
Planned features:
- Conditional put (if-not-exists, if-match)
- Atomic increment/decrement
- Watch/notify on key changes
Cache vs. Pub/Sub Integration (Future)
Section titled “Cache vs. Pub/Sub Integration (Future)”Planned feature: Pub/sub invalidation for cache consistency.
// Publish invalidates cache entryclient.publish_with_invalidation("events", "user-updated", event, vec!["cache:sessions:user123"]).await?;
// Subscribers and cache both receive updateTenant and Namespace Model
Section titled “Tenant and Namespace Model”Existence Enforcement
Section titled “Existence Enforcement”Wire protocol requirement: All data-plane operations must specify tenant and namespace.
{ "type": "publish", "tenant_id": "acme-corp", "namespace": "production", "stream": "orders", "payload": "..."}Broker validation:
The broker enforces tenant/namespace existence:
- Broker syncs metadata from control plane
- Broker maintains local registry of valid tenant/namespace pairs
- Operations for unknown tenant/namespace are rejected with
errorresponse
sequenceDiagram
participant C as Client
participant B as Broker
participant CONTROLPLANE as Control Plane
CONTROLPLANE->>B: Sync metadata (tenants, namespaces)
C->>B: publish (tenant=unknown, ...)
B-->>C: error (unknown tenant)
C->>B: publish (tenant=acme, namespace=prod, ...)
B->>B: Validate against registry
B-->>C: ok
Authorization
Section titled “Authorization”Enforced. Tenant-scoped tokens are verified at the broker, and publish, subscribe and cache operations each check a permission before doing any work.
A forwarded publish is authorized twice — at the broker the client reached and again at the shard’s owner — so routing a request through the cluster does not launder the credential it arrived with.
Enforcement points:
- Publish operations, at ingress and at the owner
- Subscribe operations
- Cache operations
- Control plane operations
Per-tenant quotas are a separate thing and are not enforced; see below.
Quota Enforcement (Planned)
Section titled “Quota Enforcement (Planned)”Not implemented. Nothing limits what a tenant can publish, subscribe to, or cache. The shape it would take:
quotas: - tenant: acme-corp namespace: production publish_rate_limit: 10000/s subscribe_connections: 100 cache_memory: 10GB stream_retention: 7dConsistency Across Components
Section titled “Consistency Across Components”Broker Internal Consistency
Section titled “Broker Internal Consistency”Within a single broker:
- Publish-fanout ordering: Messages fan out in publish order
- Cache consistency: Single-writer per key (no torn writes)
- Subscription isolation: Independent queues prevent crosstalk
Multi-Broker Consistency
Section titled “Multi-Broker Consistency”In a clustered deployment:
- Shard leadership: Only one leader per shard, and it serves only while it holds a lease. A broker that has been superseded stops acknowledging rather than discovering the fact later
- Metadata consistency: strongly consistent, because it lives in one Postgres that every control-plane instance reads and writes
- Cross-shard ordering: Not guaranteed. Ordering is per key, because a key always resolves to the same shard and a shard is one log on one leader
- Cache consistency: One owner per key, not eventual. A key hashes to a
shard, that shard has one owner, and a broker receiving an operation for a key
it does not own forwards it there — so a value written through any broker is
readable through every other, and two brokers cannot hold divergent values for
the same key. A cache write is acknowledged by its leader; a cache cannot
declare
Quorum
Failure Scenarios and Behavior
Section titled “Failure Scenarios and Behavior”Network Partition
Section titled “Network Partition”Publisher-Broker partition:
- Publisher detects connection loss (QUIC idle timeout)
- Unacknowledged publishes are lost
- Publisher must reconnect and retry
Subscriber-Broker partition:
- Subscriber detects connection loss
- Buffered events are lost
- Subscriber must reconnect and re-subscribe (starts from tail)
Broker-Control Plane partition:
- Broker continues serving with cached metadata
- New stream creation fails
- Existing streams continue operating
- Broker reconciles when connection restored
Broker Failure
Section titled “Broker Failure”Process crash:
- In-memory state is lost. Durable streams, the log-backed cache, and consumer-group positions are on disk and survive.
- Clients detect connection loss
- Clients must reconnect to recovered broker
- Subscriptions must be re-established
Planned behavior with durability:
- Durable streams can replay from last checkpoint
- Subscribers can resume from last acknowledged offset
- Cache state can be rebuilt from log
Slow Subscriber Behavior
Section titled “Slow Subscriber Behavior”Scenario: Subscriber processing slows down.
Stages:
- Buffer absorbs slowdown: Events accumulate in subscriber buffer
- Buffer fills: New events start getting dropped for that subscriber
- Other subscribers unaffected: Fast subscribers continue normally
Detection, today: the broker counts drops per subscriber queue
(felix_sub_queue_dropped_total) and logs when a subscriber falls behind.
On a durable stream the subscriber itself can detect loss from an offset gap
and resume. There is no per-subscription lag API and no automatic disconnect
of chronically slow subscribers.
Testing Semantics
Section titled “Testing Semantics”Conformance Testing
Section titled “Conformance Testing”Applications can test semantic guarantees:
Ordering test:
// Publish ordered batchlet messages = vec!["msg1", "msg2", "msg3"];use felix_wire::AckMode;let publisher = client.publisher().await?;publisher .publish_batch("test", "default", "orders", messages, AckMode::PerBatch) .await?;
// Verify subscriber receives in orderlet events = collect_events(&mut subscription, 3).await?;assert_eq!(events, vec!["msg1", "msg2", "msg3"]);Isolation test:
// Start fast and slow subscriberslet mut fast_sub = client.subscribe("test", "default", "stream").await?;let mut slow_sub = client.subscribe("test", "default", "stream").await?;
// Slow subscriber delays processingsimulate_slow_processing(&mut slow_sub);
// Verify fast subscriber still receives all messageslet fast_count = count_events(&mut fast_sub, timeout).await?;assert!(fast_count >= expected_count);Summary: Semantic Guarantees Matrix
Section titled “Summary: Semantic Guarantees Matrix”| Property | Today | Not built |
|---|---|---|
| Pub/Sub delivery | At-most-once ephemeral, at-least-once durable; idempotent producers land a re-sent publish once | Exactly-once delivery |
| Consumer groups | At-least-once, bounded redelivery, dead letters | Shard assignment across a group’s consumers |
| Message ordering | Per shard | Configurable cross-shard |
| Subscriber isolation | Yes | — |
| Cache | Routed to one owner, replicated, read-your-writes through that owner | A declared consistency level; a cache write is acknowledged by its leader |
| TTL precision | Lazy on access, against an absolute expiry | Sweeping expiry |
| Durability | Per stream: ephemeral, or Leader or Quorum acknowledgement |
— |
| Authorization | Tenant-scoped tokens, RBAC per resource, OIDC exchange | — |
| Quotas | None | Per-tenant, per-namespace |
| Multi-key operations | None | Transactions |
Recommendations
Section titled “Recommendations”Choosing Delivery Semantics
Section titled “Choosing Delivery Semantics”Use an ephemeral stream (at-most-once) when:
- Latest value is more important than history (sensor data, metrics)
- Occasional loss is acceptable (telemetry, monitoring)
- Throughput and latency matter more than guarantees
- Application implements own deduplication
Use a durable stream, or a consumer group (at-least-once) when:
- Every message matters (financial transactions, orders)
- Application can handle duplicates (idempotent processing)
- Durability matters more than latency
Exactly-once delivery is not implemented. An idempotent producer keeps a publish retry from duplicating the record; a consumer’s redelivery is still at-least-once. If duplicates are unacceptable on the consuming side — billing, accounting — the deduplication has to be in the application, keyed on something the record carries.
Cache Usage Patterns
Section titled “Cache Usage Patterns”Good cache use cases:
- Session data with TTL
- Configuration with infrequent updates
- Rate limiting counters (with planned atomic increment)
- Recently published message lookup
Poor cache use cases:
- Strongly consistent shared state requiring transactions
- Large values (> 1 MB) better served by object storage
- Frequently updated counters (better as pub/sub)
