Something went wrong. Try again.
A minimal atproto Personal Data Server, on SQLite. github.com/eth0net/manapds
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315//! Routing, what every handler is handed, and the methods small enough to//! sit beside it.
use std::sync::Arc;
use axum::{ Json, Router, extract::{DefaultBodyLimit, FromRef, State}, middleware, routing::get,};use serde::Serialize;use tower_http::{ cors::{Any, CorsLayer}, normalize_path::NormalizePath, trace::TraceLayer,};
mod admin;mod app_password;mod blob;mod identity;mod repo;mod session;mod signup;
use crate::account;use crate::config::Config;use crate::store;use crate::xrpc::{self, auth::Tokens, limit::Limits};
/// Everything a handler can ask the router for.#[derive(Clone, Debug)]pub struct Context { /// How this server was started. pub config: Arc<Config>, /// The accounts on it, and the sessions open against them. pub accounts: Arc<account::Manager>, /// The secret every session token is signed under. pub tokens: Tokens, /// The budgets every caller is held to, or `None` where they are off. pub limits: Option<Arc<Limits>>,}
impl Context { /// Opens the databases under the configured data directory. /// /// # Errors /// /// If the directory cannot be made, or a database cannot be opened or has /// been migrated past what this server reads. pub fn open(config: Config) -> Result<Self, store::Error> { let config = Arc::new(config); let directory = store::Directory::new(config.data_directory.clone()); let accounts = store::Accounts::open(&directory.accounts())?; let sequencer = store::Sequencer::open(&directory.sequencer())?; let tokens = Tokens::new(config.jwt_secret.reveal(), config.service_did.clone()); Ok(Self::new( Arc::clone(&config), Arc::new(account::Manager::new( config, accounts, sequencer, tokens.clone(), )), tokens, )) }
/// The same, around an account layer somebody else opened, which is what a /// test wants. #[must_use] pub fn new(config: Arc<Config>, accounts: Arc<account::Manager>, tokens: Tokens) -> Self { Self { limits: Limits::new(&config).map(Arc::new), config, accounts, tokens, } }}
impl FromRef<Context> for Tokens { fn from_ref(context: &Context) -> Self { context.tokens.clone() }}
impl FromRef<Context> for Arc<Config> { fn from_ref(context: &Context) -> Self { Arc::clone(&context.config) }}
impl FromRef<Context> for Arc<account::Manager> { fn from_ref(context: &Context) -> Self { Arc::clone(&context.accounts) }}
impl FromRef<Context> for Option<Arc<Limits>> { fn from_ref(context: &Context) -> Self { context.limits.clone() }}
impl From<account::Error> for xrpc::Error { fn from(error: account::Error) -> Self { use account::Error; // Storage that would not answer this time is worth saying so about: // the caller can come back, where a fault on this side it cannot. if error.busy() { return Self::new(xrpc::Status::NotEnoughResources).saying(error.to_string()); } let status = match error { Error::Credentials => xrpc::Status::AuthenticationRequired, Error::Plc(_) | Error::Repo(_) | Error::Storage(_) => { return Self::internal(error.to_string()); } _ => xrpc::Status::InvalidRequest, }; Self::new(status) .named(error.name()) .saying(error.to_string()) }}
/// Builds the whole request surface. Takes what it needs rather than opening/// it, so a test can drive the server without touching the environment or the/// disk.////// A trailing slash is trimmed before anything routes, because upstream serves/// `/xrpc/<nsid>/` and a client that sends one should not be told the method/// does not exist.#[must_use]pub fn router(context: Context) -> NormalizePath<Router> { NormalizePath::trim_trailing_slash(routes(context))}
fn routes(context: Context) -> Router { let limits = context.limits.clone(); // Only the blob route takes a body this large; every other one is a small // JSON object and keeps axum's own limit. let uploads = usize::try_from(context.config.blob_upload_limit).unwrap_or(usize::MAX); let router = Router::new() .route("/", get(root)) .route("/robots.txt", get(robots)) .route("/xrpc/_health", get(health)) .route( "/xrpc/com.atproto.server.describeServer", xrpc::query(describe_server), ) .route( "/xrpc/com.atproto.server.createInviteCode", xrpc::procedure(admin::create_invite_code), ) .route( "/xrpc/com.atproto.server.createInviteCodes", xrpc::procedure(admin::create_invite_codes), ) .route( "/xrpc/com.atproto.identity.resolveHandle", xrpc::query(identity::resolve_handle), ) .route( "/xrpc/com.atproto.server.createAccount", xrpc::procedure(signup::create), ) .route( "/xrpc/com.atproto.server.createAppPassword", xrpc::procedure(app_password::create), ) .route( "/xrpc/com.atproto.server.listAppPasswords", xrpc::query(app_password::list), ) .route( "/xrpc/com.atproto.server.revokeAppPassword", xrpc::procedure(app_password::revoke), ) .route( "/xrpc/com.atproto.repo.createRecord", xrpc::procedure(repo::create_record), ) .route( "/xrpc/com.atproto.repo.putRecord", xrpc::procedure(repo::put_record), ) .route( "/xrpc/com.atproto.repo.deleteRecord", xrpc::procedure(repo::delete_record), ) .route( "/xrpc/com.atproto.repo.applyWrites", xrpc::procedure(repo::apply_writes), ) .route( "/xrpc/com.atproto.repo.uploadBlob", xrpc::procedure(blob::upload).layer(DefaultBodyLimit::max(uploads)), ) .route("/xrpc/com.atproto.sync.getBlob", xrpc::query(blob::get)) .route( "/xrpc/com.atproto.repo.getRecord", xrpc::query(repo::get_record), ) .route( "/xrpc/com.atproto.repo.listRecords", xrpc::query(repo::list_records), ) .route( "/xrpc/com.atproto.repo.describeRepo", xrpc::query(repo::describe_repo), ) .route( "/xrpc/com.atproto.server.createSession", xrpc::procedure(session::create), ) .route( "/xrpc/com.atproto.server.getSession", xrpc::query(session::get), ) .route( "/xrpc/com.atproto.server.refreshSession", xrpc::procedure(session::refresh), ) .route( "/xrpc/com.atproto.server.deleteSession", xrpc::procedure(session::delete), ) .fallback(xrpc::fallback) .with_state(context) .layer(middleware::from_fn(xrpc::auth::private));
// CORS goes outside the budget so that a browser is told why it was // refused rather than being told nothing at all. match limits { Some(limits) => router.layer(middleware::from_fn_with_state(limits, xrpc::limit::global)), None => router, } .layer(cors()) .layer(TraceLayer::new_for_http())}
/// Anything may call this server, since every method either needs credentials/// or is public to begin with.fn cors() -> CorsLayer { CorsLayer::new() .allow_origin(Any) .allow_methods(Any) .allow_headers(Any) .expose_headers(Any) .max_age(std::time::Duration::from_hours(24))}
async fn root() -> &'static str { "This is an AT Protocol Personal Data Server.\n\nMost of it is under /xrpc/.\n"}
async fn robots() -> &'static str { "# Crawling the public API is allowed\nUser-agent: *\nAllow: /\n"}
#[derive(Debug, Serialize)]struct Health { version: &'static str,}
async fn health() -> Json<Health> { Json(Health { version: env!("CARGO_PKG_VERSION"), })}
#[derive(Debug, Serialize)]#[serde(rename_all = "camelCase")]struct DescribeServer { did: String, available_user_domains: Vec<String>, invite_code_required: bool, blob_upload_limit: u64, links: Links, contact: Contact,}
#[derive(Debug, Serialize)]#[serde(rename_all = "camelCase")]struct Links { #[serde(skip_serializing_if = "Option::is_none")] privacy_policy: Option<String>, #[serde(skip_serializing_if = "Option::is_none")] terms_of_service: Option<String>,}
#[derive(Debug, Serialize)]struct Contact { #[serde(skip_serializing_if = "Option::is_none")] email: Option<String>,}
async fn describe_server(State(config): State<Arc<Config>>) -> Json<DescribeServer> { Json(DescribeServer { did: config.service_did.clone(), available_user_domains: config.handle_domains.clone(), invite_code_required: config.invite_required, blob_upload_limit: config.blob_upload_limit, links: Links { privacy_policy: config.privacy_policy_url.clone(), terms_of_service: config.terms_of_service_url.clone(), }, contact: Contact { email: config.contact_email.clone(), }, })}