pub mod config; mod gate; use axum::{ Json, Router, extract::{DefaultBodyLimit, State}, http::StatusCode, response::{IntoResponse, Response}, routing::{get, post}, }; use bobbin_runtime::{ReqwestHttp, UnixMicros}; use gate::Gate; use jacquard_axum::service_auth::{ServiceAuthConfig, ServiceAuthError}; use jacquard_common::{deps::fluent_uri::Uri, service_auth::ServiceAuthClaims, types::string::Did}; use jacquard_identity::{JacquardResolver, resolver::PlcSource}; use k256::ecdsa::SigningKey; use serde::{Deserialize, Serialize}; use serde_json::{Value, json}; use std::{ borrow::Cow, sync::Arc, time::{Duration, SystemTime, UNIX_EPOCH}, }; use tangled_axum::atproto::{AtprotoService, ClaimsPolicy, ExtractAnyServiceAuth, Verified}; const OWNER: &str = "atproto repo:* blob:*/* rpc:* identity:* account:*?action=manage transition:generic transition:chat.bsky transition:email"; type Directory = JacquardResolver; struct RuntimeConfig { did: String, tranquil_url: String, signing_key: SigningKey, } #[derive(Clone)] pub struct App { config: Arc, http: reqwest::Client, auth: ServiceAuthConfig>, gate: Arc, } impl App { /// Validate configuration and prepare shared clients, auth policy, and quota lock. pub fn new(config: &config::Config) -> anyhow::Result { config.validate()?; let mut http = reqwest::Client::builder() .connect_timeout(Duration::from_secs(5)) .timeout(Duration::from_secs(30)) .redirect(reqwest::redirect::Policy::none()); if let Some(path) = &config.extra_ca_file { http = http.add_root_certificate(reqwest::Certificate::from_pem(&std::fs::read(path)?)?); } let http = http.build()?; let plc = format!("{}/", config.plc_url.trim_end_matches('/')); let directory = Arc::new( JacquardResolver::new(ReqwestHttp::new(http.clone()), Default::default()) .with_plc_source(PlcSource::PlcDirectory { base: Uri::parse(plc.as_str())?.to_owned(), }), ); let signing_key = config.signing_key()?; let gate = Gate::builder( http.clone(), config.did.clone(), signing_key.clone(), config.deliberi_url.trim_end_matches('/'), config.deliberi_did.clone(), ) .require_verified_email(config.require_verified_email) .build(); Ok(Self { config: Arc::new(RuntimeConfig { did: config.did.clone(), tranquil_url: config.tranquil_url.trim_end_matches('/').into(), signing_key, }), http, auth: ServiceAuthConfig::new(Did::new_owned(config.tranquil_did.clone())?, directory), gate: Arc::new(gate), }) } } impl AtprotoService for App { type Resolver = Directory; fn resolver(&self) -> &Directory { self.auth.resolver() } async fn verify_claims( &self, claims: &ServiceAuthClaims, ) -> Result { if claims.lxm.as_ref().map(|v| v.as_str()) != Some("farm.tranquil.delegation.createAccount") { return Err(ServiceAuthError::MethodBindingRequired); } ClaimsPolicy::from(&self.auth) .require_lxm(true) .check(claims, UnixMicros::new(now() * 1_000_000)) .await } } pub(crate) fn now() -> u64 { SystemTime::now() .duration_since(UNIX_EPOCH) .expect("system clock") .as_secs() } /// An xrpc error payload: `{"error": ..., "message": ...}`. pub(crate) struct Error { pub(crate) status: StatusCode, pub(crate) error: Cow<'static, str>, pub(crate) message: Option, } impl Error { pub(crate) fn code(status: StatusCode, error: &'static str) -> Self { Self { status, error: Cow::Borrowed(error), message: None, } } } impl IntoResponse for Error { fn into_response(self) -> Response { let mut body = json!({"error": self.error}); if let Some(message) = self.message { body["message"] = Value::String(message); } (self.status, Json(body)).into_response() } } pub(crate) fn upstream(e: impl std::fmt::Display) -> Error { tracing::warn!(error = %e, "upstream request failed"); Error::code(StatusCode::BAD_GATEWAY, "UpstreamUnavailable") } impl From for Error { fn from(error: reqwest::Error) -> Self { upstream(error) } } #[derive(Default, Deserialize)] struct UpstreamError { error: Option, message: Option, } /// Pass a 2xx through, relay a 4xx with the upstream's own error code so the /// caller can tell a taken handle from a delegate cap, and fail closed on 5xx. pub(crate) async fn relay(response: reqwest::Response) -> Result { let status = response.status(); if status.is_success() { return Ok(response); } if status.is_client_error() { let payload: UpstreamError = match response.json().await { Ok(payload) => payload, Err(e) => { tracing::warn!(%status, error = %e, "upstream error body unreadable"); UpstreamError::default() } }; return Err(Error { status, error: payload .error .map_or(Cow::Borrowed("UpstreamRejected"), Cow::Owned), message: payload.message, }); } Err(upstream(format!("upstream responded {status}"))) } #[derive(Deserialize)] #[serde(deny_unknown_fields, rename_all = "camelCase")] struct CreateInput { handle: String, } #[derive(Debug, Deserialize, Serialize)] #[serde(rename_all = "camelCase")] struct Account { did: String, handle: String, #[serde(default)] controller_did: String, } impl App { async fn create(&self, actor: &str, token: &str, input: CreateInput) -> Result { let cfg = &self.config; self.gate.check(actor).await?; let response = self .http .post(format!( "{}/xrpc/_delegation.createDelegatedAccount", cfg.tranquil_url )) .bearer_auth(token) .json(&json!({"handle": input.handle, "controllerScopes": OWNER})) .send() .await?; let mut account: Account = relay(response).await?.json().await?; account.controller_did = actor.to_owned(); tracing::info!(controller = actor, did = %account.did, "created delegated account"); Ok(account) } } async fn create( State(app): State, ExtractAnyServiceAuth(claims, token): ExtractAnyServiceAuth, Json(input): Json, ) -> Result, Error> { if input.handle.is_empty() || input.handle.len() > 253 || input.handle.trim() != input.handle { return Err(Error::code(StatusCode::BAD_REQUEST, "InvalidHandle")); } tokio::spawn(async move { app.create(claims.iss.as_str(), &token, input).await }) .await .map_err(upstream)? .map(Json) } async fn did_document(State(app): State) -> Json { let did = &app.config.did; let mut key = vec![0xe7, 0x01]; // secp256k1-pub multicodec key.extend_from_slice( app.config .signing_key .verifying_key() .to_encoded_point(true) .as_bytes(), ); Json(json!({ "@context": ["https://www.w3.org/ns/did/v1"], "id": did, "verificationMethod": [ { "id": format!("{did}#atproto"), "type": "Multikey", "controller": did, "publicKeyMultibase": format!("z{}", bs58::encode(key).into_string()) } ] })) } async fn health() -> Json { Json(json!({"status": "ok"})) } pub fn router(app: App) -> Router { Router::new() .route("/xrpc/_health", get(health)) .route("/.well-known/did.json", get(did_document)) .route("/xrpc/sh.tangled.delegation.createAccount", post(create)) .layer(DefaultBodyLimit::max(4096)) .with_state(app) } #[cfg(test)] mod tests;