diff --git a/Cargo.lock b/Cargo.lock index a6cbefe..4c82fb3 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -157,6 +157,12 @@ version = "0.22.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "72b3254f16251a8381aa12e40e3c4d2f0199f8c6508fbecb9d91f575e0fbb8c6" +[[package]] +name = "base64ct" +version = "1.8.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0e050f626429857a27ddccb31e0aca21356bfa709c04041aefddac081a8f068a" + [[package]] name = "bitflags" version = "1.3.2" @@ -272,6 +278,12 @@ version = "1.0.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b05b61dc5112cbb17e4b6cd61790d9845d13888356391624cbe7e41efeac1e75" +[[package]] +name = "const-oid" +version = "0.9.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c2459377285ad874054d797f3ccebf984978aa39129f6eafde5cdc8315b612f8" + [[package]] name = "cookie" version = "0.18.1" @@ -318,6 +330,33 @@ dependencies = [ "typenum", ] +[[package]] +name = "curve25519-dalek" +version = "4.1.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "97fb8b7c4503de7d6ae7b42ab72a5a59857b4c937ec27a3d4539dba95b5ab2be" +dependencies = [ + "cfg-if", + "cpufeatures", + "curve25519-dalek-derive", + "digest", + "fiat-crypto", + "rustc_version", + "subtle", + "zeroize", +] + +[[package]] +name = "curve25519-dalek-derive" +version = "0.1.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f46882e17999c6cc590af592290432be3bce0428cb0d5f8b6715e4dc7b383eb3" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.111", +] + [[package]] name = "deadpool" version = "0.12.3" @@ -353,6 +392,16 @@ dependencies = [ "tokio", ] +[[package]] +name = "der" +version = "0.7.10" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e7c1832837b905bbfb5101e07cc24c8deddf52f93225eee6ead5f4d63d53ddcb" +dependencies = [ + "const-oid", + "zeroize", +] + [[package]] name = "deranged" version = "0.5.5" @@ -390,6 +439,31 @@ version = "0.15.7" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1aaf95b3e5c8f23aa320147307562d361db0ae0d51242340f558153b4eb2439b" +[[package]] +name = "ed25519" +version = "2.2.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "115531babc129696a58c64a4fef0a8bf9e9698629fb97e9e40767d235cfbcd53" +dependencies = [ + "pkcs8", + "serde", + "signature", +] + +[[package]] +name = "ed25519-dalek" +version = "2.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "70e796c081cee67dc755e1a36a0a172b897fab85fc3f6bc48307991f64e4eca9" +dependencies = [ + "curve25519-dalek", + "ed25519", + "serde", + "sha2", + "subtle", + "zeroize", +] + [[package]] name = "encoding_rs" version = "0.8.35" @@ -427,6 +501,12 @@ version = "2.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "37909eebbb50d72f9059c3b6d82c0463f2ff062c9e95845c43a6c9c0355411be" +[[package]] +name = "fiat-crypto" +version = "0.2.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "28dea519a9695b9977216879a3ebfddf92f1c08c05d984f8996aecd6ecdc811d" + [[package]] name = "find-msvc-tools" version = "0.1.6" @@ -1091,15 +1171,19 @@ version = "0.1.0" dependencies = [ "async-trait", "axum", + "base64 0.22.1", "chrono", "deadpool-postgres", "dotenvy", + "ed25519-dalek", + "getrandom 0.3.4", "malfestio-core", "readability", "regex", "reqwest 0.12.28", "serde", "serde_json", + "sha2", "tokio", "tokio-postgres", "tower", @@ -1107,6 +1191,7 @@ dependencies = [ "tower-http", "tracing", "tracing-subscriber", + "urlencoding", "uuid", ] @@ -1414,6 +1499,16 @@ version = "0.1.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8b870d8c151b6f2fb93e84a13146138f05d02ed11c7e7c54f8826aaaf7c9f184" +[[package]] +name = "pkcs8" +version = "0.10.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f950b2377845cebe5cf8b5165cb3cc1a5e0fa5cfa3e1f7f55707d8fd82e0a7b7" +dependencies = [ + "der", + "spki", +] + [[package]] name = "pkg-config" version = "0.3.32" @@ -1721,6 +1816,15 @@ dependencies = [ "windows-sys 0.52.0", ] +[[package]] +name = "rustc_version" +version = "0.4.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cfcb3a22ef46e85b45de6ee7e79d063319ebb6594faafcf1c225ea92ab6e9b92" +dependencies = [ + "semver", +] + [[package]] name = "rustix" version = "1.1.3" @@ -1826,6 +1930,12 @@ dependencies = [ "libc", ] +[[package]] +name = "semver" +version = "1.0.27" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d767eb0aabc880b29956c35734170f26ed551a859dbd361d140cdbeca61ab1e2" + [[package]] name = "serde" version = "1.0.228" @@ -1928,6 +2038,15 @@ dependencies = [ "libc", ] +[[package]] +name = "signature" +version = "2.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "77549399552de45a898a580c1b41d445bf730df867cc44e6c0233bbc4b8329de" +dependencies = [ + "rand_core 0.6.4", +] + [[package]] name = "siphasher" version = "0.3.11" @@ -1972,6 +2091,16 @@ dependencies = [ "windows-sys 0.60.2", ] +[[package]] +name = "spki" +version = "0.7.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d91ed6c858b01f942cd56b37a94b3e0a1798290327d1236e4d9cf4eaca44d29d" +dependencies = [ + "base64ct", + "der", +] + [[package]] name = "stable_deref_trait" version = "1.2.1" @@ -2494,6 +2623,12 @@ dependencies = [ "serde", ] +[[package]] +name = "urlencoding" +version = "2.1.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "daf8dba3b7eb870caf1ddeed7bc9d2a049f3cfdfae7cb521b087cc33ae4c49da" + [[package]] name = "utf-8" version = "0.7.6" diff --git a/crates/core/src/at_uri.rs b/crates/core/src/at_uri.rs new file mode 100644 index 0000000..70cf8cf --- /dev/null +++ b/crates/core/src/at_uri.rs @@ -0,0 +1,208 @@ +//! AT-URI builder and parser for AT Protocol. +//! +//! AT-URIs are the canonical way to reference records in the AT Protocol. +//! Format: `at:////` +//! +//! - authority: DID or handle +//! - collection: NSID (e.g., "app.malfestio.deck") +//! - rkey: Record key (usually a TID) + +use std::fmt; + +/// An AT-URI representing a record in the AT Protocol network. +#[derive(Debug, Clone, PartialEq, Eq, Hash)] +pub struct AtUri { + /// The authority (DID or handle) + pub authority: String, + /// The collection NSID (e.g., "app.malfestio.deck") + pub collection: String, + /// The record key + pub rkey: String, +} + +impl AtUri { + /// Create a new AT-URI. + /// + /// # Arguments + /// + /// * `authority` - The DID or handle + /// * `collection` - The collection NSID + /// * `rkey` - The record key + pub fn new(authority: impl Into, collection: impl Into, rkey: impl Into) -> Self { + Self { authority: authority.into(), collection: collection.into(), rkey: rkey.into() } + } + + /// Create an AT-URI for a deck record. + pub fn deck(did: &str, rkey: &str) -> Self { + Self::new(did, "app.malfestio.deck", rkey) + } + + /// Create an AT-URI for a card record. + pub fn card(did: &str, rkey: &str) -> Self { + Self::new(did, "app.malfestio.card", rkey) + } + + /// Create an AT-URI for a note record. + pub fn note(did: &str, rkey: &str) -> Self { + Self::new(did, "app.malfestio.note", rkey) + } + + /// Parse an AT-URI string. + pub fn parse(s: &str) -> Result { + let s = s.strip_prefix("at://").ok_or(AtUriError::MissingScheme)?; + + let parts: Vec<&str> = s.splitn(3, '/').collect(); + if parts.len() != 3 { + return Err(AtUriError::InvalidFormat); + } + + let authority = parts[0]; + let collection = parts[1]; + let rkey = parts[2]; + + if authority.is_empty() { + return Err(AtUriError::EmptyAuthority); + } + if collection.is_empty() { + return Err(AtUriError::EmptyCollection); + } + if rkey.is_empty() { + return Err(AtUriError::EmptyRkey); + } + + if !collection.contains('.') { + return Err(AtUriError::InvalidNsid); + } + + Ok(Self { authority: authority.to_string(), collection: collection.to_string(), rkey: rkey.to_string() }) + } + + /// Check if the authority is a DID. + pub fn is_did(&self) -> bool { + self.authority.starts_with("did:") + } + + /// Check if the authority is a handle. + pub fn is_handle(&self) -> bool { + !self.is_did() + } +} + +impl fmt::Display for AtUri { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + write!(f, "at://{}/{}/{}", self.authority, self.collection, self.rkey) + } +} + +/// Error type for AT-URI parsing. +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum AtUriError { + MissingScheme, + InvalidFormat, + EmptyAuthority, + EmptyCollection, + EmptyRkey, + InvalidNsid, +} + +impl fmt::Display for AtUriError { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + AtUriError::MissingScheme => write!(f, "AT-URI must start with 'at://'"), + AtUriError::InvalidFormat => write!(f, "AT-URI must have format at://authority/collection/rkey"), + AtUriError::EmptyAuthority => write!(f, "AT-URI authority cannot be empty"), + AtUriError::EmptyCollection => write!(f, "AT-URI collection cannot be empty"), + AtUriError::EmptyRkey => write!(f, "AT-URI rkey cannot be empty"), + AtUriError::InvalidNsid => write!(f, "Collection must be a valid NSID"), + } + } +} + +impl std::error::Error for AtUriError {} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_new_at_uri() { + let uri = AtUri::new("did:plc:abc123", "app.malfestio.deck", "3k5abc123"); + assert_eq!(uri.authority, "did:plc:abc123"); + assert_eq!(uri.collection, "app.malfestio.deck"); + assert_eq!(uri.rkey, "3k5abc123"); + } + + #[test] + fn test_display() { + let uri = AtUri::new("did:plc:abc123", "app.malfestio.deck", "3k5abc123"); + assert_eq!(uri.to_string(), "at://did:plc:abc123/app.malfestio.deck/3k5abc123"); + } + + #[test] + fn test_parse_valid() { + let uri = AtUri::parse("at://did:plc:abc123/app.malfestio.deck/3k5abc123").unwrap(); + assert_eq!(uri.authority, "did:plc:abc123"); + assert_eq!(uri.collection, "app.malfestio.deck"); + assert_eq!(uri.rkey, "3k5abc123"); + } + + #[test] + fn test_parse_with_handle() { + let uri = AtUri::parse("at://alice.bsky.social/app.malfestio.note/abc123").unwrap(); + assert_eq!(uri.authority, "alice.bsky.social"); + assert!(uri.is_handle()); + assert!(!uri.is_did()); + } + + #[test] + fn test_parse_missing_scheme() { + let result = AtUri::parse("did:plc:abc123/app.malfestio.deck/3k5abc123"); + assert_eq!(result, Err(AtUriError::MissingScheme)); + } + + #[test] + fn test_parse_invalid_format() { + let result = AtUri::parse("at://did:plc:abc123/app.malfestio.deck"); + assert_eq!(result, Err(AtUriError::InvalidFormat)); + } + + #[test] + fn test_parse_empty_authority() { + let result = AtUri::parse("at:///app.malfestio.deck/rkey"); + assert_eq!(result, Err(AtUriError::EmptyAuthority)); + } + + #[test] + fn test_parse_invalid_nsid() { + let result = AtUri::parse("at://did:plc:abc123/notansid/rkey"); + assert_eq!(result, Err(AtUriError::InvalidNsid)); + } + + #[test] + fn test_roundtrip() { + let original = "at://did:plc:abc123/app.malfestio.deck/3k5abc123"; + let uri = AtUri::parse(original).unwrap(); + assert_eq!(uri.to_string(), original); + } + + #[test] + fn test_convenience_constructors() { + let deck = AtUri::deck("did:plc:abc", "tid123"); + assert_eq!(deck.collection, "app.malfestio.deck"); + + let card = AtUri::card("did:plc:abc", "tid456"); + assert_eq!(card.collection, "app.malfestio.card"); + + let note = AtUri::note("did:plc:abc", "tid789"); + assert_eq!(note.collection, "app.malfestio.note"); + } + + #[test] + fn test_is_did() { + let uri = AtUri::new("did:plc:abc123", "app.test", "rkey"); + assert!(uri.is_did()); + + let uri = AtUri::new("alice.bsky.social", "app.test", "rkey"); + assert!(!uri.is_did()); + } +} diff --git a/crates/core/src/lib.rs b/crates/core/src/lib.rs index bee65d5..80a2350 100644 --- a/crates/core/src/lib.rs +++ b/crates/core/src/lib.rs @@ -1,5 +1,7 @@ +pub mod at_uri; pub mod error; pub mod model; +pub mod tid; pub use error::{Error, Result}; pub use model::{Card, Deck, Note}; diff --git a/crates/core/src/tid.rs b/crates/core/src/tid.rs new file mode 100644 index 0000000..3a51b3a --- /dev/null +++ b/crates/core/src/tid.rs @@ -0,0 +1,161 @@ +//! TID (Timestamp Identifier) generation for AT Protocol. +//! +//! TIDs are used as record keys in the AT Protocol. They are 13-character +//! base32-sortable strings derived from timestamps with a clock identifier. +//! +//! Format: 13 characters encoding 64 bits: +//! - 53 bits: microseconds since Unix epoch +//! - 10 bits: clock identifier (for uniqueness within same microsecond) +//! - 1 bit: always 0 (reserved) + +use std::sync::atomic::{AtomicU64, Ordering}; +use std::time::{SystemTime, UNIX_EPOCH}; + +/// Base32 "sort" alphabet used by AT Protocol TIDs. +/// This alphabet maintains lexicographic sorting. +const BASE32_SORT: &[u8; 32] = b"234567abcdefghijklmnopqrstuvwxyz"; + +/// Atomic counter for clock identifier within same microsecond. +static CLOCK_ID: AtomicU64 = AtomicU64::new(0); +static LAST_TIMESTAMP: AtomicU64 = AtomicU64::new(0); + +/// Generate a new TID. +/// +/// TIDs are guaranteed to be: +/// - Unique within this process +/// - Lexicographically sortable by creation time +/// - Compatible with AT Protocol record key requirements +pub fn generate_tid() -> String { + let now = SystemTime::now() + .duration_since(UNIX_EPOCH) + .expect("Time went backwards") + .as_micros() as u64; + + let last = LAST_TIMESTAMP.load(Ordering::SeqCst); + let clock_id = if now == last { + CLOCK_ID.fetch_add(1, Ordering::SeqCst) & 0x3FF + } else { + LAST_TIMESTAMP.store(now, Ordering::SeqCst); + CLOCK_ID.store(1, Ordering::SeqCst); + 0 + }; + + let combined = (now << 11) | (clock_id << 1); + encode_base32_sort(combined) +} + +/// Encode a 64-bit value as a 13-character base32-sort string. +fn encode_base32_sort(mut value: u64) -> String { + let mut result = [0u8; 13]; + + for i in (0..13).rev() { + result[i] = BASE32_SORT[(value & 0x1F) as usize]; + value >>= 5; + } + + String::from_utf8(result.to_vec()).expect("Base32 encoding produced invalid UTF-8") +} + +/// Parse a TID string and extract the timestamp. +/// +/// Returns the Unix timestamp in microseconds, or None if invalid. +pub fn parse_tid_timestamp(tid: &str) -> Option { + if tid.len() != 13 { + return None; + } + + let decoded = decode_base32_sort(tid)?; + Some(decoded >> 11) +} + +/// Decode a base32-sort string to a 64-bit value. +fn decode_base32_sort(s: &str) -> Option { + let mut value: u64 = 0; + + for c in s.chars() { + let idx = BASE32_SORT.iter().position(|&b| b == c as u8)?; + value = (value << 5) | (idx as u64); + } + + Some(value) +} + +/// Validate that a string is a valid TID format. +pub fn is_valid_tid(tid: &str) -> bool { + if tid.len() != 13 { + return false; + } + + tid.chars().all(|c| BASE32_SORT.contains(&(c as u8))) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_tid_length() { + let tid = generate_tid(); + assert_eq!(tid.len(), 13); + } + + #[test] + fn test_tid_characters() { + let tid = generate_tid(); + for c in tid.chars() { + assert!(BASE32_SORT.contains(&(c as u8)), "Invalid character '{}' in TID", c); + } + } + + #[test] + fn test_tid_uniqueness() { + let tids: Vec = (0..100).map(|_| generate_tid()).collect(); + let mut unique = tids.clone(); + unique.sort(); + unique.dedup(); + assert_eq!(tids.len(), unique.len(), "TIDs should be unique"); + } + + #[test] + fn test_tid_sortability() { + let tid1 = generate_tid(); + std::thread::sleep(std::time::Duration::from_micros(10)); + let tid2 = generate_tid(); + + assert!(tid1 < tid2, "Later TIDs should sort after earlier ones"); + } + + #[test] + fn test_parse_tid_timestamp() { + let tid = generate_tid(); + let timestamp = parse_tid_timestamp(&tid); + assert!(timestamp.is_some()); + + let now = SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_micros() as u64; + let parsed = timestamp.unwrap(); + assert!( + now.abs_diff(parsed) < 1_000_000, + "Parsed timestamp {} too far from now {}", + parsed, + now + ); + } + + #[test] + fn test_is_valid_tid() { + let tid = generate_tid(); + assert!(is_valid_tid(&tid)); + + assert!(!is_valid_tid("short")); + assert!(!is_valid_tid("toolongstring!")); + assert!(!is_valid_tid("0123456789012")); + } + + #[test] + fn test_roundtrip_encoding() { + let value: u64 = 0x123456789ABCDEF0; + let encoded = encode_base32_sort(value); + let decoded = decode_base32_sort(&encoded); + assert_eq!(decoded, Some(value)); + } +} diff --git a/crates/server/Cargo.toml b/crates/server/Cargo.toml index bc7700e..c87dca1 100644 --- a/crates/server/Cargo.toml +++ b/crates/server/Cargo.toml @@ -6,17 +6,26 @@ edition = "2024" [dependencies] async-trait = "0.1.83" axum = "0.8.8" +base64 = "0.22" chrono = { version = "0.4.42", features = ["serde"] } deadpool-postgres = "0.14.0" dotenvy = "0.15.7" +ed25519-dalek = { version = "2.2.0", features = ["serde"] } +getrandom = { version = "0.3", features = ["std"] } malfestio-core = { version = "0.1.0", path = "../core" } readability = "0.3.0" regex = "1.12.2" reqwest = { version = "0.12.28", features = ["json"] } serde = "1.0.228" serde_json = "1.0.148" +sha2 = "0.10" tokio = { version = "1.48.0", features = ["full"] } -tokio-postgres = { version = "0.7.13", features = ["with-serde_json-1", "with-chrono-0_4", "with-uuid-1"] } +urlencoding = "2.1" +tokio-postgres = { version = "0.7.13", features = [ + "with-serde_json-1", + "with-chrono-0_4", + "with-uuid-1", +] } tower = "0.5.2" tower-cookies = "0.11.0" tower-http = { version = "0.6.8", features = ["cors", "trace"] } diff --git a/crates/server/src/lib.rs b/crates/server/src/lib.rs index fb70ae7..2c0f40a 100644 --- a/crates/server/src/lib.rs +++ b/crates/server/src/lib.rs @@ -1,6 +1,8 @@ pub mod api; pub mod db; pub mod middleware; +pub mod oauth; +pub mod pds; pub mod repository; pub mod state; diff --git a/crates/server/src/oauth/client_metadata.rs b/crates/server/src/oauth/client_metadata.rs new file mode 100644 index 0000000..d268fec --- /dev/null +++ b/crates/server/src/oauth/client_metadata.rs @@ -0,0 +1,77 @@ +//! OAuth client metadata endpoint. +//! +//! Serves the client_metadata.json for AT Protocol OAuth discovery. + +use axum::{Json, response::IntoResponse}; +use serde::Serialize; + +/// OAuth client metadata for AT Protocol. +#[derive(Serialize, Clone)] +pub struct ClientMetadata { + pub client_id: String, + pub application_type: String, + pub grant_types: Vec, + pub scope: String, + pub response_types: Vec, + pub redirect_uris: Vec, + pub client_name: String, + pub client_uri: String, + pub token_endpoint_auth_method: String, + pub dpop_bound_access_tokens: bool, +} + +impl Default for ClientMetadata { + fn default() -> Self { + Self::from_env() + } +} + +impl ClientMetadata { + /// Create client metadata from environment variables. + pub fn from_env() -> Self { + let app_url = std::env::var("APP_URL").unwrap_or_else(|_| "http://localhost:3000".to_string()); + let app_name = std::env::var("APP_NAME").unwrap_or_else(|_| "Malfestio".to_string()); + + Self { + client_id: format!("{}/oauth/client-metadata.json", app_url), + application_type: "web".to_string(), + grant_types: vec!["authorization_code".to_string(), "refresh_token".to_string()], + scope: "atproto transition:generic".to_string(), + response_types: vec!["code".to_string()], + redirect_uris: vec![format!("{}/oauth/callback", app_url)], + client_name: app_name, + client_uri: app_url, + token_endpoint_auth_method: "none".to_string(), + dpop_bound_access_tokens: true, + } + } +} + +/// Handler for `/.well-known/oauth-client-metadata` endpoint. +pub async fn client_metadata_handler() -> impl IntoResponse { + Json(ClientMetadata::from_env()) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_default_metadata() { + let meta = ClientMetadata::default(); + assert!(meta.client_id.contains("client-metadata.json")); + assert_eq!(meta.application_type, "web"); + assert!(meta.grant_types.contains(&"authorization_code".to_string())); + assert!(meta.dpop_bound_access_tokens); + } + + #[test] + fn test_metadata_serialization() { + let meta = ClientMetadata::default(); + let json = serde_json::to_string(&meta).unwrap(); + + assert!(json.contains("client_id")); + assert!(json.contains("dpop_bound_access_tokens")); + assert!(json.contains("atproto")); + } +} diff --git a/crates/server/src/oauth/dpop.rs b/crates/server/src/oauth/dpop.rs new file mode 100644 index 0000000..8b8f90c --- /dev/null +++ b/crates/server/src/oauth/dpop.rs @@ -0,0 +1,186 @@ +//! DPoP (Demonstrating Proof of Possession) implementation for OAuth 2.1. +//! +//! AT Protocol requires DPoP tokens to bind access tokens to specific clients. + +use base64::{Engine, engine::general_purpose::URL_SAFE_NO_PAD}; +use ed25519_dalek::{Signer, SigningKey, VerifyingKey}; +use serde::{Deserialize, Serialize}; +use sha2::{Digest, Sha256}; +use std::time::{SystemTime, UNIX_EPOCH}; + +/// A DPoP keypair for proof generation using Ed25519. +#[derive(Clone)] +pub struct DpopKeypair { + signing_key: SigningKey, +} + +/// DPoP proof JWT header. +#[derive(Serialize, Deserialize)] +struct DpopHeader { + typ: String, + alg: String, + jwk: DpopJwk, +} + +/// JWK representation for DPoP (Ed25519 public key). +#[derive(Serialize, Deserialize, Clone)] +pub struct DpopJwk { + kty: String, + crv: String, + x: String, +} + +/// DPoP proof JWT payload. +#[derive(Serialize, Deserialize)] +struct DpopPayload { + jti: String, + htm: String, + htu: String, + iat: u64, + #[serde(skip_serializing_if = "Option::is_none")] + ath: Option, +} + +impl DpopKeypair { + /// Generate a new random Ed25519 DPoP keypair. + pub fn generate() -> Self { + let mut rng_bytes = [0u8; 32]; + getrandom::fill(&mut rng_bytes).expect("Failed to generate random bytes"); + let signing_key = SigningKey::from_bytes(&rng_bytes); + Self { signing_key } + } + + /// Get the verifying (public) key. + pub fn verifying_key(&self) -> VerifyingKey { + self.signing_key.verifying_key() + } + + /// Get the JWK representation of the public key. + pub fn public_jwk(&self) -> DpopJwk { + let public_bytes = self.verifying_key().to_bytes(); + DpopJwk { kty: "OKP".to_string(), crv: "Ed25519".to_string(), x: URL_SAFE_NO_PAD.encode(public_bytes) } + } + + /// Generate a DPoP proof for a request. + pub fn generate_proof(&self, method: &str, url: &str, access_token: Option<&str>) -> String { + let header = DpopHeader { typ: "dpop+jwt".to_string(), alg: "EdDSA".to_string(), jwk: self.public_jwk() }; + + let now = SystemTime::now() + .duration_since(UNIX_EPOCH) + .expect("Time went backwards") + .as_secs(); + + let jti = generate_jti(); + + let ath = access_token.map(|token| { + let hash = Sha256::digest(token.as_bytes()); + URL_SAFE_NO_PAD.encode(hash) + }); + + let payload = DpopPayload { jti, htm: method.to_uppercase(), htu: url.to_string(), iat: now, ath }; + + let header_b64 = URL_SAFE_NO_PAD.encode(serde_json::to_string(&header).unwrap()); + let payload_b64 = URL_SAFE_NO_PAD.encode(serde_json::to_string(&payload).unwrap()); + + let signing_input = format!("{}.{}", header_b64, payload_b64); + + let signature = self.signing_key.sign(signing_input.as_bytes()); + let signature_b64 = URL_SAFE_NO_PAD.encode(signature.to_bytes()); + + format!("{}.{}.{}", header_b64, payload_b64, signature_b64) + } +} + +/// Generate a unique JWT ID. +fn generate_jti() -> String { + let mut bytes = [0u8; 16]; + getrandom::fill(&mut bytes).expect("Failed to generate random bytes"); + URL_SAFE_NO_PAD.encode(bytes) +} + +/// Compute the JWK thumbprint for key binding. +pub fn jwk_thumbprint(jwk: &DpopJwk) -> String { + let canonical = format!(r#"{{"crv":"{}","kty":"{}","x":"{}"}}"#, jwk.crv, jwk.kty, jwk.x); + let hash = Sha256::digest(canonical.as_bytes()); + URL_SAFE_NO_PAD.encode(hash) +} + +#[cfg(test)] +mod tests { + use super::*; + use ed25519_dalek::Verifier; + + #[test] + fn test_generate_keypair() { + let kp = DpopKeypair::generate(); + let _ = kp.verifying_key(); + } + + #[test] + fn test_keypair_uniqueness() { + let kp1 = DpopKeypair::generate(); + let kp2 = DpopKeypair::generate(); + assert_ne!(kp1.verifying_key().to_bytes(), kp2.verifying_key().to_bytes()); + } + + #[test] + fn test_public_jwk() { + let kp = DpopKeypair::generate(); + let jwk = kp.public_jwk(); + + assert_eq!(jwk.kty, "OKP"); + assert_eq!(jwk.crv, "Ed25519"); + assert!(!jwk.x.is_empty()); + assert_eq!(jwk.x.len(), 43); + } + + #[test] + fn test_generate_proof() { + let kp = DpopKeypair::generate(); + let proof = kp.generate_proof("POST", "https://example.com/token", None); + + let parts: Vec<&str> = proof.split('.').collect(); + assert_eq!(parts.len(), 3); + + let header_json = URL_SAFE_NO_PAD.decode(parts[0]).unwrap(); + let header: serde_json::Value = serde_json::from_slice(&header_json).unwrap(); + assert_eq!(header["typ"], "dpop+jwt"); + assert_eq!(header["alg"], "EdDSA"); + } + + #[test] + fn test_proof_signature_verifies() { + let kp = DpopKeypair::generate(); + let proof = kp.generate_proof("GET", "https://example.com/resource", None); + + let parts: Vec<&str> = proof.split('.').collect(); + let signing_input = format!("{}.{}", parts[0], parts[1]); + let signature_bytes = URL_SAFE_NO_PAD.decode(parts[2]).unwrap(); + + let signature = ed25519_dalek::Signature::from_slice(&signature_bytes).unwrap(); + let result = kp.verifying_key().verify(signing_input.as_bytes(), &signature); + + assert!(result.is_ok(), "Signature should verify"); + } + + #[test] + fn test_generate_proof_with_token() { + let kp = DpopKeypair::generate(); + let proof = kp.generate_proof("GET", "https://example.com/resource", Some("access_token_123")); + + let parts: Vec<&str> = proof.split('.').collect(); + let payload_json = URL_SAFE_NO_PAD.decode(parts[1]).unwrap(); + let payload: serde_json::Value = serde_json::from_slice(&payload_json).unwrap(); + + assert!(payload.get("ath").is_some()); + } + + #[test] + fn test_jwk_thumbprint() { + let kp = DpopKeypair::generate(); + let jwk = kp.public_jwk(); + let thumbprint = jwk_thumbprint(&jwk); + + assert_eq!(thumbprint.len(), 43); + } +} diff --git a/crates/server/src/oauth/flow.rs b/crates/server/src/oauth/flow.rs new file mode 100644 index 0000000..a4e669f --- /dev/null +++ b/crates/server/src/oauth/flow.rs @@ -0,0 +1,336 @@ +//! OAuth 2.1 authorization flow for AT Protocol. +//! +//! Handles the complete OAuth flow including: +//! - Authorization URL generation +//! - Token exchange with PKCE + DPoP +//! - Token refresh + +use super::dpop::DpopKeypair; +use super::pkce::{derive_code_challenge, generate_code_verifier}; +use super::resolver::{IdentityResolver, ResolveError}; +use serde::{Deserialize, Serialize}; +use std::collections::HashMap; +use std::sync::{Arc, RwLock}; + +/// OAuth session state stored during the authorization flow. +#[derive(Clone)] +pub struct OAuthSession { + /// The PKCE code verifier + pub code_verifier: String, + /// The DPoP keypair for this session + pub dpop_keypair: DpopKeypair, + /// The user's DID after resolution + pub did: Option, + /// The user's PDS URL + pub pds_url: Option, + /// When this session was created (for expiry) + pub created_at: std::time::Instant, +} + +/// OAuth tokens received from the authorization server. +#[derive(Clone, Serialize, Deserialize)] +pub struct OAuthTokens { + pub access_token: String, + pub refresh_token: Option, + pub token_type: String, + pub expires_in: Option, + pub scope: Option, +} + +/// In-memory session storage (for development). +/// In production, use a database-backed implementation. +pub type SessionStore = Arc>>; + +/// Create a new session store. +pub fn new_session_store() -> SessionStore { + Arc::new(RwLock::new(HashMap::new())) +} + +/// OAuth flow manager. +pub struct OAuthFlow { + resolver: IdentityResolver, + client: reqwest::Client, + client_id: String, + redirect_uri: String, +} + +impl OAuthFlow { + /// Create a new OAuth flow manager. + pub fn new() -> Self { + let app_url = std::env::var("APP_URL").unwrap_or_else(|_| "http://localhost:3000".to_string()); + + Self { + resolver: IdentityResolver::new(), + client: reqwest::Client::new(), + client_id: format!("{}/oauth/client-metadata.json", app_url), + redirect_uri: format!("{}/oauth/callback", app_url), + } + } + + /// Start the OAuth flow for a user handle or DID. + /// + /// Returns the authorization URL to redirect the user to. + pub async fn start_authorization( + &self, handle_or_did: &str, state: &str, sessions: &SessionStore, + ) -> Result { + let (did, pds_url) = if handle_or_did.starts_with("did:") { + let resolved = self.resolver.resolve_did(handle_or_did).await?; + (resolved.did, resolved.pds_url) + } else { + let did = self.resolver.resolve_handle(handle_or_did).await?; + let resolved = self.resolver.resolve_did(&did).await?; + (resolved.did, resolved.pds_url) + }; + + let auth_server = self.get_auth_server_metadata(&pds_url).await?; + + let code_verifier = generate_code_verifier(); + let code_challenge = derive_code_challenge(&code_verifier); + + let dpop_keypair = DpopKeypair::generate(); + + let session = OAuthSession { + code_verifier, + dpop_keypair, + did: Some(did.clone()), + pds_url: Some(pds_url), + created_at: std::time::Instant::now(), + }; + + sessions.write().unwrap().insert(state.to_string(), session); + + let auth_url = format!( + "{}?response_type=code&client_id={}&redirect_uri={}&scope={}&state={}&code_challenge={}&code_challenge_method=S256&login_hint={}", + auth_server.authorization_endpoint, + urlencoding::encode(&self.client_id), + urlencoding::encode(&self.redirect_uri), + urlencoding::encode("atproto transition:generic"), + urlencoding::encode(state), + urlencoding::encode(&code_challenge), + urlencoding::encode(&did) + ); + + Ok(auth_url) + } + + /// Exchange an authorization code for tokens. + pub async fn exchange_code( + &self, code: &str, state: &str, sessions: &SessionStore, + ) -> Result { + let session = sessions + .read() + .unwrap() + .get(state) + .cloned() + .ok_or(OAuthFlowError::SessionNotFound)?; + + let pds_url = session.pds_url.as_ref().ok_or(OAuthFlowError::SessionNotFound)?; + + let auth_server = self.get_auth_server_metadata(pds_url).await?; + + let dpop_proof = session + .dpop_keypair + .generate_proof("POST", &auth_server.token_endpoint, None); + + let response = self + .client + .post(&auth_server.token_endpoint) + .header("DPoP", dpop_proof) + .form(&[ + ("grant_type", "authorization_code"), + ("code", code), + ("redirect_uri", &self.redirect_uri), + ("client_id", &self.client_id), + ("code_verifier", &session.code_verifier), + ]) + .send() + .await + .map_err(|e| OAuthFlowError::NetworkError(e.to_string()))?; + + if !response.status().is_success() { + let error_body = response.text().await.unwrap_or_default(); + return Err(OAuthFlowError::TokenExchangeFailed(error_body)); + } + + let tokens: OAuthTokens = response + .json() + .await + .map_err(|e| OAuthFlowError::NetworkError(e.to_string()))?; + + sessions.write().unwrap().remove(state); + + Ok(tokens) + } + + /// Refresh an access token. + pub async fn refresh_token( + &self, refresh_token: &str, pds_url: &str, dpop_keypair: &DpopKeypair, + ) -> Result { + let auth_server = self.get_auth_server_metadata(pds_url).await?; + + let dpop_proof = dpop_keypair.generate_proof("POST", &auth_server.token_endpoint, None); + + let response = self + .client + .post(&auth_server.token_endpoint) + .header("DPoP", dpop_proof) + .form(&[ + ("grant_type", "refresh_token"), + ("refresh_token", refresh_token), + ("client_id", &self.client_id), + ]) + .send() + .await + .map_err(|e| OAuthFlowError::NetworkError(e.to_string()))?; + + if !response.status().is_success() { + let error_body = response.text().await.unwrap_or_default(); + return Err(OAuthFlowError::TokenRefreshFailed(error_body)); + } + + response + .json() + .await + .map_err(|e| OAuthFlowError::NetworkError(e.to_string())) + } + + /// Get authorization server metadata from PDS. + async fn get_auth_server_metadata(&self, pds_url: &str) -> Result { + // First get the protected resource metadata + let resource_url = format!("{}/.well-known/oauth-protected-resource", pds_url); + + let resource_response = self + .client + .get(&resource_url) + .timeout(std::time::Duration::from_secs(10)) + .send() + .await + .map_err(|e| OAuthFlowError::NetworkError(e.to_string()))?; + + if !resource_response.status().is_success() { + return Err(OAuthFlowError::MetadataFetchFailed(pds_url.to_string())); + } + + let resource: serde_json::Value = resource_response + .json() + .await + .map_err(|e| OAuthFlowError::NetworkError(e.to_string()))?; + + let auth_server_url = resource["authorization_servers"] + .as_array() + .and_then(|arr| arr.first()) + .and_then(|v| v.as_str()) + .ok_or_else(|| OAuthFlowError::MetadataFetchFailed(pds_url.to_string()))?; + + let auth_meta_url = format!("{}/.well-known/oauth-authorization-server", auth_server_url); + + let auth_response = self + .client + .get(&auth_meta_url) + .timeout(std::time::Duration::from_secs(10)) + .send() + .await + .map_err(|e| OAuthFlowError::NetworkError(e.to_string()))?; + + if !auth_response.status().is_success() { + return Err(OAuthFlowError::MetadataFetchFailed(auth_server_url.to_string())); + } + + auth_response + .json() + .await + .map_err(|e| OAuthFlowError::NetworkError(e.to_string())) + } +} + +impl Default for OAuthFlow { + fn default() -> Self { + Self::new() + } +} + +/// Authorization server metadata. +#[derive(Deserialize)] +pub struct AuthServerMetadata { + pub issuer: String, + pub authorization_endpoint: String, + pub token_endpoint: String, + pub pushed_authorization_request_endpoint: Option, +} + +/// Error type for OAuth flow operations. +#[derive(Debug, Clone)] +pub enum OAuthFlowError { + SessionNotFound, + NetworkError(String), + MetadataFetchFailed(String), + TokenExchangeFailed(String), + TokenRefreshFailed(String), + ResolveError(String), +} + +impl From for OAuthFlowError { + fn from(err: ResolveError) -> Self { + OAuthFlowError::ResolveError(err.to_string()) + } +} + +impl std::fmt::Display for OAuthFlowError { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + OAuthFlowError::SessionNotFound => write!(f, "OAuth session not found"), + OAuthFlowError::NetworkError(e) => write!(f, "Network error: {}", e), + OAuthFlowError::MetadataFetchFailed(url) => write!(f, "Failed to fetch metadata from {}", url), + OAuthFlowError::TokenExchangeFailed(e) => write!(f, "Token exchange failed: {}", e), + OAuthFlowError::TokenRefreshFailed(e) => write!(f, "Token refresh failed: {}", e), + OAuthFlowError::ResolveError(e) => write!(f, "Identity resolution failed: {}", e), + } + } +} + +impl std::error::Error for OAuthFlowError {} + +/// Generate a secure random state parameter. +pub fn generate_state() -> String { + use base64::{Engine, engine::general_purpose::URL_SAFE_NO_PAD}; + + let mut bytes = [0u8; 16]; + getrandom::fill(&mut bytes).expect("Failed to generate random bytes"); + URL_SAFE_NO_PAD.encode(bytes) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_generate_state() { + let state1 = generate_state(); + let state2 = generate_state(); + + assert_ne!(state1, state2); + assert_eq!(state1.len(), 22); + } + + #[test] + fn test_new_session_store() { + let store = new_session_store(); + assert!(store.read().unwrap().is_empty()); + } + + #[test] + fn test_oauth_flow_creation() { + let flow = OAuthFlow::new(); + assert!(flow.client_id.contains("client-metadata.json")); + assert!(flow.redirect_uri.contains("callback")); + } + + #[test] + fn test_oauth_flow_error_display() { + let err = OAuthFlowError::SessionNotFound; + assert!(err.to_string().contains("session not found")); + + let err = OAuthFlowError::NetworkError("timeout".to_string()); + assert!(err.to_string().contains("timeout")); + } +} diff --git a/crates/server/src/oauth/mod.rs b/crates/server/src/oauth/mod.rs new file mode 100644 index 0000000..3d00588 --- /dev/null +++ b/crates/server/src/oauth/mod.rs @@ -0,0 +1,15 @@ +//! OAuth 2.1 implementation for AT Protocol. +//! +//! This module provides the OAuth 2.1 client flow components required +//! for AT Protocol authentication: +//! +//! - PKCE (Proof Key for Code Exchange) +//! - DPoP (Demonstrating Proof of Possession) +//! - Handle/DID resolution +//! - Token management + +pub mod client_metadata; +pub mod dpop; +pub mod flow; +pub mod pkce; +pub mod resolver; diff --git a/crates/server/src/oauth/pkce.rs b/crates/server/src/oauth/pkce.rs new file mode 100644 index 0000000..c212ef7 --- /dev/null +++ b/crates/server/src/oauth/pkce.rs @@ -0,0 +1,76 @@ +//! PKCE (Proof Key for Code Exchange) implementation for OAuth 2.1. +//! +//! AT Protocol requires PKCE with S256 challenge method. + +use base64::{Engine, engine::general_purpose::URL_SAFE_NO_PAD}; +use sha2::{Digest, Sha256}; + +/// Length of the code verifier in bytes (before base64 encoding). +const CODE_VERIFIER_LENGTH: usize = 32; + +/// Generate a cryptographically random code verifier. +/// +/// The verifier is a high-entropy random string used in PKCE flow. +pub fn generate_code_verifier() -> String { + let mut bytes = [0u8; CODE_VERIFIER_LENGTH]; + getrandom::fill(&mut bytes).expect("Failed to generate random bytes"); + URL_SAFE_NO_PAD.encode(bytes) +} + +/// Derive the S256 code challenge from a code verifier. +/// +/// The challenge is the base64url-encoded SHA-256 hash of the verifier. +pub fn derive_code_challenge(verifier: &str) -> String { + let hash = Sha256::digest(verifier.as_bytes()); + URL_SAFE_NO_PAD.encode(hash) +} + +/// Verify that a code challenge matches a code verifier. +pub fn verify_challenge(verifier: &str, challenge: &str) -> bool { + derive_code_challenge(verifier) == challenge +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_generate_verifier_length() { + let verifier = generate_code_verifier(); + assert_eq!(verifier.len(), 43); + } + + #[test] + fn test_generate_verifier_uniqueness() { + let v1 = generate_code_verifier(); + let v2 = generate_code_verifier(); + assert_ne!(v1, v2); + } + + #[test] + fn test_challenge_derivation() { + let verifier = "dBjftJeZ4CVP-mB92K27uhbUJU1p1r_wW1gFWFOEjXk"; + let challenge = derive_code_challenge(verifier); + + assert!(!challenge.is_empty()); + assert_eq!(challenge.len(), 43); + } + + #[test] + fn test_verify_challenge() { + let verifier = generate_code_verifier(); + let challenge = derive_code_challenge(&verifier); + + assert!(verify_challenge(&verifier, &challenge)); + assert!(!verify_challenge(&verifier, "wrong_challenge")); + } + + #[test] + fn test_challenge_is_url_safe() { + let verifier = generate_code_verifier(); + let challenge = derive_code_challenge(&verifier); + assert!(!challenge.contains('+')); + assert!(!challenge.contains('/')); + assert!(!challenge.contains('=')); + } +} diff --git a/crates/server/src/oauth/resolver.rs b/crates/server/src/oauth/resolver.rs new file mode 100644 index 0000000..6b3f49d --- /dev/null +++ b/crates/server/src/oauth/resolver.rs @@ -0,0 +1,261 @@ +//! Handle and DID resolution for AT Protocol. +//! +//! Resolves user identities to discover their PDS (Personal Data Server). + +use serde::{Deserialize, Serialize}; + +/// Result of resolving a handle or DID. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct ResolvedIdentity { + /// The DID (always populated after resolution) + pub did: String, + /// The handle (if resolved from DID) + pub handle: Option, + /// The PDS URL for this identity + pub pds_url: String, +} + +/// Error type for resolution failures. +#[derive(Debug, Clone)] +pub enum ResolveError { + /// Handle not found + HandleNotFound(String), + /// DID not found + DidNotFound(String), + /// Network error + NetworkError(String), + /// Invalid DID format + InvalidDid(String), +} + +impl std::fmt::Display for ResolveError { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + ResolveError::HandleNotFound(h) => write!(f, "Handle not found: {}", h), + ResolveError::DidNotFound(d) => write!(f, "DID not found: {}", d), + ResolveError::NetworkError(e) => write!(f, "Network error: {}", e), + ResolveError::InvalidDid(d) => write!(f, "Invalid DID: {}", d), + } + } +} + +impl std::error::Error for ResolveError {} + +/// Resolver for AT Protocol identities. +/// +/// Handles resolution of: +/// - Handle -> DID (via DNS TXT or HTTP well-known) +/// - DID -> PDS URL (via PLC directory or did:web) +pub struct IdentityResolver { + client: reqwest::Client, + plc_directory: String, +} + +impl Default for IdentityResolver { + fn default() -> Self { + Self::new() + } +} + +impl IdentityResolver { + /// Create a new resolver with default settings. + pub fn new() -> Self { + Self { client: reqwest::Client::new(), plc_directory: "https://plc.directory".to_string() } + } + + /// Create a resolver with a custom PLC directory URL. + pub fn with_plc_directory(plc_directory: &str) -> Self { + Self { client: reqwest::Client::new(), plc_directory: plc_directory.to_string() } + } + + /// Resolve a handle to a DID. + /// + /// Tries HTTP well-known first, then falls back to DNS TXT. + pub async fn resolve_handle(&self, handle: &str) -> Result { + // Try HTTP well-known first + if let Ok(did) = self.resolve_handle_http(handle).await { + return Ok(did); + } + + // Fall back to DNS TXT (simplified - just return error for now) + Err(ResolveError::HandleNotFound(handle.to_string())) + } + + /// Resolve handle via HTTP well-known. + async fn resolve_handle_http(&self, handle: &str) -> Result { + let url = format!("https://{}/.well-known/atproto-did", handle); + + let response = self + .client + .get(&url) + .timeout(std::time::Duration::from_secs(10)) + .send() + .await + .map_err(|e| ResolveError::NetworkError(e.to_string()))?; + + if !response.status().is_success() { + return Err(ResolveError::HandleNotFound(handle.to_string())); + } + + let did = response + .text() + .await + .map_err(|e| ResolveError::NetworkError(e.to_string()))? + .trim() + .to_string(); + + if !did.starts_with("did:") { + return Err(ResolveError::HandleNotFound(handle.to_string())); + } + + Ok(did) + } + + /// Resolve a DID to its PDS URL. + pub async fn resolve_did(&self, did: &str) -> Result { + if did.starts_with("did:plc:") { + self.resolve_plc_did(did).await + } else if did.starts_with("did:web:") { + self.resolve_web_did(did).await + } else { + Err(ResolveError::InvalidDid(did.to_string())) + } + } + + /// Resolve a did:plc via the PLC directory. + async fn resolve_plc_did(&self, did: &str) -> Result { + let url = format!("{}/{}", self.plc_directory, did); + + let response = self + .client + .get(&url) + .timeout(std::time::Duration::from_secs(10)) + .send() + .await + .map_err(|e| ResolveError::NetworkError(e.to_string()))?; + + if !response.status().is_success() { + return Err(ResolveError::DidNotFound(did.to_string())); + } + + let doc: serde_json::Value = response + .json() + .await + .map_err(|e| ResolveError::NetworkError(e.to_string()))?; + + // Extract PDS URL from service array + let pds_url = doc["service"] + .as_array() + .and_then(|services| { + services.iter().find(|s| { + s["id"].as_str() == Some("#atproto_pds") || s["type"].as_str() == Some("AtprotoPersonalDataServer") + }) + }) + .and_then(|s| s["serviceEndpoint"].as_str()) + .ok_or_else(|| ResolveError::DidNotFound(did.to_string()))? + .to_string(); + + // Extract handle from alsoKnownAs + let handle = doc["alsoKnownAs"] + .as_array() + .and_then(|aka| { + aka.iter() + .find(|a| a.as_str().map(|s| s.starts_with("at://")).unwrap_or(false)) + }) + .and_then(|a| a.as_str()) + .map(|s| s.strip_prefix("at://").unwrap_or(s).to_string()); + + Ok(ResolvedIdentity { did: did.to_string(), handle, pds_url }) + } + + /// Resolve a did:web via HTTP. + async fn resolve_web_did(&self, did: &str) -> Result { + // did:web:example.com -> https://example.com/.well-known/did.json + let domain = did + .strip_prefix("did:web:") + .ok_or_else(|| ResolveError::InvalidDid(did.to_string()))?; + + let url = format!("https://{}/.well-known/did.json", domain); + + let response = self + .client + .get(&url) + .timeout(std::time::Duration::from_secs(10)) + .send() + .await + .map_err(|e| ResolveError::NetworkError(e.to_string()))?; + + if !response.status().is_success() { + return Err(ResolveError::DidNotFound(did.to_string())); + } + + let doc: serde_json::Value = response + .json() + .await + .map_err(|e| ResolveError::NetworkError(e.to_string()))?; + + let pds_url = doc["service"] + .as_array() + .and_then(|services| { + services + .iter() + .find(|s| s["type"].as_str() == Some("AtprotoPersonalDataServer")) + }) + .and_then(|s| s["serviceEndpoint"].as_str()) + .ok_or_else(|| ResolveError::DidNotFound(did.to_string()))? + .to_string(); + + Ok(ResolvedIdentity { did: did.to_string(), handle: None, pds_url }) + } +} + +/// Check if a string is a valid DID. +pub fn is_valid_did(s: &str) -> bool { + s.starts_with("did:plc:") || s.starts_with("did:web:") +} + +/// Check if a string is a valid handle. +pub fn is_valid_handle(s: &str) -> bool { + // Simple validation: contains at least one dot, no spaces + s.contains('.') && !s.contains(' ') && !s.starts_with("did:") +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_is_valid_did() { + assert!(is_valid_did("did:plc:abc123")); + assert!(is_valid_did("did:web:example.com")); + assert!(!is_valid_did("alice.bsky.social")); + assert!(!is_valid_did("did:other:xyz")); + } + + #[test] + fn test_is_valid_handle() { + assert!(is_valid_handle("alice.bsky.social")); + assert!(is_valid_handle("bob.example.com")); + assert!(!is_valid_handle("did:plc:abc123")); + assert!(!is_valid_handle("invalid handle")); + assert!(!is_valid_handle("nodots")); + } + + #[test] + fn test_resolver_creation() { + let resolver = IdentityResolver::new(); + assert_eq!(resolver.plc_directory, "https://plc.directory"); + + let custom = IdentityResolver::with_plc_directory("https://custom.plc"); + assert_eq!(custom.plc_directory, "https://custom.plc"); + } + + #[test] + fn test_resolve_error_display() { + let err = ResolveError::HandleNotFound("test.handle".to_string()); + assert!(err.to_string().contains("test.handle")); + + let err = ResolveError::InvalidDid("bad:did".to_string()); + assert!(err.to_string().contains("bad:did")); + } +} diff --git a/crates/server/src/pds/client.rs b/crates/server/src/pds/client.rs new file mode 100644 index 0000000..db0f9f4 --- /dev/null +++ b/crates/server/src/pds/client.rs @@ -0,0 +1,302 @@ +//! PDS client for XRPC operations. +//! +//! Handles communication with a user's Personal Data Server. + +use crate::oauth::dpop::DpopKeypair; +use malfestio_core::at_uri::AtUri; +use serde::{Deserialize, Serialize}; + +/// A client for interacting with a user's PDS. +pub struct PdsClient { + http_client: reqwest::Client, + pds_url: String, + access_token: String, + dpop_keypair: DpopKeypair, +} + +/// Request body for putRecord XRPC. +#[derive(Serialize)] +#[serde(rename_all = "camelCase")] +pub struct PutRecordRequest { + pub repo: String, + pub collection: String, + pub rkey: String, + pub record: serde_json::Value, + #[serde(skip_serializing_if = "Option::is_none")] + pub swap_record: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub swap_commit: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub validate: Option, +} + +/// Response from putRecord XRPC. +#[derive(Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct PutRecordResponse { + pub uri: String, + pub cid: String, +} + +/// Request body for deleteRecord XRPC. +#[derive(Serialize)] +#[serde(rename_all = "camelCase")] +pub struct DeleteRecordRequest { + pub repo: String, + pub collection: String, + pub rkey: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub swap_record: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub swap_commit: Option, +} + +/// Response from uploadBlob XRPC. +#[derive(Deserialize)] +pub struct UploadBlobResponse { + pub blob: BlobRef, +} + +/// A reference to an uploaded blob. +#[derive(Clone, Serialize, Deserialize)] +pub struct BlobRef { + #[serde(rename = "$type")] + pub blob_type: String, + #[serde(rename = "ref")] + pub cid: CidLink, + #[serde(rename = "mimeType")] + pub mime_type: String, + pub size: u64, +} + +/// A CID link. +#[derive(Clone, Serialize, Deserialize)] +pub struct CidLink { + #[serde(rename = "$link")] + pub link: String, +} + +/// Error type for PDS operations. +#[derive(Debug, Clone)] +pub enum PdsError { + NetworkError(String), + AuthError(String), + ValidationError(String), + NotFound(String), + ServerError(String), +} + +impl std::fmt::Display for PdsError { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + PdsError::NetworkError(e) => write!(f, "Network error: {}", e), + PdsError::AuthError(e) => write!(f, "Authentication error: {}", e), + PdsError::ValidationError(e) => write!(f, "Validation error: {}", e), + PdsError::NotFound(e) => write!(f, "Not found: {}", e), + PdsError::ServerError(e) => write!(f, "Server error: {}", e), + } + } +} + +impl std::error::Error for PdsError {} + +impl PdsClient { + /// Create a new PDS client. + pub fn new(pds_url: String, access_token: String, dpop_keypair: DpopKeypair) -> Self { + Self { http_client: reqwest::Client::new(), pds_url, access_token, dpop_keypair } + } + + /// Create or update a record in the repository. + /// + /// # Arguments + /// + /// * `did` - The user's DID (repository owner) + /// * `collection` - The collection NSID (e.g., "app.malfestio.deck") + /// * `rkey` - The record key (TID) + /// * `record` - The record data as JSON + pub async fn put_record( + &self, did: &str, collection: &str, rkey: &str, record: serde_json::Value, + ) -> Result { + let url = format!("{}/xrpc/com.atproto.repo.putRecord", self.pds_url); + + let dpop_proof = self.dpop_keypair.generate_proof("POST", &url, Some(&self.access_token)); + + let request = PutRecordRequest { + repo: did.to_string(), + collection: collection.to_string(), + rkey: rkey.to_string(), + record, + swap_record: None, + swap_commit: None, + validate: Some(true), + }; + + let response = self + .http_client + .post(&url) + .header("Authorization", format!("DPoP {}", self.access_token)) + .header("DPoP", dpop_proof) + .json(&request) + .send() + .await + .map_err(|e| PdsError::NetworkError(e.to_string()))?; + + self.handle_response(response).await + } + + /// Delete a record from the repository. + pub async fn delete_record(&self, did: &str, collection: &str, rkey: &str) -> Result<(), PdsError> { + let url = format!("{}/xrpc/com.atproto.repo.deleteRecord", self.pds_url); + + let dpop_proof = self.dpop_keypair.generate_proof("POST", &url, Some(&self.access_token)); + + let request = DeleteRecordRequest { + repo: did.to_string(), + collection: collection.to_string(), + rkey: rkey.to_string(), + swap_record: None, + swap_commit: None, + }; + + let response = self + .http_client + .post(&url) + .header("Authorization", format!("DPoP {}", self.access_token)) + .header("DPoP", dpop_proof) + .json(&request) + .send() + .await + .map_err(|e| PdsError::NetworkError(e.to_string()))?; + + if response.status().is_success() { + Ok(()) + } else { + let status = response.status(); + let body = response.text().await.unwrap_or_default(); + Err(self.map_error_status(status, body)) + } + } + + /// Upload a blob (media attachment) to the repository. + pub async fn upload_blob(&self, data: Vec, mime_type: &str) -> Result { + let url = format!("{}/xrpc/com.atproto.repo.uploadBlob", self.pds_url); + + let dpop_proof = self.dpop_keypair.generate_proof("POST", &url, Some(&self.access_token)); + + let response = self + .http_client + .post(&url) + .header("Authorization", format!("DPoP {}", self.access_token)) + .header("DPoP", dpop_proof) + .header("Content-Type", mime_type) + .body(data) + .send() + .await + .map_err(|e| PdsError::NetworkError(e.to_string()))?; + + if !response.status().is_success() { + let status = response.status(); + let body = response.text().await.unwrap_or_default(); + return Err(self.map_error_status(status, body)); + } + + let upload_response: UploadBlobResponse = response + .json() + .await + .map_err(|e| PdsError::NetworkError(e.to_string()))?; + + Ok(upload_response.blob) + } + + /// Handle response and parse AT-URI from success. + async fn handle_response(&self, response: reqwest::Response) -> Result { + if !response.status().is_success() { + let status = response.status(); + let body = response.text().await.unwrap_or_default(); + return Err(self.map_error_status(status, body)); + } + + let put_response: PutRecordResponse = response + .json() + .await + .map_err(|e| PdsError::NetworkError(e.to_string()))?; + + AtUri::parse(&put_response.uri).map_err(|e| PdsError::ValidationError(e.to_string())) + } + + /// Map HTTP status to PdsError. + fn map_error_status(&self, status: reqwest::StatusCode, body: String) -> PdsError { + match status.as_u16() { + 401 => PdsError::AuthError(body), + 400 => PdsError::ValidationError(body), + 404 => PdsError::NotFound(body), + _ => PdsError::ServerError(format!("{}: {}", status, body)), + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_put_record_request_serialization() { + let request = PutRecordRequest { + repo: "did:plc:abc123".to_string(), + collection: "app.malfestio.deck".to_string(), + rkey: "3k5abc123".to_string(), + record: serde_json::json!({ + "title": "Test Deck", + "createdAt": "2024-01-01T00:00:00Z" + }), + swap_record: None, + swap_commit: None, + validate: Some(true), + }; + + let json = serde_json::to_string(&request).unwrap(); + assert!(json.contains("\"repo\":\"did:plc:abc123\"")); + assert!(json.contains("\"collection\":\"app.malfestio.deck\"")); + assert!(json.contains("\"rkey\":\"3k5abc123\"")); + assert!(json.contains("\"validate\":true")); + } + + #[test] + fn test_delete_record_request_serialization() { + let request = DeleteRecordRequest { + repo: "did:plc:abc123".to_string(), + collection: "app.malfestio.deck".to_string(), + rkey: "3k5abc123".to_string(), + swap_record: None, + swap_commit: None, + }; + + let json = serde_json::to_string(&request).unwrap(); + assert!(json.contains("\"repo\":\"did:plc:abc123\"")); + assert!(!json.contains("swapRecord")); // Should be omitted when None + } + + #[test] + fn test_blob_ref_serialization() { + let blob_ref = BlobRef { + blob_type: "blob".to_string(), + cid: CidLink { link: "bafyreiabc123".to_string() }, + mime_type: "image/jpeg".to_string(), + size: 12345, + }; + + let json = serde_json::to_string(&blob_ref).unwrap(); + assert!(json.contains("\"$type\":\"blob\"")); + assert!(json.contains("\"$link\":\"bafyreiabc123\"")); + assert!(json.contains("\"mimeType\":\"image/jpeg\"")); + } + + #[test] + fn test_pds_error_display() { + let err = PdsError::AuthError("Invalid token".to_string()); + assert!(err.to_string().contains("Invalid token")); + + let err = PdsError::NetworkError("Connection refused".to_string()); + assert!(err.to_string().contains("Connection refused")); + } +} diff --git a/crates/server/src/pds/mod.rs b/crates/server/src/pds/mod.rs new file mode 100644 index 0000000..67119df --- /dev/null +++ b/crates/server/src/pds/mod.rs @@ -0,0 +1,9 @@ +//! PDS (Personal Data Server) client for AT Protocol. +//! +//! Provides record publishing operations: +//! - putRecord - Create or update records +//! - deleteRecord - Remove records +//! - uploadBlob - Upload media attachments + +pub mod client; +pub mod records; diff --git a/crates/server/src/pds/records.rs b/crates/server/src/pds/records.rs new file mode 100644 index 0000000..84a04c1 --- /dev/null +++ b/crates/server/src/pds/records.rs @@ -0,0 +1,302 @@ +//! Record serialization for AT Protocol Lexicons. +//! +//! Converts internal models to AT Protocol record format. + +use chrono::Utc; +use malfestio_core::at_uri::AtUri; +use malfestio_core::model::{Card, Deck, Note, Visibility}; +use malfestio_core::tid::generate_tid; +use serde::Serialize; +use serde_json::Value; + +/// A deck record in Lexicon format. +#[derive(Serialize)] +#[serde(rename_all = "camelCase")] +pub struct DeckRecord { + #[serde(rename = "$type")] + pub record_type: String, + pub title: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub description: Option, + pub tags: Vec, + #[serde(skip_serializing_if = "Vec::is_empty")] + pub card_refs: Vec, + #[serde(skip_serializing_if = "Vec::is_empty")] + pub source_refs: Vec, + #[serde(skip_serializing_if = "Option::is_none")] + pub license: Option, + pub created_at: String, +} + +/// A card record in Lexicon format. +#[derive(Serialize)] +#[serde(rename_all = "camelCase")] +pub struct CardRecord { + #[serde(rename = "$type")] + pub record_type: String, + pub deck_ref: String, + pub front: String, + pub back: String, + pub card_type: String, + #[serde(skip_serializing_if = "Vec::is_empty")] + pub hints: Vec, + #[serde(skip_serializing_if = "Option::is_none")] + pub media: Option, + pub created_at: String, +} + +/// Media attachment for a card. +#[derive(Serialize)] +#[serde(rename_all = "camelCase")] +pub struct CardMedia { + pub image_ref: Option, + pub audio_ref: Option, +} + +/// A note record in Lexicon format. +#[derive(Serialize)] +#[serde(rename_all = "camelCase")] +pub struct NoteRecord { + #[serde(rename = "$type")] + pub record_type: String, + pub title: String, + pub body: String, + pub tags: Vec, + #[serde(skip_serializing_if = "Vec::is_empty")] + pub links: Vec, + pub visibility: String, + pub created_at: String, +} + +/// Result of preparing a record for publishing. +pub struct PreparedRecord { + /// The record key (TID) + pub rkey: String, + /// The NSID collection + pub collection: String, + /// The serialized record + pub record: Value, +} + +impl DeckRecord { + /// Create a DeckRecord from an internal Deck model. + pub fn from_deck(deck: &Deck, card_at_uris: Vec) -> Self { + Self { + record_type: "app.malfestio.deck".to_string(), + title: deck.title.clone(), + description: if deck.description.is_empty() { None } else { Some(deck.description.clone()) }, + tags: deck.tags.clone(), + card_refs: card_at_uris, + source_refs: vec![], + license: None, + created_at: Utc::now().to_rfc3339(), + } + } +} + +impl CardRecord { + /// Create a CardRecord from an internal Card model. + pub fn from_card(card: &Card, deck_at_uri: &str) -> Self { + Self { + record_type: "app.malfestio.card".to_string(), + deck_ref: deck_at_uri.to_string(), + front: card.front.clone(), + back: card.back.clone(), + card_type: "basic".to_string(), + hints: vec![], + media: card + .media_url + .as_ref() + .map(|url| CardMedia { image_ref: Some(url.clone()), audio_ref: None }), + created_at: Utc::now().to_rfc3339(), + } + } +} + +impl NoteRecord { + /// Create a NoteRecord from an internal Note model. + pub fn from_note(note: &Note) -> Self { + Self { + record_type: "app.malfestio.note".to_string(), + title: note.title.clone(), + body: note.body.clone(), + tags: note.tags.clone(), + links: note.links.clone(), + visibility: visibility_to_string(¬e.visibility), + created_at: Utc::now().to_rfc3339(), + } + } +} + +/// Convert visibility enum to string for Lexicon. +fn visibility_to_string(visibility: &Visibility) -> String { + match visibility { + Visibility::Private => "private".to_string(), + Visibility::Unlisted => "unlisted".to_string(), + Visibility::Public => "public".to_string(), + Visibility::SharedWith(_) => "shared".to_string(), + } +} + +/// Prepare a deck for publishing to PDS. +pub fn prepare_deck_record(deck: &Deck, card_at_uris: Vec) -> PreparedRecord { + let record = DeckRecord::from_deck(deck, card_at_uris); + PreparedRecord { + rkey: generate_tid(), + collection: "app.malfestio.deck".to_string(), + record: serde_json::to_value(record).expect("Failed to serialize deck record"), + } +} + +/// Prepare a card for publishing to PDS. +pub fn prepare_card_record(card: &Card, deck_at_uri: &str) -> PreparedRecord { + let record = CardRecord::from_card(card, deck_at_uri); + PreparedRecord { + rkey: generate_tid(), + collection: "app.malfestio.card".to_string(), + record: serde_json::to_value(record).expect("Failed to serialize card record"), + } +} + +/// Prepare a note for publishing to PDS. +pub fn prepare_note_record(note: &Note) -> PreparedRecord { + let record = NoteRecord::from_note(note); + PreparedRecord { + rkey: generate_tid(), + collection: "app.malfestio.note".to_string(), + record: serde_json::to_value(record).expect("Failed to serialize note record"), + } +} + +/// Generate an AT-URI for a record. +pub fn make_at_uri(did: &str, collection: &str, rkey: &str) -> AtUri { + AtUri::new(did, collection, rkey) +} + +#[cfg(test)] +mod tests { + use super::*; + + fn sample_deck() -> Deck { + Deck { + id: "deck-123".to_string(), + owner_did: "did:plc:abc123".to_string(), + title: "Test Deck".to_string(), + description: "A test deck".to_string(), + tags: vec!["test".to_string(), "sample".to_string()], + visibility: Visibility::Public, + published_at: None, + fork_of: None, + } + } + + fn sample_card() -> Card { + Card { + id: "card-123".to_string(), + owner_did: "did:plc:abc123".to_string(), + deck_id: "deck-123".to_string(), + front: "What is the capital of France?".to_string(), + back: "Paris".to_string(), + media_url: None, + } + } + + fn sample_note() -> Note { + Note { + id: "note-123".to_string(), + owner_did: "did:plc:abc123".to_string(), + title: "Test Note".to_string(), + body: "This is a test note with **markdown**.".to_string(), + tags: vec!["notes".to_string()], + visibility: Visibility::Public, + published_at: None, + links: vec![], + } + } + + #[test] + fn test_deck_record_from_deck() { + let deck = sample_deck(); + let record = DeckRecord::from_deck(&deck, vec![]); + + assert_eq!(record.record_type, "app.malfestio.deck"); + assert_eq!(record.title, "Test Deck"); + assert_eq!(record.description, Some("A test deck".to_string())); + assert_eq!(record.tags.len(), 2); + } + + #[test] + fn test_deck_record_serialization() { + let deck = sample_deck(); + let record = DeckRecord::from_deck(&deck, vec!["at://did:plc:abc/app.malfestio.card/tid1".to_string()]); + + let json = serde_json::to_string(&record).unwrap(); + assert!(json.contains("\"$type\":\"app.malfestio.deck\"")); + assert!(json.contains("\"title\":\"Test Deck\"")); + assert!(json.contains("cardRefs")); + } + + #[test] + fn test_card_record_from_card() { + let card = sample_card(); + let deck_uri = "at://did:plc:abc123/app.malfestio.deck/tid123"; + let record = CardRecord::from_card(&card, deck_uri); + + assert_eq!(record.record_type, "app.malfestio.card"); + assert_eq!(record.deck_ref, deck_uri); + assert_eq!(record.front, "What is the capital of France?"); + assert_eq!(record.back, "Paris"); + } + + #[test] + fn test_note_record_from_note() { + let note = sample_note(); + let record = NoteRecord::from_note(¬e); + + assert_eq!(record.record_type, "app.malfestio.note"); + assert_eq!(record.title, "Test Note"); + assert_eq!(record.visibility, "public"); + } + + #[test] + fn test_prepare_deck_record() { + let deck = sample_deck(); + let prepared = prepare_deck_record(&deck, vec![]); + + assert_eq!(prepared.collection, "app.malfestio.deck"); + assert_eq!(prepared.rkey.len(), 13); // TID length + assert!(prepared.record.is_object()); + } + + #[test] + fn test_prepare_card_record() { + let card = sample_card(); + let prepared = prepare_card_record(&card, "at://did:plc:abc/app.malfestio.deck/tid"); + + assert_eq!(prepared.collection, "app.malfestio.card"); + assert_eq!(prepared.rkey.len(), 13); + } + + #[test] + fn test_prepare_note_record() { + let note = sample_note(); + let prepared = prepare_note_record(¬e); + + assert_eq!(prepared.collection, "app.malfestio.note"); + assert_eq!(prepared.rkey.len(), 13); + } + + #[test] + fn test_make_at_uri() { + let uri = make_at_uri("did:plc:abc123", "app.malfestio.deck", "3k5abc123"); + assert_eq!(uri.to_string(), "at://did:plc:abc123/app.malfestio.deck/3k5abc123"); + } + + #[test] + fn test_visibility_to_string() { + assert_eq!(visibility_to_string(&Visibility::Private), "private"); + assert_eq!(visibility_to_string(&Visibility::Public), "public"); + assert_eq!(visibility_to_string(&Visibility::Unlisted), "unlisted"); + assert_eq!(visibility_to_string(&Visibility::SharedWith(vec![])), "shared"); + } +} diff --git a/docs/at-notes.md b/docs/at-notes.md new file mode 100644 index 0000000..c25d58d --- /dev/null +++ b/docs/at-notes.md @@ -0,0 +1,103 @@ +# AT Protocol Research Notes + +## OAuth 2.1 Specification + +AT Protocol uses a specific profile of OAuth 2.1 for client↔PDS authorization. + +### Required Components + +- **Client Metadata Endpoint**: Serve `client_metadata.json` at a public HTTPS URL (this URL becomes the `client_id`) + + ```json + { + "client_id": "https://your-app.com/oauth/client-metadata.json", + "application_type": "web", + "grant_types": ["authorization_code", "refresh_token"], + "scope": "atproto transition:generic", + "response_types": ["code"], + "redirect_uris": ["https://your-app.com/oauth/callback"], + "client_name": "Malfestio", + "client_uri": "https://your-app.com" + } + ``` + +- **PKCE (Mandatory)**: Generate `code_verifier` and `code_challenge` (S256 only) +- **DPoP (Mandatory)**: Bind tokens to client instances with proof-of-possession JWTs +- **Handle/DID Resolution**: Resolve user identity to discover their PDS +- **Token Exchange**: Authorization code flow with token refresh + +## Record Publishing + +### XRPC Endpoints + +- `com.atproto.repo.putRecord` — Create or update records +- `com.atproto.repo.deleteRecord` — Remove records +- `com.atproto.repo.uploadBlob` — Upload media attachments + +### Record Keys + +Use TID (timestamp-based identifiers) per Lexicon spec. + +### AT-URIs + +Format: `at:////` + +Example: `at://did:plc:abc123/app.malfestio.deck/3k5abc123` + +## Firehose Consumption + +For social features (trending, discovery, feeds): + +- **WebSocket Connection**: Subscribe to `com.atproto.sync.subscribeRepos` from a Relay +- **CBOR Decoding**: Parse incoming events (or use Jetstream for JSON) +- **Cursor Management**: Track position for reconnection + +## AppView Pattern + +Index network-wide records to power discovery features: + +- Index `app.malfestio.*` records from firehose +- Implement `getFeedSkeleton` for custom algorithmic feeds +- Hydration service combines skeletons with full content from PDSes + +## Well-Known Endpoints + +- `/.well-known/atproto-did` — Domain verification for handle claims +- `/.well-known/oauth-protected-resource` — PDS OAuth metadata +- `/.well-known/oauth-authorization-server` — Auth server metadata + +## Patterns from Real AT Protocol Apps + +### plyr.fm (Music) + +- OAuth 2.1 via `@atproto/oauth-client` library +- Records synced to PDS: tracks, likes, playlists +- Separate moderation service (Rust labeler) +- Data ownership: "tracks, likes, playlists synced to your PDS as ATProto records" + +### leaflet.pub (Writing) + +- React/Next.js frontend with Supabase + Replicache for sync +- Bluesky integration via dedicated `lexicons/` and `appview/` directories +- Publications posted to Bluesky + +### wisp.place (Static Sites) + +- Stores site files as `place.wisp.fs` records in user's PDS +- Firehose consumer to index and serve sites +- CDN layer caches content from PDS + +### Common Patterns + +1. Local database for fast queries + PDS for portable, signed records +2. Firehose consumption for discovery/aggregation +3. OAuth 2.1 for production auth (app passwords only for development) +4. Lexicons define the public contract; internal state stays private + +## References + +- [AT Protocol OAuth Spec](https://atproto.com/specs/oauth) +- [Lexicon Schema Language](https://atproto.com/specs/lexicon) +- [Repository & XRPC](https://atproto.com/specs/xrpc) +- [Feed Generator Starter Kit](https://github.com/bluesky-social/feed-generator) +- [atproto TypeScript SDK](https://github.com/bluesky-social/atproto) diff --git a/docs/todo.md b/docs/todo.md index 05b642e..c90bcb2 100644 --- a/docs/todo.md +++ b/docs/todo.md @@ -84,10 +84,6 @@ - A user can authenticate via OAuth, create a deck, and see it in their PDS repository. -#### Notes - -- See [docs/at.md](at.md) for full AT Protocol integration research. - ### Milestone G - Study Engine (SRS) + Daily Review UX #### Deliverables