Skip to main content

cairn_mod/server/
mod.rs

1//! HTTP server: XRPC endpoints over axum.
2//!
3//! Currently exposes one route — `GET /xrpc/com.atproto.label.subscribeLabels`
4//! (§F4, WebSocket). Other endpoints land in their respective issues: the
5//! binary wire-up (main.rs serving a long-lived listener) is parked until
6//! #8 (queryLabels) so both public endpoints compose into one router.
7//!
8//! Design notes:
9//!
10//! - [`SubscribeConfig`] centralizes every runtime knob §F4 mentions: per-IP
11//!   cap, global cap, batch size, ping/pong cadence, retention window.
12//!   `Default` values match §F4 defaults so operators only override what
13//!   their deployment requires.
14//! - The [`router`] constructor takes ownership of a [`Pool<Sqlite>`] and a
15//!   [`WriterHandle`]; both clone cheaply. The limiter is created here and
16//!   wrapped in `Arc` so tests can inspect it with [`Limiter::snapshot`].
17//! - [`serve`] is a thin helper over `axum::serve`, mostly to keep
18//!   integration tests from reimplementing the listen/accept dance.
19
20use std::net::SocketAddr;
21use std::sync::Arc;
22use std::time::Duration;
23
24use axum::Extension;
25use axum::Router;
26use axum::http::Method;
27use axum::routing::get;
28use sqlx::{Pool, Sqlite};
29use tokio::net::TcpListener;
30use tower_http::cors::{Any, CorsLayer};
31
32use crate::error::Result;
33use crate::writer::WriterHandle;
34
35pub mod admin;
36mod create_report;
37pub mod did_document;
38pub mod limits;
39mod query;
40pub(crate) mod subscribe;
41pub mod wellknown;
42pub mod xrpc;
43
44pub use admin::{AdminConfig, admin_router};
45pub use create_report::{CreateReportConfig, create_report_router};
46pub use did_document::did_document_router;
47pub use limits::Limiter;
48pub use wellknown::wellknown_router;
49
50/// Tunables for the subscribeLabels endpoint. All defaults match §F4.
51#[derive(Debug, Clone)]
52pub struct SubscribeConfig {
53    /// Maximum concurrent subscribers from one IP (§F4 default 8).
54    pub per_ip_cap: usize,
55    /// Maximum concurrent subscribers overall (§F4 default 256).
56    pub global_cap: usize,
57    /// Replay batch size in rows (§F4 "up to 1000 sequence rows per batch").
58    pub batch_size: i64,
59    /// Ping cadence (§F4 "30s").
60    pub ping_interval: Duration,
61    /// Pong silence that closes the connection (§F4 "90s").
62    pub pong_timeout: Duration,
63    /// Rolling retention window. `None` disables the read-side floor
64    /// (useful in tests and when an operator wants unbounded replay).
65    /// Default per §F4: 180 days.
66    pub retention_days: Option<u32>,
67}
68
69impl Default for SubscribeConfig {
70    fn default() -> Self {
71        Self {
72            per_ip_cap: 8,
73            global_cap: 256,
74            batch_size: 1000,
75            ping_interval: Duration::from_secs(30),
76            pong_timeout: Duration::from_secs(90),
77            retention_days: Some(180),
78        }
79    }
80}
81
82/// Build a router exposing the public label endpoints. The caller owns
83/// the pool and writer; dropping all cloned `WriterHandle`s signals
84/// shutdown.
85///
86/// CORS is applied to `queryLabels` only (§F3 "accepts requests from any
87/// origin"). subscribeLabels is WebSocket — browsers don't apply CORS to
88/// WS connections the same way, and the WS handshake rejects cross-
89/// origin reads at the Sec-WebSocket-* layer. Credentials are never
90/// echoed (`allow_credentials` defaults to false on `CorsLayer`).
91pub fn router(pool: Pool<Sqlite>, writer: WriterHandle, config: SubscribeConfig) -> Router {
92    let limiter = Limiter::new(config.global_cap, config.per_ip_cap);
93    let state = subscribe::AppState {
94        pool,
95        writer,
96        limiter,
97        config: Arc::new(config),
98    };
99
100    let query_cors = CorsLayer::new()
101        .allow_origin(Any)
102        .allow_methods([Method::GET, Method::OPTIONS]);
103
104    Router::new()
105        .route(
106            "/xrpc/com.atproto.label.subscribeLabels",
107            get(subscribe::handler),
108        )
109        .route(
110            "/xrpc/com.atproto.label.queryLabels",
111            get(query::handler).layer(query_cors),
112        )
113        .layer(Extension(state))
114}
115
116/// Serve `router` on `addr`. Returns when the listener errors or the
117/// surrounding runtime shuts down.
118pub async fn serve(router: Router, addr: SocketAddr) -> Result<()> {
119    let listener = TcpListener::bind(addr).await?;
120    axum::serve(
121        listener,
122        router.into_make_service_with_connect_info::<SocketAddr>(),
123    )
124    .await?;
125    Ok(())
126}