From 4a7b954e7726ffbd7cce7ad73ab53aa5dde58668 Mon Sep 17 00:00:00 2001 From: dawn <90008@klbr.net> Date: Sat, 12 Sep 2026 05:50:42 +0300 Subject: [PATCH] [backfill] announce identity changes, not just account changes --- src/backfill/worker/process.rs | 49 +++++++++++++----- src/types.rs | 95 ++++++++++++++++++++++++++++++---- 2 files changed, 120 insertions(+), 24 deletions(-) diff --git a/src/backfill/worker/process.rs b/src/backfill/worker/process.rs index 1f78f3a..c65d5ea 100644 --- a/src/backfill/worker/process.rs +++ b/src/backfill/worker/process.rs @@ -31,7 +31,9 @@ use crate::types::{Commit, GaugeState, RepoState, RepoStatus, ResyncState}; use crate::util::{parse_retry_after, throttle::ThrottleHandle}; #[cfg(feature = "indexer_stream")] -use crate::types::{AccountEvt, BroadcastEvent}; +use crate::types::{AccountEvt, BroadcastEvent, IdentityEvt}; +#[cfg(feature = "indexer_stream")] +use jacquard_common::types::string::Handle; #[cfg(feature = "indexer_stream")] use std::sync::atomic::Ordering; @@ -82,7 +84,7 @@ pub(super) async fn process_did( }; #[cfg(feature = "indexer_stream")] - let emit_identity = |status: &RepoStatus, active: bool| { + let emit_account = |status: &RepoStatus, active: bool| { let status = match status { RepoStatus::Deactivated => "deactivated", RepoStatus::Takendown => "takendown", @@ -104,6 +106,22 @@ pub(super) async fn process_did( .send(ops::make_account_event(db, evt)); }; + // an identity announcement is ephemeral, and only consumers that already + // hold a handle for this repo can act on it: they swap the handle they + // cached instead of waiting for a resolve of their own. + #[cfg(feature = "indexer_stream")] + let emit_identity = |handle: Option| { + let evt = IdentityEvt { + did: did.clone(), + handle, + }; + let _ = app_state + .db + .stream + .event_tx + .send(ops::make_identity_event(db, evt)); + }; + if strategy != BackfillStrategy::Full { let filter = app_state.filter.load(); let sparse_supported = !filter.collections.is_empty() @@ -122,11 +140,12 @@ pub(super) async fn process_did( { Ok(SparseBackfillResult::Imported(sparse)) => { #[cfg(feature = "indexer_stream")] - if sparse.state.active != previous_state.active - || sparse.state.status != previous_state.status - || previous_state.pds.is_none() - { - emit_identity(&sparse.state.status, sparse.state.active); + if sparse.state.account_moved_from(&previous_state) { + emit_account(&sparse.state.status, sparse.state.active); + } + #[cfg(feature = "indexer_stream")] + if sparse.state.identity_moved_from(&previous_state) { + emit_identity(sparse.state.handle.clone()); } trace!( @@ -198,7 +217,7 @@ pub(super) async fn process_did( warn!(?status, "repo is inactive, stopping backfill"); #[cfg(feature = "indexer_stream")] - emit_identity(&status, false); + emit_account(&status, false); let resync_state = ResyncState::Gone { status: status.clone(), @@ -244,13 +263,15 @@ pub(super) async fn process_did( }; drop(admission_permit); - // emit identity event so any consumers know, but only if something changed + // tell consumers what moved: the account lifecycle, and the identity we + // just re-resolved. either may have changed without the other moving. #[cfg(feature = "indexer_stream")] - if state.active != previous_state.active - || state.status != previous_state.status - || previous_state.pds.is_none() - { - emit_identity(&state.status, state.active); + if state.account_moved_from(&previous_state) { + emit_account(&state.status, state.active); + } + #[cfg(feature = "indexer_stream")] + if state.identity_moved_from(&previous_state) { + emit_identity(state.handle.clone()); } trace!( diff --git a/src/types.rs b/src/types.rs index 060e65e..5ae6e7f 100644 --- a/src/types.rs +++ b/src/types.rs @@ -1,13 +1,13 @@ use std::fmt::Display; -use jacquard_common::types::string::Did; +use jacquard_common::types::string::{Did, Handle}; use jacquard_common::{CowStr, IntoStatic}; use jacquard_repo::commit::Commit as AtpCommit; #[cfg(feature = "indexer")] use serde::{Deserialize, Serialize}; use smol_str::ToSmolStr; -use crate::db::types::DbTid; +use crate::db::types::{DbTid, DidKey}; use crate::resolver::MiniDoc; #[cfg(any(feature = "indexer_stream", feature = "relay", feature = "jetstream"))] @@ -158,6 +158,16 @@ mod indexer { #[cfg(feature = "indexer")] pub(crate) use indexer::*; +/// the identity a DID doc owns. it is the only part of a repo state that +/// resolving a doc rewrites, so consumers that cache it have to be told when it +/// moves instead of discovering it later. +#[derive(Debug, Clone, PartialEq)] +pub(crate) struct DocIdentity<'i> { + handle: Option, + pds: Option>, + signing_key: Option>, +} + impl<'i> RepoState<'i> { pub fn backfilling() -> Self { Self { @@ -210,15 +220,34 @@ impl<'i> RepoState<'i> { self.last_updated_at = chrono::Utc::now().timestamp(); } + /// the doc-owned identity, the whole of what a resolve can rewrite. + pub(crate) fn identity(&self) -> DocIdentity<'i> { + DocIdentity { + handle: self.handle.clone(), + pds: self.pds.clone(), + signing_key: self.signing_key.clone(), + } + } + + /// the identity moved, so anyone holding a cached copy must re-resolve it. + #[cfg(any(test, feature = "indexer_stream"))] + pub(crate) fn identity_moved_from(&self, previous: &Self) -> bool { + self.identity() != previous.identity() + } + + /// the account lifecycle moved. a first sighting of a pds counts as a move: + /// nothing downstream keyed off an identity it never saw. + #[cfg(any(test, feature = "indexer_stream"))] + pub(crate) fn account_moved_from(&self, previous: &Self) -> bool { + self.active != previous.active || self.status != previous.status || previous.pds.is_none() + } + pub fn update_from_doc(&mut self, doc: MiniDoc) -> bool { - let new_signing_key = doc.key.map(From::from); - let changed = self.pds.as_deref() != Some(doc.pds.as_str()) - || self.handle != doc.handle - || self.signing_key != new_signing_key; + let previous = self.identity(); self.pds = Some(CowStr::Owned(doc.pds.to_smolstr())); self.handle = doc.handle; - self.signing_key = new_signing_key; - changed + self.signing_key = doc.key.map(From::from); + self.identity() != previous } } @@ -265,10 +294,56 @@ pub(crate) enum ResyncState { #[cfg(test)] mod tests { use super::*; - use crate::db::types::DidKey; - use jacquard_common::types::string::Handle; use miette::IntoDiagnostic; + #[test] + fn identity_moves_only_for_doc_owned_fields() -> miette::Result<()> { + let mut previous = RepoState::synced(); + previous.pds = Some(CowStr::Borrowed("https://pds.example")); + previous.handle = Some(Handle::new_static("alice.test").into_diagnostic()?); + + let mut state = previous.clone(); + state.advance_message_time(1_000); + state.touch(); + state.root = previous.root.clone(); + assert!(!state.identity_moved_from(&previous)); + assert!(!state.account_moved_from(&previous)); + + let mut handle_moved = previous.clone(); + handle_moved.handle = Some(Handle::new_static("carol.test").into_diagnostic()?); + assert!(handle_moved.identity_moved_from(&previous)); + assert!(!handle_moved.account_moved_from(&previous)); + + let mut pds_moved = previous.clone(); + pds_moved.pds = Some(CowStr::Borrowed("https://other.example")); + assert!(pds_moved.identity_moved_from(&previous)); + + let mut key_moved = previous.clone(); + key_moved.signing_key = Some(DidKey::from_did_key( + "did:key:zQ3shokFTS3brHcDQrn82RUDfCZESWL1ZdCEJwekUDPQiYBme", + )?); + assert!(key_moved.identity_moved_from(&previous)); + + let mut account_moved = previous.clone(); + account_moved.status = RepoStatus::Throttled; + assert!(account_moved.account_moved_from(&previous)); + assert!(!account_moved.identity_moved_from(&previous)); + + Ok(()) + } + + #[test] + fn first_sighting_of_a_pds_moves_the_account() { + let mut previous = RepoState::backfilling(); + previous.pds = None; + + let mut state = previous.clone(); + state.pds = Some(CowStr::Borrowed("https://pds.example")); + + assert!(state.account_moved_from(&previous)); + assert!(state.identity_moved_from(&previous)); + } + #[test] fn identity_dedupe_does_not_depend_on_commit_clock() { let mut state = RepoState::backfilling(); -- 2.51.2