diff --git a/apps/aqua/src/api/mod.rs b/apps/aqua/src/api/mod.rs index 5466ecd..793f868 100644 --- a/apps/aqua/src/api/mod.rs +++ b/apps/aqua/src/api/mod.rs @@ -126,12 +126,6 @@ pub async fn get_car_import_job_status( } } -#[derive(Debug, Serialize, Deserialize)] -pub struct CarImportRequest { - pub import_id: Option, - pub description: Option, -} - #[derive(Debug, Serialize, Deserialize)] pub struct CarImportResponse { pub import_id: String, diff --git a/apps/aqua/src/db.rs b/apps/aqua/src/db.rs index 2b03656..433a93b 100644 --- a/apps/aqua/src/db.rs +++ b/apps/aqua/src/db.rs @@ -26,7 +26,7 @@ pub async fn init_pool() -> anyhow::Result { let pool = PgPoolOptions::new() .max_connections(50) .min_connections(1) - .max_lifetime(std::time::Duration::from_secs(10)) + .max_lifetime(std::time::Duration::from_secs(30 * 60)) .connect(&env::var("DATABASE_URL")?) .await?; info!(target: "db", "Connected to the database!"); diff --git a/apps/aqua/src/repos/feed_play.rs b/apps/aqua/src/repos/feed_play.rs index bd117c7..ae47788 100644 --- a/apps/aqua/src/repos/feed_play.rs +++ b/apps/aqua/src/repos/feed_play.rs @@ -61,11 +61,15 @@ impl FeedPlayRepo for PgDataSource { profile.did, profile.handle, profile.display_name, profile.avatar ORDER BY p.processed_time desc "#, - &uri.to_string() + uri ) - .fetch_one(&self.db) + .fetch_optional(&self.db) .await?; + let Some(row) = row else { + return Ok(None); + }; + let artists: Vec = match row.artists { Some(value) => from_json_value::>(value).unwrap_or_default(), None => vec![], diff --git a/apps/aqua/src/xrpc/actor.rs b/apps/aqua/src/xrpc/actor.rs index b23cf43..bbdd4c7 100644 --- a/apps/aqua/src/xrpc/actor.rs +++ b/apps/aqua/src/xrpc/actor.rs @@ -26,16 +26,12 @@ pub async fn get_actor( axum::extract::Query(query): axum::extract::Query, ) -> Result { let repo = &ctx.db; // assuming ctx.db is Box - let identity = &query.actor; + let identity = query + .actor + .as_deref() + .ok_or_else(|| (StatusCode::BAD_REQUEST, "actor is required".to_string()))?; - if identity.is_none() { - return Err((StatusCode::BAD_REQUEST, "actor is required".to_string())); - } - - match repo - .get_actor_profile(identity.as_ref().expect("actor is not none").as_str()) - .await - { + match repo.get_actor_profile(identity).await { Ok(Some(profile)) => Ok(axum::Json(GetProfileResponse { profile: profile.into_static(), })), diff --git a/services/cadet/src/bin/tap_backfill.rs b/services/cadet/src/bin/tap_backfill.rs index 5ff7405..401ac8a 100644 --- a/services/cadet/src/bin/tap_backfill.rs +++ b/services/cadet/src/bin/tap_backfill.rs @@ -142,7 +142,7 @@ async fn main() -> anyhow::Result<()> { Ok(event) => match ingestor.ingest(event).await { Ok(()) => { processed += 1; - if processed % 100 == 0 { + if processed.is_multiple_of(100) { info!( "Processed {} TAP Teal records (skipped {})", processed, skipped diff --git a/services/cadet/src/db.rs b/services/cadet/src/db.rs index 2b03656..433a93b 100644 --- a/services/cadet/src/db.rs +++ b/services/cadet/src/db.rs @@ -26,7 +26,7 @@ pub async fn init_pool() -> anyhow::Result { let pool = PgPoolOptions::new() .max_connections(50) .min_connections(1) - .max_lifetime(std::time::Duration::from_secs(10)) + .max_lifetime(std::time::Duration::from_secs(30 * 60)) .connect(&env::var("DATABASE_URL")?) .await?; info!(target: "db", "Connected to the database!"); diff --git a/services/cadet/src/ingestors/car/car_import.rs b/services/cadet/src/ingestors/car/car_import.rs index ee6b11d..57799c8 100644 --- a/services/cadet/src/ingestors/car/car_import.rs +++ b/services/cadet/src/ingestors/car/car_import.rs @@ -120,12 +120,12 @@ impl CarImportIngestor { info!("Extracted {} records from MST", records.len()); // Process each record through the appropriate ingestor - let mut processed_count = 0; + let mut processed_count = 0_usize; for record in records { match self.process_extracted_record(&record, import_id, did).await { Ok(()) => { processed_count += 1; - if processed_count % 10 == 0 { + if processed_count.is_multiple_of(10) { info!("Processed {} records so far", processed_count); } } diff --git a/services/cadet/src/resolve.rs b/services/cadet/src/resolve.rs index 7897975..df31807 100644 --- a/services/cadet/src/resolve.rs +++ b/services/cadet/src/resolve.rs @@ -5,31 +5,26 @@ use serde::{Deserialize, Serialize}; use anyhow::{anyhow, Result}; -// should be same as regex /^did:[a-z]+:[\S\s]+/ -fn is_did(did: &str) -> bool { - let parts: Vec<&str> = did.split(':').collect(); - - if parts.len() != 3 { - // must have exactly 3 parts: "did", method, and identifier - return false; - } - - if parts[0] != "did" { - // first part must be "did" - return false; - } - - if !parts[1].chars().all(|c| c.is_ascii_lowercase()) || parts[1].is_empty() { - // method must be all lowercase - return false; - } - - if parts[2].is_empty() { - // identifier can't be empty - return false; - } +// This deliberately only checks the DID's envelope. Method-specific validation +// belongs to the relevant resolver. In particular, a did:web identifier may +// contain additional colons for path components. +fn did_parts(did: &str) -> Option<(&str, &str)> { + let mut parts = did.splitn(3, ':'); + let prefix = parts.next()?; + let method = parts.next()?; + let identifier = parts.next()?; + + (prefix == "did" + && !method.is_empty() + && method + .chars() + .all(|character| character.is_ascii_lowercase()) + && !identifier.is_empty()) + .then_some((method, identifier)) +} - true +fn is_did(did: &str) -> bool { + did_parts(did).is_some() } fn is_valid_domain(domain: &str) -> bool { @@ -80,14 +75,8 @@ async fn resolve_handle(handle: &str, resolver_app_view: &str) -> Result Result { - // get the specific did spec - // did:plc:abcd1e -> plc - let parts: Vec<&str> = did.split(':').collect(); - let spec = parts[1]; - if spec.is_empty() { - return Err(anyhow!("Empty spec in DID: {}", did)); - } - match spec { + let (method, _) = did_parts(did).ok_or_else(|| anyhow!("Invalid DID: {did}"))?; + match method { "plc" => { let res: DidDocument = reqwest::get(format!("https://plc.directory/{}", did)) .await? @@ -97,11 +86,7 @@ async fn get_did_doc(did: &str) -> Result { Ok(res) } "web" => { - if !is_valid_domain(parts[2]) { - todo!("Error for domain in did:web is not valid"); - }; - let ident = parts[2]; - let res = reqwest::get(format!("https://{}/.well-known/did.json", ident)) + let res = reqwest::get(did_web_document_url(did)?) .await? .error_for_status()? .json() @@ -109,10 +94,40 @@ async fn get_did_doc(did: &str) -> Result { Ok(res) } - _ => todo!("Identifier not supported"), + _ => Err(anyhow!("Unsupported DID method: {method}")), } } +fn did_web_document_url(did: &str) -> Result { + let (method, identifier) = did_parts(did).ok_or_else(|| anyhow!("Invalid DID: {did}"))?; + if method != "web" { + return Err(anyhow!("Expected a did:web DID: {did}")); + } + + let mut components = identifier.split(':'); + let host = components + .next() + .ok_or_else(|| anyhow!("Invalid did:web DID: {did}"))?; + if !is_valid_domain(host) { + return Err(anyhow!("Invalid did:web host: {host}")); + } + + let path_components = components.collect::>(); + if path_components + .iter() + .any(|component| component.is_empty() || component.contains(['/', '?', '#'])) + { + return Err(anyhow!("Invalid did:web path in DID: {did}")); + } + + let path = if path_components.is_empty() { + ".well-known/did.json".to_owned() + } else { + format!("{}/did.json", path_components.join("/")) + }; + Ok(format!("https://{host}/{path}")) +} + fn get_pds_endpoint(doc: &DidDocument) -> Option { get_service_endpoint(doc, "#atproto_pds", "AtprotoPersonalDataServer") } @@ -129,30 +144,23 @@ fn get_service_endpoint( } pub async fn resolve_identity(id: &str, resolver_app_view: &str) -> Result { - // is our identifier a did let did = if is_did(id) { - id + id.to_owned() } else { - // our id must be either invalid or a handle - if let Ok(res) = resolve_handle(id, resolver_app_view).await { - &res.clone() - } else { - todo!("Error type for could not resolve handle") - } + resolve_handle(id, resolver_app_view) + .await + .map_err(|error| anyhow!("Failed to resolve handle {id}: {error}"))? }; - let doc = get_did_doc(did).await?; - let pds = get_pds_endpoint(&doc); - - if pds.is_none() { - todo!("Error for could not find PDS") - } + let doc = get_did_doc(&did).await?; + let pds = get_pds_endpoint(&doc) + .ok_or_else(|| anyhow!("No AT Protocol PDS service found for DID: {did}"))?; Ok(ResolvedIdentity { - did: did.to_owned(), + did, doc, identity: id.to_owned(), - pds: pds.unwrap().service_endpoint, + pds: pds.service_endpoint, }) } @@ -212,6 +220,7 @@ fn test_match_did() { assert!(!is_did("did::123")); // empty method assert!(!is_did("notdid:example:123")); // doesn't start with did assert!(!is_did("did:example")); // missing identifier part + assert!(is_did("did:web:example.com:users:alice")); } #[test] @@ -227,3 +236,17 @@ fn test_valid_domain() { assert!(!is_valid_domain("-example.com")); // starts with hyphen assert!(!is_valid_domain("example-.com")); // ends with hyphen } + +#[test] +fn test_did_web_document_url() { + assert_eq!( + did_web_document_url("did:web:example.com").unwrap(), + "https://example.com/.well-known/did.json" + ); + assert_eq!( + did_web_document_url("did:web:example.com:users:alice").unwrap(), + "https://example.com/users/alice/did.json" + ); + assert!(did_web_document_url("did:web:example.com::alice").is_err()); + assert!(did_web_document_url("did:web:not_a_domain").is_err()); +} diff --git a/services/satellite/src/counts.rs b/services/satellite/src/counts.rs index 5ab0f06..dc4d85b 100644 --- a/services/satellite/src/counts.rs +++ b/services/satellite/src/counts.rs @@ -4,6 +4,8 @@ use axum::{ Json, }; use serde::{Deserialize, Serialize}; +use serde_json::Value; +use sqlx::types::Json as SqlxJson; use sqlx::FromRow; use uuid::Uuid; @@ -51,8 +53,7 @@ pub struct Play { pub release_mbid: Option, pub duration: Option, pub uri: Option, - // MASSIVE HUGE HACK - pub artists: Option, + pub artists: SqlxJson, } #[derive(FromRow, Debug, Deserialize, Serialize)] @@ -85,10 +86,15 @@ pub async fn get_latest_plays( SELECT p.did, p.track_name, - STRING_AGG( - ptae.artist_name || '|' || COALESCE(TEXT(ae.mbid), ''), - ',' - ORDER BY ptae.artist_name + COALESCE( + JSON_AGG( + JSON_BUILD_OBJECT( + 'artist_name', ptae.artist_name, + 'artist_mbid', ae.mbid + ) + ORDER BY ptae.artist_name + ) FILTER (WHERE ptae.artist_name IS NOT NULL), + '[]'::json ) AS artists, p.release_name, p.duration, @@ -110,28 +116,16 @@ pub async fn get_latest_plays( match result { Ok(counts) => { - let fin: Vec = counts + let fin = counts .into_iter() - .map(|play| -> PlayReturn { - let artists = play - .artists - .unwrap_or_default() - .split(',') - .filter(|artist| !artist.is_empty()) - .map(|artist| { - let mut parts = artist.split('|'); - Artist { - artist_name: parts - .next() - .expect("Artist name is required") - .to_string(), - artist_mbid: parts - .next() - .and_then(|mbid| Uuid::parse_str(mbid).ok()), - } - }) - .collect(); - PlayReturn { + .map(|play| { + let artists = serde_json::from_value(play.artists.0).map_err(|error| { + ( + StatusCode::INTERNAL_SERVER_ERROR, + format!("Invalid artist data returned by database: {error}"), + ) + })?; + Ok(PlayReturn { did: play.did.to_string(), track_name: play.track_name, recording_mbid: play.recording_mbid, @@ -140,9 +134,9 @@ pub async fn get_latest_plays( duration: play.duration, uri: play.uri, artists, - } + }) }) - .collect(); + .collect::, _>>()?; Ok(Json(fin)) } @@ -152,3 +146,22 @@ pub async fn get_latest_plays( )), } } + +#[cfg(test)] +mod tests { + use super::Artist; + use uuid::Uuid; + + #[test] + fn deserializes_artist_names_with_legacy_delimiters() { + let mbid = Uuid::nil(); + let artists: Vec = serde_json::from_value(serde_json::json!([ + {"artist_name": "AC/DC, Live | Remastered", "artist_mbid": mbid} + ])) + .unwrap(); + + assert_eq!(artists.len(), 1); + assert_eq!(artists[0].artist_name, "AC/DC, Live | Remastered"); + assert_eq!(artists[0].artist_mbid, Some(mbid)); + } +} diff --git a/services/satellite/src/db.rs b/services/satellite/src/db.rs index 2b03656..433a93b 100644 --- a/services/satellite/src/db.rs +++ b/services/satellite/src/db.rs @@ -26,7 +26,7 @@ pub async fn init_pool() -> anyhow::Result { let pool = PgPoolOptions::new() .max_connections(50) .min_connections(1) - .max_lifetime(std::time::Duration::from_secs(10)) + .max_lifetime(std::time::Duration::from_secs(30 * 60)) .connect(&env::var("DATABASE_URL")?) .await?; info!(target: "db", "Connected to the database!"); diff --git a/todo.md b/todo.md index 4d07449..630b36e 100644 --- a/todo.md +++ b/todo.md @@ -12,6 +12,7 @@ Last synced with GitHub and Linear issues: 2026-06-14. - Use `pnpm tunnel:up`, `pnpm tunnel:down`, `pnpm tunnel:status`, `pnpm tunnel:logs`, and `pnpm tunnel:verify` for the stable preview. - The public Amethyst feed must use only live Aqua XRPC data. Do not add seeded, mocked, demo, or backup play data. - Public preview refreshed on 2026-06-15 by rebuilding Amethyst, Aqua, and Cadet images, recreating the Compose preview stack, and verifying `https://sigilyph.teal.fm/client-metadata.json` plus latest plays XRPC. +- Focused code-health pass (2026-07-10): fixed `did:web` path resolution and resolver error handling in Cadet, avoided 10-second Postgres connection churn in Aqua/Cadet/Satellite, returned proper 404s for missing feed plays, and preserved delimiter-containing artist names in Satellite's latest-play response. Verified with targeted Rust checks and Cadet resolver tests. ## Local Open Work