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 health;
39pub mod limits;
40pub mod public;
41mod query;
42pub(crate) mod strike_state;
43pub(crate) mod subscribe;
44pub mod wellknown;
45pub mod xrpc;
46
47pub use admin::{AdminConfig, admin_router};
48pub use create_report::{CreateReportConfig, create_report_router};
49pub use public::public_router;
50pub use subscribe::current_retention_floor;
51
52/// Sweep-execution policy for the subscribeLabels retention sweep (§F4
53/// retention task). Distinct from [`SubscribeConfig::retention_days`] —
54/// that field is the cutoff source of truth (read-side floor + write-
55/// side cutoff); this struct holds *how* the sweep runs, not *when*
56/// the cutoff lies. Operators tune the two independently:
57///
58/// - `[subscribe].retention_days` (logical owner: [`SubscribeConfig`])
59/// sets the rolling window ("180 days back").
60/// - `[retention]` (logical owner: this struct) sets sweep schedule
61/// and batching ("run at 04:00 UTC, 1000 rows per transaction").
62///
63/// Defaults match §F4 prose: sweep enabled, runs daily at 04:00 UTC,
64/// 1000-row batches.
65#[derive(Debug, Clone)]
66pub struct RetentionConfig {
67 /// Master toggle for the scheduled sweep. When `false`, no
68 /// scheduled sweep runs; the manual CLI / admin-XRPC path still
69 /// works (operators retain explicit control). Default `true`.
70 pub sweep_enabled: bool,
71 /// UTC hour-of-day (0..=23) at which the scheduled sweep fires.
72 /// Default 4 (04:00 UTC — quiet hour for most operator regions,
73 /// matches §F4 example).
74 pub sweep_run_at_utc_hour: u8,
75 /// Rows per DELETE transaction. Larger batches are throughput-
76 /// efficient but hold the writer for longer; smaller batches let
77 /// normal label writes interleave with finer granularity. Default
78 /// 1000. Distinct from [`SubscribeConfig::batch_size`] (replay
79 /// batching is latency-sensitive, sweep batching is throughput-
80 /// sensitive — tune independently).
81 pub sweep_batch_size: i64,
82}
83
84impl Default for RetentionConfig {
85 fn default() -> Self {
86 Self {
87 sweep_enabled: true,
88 sweep_run_at_utc_hour: 4,
89 sweep_batch_size: 1000,
90 }
91 }
92}
93pub use did_document::did_document_router;
94pub use health::health_router;
95pub use limits::Limiter;
96pub use wellknown::wellknown_router;
97
98/// Tunables for the subscribeLabels endpoint. All defaults match §F4.
99#[derive(Debug, Clone)]
100pub struct SubscribeConfig {
101 /// Maximum concurrent subscribers from one IP (§F4 default 8).
102 pub per_ip_cap: usize,
103 /// Maximum concurrent subscribers overall (§F4 default 256).
104 pub global_cap: usize,
105 /// Replay batch size in rows (§F4 "up to 1000 sequence rows per batch").
106 pub batch_size: i64,
107 /// Ping cadence (§F4 "30s").
108 pub ping_interval: Duration,
109 /// Pong silence that closes the connection (§F4 "90s").
110 pub pong_timeout: Duration,
111 /// Rolling retention window. `None` disables the read-side floor
112 /// (useful in tests and when an operator wants unbounded replay).
113 /// Default per §F4: 180 days.
114 pub retention_days: Option<u32>,
115}
116
117impl Default for SubscribeConfig {
118 fn default() -> Self {
119 Self {
120 per_ip_cap: 8,
121 global_cap: 256,
122 batch_size: 1000,
123 ping_interval: Duration::from_secs(30),
124 pong_timeout: Duration::from_secs(90),
125 retention_days: Some(180),
126 }
127 }
128}
129
130/// Build a router exposing the public label endpoints. The caller owns
131/// the pool and writer; dropping all cloned `WriterHandle`s signals
132/// shutdown.
133///
134/// CORS is applied to `queryLabels` only (§F3 "accepts requests from any
135/// origin"). subscribeLabels is WebSocket — browsers don't apply CORS to
136/// WS connections the same way, and the WS handshake rejects cross-
137/// origin reads at the Sec-WebSocket-* layer. Credentials are never
138/// echoed (`allow_credentials` defaults to false on `CorsLayer`).
139pub fn router(pool: Pool<Sqlite>, writer: WriterHandle, config: SubscribeConfig) -> Router {
140 let limiter = Limiter::new(config.global_cap, config.per_ip_cap);
141 let state = subscribe::AppState {
142 pool,
143 writer,
144 limiter,
145 config: Arc::new(config),
146 };
147
148 let query_cors = CorsLayer::new()
149 .allow_origin(Any)
150 .allow_methods([Method::GET, Method::OPTIONS]);
151
152 Router::new()
153 .route(
154 "/xrpc/com.atproto.label.subscribeLabels",
155 get(subscribe::handler),
156 )
157 .route(
158 "/xrpc/com.atproto.label.queryLabels",
159 get(query::handler).layer(query_cors),
160 )
161 .layer(Extension(state))
162}
163
164/// Serve `router` on `addr`. Returns when the listener errors or the
165/// surrounding runtime shuts down.
166pub async fn serve(router: Router, addr: SocketAddr) -> Result<()> {
167 let listener = TcpListener::bind(addr).await?;
168 axum::serve(
169 listener,
170 router.into_make_service_with_connect_info::<SocketAddr>(),
171 )
172 .await?;
173 Ok(())
174}