From 152d9ac0798c1e4c4ca033ad15a588d79f9ab117 Mon Sep 17 00:00:00 2001 From: Timothy Quilling Date: Wed, 5 Nov 2025 17:30:10 -0500 Subject: [PATCH] feat: add account_created_at --- Cargo.lock | 1 + .../database_writer/operations/executor.rs | 2 + .../src/database_writer/operations/types.rs | 1 + consumer/src/db/actor.rs | 84 +++++++++++++++++-- consumer/src/db/allowlist.rs | 1 + consumer/src/workers/handle_resolver.rs | 53 ++++++++++-- consumer/src/workers/jetstream/processing.rs | 2 + did-resolver/Cargo.toml | 1 + did-resolver/src/lib.rs | 34 ++++++++ did-resolver/src/types.rs | 8 ++ .../down.sql | 3 + .../up.sql | 8 ++ parakeet-db/src/schema.rs | 1 + parakeet/src/hydration/profile/builders.rs | 23 ++--- parakeet/src/loaders/profile.rs | 14 ++-- parakeet/src/xrpc/app_bsky/actor.rs | 2 +- .../src/xrpc/app_bsky/notification/mod.rs | 10 +-- parakeet/src/xrpc/app_bsky/unspecced/mod.rs | 2 +- 18 files changed, 214 insertions(+), 36 deletions(-) create mode 100644 migrations/2025-11-05-120000_add_account_created_at_to_actors/down.sql create mode 100644 migrations/2025-11-05-120000_add_account_created_at_to_actors/up.sql diff --git a/Cargo.lock b/Cargo.lock index 065db8a3..e40c2461 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1009,6 +1009,7 @@ dependencies = [ name = "did-resolver" version = "0.1.0" dependencies = [ + "chrono", "hickory-resolver", "reqwest", "serde", diff --git a/consumer/src/database_writer/operations/executor.rs b/consumer/src/database_writer/operations/executor.rs index c7a7f7f1..c4c8e5fa 100644 --- a/consumer/src/database_writer/operations/executor.rs +++ b/consumer/src/database_writer/operations/executor.rs @@ -573,6 +573,7 @@ pub async fn execute_operation( status, sync_state, handle, + account_created_at, timestamp, } => { db::actor_upsert( @@ -581,6 +582,7 @@ pub async fn execute_operation( status.as_ref(), &sync_state, handle.as_deref(), + account_created_at.as_ref(), timestamp, ) .await?; diff --git a/consumer/src/database_writer/operations/types.rs b/consumer/src/database_writer/operations/types.rs index 0e9adb6b..39e41b37 100644 --- a/consumer/src/database_writer/operations/types.rs +++ b/consumer/src/database_writer/operations/types.rs @@ -121,6 +121,7 @@ pub enum DatabaseOperation { status: Option, sync_state: parakeet_db::types::ActorSyncState, handle: Option, + account_created_at: Option>, timestamp: DateTime, }, diff --git a/consumer/src/db/actor.rs b/consumer/src/db/actor.rs index b54b7b8b..e646351f 100644 --- a/consumer/src/db/actor.rs +++ b/consumer/src/db/actor.rs @@ -11,15 +11,37 @@ pub async fn actor_upsert( status: Option<&ActorStatus>, sync_state: &ActorSyncState, handle: Option<&str>, + account_created_at: Option<&DateTime>, time: DateTime, ) -> Result { // Allow allowlist states (synced, dirty, processing) to flow freely // Allow upgrading from partial to allowlist states // Never downgrade from allowlist states to partial + // + // Account created_at is updated if provided and current value is NULL + // This allows enrichment during handle resolution without overwriting existing values - match (status, handle) { - (Some(status), Some(handle)) => { - // Both status and handle provided + match (status, handle, account_created_at) { + (Some(status), Some(handle), Some(created_at)) => { + // All three provided + conn.execute( + "INSERT INTO actors (did, status, handle, sync_state, account_created_at, last_indexed) VALUES ($1, $2, $3, $4, $5, $6) + ON CONFLICT (did) DO UPDATE SET + status=EXCLUDED.status, + handle=EXCLUDED.handle, + sync_state=CASE + WHEN actors.sync_state IN ('synced', 'dirty', 'processing') AND EXCLUDED.sync_state = 'partial' + THEN actors.sync_state + ELSE EXCLUDED.sync_state + END, + account_created_at=COALESCE(actors.account_created_at, EXCLUDED.account_created_at), + last_indexed=EXCLUDED.last_indexed", + &[&did, &status, &handle, &sync_state, &created_at, &time], + ) + .await + } + (Some(status), Some(handle), None) => { + // Status and handle, no created_at conn.execute( "INSERT INTO actors (did, status, handle, sync_state, last_indexed) VALUES ($1, $2, $3, $4, $5) ON CONFLICT (did) DO UPDATE SET @@ -35,7 +57,24 @@ pub async fn actor_upsert( ) .await } - (Some(status), None) => { + (Some(status), None, Some(created_at)) => { + // Status and created_at, no handle + conn.execute( + "INSERT INTO actors (did, status, sync_state, account_created_at, last_indexed) VALUES ($1, $2, $3, $4, $5) + ON CONFLICT (did) DO UPDATE SET + status=EXCLUDED.status, + sync_state=CASE + WHEN actors.sync_state IN ('synced', 'dirty', 'processing') AND EXCLUDED.sync_state = 'partial' + THEN actors.sync_state + ELSE EXCLUDED.sync_state + END, + account_created_at=COALESCE(actors.account_created_at, EXCLUDED.account_created_at), + last_indexed=EXCLUDED.last_indexed", + &[&did, &status, &sync_state, &created_at, &time], + ) + .await + } + (Some(status), None, None) => { // Only status provided conn.execute( "INSERT INTO actors (did, status, sync_state, last_indexed) VALUES ($1, $2, $3, $4) @@ -51,7 +90,24 @@ pub async fn actor_upsert( ) .await } - (None, Some(handle)) => { + (None, Some(handle), Some(created_at)) => { + // Handle and created_at, no status + conn.execute( + "INSERT INTO actors (did, handle, sync_state, account_created_at, last_indexed) VALUES ($1, $2, $3, $4, $5) + ON CONFLICT (did) DO UPDATE SET + handle=EXCLUDED.handle, + sync_state=CASE + WHEN actors.sync_state IN ('synced', 'dirty', 'processing') AND EXCLUDED.sync_state = 'partial' + THEN actors.sync_state + ELSE EXCLUDED.sync_state + END, + account_created_at=COALESCE(actors.account_created_at, EXCLUDED.account_created_at), + last_indexed=EXCLUDED.last_indexed", + &[&did, &handle, &sync_state, &created_at, &time], + ) + .await + } + (None, Some(handle), None) => { // Only handle provided conn.execute( "INSERT INTO actors (did, handle, sync_state, last_indexed) VALUES ($1, $2, $3, $4) @@ -67,7 +123,23 @@ pub async fn actor_upsert( ) .await } - (None, None) => { + (None, None, Some(created_at)) => { + // Only created_at provided + conn.execute( + "INSERT INTO actors (did, sync_state, account_created_at, last_indexed) VALUES ($1, $2, $3, $4) + ON CONFLICT (did) DO UPDATE SET + sync_state=CASE + WHEN actors.sync_state IN ('synced', 'dirty', 'processing') AND EXCLUDED.sync_state = 'partial' + THEN actors.sync_state + ELSE EXCLUDED.sync_state + END, + account_created_at=COALESCE(actors.account_created_at, EXCLUDED.account_created_at), + last_indexed=EXCLUDED.last_indexed", + &[&did, &sync_state, &created_at, &time], + ) + .await + } + (None, None, None) => { // Neither provided - just ensure actor exists with sync_state conn.execute( "INSERT INTO actors (did, sync_state, last_indexed) VALUES ($1, $2, $3) diff --git a/consumer/src/db/allowlist.rs b/consumer/src/db/allowlist.rs index 60c03c6c..920e7c4f 100644 --- a/consumer/src/db/allowlist.rs +++ b/consumer/src/db/allowlist.rs @@ -172,6 +172,7 @@ pub async fn ensure_allowlist_actors( Some(&ActorStatus::Active), &ActorSyncState::Dirty, None, // no handle + None, // no account_created_at (enriched during handle resolution) now, ) .await diff --git a/consumer/src/workers/handle_resolver.rs b/consumer/src/workers/handle_resolver.rs index f389bf9a..53e2216e 100644 --- a/consumer/src/workers/handle_resolver.rs +++ b/consumer/src/workers/handle_resolver.rs @@ -118,14 +118,32 @@ impl HandleResolutionWorker { let cache_key = format!("did_handle:{}", did); let cached_handle: Option = self.redis.get(&cache_key).await?; - let handle = if let Some(cached) = cached_handle { + let (handle, account_created_at) = if let Some(cached) = cached_handle { counter!("handle_resolver.cache_hit").increment(1); - Some(cached) + // Cache hit - we have the handle but not the account creation time + // Fetch account creation time separately (not cached, but only fetched once per actor typically) + let created_at = match self.resolver.get_plc_creation_time(did).await { + Ok(Some(timestamp)) => { + counter!("handle_resolver.plc_creation_fetch_success").increment(1); + Some(timestamp) + } + Ok(None) => { + debug!("No PLC creation time found for {}", did); + counter!("handle_resolver.plc_creation_not_found").increment(1); + None + } + Err(e) => { + debug!("Failed to fetch PLC creation time for {}: {}", did, e); + counter!("handle_resolver.plc_creation_error").increment(1); + None + } + }; + (Some(cached), created_at) } else { // 2. Cache miss - resolve from PLC directory counter!("handle_resolver.cache_miss").increment(1); - match self.resolver.resolve_did(did).await { + let (handle, created_at) = match self.resolver.resolve_did(did).await { Ok(Some(doc)) => { // Extract handle from DID document's alsoKnownAs field let handle = doc @@ -148,19 +166,39 @@ impl HandleResolutionWorker { counter!("handle_resolver.no_handle").increment(1); } - handle + // Fetch account creation time from PLC audit log + let created_at = match self.resolver.get_plc_creation_time(did).await { + Ok(Some(timestamp)) => { + counter!("handle_resolver.plc_creation_fetch_success").increment(1); + Some(timestamp) + } + Ok(None) => { + debug!("No PLC creation time found for {}", did); + counter!("handle_resolver.plc_creation_not_found").increment(1); + None + } + Err(e) => { + debug!("Failed to fetch PLC creation time for {}: {}", did, e); + counter!("handle_resolver.plc_creation_error").increment(1); + None + } + }; + + (handle, created_at) } Ok(None) => { debug!("DID document not found for {}", did); counter!("handle_resolver.did_not_found").increment(1); - None + (None, None) } Err(e) => { debug!("Failed to resolve DID for {}: {}", did, e); counter!("handle_resolver.did_resolution_error").increment(1); - None + (None, None) } - } + }; + + (handle, created_at) }; // 4. Send handle update operation to batch writer (bounded channel with backpressure) @@ -172,6 +210,7 @@ impl HandleResolutionWorker { status: None, sync_state: parakeet_db::types::ActorSyncState::Partial, handle, + account_created_at, timestamp: chrono::Utc::now(), }], cache_invalidations: vec![], diff --git a/consumer/src/workers/jetstream/processing.rs b/consumer/src/workers/jetstream/processing.rs index 30768d8d..7755bb7f 100644 --- a/consumer/src/workers/jetstream/processing.rs +++ b/consumer/src/workers/jetstream/processing.rs @@ -202,6 +202,7 @@ pub fn process_identity_event( status: None, sync_state, handle, // Pass the Option directly + account_created_at: None, // Jetstream doesn't provide creation time, enriched during handle resolution timestamp, }]; @@ -244,6 +245,7 @@ pub fn process_account_event( status: Some(status), sync_state, handle: None, + account_created_at: None, // Jetstream doesn't provide creation time, enriched during handle resolution timestamp, }]; diff --git a/did-resolver/Cargo.toml b/did-resolver/Cargo.toml index 5d58b37c..1a50c4fb 100644 --- a/did-resolver/Cargo.toml +++ b/did-resolver/Cargo.toml @@ -4,6 +4,7 @@ version = "0.1.0" edition = "2021" [dependencies] +chrono = { version = "0.4", features = ["serde"] } hickory-resolver = "0.24.2" reqwest = { version = "0.12.12", features = ["json", "native-tls"] } serde = { version = "1.0.217", features = ["derive"] } diff --git a/did-resolver/src/lib.rs b/did-resolver/src/lib.rs index 0786b08f..5435e76f 100644 --- a/did-resolver/src/lib.rs +++ b/did-resolver/src/lib.rs @@ -79,6 +79,40 @@ impl Resolver { Ok(Some(did_doc)) } + /// Get account creation timestamp from PLC audit log + /// Returns the timestamp of the first operation (where prev is null) + /// Only works for did:plc DIDs - returns None for other DID methods + pub async fn get_plc_creation_time( + &self, + did: &str, + ) -> Result>, Error> { + // Only fetch for did:plc + if !did.starts_with("did:plc:") { + return Ok(None); + } + + let res = self + .client + .get(format!("{}/{did}/log/audit", self.plc)) + .send() + .await?; + + let status = res.status(); + + if status.is_server_error() { + return Err(Error::ServerError); + } + + if status == StatusCode::NOT_FOUND || status == StatusCode::GONE { + return Ok(None); + } + + let audit_log: Vec = res.json().await?; + + // First entry in the audit log is the account creation + Ok(audit_log.first().map(|entry| entry.created_at)) + } + async fn resolve_did_web(&self, id: &str) -> Result, Error> { let res = match self .client diff --git a/did-resolver/src/types.rs b/did-resolver/src/types.rs index c7399c4b..4dc509c1 100644 --- a/did-resolver/src/types.rs +++ b/did-resolver/src/types.rs @@ -1,5 +1,13 @@ use serde::{Deserialize, Serialize}; +#[derive(Debug, Deserialize, Serialize)] +#[serde(rename_all = "camelCase")] +pub struct PlcAuditLogEntry { + pub did: String, + pub created_at: chrono::DateTime, + // Other fields exist but we only need createdAt +} + #[derive(Debug, Deserialize, Serialize)] #[serde(rename_all = "camelCase")] pub struct DidDocument { diff --git a/migrations/2025-11-05-120000_add_account_created_at_to_actors/down.sql b/migrations/2025-11-05-120000_add_account_created_at_to_actors/down.sql new file mode 100644 index 00000000..2b92ccf7 --- /dev/null +++ b/migrations/2025-11-05-120000_add_account_created_at_to_actors/down.sql @@ -0,0 +1,3 @@ +-- Remove account_created_at column from actors table +DROP INDEX IF EXISTS idx_actors_account_created_at; +ALTER TABLE actors DROP COLUMN account_created_at; diff --git a/migrations/2025-11-05-120000_add_account_created_at_to_actors/up.sql b/migrations/2025-11-05-120000_add_account_created_at_to_actors/up.sql new file mode 100644 index 00000000..924995e2 --- /dev/null +++ b/migrations/2025-11-05-120000_add_account_created_at_to_actors/up.sql @@ -0,0 +1,8 @@ +-- Add account_created_at column to actors table +-- This represents when the account was created (from PLC directory first operation) +-- Nullable because we may not know it for all actors initially (especially stubs) +-- Will be populated during handle resolution +ALTER TABLE actors ADD COLUMN account_created_at TIMESTAMP WITH TIME ZONE; + +-- Create index for potential queries by account age +CREATE INDEX idx_actors_account_created_at ON actors(account_created_at) WHERE account_created_at IS NOT NULL; diff --git a/parakeet-db/src/schema.rs b/parakeet-db/src/schema.rs index 902e9280..78e255a7 100644 --- a/parakeet-db/src/schema.rs +++ b/parakeet-db/src/schema.rs @@ -124,6 +124,7 @@ diesel::table! { repo_rev -> Nullable, repo_cid -> Nullable, last_indexed -> Nullable, + account_created_at -> Nullable, } } diff --git a/parakeet/src/hydration/profile/builders.rs b/parakeet/src/hydration/profile/builders.rs index 4020a482..6bbf6c26 100644 --- a/parakeet/src/hydration/profile/builders.rs +++ b/parakeet/src/hydration/profile/builders.rs @@ -115,7 +115,7 @@ fn build_status(status: crate::loaders::EnrichedStatus, cdn: &BskyCdn) -> Option } pub(super) fn build_basic( - (did, handle, profile, chat_decl, is_labeler, status, notif_decl): ProfileLoaderRet, + (did, handle, account_created_at, profile, chat_decl, is_labeler, status, notif_decl): ProfileLoaderRet, stats: Option, labels: Vec, verifications: Option>, @@ -141,8 +141,9 @@ pub(super) fn build_basic( }) .unwrap_or_else(|| (None, None, None)); - // Use profile.created_at if available, fall back to Utc::now() - let created_at = profile.as_ref().map(|p| p.created_at).unwrap_or_else(Utc::now); + // Use actor.account_created_at if available (from PLC directory), fall back to Utc::now() + // This is the actual account creation time, not when the profile record was created + let created_at = account_created_at.unwrap_or_else(Utc::now); ProfileViewBasic { did, @@ -160,7 +161,7 @@ pub(super) fn build_basic( } pub(super) fn build_profile( - (did, handle, profile, chat_decl, is_labeler, status, notif_decl): ProfileLoaderRet, + (did, handle, account_created_at, profile, chat_decl, is_labeler, status, notif_decl): ProfileLoaderRet, stats: Option, labels: Vec, verifications: Option>, @@ -187,9 +188,9 @@ pub(super) fn build_profile( }) .unwrap_or_else(|| (None, None, None, None)); - // Use profile.created_at for both createdAt and indexedAt (if profile exists) - // Fall back to Utc::now() if no profile record exists yet - let created_at = profile.as_ref().map(|p| p.created_at).unwrap_or_else(Utc::now); + // Use actor.account_created_at for both createdAt and indexedAt + // Fall back to Utc::now() if account creation time not available + let created_at = account_created_at.unwrap_or_else(Utc::now); ProfileView { did, @@ -209,7 +210,7 @@ pub(super) fn build_profile( } pub(super) fn build_detailed( - (did, handle, profile, chat_decl, is_labeler, status, notif_decl): ProfileLoaderRet, + (did, handle, account_created_at, profile, chat_decl, is_labeler, status, notif_decl): ProfileLoaderRet, stats: Option, labels: Vec, verifications: Option>, @@ -241,9 +242,9 @@ pub(super) fn build_detailed( }) .unwrap_or_else(|| (None, None, None, None, None, None)); - // Use profile.created_at for both createdAt and indexedAt (if profile exists) - // Fall back to Utc::now() if no profile record exists yet - let created_at = profile.as_ref().map(|p| p.created_at).unwrap_or_else(Utc::now); + // Use actor.account_created_at for both createdAt and indexedAt + // Fall back to Utc::now() if account creation time not available + let created_at = account_created_at.unwrap_or_else(Utc::now); ProfileViewDetailed { did, diff --git a/parakeet/src/loaders/profile.rs b/parakeet/src/loaders/profile.rs index 3f61871e..d5c70789 100644 --- a/parakeet/src/loaders/profile.rs +++ b/parakeet/src/loaders/profile.rs @@ -16,6 +16,7 @@ pub fn build_profiles_batch_query(dids_list: &str) -> String { "SELECT a.did, a.handle, + a.account_created_at, p.actor_id, p.cid, p.created_at, @@ -171,9 +172,10 @@ impl BatchFn for HandleLoader { pub struct ProfileLoader(pub(super) Pool, pub(super) redis::aio::MultiplexedConnection); pub type ProfileLoaderRet = ( - String, // did - always present (from actors table) - Option, // handle - from actors table - Option, // profile - optional (may not exist) + String, // did - always present (from actors table) + Option, // handle - from actors table + Option>, // account_created_at - from actors table + Option, // profile - optional (may not exist) Option, bool, Option, @@ -198,6 +200,8 @@ impl BatchFn for ProfileLoader { did: String, #[diesel(sql_type = diesel::sql_types::Nullable)] handle: Option, + #[diesel(sql_type = diesel::sql_types::Nullable)] + account_created_at: Option>, #[diesel(sql_type = diesel::sql_types::Nullable)] actor_id: Option, #[diesel(sql_type = diesel::sql_types::Nullable)] @@ -344,7 +348,7 @@ impl BatchFn for ProfileLoader { let is_labeler = row.labeler_actor_id.is_some(); let status = status_map.get(&row.did).cloned(); - let val = (row.did.clone(), row.handle, profile, chat_decl, is_labeler, status, notif_decl); + let val = (row.did.clone(), row.handle, row.account_created_at, profile, chat_decl, is_labeler, status, notif_decl); (row.did, val) }, @@ -353,7 +357,7 @@ impl BatchFn for ProfileLoader { // Enqueue missing profiles for background fetching using Redis let missing_profiles: Vec = results .iter() - .filter(|(_, (_, _, profile, _, _, _, _))| profile.is_none()) + .filter(|(_, (_, _, _, profile, _, _, _, _))| profile.is_none()) .map(|(did, _)| format!("at://{}/app.bsky.actor.profile/self", did)) .collect(); diff --git a/parakeet/src/xrpc/app_bsky/actor.rs b/parakeet/src/xrpc/app_bsky/actor.rs index 3ae0cb95..21b16d5b 100644 --- a/parakeet/src/xrpc/app_bsky/actor.rs +++ b/parakeet/src/xrpc/app_bsky/actor.rs @@ -319,7 +319,7 @@ pub async fn get_suggestions( // Check if has profile (display name or description) profiles .get(did) - .and_then(|p| p.2.as_ref()) // .2 is Option + .and_then(|p| p.3.as_ref()) // .3 is Option .map(|prof| { prof.display_name.is_some() || prof.description.is_some() }) diff --git a/parakeet/src/xrpc/app_bsky/notification/mod.rs b/parakeet/src/xrpc/app_bsky/notification/mod.rs index 297a26d6..9306f2df 100644 --- a/parakeet/src/xrpc/app_bsky/notification/mod.rs +++ b/parakeet/src/xrpc/app_bsky/notification/mod.rs @@ -146,7 +146,7 @@ pub async fn list_notifications( let profile_state = profile_states.get(¬if.author_did); // Build associated object - let associated = if let Some((_, _, _, chat_decl, _, _, notif_decl)) = profile_data { + let associated = if let Some((_, _, _, _, chat_decl, _, _, notif_decl)) = profile_data { let mut assoc = serde_json::Map::new(); if let Some(chat) = chat_decl { @@ -224,15 +224,15 @@ pub async fn list_notifications( let author = NotificationAuthor { did: notif.author_did.clone(), handle: profile_data - .and_then(|(_, h, _, _, _, _, _)| h.clone()) + .and_then(|(_, h, _, _, _, _, _, _)| h.clone()) .unwrap_or_else(|| { // Fallback to handle.invalid if profile not yet indexed "handle.invalid".to_string() }), - display_name: profile_data.and_then(|(_, _, p, _, _, _, _)| { + display_name: profile_data.and_then(|(_, _, _, p, _, _, _, _)| { p.as_ref().and_then(|prof| prof.display_name.clone()) }), - avatar: profile_data.and_then(|(did, _, p, _, _, _, _)| { + avatar: profile_data.and_then(|(did, _, _, p, _, _, _, _)| { p.as_ref().and_then(|prof| { prof.avatar_cid.as_ref().and_then(|cid_digest| { parakeet_db::cid_util::digest_to_blob_cid_string(cid_digest) @@ -243,7 +243,7 @@ pub async fn list_notifications( associated, viewer, labels: Some(Vec::new()), // Empty labels array for compatibility - description: profile_data.and_then(|(_, _, p, _, _, _, _)| { + description: profile_data.and_then(|(_, _, _, p, _, _, _, _)| { p.as_ref().and_then(|prof| prof.description.clone()) }), // TODO: Profile timestamps should be fetched from records_literal table via record_id FK diff --git a/parakeet/src/xrpc/app_bsky/unspecced/mod.rs b/parakeet/src/xrpc/app_bsky/unspecced/mod.rs index 5b82ed6a..0ba63f64 100644 --- a/parakeet/src/xrpc/app_bsky/unspecced/mod.rs +++ b/parakeet/src/xrpc/app_bsky/unspecced/mod.rs @@ -256,7 +256,7 @@ pub async fn get_suggested_users( // Check if has profile (display name or description) profiles .get(did) - .and_then(|p| p.2.as_ref()) + .and_then(|p| p.3.as_ref()) // .3 is Option .map(|prof| prof.display_name.is_some() || prof.description.is_some()) .unwrap_or(false) }) -- 2.51.2