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.
Requirements
Section titled “Requirements”- 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
msgpackpackage, installed automatically with the client. - Go: Go 1.23 or newer. No third-party dependencies. Every network method
takes a
context.Contextas its first argument (the snippets below usectx; passcontext.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 optionalCancellationTokenas its last argument.
For reconnect behavior and current limits, see reconnects.
Discovery during failover
Section titled “Discovery during failover”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.
Connect
Section titled “Connect”use fibril_client::ClientOptions;
let client = ClientOptions::new() .auth("fibril", "fibril") .connect("127.0.0.1:9876") .await?;import { ClientOptions } from "@fibril/client";
const client = await new ClientOptions() .withAuth("fibril", "fibril") .connect("127.0.0.1:9876");from fibril import ClientOptions
client = await ClientOptions().with_auth("fibril", "fibril").connect("127.0.0.1:9876")from fibril import ClientOptionsfrom fibril.blocking import BlockingClient
client = BlockingClient.connect( "127.0.0.1:9876", ClientOptions().with_auth("fibril", "fibril"))import ( "context"
fibril "github.com/Axmouth/fibril/clients/go")
ctx := context.Background()client, err := fibril.Dial(ctx, "127.0.0.1:9876", fibril.ClientOptions{ Credentials: &fibril.Credentials{Username: "fibril", Password: "fibril"},})using Fibril;
var client = await Client.ConnectAsync("127.0.0.1:9876", new ClientOptions{ Credentials = new Credentials("fibril", "fibril"),});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.pemor 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?;import { ClientOptions } from "@fibril/client";
const client = await new ClientOptions() .withAuth("fibril", "fibril") .withTlsCaPath("/var/lib/fibril/tls/ca.pem") .connect("broker.internal:9876");from fibril import ClientOptions
client = await ( ClientOptions() .with_auth("fibril", "fibril") .with_tls_ca_path("/var/lib/fibril/tls/ca.pem") .connect("broker.internal:9876"))client, err := fibril.Dial(ctx, "broker.internal:9876", fibril.ClientOptions{ Credentials: &fibril.Credentials{Username: "fibril", Password: "fibril"}, TLS: &fibril.TLSOptions{CAFile: "/var/lib/fibril/tls/ca.pem"},})The Go trust order matches the others: CAFingerprint pins the SHA-256
fingerprint, else CAFile trusts a PEM, else an empty TLSOptions{} uses the OS
roots. A client certificate is ClientCertFile with ClientKeyFile.
var client = await Client.ConnectAsync("broker.internal:9876", new ClientOptions{ Credentials = new Credentials("fibril", "fibril"), Tls = new TlsOptions { CaFile = "/var/lib/fibril/tls/ca.pem" },});The C# trust order matches the others: CaFingerprint pins the SHA-256
fingerprint, else CaFile trusts a PEM, else a bare new TlsOptions() uses the
OS roots. A client certificate is ClientCertFile with ClientKeyFile.
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.
Reconnect Policy
Section titled “Reconnect Policy”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?;import { ClientOptions, type ReconcilePolicy } from "@fibril/client"
const policy: ReconcilePolicy = "restore"const client = await new ClientOptions() .withReconnectReconcilePolicy(policy) .connect("127.0.0.1:9876")from fibril import ClientOptions
client = await ClientOptions( reconnect_reconcile_policy="restore").connect("127.0.0.1:9876")from fibril import ClientOptionsfrom fibril.blocking import BlockingClient
client = BlockingClient.connect( "127.0.0.1:9876", ClientOptions(reconnect_reconcile_policy="restore"),)client, err := fibril.Dial(ctx, "127.0.0.1:9876", fibril.ClientOptions{ ReconcilePolicy: fibril.ReconcileRestore,})var client = await Client.ConnectAsync("127.0.0.1:9876", new ClientOptions{ ReconcilePolicy = ReconcilePolicy.Restore,});Publish
Section titled “Publish”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?;const publisher = client.publisher("email.send");
await publisher.publish("hello");
const offset = await publisher.publishConfirmed("needs an offset");
const confirmation = await publisher.publishWithConfirmation("pipeline me");const pipelinedOffset = await confirmation.confirmed();publisher = client.publisher("email.send")
await publisher.publish("hello")
offset = await publisher.publish_confirmed("needs an offset")
confirmation = await publisher.publish_with_confirmation("pipeline me")pipelined_offset = await confirmation.confirmed()publisher = client.publisher("email.send")
publisher.publish("hello")
offset = publisher.publish_confirmed("needs an offset")
confirmation = publisher.publish_with_confirmation("pipeline me")pipelined_offset = confirmation.confirmed()publisher := client.Publisher("email.send")
publisher.Publish(ctx, fibril.Text("hello"))
offset, err := publisher.PublishConfirmed(ctx, fibril.Text("needs an offset"))
confirmation, err := publisher.PublishWithConfirmation(ctx, fibril.Text("pipeline me"))pipelinedOffset, err := confirmation.Confirmed(ctx)var publisher = client.Publisher("email.send");
await publisher.PublishAsync(Message.Text("hello"));
var offset = await publisher.PublishConfirmedAsync(Message.Text("needs an offset"));
var confirmation = await publisher.PublishWithConfirmationAsync(Message.Text("pipeline me"));var pipelinedOffset = await confirmation.Confirmed();Delayed Publish
Section titled “Delayed Publish”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)).
await publisher.publishDelayed("send later", 30_000);
const offset = await publisher.publishDelayedConfirmed( "send later and return offset", new Date(Date.now() + 30_000),);TypeScript numeric delays are milliseconds; pass { seconds: 30 }, { ms: 250 }, or { minutes: 5 } to make the unit explicit, or a Date for an absolute deadline.
await publisher.publish_delayed("send later", 30)
offset = await publisher.publish_delayed_confirmed( "send later and return offset", 30)publisher.publish_delayed("send later", 30)
offset = publisher.publish_delayed_confirmed("send later and return offset", 30)Python delays are seconds (a float or a datetime.timedelta).
publisher.PublishDelayed(ctx, fibril.Text("send later"), 30*time.Second)
offset, err := publisher.PublishDelayedConfirmed(ctx, fibril.Text("send later and return offset"), 30*time.Second)Go delays are a time.Duration.
await publisher.PublishDelayedAsync(Message.Text("send later"), TimeSpan.FromSeconds(30));
var offset = await publisher.PublishDelayedConfirmedAsync( Message.Text("send later and return offset"), TimeSpan.FromSeconds(30));C# delays are a TimeSpan.
Delayed Retry
Section titled “Delayed Retry”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.
const msg = await sub.recv();await msg.retryAfter(30_000);TypeScript numeric delays are milliseconds; pass { seconds: 30 }, { ms: 250 }, or { minutes: 5 } to make the unit explicit, or a Date for an absolute retry deadline.
msg = await sub.recv()await msg.retry_after(30)msg = sub.recv()msg.retry_after(30)Python delays are seconds (a float or a datetime.timedelta).
msg := <-sub.Deliveriesmsg.RetryAfter(30 * time.Second)Retry() requeues immediately. RetryAfter(d) takes a time.Duration.
await foreach (var msg in sub.Deliveries()){ msg.RetryAfter(TimeSpan.FromSeconds(30));}Retry() requeues immediately. RetryAfter(d) takes a TimeSpan.
Message Payloads
Section titled “Message Payloads”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?;import { NewMessage } from "@fibril/client";
await publisher.publish( NewMessage.json({ id: 42, kind: "welcome", }),);
await publisher.publish( NewMessage.text("plain text") .header("x-trace-id", "abc123"),);
await publisher.publish(NewMessage.raw(new Uint8Array([1, 2, 3])));from fibril import NewMessage
await publisher.publish(NewMessage.json({"id": 42, "kind": "welcome"}))
await publisher.publish( NewMessage.text("plain text").header("x-trace-id", "abc123"))
await publisher.publish(NewMessage.raw(bytes([1, 2, 3])))from fibril import NewMessage
publisher.publish(NewMessage.json({"id": 42, "kind": "welcome"}))
publisher.publish( NewMessage.text("plain text").header("x-trace-id", "abc123"))
publisher.publish(NewMessage.raw(bytes([1, 2, 3])))job, err := fibril.JSON(map[string]any{"id": 42, "kind": "welcome"})publisher.Publish(ctx, job)
publisher.Publish(ctx, fibril.Text("plain text").WithHeader("x-trace-id", "abc123"))
publisher.Publish(ctx, fibril.Raw([]byte{1, 2, 3}))fibril.Msgpack(bytes) and fibril.Custom(contentType, bytes) tag already-encoded
payloads, so the Go client needs no msgpack dependency.
await publisher.PublishAsync(Message.Json(new { id = 42, kind = "welcome" }));
await publisher.PublishAsync(Message.Text("plain text").WithHeader("x-trace-id", "abc123"));
await publisher.PublishAsync(Message.Raw(new byte[] { 1, 2, 3 }));Message.Msgpack(bytes) and Message.Custom(contentType, bytes) tag
already-encoded payloads, so the C# client needs no MessagePack dependency.
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 Acknowledgements
Section titled “Manual Acknowledgements”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?; } }}const sub = await client .subscribe("email.send") .group("workers") .prefetch(32) .sub();
for await (const msg of sub) { const body = msg.text();
try { await sendEmail(body); await msg.complete(); } catch { await msg.retry(); }}sub = await client.subscribe("email.send").group("workers").prefetch(32).sub()
async for msg in sub: body = msg.text()
try: await send_email(body) await msg.complete() except Exception: await msg.retry()sub = client.subscribe("email.send").group("workers").prefetch(32).sub()
for msg in sub: body = msg.text()
try: send_email(body) msg.complete() except Exception: msg.retry()group := "workers"sub, err := client.SubscribeTopic(ctx, "email.send", fibril.TopicSubscribeOptions{Group: &group, Prefetch: 32, AutoAck: false})
for msg := range sub.Deliveries { body := msg.Text() if err := sendEmail(body); err == nil { msg.Complete() } else { msg.Retry() }}var sub = await client.SubscribeTopicAsync("email.send", group: "workers", prefetch: 32);
await foreach (var msg in sub.Deliveries()){ var body = msg.Text(); if (await SendEmail(body)) { msg.Complete(); } else { msg.Retry(); }}retry() requeues immediately. Use retry_after(..) in Rust and Python, retryAfter(..) in TypeScript, or RetryAfter(..) in Go and C# for delayed retry.
Dead Letter Workflow
Section titled “Dead Letter Workflow”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?;import { QueueConfig } from "@fibril/client";
await client.declareQueue( new QueueConfig("email.send") .group("workers") .useGlobalDeadLetterQueue() .maxRetries(3),);from fibril import QueueConfig
await client.declare_queue( QueueConfig("email.send") .group("workers") .use_global_dead_letter_queue() .max_retries(3))from fibril import QueueConfig
client.declare_queue( QueueConfig("email.send") .group("workers") .use_global_dead_letter_queue() .max_retries(3))client.DeclareQueue(ctx, fibril.NewQueueConfig("email.send"). Group("workers"). Dlq(fibril.DlqPolicy{Kind: fibril.DlqGlobal}). DlqMaxRetries(3))await client.DeclareQueueAsync("email.send", new QueueConfig{ Group = "workers", DeadLetter = DeadLetterPolicy.Global, MaxRetries = 3,});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?; }}const publisher = client.publisherGrouped("email.send", "workers");await publisher.publishConfirmed(NewMessage.json(job));
const sub = await client .subscribe("email.send") .group("workers") .prefetch(32) .sub();
for await (const msg of sub) { if (await process(msg.deserialize())) { await msg.complete(); } else { await msg.retry(); }}publisher = client.publisher_grouped("email.send", "workers")await publisher.publish_confirmed(NewMessage.json(job))
sub = await client.subscribe("email.send").group("workers").prefetch(32).sub()
async for msg in sub: if await process(msg.deserialize()): await msg.complete() else: await msg.retry()publisher = client.publisher_grouped("email.send", "workers")publisher.publish_confirmed(NewMessage.json(job))
sub = client.subscribe("email.send").group("workers").prefetch(32).sub()
for msg in sub: if process(msg.deserialize()): msg.complete() else: msg.retry()publisher := client.PublisherGrouped("email.send", "workers")job, _ := fibril.JSON(jobData)publisher.PublishConfirmed(ctx, job)
group := "workers"sub, err := client.SubscribeTopic(ctx, "email.send", fibril.TopicSubscribeOptions{Group: &group, Prefetch: 32, AutoAck: false})
for msg := range sub.Deliveries { var j Job if msg.JSON(&j) == nil && process(j) { msg.Complete() } else { msg.Retry() }}var publisher = client.PublisherGrouped("email.send", "workers");await publisher.PublishConfirmedAsync(Message.Json(jobData));
var sub = await client.SubscribeTopicAsync("email.send", group: "workers", prefetch: 32);
await foreach (var msg in sub.Deliveries()){ var job = msg.Json<Job>(); if (job is not null && Process(job)) { msg.Complete(); } else { msg.Retry(); }}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.
Decode Messages
Section titled “Decode Messages”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()?;type Job = { id: number;};
const job = msg.deserialize<Job>();const raw = msg.raw();const text = msg.text();job = msg.deserialize()raw = msg.raw()text = msg.text()type Job struct { ID uint64 `json:"id"`}
var job Joberr := msg.JSON(&job)raw := msg.Payloadtext := msg.Text()Go decode is explicit: JSON(&v) for JSON, Payload for raw bytes, Text() for text.
record Job(ulong Id);
var job = msg.Json<Job>();byte[] raw = msg.Payload;string text = msg.Text();C# decode is explicit: Json<T>() for JSON, Payload for raw bytes, Text() for text.
Auto Ack
Section titled “Auto Ack”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>()?);}const sub = await client .subscribe("metrics") .prefetch(128) .subAutoAck();
for await (const msg of sub) { observe(msg.deserialize<Metric>());}sub = await client.subscribe("metrics").prefetch(128).sub_auto_ack()
async for msg in sub: observe(msg.deserialize())sub = client.subscribe("metrics").prefetch(128).sub_auto_ack()
for msg in sub: observe(msg.deserialize())sub, err := client.SubscribeTopic(ctx, "metrics", fibril.TopicSubscribeOptions{Prefetch: 128, AutoAck: true})
for msg := range sub.Deliveries { var m Metric msg.JSON(&m) observe(m)}var sub = await client.SubscribeTopicAsync("metrics", prefetch: 128, autoAck: true);
await foreach (var msg in sub.Deliveries()){ observe(msg.Json<Metric>());}Plexus streams (fan-out)
Section titled “Plexus streams (fan-out)”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}import { StreamConfig } from "@fibril/client";
await client.declarePlexus( new StreamConfig("events").partitions(4).retainRecords(1_000_000),);
await client.publisher("events").publish("hello");
const sub = await client .stream("events") .durable("analytics") .filter("region", "eu-*") .sub();
for await (const msg of sub) { handle(msg.text()); await msg.complete(); // advances the durable cursor}from fibril import StreamConfig
await client.declare_plexus(StreamConfig("events").partitions(4).retain_records(1_000_000))
await client.publisher("events").publish("hello")
sub = await client.stream("events").durable("analytics").filter("region", "eu-*").sub()
async for msg in sub: handle(msg.text()) await msg.complete() # advances the durable cursorfrom fibril import StreamConfig
client.declare_plexus(StreamConfig("events").partitions(4).retain_records(1_000_000))
client.publisher("events").publish("hello")
sub = client.stream("events").durable("analytics").filter("region", "eu-*").sub()
for msg in sub: handle(msg.text()) msg.complete() # advances the durable cursorclient.DeclarePlexus(ctx, fibril.NewStreamConfig("events"). PartitionCount(4). Retention(fibril.NewStreamRetention().RetainRecords(1_000_000)))
// Publish uses the normal publisher.client.Publisher("events").Publish(ctx, fibril.Text("hello"))
durable := "analytics"sub, err := client.SubscribeStreamTopic(ctx, "events", fibril.StreamSubscribeOptions{ DurableName: &durable, Filter: []fibril.StreamFilter{{Key: "region", Pattern: "eu-*"}}, Prefetch: 16,})
for msg := range sub.Deliveries { handle(msg.Text()) msg.Complete() // advances the durable cursor}await client.DeclarePlexusAsync("events", new StreamConfig{ PartitionCount = 4, Retention = new StreamRetentionPolicy { RetainRecords = 1_000_000 },});
// Publish uses the normal publisher.await client.Publisher("events").PublishAsync(Message.Text("hello"));
var sub = await client.SubscribeStreamTopicAsync("events", new StreamSubscribeOptions{ DurableName = "analytics", Filter = new[] { new StreamHeaderFilter("region", "eu-*") }, Prefetch = 16,});
await foreach (var msg in sub.Deliveries()){ handle(msg.Text()); msg.Complete(); // 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.
Pattern subscribe and discovery
Section titled “Pattern subscribe and discovery”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?;}const sub = await client.routing().subscribePattern("events.*").sub();
for await (const { source, message } of sub) { handle(source.topic, message.text()); await message.complete();}sub = await client.routing().subscribe_pattern("events.*").sub()
async for item in sub: handle(item.source.topic, item.message.text()) await item.message.complete()sub = client.routing().subscribe_pattern("events.*").sub()
for item in sub: handle(item.source.topic, item.message.text()) item.message.complete()routing := client.Routing()sub, err := routing.SubscribePattern(ctx, "events.*", fibril.PatternSubscribeOptions{})
// Each delivery carries its source topic in msg.Topic.for msg := range sub.Deliveries { handle(msg.Topic, msg.Text()) msg.Complete()}var sub = await client.Routing.SubscribePatternAsync("events.*");
// Each delivery carries its source topic in msg.Topic.await foreach (var msg in sub.Deliveries()){ handle(msg.Topic, msg.Text()); msg.Complete();}Shutdown
Section titled “Shutdown”client.shutdown().await;await client.shutdown();client.Shutdown()await client.ShutdownAsync();ShutdownAsync() is an alias for DisposeAsync(). Client is IAsyncDisposable,
so await using var client = ... also disposes it automatically at the end of scope.
Blocking Python client
Section titled “Blocking Python client”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()