Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
8.5 kB · 269 lines
Rust
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270pub 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<ReqwestHttp>;
struct RuntimeConfig { did: String, tranquil_url: String, signing_key: SigningKey,}
#[derive(Clone)]pub struct App { config: Arc<RuntimeConfig>, http: reqwest::Client, auth: ServiceAuthConfig<Arc<Directory>>, gate: Arc<Gate>,}
impl App { /// Validate configuration and prepare shared clients, auth policy, and quota lock. pub fn new(config: &config::Config) -> anyhow::Result<Self> { 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<Verified, ServiceAuthError> { 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<String>,}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<reqwest::Error> for Error { fn from(error: reqwest::Error) -> Self { upstream(error) }}
#[derive(Default, Deserialize)]struct UpstreamError { error: Option<String>, message: Option<String>,}
/// 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<reqwest::Response, Error> { 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<Account, Error> { 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<App>, ExtractAnyServiceAuth(claims, token): ExtractAnyServiceAuth, Json(input): Json<CreateInput>,) -> Result<Json<Account>, 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<App>) -> Json<Value> { 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<Value> { 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;