diff --git a/crates/appview/src/firehose.rs b/crates/appview/src/firehose.rs index e08de65..5c9703c 100644 --- a/crates/appview/src/firehose.rs +++ b/crates/appview/src/firehose.rs @@ -10,6 +10,7 @@ use anyhow::Result; use at_firehose::JetstreamEvent; +use serde_json::Value; use sqlx::PgPool; use std::sync::atomic::{AtomicBool, AtomicI64, AtomicU64, Ordering}; use std::sync::Arc; @@ -133,13 +134,17 @@ impl IndexHandler { } }, "identity" => { - trace!(did = %ev.did, "identity event (logged only)"); - let _ = handle_identity(&ev); + if let Err(e) = handle_identity(&self.db, &ev).await { + warn!(error = %e, did = %ev.did, "handle_identity failed"); + return Ok(()); // don't advance cursor; let next replay retry + } true } "account" => { - trace!(did = %ev.did, "account event (logged only)"); - let _ = handle_account(&ev); + if let Err(e) = handle_account(&self.db, &ev).await { + warn!(error = %e, did = %ev.did, "handle_account failed"); + return Ok(()); // don't advance cursor; let next replay retry + } true } other => { @@ -164,16 +169,67 @@ impl IndexHandler { } } -fn handle_identity(_ev: &JetstreamEvent) -> Result<()> { - info!("identity change (DID doc rotation)"); +/// `identity` event — Jetstream tells us a DID's handle changed. +/// +/// The Jetstream payload includes `identity.handle` (the *current* +/// handle, since the event fires after every handle change) and +/// optionally `identity.did` (the DID — redundant with the outer +/// `ev.did` but we accept both). We pull the handle out and run it +/// through `indexer::backfill_handle` so every existing post row for +/// that DID gets the new value. The COALESCE guard inside +/// `indexer::PostRow::from_record` keeps empty strings from +/// clobbering this backfilled value when a later `commit` event +/// arrives. +async fn handle_identity(db: &PgPool, ev: &JetstreamEvent) -> Result<()> { + let handle = extract_handle(&ev.identity).or_else(|| extract_handle(&ev.account)); + let Some(handle) = handle else { + // Some identity events carry only a DID-doc rotation signal + // with no handle payload — those are uninteresting for our + // purpose. Advance the cursor anyway. + debug!(did = %ev.did, "identity event without a usable handle payload"); + return Ok(()); + }; + let rows = indexer::backfill_handle(db, &ev.did, &handle).await?; + info!( + did = %ev.did, + handle = %handle, + rows_updated = rows, + "backfilled handle on posts" + ); Ok(()) } -fn handle_account(_ev: &JetstreamEvent) -> Result<()> { - info!("account change (active/-status)"); +/// `account` event — Jetstream tells us an account's active/deactive +/// status changed. We mirror the handle-backfill behaviour in case +/// the `account` payload carries the verified handle alongside +/// `active`; many real-world identities show the handle there even +/// when no `identity` event was emitted. +async fn handle_account(db: &PgPool, ev: &JetstreamEvent) -> Result<()> { + let Some(handle) = extract_handle(&ev.account) else { + return Ok(()); + }; + let rows = indexer::backfill_handle(db, &ev.did, &handle).await?; + info!( + did = %ev.did, + handle = %handle, + rows_updated = rows, + "backfilled handle on posts (account event)" + ); Ok(()) } +/// Pull a handle string out of a Jetstream event fragment. Returns +/// `None` if the fragment is absent or doesn't carry a usable +/// `handle` string field. +fn extract_handle(fragment: &Option) -> Option { + fragment + .as_ref() + .and_then(|v| v.get("handle")) + .and_then(|v| v.as_str()) + .map(|s| s.trim().to_string()) + .filter(|s| !s.is_empty()) +} + /// Spawn the background task that drains the cursor-flush channel and /// writes the running maximum to the DB. Returns when the receiver is /// dropped (i.e. the main process is shutting down).