= Examples
:sectanchors:
:toc: left
:toclevels: 2
== Chat Server (~30 lines)
[source,rust]
----
use std::sync::Arc;
use rifts::RiftServer;
use tokio::sync::Notify;
#[tokio::main]
async fn main() -> rifts::Result<()> {
let shutdown = Arc::new(Notify::new());
RiftServer::builder()
.websocket_transport()
.build()?
.run("0.0.0.0:9000".parse().unwrap(), shutdown)
.await?;
Ok(())
}
----
== Authenticated Pub/Sub
[source,rust]
----
use std::sync::Arc;
use rifts::{
RiftServer, ServerConfig, AuthMode,
session::{TokenAuth, AuthContext, ClientId},
};
use tokio::sync::Notify;
#[tokio::main]
async fn main() -> rifts::Result<()> {
let auth = Arc::new(TokenAuth::new());
auth.register("admin-token", AuthContext {
client_id: ClientId::new("admin"),
claims: serde_json::json!({"role": "admin"}),
mode: AuthMode::Bearer,
hints: Default::default(),
});
let shutdown = Arc::new(Notify::new());
RiftServer::builder()
.websocket_transport()
.auth(auth)
.config(ServerConfig::default())
.build()?
.run("0.0.0.0:9000".parse().unwrap(), shutdown)
.await?;
Ok(())
}
----
== Sled Persistence
[source,rust]
----
// Cargo.toml
// rifts = { version = "0.1", features = ["sled"] }
use std::sync::Arc;
use std::time::Duration;
use rifts::{
RiftServer, TopicProfile,
broker::InMemoryBroker,
storage::{
SledEngine, SledOffsetStore, SledLogStore,
SledDedupeStore, SledSnapshotStore,
},
};
use tokio::sync::Notify;
#[tokio::main]
async fn main() -> rifts::Result<()> {
let db = sled::Config::new()
.path("/var/lib/rifts/broker")
.open()?;
let broker = InMemoryBroker::with_stores(
TopicProfile::default(),
Duration::from_secs(60),
65_536,
SledOffsetStore::new(SledEngine::new(db.open_tree(b"offsets")?)),
SledLogStore::new(SledEngine::new(db.open_tree(b"log")?)),
SledDedupeStore::new(SledEngine::new(db.open_tree(b"dedupe")?)),
SledSnapshotStore::new(SledEngine::new(db.open_tree(b"snapshots")?)),
);
let shutdown = Arc::new(Notify::new());
RiftServer::builder()
.websocket_transport()
.broker(Arc::new(broker))
.build()?
.run("0.0.0.0:9000".parse().unwrap(), shutdown)
.await?;
Ok(())
}
----
== Actor‑Based Broker (per‑topic parallelism)
[source,rust]
----
use std::sync::Arc;
use std::time::Duration;
use rifts::{
RiftServer, TopicProfile,
actor::TopicRegistry,
broker::ActorBroker,
storage::{
MemoryOffsetStore, MemoryLogStore,
MemoryDedupeStore, MemorySnapshotStore,
},
};
use tokio::sync::Notify;
#[tokio::main]
async fn main() -> rifts::Result<()> {
let registry = Arc::new(TopicRegistry::new(
Arc::new(MemoryOffsetStore::new()),
Arc::new(MemoryLogStore::new()),
Arc::new(MemoryDedupeStore::new()),
Arc::new(MemorySnapshotStore::new()),
TopicProfile::default(),
Duration::from_secs(60),
));
let broker = ActorBroker::new(registry);
let shutdown = Arc::new(Notify::new());
RiftServer::builder()
.websocket_transport()
.broker(Arc::new(broker))
.build()?
.run("0.0.0.0:9000".parse().unwrap(), shutdown)
.await?;
Ok(())
}
----
== RemoteBroker (Gateway ↔ External Broker Node)
[source,rust]
----
use std::sync::Arc;
use rifts::{RiftServer, broker::RemoteBroker};
use tokio::sync::Notify;
#[tokio::main]
async fn main() -> rifts::Result<()> {
let broker = RemoteBroker::connect("192.168.1.10:9200".parse()?).await?;
let shutdown = Arc::new(Notify::new());
RiftServer::builder()
.websocket_transport()
.broker(Arc::new(broker))
.build()?
.run("0.0.0.0:9000".parse().unwrap(), shutdown)
.await?;
Ok(())
}
----
The broker node must speak the `WireMsg` CBOR protocol defined in
`src/broker/wire.rs`. Any language / stack can implement it.
== Direct Broker Usage (no network)
[source,rust]
----
use std::sync::Arc;
use std::time::Duration;
use bytes::Bytes;
use rifts::{
broker::{InMemoryBroker, Broker, SubscribeIntent},
frame::{Frame, FrameType, Codec},
TopicProfile, RetentionPolicy,
};
#[tokio::main]
async fn main() -> rifts::Result<()> {
let profile = TopicProfile {
retention: RetentionPolicy::Count(100),
..TopicProfile::default()
};
let broker = InMemoryBroker::new(profile, Duration::from_secs(60), 65_536);
let frame = Frame {
frame_type: FrameType::Data,
codec: Codec::Json,
topic: Some("test".into()),
message_id: Some("msg-1".into()),
payload: Some(Bytes::from_static(b"hello")),
..Frame::default()
};
let outcome = broker.publish(&frame).await?;
println!("published at offset {}", outcome.offset);
Ok(())
}
----
== Axum WebSocket Adapter
[source,rust]
----
// Cargo.toml
// rifts = { version = "0.1", default-features = false, features = ["axum"] }
use axum::{Router, routing::get, extract::ws::WebSocketUpgrade, response::IntoResponse};
use rifts::transport::axum::AxumWsConnection;
async fn ws_upgrade(ws: WebSocketUpgrade) -> impl IntoResponse {
ws.on_upgrade(|socket| async move {
let conn = AxumWsConnection::new(socket);
// Pass conn to your Rift connection handler.
})
}
#[tokio::main]
async fn main() {
let app = Router::new().route("/ws", get(ws_upgrade));
let listener = tokio::net::TcpListener::bind("0.0.0.0:3000").await.unwrap();
axum::serve(listener, app).await.unwrap();
}
----
== Custom AuthProvider
[source,rust]
----
use async_trait::async_trait;
use rifts::{
session::{AuthProvider, AuthContext, AuthHints, ClientId},
AuthMode, Result, RiftError,
};
struct MyAuth {
jwks_url: String,
}
#[async_trait]
impl AuthProvider for MyAuth {
async fn authenticate(&self, mode: AuthMode, token: Option<&str>)
-> Result<AuthContext>
{
let token = token.ok_or_else(||
RiftError::Auth(rifts::error::AuthReject::Required)
)?;
let claims = validate_jwt(token, &self.jwks_url).await?;
Ok(AuthContext {
client_id: ClientId::new(claims.sub),
claims: serde_json::to_value(&claims).unwrap(),
mode,
hints: AuthHints::default(),
})
}
async fn revoke(&self, _client_id: &ClientId) -> Result<()> {
Ok(())
}
}
----