From c3f05f1010e75cd37071eaef50aadea8cc844e5c Mon Sep 17 00:00:00 2001 From: phil Date: Tue, 14 Jul 2026 13:16:50 -0400 Subject: [PATCH] pull engine key prefixing into engine wrapper so we can share an engine with the app, but hubble-sync layer prefixing on top of that instead of providing a different way in --- hubble-sync-fjall/examples/firehose_smoke.rs | 10 +- hubble-sync-fjall/src/lib.rs | 65 ++--- hubble-sync-rocksdb/src/lib.rs | 99 +++----- hubble-sync/readme.md | 9 +- hubble-sync/src/config.rs | 10 + hubble-sync/src/firehose/subscribe_repos.rs | 6 +- hubble-sync/src/hubble_sync.rs | 7 +- hubble-sync/src/lib.rs | 4 +- hubble-sync/src/pending_identity_scheduler.rs | 5 +- hubble-sync/src/repo_actor/actor.rs | 16 +- .../src/repo_actor/identity_initial.rs | 4 +- hubble-sync/src/repo_actor/repo_registry.rs | 12 +- .../src/repo_actor/task_processor/mod.rs | 13 +- hubble-sync/src/resync_scheduler.rs | 4 +- hubble-sync/src/storage/engine/mem.rs | 12 + hubble-sync/src/storage/engine/mod.rs | 2 + hubble-sync/src/storage/engine/prefixed.rs | 224 ++++++++++++++++++ hubble-sync/src/sync_handle.rs | 29 ++- hubble/readme.md | 7 + hubble/src/sync.rs | 4 +- 20 files changed, 369 insertions(+), 173 deletions(-) create mode 100644 hubble-sync/src/storage/engine/prefixed.rs diff --git a/hubble-sync-fjall/examples/firehose_smoke.rs b/hubble-sync-fjall/examples/firehose_smoke.rs index 9d022a5..1d4a70f 100644 --- a/hubble-sync-fjall/examples/firehose_smoke.rs +++ b/hubble-sync-fjall/examples/firehose_smoke.rs @@ -27,7 +27,7 @@ use hubble_sync::{ AccountStatus, AppResult, Commit, ConsumerAppError, Did, HubbleSync, RepoContext, ResyncData, SimpleSyncConsumer, StorageEngine, SyncConfig, }; -use hubble_sync_fjall::{FjallEngine, FjallEngineConfig}; +use hubble_sync_fjall::FjallEngine; use metrics_exporter_prometheus::{Matcher, PrometheusBuilder}; use tracing_subscriber::EnvFilter; @@ -132,7 +132,7 @@ async fn main() -> Result<(), Box> { start_metrics(); // Temp fjall database, cleaned up on drop. - let db = if let Some(path) = std::env::var("HUBBLE_DB_PATH").ok() { + let db = if let Ok(path) = std::env::var("HUBBLE_DB_PATH") { Database::builder(path) } else { Database::builder(std::env::temp_dir().join("hubble-sync-fjall-smoke")).temporary(true) @@ -140,11 +140,11 @@ async fn main() -> Result<(), Box> { .open()?; let point = db.keyspace("point", KeyspaceCreateOptions::default)?; let queue = db.keyspace("queue", KeyspaceCreateOptions::default)?; - let storage = FjallEngine::new(db, point, queue, &FjallEngineConfig::default()); + let storage = FjallEngine::new(db, point, queue); let mut config = SyncConfig::default(); if let Ok(relay) = std::env::var("HUBBLE_SUBSCRIBE") { - config.subscribe_host = relay; + config.upstream.hostname = relay; } let duration_secs: u64 = std::env::var("HUBBLE_DURATION_SECS") .ok() @@ -152,7 +152,7 @@ async fn main() -> Result<(), Box> { .unwrap_or(10); tracing::info!( - subscribe_host = %config.subscribe_host, + subscribe_host = %config.upstream.hostname, duration_secs, "hubble-sync smoke test starting" ); diff --git a/hubble-sync-fjall/src/lib.rs b/hubble-sync-fjall/src/lib.rs index daccd37..e801574 100644 --- a/hubble-sync-fjall/src/lib.rs +++ b/hubble-sync-fjall/src/lib.rs @@ -1,8 +1,9 @@ //! fjall backend for hubble-sync //! -//! app provides keyspaces (one point-oriented, one queue-oriented), which it's -//! allowed to share: -//! - this engine will only ever read/write keys under its configured prefix +//! app provides keyspaces (one point-oriented, one queue-oriented, one for +//! counts), which it's allowed to share: +//! +//! - hubble-sync will only ever read/write keys under its configured prefix //! - the app MUST NOT ever write its own keys under that prefix use fjall::util::prefixed_range; @@ -15,36 +16,20 @@ pub struct FjallError(#[from] fjall::Error); impl StorageError for FjallError {} -/// key prefix that hubble-sync's data will be written under, for both -/// keyspaces. if sharing a keyspace with app data, you MUST NEVER write -/// keys under hubble-sync's prefix. -pub const DEFAULT_SYNC_PREFIX: &[u8] = &[0x00]; - /// Fjall storage engine implementation for hubble-sync. #[derive(Clone)] pub struct FjallEngine { db: Database, point: Keyspace, queue: Keyspace, - prefix: &'static [u8], -} - -#[derive(Debug, Default)] -pub struct FjallEngineConfig { - pub prefix: Option<&'static [u8]>, } type FjallResult = Result; type FjallScanItem = FjallResult<(Vec, Vec)>; impl FjallEngine { - pub fn new(db: Database, point: Keyspace, queue: Keyspace, config: &FjallEngineConfig) -> Self { - Self { - db, - point, - queue, - prefix: config.prefix.unwrap_or(DEFAULT_SYNC_PREFIX), - } + pub fn new(db: Database, point: Keyspace, queue: Keyspace) -> Self { + Self { db, point, queue } } pub fn scan_keyspace( @@ -53,13 +38,12 @@ impl FjallEngine { scan_prefix: &[u8], from_suffix: &[u8], ) -> Box + '_> { - let full_prefix = prefixed(self.prefix, scan_prefix); - let full_prefix_len = full_prefix.len(); - let range = prefixed_range(full_prefix, from_suffix..); + let prefix_len = scan_prefix.len(); + let range = prefixed_range(scan_prefix, from_suffix..); Box::new(ks.range(range).map(move |guard| { let (k, v) = guard.into_inner()?; // STRIPS THE FULL PREFIX! - Ok((k.as_ref()[full_prefix_len..].to_vec(), v.as_ref().to_vec())) + Ok((k.as_ref()[prefix_len..].to_vec(), v.as_ref().to_vec())) })) } } @@ -73,22 +57,15 @@ impl StorageEngine for FjallEngine { inner: self.db.batch(), point: self.point.clone(), queue: self.queue.clone(), - prefix: self.prefix, } } fn get(&self, k: &[u8]) -> FjallResult>> { - Ok(self - .point - .get(prefixed(self.prefix, k))? - .map(|s| s.as_ref().to_vec())) + Ok(self.point.get(k)?.map(|s| s.as_ref().to_vec())) } fn get_queue(&self, k: &[u8]) -> FjallResult>> { - Ok(self - .queue - .get(prefixed(self.prefix, k))? - .map(|s| s.as_ref().to_vec())) + Ok(self.queue.get(k)?.map(|s| s.as_ref().to_vec())) } fn get_counter(&self, key: &[u8]) -> FjallResult { @@ -120,13 +97,10 @@ pub struct FjallBatch { inner: OwnedWriteBatch, point: Keyspace, queue: Keyspace, - prefix: &'static [u8], } impl FjallBatch { /// directly access the underlying batch - /// - /// unlike the trait methods, does not add the hubble-sync prefix to keys! pub fn inner_mut(&mut self) -> &mut OwnedWriteBatch { &mut self.inner } @@ -139,32 +113,25 @@ impl StorageBatch for FjallBatch { } fn put(&mut self, k: &[u8], v: &[u8]) { - self.inner.insert(&self.point, prefixed(self.prefix, k), v); + self.inner.insert(&self.point, k, v); } fn put_queue(&mut self, k: &[u8], v: &[u8]) { - self.inner.insert(&self.queue, prefixed(self.prefix, k), v); + self.inner.insert(&self.queue, k, v); } fn delete(&mut self, k: &[u8]) { - self.inner.remove(&self.point, prefixed(self.prefix, k)); + self.inner.remove(&self.point, k); } fn delete_queue(&mut self, k: &[u8]) { - self.inner.remove(&self.queue, prefixed(self.prefix, k)); + self.inner.remove(&self.queue, k); } fn put_counter(&mut self, _k: &[u8], _n: i64) {} fn delete_counter(&mut self, _k: &[u8]) {} fn increment_counter(&mut self, _k: &[u8], _d: i64) {} } -fn prefixed(prefix: &[u8], k: &[u8]) -> Vec { - let mut out = Vec::with_capacity(prefix.len() + k.len()); - out.extend_from_slice(prefix); - out.extend_from_slice(k); - out -} - #[cfg(test)] mod tests { use fjall::KeyspaceCreateOptions; @@ -191,7 +158,7 @@ mod tests { let queue = db .keyspace("queue", KeyspaceCreateOptions::default) .expect("open queue keyspace"); - FjallEngine::new(db, point, queue, &FjallEngineConfig::default()) + FjallEngine::new(db, point, queue) } #[test] diff --git a/hubble-sync-rocksdb/src/lib.rs b/hubble-sync-rocksdb/src/lib.rs index 0414504..5b064ba 100644 --- a/hubble-sync-rocksdb/src/lib.rs +++ b/hubble-sync-rocksdb/src/lib.rs @@ -1,8 +1,9 @@ //! rocksdb backend for hubble-sync //! -//! app provides column families (one point-oriented, one queue), which it's -//! allowed to share: -//! - this engine will only ever read/write keys under its configured prefix +//! app provides column families (one point-oriented, one queue, one counter), +//! which it's allowed to share: +//! +//! - hubble-sync will only ever read/write keys under its configured prefix //! - the app MUST NOT ever write its own keys under that prefix use std::path::Path; @@ -31,29 +32,23 @@ impl StorageError for Error {} /// point-read-optimized column family for hubble-sync to use. /// it's fine (+ recommended) to share this CF with point-oriented app data, but -/// you MUST ensure that app data NEVER writes keys under `sync_prefix`. +/// apps must never write keys under `hubble_sync_storage_prefix`. pub const DEFAULT_POINT_CF: &str = "default"; /// queue-optimized column family for hubble-sync to use. /// it's fine (+ recommended) to share this CF with queue-oriented app data, but -/// you MUST ensure that app data NEVER writes keys under `sync_prefix`. +/// apps must never write keys under `hubble_sync_storage_prefix`. pub const DEFAULT_QUEUE_CF: &str = "queue"; /// merge-operator CF for read-free counters for hubble-sync to use. /// it's fine (+ recommended) to share this CF with any counters you might use -/// in your app, but you MUST NOT write keys under `sync_prefix`. +/// in your app, but apps must never write under `hubble_sync_storage_prefix`. pub const DEFAULT_COUNT_CF: &str = "counter"; -/// key prefix that hubble-sync's data will be written under, for both column -/// families. if sharing a column family for app data, you MUST NEVER write keys -/// under hubble-sync's prefix. -pub const DEFAULT_SYNC_PREFIX: &[u8] = &[0x00]; - /// RocksDB storage engine implementation for hubble-sync #[derive(Debug, Clone)] pub struct RocksEngine { db: Arc, - prefix: &'static [u8], point_cf: &'static str, queue_cf: &'static str, counter_cf: &'static str, @@ -61,7 +56,6 @@ pub struct RocksEngine { #[derive(Debug, Default)] pub struct RocksEngineConfig { - pub prefix: Option<&'static [u8]>, pub point_cf: Option<&'static str>, pub queue_cf: Option<&'static str>, pub counter_cf: Option<&'static str>, @@ -85,24 +79,14 @@ impl Default for EasySetup { impl RocksEngine { /// simple mode: let the engine set everything up and hand it to you - pub fn new_easy( - path: impl AsRef, - prefix: &'static [u8], - setup: &EasySetup, - ) -> RocksResult<(Self, Arc)> { + pub fn new_easy(path: impl AsRef, setup: &EasySetup) -> RocksResult<(Self, Arc)> { let cache = Cache::new_lru_cache(setup.block_cache_mb * 2_usize.pow(20)); let cfs = recommended_cf_descriptors(&cache); let mut db_opts = Options::default(); db_opts.create_if_missing(true); db_opts.create_missing_column_families(true); let db = Arc::new(DB::open_cf_descriptors(&db_opts, path, cfs)?); - let engine = Self::new( - db.clone(), - &RocksEngineConfig { - prefix: Some(prefix), - ..Default::default() - }, - )?; + let engine = Self::new(db.clone(), &RocksEngineConfig::default())?; Ok((engine, db)) } @@ -116,7 +100,6 @@ impl RocksEngine { .ok_or(Error::MissingCf(counter_cf))?; Ok(Self { db, - prefix: config.prefix.unwrap_or(DEFAULT_SYNC_PREFIX), point_cf, queue_cf, counter_cf, @@ -129,25 +112,25 @@ impl RocksEngine { scan_prefix: &[u8], from_suffix: &[u8], ) -> Box + '_> { - let full_prefix = prefixed(self.prefix, scan_prefix); - - let mut start = Vec::with_capacity(full_prefix.len() + from_suffix.len()); - start.extend_from_slice(&full_prefix); + let mut start = Vec::with_capacity(scan_prefix.len() + from_suffix.len()); + start.extend_from_slice(scan_prefix); start.extend_from_slice(from_suffix); let mut opts = ReadOptions::default(); opts.set_iterate_lower_bound(start); - if let Some(end) = prefix_end_exclusive(&full_prefix) { + if let Some(end) = prefix_end_exclusive(scan_prefix) { opts.set_iterate_upper_bound(end); } + // owned: the returned iterator outlives this call + let scan_prefix = scan_prefix.to_vec(); let iter = self .db .iterator_cf_opt(&cf, opts, IteratorMode::Start) .map(move |r| { r.map(|(k, v)| { - // STRIPS THE FULL PREFIX! - let k = unprefixed(&full_prefix, &k).to_vec(); + // STRIPS THE PREFIX! + let k = unprefixed(&scan_prefix, &k); let v = v.into_vec(); (k, v) }) @@ -169,18 +152,17 @@ impl StorageEngine for RocksEngine { point_cf: self.point_cf, queue_cf: self.queue_cf, counter_cf: self.counter_cf, - prefix: self.prefix, } } fn get(&self, k: &[u8]) -> RocksResult>> { let cf = self.db.cf_handle(self.point_cf).expect("point cf to exist"); - Ok(self.db.get_cf(&cf, prefixed(self.prefix, k))?) + Ok(self.db.get_cf(&cf, k)?) } fn get_queue(&self, k: &[u8]) -> RocksResult>> { let cf = self.db.cf_handle(self.queue_cf).expect("queue cf to exist"); - Ok(self.db.get_cf(&cf, prefixed(self.prefix, k))?) + Ok(self.db.get_cf(&cf, k)?) } fn get_counter(&self, k: &[u8]) -> RocksResult { @@ -188,7 +170,7 @@ impl StorageEngine for RocksEngine { .db .cf_handle(self.counter_cf) .expect("counter cf to exist"); - let Some(bytes) = self.db.get_cf(&cf, prefixed(self.prefix, k))? else { + let Some(bytes) = self.db.get_cf(&cf, k)? else { return Ok(0); }; let Ok(arr) = bytes.try_into() else { @@ -230,7 +212,6 @@ impl StorageEngine for RocksEngine { pub struct RocksBatch { db: Arc, inner: WriteBatch, - prefix: &'static [u8], point_cf: &'static str, queue_cf: &'static str, counter_cf: &'static str, @@ -252,57 +233,43 @@ impl StorageBatch for RocksBatch { } fn put(&mut self, k: &[u8], v: &[u8]) { let cf = self.db.cf_handle(self.point_cf).expect("point CF missing"); - let key = prefixed(self.prefix, k); - self.inner.put_cf(&cf, key, v); + self.inner.put_cf(&cf, k, v); } fn put_queue(&mut self, k: &[u8], v: &[u8]) { let cf = self.db.cf_handle(self.queue_cf).expect("queue CF missing"); - let key = prefixed(self.prefix, k); - self.inner.put_cf(&cf, key, v); + self.inner.put_cf(&cf, k, v); } fn delete(&mut self, k: &[u8]) { let cf = self.db.cf_handle(self.point_cf).expect("point CF missing"); - let key = prefixed(self.prefix, k); - self.inner.delete_cf(&cf, key); + self.inner.delete_cf(&cf, k); } fn delete_queue(&mut self, k: &[u8]) { let cf = self.db.cf_handle(self.queue_cf).expect("queue CF missing"); - let key = prefixed(self.prefix, k); - self.inner.delete_cf(&cf, key); + self.inner.delete_cf(&cf, k); } fn put_counter(&mut self, k: &[u8], v: i64) { let cf = self .db .cf_handle(self.counter_cf) .expect("counter CF missing"); - let key = prefixed(self.prefix, k); - self.inner.put_cf(&cf, key, v.to_be_bytes()); + self.inner.put_cf(&cf, k, v.to_be_bytes()); } fn delete_counter(&mut self, k: &[u8]) { let cf = self .db .cf_handle(self.counter_cf) .expect("counter CF missing"); - let key = prefixed(self.prefix, k); - self.inner.delete_cf(&cf, key); + self.inner.delete_cf(&cf, k); } fn increment_counter(&mut self, k: &[u8], v: i64) { let cf = self .db .cf_handle(self.counter_cf) .expect("counter CF missing"); - let key = prefixed(self.prefix, k); - self.inner.merge_cf(&cf, key, v.to_be_bytes()); + self.inner.merge_cf(&cf, k, v.to_be_bytes()); } } -fn prefixed(prefix: &[u8], k: &[u8]) -> Vec { - let mut out = Vec::with_capacity(prefix.len() + k.len()); - out.extend_from_slice(prefix); - out.extend_from_slice(k); - out -} - fn unprefixed(prefix: &[u8], k: &[u8]) -> Vec { assert!( k.starts_with(prefix), @@ -408,12 +375,8 @@ mod tests { /// unwind). the `TempDir` is returned so the caller keeps it alive. fn temp_engine() -> (RocksEngine, Arc, tempfile::TempDir) { let dir = tempfile::tempdir().expect("tempdir"); - let (eng, db) = RocksEngine::new_easy( - dir.path(), - DEFAULT_SYNC_PREFIX, - &EasySetup { block_cache_mb: 8 }, - ) - .expect("open rocks"); + let (eng, db) = RocksEngine::new_easy(dir.path(), &EasySetup { block_cache_mb: 8 }) + .expect("open rocks"); (eng, db, dir) } @@ -544,12 +507,10 @@ mod tests { #[test] fn get_counter_errors_on_corrupt_value() { - // a non-8-byte value under our prefix is corruption; the read fails loud. + // a non-8-byte counter value is corruption; the read fails loud. let (eng, db, _dir) = temp_engine(); let cf = db.cf_handle(DEFAULT_COUNT_CF).expect("counter cf"); - let mut key = DEFAULT_SYNC_PREFIX.to_vec(); - key.extend_from_slice(b"k"); - db.put_cf(&cf, key, b"not-8-bytes").unwrap(); + db.put_cf(&cf, b"k", b"not-8-bytes").unwrap(); assert!(matches!( eng.get_counter(b"k").unwrap_err(), Error::Integrity { .. } diff --git a/hubble-sync/readme.md b/hubble-sync/readme.md index 8166153..067e4d1 100644 --- a/hubble-sync/readme.md +++ b/hubble-sync/readme.md @@ -233,7 +233,7 @@ a jetstream impl based on hubble-sync might add a new emitted event type (somewh ## actual todos - [x] be nice to all pdses except mushrooms (10+req/sec + high concurrency kills everyone else) - - [ ] ..triple check this + - [x] ..triple check this - [x] we're re-dispatching resyncs per host many times, not advancing to that host's next item - [x] we don't have repo discovery / crawling yet - [x] we don't have local-inactive state tracking yet (distinguishable from upstream) @@ -242,15 +242,15 @@ a jetstream impl based on hubble-sync might add a new emitted event type (somewh - [x] fixed via timeout on acquisition -- retry preacquires - [ ] add DesyncReason::NotFoundUpstream -- [ ] we're getting bogged down on big-mem token acquisition +- [x] we're getting bogged down on big-mem token acquisition - [x] for pds-upstream: need direct upstream getRepo, bc otherwise people who have migrated get requested from other PDSes get resynced from their new place (we only want to sync this pds) - [x] exclude any non-`active` repos from the upstream crawl (though maybe it would be nice to record their account status?) - [ ] should check upstream for account status? eg., from crawl-discovered repos + - this is slightly deferrable since we request from upstream now -- [ ] fix public storage engine apis to not access under hubble's prefix (eek) -- [ ] commit: replay-vs-live flag +- [x] fix public storage engine apis to not access under hubble's prefix (eek) - [ ] guard against ssrf - [ ] sync new repos directly (no-resync) from firehose when they start from definitely-empty @@ -261,3 +261,4 @@ a jetstream impl based on hubble-sync might add a new emitted event type (somewh - [ ] figure out output compression (just do it in the reverse proxy?) - [ ] cache CARs (reverse-proxy with rev-checks reaching back to us?) - [ ] add plc export stream; plc tombstoning +- [ ] commit: replay-vs-live flag diff --git a/hubble-sync/src/config.rs b/hubble-sync/src/config.rs index ab8c68b..328a66d 100644 --- a/hubble-sync/src/config.rs +++ b/hubble-sync/src/config.rs @@ -86,6 +86,15 @@ pub struct SyncConfig { /// /// default: 10s pub reactive_permit_wait_timeout: Duration, + /// key prefix for hubble-sync's state, in every CF/keyspace provided + /// + /// must be exclusive: your app must never write keys under this prefix. you + /// can use `PrefixedEngine` to make sure all your keys are under a + /// different prefix, or use separate column families/keyspaces entirely, or + /// just use care not to trample. + /// + /// default: `[0x00]` + pub hubble_sync_storage_prefix: &'static [u8], /// upstream (subscribeRepos source) config pub upstream: UpstreamConfig, /// firehose subscriber tuning @@ -107,6 +116,7 @@ impl Default for SyncConfig { pending_identity_queue_limit: 32_768.try_into().unwrap(), scheduled_resolve_limit: 3.try_into().unwrap(), reactive_permit_wait_timeout: Duration::from_secs(10), + hubble_sync_storage_prefix: &[0x00], upstream: UpstreamConfig::default(), firehose: FirehoseConfig::default(), hosts: HostRegistryConfig::default(), diff --git a/hubble-sync/src/firehose/subscribe_repos.rs b/hubble-sync/src/firehose/subscribe_repos.rs index bada6b6..2eb3d60 100644 --- a/hubble-sync/src/firehose/subscribe_repos.rs +++ b/hubble-sync/src/firehose/subscribe_repos.rs @@ -33,7 +33,7 @@ use crate::repo_actor::{RepoMessage, RepoSendError}; use crate::storage::engine::{StorageBatch, StorageError}; use crate::storage::firehose_cursor::CursorState; use crate::storage::{LoadError, StorageEngine}; -use crate::{CancelExt, Did, Host, RepoRegistry, SyncConsumer}; +use crate::{CancelExt, Did, Host, PrefixedEngine, RepoRegistry, SyncConsumer}; #[derive(Debug, Clone)] pub struct FirehoseConfig { @@ -129,7 +129,7 @@ struct SubProgress { pub struct FirehoseSubscriber, R: Resolve> { host: Arc, registry: Arc>, - storage: S, + storage: PrefixedEngine, config: FirehoseConfig, } @@ -137,7 +137,7 @@ impl, R: Resolve> FirehoseSubscrib pub fn new( host: Arc, registry: Arc>, - storage: S, + storage: PrefixedEngine, config: FirehoseConfig, ) -> Self { Self { diff --git a/hubble-sync/src/hubble_sync.rs b/hubble-sync/src/hubble_sync.rs index b7e87f2..6cfd452 100644 --- a/hubble-sync/src/hubble_sync.rs +++ b/hubble-sync/src/hubble_sync.rs @@ -18,7 +18,7 @@ use crate::resync_scheduler; use crate::{ CrawlState, FirehoseError, FirehoseSubscriber, HostRegistry, HubbleSyncResolver, LoadError, - PendingScheduler, RepoCountsByState, RepoCountsByStatus, RepoRegistry, Resolve, + PendingScheduler, PrefixedEngine, RepoCountsByState, RepoCountsByStatus, RepoRegistry, Resolve, ResyncScheduler, StorageEngine, StorageError, SyncConfig, SyncConsumer, SyncHandle, }; @@ -51,7 +51,7 @@ where R: Resolve, { cancel: CancellationToken, - storage: S, + storage: PrefixedEngine, hosts: Arc, repo_registry: Arc>, resync_scheduler: Arc, @@ -109,6 +109,7 @@ where cancel: CancellationToken, config: &SyncConfig, ) -> Self { + let storage = PrefixedEngine::new(storage, config.hubble_sync_storage_prefix); let resync_scheduler = Arc::new(ResyncScheduler::new()); let pending_scheduler = Arc::new(PendingScheduler::new(config.pending_identity_queue_limit)); @@ -140,7 +141,7 @@ where self } - pub fn crawl_state(&self, strategy_id: &str) -> CrawlState { + pub fn crawl_state(&self, strategy_id: &str) -> CrawlState> { CrawlState::new(self.storage.clone(), strategy_id) } diff --git a/hubble-sync/src/lib.rs b/hubble-sync/src/lib.rs index 57f4797..2d54719 100644 --- a/hubble-sync/src/lib.rs +++ b/hubble-sync/src/lib.rs @@ -46,7 +46,9 @@ pub use repo_actor::{ pub use repo_slot::{RepoSlot, RepoSlots, Slot, SlotDecodeError}; pub use resync::ResyncData; pub use resync_scheduler::ResyncScheduler; -pub use storage::engine::{StorageBatch, StorageEngine, StorageError}; +pub use storage::engine::{ + PrefixedBatch, PrefixedEngine, StorageBatch, StorageEngine, StorageError, +}; pub use storage::repo::{ AccountStatus, AccountStatusEvent, AccountStatusKind, ModAction, Repo, RepoCountsByState, RepoCountsByStatus, diff --git a/hubble-sync/src/pending_identity_scheduler.rs b/hubble-sync/src/pending_identity_scheduler.rs index 0fe0b77..a2bf4b5 100644 --- a/hubble-sync/src/pending_identity_scheduler.rs +++ b/hubble-sync/src/pending_identity_scheduler.rs @@ -38,7 +38,8 @@ use crate::metrics::{ use crate::repo_actor::{RepoMessage, RepoRegistry}; use crate::storage::repo::PendingIdentityQueueEntry; use crate::{ - CancelExt, Did, LoadError, RepoSendError, Resolve, StorageBatch, StorageEngine, SyncConsumer, + CancelExt, Did, LoadError, PrefixedEngine, RepoSendError, Resolve, StorageBatch, StorageEngine, + SyncConsumer, }; const IDLE_POLL: Duration = Duration::from_secs(15); @@ -319,7 +320,7 @@ pub fn bootstrap( pub async fn drive, R: Resolve>( scheduler: Arc, resolve_limit: Arc, - storage: S, + storage: PrefixedEngine, repo_registry: Arc>, cancel: CancellationToken, ) -> Result<(), LoadError> { diff --git a/hubble-sync/src/repo_actor/actor.rs b/hubble-sync/src/repo_actor/actor.rs index f15182d..4eebe8b 100644 --- a/hubble-sync/src/repo_actor/actor.rs +++ b/hubble-sync/src/repo_actor/actor.rs @@ -67,7 +67,9 @@ use crate::firehose::FirehoseEvent; use crate::identity::{Resolve, ResolvedIdentity}; use crate::storage::StorageEngine; use crate::storage::repo::{Awoken, Repo}; -use crate::{Did, Host, HostRegistry, ModerateEvent, ModerateOutcome, SyncConsumer}; +use crate::{ + Did, Host, HostRegistry, ModerateEvent, ModerateOutcome, PrefixedEngine, SyncConsumer, +}; /// some work the system can request the repo actor to do /// @@ -259,7 +261,7 @@ impl Intake { } pub(crate) struct RepoActorContext, R: Resolve> { - pub storage: S, + pub storage: PrefixedEngine, pub hosts: Arc, pub upstream: Arc, pub resolver: Arc, @@ -325,7 +327,7 @@ impl, R: Resolve> RepoActor(&storage, &hosts, did, now) + Repo::wake::<_, A::CommitState, A::InfoState>(&storage, &hosts, did, now) }) .await .expect("repo wake task not to panic") @@ -588,14 +590,14 @@ mod tests { } fn setup() -> ( - MemEngine, + PrefixedEngine, Arc, Arc, Arc, ) { let hosts = HostRegistry::new_default(); ( - MemEngine::new(), + MemEngine::new_prefixed(), hosts.clone(), Arc::new(HubbleSyncResolver::new( "@bad-example.com", @@ -610,7 +612,7 @@ mod tests { /// with a fresh resolved identity. Without this, the bootstrap path /// fires `ResolveIdentity`, which attempts a real HTTP call we can't /// service in unit tests (and would blow past the test's yield deadline). - fn install_resolved_repo(storage: &MemEngine, hosts: &HostRegistry, did: &Did) { + fn install_resolved_repo(storage: &PrefixedEngine, hosts: &HostRegistry, did: &Did) { let pds = hosts.get("pds.example.com").expect("interned"); let now = SystemTime::now(); let info = RepoInfo { @@ -659,7 +661,7 @@ mod tests { } fn test_context, R: Resolve>( - storage: S, + storage: PrefixedEngine, hosts: Arc, resolver: Arc, consumer: Arc, diff --git a/hubble-sync/src/repo_actor/identity_initial.rs b/hubble-sync/src/repo_actor/identity_initial.rs index 5814abb..76fc0aa 100644 --- a/hubble-sync/src/repo_actor/identity_initial.rs +++ b/hubble-sync/src/repo_actor/identity_initial.rs @@ -13,14 +13,14 @@ use tokio::task::spawn_blocking; use tokio_util::sync::CancellationToken; use crate::identity::Resolve; +use crate::storage::engine::{PrefixedEngine, StorageBatch, StorageEngine}; use crate::storage::repo::{PendingIdentity, Repo}; -use crate::storage::{StorageEngine, engine::StorageBatch}; use crate::{Host, SyncConsumer}; use super::{ProcessError, TaskProcessor}; pub(super) struct InitialResolve, R: Resolve> { - pub(super) storage: S, + pub(super) storage: PrefixedEngine, pub(super) resolver: Arc, pub(super) pending: PendingIdentity, pub(super) upstream: Arc, diff --git a/hubble-sync/src/repo_actor/repo_registry.rs b/hubble-sync/src/repo_actor/repo_registry.rs index e516590..7b78d8c 100644 --- a/hubble-sync/src/repo_actor/repo_registry.rs +++ b/hubble-sync/src/repo_actor/repo_registry.rs @@ -27,7 +27,7 @@ use crate::metrics::{ REPO_REGISTRY_ACTIVE, REPO_REGISTRY_EVICTION_STUCK_TOTAL, REPO_REGISTRY_EVICTIONS_TOTAL, REPO_REGISTRY_LOADS_TOTAL, }; -use crate::{Did, Host, HostRegistry, StorageEngine, SyncConfig, SyncConsumer}; +use crate::{Did, Host, HostRegistry, PrefixedEngine, StorageEngine, SyncConfig, SyncConsumer}; #[derive(Debug, thiserror::Error)] pub enum RepoSendError { @@ -71,7 +71,7 @@ pub trait RepoSender: Send + Sync { } pub struct RepoRegistry, R: Resolve> { - storage: S, + storage: PrefixedEngine, hosts: Arc, upstream: Arc, resolver: Arc, @@ -99,7 +99,7 @@ impl, R: Resolve> fmt::Debug impl, R: Resolve> RepoRegistry { pub fn new( - storage: S, + storage: PrefixedEngine, hosts: Arc, resolver: Arc, consumer_app: Arc, @@ -396,10 +396,10 @@ mod tests { max_repo_actors: usize, ) -> ( RepoRegistry, - MemEngine, + PrefixedEngine, Arc, ) { - let storage = MemEngine::new(); + let storage = MemEngine::new_prefixed(); let hosts = HostRegistry::new_default(); let plc_url = "https://plc.directory"; let resolver = Arc::new(HubbleSyncResolver::new( @@ -431,7 +431,7 @@ mod tests { /// with a fresh resolved identity. Without this, the bootstrap path /// fires `ResolveIdentity`, which attempts a real HTTP call we can't /// service in unit tests. - fn install_resolved_repo(storage: &MemEngine, hosts: &HostRegistry, did: &Did) { + fn install_resolved_repo(storage: &PrefixedEngine, hosts: &HostRegistry, did: &Did) { let pds = hosts.get("pds.example.com").expect("interned"); let now = SystemTime::now(); let info = RepoInfo { diff --git a/hubble-sync/src/repo_actor/task_processor/mod.rs b/hubble-sync/src/repo_actor/task_processor/mod.rs index 362140a..6b18768 100644 --- a/hubble-sync/src/repo_actor/task_processor/mod.rs +++ b/hubble-sync/src/repo_actor/task_processor/mod.rs @@ -28,7 +28,7 @@ use crate::metrics::{ use crate::resync::{ BigRepoPermits, ResyncData, ResyncError, Resyncable, TransientResyncError, load_repo, }; -use crate::storage::engine::{StorageBatch, StorageEngine}; +use crate::storage::engine::{PrefixedBatch, PrefixedEngine, StorageBatch, StorageEngine}; use crate::storage::repo::{ AccountStatus, AccountStatusEvent, AccountStatusUpstream, DesyncReason, Repo, ResyncInfo, }; @@ -179,7 +179,7 @@ impl<'a> ResyncContext<'a> { } pub(super) struct TaskProcessor, R: Resolve> { - pub(super) storage: S, + pub(super) storage: PrefixedEngine, pub(super) resolver: Arc, pub(super) consumer_app: Arc, pub(super) cancel: CancellationToken, @@ -314,7 +314,7 @@ impl, R: Resolve> TaskProcessor { error!("handle consumer app fatal error"); @@ -746,7 +746,7 @@ impl, R: Resolve> TaskProcessor { error!(%err, "consumer app's apply_resync"); @@ -1065,7 +1065,7 @@ impl, R: Resolve> TaskProcessor(&mut self, now: SystemTime, mutate: F) -> PEResult where - F: FnOnce(&mut Repo, &mut S::Batch) + Send + 'static, + F: FnOnce(&mut Repo, &mut PrefixedBatch) + Send + 'static, { let storage = self.storage.clone(); let consumer_app = self.consumer_app.clone(); @@ -1095,7 +1095,8 @@ impl, R: Resolve> TaskProcessor {} Err(ConsumerAppError::Desynchronize { reason }) => { let reason = DesyncReason::AppRequested { diff --git a/hubble-sync/src/resync_scheduler.rs b/hubble-sync/src/resync_scheduler.rs index 08763ef..6de20fa 100644 --- a/hubble-sync/src/resync_scheduler.rs +++ b/hubble-sync/src/resync_scheduler.rs @@ -48,7 +48,7 @@ use crate::metrics::{ }; use crate::repo_actor::{RepoMessage, RepoRegistry}; use crate::storage::repo::NextQueuedByHost; -use crate::{CancelExt, Host, Hostname, LoadError, Resolve, StorageEngine}; +use crate::{CancelExt, Host, Hostname, LoadError, PrefixedEngine, Resolve, StorageEngine}; use crate::{Did, HostRegistry, SyncConsumer}; const IDLE_POLL: Duration = Duration::from_secs(1); @@ -336,7 +336,7 @@ async fn rebootstrap_if_due( /// TODO: probably goes on ResyncScheduler impl pub async fn drive, R: Resolve>( scheduler: Arc, - storage: S, + storage: PrefixedEngine, host_registry: Arc, repo_registry: Arc>, dispatch_qps: NonZeroU32, diff --git a/hubble-sync/src/storage/engine/mem.rs b/hubble-sync/src/storage/engine/mem.rs index c6e2412..e74bb12 100644 --- a/hubble-sync/src/storage/engine/mem.rs +++ b/hubble-sync/src/storage/engine/mem.rs @@ -6,10 +6,16 @@ use std::collections::BTreeMap; use std::sync::{Arc, Mutex}; +use super::prefixed::PrefixedEngine; use super::types::{StorageBatch, StorageEngine, StorageError}; type Pair = (Vec, Vec); +/// prefix for tests exercising hubble-sync internals through a +/// [`PrefixedEngine`] -- deliberately not the default, to catch any +/// hardcoded-prefix assumptions +pub(crate) const TEST_PREFIX: &[u8] = b"~test|"; + #[derive(Debug, Clone, Default)] pub struct MemEngine { inner: Arc>, @@ -32,6 +38,12 @@ impl MemEngine { pub fn new() -> Self { Self::default() } + + /// a fresh engine wrapped under [`TEST_PREFIX`], the way hubble-sync + /// internals always see storage + pub(crate) fn new_prefixed() -> PrefixedEngine { + PrefixedEngine::new(Self::new(), TEST_PREFIX) + } } impl StorageEngine for MemEngine { diff --git a/hubble-sync/src/storage/engine/mod.rs b/hubble-sync/src/storage/engine/mod.rs index 8c58776..ee79ef0 100644 --- a/hubble-sync/src/storage/engine/mod.rs +++ b/hubble-sync/src/storage/engine/mod.rs @@ -1,7 +1,9 @@ +mod prefixed; mod types; #[cfg(test)] pub(crate) mod mem; +pub use prefixed::{PrefixedBatch, PrefixedEngine}; pub(super) use types::Pair; pub use types::{StorageBatch, StorageEngine, StorageError}; diff --git a/hubble-sync/src/storage/engine/prefixed.rs b/hubble-sync/src/storage/engine/prefixed.rs new file mode 100644 index 0000000..8a26094 --- /dev/null +++ b/hubble-sync/src/storage/engine/prefixed.rs @@ -0,0 +1,224 @@ +//! wrap any StorageEngine to prefix every key +//! +//! hubble-sync uses this to keep its internal state under a strict prefx. apps +//! have to avoid writing under hubble's prefix, and they can use this to help +//! prevent that too, putting all their keys under their own disjoint prefix. + +use super::types::{Pair, StorageBatch, StorageEngine, StorageError}; +use std::time::Duration; + +#[derive(Debug, Clone)] +pub struct PrefixedEngine { + inner: E, + prefix: &'static [u8], +} + +/// wrap a storage engine, prefixing keys on every access +impl PrefixedEngine { + pub fn new(inner: E, prefix: &'static [u8]) -> Self { + Self { inner, prefix } + } + + /// access the wrapped engine + /// + /// same storage, unprefixed + pub fn inner(&self) -> &E { + &self.inner + } + + fn prefixed(&self, k: &[u8]) -> Vec { + [self.prefix, k].concat() + } +} + +impl StorageEngine for PrefixedEngine { + type Error = E::Error; + type Batch = PrefixedBatch; + + fn batch(&self) -> Self::Batch { + PrefixedBatch { + inner: self.inner.batch(), + prefix: self.prefix, + } + } + + fn get(&self, k: &[u8]) -> Result>, Self::Error> { + self.inner.get(&self.prefixed(k)) + } + fn get_queue(&self, k: &[u8]) -> Result>, Self::Error> { + self.inner.get_queue(&self.prefixed(k)) + } + fn get_counter(&self, k: &[u8]) -> Result { + self.inner.get_counter(&self.prefixed(k)) + } + + fn scan_from( + &self, + prefix: &[u8], + from_suffix: &[u8], + ) -> Box> + '_> { + self.inner.scan_from(&self.prefixed(prefix), from_suffix) + } + + fn scan_from_queue( + &self, + prefix: &[u8], + from_suffix: &[u8], + ) -> Box> + '_> { + self.inner + .scan_from_queue(&self.prefixed(prefix), from_suffix) + } + + fn maintenance(&self) -> Result, Self::Error> { + self.inner.maintenance() + } +} + +pub struct PrefixedBatch { + inner: B, + prefix: &'static [u8], +} + +impl PrefixedBatch { + /// same atomic batch, unprefixed + pub fn inner_mut(&mut self) -> &mut B { + &mut self.inner + } + + fn prefixed(&self, k: &[u8]) -> Vec { + [self.prefix, k].concat() + } +} + +impl> StorageBatch for PrefixedBatch { + fn commit(self) -> Result<(), E> { + self.inner.commit() + } + fn put(&mut self, k: &[u8], v: &[u8]) { + self.inner.put(&self.prefixed(k), v) + } + fn delete(&mut self, k: &[u8]) { + self.inner.delete(&self.prefixed(k)) + } + fn put_queue(&mut self, k: &[u8], v: &[u8]) { + self.inner.put_queue(&self.prefixed(k), v) + } + fn delete_queue(&mut self, k: &[u8]) { + self.inner.delete_queue(&self.prefixed(k)) + } + fn put_counter(&mut self, k: &[u8], count: i64) { + self.inner.put_counter(&self.prefixed(k), count) + } + fn increment_counter(&mut self, k: &[u8], delta: i64) { + self.inner.increment_counter(&self.prefixed(k), delta) + } + fn delete_counter(&mut self, k: &[u8]) { + self.inner.delete_counter(&self.prefixed(k)) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::storage::engine::mem::{MemBatch, MemEngine}; + + const P: &[u8] = b"~p|"; + + fn wrapped() -> PrefixedEngine { + PrefixedEngine::new(MemEngine::new(), P) + } + + fn commit_batch(eng: &PrefixedEngine, f: impl FnOnce(&mut PrefixedBatch)) { + let mut b = eng.batch(); + f(&mut b); + b.commit().unwrap(); + } + + #[test] + fn roundtrip_lands_under_prefix_in_every_cf() { + let eng = wrapped(); + commit_batch(&eng, |b| { + b.put(b"k", b"point"); + b.put_queue(b"k", b"queue"); + b.increment_counter(b"k", 7); + }); + // visible through the wrapper at the logical key + assert_eq!(eng.get(b"k").unwrap(), Some(b"point".to_vec())); + assert_eq!(eng.get_queue(b"k").unwrap(), Some(b"queue".to_vec())); + assert_eq!(eng.get_counter(b"k").unwrap(), 7); + // physically stored under the prefix, not at the raw key + assert_eq!(eng.inner().get(b"k").unwrap(), None); + assert_eq!(eng.inner().get(b"~p|k").unwrap(), Some(b"point".to_vec())); + assert_eq!(eng.inner().get_counter(b"k").unwrap(), 0); + assert_eq!(eng.inner().get_counter(b"~p|k").unwrap(), 7); + } + + #[test] + fn scan_strips_only_the_callers_prefix_and_stays_in_bounds() { + let eng = wrapped(); + commit_batch(&eng, |b| { + b.put(b"px|a", b"1"); + b.put(b"px|b", b"2"); + b.put(b"py|x", b"other"); // sibling logical prefix must not leak + }); + // a raw key that looks prefixed-ish must not leak into wrapped scans + let mut raw = eng.inner().batch(); + raw.put(b"px|raw", b"unprefixed"); + raw.commit().unwrap(); + + let got: Vec<_> = eng.scan_from(b"px|", b"").map(|r| r.unwrap()).collect(); + assert_eq!( + got, + vec![ + (b"a".to_vec(), b"1".to_vec()), + (b"b".to_vec(), b"2".to_vec()) + ] + ); + } + + #[test] + fn inner_mut_writes_raw_keys_in_the_same_atomic_batch() { + let eng = wrapped(); + commit_batch(&eng, |b| { + b.put(b"ours", b"prefixed"); + b.inner_mut().put(b"theirs", b"raw"); + }); + assert_eq!(eng.get(b"ours").unwrap(), Some(b"prefixed".to_vec())); + assert_eq!(eng.inner().get(b"theirs").unwrap(), Some(b"raw".to_vec())); + // and neither crossed into the other's keyspace + assert_eq!(eng.get(b"theirs").unwrap(), None); + assert_eq!(eng.inner().get(b"ours").unwrap(), None); + } + + #[test] + fn distinct_prefixes_over_one_engine_are_disjoint() { + let raw = MemEngine::new(); + let a = PrefixedEngine::new(raw.clone(), b"a|"); + let b = PrefixedEngine::new(raw, b"b|"); + + let mut batch = a.batch(); + batch.put(b"k", b"from-a"); + batch.commit().unwrap(); + let mut batch = b.batch(); + batch.put(b"k", b"from-b"); + batch.commit().unwrap(); + + assert_eq!(a.get(b"k").unwrap(), Some(b"from-a".to_vec())); + assert_eq!(b.get(b"k").unwrap(), Some(b"from-b".to_vec())); + } + + #[test] + fn nesting_composes() { + let outer = PrefixedEngine::new(wrapped(), b"in|"); + let mut batch = outer.batch(); + batch.put(b"k", b"v"); + batch.commit().unwrap(); + + assert_eq!(outer.get(b"k").unwrap(), Some(b"v".to_vec())); + // outer's prefix applies first, then the wrapped engine's own + assert_eq!( + outer.inner().inner().get(b"~p|in|k").unwrap(), + Some(b"v".to_vec()) + ); + } +} diff --git a/hubble-sync/src/sync_handle.rs b/hubble-sync/src/sync_handle.rs index e2c7e45..5daafef 100644 --- a/hubble-sync/src/sync_handle.rs +++ b/hubble-sync/src/sync_handle.rs @@ -13,22 +13,23 @@ use tracing::{debug, error, info}; use crate::storage::moderation_log::ModerateEvent; use crate::storage::repo::{AccountStatus, AccountStatusEvent, ModAction, SyncStatus}; use crate::{ - Did, HostRegistry, LoadError, ModerateOutcome, Repo, RepoContext, RepoCountsByState, - RepoCountsByStatus, RepoMessage, RepoSendError, RepoSender, RepoSlot, StorageEngine, Tid, + Did, HostRegistry, LoadError, ModerateOutcome, PrefixedEngine, Repo, RepoContext, + RepoCountsByState, RepoCountsByStatus, RepoMessage, RepoSendError, RepoSender, RepoSlot, + StorageEngine, Tid, }; /// read-only view into hubble-sync's state about repos #[derive(Clone)] pub struct SyncHandle { - storage: S, + storage: PrefixedEngine, hosts: Arc, - endpoints: SyncEndpoints, + endpoints: SyncEndpoints>, slots: PhantomData<(C, I)>, } impl SyncHandle { pub(crate) fn new( - storage: S, + storage: PrefixedEngine, hosts: Arc, repo_sender: Arc, ) -> Self { @@ -50,12 +51,12 @@ impl SyncHandle { Ok(Some(RepoView { repo, commit, info })) } - pub fn endpoints(&self) -> SyncEndpoints { + pub fn endpoints(&self) -> SyncEndpoints> { self.endpoints.clone() } pub fn storage(&self) -> &S { - &self.storage + self.storage.inner() } } @@ -457,8 +458,11 @@ mod tests { /// seed a repo's `ri|` info (always) and, if `rev` is `Some`, its `si|` sync /// state — i.e. a complete repo that `list_synced` will surface. - fn seed_repo( - eng: &MemEngine, + /// + /// generic so it can seed through either a raw engine (SyncEndpoints + /// tests) or the PrefixedEngine view a SyncHandle reads from. + fn seed_repo( + eng: &S, hosts: &HostRegistry, did: &Did, account_status: AccountStatus, @@ -696,13 +700,14 @@ mod tests { #[test] fn get_returns_none_for_missing_repo() { - let handle = SyncHandle::<_, (), ()>::new(MemEngine::new(), hosts(), FakeSender::new()); + let handle = + SyncHandle::<_, (), ()>::new(MemEngine::new_prefixed(), hosts(), FakeSender::new()); assert!(handle.get(&plc_did("ghost")).unwrap().is_none()); } #[test] fn get_returns_repo_view_with_context() { - let (eng, hs, fs) = (MemEngine::new(), hosts(), FakeSender::new()); + let (eng, hs, fs) = (MemEngine::new_prefixed(), hosts(), FakeSender::new()); let did = plc_did("alice"); seed_repo( &eng, @@ -721,7 +726,7 @@ mod tests { #[test] fn get_with_unit_slots_has_no_slot_data() { // unit slots are never present, so the view carries no app-owned data. - let (eng, hs, fs) = (MemEngine::new(), hosts(), FakeSender::new()); + let (eng, hs, fs) = (MemEngine::new_prefixed(), hosts(), FakeSender::new()); let did = plc_did("noslots"); seed_repo( &eng, diff --git a/hubble/readme.md b/hubble/readme.md index 64cff0d..ee89e33 100644 --- a/hubble/readme.md +++ b/hubble/readme.md @@ -1,3 +1,10 @@ # hubble just getting started. see [./hacking.md](./hacking.md) for now + + +--- + +- [ ] resync "size" is from repo-stream, which counts block bytes, not download size, which is a bit less publicly-useful. we could either count over the raw byte stream or extend repo-stream to expose input byte counts too. + +- [ ] index blob links diff --git a/hubble/src/sync.rs b/hubble/src/sync.rs index 66ce277..10e7f1a 100644 --- a/hubble/src/sync.rs +++ b/hubble/src/sync.rs @@ -29,7 +29,7 @@ impl StatefulSyncConsumer for Hubble { let generation = slots.info.get().map_or(0, |i| i.resync_generation); let records = self.db.cf_handle("records").expect("records cf"); - let wb = batch.inner_mut(); // TODO: provide unprefixed-enforcing apis from hubble-sync storage engine + let wb = batch.inner_mut(); for op in commit.added().chain(commit.updated()) { let key = record::key(ctx.did(), generation, &op.path); @@ -46,7 +46,7 @@ impl StatefulSyncConsumer for Hubble { } let diff = commit.added().count() as i64 - removed; - // TODO oops we're writing under hubble's prefix with this + batch.increment_counter(&counts::global_key(), diff); batch.increment_counter(&counts::key(ctx.did()), diff); -- 2.51.2