rifts 0.2.0

Rift Realtime Protocol / 1.0 — server-side implementation
Documentation
= 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(())
    }
}
----