use std::net::SocketAddr;
use std::sync::Arc;
use std::time::Duration;
use axum::Extension;
use axum::Router;
use axum::http::Method;
use axum::routing::get;
use sqlx::{Pool, Sqlite};
use tokio::net::TcpListener;
use tower_http::cors::{Any, CorsLayer};
use crate::error::Result;
use crate::writer::WriterHandle;
pub mod admin;
mod create_report;
pub mod did_document;
pub mod health;
pub mod limits;
pub mod public;
mod query;
pub(crate) mod strike_state;
pub(crate) mod subscribe;
pub mod wellknown;
pub mod xrpc;
pub use admin::{AdminConfig, admin_router};
pub use create_report::{CreateReportConfig, create_report_router};
pub use public::public_router;
pub use subscribe::current_retention_floor;
#[derive(Debug, Clone)]
pub struct RetentionConfig {
pub sweep_enabled: bool,
pub sweep_run_at_utc_hour: u8,
pub sweep_batch_size: i64,
}
impl Default for RetentionConfig {
fn default() -> Self {
Self {
sweep_enabled: true,
sweep_run_at_utc_hour: 4,
sweep_batch_size: 1000,
}
}
}
pub use did_document::did_document_router;
pub use health::health_router;
pub use limits::Limiter;
pub use wellknown::wellknown_router;
#[derive(Debug, Clone)]
pub struct SubscribeConfig {
pub per_ip_cap: usize,
pub global_cap: usize,
pub batch_size: i64,
pub ping_interval: Duration,
pub pong_timeout: Duration,
pub retention_days: Option<u32>,
}
impl Default for SubscribeConfig {
fn default() -> Self {
Self {
per_ip_cap: 8,
global_cap: 256,
batch_size: 1000,
ping_interval: Duration::from_secs(30),
pong_timeout: Duration::from_secs(90),
retention_days: Some(180),
}
}
}
pub fn router(pool: Pool<Sqlite>, writer: WriterHandle, config: SubscribeConfig) -> Router {
let limiter = Limiter::new(config.global_cap, config.per_ip_cap);
let state = subscribe::AppState {
pool,
writer,
limiter,
config: Arc::new(config),
};
let query_cors = CorsLayer::new()
.allow_origin(Any)
.allow_methods([Method::GET, Method::OPTIONS]);
Router::new()
.route(
"/xrpc/com.atproto.label.subscribeLabels",
get(subscribe::handler),
)
.route(
"/xrpc/com.atproto.label.queryLabels",
get(query::handler).layer(query_cors),
)
.layer(Extension(state))
}
pub async fn serve(router: Router, addr: SocketAddr) -> Result<()> {
let listener = TcpListener::bind(addr).await?;
axum::serve(
listener,
router.into_make_service_with_connect_info::<SocketAddr>(),
)
.await?;
Ok(())
}