Skip to content

Client usage

Fibril has Rust, TypeScript, Python, Go, and C# clients for the same core broker surface:

  • connect with optional auth
  • publish with or without confirmation
  • delayed publish
  • manual acknowledgements
  • auto-ack convenience
  • message helpers for msgpack, JSON, raw bytes, and text payloads

The examples below are source-tree examples for the current pre-alpha API. They are useful for building against this repository, but not yet a stable package contract.

Parity note: all five clients cover the branch’s partition topology, owner-redirect, partition-key routing, exclusive consumer-group surface, and message TTL. The Rust client is the reference implementation. The TypeScript, Python, Go, and C# clients mirror its behavior. The Python client also ships a synchronous facade (fibril.blocking.BlockingClient) over the same async core. The Go client presents ordinary blocking, goroutine-safe calls and delivers on channels, so it has no separate async and sync split. The C# client is async-native (Task-based) and takes an optional CancellationToken on every network call.

  • Rust: the workspace toolchain (current stable Rust).
  • TypeScript: Node.js 18 or newer, ESM.
  • Python: Python 3.11 or newer. The only runtime dependency is the PyPI msgpack package, installed automatically with the client.
  • Go: Go 1.23 or newer. No third-party dependencies. Every network method takes a context.Context as its first argument (the snippets below use ctx; pass context.Background() when you have no deadline or cancellation).
  • C#: .NET 8 or newer (the library targets net8.0). No third-party dependencies. Every network method takes an optional CancellationToken as its last argument.

For reconnect behavior and current limits, see reconnects.

Configure additional trusted broker addresses when the initial node may be lost. Topology refresh tries the initial address, configured discovery addresses, known owners and pooled connections. Each connection uses the same authentication and TLS settings. Initial connection still requires the explicit connect address to be reachable; these addresses provide fallback discovery after connection.

| Client | Connection option | | --- | --- | | Rust | .discovery_endpoints(["node-b:9876".into(), "node-c:9876".into()]) | | TypeScript | new ClientOptions({ discoveryEndpoints: ["node-b:9876", "node-c:9876"] }) | | Python | ClientOptions(discovery_endpoints=("node-b:9876", "node-c:9876")) | | Go | ClientOptions{DiscoveryEndpoints: []string{"node-b:9876", "node-c:9876"}} | | C# | new ClientOptions { DiscoveryEndpoints = new[] { "node-b:9876", "node-c:9876" } } |

Topology currently advertises partition owners, so a client connected to the only owner it knows cannot discover an unknown survivor without a fallback address. Supervised subscriptions retry temporary ownership and server errors while the new owner recovers. Interrupted publish confirmations still have an unknown outcome; application retries require the usual duplicate handling.

use fibril_client::ClientOptions;
let client = ClientOptions::new()
.auth("fibril", "fibril")
.connect("127.0.0.1:9876")
.await?;

By default, the clients make one automatic reconnect attempt before a new operation if the previous engine is already closed. Disable that with disable_auto_reconnect() in Rust and Python, or disableAutoReconnect() in TypeScript. After a successful resume, the clients send their known subscription metadata to the broker and read the reconciliation result. Subscriptions that the broker confirms with keep continue on the existing stream. This does not replay in-flight publish operations. See reconnects for the current limits.

The default reconciliation policy is conservative. Rust can opt into restoring missing server-side subscriptions with reconnect_reconcile_policy(ReconcilePolicy::Restore). TypeScript can do the same with withReconnectReconcilePolicy("restore"), and Python with ClientOptions(reconnect_reconcile_policy="restore").

When the broker serves TLS (tls.enabled = true in the broker config), enable TLS in the client options. Three trust paths, in the order the client resolves them:

  • tls_ca_fingerprint(...) pins the SHA-256 fingerprint printed in the broker startup log. No files needed, survives leaf certificate rotation under the same CA.
  • tls_ca_path(...) trusts the CA certificate(s) in a PEM file, e.g. the broker’s generated <data_dir>/tls/ca.pem or an internal CA.
  • bare tls() uses the OS trust store, for brokers with publicly issued certificates.
use fibril_client::ClientOptions;
let client = ClientOptions::new()
.auth("fibril", "fibril")
.tls_ca_path("/var/lib/fibril/tls/ca.pem")
.connect("broker.internal:9876")
.await?;

On brokers with tls.client_auth, a client certificate replaces the password: a verified certificate whose identity names an existing broker user authenticates the connection with no auth call at all. Issue one from the deployment CA with fibrilctl cert issue <identity>, then present it with tls_client_cert(cert, key) (TypeScript withTlsClientCert, Python with_tls_client_cert). See TLS across nodes for the broker-side setup.

Misconfiguration surfaces as distinct typed errors because the fixes differ: TlsRequiredByBroker (the broker reported that this plaintext client must enable TLS), TlsNotSupportedByBroker (the handshake ended early, the broker listener is probably plaintext), TlsCertificateUntrusted (trust configuration: fix the CA path or the fingerprint), TlsClientCertificateRequired (the broker demands a client certificate this client did not present), and TlsConfig for client-side option problems. The TypeScript and Python names are TlsRequiredByBrokerError and so on. None of them are retried automatically.

One Python caveat: pinning the broker’s leaf fingerprint works with only the standard library. Pinning a CA fingerprint (so the leaf can rotate under the same CA) path-validates the presented chain, which needs Python 3.13+ (to expose the chain) and the optional cryptography package (pip install fibril[tls-pin]); otherwise pin the leaf certificate or use with_tls_ca_path.

The conservative default is the safest operational behavior. It keeps matching subscriptions, closes client streams the broker cannot prove are still valid, and drops server-side subscriptions the client no longer reports.

Use restore mode only when you want the client to ask the broker to recreate client-owned subscriptions that are missing after a successful resume. Restored subscriptions may receive a new server subscription id, which the clients remap internally.

use fibril_client::{ClientOptions, ReconcilePolicy};
let client = ClientOptions::new()
.reconnect_reconcile_policy(ReconcilePolicy::Restore)
.connect("127.0.0.1:9876")
.await?;

Plain publish uses the common unconfirmed path. Confirmed publish waits for the broker to acknowledge the stored message and returns the topic offset. For pipelined confirmed publishing, send with a confirmation handle and await those handles later.

let publisher = client.publisher("email.send")?;
publisher
.publish("hello")
.await?;
let offset = publisher
.publish_confirmed("needs an offset")
.await?;
let confirmation = publisher
.publish_with_confirmation("pipeline me")
.await?;
let pipelined_offset = confirmation.confirmed().await?;

Delayed publish uses a distinct protocol frame and stores a not_before Unix-millisecond deadline. It does not add delay bytes to the common publish frame.

use std::time::Duration;
publisher
.publish_delayed("send later", Duration::from_secs(30))
.await?;
let offset = publisher
.publish_delayed_confirmed("send later and return offset", Duration::from_secs(30))
.await?;

Rust delays are a std::time::Duration, so the unit is always explicit (write Duration::from_secs(30) or Duration::from_millis(250)).

Manual-ack consumers can requeue a message after a delay instead of making it ready immediately.

use std::time::Duration;
let msg = sub.recv().await.expect("message");
msg.retry_after(Duration::from_secs(30)).await?;

Rust delays are a std::time::Duration, matching delayed publish.

Plain values are encoded as msgpack by default. Use explicit message helpers when you want JSON, raw bytes, text, or custom headers.

use fibril_client::NewMessage;
publisher
.publish(NewMessage::json(&serde_json::json!({
"id": 42,
"kind": "welcome",
}))?)
.await?;
publisher
.publish(
NewMessage::text("plain text")
.header("x-trace-id", "abc123"),
)
.await?;
publisher
.publish(NewMessage::raw(vec![1, 2, 3]))
.await?;

Fibril reserves two header namespaces for system metadata, and the broker rejects a publish that sets either: fibril.* for broker and client protocol metadata (for example the reliable-publisher dedup ids), and stroma.* for metadata owned by the storage and queue-state layer (Stroma), such as the stroma.dlq.* provenance stamped onto a dead-lettered message. User code should not set headers with those prefixes. See metadata policy for the development note.

Manual acknowledgement is the primary processing model. Messages are leased while inflight. Consumers settle each message explicitly.

let mut sub = client
.subscribe("email.send")?
.group("workers")?
.prefetch(32)
.sub()
.await?;
while let Some(msg) = sub.recv().await {
let body = msg.text()?;
match send_email(body).await {
Ok(()) => {
msg.complete().await?;
}
Err(_) => {
msg.retry().await?;
}
}
}

retry() requeues immediately. Use retry_after(..) in Rust and Python, retryAfter(..) in TypeScript, or RetryAfter(..) in Go and C# for delayed retry.

Dead-letter routing has two parts:

  • Operators configure the global DLQ target through the admin UI/API.
  • Applications may declare per-queue retry and DLQ policy through the client or fibrilctl.

For example, first configure a global DLQ target:

PUT /admin/api/global-dlq
{
"expected_version": 0,
"target": {
"topic": "_dlq.email",
"group": null
}
}

Then configure a source queue to use that global target after retries are exhausted:

use fibril_client::QueueConfig;
client
.declare_queue(
QueueConfig::new("email.send")?
.group("workers")?
.use_global_dead_letter_queue()
.max_retries(3),
)
.await?;

Application clients keep using normal publish/consume code:

let publisher = client.publisher_grouped("email.send", "workers")?;
publisher.publish_confirmed(NewMessage::json(&job)?).await?;
let mut sub = client
.subscribe("email.send")?
.group("workers")?
.prefetch(32)
.sub()
.await?;
while let Some(msg) = sub.recv().await {
if process(msg.deserialize()?).await.is_ok() {
msg.complete().await?;
} else {
msg.retry().await?;
}
}

Once the configured retry limit is exhausted, the broker routes the failed message to _dlq.email. Subscribe to that queue like any other queue when you want inspection or replay tooling.

In the clients that dispatch on content type, deserialize() picks the decoder from the message content-type (a missing content type defaults to msgpack). Go and C# do not dispatch on content type: you decode explicitly with JSON/ Json<T> for JSON, Payload for raw bytes, and Text for text (see the Go and C# tabs).

#[derive(serde::Deserialize)]
struct Job {
id: u64,
}
let job: Job = msg.deserialize()?;
let raw: &[u8] = msg.raw();
let text: &str = msg.text()?;

Auto-ack subscriptions are a convenience mode that settles each delivery server-side as the broker sends it, so the consumer just reads messages. Use manual ack when processing correctness matters.

let mut sub = client
.subscribe("metrics")?
.prefetch(128)
.sub_auto_ack()
.await?;
while let Some(msg) = sub.recv().await {
observe(msg.deserialize::<Metric>()?);
}

A Plexus stream delivers every record to every consumer. Declare it with declare_plexus, publish with the ordinary publisher (the broker routes by channel kind), and consume with stream(...). Name the subscription to get a durable cursor that resumes after a restart and advances on ack. A stream subscription reads all partitions and fans them in.

use fibril_client::StreamConfig;
client
.declare_plexus(StreamConfig::new("events")?.partitions(4).retain_records(1_000_000))
.await?;
// Publish uses the normal publisher.
client.publisher("events")?.publish("hello").await?;
let mut sub = client
.stream("events")?
.durable("analytics")
.filter("region", "eu-*")
.sub()
.await?;
while let Some(msg) = sub.recv().await {
handle(msg.text()?);
msg.complete().await?; // advances the durable cursor
}

For a live tail without a durable cursor, drop .durable(...) and pick a start position: from_latest() / fromLatest() (the default), from_earliest(), from_offset(..), from_last(..), or from_time(..). Use sub_auto_ack() / subAutoAck() when you do not need explicit settlement.

Discovery is an opt-in surface: call routing() to get a routing view of the same connection, then subscribe_pattern(glob) to fan in across every work queue whose topic matches, or subscribe_stream_pattern(glob) for Plexus streams. The subscription keeps attaching channels that start matching later, so queues or streams declared after you subscribe are picked up without a reconnect. Each delivery is paired with the channel it came from. The glob is the same *-wildcard grammar as the header filter (no regex), and "*" matches everything. Manual (sub / .sub()) and auto-ack (sub_auto_ack / .subAutoAck()) variants both exist, and the routing view still publishes and subscribes normally, so it composes with reliable publishing.

let mut sub = client
.routing()
.subscribe_pattern("events.*")
.sub()
.await?;
while let Some((source, msg)) = sub.recv().await {
handle(&source.topic, msg.text()?);
msg.complete().await?;
}
client.shutdown().await;

For code not built on asyncio, the Python package ships a synchronous facade. It runs the async client on a background event-loop thread and bridges each call, so it is the same client behind one source of truth rather than a second implementation. Subscriptions become ordinary iterators. The Async/Blocking toggle on the examples above shows the synchronous form of each call, and the full shape is:

from fibril.blocking import BlockingClient
with BlockingClient.connect("127.0.0.1:9876") as client:
client.publisher("email.send").publish_confirmed({"body": "hello"})
sub = client.subscribe("email.send").group("workers").sub()
for msg in sub:
process(msg.deserialize())
msg.complete()
sub.close()