diff --git a/hubble-sync/src/identity/mod.rs b/hubble-sync/src/identity/mod.rs index e721adb..b9ba301 100644 --- a/hubble-sync/src/identity/mod.rs +++ b/hubble-sync/src/identity/mod.rs @@ -1,7 +1,9 @@ mod did; +mod pending_backoff; mod resolver; mod validity; pub use did::{Did, DidError, DidMethod}; +pub use pending_backoff::pending_resolve_backoff; pub use resolver::{ResolutionError, ResolvedIdentity, Resolver}; pub use validity::Validity; diff --git a/hubble-sync/src/identity/pending_backoff.rs b/hubble-sync/src/identity/pending_backoff.rs new file mode 100644 index 0000000..de24505 --- /dev/null +++ b/hubble-sync/src/identity/pending_backoff.rs @@ -0,0 +1,30 @@ +//! failure stuff for when we encounter a new DID we'd like to sync + +use crate::DidMethod; +use std::time::Duration; + +/// TODO: i'm pretty sure we should, eventually, stop trying (return option?) +pub fn pending_resolve_backoff(method: DidMethod, attempts: u32) -> Duration { + let seconds = match method { + // plc errors are weird. directory down? + // + // TODO: verify never-heard-of vs deleted (which might get undone) + DidMethod::Plc => match attempts { + 0 => 20, + 1 => 40, + 2 => 100, + 3..=6 => 600, + _ => 3600, // for now we never stop trying, every 1h, forever + }, + // did:web errors are less weird! servers go down all the time! + // + // backoff schedule here is pretty arbitrary + DidMethod::Web => match attempts { + 0..=1 => 60, + 2..=4 => 600, + 5..=10 => 3600, + _ => 86_400, // once per day forevermore + }, + }; + Duration::from_secs(seconds) +} diff --git a/hubble-sync/src/identity/resolver.rs b/hubble-sync/src/identity/resolver.rs index 8e44421..e2fc31d 100644 --- a/hubble-sync/src/identity/resolver.rs +++ b/hubble-sync/src/identity/resolver.rs @@ -8,7 +8,7 @@ use jacquard_identity::resolver::{ }; use reqwest::StatusCode; -use crate::storage::repo_state::RepoIdentity; +use crate::storage::repo::RepoIdentity; use crate::{Did, DidMethod, HostRegistry}; const PLC_RETRY_AFTER: Duration = Duration::from_secs(3); diff --git a/hubble-sync/src/identity/validity.rs b/hubble-sync/src/identity/validity.rs index be70bd8..3f01a58 100644 --- a/hubble-sync/src/identity/validity.rs +++ b/hubble-sync/src/identity/validity.rs @@ -2,7 +2,7 @@ use std::time::Duration; -use crate::storage::repo_state::RepoIdentity; +use crate::storage::repo::RepoIdentity; use crate::{Did, DidMethod}; /// no refresh triggered for this period. @@ -14,12 +14,6 @@ const DID_PLC_EXPIRE_AT: Duration = Duration::from_secs(3 * 86_400); /// no refresh triggered for this period. did:webs never hard-expire. const DID_WEB_VALID_FOR: Duration = Duration::from_secs(86_400); -#[derive(Debug, thiserror::Error)] -pub enum ValiditityError { - #[error("unknown did method: {0:?}")] - UnknownMethod(Did), -} - pub enum Validity { Valid, Stale, @@ -27,14 +21,14 @@ pub enum Validity { } impl Validity { - pub fn for_identity(did: &Did, identity: &RepoIdentity) -> Result { + pub fn for_identity(did: &Did, identity: &RepoIdentity) -> Validity { let age = identity.resolved_at.elapsed().unwrap_or(Duration::ZERO); match did.method() { - DidMethod::Plc if age < DID_PLC_VALID_FOR => Ok(Self::Valid), - DidMethod::Plc if age < DID_PLC_EXPIRE_AT => Ok(Self::Stale), - DidMethod::Plc => Ok(Self::Expired), - DidMethod::Web if age < DID_WEB_VALID_FOR => Ok(Self::Valid), - DidMethod::Web => Ok(Self::Stale), + DidMethod::Plc if age < DID_PLC_VALID_FOR => Self::Valid, + DidMethod::Plc if age < DID_PLC_EXPIRE_AT => Self::Stale, + DidMethod::Plc => Self::Expired, + DidMethod::Web if age < DID_WEB_VALID_FOR => Self::Valid, + DidMethod::Web => Self::Stale, } } @@ -44,4 +38,11 @@ impl Validity { Self::Stale | Self::Expired => true, } } + + pub fn must_refresh(&self) -> bool { + match self { + Self::Valid | Self::Stale => false, + Self::Expired => true, + } + } } diff --git a/hubble-sync/src/lib.rs b/hubble-sync/src/lib.rs index 752eaee..52faee0 100644 --- a/hubble-sync/src/lib.rs +++ b/hubble-sync/src/lib.rs @@ -4,6 +4,7 @@ mod firehose; mod host; mod identity; mod metrics; +mod pending_identity_scheduler; mod repo_actor; mod resync_scheduler; mod storage; @@ -20,7 +21,6 @@ pub use identity::{Did, DidMethod}; pub use metrics::describe_metrics; pub use repo_actor::{RepoContext, RepoRegistry, RepoSendError, RepoTask}; pub use storage::engine::{StorageBatch, StorageEngine, StorageError}; -pub use storage::repo::Repo; -pub use storage::repo_state::AccountStatus; +pub use storage::repo::{AccountStatus, Repo}; pub use sync_consumer::{AppResult, ConsumerAppError, SyncConsumer}; pub use tid::Tid; diff --git a/hubble-sync/src/metrics.rs b/hubble-sync/src/metrics.rs index 8b8540a..55a58bd 100644 --- a/hubble-sync/src/metrics.rs +++ b/hubble-sync/src/metrics.rs @@ -47,6 +47,13 @@ pub(crate) const RESYNC_SCHEDULER_DISPATCHED_TOTAL: &str = /// hosts currently scheduled in the round-robin pub(crate) const RESYNC_SCHEDULER_SCHEDULED: &str = "hubble_sync_resync_scheduler_scheduled"; +// ///// pending identity resolution + +pub(crate) const PENDING_SCHEDULER_DISPATCHED_TOTAL: &str = + "hubble_sync_pending_scheduler_dispatched_total"; + +pub(crate) const PENDING_SCHEDULER_SCHEDULED: &str = "hubble_sync_pending_scheduler_scheduled"; + // ///// firehose /// usable subscribeRepos messages. `kind` = `commit` | `sync` | `identity` | `account` | `info` | `unknown` @@ -119,11 +126,21 @@ pub fn describe_metrics() { ); describe_counter!( RESYNC_SCHEDULER_DISPATCHED_TOTAL, - "in-memory scheduler dispatches" + "in-memory resync scheduler dispatches" ); describe_gauge!( RESYNC_SCHEDULER_SCHEDULED, - "hosts currently scheduled in the round-robin" + "hosts currently scheduled in the resync round-robin" + ); + + // pending identity resolution retries + describe_counter!( + PENDING_SCHEDULER_DISPATCHED_TOTAL, + "in-memory identity resolution retry dispatches" + ); + describe_gauge!( + PENDING_SCHEDULER_SCHEDULED, + "identities currently scheduled in the retry queue" ); // firehose diff --git a/hubble-sync/src/pending_identity_scheduler.rs b/hubble-sync/src/pending_identity_scheduler.rs new file mode 100644 index 0000000..740e44c --- /dev/null +++ b/hubble-sync/src/pending_identity_scheduler.rs @@ -0,0 +1,105 @@ +//! in-memory scheduler over the pending identity resolution retry queue +//! +//! the scheduler is typically pretty quiet, since identity resolution usually +//! succeeds. + +use std::collections::{BTreeSet, HashMap}; +use std::fmt; +use std::sync::Mutex; +use std::time::{Duration, SystemTime}; + +use metrics::{counter, gauge}; + +use crate::Did; +use crate::metrics::{PENDING_SCHEDULER_DISPATCHED_TOTAL, PENDING_SCHEDULER_SCHEDULED}; +use crate::storage::{LoadError, StorageEngine, repo::PendingIdentityQueueEntry}; + +#[derive(Default)] +pub struct PendingScheduler { + inner: Mutex, +} + +#[derive(Default)] +struct Schedule { + ordered: BTreeSet<(SystemTime, Did)>, + by_did: HashMap, +} + +impl PendingScheduler { + pub fn new() -> Self { + Self::default() + } + + fn rm(s: &mut Schedule, did: &Did) { + if let Some(t) = s.by_did.remove(did) { + s.ordered.remove(&(t, did.clone())); + } + } + + pub fn upsert(&self, did: Did, retry_at: SystemTime) { + let mut s = self.inner.lock().expect("pending scheduler lock"); + Self::rm(&mut s, &did); + s.ordered.insert((retry_at, did.clone())); + s.by_did.insert(did, retry_at); + gauge!(PENDING_SCHEDULER_SCHEDULED).set(s.by_did.len() as f64); + } + + pub fn remove(&self, did: &Did) { + let mut s = self.inner.lock().expect("pending scheduler lock"); + Self::rm(&mut s, did); + gauge!(PENDING_SCHEDULER_SCHEDULED).set(s.by_did.len() as f64); + } + + pub fn next(&self, now: SystemTime) -> Option { + let mut s = self.inner.lock().expect("pending scheduler lock"); + let (_, did) = s.ordered.pop_first()?; + s.by_did.remove(&did); + gauge!(PENDING_SCHEDULER_SCHEDULED).set(s.by_did.len() as f64); + counter!(PENDING_SCHEDULER_DISPATCHED_TOTAL).increment(1); + Some(did) + } + + pub fn wait(&self, now: SystemTime) -> Option { + let s = self.inner.lock().expect("pending scheduler lock"); + let (t, _) = s.ordered.first()?; + Some(t.duration_since(now).unwrap_or(Duration::ZERO)) + } + + pub fn len(&self) -> usize { + self.inner + .lock() + .expect("pending scheduler lock") + .by_did + .len() + } + + pub fn is_empty(&self) -> bool { + self.inner + .lock() + .expect("pending scheduler lock") + .by_did + .is_empty() + } +} + +impl fmt::Debug for PendingScheduler { + fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result { + f.debug_struct("PendingScheduler") + .field("size", &self.inner.try_lock().map(|g| g.by_did.len())) + .finish_non_exhaustive() + } +} + +/// Initial in-memory hydration from the db +pub fn bootstrap( + scheduler: &PendingScheduler, + engine: &S, +) -> Result> { + let mut count = 0; + for entry in PendingIdentityQueueEntry::scan(engine) { + let entry = entry?; + scheduler.upsert(entry.did, entry.due); + count += 1; + } + Ok(count) +} diff --git a/hubble-sync/src/repo_actor/actor.rs b/hubble-sync/src/repo_actor/actor.rs index faa9423..5b67396 100644 --- a/hubble-sync/src/repo_actor/actor.rs +++ b/hubble-sync/src/repo_actor/actor.rs @@ -56,12 +56,13 @@ use std::time::{Instant, SystemTime}; use tokio::sync::mpsc::{self, error::TrySendError}; use tokio::sync::oneshot; +use tokio::task::spawn_blocking; -use super::{EvictableState, Process, ProcessError}; +use super::{Bootstrap, BootstrapOutcome, EvictableState, Process, ProcessError, Refresh}; use crate::firehose::FirehoseEvent; use crate::identity::{ResolvedIdentity, Resolver}; -use crate::storage::repo::Repo; -use crate::storage::{LoadError, StorageEngine}; +use crate::storage::StorageEngine; +use crate::storage::repo::{Awoken, Repo}; use crate::{AccountStatus, Did, HostRegistry, SyncConsumer}; /// some work the system can send into the repo actor @@ -225,48 +226,95 @@ impl> RepoActor { // TODO: would there be any value in spawning actors in one big joinset? // could be awaited (with timeout) on shutdown etc. tokio::spawn(async move { - let repo = Self::hydrate(did.clone(), storage.clone(), hosts.clone(), now) - .await - .expect("fallible hydration is todo! how do we indicate that back?"); - - let mut intake = Intake::new(rx, queue_limit); - - if repo.needs_identity_refresh() { - intake.intake(Task::ResolveIdentity { reply: None }); - } + let awoken = { + let storage = storage.clone(); + let hosts = hosts.clone(); + let did = did.clone(); + let Ok(awoken) = spawn_blocking(move || Repo::wake(&*storage, &hosts, did, now)) + .await + .expect("repo wake task not to panic") + .inspect_err(|e| tracing::error!(?e, "wake failed")) + else { + return; + }; + awoken + }; - let process = Process { - storage, - hosts, - resolver, - consumer_app, - revived_at, - repo: Some(repo), + let intake = Intake::new(rx, queue_limit); + + let process = match awoken { + Awoken::Resolved(repo) => { + // TODO: might still use something here for stale not expired + // if repo.needs_identity_refresh() { + // intake.intake(Task::ResolveIdentity { reply: None }); + // } + Process { + storage, + hosts, + resolver, + consumer_app, + revived_at, + repo: Some(repo), + } + } + Awoken::Pending(pending) => { + let Ok(outcome) = Bootstrap { + storage: storage.clone(), + hosts: hosts.clone(), + resolver: resolver.clone(), + pending, + } + .run(consumer_app.clone(), revived_at, now) + .await + .inspect_err(|e| tracing::error!(?e, "bootstrap storage error")) else { + // TODO: a storage error should be fatal + return; + }; + match outcome { + BootstrapOutcome::GetUp(p) => p, + BootstrapOutcome::Snooze | BootstrapOutcome::Retry => return, + } + } + Awoken::Refreshing { repo, pending } => { + let mut process = Process { + storage: storage.clone(), + hosts: hosts.clone(), + resolver: resolver.clone(), + consumer_app: consumer_app.clone(), + revived_at, + repo: Some(repo), + }; + // attempt refresh: failure leaves it stale + let repo_ref = process.repo.as_mut().expect("repo present"); + let Ok(_) = Refresh { + storage: storage.clone(), + hosts: hosts.clone(), + resolver: resolver.clone(), + pending, + } + .run(repo_ref, now) + .await + .inspect_err(|e| tracing::error!(?e, "refresh storage error")) else { + // TODO: a storage error should be fatal + return; + }; + process + } }; - let actor = Self { + let _ = Self { intake, process, state, exit_guard, - }; - actor.run().await + } + .run() + .await + .inspect_err(|e| tracing::error!("maybe should be panic? {e}")); }); lifeline } - /// load existing state from the database - async fn hydrate( - did: Did, - storage: Arc, - hosts: Arc, - now: SystemTime, - ) -> Result> { - tokio::task::spawn_blocking(move || Repo::load_or_create(&*storage, &hosts, did, now)) - .await - .expect("spawn_blocking") - } - async fn run(mut self) -> Result<(), ProcessError> { while self.run_step().await? {} while let Some(t) = self.intake.queue.pop_front() { @@ -347,7 +395,7 @@ mod tests { use crate::RepoContext; use crate::StorageBatch; use crate::storage::engine::mem::MemEngine; - use crate::storage::repo_state::{RepoIdentity, RepoInfo, SyncStatus}; + use crate::storage::repo::{RepoIdentity, RepoInfo, SyncStatus}; use crate::sync_consumer::AppResult; /// No-op consumer. The actor's `Process::process` currently @@ -412,12 +460,12 @@ mod tests { let info = RepoInfo { sync_status: SyncStatus::Synchronized, account_status: AccountStatus::Active, - identity: Some(RepoIdentity { + identity: RepoIdentity { pds_host: pds, signing_key: vec![0u8; 32], supposed_handle: "alice.example".to_string(), resolved_at: now, - }), + }, first_seen_at: now, resyncs: 0, pds_changes: 0, diff --git a/hubble-sync/src/repo_actor/bootstrap.rs b/hubble-sync/src/repo_actor/bootstrap.rs new file mode 100644 index 0000000..e7999ca --- /dev/null +++ b/hubble-sync/src/repo_actor/bootstrap.rs @@ -0,0 +1,95 @@ +//! when we see brand new DIDs, we need to resolve their identity etc. +//! +//! when waking the repo produces `Awoken::Pending`, we end up here, before +//! it can go to its normal task loop. + +use std::sync::Arc; +use std::time::{Instant, SystemTime}; + +use tokio::task::spawn_blocking; + +use crate::identity::Resolver; +use crate::storage::repo::{PendingIdentity, Repo}; +use crate::storage::{StorageEngine, engine::StorageBatch}; +use crate::{HostRegistry, SyncConsumer}; + +use super::process_task::{Process, ProcessError}; + +pub(super) struct Bootstrap { + pub(super) storage: Arc, + pub(super) hosts: Arc, + pub(super) resolver: Arc, + pub(super) pending: PendingIdentity, +} + +pub(super) enum BootstrapOutcome> { + /// not time to retry yet, nothing happened + Snooze, + /// tried and failed and should try again + Retry, + /// aw yea we are up! + GetUp(Process), +} + +impl Bootstrap { + pub(super) async fn run>( + self, + consumer_app: Arc, + revived_at: Instant, + now: SystemTime, + ) -> Result, ProcessError> { + let Self { + storage, + hosts, + resolver, + mut pending, + } = self; + + if !pending.is_ready(now) { + return Ok(BootstrapOutcome::Snooze); + } + + let did = pending.did.clone(); + + match resolver.resolve(&did, &hosts, now).await { + Ok(identity) => { + let storage_inner = storage.clone(); + let did = did.clone(); + + let repo = spawn_blocking(move || { + let mut batch = storage_inner.batch(); + let repo = Repo::create_resolved(did.clone(), identity, now, &mut batch); + pending.delete(&mut batch); + batch.commit()?; + Ok(repo) + }) + .await + .expect("storage not to panic") + .map_err(ProcessError::Storage)?; + + Ok(BootstrapOutcome::GetUp(Process { + storage, + hosts, + resolver, + consumer_app, + revived_at, + repo: Some(repo), + })) + } + Err(err) => { + // TODO: handle different kinds of identity error?? + let storage = storage.clone(); + spawn_blocking(move || { + let mut batch = storage.batch(); + pending.store_failed(&format!("failed to resolve: {err}"), now, &mut batch); + batch.commit()?; + Ok(()) + }) + .await + .expect("storage not to panic") + .map_err(ProcessError::Storage)?; + Ok(BootstrapOutcome::Retry) + } + } + } +} diff --git a/hubble-sync/src/repo_actor/mod.rs b/hubble-sync/src/repo_actor/mod.rs index 477111c..111c32f 100644 --- a/hubble-sync/src/repo_actor/mod.rs +++ b/hubble-sync/src/repo_actor/mod.rs @@ -1,8 +1,10 @@ //! Repository actor, RepoState, RepoTask mod actor; +mod bootstrap; mod evictable_state; mod process_task; +mod refresh_identity; mod repo_registry; pub use actor::RepoTask; @@ -10,5 +12,7 @@ pub use process_task::RepoContext; pub use repo_registry::{RepoRegistry, RepoSendError}; use actor::Task; +use bootstrap::{Bootstrap, BootstrapOutcome}; use evictable_state::EvictableState; use process_task::{Process, ProcessError}; +use refresh_identity::Refresh; diff --git a/hubble-sync/src/repo_actor/process_task.rs b/hubble-sync/src/repo_actor/process_task.rs index db8afb8..d929f6f 100644 --- a/hubble-sync/src/repo_actor/process_task.rs +++ b/hubble-sync/src/repo_actor/process_task.rs @@ -6,12 +6,10 @@ use tokio::task::spawn_blocking; use super::Task; use crate::HostRegistry; -use crate::identity::{ResolutionError, ResolvedIdentity, Resolver}; +use crate::identity::{ResolvedIdentity, Resolver}; use crate::storage::StorageEngine; use crate::storage::engine::StorageBatch; -use crate::storage::repo::Repo; -use crate::storage::repo::repo_state_account::{AccountStatus, RepoIdentity}; -use crate::storage::repo::repo_state_account_desync::DesyncReason; +use crate::storage::repo::{AccountStatus, Repo, RepoIdentity}; use crate::{ConsumerAppError, Did, Host, StorageError, SyncConsumer, Tid}; #[derive(Debug, thiserror::Error)] @@ -35,8 +33,8 @@ impl<'a> RepoContext<'a> { pub fn status(&self) -> &AccountStatus { &self.0.info().account_status } - pub fn pds(&self) -> Option<&Arc> { - self.0.info().identity.as_ref().map(|i| &i.pds_host) + pub fn pds(&self) -> &Arc { + &self.0.info().identity.pds_host } pub fn is_active(&self) -> bool { matches!(self.0.info().account_status, AccountStatus::Active) @@ -118,7 +116,7 @@ impl> Process { let updated_repo = match outcome { Ok(ref identity) => self.commit_resolved(repo, identity.clone()).await, - Err(ref err) => self.commit_resolved_fail(repo, err.clone()).await, + Err(ref err) => todo!("handle identity resolution failure from task: {err}"), }?; self.repo = Some(updated_repo); @@ -147,23 +145,4 @@ impl> Process { Ok(updated_repo) } - - async fn commit_resolved_fail( - &mut self, - mut repo: Repo, - err: ResolutionError, - ) -> Result { - tracing::info!(?err, "scheduling retry for failed identity resolve"); - let storage = self.storage.clone(); - let now = SystemTime::now(); - let updated_repo = spawn_blocking(move || { - let mut batch = storage.batch(); - repo.desynchronize(DesyncReason::UnresolvableIdentity, now, &mut batch); - batch.commit().map_err(ProcessError::Storage)?; - Ok(repo) - }) - .await - .expect("update not to panic")?; - Ok(updated_repo) - } } diff --git a/hubble-sync/src/repo_actor/refresh_identity.rs b/hubble-sync/src/repo_actor/refresh_identity.rs new file mode 100644 index 0000000..57ba887 --- /dev/null +++ b/hubble-sync/src/repo_actor/refresh_identity.rs @@ -0,0 +1,80 @@ +use std::sync::Arc; +use std::time::SystemTime; + +use tokio::task::spawn_blocking; + +use super::ProcessError; +use crate::identity::Resolver; +use crate::storage::repo::{PendingIdentity, Repo}; +use crate::{HostRegistry, StorageBatch, StorageEngine}; + +pub(super) struct Refresh { + pub(super) storage: Arc, + pub(super) hosts: Arc, + pub(super) resolver: Arc, + pub(super) pending: PendingIdentity, +} + +pub(super) enum RefreshOutcome { + Snoozed, + Failed, + Refreshed, +} + +impl Refresh { + pub(super) async fn run( + self, + repo: &mut Repo, + now: SystemTime, + ) -> Result> { + let Self { + storage, + hosts, + resolver, + mut pending, + } = self; + + if !pending.is_ready(now) { + return Ok(RefreshOutcome::Snoozed); + } + + let did = repo.did().clone(); + + match resolver.resolve(&did, &hosts, now).await { + Ok(identity) => { + let mut repo_inner = repo.clone(); + let storage_inner = storage.clone(); + + *repo = spawn_blocking(move || { + let mut batch = storage_inner.batch(); + repo_inner.set_identity(identity, &mut batch); + pending.delete(&mut batch); // TODO is this actually noop for in-mem pending or are we dropping a tombstone?? + batch.commit()?; + Ok(repo_inner) + }) + .await + .expect("storage not to panic") + .map_err(ProcessError::Storage)?; + + Ok(RefreshOutcome::Refreshed) + } + Err(err) => { + let mut repo_inner = repo.clone(); + let storage_inner = storage.clone(); + *repo = spawn_blocking(move || { + let mut batch = storage_inner.batch(); + pending.store_failed(&format!("resolve: {err}"), now, &mut batch); + // cross-schedule coupling: if the repo is Desync, align + // its resync due_at with the new pending next_try_at. + repo_inner.delay_resync_until(pending.next_try_at(), &mut batch); + batch.commit()?; + Ok(repo_inner) + }) + .await + .expect("storage not to panic") + .map_err(ProcessError::Storage)?; + Ok(RefreshOutcome::Failed) + } + } + } +} diff --git a/hubble-sync/src/repo_actor/repo_registry.rs b/hubble-sync/src/repo_actor/repo_registry.rs index 50a862c..8bc8f84 100644 --- a/hubble-sync/src/repo_actor/repo_registry.rs +++ b/hubble-sync/src/repo_actor/repo_registry.rs @@ -228,7 +228,7 @@ mod tests { use crate::StorageBatch; use crate::storage::engine::mem::MemEngine; - use crate::storage::repo_state::{RepoIdentity, RepoInfo, SyncStatus}; + use crate::storage::repo::{RepoIdentity, RepoInfo, SyncStatus}; use crate::sync_consumer::AppResult; use crate::{AccountStatus, RepoContext}; @@ -297,12 +297,12 @@ mod tests { let info = RepoInfo { sync_status: SyncStatus::Synchronized, account_status: AccountStatus::Active, - identity: Some(RepoIdentity { + identity: RepoIdentity { pds_host: pds, signing_key: vec![0u8; 32], supposed_handle: "alice.example".to_string(), resolved_at: now, - }), + }, first_seen_at: now, resyncs: 0, pds_changes: 0, diff --git a/hubble-sync/src/resync_scheduler.rs b/hubble-sync/src/resync_scheduler.rs index aa337e2..3ee2287 100644 --- a/hubble-sync/src/resync_scheduler.rs +++ b/hubble-sync/src/resync_scheduler.rs @@ -22,7 +22,8 @@ use metrics::{counter, gauge}; use crate::host::{Host, Hostname}; use crate::metrics::{RESYNC_SCHEDULER_DISPATCHED_TOTAL, RESYNC_SCHEDULER_SCHEDULED}; -use crate::storage::{LoadError, resync_queue}; +use crate::storage::LoadError; +use crate::storage::repo::NextQueuedByHost; use crate::{Did, HostRegistry}; /// in-memory round-robin-by-host resync scheduler @@ -166,7 +167,7 @@ pub fn bootstrap( registry: &HostRegistry, ) -> Result> { let mut count = 0; - for next_queued in resync_queue::NextQueuedByHost::new(engine, registry) { + for next_queued in NextQueuedByHost::new(engine, registry) { scheduler.insert(next_queued?); count += 1; } diff --git a/hubble-sync/src/storage/mod.rs b/hubble-sync/src/storage/mod.rs index 4e6bf4b..0b4b45a 100644 --- a/hubble-sync/src/storage/mod.rs +++ b/hubble-sync/src/storage/mod.rs @@ -8,14 +8,21 @@ use crate::host::HostnameError; use crate::identity::DidError; pub use engine::StorageEngine; -pub use repo::repo_index_resync as resync_queue; -pub use repo::repo_state_account as repo_state; -const PREFIX_FIREHOSE_CURSOR: &[u8] = b"fc/"; -const PREFIX_HOST_INFO: &[u8] = b"hi/"; -const PREFIX_REPO_ACCOUNT_SYNC_STATE: &[u8] = b"rc/"; -const PREFIX_REPO_INFO: &[u8] = b"ri/"; -const PREFIX_RESYNC_QUEUE: &[u8] = b"rq/"; +/// key prefix, type-aliased to make it hard to make a wrong-length mistake +type P = &'static [u8; 3]; + +const PREFIX_FIREHOSE_CURSOR: P = b"fc|"; + +const PREFIX_HOST_INFO: P = b"hi|"; + +const PREFIX_REPO_INFO: P = b"ri|"; +const PREFIX_REPO_INFO_IDX_RESYNC: P = b"rq|"; + +const PREFIX_REPO_PENDING: P = b"pi|"; +const PREFIX_REPO_PENDING_IDX_RETRY: P = b"pq|"; + +const PREFIX_REPO_PREV: P = b"si|"; /// any kind of persistent data loading problem /// @@ -73,7 +80,7 @@ impl From for LoadError { } /// str-like split_once for slices until `slice_split_once` lands -fn split_once<'a, T: PartialEq>(who: &'a [T], what: &'a T) -> Option<(&'a [T], &'a [T])> { +fn slice_split_once<'a, T: PartialEq>(who: &'a [T], what: &'a T) -> Option<(&'a [T], &'a [T])> { let found_index = who.iter().position(|x| x == what)?; Some((&who[..found_index], &who[found_index + 1..])) } diff --git a/hubble-sync/src/storage/repo/repo_state_account.rs b/hubble-sync/src/storage/repo/info.rs similarity index 67% rename from hubble-sync/src/storage/repo/repo_state_account.rs rename to hubble-sync/src/storage/repo/info.rs index 5fed5dc..08dff1a 100644 --- a/hubble-sync/src/storage/repo/repo_state_account.rs +++ b/hubble-sync/src/storage/repo/info.rs @@ -2,34 +2,28 @@ //! //! see `repo` for public-facing repo view of combined state. //! -//! "ri/"|| => +//! "ri|"|| => //! //! todo: rate-limit resyncs per-did (persisted?, maybe?) //! todo: per-upstream account status, and moderated account status //! todo: do we have last-desync-reason kept around? -use crate::storage::repo::repo_state_account_desync::DesyncReason; -use crate::storage::repo::repo_state_account_desync::Desynchronized; -use crate::storage::repo::repo_state_account_desync::ResyncInfo; use std::sync::Arc; -use std::time::SystemTime; +use std::time::{Duration, SystemTime}; use dasl::drisl; use serde::{Deserialize, Serialize}; +use super::{LoadError, PREFIX_REPO_INFO, unix_ms_u64}; use crate::host::{Host, HostnameError}; -use crate::resync_scheduler::QueuedResync; -use crate::{Did, HostRegistry, StorageBatch, StorageEngine, StorageError}; - -use super::unix_ms_u64; -use super::{LoadError, PREFIX_REPO_INFO}; +use crate::{Did, DidMethod, HostRegistry, StorageBatch, StorageEngine, StorageError, Tid}; /// raw repo stuff -- see [`RepoInfo`] for the nice type #[derive(Debug, Clone, Serialize, Deserialize)] struct DbRepoInfo { sync_status: SyncStatus, account_status: AccountStatus, - identity: Option, + identity: DbRepoIdentity, #[serde(with = "unix_ms_u64")] first_seen_at: SystemTime, resyncs: u32, @@ -51,9 +45,9 @@ impl DbRepoInfo { /// nice repo stuff (interned host) #[derive(Debug, Clone)] pub struct RepoInfo { - pub sync_status: SyncStatus, pub account_status: AccountStatus, - pub identity: Option, + pub sync_status: SyncStatus, + pub identity: RepoIdentity, pub first_seen_at: SystemTime, pub resyncs: u32, pub pds_changes: u32, @@ -78,54 +72,6 @@ impl RepoInfo { Ok(Some(info)) } - pub fn load_or_create( - storage: &S, - registry: &HostRegistry, - did: &Did, - now: SystemTime, - ) -> Result> { - if let Some(existing) = Self::load(storage, registry, did)? { - return Ok(existing); - } - - // TODO: i still kind of feel like we might need a global clock to avoid - // scheduling resyncs before the host-earliest threshold... - - let desynchronized = Desynchronized { - reason: DesyncReason::FirstSeen, - resync_attempts: 0, - due_at: now, - }; - - let mut batch = storage.batch(); - - // add this to the batch - QueuedResync { - did: did.clone(), - due: desynchronized.due_at, - host: None, - } - .store(&mut batch); - - // and this, together - let new_me = Self { - sync_status: SyncStatus::Desynchronized(desynchronized), - account_status: AccountStatus::Active, - identity: None, - first_seen_at: now, - resyncs: 0, - pds_changes: 0, - handle_changes: 0, - last_resync: None, - }; - - new_me.clone().store(&mut batch, did); - - batch.commit().map_err(LoadError::Storage)?; - - Ok(new_me) - } - pub fn store(&self, batch: &mut B, did: &Did) where E: StorageError, @@ -140,19 +86,16 @@ impl RepoInfo { registry: &HostRegistry, db_info: DbRepoInfo, ) -> Result { - let identity = match db_info.identity { - Some(id) => Some(RepoIdentity { + let id = db_info.identity; + Ok(Self { + sync_status: db_info.sync_status, + account_status: db_info.account_status, + identity: RepoIdentity { pds_host: registry.get(&id.pds)?, signing_key: id.signing_key, supposed_handle: id.supposed_handle, resolved_at: id.resolved_at, - }), - None => None, - }; - Ok(Self { - sync_status: db_info.sync_status, - account_status: db_info.account_status, - identity, + }, first_seen_at: db_info.first_seen_at, resyncs: db_info.resyncs, pds_changes: db_info.pds_changes, @@ -164,16 +107,16 @@ impl RepoInfo { impl From for DbRepoInfo { fn from(info: RepoInfo) -> Self { - let identity = info.identity.map(|id| DbRepoIdentity { - pds: id.pds_host.name().as_str().to_string(), - signing_key: id.signing_key, - supposed_handle: id.supposed_handle, - resolved_at: id.resolved_at, - }); + let id = info.identity; Self { sync_status: info.sync_status, account_status: info.account_status, - identity, + identity: DbRepoIdentity { + pds: id.pds_host.name().as_str().to_string(), + signing_key: id.signing_key, + supposed_handle: id.supposed_handle, + resolved_at: id.resolved_at, + }, first_seen_at: info.first_seen_at, resyncs: info.resyncs, pds_changes: info.pds_changes, @@ -202,6 +145,17 @@ pub struct RepoIdentity { pub resolved_at: SystemTime, } +#[derive(Debug, Clone, Serialize, Deserialize)] +pub enum AccountStatus { + Active, + Deleted, + Deactivated, + Suspended, + Takendown, + /// fallback for unrecognized non-active status + Inactive(String), +} + #[derive(Debug, Clone, Serialize, Deserialize)] pub enum SyncStatus { Synchronized, @@ -216,14 +170,113 @@ pub enum SyncStatus { } #[derive(Debug, Clone, Serialize, Deserialize)] -pub enum AccountStatus { - Active, - Deleted, - Deactivated, - Suspended, - Takendown, - /// fallback for unrecognized non-active status - Inactive(String), +pub struct Desynchronized { + pub reason: DesyncReason, + pub resync_attempts: u32, + #[serde(with = "unix_ms_u64")] + pub due_at: SystemTime, +} + +impl Desynchronized { + pub fn backoff(reason: &DesyncReason, attempts: u32, method: DidMethod) -> Option { + match reason { + DesyncReason::UnresolvableIdentity => match method { + DidMethod::Plc => match attempts { + ..=1 => Some(60), + 2 => Some(300), + _ => Some(3600), // todo *definitely* need to handle not-found + }, + DidMethod::Web => match attempts { + ..=1 => Some(60), + 2..=4 => Some(300), + 5 => Some(3600), + _ => Some(86_400), // once per day forever (!) + }, + }, + _ => unimplemented!("todo"), + } + .map(Duration::from_secs) + } +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub enum DesyncReason { + /// from the firehose, or backfill discovery, or manually added + FirstSeen, + /// could not resolve a repo from discovery or in some other scenarios + /// + /// we can't accept any data if we can't get the signing key to verify it + UnresolvableIdentity, + /// a #sync event on the firehose + /// + /// #sync events are ignored when `rev` and `prev_data` fields don't change + /// + /// #sync events with a change also don't *necessarily* lead to + /// desynchronized state -- the actor is allowed to attempt a resync inline. + /// + /// but, if it doesn't, or if that fails, this was the source of that. + FirehoseSync, + /// sync1.1 proof validation failure on a #commit + /// + /// commits dropped otherwise (signature fail, old `rev`,) don't resync. + FirehoseFail { + #[serde(with = "crate::tid::tid_raw_u64_opt")] + bad_commit_rev: Option, // TODO Tid + }, + /// repo-level or host-level firehose rate-limits were applied + Throttled { exceeded: Option }, + /// pds is not emitting sync1.1 + Sync11Lax, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct ResyncInfo { + #[serde(with = "unix_ms_u64")] + pub at: SystemTime, + /// resync data size in bytes, if available + /// + /// for the default implementation, this should be the decoded raw car size + pub size: Option, + /// total repo records, if available + pub records: Option, + /// resync data source + /// + /// usually a pds, but could be a fallback, etc. (manual upload??). use + /// accessor methods (which enforce host interning) + from_host: Option, + /// how long the resync took, if available + /// + /// use an accessor for the nicely-typed value + duration_ms: Option, +} + +impl ResyncInfo { + pub fn new( + at: SystemTime, + size: Option, + records: Option, + from_host: Option<&Host>, + duration: Option, + ) -> Self { + Self { + at, + size, + records, + from_host: from_host.map(|h| h.name().as_str().to_string()), + duration_ms: duration.map(|d| d.as_millis() as u32), + } + } + pub fn host(&self, registry: &HostRegistry) -> Result>, HostnameError> { + let Some(ref hostname) = self.from_host else { + return Ok(None); + }; + let host = registry.get(hostname)?; + Ok(Some(host)) + } + pub fn duration(&self) -> Option { + let ms = self.duration_ms?; + Some(Duration::from_millis(ms.into())) + } } #[cfg(test)] @@ -245,16 +298,16 @@ mod tests { fn registry() -> HostRegistry { HostRegistry::new(Default::default()) } - fn sample_info(host: Option>) -> RepoInfo { + fn sample_info(host: Arc) -> RepoInfo { RepoInfo { sync_status: SyncStatus::Synchronized, account_status: AccountStatus::Active, - identity: host.map(|h| RepoIdentity { - pds_host: h, + identity: RepoIdentity { + pds_host: host, signing_key: vec![1, 2, 3, 4], supposed_handle: "alice.example".into(), resolved_at: t(1_000), - }), + }, first_seen_at: t(1_000), resyncs: 0, pds_changes: 0, @@ -271,31 +324,11 @@ mod tests { assert!(got.is_none()); } - #[test] - fn repo_info_round_trip_without_identity() { - let eng = MemEngine::new(); - let reg = registry(); - let info = sample_info(None); - let mut b = eng.batch(); - info.store(&mut b, &Did::raw("did:plc:abc")); - b.commit().unwrap(); - - let got = RepoInfo::load(&eng, ®, &Did::raw("did:plc:abc")) - .unwrap() - .expect("entry present"); - assert!(got.identity.is_none()); - assert_eq!(got.first_seen_at, t(1_000)); - assert_eq!(got.resyncs, 0); - assert!(matches!(got.sync_status, SyncStatus::Synchronized)); - assert!(matches!(got.account_status, AccountStatus::Active)); - assert!(got.last_resync.is_none()); - } - #[test] fn repo_info_round_trip_with_identity_interns_host() { let eng = MemEngine::new(); let reg = registry(); - let info = sample_info(Some(h("example.com"))); + let info = sample_info(h("example.com")); let mut b = eng.batch(); info.store(&mut b, &Did::raw("did:plc:abc")); b.commit().unwrap(); @@ -303,7 +336,7 @@ mod tests { let got = RepoInfo::load(&eng, ®, &Did::raw("did:plc:abc")) .unwrap() .expect("entry present"); - let id = got.identity.expect("identity preserved"); + let id = got.identity; assert_eq!(id.pds_host.name().as_str(), "example.com"); assert_eq!(id.signing_key, vec![1, 2, 3, 4]); assert_eq!(&*id.supposed_handle, "alice.example"); @@ -318,7 +351,7 @@ mod tests { let eng = MemEngine::new(); let reg = registry(); let host = h("example.com"); - let mut info = sample_info(Some(host.clone())); + let mut info = sample_info(host.clone()); info.sync_status = SyncStatus::Desynchronized(Desynchronized { reason: DesyncReason::FirehoseFail { bad_commit_rev: Some(42.try_into().unwrap()), @@ -377,12 +410,12 @@ mod tests { let db_info = DbRepoInfo { sync_status: SyncStatus::Synchronized, account_status: AccountStatus::Active, - identity: Some(DbRepoIdentity { + identity: DbRepoIdentity { pds: "user@example.com".into(), signing_key: vec![], supposed_handle: "x.example".into(), resolved_at: t(1000), - }), + }, first_seen_at: t(0), resyncs: 0, pds_changes: 0, diff --git a/hubble-sync/src/storage/repo/repo_index_resync.rs b/hubble-sync/src/storage/repo/info_idx_resync.rs similarity index 94% rename from hubble-sync/src/storage/repo/repo_index_resync.rs rename to hubble-sync/src/storage/repo/info_idx_resync.rs index 86e4229..dd79358 100644 --- a/hubble-sync/src/storage/repo/repo_index_resync.rs +++ b/hubble-sync/src/storage/repo/info_idx_resync.rs @@ -9,7 +9,7 @@ //! if this was sql, we'd hold the `due_millis` as a field on the repo table. //! here it only exists in the index. //! -//! "rq/"||||NUL|||| => [empty value] +//! "rq|"||||NUL|||| => [] //! //! the index is covering: the resync scheduler only need the DID to send a //! `ScheduledResync` task to its actor. @@ -29,8 +29,10 @@ use crate::metrics::{RESYNC_QUEUE_DEQUEUED_TOTAL, RESYNC_QUEUE_ENQUEUED_TOTAL}; use crate::resync_scheduler::QueuedResync; use crate::{Did, Host, HostRegistry, StorageBatch, StorageEngine, StorageError}; -use super::repo_state_account::{AccountStatus, RepoInfo, SyncStatus}; -use super::{DecodeError, LoadError, PREFIX_RESYNC_QUEUE, split_once}; +use super::{ + AccountStatus, DecodeError, LoadError, PREFIX_REPO_INFO_IDX_RESYNC, RepoInfo, SyncStatus, + slice_split_once, +}; fn encode_key(host: Option<&Host>, did: &Did, due: SystemTime) -> Vec { let host_bytes = host.map(|h| h.name().as_str().as_bytes()).unwrap_or(&[]); @@ -41,9 +43,10 @@ fn encode_key(host: Option<&Host>, did: &Did, due: SystemTime) -> Vec { .expect("ancient resync due-at") .as_millis() as u64; - let mut k = - Vec::with_capacity(PREFIX_RESYNC_QUEUE.len() + host_bytes.len() + 1 + did.len() + 8); - k.extend_from_slice(PREFIX_RESYNC_QUEUE); + let mut k = Vec::with_capacity( + PREFIX_REPO_INFO_IDX_RESYNC.len() + host_bytes.len() + 1 + did.len() + 8, + ); + k.extend_from_slice(PREFIX_REPO_INFO_IDX_RESYNC); k.extend_from_slice(host_bytes); k.push(0x00); k.extend_from_slice(&due_at_ms.to_be_bytes()); @@ -52,7 +55,7 @@ fn encode_key(host: Option<&Host>, did: &Did, due: SystemTime) -> Vec { } fn decode_key(reg: &HostRegistry, k: &[u8]) -> Result { - let (host_bytes, rest) = split_once(k, &0x00).ok_or(DecodeError::MissingNullSeparator)?; + let (host_bytes, rest) = slice_split_once(k, &0x00).ok_or(DecodeError::MissingNullSeparator)?; let host = if host_bytes.is_empty() { None @@ -69,8 +72,8 @@ fn build_host_prefix(host: Option<&Host>) -> Vec { let host_bytes = host .map(|h| h.name().as_str().as_bytes()) .unwrap_or_default(); - let mut p = Vec::with_capacity(PREFIX_RESYNC_QUEUE.len() + host_bytes.len() + 1); - p.extend_from_slice(PREFIX_RESYNC_QUEUE); + let mut p = Vec::with_capacity(PREFIX_REPO_INFO_IDX_RESYNC.len() + host_bytes.len() + 1); + p.extend_from_slice(PREFIX_REPO_INFO_IDX_RESYNC); p.extend_from_slice(host_bytes); p.push(0x00); p @@ -114,7 +117,7 @@ enum Index { #[derive(Debug, Clone, PartialEq)] struct Partial { - host: Option>, + host: Arc, due: SystemTime, } @@ -122,7 +125,7 @@ impl Index { fn should(info: &RepoInfo) -> Self { match (&info.account_status, &info.sync_status) { (AccountStatus::Active, SyncStatus::Desynchronized(d)) => Self::Should(Partial { - host: info.identity.as_ref().map(|i| i.pds_host.clone()), + host: info.identity.pds_host.clone(), due: d.due_at, }), _ => Self::ShouldNot, @@ -132,6 +135,7 @@ impl Index { impl QueuedResync { fn from_partial(did: Did, Partial { host, due }: Partial) -> Self { + let host = Some(host); Self { did, host, due } } @@ -216,7 +220,7 @@ impl<'a, S: StorageEngine> NextQueuedByHost<'a, S> { fn next_inner(&mut self) -> Result, LoadError> { let Some((k, _)) = self .engine - .scan_from_queue(PREFIX_RESYNC_QUEUE, &self.next_scan_suffix) + .scan_from_queue(PREFIX_REPO_INFO_IDX_RESYNC, &self.next_scan_suffix) .next() .transpose() .map_err(LoadError::Storage)? @@ -438,7 +442,7 @@ mod tests { .collect::, _>>() .expect("all ok"); assert_eq!(items.len(), 3); - // order is by hostname bytewise (the key prefix after PREFIX_RESYNC_QUEUE) + // order is by hostname bytewise (the key prefix after PREFIX_REPO_INFO_IDX_RESYNC) assert_eq!(items[0].host.as_ref().unwrap().name().as_str(), "a.com"); assert_eq!(items[1].host.as_ref().unwrap().name().as_str(), "b.com"); assert_eq!(items[2].host.as_ref().unwrap().name().as_str(), "c.com"); @@ -494,11 +498,11 @@ mod tests { let eng = MemEngine::new(); let reg = registry(); // write a raw key with no NUL separator after the prefix. the engine - // strips PREFIX_RESYNC_QUEUE during scan, so the iterator sees just + // strips PREFIX_REPO_INFO_IDX_RESYNC during scan, so the iterator sees just // the bad bytes and decode_key fails at the host-separator step. let mut b = eng.batch(); let mut bad_key = Vec::new(); - bad_key.extend_from_slice(PREFIX_RESYNC_QUEUE); + bad_key.extend_from_slice(PREFIX_REPO_INFO_IDX_RESYNC); bad_key.extend_from_slice(b"bytes_without_a_separator"); b.put_queue(&bad_key, &[]); b.commit().unwrap(); diff --git a/hubble-sync/src/storage/repo/mod.rs b/hubble-sync/src/storage/repo/mod.rs index 2e14ffa..8751538 100644 --- a/hubble-sync/src/storage/repo/mod.rs +++ b/hubble-sync/src/storage/repo/mod.rs @@ -1,9 +1,210 @@ -pub mod repo_impl; -pub mod repo_index_resync; -pub mod repo_state_account; -pub mod repo_state_account_desync; -pub mod repo_state_sync; +//! the aggregate a repo persisted state +//! +//! canonical state is in two pieces: account and sync +//! +//! - sync (see `repo_state_sync`): minimal state for sync1.1 inductive proof, +//! updates on every repo commit. +//! - account (see `repo_state_account`): everything else: account status, +//! resolved identity, overall sync state, etc. +//! +//! additionally there is one index (see `repo_index_resync`) which is a host- +//! partitioned, time-ordered index of every active repo in desynchronized +//! sync state. +//! +//! this module provides a higher-level interface to repo state, maintaining +//! consistency across all three key ranges through any changes. + +mod info; +mod info_idx_resync; +mod pending; +mod pending_idx_retry; +mod prev; use super::*; -pub use repo_impl::Repo; +pub use info::{AccountStatus, DesyncReason, Desynchronized, RepoIdentity, RepoInfo, SyncStatus}; +pub use info_idx_resync::NextQueuedByHost; +pub use pending::PendingIdentity; +pub use pending_idx_retry::PendingIdentityQueueEntry; +pub use prev::AccountSyncState; + +use std::time::SystemTime; + +use tracing::warn; + +use crate::identity::Validity; +use crate::resync_scheduler::QueuedResync; +use crate::{Did, HostRegistry, StorageBatch, StorageError}; + +pub enum Awoken { + Resolved(Repo), + Pending(PendingIdentity), + Refreshing { + repo: Repo, + pending: PendingIdentity, + }, +} + +#[derive(Debug, Clone)] +pub struct Repo { + did: Did, + info: RepoInfo, + sync: Option, +} + +impl Repo { + pub fn wake( + storage: &S, + hosts: &HostRegistry, + did: Did, + now: SystemTime, + ) -> Result> { + // happy path: existing repos + if let Some(info) = RepoInfo::load(storage, hosts, &did)? { + let sync = AccountSyncState::load(storage, &did)?; + let repo = Self { + did: did.clone(), + info, + sync, + }; + if repo.requires_identity_refresh() { + // if we cannot move forward with the existing identity (hard-expired) + let pending = PendingIdentity::load(storage, &did)? + .unwrap_or_else(|| PendingIdentity::new(did, now)); + return Ok(Awoken::Refreshing { repo, pending }); + } + return Ok(Awoken::Resolved(repo)); + } + // back again path: initial identity resolution pending + if let Some(p) = PendingIdentity::load(storage, &did)? { + return Ok(Awoken::Pending(p)); + } + // first sight, brand new, for this identity + Ok(Awoken::Pending(PendingIdentity::new(did, now))) + } + + /// bootstrap-promote constructor (pending resolved) + pub fn create_resolved>( + did: Did, + identity: RepoIdentity, + now: SystemTime, + batch: &mut B, + ) -> Self { + let info = RepoInfo { + sync_status: SyncStatus::Synchronized, // transitioned in a sec before commit + account_status: AccountStatus::Active, + identity, + first_seen_at: now, + resyncs: 0, + pds_changes: 0, + handle_changes: 0, + last_resync: None, + }; + let sync = None; + let mut repo = Self { did, info, sync }; + + let actual_status = SyncStatus::Desynchronized(Desynchronized { + reason: DesyncReason::FirstSeen, + resync_attempts: 0, + due_at: now, + }); + // the promised transition (does resync index maintenance) + repo.set_sync_status(actual_status, batch); + repo + } + + pub fn did(&self) -> &Did { + &self.did + } + /// TODO: expose these or just make accessors for everything? + pub fn info(&self) -> &RepoInfo { + &self.info + } + pub fn sync_state(&self) -> Option<&AccountSyncState> { + self.sync.as_ref() + } + + pub fn desynchronize>( + &mut self, + reason: DesyncReason, + now: SystemTime, + batch: &mut B, + ) { + let method = self.did().method(); + let attempts = match &self.info().sync_status { + SyncStatus::Desynchronized(d) => d.resync_attempts.saturating_add(1), + _ => 1, + }; + let new_status = match Desynchronized::backoff(&reason, attempts, method) { + Some(backoff) => SyncStatus::Desynchronized(Desynchronized { + reason, + resync_attempts: attempts, + due_at: now + backoff, + }), + None => { + warn!( + "danger: desynchronizing -> repo Gone due to None backoff (i don't have brain to think this through rn)" + ); + SyncStatus::Gone + } + }; + self.set_sync_status(new_status, batch); + } + + /// bump a resync's due_at (used to match an identity retry) + pub fn delay_resync_until>( + &mut self, + until: SystemTime, + batch: &mut B, + ) { + if let SyncStatus::Desynchronized(d) = &self.info.sync_status { + let mut delayed = d.clone(); + delayed.due_at = until; + self.set_sync_status(SyncStatus::Desynchronized(delayed), batch); + } + } + + pub fn set_account_status>( + &mut self, + new_status: AccountStatus, + batch: &mut B, + ) -> AccountStatus { + let prev = self.info.clone(); + self.info.account_status = new_status; + QueuedResync::reconcile_for(&self.did, &prev, &self.info, batch); + self.info.store(batch, &self.did); + prev.account_status + } + + pub fn set_sync_status>( + &mut self, + new_status: SyncStatus, + batch: &mut B, + ) -> SyncStatus { + let prev = self.info.clone(); + self.info.sync_status = new_status; + QueuedResync::reconcile_for(&self.did, &prev, &self.info, batch); + self.info.store(batch, &self.did); + prev.sync_status + } + + pub fn set_identity>( + &mut self, + new_id: RepoIdentity, + batch: &mut B, + ) -> RepoIdentity { + let prev = self.info.clone(); + self.info.identity = new_id; + QueuedResync::reconcile_for(&self.did, &prev, &self.info, batch); + self.info.store(batch, &self.did); + prev.identity + } + + pub fn needs_identity_refresh(&self) -> bool { + Validity::for_identity(&self.did, &self.info.identity).should_refresh() + } + + pub fn requires_identity_refresh(&self) -> bool { + Validity::for_identity(&self.did, &self.info.identity).must_refresh() + } +} diff --git a/hubble-sync/src/storage/repo/pending.rs b/hubble-sync/src/storage/repo/pending.rs new file mode 100644 index 0000000..accc42c --- /dev/null +++ b/hubble-sync/src/storage/repo/pending.rs @@ -0,0 +1,130 @@ +//! DIDs we've seen but couldn't resolve +//! +//! (or is it haven't-yet-resolved? will we always write early?) +//! +//! when resolution succeeds, deleted in an atomic batch with the repo info put +//! +//! "pi|"|| => + +use std::time::SystemTime; + +use dasl::drisl; +use serde::{Deserialize, Serialize}; + +use crate::identity::pending_resolve_backoff; +use crate::{Did, StorageBatch, StorageEngine, StorageError}; + +use super::{LoadError, PREFIX_REPO_PENDING, PendingIdentityQueueEntry, unix_ms_u64}; + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct DbPendingValue { + #[serde(with = "unix_ms_u64")] + first_seen: SystemTime, + #[serde(with = "unix_ms_u64")] + next_try_at: SystemTime, + attempts: u32, + last_error: Option, +} + +#[derive(Debug, Clone)] +pub struct PendingIdentity { + pub did: Did, + first_seen: SystemTime, + next_try_at: SystemTime, + attempts: u32, + last_error: Option, +} + +impl PendingIdentity { + pub fn new(did: Did, now: SystemTime) -> Self { + Self { + did, + first_seen: now, + next_try_at: now, + attempts: 0, + last_error: None, + } + } + + pub fn is_ready(&self, now: SystemTime) -> bool { + now >= self.next_try_at + } + + pub fn next_try_at(&self) -> SystemTime { + self.next_try_at + } + + fn key(did: &Did) -> Vec { + let did_bytes = did.as_str().as_bytes(); + let mut k = Vec::with_capacity(PREFIX_REPO_PENDING.len() + did_bytes.len()); + k.extend_from_slice(PREFIX_REPO_PENDING); + k.extend_from_slice(did_bytes); + k + } + + fn as_retry_entry(&self) -> Option { + (self.attempts > 0).then_some(PendingIdentityQueueEntry { + did: self.did.clone(), + due: self.next_try_at, + }) + } + + pub fn load( + storage: &S, + did: &Did, + ) -> Result, LoadError> { + let Some(bytes) = storage.get(&Self::key(did)).map_err(LoadError::Storage)? else { + return Ok(None); + }; + let db_val: DbPendingValue = drisl::from_slice(&bytes)?; + Ok(Some(Self { + did: did.clone(), + first_seen: db_val.first_seen, + attempts: db_val.attempts, + next_try_at: db_val.next_try_at, + last_error: db_val.last_error, + })) + } + + fn store>(&self, batch: &mut B) { + let db_val = DbPendingValue { + first_seen: self.first_seen, + attempts: self.attempts, + next_try_at: self.next_try_at, + last_error: self.last_error.clone(), + }; + let bytes = drisl::to_vec(&db_val).expect("drisl to_vec infallibe"); + batch.put(&Self::key(&self.did), &bytes); + } + + /// pending id resolution failed: get set up for the retry + /// + /// insert the retry entry, fix up the queue index (atomically) + pub fn store_failed>( + &mut self, + reason: &str, + now: SystemTime, + batch: &mut B, + ) { + // if we have saved a queue entry (attempt 0 is in-mem) + let prior_retry = self.as_retry_entry(); + + // bookkeep + self.attempts = self.attempts.saturating_add(1); + self.next_try_at = now + pending_resolve_backoff(self.did.method(), self.attempts); + self.last_error = Some(reason.to_string()); + + // i guess we're not doing fancy index reconciliation like resync queue + if let Some(r) = prior_retry { + r.delete(batch); + } + self.as_retry_entry() + .expect("just incremented attempts") + .insert(batch); + self.store(batch); + } + + pub fn delete>(self, batch: &mut B) { + batch.delete(&Self::key(&self.did)); + } +} diff --git a/hubble-sync/src/storage/repo/pending_idx_retry.rs b/hubble-sync/src/storage/repo/pending_idx_retry.rs new file mode 100644 index 0000000..e028b93 --- /dev/null +++ b/hubble-sync/src/storage/repo/pending_idx_retry.rs @@ -0,0 +1,73 @@ +//! Index over pending identity resolution retries +//! +//! "pq|"|||| => [] + +use std::time::{Duration, SystemTime, UNIX_EPOCH}; + +use super::{DecodeError, LoadError, PREFIX_REPO_PENDING_IDX_RETRY}; +use crate::{Did, Host, StorageBatch, StorageEngine, StorageError}; + +pub struct PendingIdentityQueueEntry { + pub did: Did, + pub due: SystemTime, +} + +impl PendingIdentityQueueEntry { + fn key(did: &Did, due: SystemTime) -> Vec { + let due_ms = due + .duration_since(UNIX_EPOCH) + .expect("post-epoch pending due") + .as_millis() as u64; + let did_bytes = did.as_str().as_bytes(); + let mut k = Vec::with_capacity(PREFIX_REPO_PENDING_IDX_RETRY.len() + 8 + did_bytes.len()); + k.extend_from_slice(PREFIX_REPO_PENDING_IDX_RETRY); + k.extend_from_slice(&due_ms.to_be_bytes()); + k.extend_from_slice(did_bytes); + k + } + + /// a key that is strictly after all items under this host + /// + /// does not write the top-level prefix (it's for prefix-scans under it) + fn build_host_post_suffix(host: Option<&Host>) -> Vec { + let host_bytes = host + .map(|h| h.name().as_str().as_bytes()) + .unwrap_or(&[0x00]); // [0x00] is after [] + let mut p = Vec::with_capacity(host_bytes.len() + 1); + p.extend_from_slice(host_bytes); + p.push(0x01); + p + } + + fn decode_key(k_unprefixed: &[u8]) -> Result { + let (due_bytes, did_bytes) = k_unprefixed + .split_first_chunk::<8>() + .ok_or(DecodeError::InputTooShort)?; + let due_ms = u64::from_be_bytes(*due_bytes); + let due = UNIX_EPOCH + Duration::from_millis(due_ms); + let did = str::from_utf8(did_bytes) + .map_err(|err| DecodeError::NotUtf8 { what: "did", err })? + .try_into() + .map_err(DecodeError::BadDid)?; + Ok(Self { did, due }) + } + + pub fn insert>(&self, batch: &mut B) { + batch.put(&Self::key(&self.did, self.due), &[]); + } + + pub fn delete>(self, batch: &mut B) { + batch.delete(&Self::key(&self.did, self.due)); + } + + pub fn scan( + engine: &S, + ) -> impl Iterator>> { + engine + .scan_from_queue(PREFIX_REPO_PENDING_IDX_RETRY, &[]) + .map(|r| { + let (k, _) = r.map_err(LoadError::Storage)?; + Ok(Self::decode_key(&k)?) + }) + } +} diff --git a/hubble-sync/src/storage/repo/repo_state_sync.rs b/hubble-sync/src/storage/repo/prev.rs similarity index 95% rename from hubble-sync/src/storage/repo/repo_state_sync.rs rename to hubble-sync/src/storage/repo/prev.rs index e43b0c5..84871f8 100644 --- a/hubble-sync/src/storage/repo/repo_state_sync.rs +++ b/hubble-sync/src/storage/repo/prev.rs @@ -2,7 +2,7 @@ //! //! see `repo` for public-facing repo view of combined state. //! -//! "rc/"|| => +//! "pi|"|| => //! |||||||||| //! //! manual ser/de for mimimal size @@ -16,7 +16,7 @@ use crate::{Did, StorageBatch, StorageEngine, StorageError, Tid}; -use super::{DecodeError, LoadError, PREFIX_REPO_ACCOUNT_SYNC_STATE}; +use super::{DecodeError, LoadError, PREFIX_REPO_PREV}; #[derive(Debug, Clone)] pub struct AccountSyncState { @@ -31,8 +31,8 @@ pub struct AccountSyncState { impl AccountSyncState { fn key(did: &Did) -> Vec { let did = did.as_str().as_bytes(); - let mut k = Vec::with_capacity(PREFIX_REPO_ACCOUNT_SYNC_STATE.len() + did.len()); - k.extend_from_slice(PREFIX_REPO_ACCOUNT_SYNC_STATE); + let mut k = Vec::with_capacity(PREFIX_REPO_PREV.len() + did.len()); + k.extend_from_slice(PREFIX_REPO_PREV); k.extend_from_slice(did); k } @@ -106,10 +106,10 @@ impl From<&AccountSyncState> for Vec { #[cfg(test)] mod tests { + use super::AccountSyncState; use super::*; use crate::StorageBatch; use crate::storage::engine::mem::MemEngine; - use crate::storage::repo::repo_state_sync::AccountSyncState; /// build a 36-byte v1 dag-cbor sha256 CID stuffed with `seed` for the /// digest. handles dasl::cid's raw-bytes constructor for tests where diff --git a/hubble-sync/src/storage/repo/repo_impl.rs b/hubble-sync/src/storage/repo/repo_impl.rs deleted file mode 100644 index b62e81d..0000000 --- a/hubble-sync/src/storage/repo/repo_impl.rs +++ /dev/null @@ -1,135 +0,0 @@ -//! the aggregate a repo persisted state -//! -//! canonical state is in two pieces: account and sync -//! -//! - sync (see `repo_state_sync`): minimal state for sync1.1 inductive proof, -//! updates on every repo commit. -//! - account (see `repo_state_account`): everything else: account status, -//! resolved identity, overall sync state, etc. -//! -//! additionally there is one index (see `repo_index_resync`) which is a host- -//! partitioned, time-ordered index of every active repo in desynchronized -//! sync state. -//! -//! this module provides a higher-level interface to repo state, maintaining -//! consistency across all three key ranges through any changes. - -use std::time::SystemTime; - -use tracing::warn; - -use super::repo_state_account::{AccountStatus, RepoIdentity, RepoInfo, SyncStatus}; -use super::repo_state_account_desync::{DesyncReason, Desynchronized}; -use super::repo_state_sync::AccountSyncState; -use super::{LoadError, StorageEngine}; -use crate::identity::Validity; -use crate::resync_scheduler::QueuedResync; -use crate::{Did, HostRegistry, StorageBatch, StorageError}; - -#[derive(Debug, Clone)] -pub struct Repo { - did: Did, - info: RepoInfo, - sync: Option, -} - -impl Repo { - pub fn load_or_create( - storage: &S, - hosts: &HostRegistry, - did: Did, - now: SystemTime, - ) -> Result> { - let info = RepoInfo::load_or_create(storage, hosts, &did, now)?; - let sync = AccountSyncState::load(storage, &did)?; - Ok(Self { did, info, sync }) - } - - pub fn did(&self) -> &Did { - &self.did - } - /// TODO: expose these or just make accessors for everything? - pub fn info(&self) -> &RepoInfo { - &self.info - } - pub fn sync_state(&self) -> Option<&AccountSyncState> { - self.sync.as_ref() - } - - pub fn desynchronize>( - &mut self, - reason: DesyncReason, - now: SystemTime, - batch: &mut B, - ) { - let method = self.did().method(); - let attempts = match &self.info().sync_status { - SyncStatus::Desynchronized(d) => d.resync_attempts.saturating_add(1), - _ => 1, - }; - let new_status = match Desynchronized::backoff(&reason, attempts, method) { - Some(backoff) => SyncStatus::Desynchronized(Desynchronized { - reason, - resync_attempts: attempts, - due_at: now + backoff, - }), - None => { - warn!( - "danger: desynchronizing -> repo Gone due to None backoff (i don't have brain to think this through rn)" - ); - SyncStatus::Gone - } - }; - self.set_sync_status(new_status, batch); - } - - pub fn set_account_status>( - &mut self, - new_status: AccountStatus, - batch: &mut B, - ) -> AccountStatus { - let prev = self.info.clone(); - self.info.account_status = new_status; - QueuedResync::reconcile_for(&self.did, &prev, &self.info, batch); - self.info.store(batch, &self.did); - prev.account_status - } - - pub fn set_sync_status>( - &mut self, - new_status: SyncStatus, - batch: &mut B, - ) -> SyncStatus { - let prev = self.info.clone(); - self.info.sync_status = new_status; - QueuedResync::reconcile_for(&self.did, &prev, &self.info, batch); - self.info.store(batch, &self.did); - prev.sync_status - } - - pub fn set_identity>( - &mut self, - new_id: RepoIdentity, - batch: &mut B, - ) -> Option { - let prev = self.info.clone(); - self.info.identity = Some(new_id); - QueuedResync::reconcile_for(&self.did, &prev, &self.info, batch); - self.info.store(batch, &self.did); - prev.identity - } - - pub fn needs_identity_refresh(&self) -> bool { - let Some(identity) = &self.info.identity else { - return true; - }; - let validity = match Validity::for_identity(&self.did, identity) { - Ok(v) => v, - Err(did) => { - tracing::warn!("can't resolve unrecognized DID method for {did}"); - return false; - } - }; - validity.should_refresh() - } -} diff --git a/hubble-sync/src/storage/repo/repo_state_account_desync.rs b/hubble-sync/src/storage/repo/repo_state_account_desync.rs deleted file mode 100644 index 0e47fa8..0000000 --- a/hubble-sync/src/storage/repo/repo_state_account_desync.rs +++ /dev/null @@ -1,123 +0,0 @@ -//! broken-out account-info state subtypes for desync -//! -//! (this module does not represent a distinct keyspace in the database) - -use std::sync::Arc; -use std::time::{Duration, SystemTime}; - -use serde::{Deserialize, Serialize}; - -use crate::host::{Host, HostnameError}; -use crate::{DidMethod, HostRegistry, Tid}; - -use super::unix_ms_u64; - -#[derive(Debug, Clone, Serialize, Deserialize)] -pub struct Desynchronized { - pub reason: DesyncReason, - pub resync_attempts: u32, - #[serde(with = "unix_ms_u64")] - pub due_at: SystemTime, -} - -impl Desynchronized { - pub fn backoff(reason: &DesyncReason, attempts: u32, method: DidMethod) -> Option { - match reason { - DesyncReason::UnresolvableIdentity => match method { - DidMethod::Plc => match attempts { - ..=1 => Some(60), - 2 => Some(300), - _ => Some(3600), // todo *definitely* need to handle not-found - }, - DidMethod::Web => match attempts { - ..=1 => Some(60), - 2..=4 => Some(300), - 5 => Some(3600), - _ => Some(86_400), // once per day forever (!) - }, - }, - _ => unimplemented!("todo"), - } - .map(Duration::from_secs) - } -} - -#[derive(Debug, Clone, Serialize, Deserialize)] -pub enum DesyncReason { - /// from the firehose, or backfill discovery, or manually added - FirstSeen, - /// could not resolve a repo from discovery or in some other scenarios - /// - /// we can't accept any data if we can't get the signing key to verify it - UnresolvableIdentity, - /// a #sync event on the firehose - /// - /// #sync events are ignored when `rev` and `prev_data` fields don't change - /// - /// #sync events with a change also don't *necessarily* lead to - /// desynchronized state -- the actor is allowed to attempt a resync inline. - /// - /// but, if it doesn't, or if that fails, this was the source of that. - FirehoseSync, - /// sync1.1 proof validation failure on a #commit - /// - /// commits dropped otherwise (signature fail, old `rev`,) don't resync. - FirehoseFail { - #[serde(with = "crate::tid::tid_raw_u64_opt")] - bad_commit_rev: Option, // TODO Tid - }, - /// repo-level or host-level firehose rate-limits were applied - Throttled { exceeded: Option }, - /// pds is not emitting sync1.1 - Sync11Lax, -} - -#[derive(Debug, Clone, Serialize, Deserialize)] -pub struct ResyncInfo { - #[serde(with = "unix_ms_u64")] - pub at: SystemTime, - /// resync data size in bytes, if available - /// - /// for the default implementation, this should be the decoded raw car size - pub size: Option, - /// total repo records, if available - pub records: Option, - /// resync data source - /// - /// usually a pds, but could be a fallback, etc. (manual upload??). use - /// accessor methods (which enforce host interning) - from_host: Option, - /// how long the resync took, if available - /// - /// use an accessor for the nicely-typed value - duration_ms: Option, -} - -impl ResyncInfo { - pub fn new( - at: SystemTime, - size: Option, - records: Option, - from_host: Option<&Host>, - duration: Option, - ) -> Self { - Self { - at, - size, - records, - from_host: from_host.map(|h| h.name().as_str().to_string()), - duration_ms: duration.map(|d| d.as_millis() as u32), - } - } - pub fn host(&self, registry: &HostRegistry) -> Result>, HostnameError> { - let Some(ref hostname) = self.from_host else { - return Ok(None); - }; - let host = registry.get(hostname)?; - Ok(Some(host)) - } - pub fn duration(&self) -> Option { - let ms = self.duration_ms?; - Some(Duration::from_millis(ms.into())) - } -}