diff --git a/hubble-pds/src/hostname.rs b/hubble-pds/src/hostname.rs index 06d49f0..538d8d9 100644 --- a/hubble-pds/src/hostname.rs +++ b/hubble-pds/src/hostname.rs @@ -10,6 +10,12 @@ use std::str::FromStr; use reqwest::Url; +/// Normalized atproto PDS hostname. +/// +/// Usually should *not* be passed around. Functions should usually not accept +/// this type as a parameter, and structs should almost never hold one directly. +/// Instead, wherever we're talking about a host, pass the `Host` in (`&Host` or +/// `Arc`, cheap both ways). #[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord)] pub struct Hostname(String); diff --git a/hubble-sync/readme.md b/hubble-sync/readme.md index 2e39042..1ce4741 100644 --- a/hubble-sync/readme.md +++ b/hubble-sync/readme.md @@ -261,7 +261,13 @@ a jetstream impl based on hubble-sync might add a new emitted event type (somewh - [ ] gauge: rss - [ ] ulimit helper +- [ ] set an app name for the user-agent, not just a contact info + - [ ] guard against ssrf + +- [ ] deep-crawl strategy (listHosts -> listRepos) + - [ ] we'll want to add the upstream status checks for this probably + - [ ] sync new repos directly (no-resync) from firehose when they start from definitely-empty - [ ] for firehose-discovered repos, we can (and should) use the car slice to predict if it's big! (from mst height) - [ ] the firehose high water seq mark is nice but... i think we can just persist a sequence a few seconds behind for the same effect without acks-per-message flying around? diff --git a/hubble-sync/src/config.rs b/hubble-sync/src/config.rs index 328a66d..56c4771 100644 --- a/hubble-sync/src/config.rs +++ b/hubble-sync/src/config.rs @@ -5,7 +5,7 @@ use std::path::PathBuf; use std::time::Duration; use crate::firehose::FirehoseConfig; -use crate::host::HostRegistryConfig; +use crate::host::{Host, HostRegistryConfig}; #[derive(Debug, Clone)] pub struct SyncConfig { @@ -95,6 +95,10 @@ pub struct SyncConfig { /// /// default: `[0x00]` pub hubble_sync_storage_prefix: &'static [u8], + /// policy for which repos get synchronized + /// + /// default: SyncScope::Everything + pub sync_scope: SyncScope, /// upstream (subscribeRepos source) config pub upstream: UpstreamConfig, /// firehose subscriber tuning @@ -117,6 +121,7 @@ impl Default for SyncConfig { scheduled_resolve_limit: 3.try_into().unwrap(), reactive_permit_wait_timeout: Duration::from_secs(10), hubble_sync_storage_prefix: &[0x00], + sync_scope: SyncScope::default(), upstream: UpstreamConfig::default(), firehose: FirehoseConfig::default(), hosts: HostRegistryConfig::default(), @@ -160,3 +165,30 @@ impl UpstreamKind { matches!(self, UpstreamKind::Relay) } } + +/// policy for repo tracking +/// +/// right now hacked for one non-everything policy, but later this will be +/// more consumer-app-configurable (and work together with discovery strategies) +#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)] +pub enum SyncScope { + /// track only bsky-pds repos + /// + /// bluesky has reasonable rate-limiting and other defensive config; many + /// indie pdses run their machines close to the margin and without much + /// defense. so, it can be useful to run non-permanent tests only against + /// bluesky (in particular because it's still *nearly* full-network alone). + BskyPdsOnly, + /// track everything + #[default] + Everything, +} + +impl SyncScope { + pub fn allows_host(&self, host: &Host) -> bool { + match self { + SyncScope::BskyPdsOnly => host.name().is_bsky(), + SyncScope::Everything => true, + } + } +} diff --git a/hubble-sync/src/host/mod.rs b/hubble-sync/src/host/mod.rs index 887576d..d27de9b 100644 --- a/hubble-sync/src/host/mod.rs +++ b/hubble-sync/src/host/mod.rs @@ -43,7 +43,7 @@ use tracing::{debug, info}; use crate::metrics::{HOST_BACKOFFS_TOTAL, HOST_SYNC_COMPLIANCE_UPGRADES_TOTAL}; use crate::storage::host_info; -use crate::{StorageBatch, StorageEngine, StorageError}; +use crate::{StorageBatch, StorageEngine, StorageError, SyncScope}; const HOST_CONNECT_TIMEOUT: Duration = Duration::from_secs(15); const HOST_REQUEST_TIMEOUT: Duration = Duration::from_secs(30); @@ -161,10 +161,11 @@ impl BackoffCell { /// to the same host use the same underlying resource, including self-rate- /// limiting and request concurrency limits. /// -/// we wrap some `com.atproto.sync.*` calls directly for convenience -/// -/// TODO: expose generic request wrapper -/// TODO: expose a generic jacquard xrpc request wrapper +/// we should almost always hold on to whole `Host`s, not `Hostnames`. ff we're +/// talking about a host somewhere then it's useful to keep it alive. in +/// particular, the norm is to pass `&Host` or `Arc` to functions, and +/// hold `Arc, @@ -172,6 +173,7 @@ pub struct Host { request_concurrency: Arc, backoff: BackoffCell, sync_compliance: ComplianceCell, + in_scope: bool, // temp: (maybe) host-based sync policy } impl Host { @@ -181,21 +183,32 @@ impl Host { client: Arc, qps: NonZeroU32, concurrency: usize, + sync_scope: SyncScope, ) -> Self { - Self { + let mut me = Self { name, client, limiter: RateLimiter::direct(Quota::per_second(qps)), request_concurrency: Arc::new(Semaphore::new(concurrency)), backoff: BackoffCell::new(), sync_compliance: ComplianceCell::new_unset(), + in_scope: true, + }; + if !sync_scope.allows_host(&me) { + // slightly awk, oh well + me.in_scope = false; } + me } pub fn name(&self) -> &Hostname { &self.name } + pub fn in_scope(&self) -> bool { + self.in_scope + } + pub fn jacquard_uri_base(&self, scheme: &str) -> Result, String> { // ...could maybe use the builder stuff, but it's pretty verbose let base = format!("{scheme}://{}", self.name().as_str()); @@ -289,6 +302,7 @@ impl Host { hosts.client_handle(), DEFAULT_HOST_QPS, DEFAULT_HOST_CONCURRENCY, + SyncScope::Everything, ) } diff --git a/hubble-sync/src/host/registry.rs b/hubble-sync/src/host/registry.rs index 17b4376..e68590f 100644 --- a/hubble-sync/src/host/registry.rs +++ b/hubble-sync/src/host/registry.rs @@ -11,15 +11,17 @@ use std::sync::{Arc, Mutex, Weak}; use metrics::counter; use tokio_util::sync::CancellationToken; +use crate::SyncScope; use crate::config::UpstreamConfig; use crate::metrics::HOST_REGISTRY_LOADS_TOTAL; #[derive(Debug, Clone, Default)] pub struct HostRegistryConfig { - host_qps: Option, - host_bsky_qps: Option, - host_concurrency: Option, - host_bsky_concurrency: Option, + pub host_qps: Option, + pub host_bsky_qps: Option, + pub host_concurrency: Option, + pub host_bsky_concurrency: Option, + pub sync_scope: SyncScope, // temp: (maybe) host-based sync policy } pub struct HostRegistry { @@ -30,6 +32,7 @@ pub struct HostRegistry { host_concurrency: usize, host_bsky_concurrency: usize, upstream: UpstreamConfig, + sync_scope: SyncScope, } impl HostRegistry { @@ -38,6 +41,7 @@ impl HostRegistry { upstream: UpstreamConfig, ua: &str, cancel: CancellationToken, + sync_scope: SyncScope, ) -> Arc { let client = reqwest::Client::builder() .user_agent(ua) @@ -61,6 +65,7 @@ impl HostRegistry { .host_bsky_concurrency .unwrap_or(DEFAULT_BSKY_HOST_CONCURRENCY), upstream, + sync_scope, } }) } @@ -72,6 +77,7 @@ impl HostRegistry { Default::default(), "", CancellationToken::new(), + SyncScope::Everything, ) } @@ -119,6 +125,7 @@ impl HostRegistry { self.client.clone(), qps, concurrency, + self.sync_scope, )); map.insert(hostname, Arc::downgrade(&host)); Ok(host) diff --git a/hubble-sync/src/hubble_sync.rs b/hubble-sync/src/hubble_sync.rs index b81cbfa..366bf6d 100644 --- a/hubble-sync/src/hubble_sync.rs +++ b/hubble-sync/src/hubble_sync.rs @@ -89,6 +89,7 @@ where config.upstream.clone(), &ua, cancel.clone(), + config.sync_scope, ); let resolver = HubbleSyncResolver::new(&ua, &config.plc_url, hosts.clone()); Self::with_resolver(consumer_app, storage, hosts, resolver, cancel, config) diff --git a/hubble-sync/src/lib.rs b/hubble-sync/src/lib.rs index f3776e9..b986d96 100644 --- a/hubble-sync/src/lib.rs +++ b/hubble-sync/src/lib.rs @@ -27,7 +27,7 @@ use host::Hostname; pub use cid::{Cid as DaslCid, CidParseError}; // VENDORED for now pub use commit::{Commit, CommitObject, Op, OpKind}; -pub use config::{SyncConfig, UpstreamConfig}; +pub use config::{SyncConfig, SyncScope, UpstreamConfig}; pub use firehose::{ FirehoseAck, FirehoseConfig, FirehoseError, FirehoseEvent, FirehosePayload, FirehoseSubscriber, }; diff --git a/hubble-sync/src/repo_actor/identity_refresh.rs b/hubble-sync/src/repo_actor/identity_refresh.rs index ac8150c..4ddd609 100644 --- a/hubble-sync/src/repo_actor/identity_refresh.rs +++ b/hubble-sync/src/repo_actor/identity_refresh.rs @@ -54,7 +54,7 @@ impl IdentityRefresh { *repo = spawn_blocking(move || { let mut batch = storage_inner.batch(); - repo_inner.set_identity(identity, &mut batch); + repo_inner.set_identity(identity, now, &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) diff --git a/hubble-sync/src/repo_actor/task_processor/mod.rs b/hubble-sync/src/repo_actor/task_processor/mod.rs index f6b4ff6..e7f641f 100644 --- a/hubble-sync/src/repo_actor/task_processor/mod.rs +++ b/hubble-sync/src/repo_actor/task_processor/mod.rs @@ -285,9 +285,14 @@ impl, R: Resolve> TaskProcessor "drop_out_of_scope").increment(1); + return Ok(()); + } if !self.repo().is_synchronized() { - trace!("dropping commit for desynchronized account"); + trace!("dropping #commit for desynchronized account"); counter!(COMMIT_OUTCOMES_TOTAL, "outcome" => "drop_desynced").increment(1); return Ok(()); } @@ -402,7 +407,12 @@ impl, R: Resolve> TaskProcessor PEResult { trace!("started processing #sync event"); - // ...just drop if we know we aren't in sync anyway + // ...just drop if we know we aren't in sync (or scope) anyway + if !self.repo().is_in_scope() { + trace!("dropping #sync for out-of-scope repo"); + counter!(COMMIT_OUTCOMES_TOTAL, "outcome" => "drop_out_of_scope").increment(1); + return Ok(()); + } if !self.repo().is_synchronized() { trace!("dropping #sync for desynchronized account"); counter!(SYNC_OUTCOMES_TOTAL, "outcome" => "drop_desynced").increment(1); @@ -1065,7 +1075,7 @@ impl, R: Resolve> TaskProcessor PEResult { let mut batch = storage.batch(); - updating.set_identity(identity, &mut batch); + updating.set_identity(identity, now, &mut batch); batch.commit().map_err(ProcessError::Storage)?; Ok(updating) }) diff --git a/hubble-sync/src/storage/repo/info.rs b/hubble-sync/src/storage/repo/info.rs index f74942e..8ff24f4 100644 --- a/hubble-sync/src/storage/repo/info.rs +++ b/hubble-sync/src/storage/repo/info.rs @@ -292,6 +292,10 @@ impl From for AccountStatus { pub enum SyncStatus { Synchronized, Desynchronized(Desynchronized), + OutOfScope { + #[serde(with = "unix_ms_u64")] + since: SystemTime, + }, /// warning with this one: we aren't sync-guaranteed to see a reactivation, /// which could leave the repo accidentally in this state forever. /// if we see a firehose event from a deactivate account, we should probably @@ -310,6 +314,7 @@ impl SyncStatus { match self { Self::Synchronized => SyncStatusKind::Synchronized, Self::Desynchronized(_) => SyncStatusKind::Desynchronized, + Self::OutOfScope { .. } => SyncStatusKind::OutOfScope, Self::Deactivated => SyncStatusKind::Deactivated, Self::Gone { .. } => SyncStatusKind::Gone, } @@ -324,6 +329,7 @@ impl SyncStatus { pub enum SyncStatusKind { Synchronized, Desynchronized, + OutOfScope, Deactivated, Gone, } @@ -333,6 +339,7 @@ impl SyncStatusKind { match self { Self::Synchronized => "synchronized", Self::Desynchronized => "desynchronized", + Self::OutOfScope => "outOfScope", Self::Deactivated => "deactivated", Self::Gone => "gone", } @@ -411,6 +418,7 @@ impl Desynchronized { }, }, DesyncReason::FirstSeen => Some(Self::backoff_from_zero(1.5, attempts)), + DesyncReason::CameIntoScope => Some(Self::backoff_from_zero(1.5, attempts)), DesyncReason::FutureRev { rev } => { let now = SystemTime::now(); let (t, _) = (*rev).into(); @@ -447,6 +455,7 @@ impl Desynchronized { use DesyncReason as R; match self.reason { R::FirstSeen => matches!(o, R::FirstSeen), + R::CameIntoScope => matches!(o, R::CameIntoScope), R::UnresolvableIdentity => matches!(o, R::UnresolvableIdentity), R::FirehoseSync { .. } => matches!(o, R::FirehoseSync { .. }), R::FirehoseFail { .. } => matches!(o, R::FirehoseFail { .. }), @@ -473,6 +482,8 @@ impl Desynchronized { pub enum DesyncReason { /// from the firehose, or backfill discovery, or manually added FirstSeen, + /// a previously-out-of-scope repo became in-scope + CameIntoScope, /// 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 @@ -521,6 +532,7 @@ impl DesyncReason { pub fn kind(&self) -> DesyncReasonKind { match self { Self::FirstSeen => DesyncReasonKind::FirstSeen, + Self::CameIntoScope => DesyncReasonKind::CameIntoScope, Self::UnresolvableIdentity => DesyncReasonKind::UnresolvableIdentity, Self::FirehoseSync { .. } => DesyncReasonKind::FirehoseSync, Self::FirehoseFail { .. } => DesyncReasonKind::FirehoseFail, @@ -544,6 +556,7 @@ impl DesyncReason { #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum DesyncReasonKind { FirstSeen, + CameIntoScope, UnresolvableIdentity, FirehoseSync, FirehoseFail, @@ -559,6 +572,7 @@ impl DesyncReasonKind { /// every variant, for iterating count buckets / emitting gauges pub const ALL: &'static [DesyncReasonKind] = &[ Self::FirstSeen, + Self::CameIntoScope, Self::UnresolvableIdentity, Self::FirehoseSync, Self::FirehoseFail, @@ -573,6 +587,7 @@ impl DesyncReasonKind { pub fn name(&self) -> &'static str { match self { Self::FirstSeen => "first_seen", + Self::CameIntoScope => "came_into_scope", Self::UnresolvableIdentity => "unresolvable_identity", Self::FirehoseSync => "firehose_sync", Self::FirehoseFail => "firehose_fail", diff --git a/hubble-sync/src/storage/repo/info_idx_state_count.rs b/hubble-sync/src/storage/repo/info_idx_state_count.rs index e5eec58..39ed744 100644 --- a/hubble-sync/src/storage/repo/info_idx_state_count.rs +++ b/hubble-sync/src/storage/repo/info_idx_state_count.rs @@ -90,6 +90,7 @@ impl SyncStatus { #[derive(Debug, Clone, serde::Serialize)] pub struct RepoCountsByState { pub synchronized: u64, + pub out_of_scope: u64, /// total across every desync reason pub desynchronized: u64, pub deactivated: u64, @@ -136,6 +137,7 @@ impl RepoCountsByState { } Ok(Self { synchronized: Self::load_state(storage, SyncStatusKind::Synchronized)?, + out_of_scope: Self::load_state(storage, SyncStatusKind::OutOfScope)?, desynchronized, deactivated: Self::load_state(storage, SyncStatusKind::Deactivated)?, gone: Self::load_state(storage, SyncStatusKind::Gone)?, diff --git a/hubble-sync/src/storage/repo/mod.rs b/hubble-sync/src/storage/repo/mod.rs index 8dc4bb8..e35e5c3 100644 --- a/hubble-sync/src/storage/repo/mod.rs +++ b/hubble-sync/src/storage/repo/mod.rs @@ -135,10 +135,12 @@ impl Repo { let sync = None; let mut repo = Self { did, info, sync }; - let actual_status = - SyncStatus::Desynchronized(Desynchronized::new(DesyncReason::FirstSeen, now, 0, None)); - // the promised transition (does resync index maintenance) - repo.set_sync_status(actual_status, batch); + let actual_status = if repo.is_in_scope() { + SyncStatus::Desynchronized(Desynchronized::new(DesyncReason::FirstSeen, now, 0, None)) + } else { + SyncStatus::OutOfScope { since: now } + }; + repo.init_sync_status(actual_status, batch); AccountStatus::reconcile_for(None, Some(&repo.info.account_status()), batch); repo } @@ -237,6 +239,10 @@ impl Repo { self.info.is_active() } + pub fn is_in_scope(&self) -> bool { + self.info.identity.pds_host.in_scope() + } + pub fn is_upstream_active(&self) -> bool { self.info.upstream_status.is_active() } @@ -372,6 +378,13 @@ impl Repo { now: SystemTime, batch: &mut B, ) { + // can't transition sync state if we're out-of-scope without risking + // accidentally bringing a repo out of out-of-scope + if let SyncStatus::OutOfScope { since } = self.info.sync_status { + debug!(did = %self.did, ?reason, out_since = ?since, "ignoring desynchronize for out-of-scope repo"); + return; + } + let method = self.did().method(); // TODO: resync_error should never be None except on the first desync // ..we could probably make that a bit more type-explicit @@ -495,6 +508,19 @@ impl Repo { .insert(batch); } + /// like set_sync_status but when no previous status exists (create_resolved) + fn init_sync_status>( + &mut self, + initial: SyncStatus, + batch: &mut B, + ) { + let prev = self.info.clone(); + self.info.sync_status = initial; + QueuedResync::reconcile_for(&self.did, &prev, &self.info, batch); + SyncStatus::reconcile_for(None, Some(&self.info.sync_status), batch); + self.info.store(&self.did, batch); + } + pub fn set_sync_status>( &mut self, new_status: SyncStatus, @@ -502,27 +528,8 @@ impl Repo { ) -> SyncStatus { let prev = self.info.clone(); self.info.sync_status = new_status; - - // hack: non-persisted sync_status: SyncStatus::Synchronized in create_resolved QueuedResync::reconcile_for(&self.did, &prev, &self.info, batch); - - // hack: workaround for non-persisted create_resolved status (avoid decrement) - let is_being_created = matches!( - (&prev.sync_status, &self.info.sync_status), - ( - SyncStatus::Synchronized, - SyncStatus::Desynchronized(Desynchronized { - reason: DesyncReason::FirstSeen, - .. - }), - ) - ); - if is_being_created { - SyncStatus::reconcile_for(None, Some(&self.info.sync_status), batch); - } else { - SyncStatus::reconcile_for(Some(&prev.sync_status), Some(&self.info.sync_status), batch); - } - + SyncStatus::reconcile_for(Some(&prev.sync_status), Some(&self.info.sync_status), batch); self.info.store(&self.did, batch); prev.sync_status } @@ -530,15 +537,44 @@ impl Repo { pub fn set_identity>( &mut self, new_id: RepoIdentity, + now: SystemTime, 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(&self.did, batch); + self.reconcile_scope(now, batch); prev.identity } + // TODO: follow app style, don't do an ad-hoc manual reconcile helper here + fn reconcile_scope>( + &mut self, + now: SystemTime, + batch: &mut B, + ) { + match (&self.info.sync_status, self.is_in_scope()) { + (SyncStatus::OutOfScope { .. }, true) => { + // departing out-of-scope (coming into scope) + let status = SyncStatus::Desynchronized(Desynchronized::new( + DesyncReason::CameIntoScope, + now, + 0, + None, + )); + self.set_sync_status(status, batch); + } + (SyncStatus::OutOfScope { .. }, false) => {} // noop + (SyncStatus::Gone { .. }, _) => {} // leave gone + (_, false) => { + // became out of scope (leaving in-scope) + self.set_sync_status(SyncStatus::OutOfScope { since: now }, batch); + } + (_, true) => {} // in scope and tracked + } + } + // TODO: rename to "should" or something pub fn needs_identity_refresh(&self, now: SystemTime) -> bool { let age = now diff --git a/hubble-sync/src/sync_handle.rs b/hubble-sync/src/sync_handle.rs index 5daafef..2efa0c8 100644 --- a/hubble-sync/src/sync_handle.rs +++ b/hubble-sync/src/sync_handle.rs @@ -201,6 +201,7 @@ impl SyncEndpoints { use {AccountStatus as A, SyncStatus as S}; match (acct, sync) { (A::Active, S::Synchronized) => (true, None), + (A::Active, S::OutOfScope { .. }) => (true, Some("outOfScope")), (A::Active, S::Desynchronized(d)) if matches!(d.reason, DesyncReason::Throttled { .. }) => { diff --git a/hubble/src/generated_lexicons/blue_microcosm/hubble/get_repo_info.rs b/hubble/src/generated_lexicons/blue_microcosm/hubble/get_repo_info.rs index daa173c..9712778 100644 --- a/hubble/src/generated_lexicons/blue_microcosm/hubble/get_repo_info.rs +++ b/hubble/src/generated_lexicons/blue_microcosm/hubble/get_repo_info.rs @@ -386,6 +386,7 @@ pub enum SyncStateDesyncReason< S: jacquard_common::BosStr = jacquard_common::DefaultStr, > { FirstSeen, + CameIntoScope, UnresolvableIdentity, FirehoseCommitVerificationFail, FirehoseCommitFutureRev, @@ -402,6 +403,7 @@ impl SyncStateDesyncReason { pub fn as_str(&self) -> &str { match self { Self::FirstSeen => "firstSeen", + Self::CameIntoScope => "cameIntoScope", Self::UnresolvableIdentity => "unresolvableIdentity", Self::FirehoseCommitVerificationFail => "firehoseCommitVerificationFail", Self::FirehoseCommitFutureRev => "firehoseCommitFutureRev", @@ -420,6 +422,7 @@ impl SyncStateDesyncReason { pub fn from_value(s: S) -> Self { match s.as_ref() { "firstSeen" => Self::FirstSeen, + "cameIntoScope" => Self::CameIntoScope, "unresolvableIdentity" => Self::UnresolvableIdentity, "firehoseCommitVerificationFail" => Self::FirehoseCommitVerificationFail, "firehoseCommitFutureRev" => Self::FirehoseCommitFutureRev, @@ -483,6 +486,7 @@ where fn into_static(self) -> Self::Output { match self { SyncStateDesyncReason::FirstSeen => SyncStateDesyncReason::FirstSeen, + SyncStateDesyncReason::CameIntoScope => SyncStateDesyncReason::CameIntoScope, SyncStateDesyncReason::UnresolvableIdentity => { SyncStateDesyncReason::UnresolvableIdentity } @@ -517,6 +521,7 @@ where #[derive(Debug, Clone, PartialEq, Eq, Hash)] pub enum SyncStateState { Synchronized, + OutOfScope, Desynchronized, Pending, Other(S), @@ -526,6 +531,7 @@ impl SyncStateState { pub fn as_str(&self) -> &str { match self { Self::Synchronized => "synchronized", + Self::OutOfScope => "outOfScope", Self::Desynchronized => "desynchronized", Self::Pending => "pending", Self::Other(s) => s.as_ref(), @@ -535,6 +541,7 @@ impl SyncStateState { pub fn from_value(s: S) -> Self { match s.as_ref() { "synchronized" => Self::Synchronized, + "outOfScope" => Self::OutOfScope, "desynchronized" => Self::Desynchronized, "pending" => Self::Pending, _ => Self::Other(s), @@ -589,6 +596,7 @@ where fn into_static(self) -> Self::Output { match self { SyncStateState::Synchronized => SyncStateState::Synchronized, + SyncStateState::OutOfScope => SyncStateState::OutOfScope, SyncStateState::Desynchronized => SyncStateState::Desynchronized, SyncStateState::Pending => SyncStateState::Pending, SyncStateState::Other(v) => SyncStateState::Other(v.into_static()), diff --git a/hubble/src/serve/get_repo_info.rs b/hubble/src/serve/get_repo_info.rs index 0184211..64d3823 100644 --- a/hubble/src/serve/get_repo_info.rs +++ b/hubble/src/serve/get_repo_info.rs @@ -65,13 +65,15 @@ pub(super) async fn get_repo_info( )); } - let (sync_state_token, desync) = match ctx.sync_status() { + let (sync_state_name, desync) = match ctx.sync_status() { SyncStatus::Synchronized => (SyncStateState::Synchronized, None), + SyncStatus::OutOfScope { .. } => (SyncStateState::OutOfScope, None), SyncStatus::Desynchronized(d) if ctx.rev().is_some() => { (SyncStateState::Desynchronized, Some(d)) } // desynchronized with nothing synced yet: still working on the first copy SyncStatus::Desynchronized(d) => (SyncStateState::Pending, Some(d)), + // TODO: promote to known values or map to something else SyncStatus::Deactivated => (SyncStateState::Other("deactivated".into()), None), SyncStatus::Gone { .. } => (SyncStateState::Other("gone".into()), None), }; @@ -82,7 +84,7 @@ pub(super) async fn get_repo_info( rev: ctx .rev() .map(|rev| LexTid::new(rev.to_string()).expect("stored revs are valid TIDs")), - state: sync_state_token, + state: sync_state_name, extra_data: None, }; @@ -165,6 +167,7 @@ fn desync_reason(kind: DesyncReasonKind) -> SyncStateDesyncReason { use DesyncReasonKind as K; match kind { K::FirstSeen => SyncStateDesyncReason::FirstSeen, + K::CameIntoScope => SyncStateDesyncReason::CameIntoScope, K::UnresolvableIdentity => SyncStateDesyncReason::UnresolvableIdentity, K::FirehoseFail => SyncStateDesyncReason::FirehoseCommitVerificationFail, K::FutureRev => SyncStateDesyncReason::FirehoseCommitFutureRev, diff --git a/hubble/src/store/record.rs b/hubble/src/store/record.rs index cac0a9c..2ae07dc 100644 --- a/hubble/src/store/record.rs +++ b/hubble/src/store/record.rs @@ -2,7 +2,12 @@ //! //! keys are in the "records" column family: //! -//! `"r|" || || NUL || || ` +//! `"r|" || || NUL || || ` +//! +//! the generation marker u8 overflows if we ever do more than 256 resyncs of a +//! single repo, and this is fine -- there is only ever an `active` generation +//! and a `next` one at the next-ascending (mod overflow) generation value, so +//! they cannot collide. use hubble_sync::Did; use hubble_sync_rocksdb::prefix_end_exclusive; diff --git a/lexicons/blue_microcosm/hubble/getRepoInfo.json b/lexicons/blue_microcosm/hubble/getRepoInfo.json index 0178d2d..f5318f1 100644 --- a/lexicons/blue_microcosm/hubble/getRepoInfo.json +++ b/lexicons/blue_microcosm/hubble/getRepoInfo.json @@ -65,13 +65,19 @@ "properties": { "state": { "type": "string", - "knownValues": ["synchronized", "desynchronized", "pending"] + "knownValues": [ + "synchronized", + "outOfScope", + "desynchronized", + "pending" + ] }, "rev": { "type": "string", "format": "tid" }, "desyncReason": { "type": "string", "knownValues": [ "firstSeen", + "cameIntoScope", "unresolvableIdentity", "firehoseCommitVerificationFail", "firehoseCommitFutureRev",