cnats: NATS-native chat with Leptos SSR and Kanidm SSO
Initial import plus deployment packaging: multi-stage Dockerfile (cargo-leptos build -> debian-slim runtime), .dockerignore, and a dev-only docker-compose (app + local NATS). Production deployment lives in the infrastructure repo. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
+349
@@ -0,0 +1,349 @@
|
||||
use leptos::prelude::*;
|
||||
use leptos_meta::{provide_meta_context, MetaTags, Stylesheet, Title};
|
||||
use leptos_router::{
|
||||
components::{Route, Router, Routes, A},
|
||||
hooks::use_params_map,
|
||||
path,
|
||||
};
|
||||
|
||||
use crate::auth::{current_user, User};
|
||||
use crate::chat::{is_valid_room, room_subject, ChatMessage, SendMessage, DEFAULT_ROOM, ROOMS};
|
||||
|
||||
pub fn shell(options: LeptosOptions) -> impl IntoView {
|
||||
view! {
|
||||
<!DOCTYPE html>
|
||||
<html lang="en">
|
||||
<head>
|
||||
<meta charset="utf-8"/>
|
||||
<meta name="viewport" content="width=device-width, initial-scale=1"/>
|
||||
<link rel="icon" href="/favicon.svg" type="image/svg+xml"/>
|
||||
<link rel="preconnect" href="https://fonts.googleapis.com"/>
|
||||
<link rel="preconnect" href="https://fonts.gstatic.com" crossorigin=""/>
|
||||
<link
|
||||
href="https://fonts.googleapis.com/css2?family=Archivo:wght@400;500;600;700;900&family=IBM+Plex+Mono:ital,wght@0,400;0,500;0,600;1,400&display=swap"
|
||||
rel="stylesheet"
|
||||
/>
|
||||
<AutoReload options=options.clone()/>
|
||||
<HydrationScripts options/>
|
||||
<MetaTags/>
|
||||
</head>
|
||||
<body>
|
||||
<App/>
|
||||
</body>
|
||||
</html>
|
||||
}
|
||||
}
|
||||
|
||||
#[component]
|
||||
pub fn App() -> impl IntoView {
|
||||
provide_meta_context();
|
||||
|
||||
view! {
|
||||
<Stylesheet id="leptos" href="/pkg/cnats.css"/>
|
||||
<Title text="cnats — chat over the bus"/>
|
||||
<Router>
|
||||
<Routes fallback=|| view! { <NotFound/> }>
|
||||
<Route path=path!("") view=ChatPage/>
|
||||
<Route path=path!("/r/:room") view=ChatPage/>
|
||||
</Routes>
|
||||
</Router>
|
||||
}
|
||||
}
|
||||
|
||||
#[component]
|
||||
fn ChatPage() -> impl IntoView {
|
||||
let params = use_params_map();
|
||||
let room = Memo::new(move |_| {
|
||||
params.with(|p| {
|
||||
p.get("room")
|
||||
.filter(|r| is_valid_room(r))
|
||||
.unwrap_or_else(|| DEFAULT_ROOM.to_string())
|
||||
})
|
||||
});
|
||||
let user = Resource::new(|| (), |_| current_user());
|
||||
|
||||
view! {
|
||||
<Suspense fallback=move || {
|
||||
view! {
|
||||
<main class="gate">
|
||||
<p class="gate-loading">"handshaking with the bus…"</p>
|
||||
</main>
|
||||
}
|
||||
}>
|
||||
{move || {
|
||||
user.get()
|
||||
.map(|res| match res {
|
||||
Ok(Some(u)) => view! { <ChatShell room user=u/> }.into_any(),
|
||||
_ => view! { <LoginGate/> }.into_any(),
|
||||
})
|
||||
}}
|
||||
</Suspense>
|
||||
}
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// unauthenticated: the gate
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
#[component]
|
||||
fn LoginGate() -> impl IntoView {
|
||||
view! {
|
||||
<main class="gate">
|
||||
<div class="gate-card">
|
||||
<div class="gate-badge">"MESSAGE BUS · AUTH REQUIRED"</div>
|
||||
<h1 class="gate-title">
|
||||
"CN" <span class="gate-title-accent">"ATS"</span>
|
||||
</h1>
|
||||
<p class="gate-sub">
|
||||
"Realtime chat carried on NATS subjects. Identity issued by your Kanidm realm — no separate passwords, no local accounts."
|
||||
</p>
|
||||
<a class="gate-btn" href="/auth/login" rel="external">
|
||||
<span class="gate-btn-glyph">"⏻"</span>
|
||||
"Sign in with Kanidm"
|
||||
</a>
|
||||
<dl class="gate-meta">
|
||||
<div><dt>"subjects"</dt><dd><code>"chat.room.*"</code></dd></div>
|
||||
<div><dt>"downlink"</dt><dd><code>"SSE"</code></dd></div>
|
||||
<div><dt>"identity"</dt><dd><code>"OIDC + PKCE"</code></dd></div>
|
||||
</dl>
|
||||
</div>
|
||||
</main>
|
||||
}
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// authenticated: the console
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
#[derive(Clone, Copy, PartialEq, Eq)]
|
||||
enum BusStatus {
|
||||
Connecting,
|
||||
Live,
|
||||
Offline,
|
||||
}
|
||||
|
||||
impl BusStatus {
|
||||
fn label(self) -> &'static str {
|
||||
match self {
|
||||
BusStatus::Connecting => "SYN…",
|
||||
BusStatus::Live => "LIVE",
|
||||
BusStatus::Offline => "RETRY",
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[component]
|
||||
fn ChatShell(room: Memo<String>, user: User) -> impl IntoView {
|
||||
let messages = RwSignal::new(Vec::<ChatMessage>::new());
|
||||
let status = RwSignal::new(BusStatus::Connecting);
|
||||
let list_ref = NodeRef::<leptos::html::Div>::new();
|
||||
|
||||
// Subscribe to the room's SSE feed (browser only). Reconnects whenever
|
||||
// the room changes; closes the previous stream first.
|
||||
#[cfg(feature = "hydrate")]
|
||||
{
|
||||
let es_handle = StoredValue::new_local(None::<web_sys::EventSource>);
|
||||
Effect::new(move |_| {
|
||||
let room = room.get();
|
||||
es_handle.update_value(|es| {
|
||||
if let Some(es) = es.take() {
|
||||
es.close();
|
||||
}
|
||||
});
|
||||
messages.set(Vec::new());
|
||||
status.set(BusStatus::Connecting);
|
||||
es_handle.set_value(open_event_source(&room, messages, status));
|
||||
});
|
||||
on_cleanup(move || {
|
||||
es_handle.update_value(|es| {
|
||||
if let Some(es) = es.take() {
|
||||
es.close();
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
// Pin the stream to the bottom as messages arrive.
|
||||
Effect::new(move |_| {
|
||||
messages.track();
|
||||
if let Some(el) = list_ref.get() {
|
||||
el.set_scroll_top(el.scroll_height());
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
let draft = RwSignal::new(String::new());
|
||||
let send = ServerAction::<SendMessage>::new();
|
||||
let on_submit = move |ev: leptos::ev::SubmitEvent| {
|
||||
ev.prevent_default();
|
||||
let text = draft.get();
|
||||
if text.trim().is_empty() {
|
||||
return;
|
||||
}
|
||||
send.dispatch(SendMessage {
|
||||
room: room.get(),
|
||||
text,
|
||||
});
|
||||
draft.set(String::new());
|
||||
};
|
||||
|
||||
let me = user.username.clone();
|
||||
|
||||
view! {
|
||||
<div class="console">
|
||||
<nav class="rail">
|
||||
<div class="wordmark">
|
||||
"CN" <span class="wordmark-accent">"ATS"</span>
|
||||
<span class="wordmark-sub">"chat over the bus"</span>
|
||||
</div>
|
||||
|
||||
<div class="rail-label">"SUBJECTS"</div>
|
||||
<div class="rooms">
|
||||
{ROOMS
|
||||
.iter()
|
||||
.map(|(name, desc)| {
|
||||
let name = *name;
|
||||
let desc = *desc;
|
||||
view! {
|
||||
<A
|
||||
href=format!("/r/{name}")
|
||||
attr:class=move || {
|
||||
if room.get() == name { "room-link active" } else { "room-link" }
|
||||
}
|
||||
>
|
||||
<span class="room-hash">"❯"</span>
|
||||
<span class="room-text">
|
||||
<span class="room-name">{name}</span>
|
||||
<span class="room-desc">{desc}</span>
|
||||
</span>
|
||||
</A>
|
||||
}
|
||||
})
|
||||
.collect_view()}
|
||||
</div>
|
||||
|
||||
<div class="rail-foot">
|
||||
<div class="user-chip">
|
||||
<span class="user-avatar">
|
||||
{user.display_name.chars().next().unwrap_or('?').to_uppercase().to_string()}
|
||||
</span>
|
||||
<span class="user-names">
|
||||
<span class="user-display">{user.display_name.clone()}</span>
|
||||
<span class="user-handle">{format!("@{}", user.username)}</span>
|
||||
</span>
|
||||
</div>
|
||||
<a class="logout" href="/auth/logout" rel="external" title="sign out">
|
||||
"EOT ⏏"
|
||||
</a>
|
||||
</div>
|
||||
</nav>
|
||||
|
||||
<main class="deck">
|
||||
<header class="topbar">
|
||||
<div class="topbar-room">
|
||||
<h1>{move || room.get()}</h1>
|
||||
<code class="subject-code">{move || room_subject(&room.get())}</code>
|
||||
</div>
|
||||
<div
|
||||
class="topbar-status"
|
||||
class:live=move || status.get() == BusStatus::Live
|
||||
class:offline=move || status.get() == BusStatus::Offline
|
||||
>
|
||||
<span class="led"></span>
|
||||
<span class="status-text">{move || status.get().label()}</span>
|
||||
</div>
|
||||
</header>
|
||||
|
||||
<div class="stream" node_ref=list_ref>
|
||||
<Show when=move || messages.with(|m| m.is_empty())>
|
||||
<div class="stream-empty">
|
||||
<span class="empty-glyph">"[ ∅ ]"</span>
|
||||
<p>"no traffic on this subject yet — say something"</p>
|
||||
</div>
|
||||
</Show>
|
||||
<For
|
||||
each=move || messages.get()
|
||||
key=|m| m.id.clone()
|
||||
children=move |m: ChatMessage| {
|
||||
let mine = m.username == me;
|
||||
view! {
|
||||
<article class="msg" class:mine=mine>
|
||||
<span class="msg-time">{m.time}</span>
|
||||
<span class="msg-user">{m.display_name}</span>
|
||||
<span class="msg-text">{m.text}</span>
|
||||
</article>
|
||||
}
|
||||
}
|
||||
/>
|
||||
</div>
|
||||
|
||||
<form class="composer" on:submit=on_submit>
|
||||
<span class="composer-prompt">"❯"</span>
|
||||
<input
|
||||
type="text"
|
||||
class="composer-input"
|
||||
placeholder=move || format!("PUB {} …", room_subject(&room.get()))
|
||||
prop:value=move || draft.get()
|
||||
on:input=move |ev| draft.set(event_target_value(&ev))
|
||||
autocomplete="off"
|
||||
maxlength="2000"
|
||||
/>
|
||||
<button type="submit" class="composer-send" disabled=move || send.pending().get()>
|
||||
"PUB ↵"
|
||||
</button>
|
||||
</form>
|
||||
</main>
|
||||
</div>
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(feature = "hydrate")]
|
||||
fn open_event_source(
|
||||
room: &str,
|
||||
messages: RwSignal<Vec<ChatMessage>>,
|
||||
status: RwSignal<BusStatus>,
|
||||
) -> Option<web_sys::EventSource> {
|
||||
use wasm_bindgen::{prelude::Closure, JsCast};
|
||||
use web_sys::{EventSource, MessageEvent};
|
||||
|
||||
let es = EventSource::new(&format!("/sse/{room}")).ok()?;
|
||||
|
||||
let on_open = Closure::<dyn FnMut()>::new(move || status.set(BusStatus::Live));
|
||||
es.set_onopen(Some(on_open.as_ref().unchecked_ref()));
|
||||
on_open.forget();
|
||||
|
||||
let on_error = Closure::<dyn FnMut()>::new(move || status.set(BusStatus::Offline));
|
||||
es.set_onerror(Some(on_error.as_ref().unchecked_ref()));
|
||||
on_error.forget();
|
||||
|
||||
let on_message = Closure::<dyn FnMut(MessageEvent)>::new(move |ev: MessageEvent| {
|
||||
if let Some(data) = ev.data().as_string() {
|
||||
if let Ok(msg) = serde_json::from_str::<ChatMessage>(&data) {
|
||||
messages.update(|m| {
|
||||
m.push(msg);
|
||||
// keep the DOM bounded on long-lived tabs
|
||||
if m.len() > 500 {
|
||||
m.remove(0);
|
||||
}
|
||||
});
|
||||
}
|
||||
}
|
||||
});
|
||||
es.set_onmessage(Some(on_message.as_ref().unchecked_ref()));
|
||||
on_message.forget();
|
||||
|
||||
Some(es)
|
||||
}
|
||||
|
||||
#[component]
|
||||
fn NotFound() -> impl IntoView {
|
||||
view! {
|
||||
<main class="gate">
|
||||
<div class="gate-card">
|
||||
<div class="gate-badge">"ERR · NO RESPONDERS"</div>
|
||||
<h1 class="gate-title">"404"</h1>
|
||||
<p class="gate-sub">"No subscribers on this subject."</p>
|
||||
<a class="gate-btn" href="/">"⟵ back to the lobby"</a>
|
||||
</div>
|
||||
</main>
|
||||
}
|
||||
}
|
||||
+24
@@ -0,0 +1,24 @@
|
||||
use leptos::prelude::*;
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
/// The authenticated user, as established by the Kanidm OIDC flow and
|
||||
/// stored in the server-side session.
|
||||
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
|
||||
pub struct User {
|
||||
pub sub: String,
|
||||
pub username: String,
|
||||
pub display_name: String,
|
||||
}
|
||||
|
||||
pub const SESSION_USER_KEY: &str = "user";
|
||||
|
||||
/// Returns the currently signed-in user, if any.
|
||||
#[server]
|
||||
pub async fn current_user() -> Result<Option<User>, ServerFnError> {
|
||||
let session: tower_sessions::Session = leptos_axum::extract().await?;
|
||||
let user = session
|
||||
.get::<User>(SESSION_USER_KEY)
|
||||
.await
|
||||
.map_err(|e| ServerFnError::new(e.to_string()))?;
|
||||
Ok(user)
|
||||
}
|
||||
+76
@@ -0,0 +1,76 @@
|
||||
use leptos::prelude::*;
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
/// Rooms available in the UI. Each maps to the NATS subject
|
||||
/// `chat.room.<name>`, so any other NATS client on the bus can join in.
|
||||
pub const ROOMS: &[(&str, &str)] = &[
|
||||
("lobby", "general traffic"),
|
||||
("dev", "build & ship"),
|
||||
("ops", "incidents & infra"),
|
||||
("random", "off the record"),
|
||||
];
|
||||
|
||||
pub const DEFAULT_ROOM: &str = "lobby";
|
||||
|
||||
pub fn is_valid_room(room: &str) -> bool {
|
||||
ROOMS.iter().any(|(name, _)| *name == room)
|
||||
}
|
||||
|
||||
pub fn room_subject(room: &str) -> String {
|
||||
format!("chat.room.{room}")
|
||||
}
|
||||
|
||||
/// A single chat message as it travels over NATS (JSON-encoded payload).
|
||||
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
|
||||
pub struct ChatMessage {
|
||||
pub id: String,
|
||||
pub room: String,
|
||||
pub username: String,
|
||||
pub display_name: String,
|
||||
pub text: String,
|
||||
/// Pre-formatted UTC wall-clock time, e.g. "14:03:27".
|
||||
pub time: String,
|
||||
}
|
||||
|
||||
/// Publishes a message to the room's NATS subject. Requires a signed-in
|
||||
/// session; the sender identity always comes from the session, never from
|
||||
/// the client.
|
||||
#[server]
|
||||
pub async fn send_message(room: String, text: String) -> Result<(), ServerFnError> {
|
||||
use crate::auth::{User, SESSION_USER_KEY};
|
||||
use crate::server::AppState;
|
||||
|
||||
let text = text.trim().to_string();
|
||||
if text.is_empty() || text.len() > 2000 {
|
||||
return Err(ServerFnError::new("message must be 1..=2000 characters"));
|
||||
}
|
||||
if !is_valid_room(&room) {
|
||||
return Err(ServerFnError::new("unknown room"));
|
||||
}
|
||||
|
||||
let session: tower_sessions::Session = leptos_axum::extract().await?;
|
||||
let Some(user) = session
|
||||
.get::<User>(SESSION_USER_KEY)
|
||||
.await
|
||||
.map_err(|e| ServerFnError::new(e.to_string()))?
|
||||
else {
|
||||
return Err(ServerFnError::new("not signed in"));
|
||||
};
|
||||
|
||||
let state = expect_context::<AppState>();
|
||||
let msg = ChatMessage {
|
||||
id: uuid::Uuid::new_v4().to_string(),
|
||||
room: room.clone(),
|
||||
username: user.username,
|
||||
display_name: user.display_name,
|
||||
text,
|
||||
time: chrono::Utc::now().format("%H:%M:%S").to_string(),
|
||||
};
|
||||
let payload = serde_json::to_vec(&msg).map_err(|e| ServerFnError::new(e.to_string()))?;
|
||||
state
|
||||
.nats
|
||||
.publish(room_subject(&room), payload.into())
|
||||
.await
|
||||
.map_err(|e| ServerFnError::new(format!("nats publish failed: {e}")))?;
|
||||
Ok(())
|
||||
}
|
||||
+13
@@ -0,0 +1,13 @@
|
||||
pub mod app;
|
||||
pub mod auth;
|
||||
pub mod chat;
|
||||
|
||||
#[cfg(feature = "ssr")]
|
||||
pub mod server;
|
||||
|
||||
#[cfg(feature = "hydrate")]
|
||||
#[wasm_bindgen::prelude::wasm_bindgen]
|
||||
pub fn hydrate() {
|
||||
console_error_panic_hook::set_once();
|
||||
leptos::mount::hydrate_body(app::App);
|
||||
}
|
||||
+97
@@ -0,0 +1,97 @@
|
||||
#[cfg(feature = "ssr")]
|
||||
#[tokio::main]
|
||||
async fn main() -> anyhow::Result<()> {
|
||||
use axum::{
|
||||
body::Body,
|
||||
extract::State,
|
||||
http::Request,
|
||||
response::IntoResponse,
|
||||
routing::{any, get},
|
||||
Router,
|
||||
};
|
||||
use cnats::app::{shell, App};
|
||||
use cnats::server::{oidc, sse, AppState};
|
||||
use leptos::prelude::*;
|
||||
use leptos_axum::{generate_route_list, LeptosRoutes};
|
||||
use std::sync::Arc;
|
||||
use tower_sessions::{MemoryStore, SessionManagerLayer};
|
||||
|
||||
dotenvy::dotenv().ok();
|
||||
tracing_subscriber::fmt()
|
||||
.with_env_filter(
|
||||
tracing_subscriber::EnvFilter::try_from_default_env()
|
||||
.unwrap_or_else(|_| "info,cnats=debug".into()),
|
||||
)
|
||||
.init();
|
||||
|
||||
let conf = get_configuration(None)?;
|
||||
let leptos_options = conf.leptos_options;
|
||||
let addr = leptos_options.site_addr;
|
||||
let routes = generate_route_list(App);
|
||||
|
||||
let nats_url =
|
||||
std::env::var("NATS_URL").unwrap_or_else(|_| "nats://127.0.0.1:4222".to_string());
|
||||
tracing::info!(%nats_url, "connecting to NATS");
|
||||
let nats = async_nats::connect(&nats_url).await?;
|
||||
|
||||
let oidc_state = Arc::new(oidc::Oidc::from_env().await?);
|
||||
|
||||
let state = AppState {
|
||||
leptos_options: leptos_options.clone(),
|
||||
nats,
|
||||
oidc: oidc_state,
|
||||
};
|
||||
|
||||
// Dev-friendly defaults: in-memory sessions, secure cookies only when
|
||||
// COOKIE_SECURE=true (set it behind TLS in production).
|
||||
let cookie_secure = std::env::var("COOKIE_SECURE")
|
||||
.map(|v| v == "true" || v == "1")
|
||||
.unwrap_or(false);
|
||||
let session_layer = SessionManagerLayer::new(MemoryStore::default())
|
||||
.with_secure(cookie_secure)
|
||||
.with_name("cnats_session");
|
||||
|
||||
async fn server_fn_handler(
|
||||
State(state): State<AppState>,
|
||||
request: Request<Body>,
|
||||
) -> impl IntoResponse {
|
||||
leptos_axum::handle_server_fns_with_context(
|
||||
move || provide_context(state.clone()),
|
||||
request,
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
let app = Router::new()
|
||||
.route("/auth/login", get(oidc::login))
|
||||
.route("/auth/callback", get(oidc::callback))
|
||||
.route("/auth/logout", get(oidc::logout))
|
||||
.route("/sse/{room}", get(sse::room_events))
|
||||
.route("/api/{*fn_name}", any(server_fn_handler))
|
||||
.leptos_routes_with_context(
|
||||
&state,
|
||||
routes,
|
||||
{
|
||||
let state = state.clone();
|
||||
move || provide_context(state.clone())
|
||||
},
|
||||
{
|
||||
let leptos_options = leptos_options.clone();
|
||||
move || shell(leptos_options.clone())
|
||||
},
|
||||
)
|
||||
.fallback(leptos_axum::file_and_error_handler::<AppState, _>(shell))
|
||||
.layer(session_layer)
|
||||
.with_state(state);
|
||||
|
||||
tracing::info!("listening on http://{addr}");
|
||||
let listener = tokio::net::TcpListener::bind(&addr).await?;
|
||||
axum::serve(listener, app.into_make_service()).await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[cfg(not(feature = "ssr"))]
|
||||
fn main() {
|
||||
// The browser build is a cdylib; this stub only exists so `cargo check`
|
||||
// without features still succeeds.
|
||||
}
|
||||
@@ -0,0 +1,19 @@
|
||||
pub mod oidc;
|
||||
pub mod sse;
|
||||
|
||||
use axum::extract::FromRef;
|
||||
use leptos::prelude::LeptosOptions;
|
||||
use std::sync::Arc;
|
||||
|
||||
#[derive(Clone)]
|
||||
pub struct AppState {
|
||||
pub leptos_options: LeptosOptions,
|
||||
pub nats: async_nats::Client,
|
||||
pub oidc: Arc<oidc::Oidc>,
|
||||
}
|
||||
|
||||
impl FromRef<AppState> for LeptosOptions {
|
||||
fn from_ref(state: &AppState) -> Self {
|
||||
state.leptos_options.clone()
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,209 @@
|
||||
//! OIDC authorization-code + PKCE flow against Kanidm.
|
||||
//!
|
||||
//! Kanidm serves per-client OIDC discovery documents at
|
||||
//! `<KANIDM_URL>/oauth2/openid/<client_id>/.well-known/openid-configuration`,
|
||||
//! and enforces PKCE, so this module always sends a S256 challenge.
|
||||
|
||||
use axum::{
|
||||
extract::{Query, State},
|
||||
http::StatusCode,
|
||||
response::Redirect,
|
||||
};
|
||||
use openidconnect::{
|
||||
core::{CoreAuthenticationFlow, CoreClient, CoreProviderMetadata},
|
||||
AuthorizationCode, ClientId, ClientSecret, CsrfToken, EndpointMaybeSet, EndpointNotSet,
|
||||
EndpointSet, IssuerUrl, Nonce, PkceCodeChallenge, PkceCodeVerifier, RedirectUrl, Scope,
|
||||
TokenResponse,
|
||||
};
|
||||
use serde::Deserialize;
|
||||
use tower_sessions::Session;
|
||||
|
||||
use crate::auth::{User, SESSION_USER_KEY};
|
||||
|
||||
use super::AppState;
|
||||
|
||||
type OidcClient = CoreClient<
|
||||
EndpointSet, // auth endpoint
|
||||
EndpointNotSet, // device auth
|
||||
EndpointNotSet, // introspection
|
||||
EndpointNotSet, // revocation
|
||||
EndpointMaybeSet, // token endpoint (from discovery)
|
||||
EndpointMaybeSet, // userinfo endpoint (from discovery)
|
||||
>;
|
||||
|
||||
pub struct Oidc {
|
||||
client: OidcClient,
|
||||
http: openidconnect::reqwest::Client,
|
||||
}
|
||||
|
||||
const PKCE_KEY: &str = "oidc_pkce_verifier";
|
||||
const CSRF_KEY: &str = "oidc_csrf_state";
|
||||
const NONCE_KEY: &str = "oidc_nonce";
|
||||
|
||||
impl Oidc {
|
||||
/// Discovers the provider and builds the client from environment:
|
||||
/// `KANIDM_URL`, `OAUTH2_CLIENT_ID`, `OAUTH2_CLIENT_SECRET`, `PUBLIC_URL`.
|
||||
pub async fn from_env() -> anyhow::Result<Self> {
|
||||
let kanidm_url = require_env("KANIDM_URL")?;
|
||||
let client_id = require_env("OAUTH2_CLIENT_ID")?;
|
||||
let client_secret = require_env("OAUTH2_CLIENT_SECRET")?;
|
||||
let public_url = require_env("PUBLIC_URL")?;
|
||||
|
||||
let issuer = IssuerUrl::new(format!(
|
||||
"{}/oauth2/openid/{}",
|
||||
kanidm_url.trim_end_matches('/'),
|
||||
client_id
|
||||
))?;
|
||||
let redirect = RedirectUrl::new(format!(
|
||||
"{}/auth/callback",
|
||||
public_url.trim_end_matches('/')
|
||||
))?;
|
||||
|
||||
// Never follow redirects when talking to the IdP (SSRF hygiene).
|
||||
let http = openidconnect::reqwest::ClientBuilder::new()
|
||||
.redirect(openidconnect::reqwest::redirect::Policy::none())
|
||||
.build()?;
|
||||
|
||||
tracing::info!(issuer = %issuer.as_str(), "discovering OIDC provider");
|
||||
let metadata = CoreProviderMetadata::discover_async(issuer, &http).await?;
|
||||
let client = CoreClient::from_provider_metadata(
|
||||
metadata,
|
||||
ClientId::new(client_id),
|
||||
Some(ClientSecret::new(client_secret)),
|
||||
)
|
||||
.set_redirect_uri(redirect);
|
||||
|
||||
Ok(Self { client, http })
|
||||
}
|
||||
}
|
||||
|
||||
fn require_env(name: &str) -> anyhow::Result<String> {
|
||||
std::env::var(name).map_err(|_| anyhow::anyhow!("missing required env var {name}"))
|
||||
}
|
||||
|
||||
type HandlerError = (StatusCode, String);
|
||||
|
||||
fn internal(err: impl std::fmt::Display) -> HandlerError {
|
||||
tracing::error!("oidc error: {err}");
|
||||
(
|
||||
StatusCode::INTERNAL_SERVER_ERROR,
|
||||
"authentication failed; see server logs".to_string(),
|
||||
)
|
||||
}
|
||||
|
||||
/// GET /auth/login — stash PKCE/state/nonce in the session and bounce to Kanidm.
|
||||
pub async fn login(State(state): State<AppState>, session: Session) -> Result<Redirect, HandlerError> {
|
||||
let (pkce_challenge, pkce_verifier) = PkceCodeChallenge::new_random_sha256();
|
||||
|
||||
let (auth_url, csrf_state, nonce) = state
|
||||
.oidc
|
||||
.client
|
||||
.authorize_url(
|
||||
CoreAuthenticationFlow::AuthorizationCode,
|
||||
CsrfToken::new_random,
|
||||
Nonce::new_random,
|
||||
)
|
||||
.add_scope(Scope::new("openid".to_string()))
|
||||
.add_scope(Scope::new("profile".to_string()))
|
||||
.add_scope(Scope::new("email".to_string()))
|
||||
.set_pkce_challenge(pkce_challenge)
|
||||
.url();
|
||||
|
||||
session
|
||||
.insert(PKCE_KEY, pkce_verifier.secret())
|
||||
.await
|
||||
.map_err(internal)?;
|
||||
session
|
||||
.insert(CSRF_KEY, csrf_state.secret())
|
||||
.await
|
||||
.map_err(internal)?;
|
||||
session
|
||||
.insert(NONCE_KEY, nonce.secret())
|
||||
.await
|
||||
.map_err(internal)?;
|
||||
|
||||
Ok(Redirect::to(auth_url.as_str()))
|
||||
}
|
||||
|
||||
#[derive(Deserialize)]
|
||||
pub struct CallbackParams {
|
||||
code: String,
|
||||
state: String,
|
||||
}
|
||||
|
||||
/// GET /auth/callback — verify state, exchange the code, verify the ID token,
|
||||
/// and store the user in the session.
|
||||
pub async fn callback(
|
||||
State(state): State<AppState>,
|
||||
session: Session,
|
||||
Query(params): Query<CallbackParams>,
|
||||
) -> Result<Redirect, HandlerError> {
|
||||
let stored_csrf: Option<String> = session.remove(CSRF_KEY).await.map_err(internal)?;
|
||||
let pkce_verifier: Option<String> = session.remove(PKCE_KEY).await.map_err(internal)?;
|
||||
let nonce: Option<String> = session.remove(NONCE_KEY).await.map_err(internal)?;
|
||||
|
||||
let (Some(stored_csrf), Some(pkce_verifier), Some(nonce)) =
|
||||
(stored_csrf, pkce_verifier, nonce)
|
||||
else {
|
||||
return Err((
|
||||
StatusCode::BAD_REQUEST,
|
||||
"no login in progress; start again at /auth/login".to_string(),
|
||||
));
|
||||
};
|
||||
if params.state != stored_csrf {
|
||||
return Err((
|
||||
StatusCode::BAD_REQUEST,
|
||||
"state mismatch; start again at /auth/login".to_string(),
|
||||
));
|
||||
}
|
||||
|
||||
let oidc = &state.oidc;
|
||||
let token_response = oidc
|
||||
.client
|
||||
.exchange_code(AuthorizationCode::new(params.code))
|
||||
.map_err(internal)?
|
||||
.set_pkce_verifier(PkceCodeVerifier::new(pkce_verifier))
|
||||
.request_async(&oidc.http)
|
||||
.await
|
||||
.map_err(internal)?;
|
||||
|
||||
let id_token = token_response
|
||||
.id_token()
|
||||
.ok_or_else(|| internal("provider returned no ID token"))?;
|
||||
let claims = id_token
|
||||
.claims(&oidc.client.id_token_verifier(), &Nonce::new(nonce))
|
||||
.map_err(internal)?;
|
||||
|
||||
let username = claims
|
||||
.preferred_username()
|
||||
.map(|u| u.as_str().to_string())
|
||||
.or_else(|| claims.email().map(|e| e.as_str().to_string()))
|
||||
.unwrap_or_else(|| claims.subject().as_str().to_string());
|
||||
let display_name = claims
|
||||
.name()
|
||||
.and_then(|n| n.get(None))
|
||||
.map(|n| n.as_str().to_string())
|
||||
.unwrap_or_else(|| username.clone());
|
||||
|
||||
let user = User {
|
||||
sub: claims.subject().as_str().to_string(),
|
||||
username,
|
||||
display_name,
|
||||
};
|
||||
|
||||
// Rotate the session id on privilege change, then store the user.
|
||||
session.cycle_id().await.map_err(internal)?;
|
||||
session
|
||||
.insert(SESSION_USER_KEY, &user)
|
||||
.await
|
||||
.map_err(internal)?;
|
||||
|
||||
tracing::info!(user = %user.username, "signed in");
|
||||
Ok(Redirect::to("/"))
|
||||
}
|
||||
|
||||
/// GET /auth/logout — drop the session.
|
||||
pub async fn logout(session: Session) -> Result<Redirect, HandlerError> {
|
||||
session.flush().await.map_err(internal)?;
|
||||
Ok(Redirect::to("/"))
|
||||
}
|
||||
@@ -0,0 +1,58 @@
|
||||
//! Server-Sent Events bridge: one NATS subscription per connected browser,
|
||||
//! scoped to a single room subject.
|
||||
|
||||
use std::{convert::Infallible, time::Duration};
|
||||
|
||||
use axum::{
|
||||
extract::{Path, State},
|
||||
http::StatusCode,
|
||||
response::sse::{Event, KeepAlive, Sse},
|
||||
};
|
||||
use futures::{Stream, StreamExt};
|
||||
use tower_sessions::Session;
|
||||
|
||||
use crate::auth::{User, SESSION_USER_KEY};
|
||||
use crate::chat::{is_valid_room, room_subject};
|
||||
|
||||
use super::AppState;
|
||||
|
||||
/// GET /sse/{room} — stream the room's NATS subject to the browser.
|
||||
pub async fn room_events(
|
||||
Path(room): Path<String>,
|
||||
State(state): State<AppState>,
|
||||
session: Session,
|
||||
) -> Result<Sse<impl Stream<Item = Result<Event, Infallible>>>, (StatusCode, &'static str)> {
|
||||
let signed_in = session
|
||||
.get::<User>(SESSION_USER_KEY)
|
||||
.await
|
||||
.ok()
|
||||
.flatten()
|
||||
.is_some();
|
||||
if !signed_in {
|
||||
return Err((StatusCode::UNAUTHORIZED, "sign in first"));
|
||||
}
|
||||
if !is_valid_room(&room) {
|
||||
return Err((StatusCode::NOT_FOUND, "unknown room"));
|
||||
}
|
||||
|
||||
let subscriber = state
|
||||
.nats
|
||||
.subscribe(room_subject(&room))
|
||||
.await
|
||||
.map_err(|e| {
|
||||
tracing::error!("nats subscribe failed: {e}");
|
||||
(StatusCode::BAD_GATEWAY, "message bus unavailable")
|
||||
})?;
|
||||
|
||||
let stream = subscriber.map(|msg| {
|
||||
Ok(Event::default()
|
||||
.event("message")
|
||||
.data(String::from_utf8_lossy(&msg.payload).into_owned()))
|
||||
});
|
||||
|
||||
Ok(Sse::new(stream).keep_alive(
|
||||
KeepAlive::new()
|
||||
.interval(Duration::from_secs(15))
|
||||
.text("ping"),
|
||||
))
|
||||
}
|
||||
Reference in New Issue
Block a user