diff --git a/src/api/xrpc/get_latest_commit.rs b/src/api/xrpc/get_latest_commit.rs index 7e6f10e..91e8e09 100644 --- a/src/api/xrpc/get_latest_commit.rs +++ b/src/api/xrpc/get_latest_commit.rs @@ -41,7 +41,7 @@ pub async fn handle( } // return whatever we last recorded; if we haven't synced at all yet, we have nothing to give - let (Some(data), Some(rev_db)) = (state.data, state.rev) else { + let Some(commit) = state.root else { return Err(XrpcErrorResponse { status: StatusCode::NOT_FOUND, error: XrpcError::Xrpc(GetLatestCommitError::RepoNotFound(Some(CowStr::Borrowed( @@ -51,8 +51,8 @@ pub async fn handle( }; Ok(Json(GetLatestCommitOutput { - cid: Cid::from(data), - rev: rev_db.to_tid(), + cid: Cid::from(commit.data), + rev: commit.rev.to_tid(), extra_data: None, })) } diff --git a/src/api/xrpc/get_repo_status.rs b/src/api/xrpc/get_repo_status.rs index 889ea50..3f28534 100644 --- a/src/api/xrpc/get_repo_status.rs +++ b/src/api/xrpc/get_repo_status.rs @@ -32,7 +32,7 @@ pub async fn handle( let (active, status) = repo_status_to_api(state.status); // rev is only meaningful when the repo is active and has been synced at least once - let rev = active.then(|| state.rev.map(|r| r.to_tid())).flatten(); + let rev = active.then(|| state.root.map(|c| c.rev.to_tid())).flatten(); Ok(Json(GetRepoStatusOutput { active, diff --git a/src/api/xrpc/list_repos.rs b/src/api/xrpc/list_repos.rs index c7ec0b2..efeb759 100644 --- a/src/api/xrpc/list_repos.rs +++ b/src/api/xrpc/list_repos.rs @@ -34,7 +34,7 @@ pub async fn handle( let (did, state) = item?; // skip repos that haven't been synced at least once - let (Some(data), Some(rev_db)) = (state.data, state.rev) else { + let Some(commit) = state.root else { continue; }; @@ -42,8 +42,8 @@ pub async fn handle( repos.push(Repo { active: Some(active), did: did.clone(), - head: Cid::from(data), - rev: rev_db.to_tid(), + head: Cid::from(commit.data), + rev: commit.rev.to_tid(), status, extra_data: None, }); diff --git a/src/backfill/mod.rs b/src/backfill/mod.rs index b793c9a..ddd75f5 100644 --- a/src/backfill/mod.rs +++ b/src/backfill/mod.rs @@ -1,12 +1,12 @@ -use crate::db::types::{DbAction, DbRkey, DbTid, TrimmedDid}; +use crate::db::types::{DbAction, DbRkey, TrimmedDid}; use crate::db::{self, Db, keys, ser_repo_state}; use crate::filter::FilterMode; use crate::ops; use crate::resolver::ResolverError; use crate::state::AppState; use crate::types::{ - AccountEvt, BroadcastEvent, GaugeState, RepoState, RepoStatus, ResyncErrorKind, ResyncState, - StoredData, StoredEvent, + AccountEvt, BroadcastEvent, Commit, GaugeState, RepoState, RepoStatus, ResyncErrorKind, + ResyncState, StoredData, StoredEvent, }; use fjall::Slice; @@ -545,6 +545,8 @@ async fn process_did<'i>( trace!("signature verified"); } + let root_commit = Commit::from(root_commit); + // 5. walk mst let start = Instant::now(); let mst: Mst = Mst::load(store, root_commit.data, None); @@ -664,7 +666,7 @@ async fn process_did<'i>( let evt = StoredEvent { live: false, did: TrimmedDid::from(&did), - rev: DbTid::from(&rev), + rev, collection: CowStr::Borrowed(collection), rkey, action, @@ -700,7 +702,7 @@ async fn process_did<'i>( let evt = StoredEvent { live: false, did: TrimmedDid::from(&did), - rev: DbTid::from(&rev), + rev, collection: CowStr::Borrowed(&collection), rkey, action: DbAction::Delete, @@ -720,8 +722,7 @@ async fn process_did<'i>( // 6. update data, status is updated in worker shard state.tracked = true; - state.rev = Some((&rev).into()); - state.data = Some(root_commit.data); + state.root = Some(root_commit); state.touch(); batch.insert( diff --git a/src/control/repos.rs b/src/control/repos.rs index 6844762..8235fdd 100644 --- a/src/control/repos.rs +++ b/src/control/repos.rs @@ -389,12 +389,16 @@ impl ReposControl { } pub(crate) fn repo_state_to_info(did: Did<'static>, s: RepoState<'_>) -> RepoInfo { + let (rev, data) = s + .root + .map(|c| (Some(c.rev.to_tid()), Some(c.data))) + .unwrap_or_default(); RepoInfo { did, status: s.status, tracked: s.tracked, - rev: s.rev.map(|r| r.to_tid()), - data: s.data, + rev, + data, handle: s.handle.map(|h| h.into_static()), pds: s.pds.and_then(|p| p.parse().ok()), signing_key: s.signing_key.map(|k| k.into_static()), diff --git a/src/db/migration/mod.rs b/src/db/migration/mod.rs index 8e0deed..34cc62b 100644 --- a/src/db/migration/mod.rs +++ b/src/db/migration/mod.rs @@ -5,12 +5,15 @@ use crate::db::Db; use crate::db::keys::VERSIONING_KEY; mod v1; +mod v2; type MigrationFn = fn(&Db, &mut OwnedWriteBatch) -> Result<()>; /// ordered list of migrations. migration at index `i` upgrades the schema from version `i` to `i+1`. -const MIGRATIONS: &[(&str, MigrationFn)] = - &[("stable_firehose_cursors", v1::stable_firehose_cursors)]; +const MIGRATIONS: &[(&str, MigrationFn)] = &[ + ("stable_firehose_cursors", v1::stable_firehose_cursors), + ("repo_state_root_commit", v2::repo_state_root_commit), +]; fn read_version(db: &Db) -> Result { db.counts diff --git a/src/db/migration/v2.rs b/src/db/migration/v2.rs new file mode 100644 index 0000000..d9ee23e --- /dev/null +++ b/src/db/migration/v2.rs @@ -0,0 +1,68 @@ +use bytes::Bytes; +use cid::Cid as IpldCid; +use fjall::OwnedWriteBatch; +use jacquard_common::{CowStr, types::string::Handle}; +use miette::{Context, IntoDiagnostic, Result}; +use serde::{Deserialize, Serialize}; + +use crate::db::{ + Db, + types::{DbTid, DidKey}, +}; +use crate::types::v2::*; + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(bound(deserialize = "'i: 'de"))] +pub(crate) struct OldRepoState<'i> { + pub status: RepoStatus, // from v2, old is same as new + pub rev: Option, + pub data: Option, + pub last_message_time: Option, + pub last_updated_at: i64, + pub tracked: bool, + pub index_id: u64, + #[serde(borrow)] + pub signing_key: Option>, + #[serde(borrow)] + pub pds: Option>, + #[serde(borrow)] + pub handle: Option>, +} + +pub(super) fn repo_state_root_commit(db: &Db, batch: &mut OwnedWriteBatch) -> Result<()> { + for item in db.repos.iter() { + let (k, v) = item.into_inner().into_diagnostic()?; + let old: OldRepoState = rmp_serde::from_slice(&v) + .into_diagnostic() + .wrap_err("invalid old repo state")?; + let new = RepoState { + root: match (old.rev, old.data) { + (Some(rev), Some(data)) => Some(Commit { + version: -1, + rev, + data, + prev: None, + sig: Bytes::new(), + }), + _ => None, + }, + status: old.status, + handle: old.handle, + index_id: old.index_id, + last_message_time: old.last_message_time, + last_updated_at: old.last_updated_at, + pds: old.pds, + signing_key: old.signing_key, + tracked: old.tracked, + }; + batch.insert( + &db.repos, + k, + rmp_serde::to_vec(&new) + .into_diagnostic() + .wrap_err("cant serialize new repo state")?, + ); + } + + Ok(()) +} diff --git a/src/ingest/worker.rs b/src/ingest/worker.rs index 1376f51..9eaff30 100644 --- a/src/ingest/worker.rs +++ b/src/ingest/worker.rs @@ -435,19 +435,20 @@ impl FirehoseWorker { repo_state.advance_message_time(commit.time.0.timestamp_millis()); // skip replayed events (already seen revision) - if matches!(repo_state.rev, Some(ref rev) if commit.rev.as_str() <= rev.to_tid().as_str()) { + if matches!(repo_state.root, Some(ref root) if commit.rev.as_str() <= root.rev.to_tid().as_str()) + { debug!( did = %did, commit_rev = %commit.rev, - state_rev = %repo_state.rev.as_ref().map(|r| r.to_tid()).expect("we checked in if"), + state_rev = %repo_state.root.as_ref().map(|c| c.rev.to_tid()).expect("we checked in if"), "skipping replayed event" ); return Ok(RepoProcessResult::Ok(repo_state)); } - if let (Some(repo), Some(prev_commit)) = (&repo_state.data, &commit.prev_data) - && repo - != &prev_commit + if let (Some(repo_commit), Some(prev_commit)) = (&repo_state.root, &commit.prev_data) + && repo_commit.data + != prev_commit .0 .to_ipld() .into_diagnostic() @@ -455,7 +456,7 @@ impl FirehoseWorker { { warn!( did = %did, - repo = %repo, + repo = %repo_commit.data, prev_commit = %prev_commit.0, "gap detected, triggering backfill" ); @@ -508,15 +509,13 @@ impl FirehoseWorker { match ops::verify_sync_event(sync.blocks.as_ref(), Self::fetch_key(ctx, did)?.as_ref()) { Ok((root, rev)) => { - if let Some(current_data) = &repo_state.data { - if current_data == &root.to_ipld().expect("valid cid") { + if let Some(current_commit) = &repo_state.root { + if current_commit.data == root.to_ipld().expect("valid cid") { debug!(did = %did, "skipping noop sync"); return Ok(RepoProcessResult::Ok(repo_state)); } - } - if let Some(current_rev) = &repo_state.rev { - if rev.as_str() <= current_rev.to_tid().as_str() { + if rev.as_str() <= current_commit.rev.to_tid().as_str() { debug!(did = %did, "skipping replayed sync"); return Ok(RepoProcessResult::Ok(repo_state)); } diff --git a/src/ops.rs b/src/ops.rs index e7b48d2..bf31437 100644 --- a/src/ops.rs +++ b/src/ops.rs @@ -256,17 +256,16 @@ pub fn apply_commit<'commit, 's>( .get(&parsed.root) .ok_or_else(|| miette::miette!("root block missing from CAR"))?; - let repo_commit = jacquard_repo::commit::Commit::from_cbor(root_bytes).into_diagnostic()?; + let root_commit = jacquard_repo::commit::Commit::from_cbor(root_bytes).into_diagnostic()?; if let Some(key) = signing_key { - repo_commit + root_commit .verify(key) .map_err(|e| miette::miette!("signature verification failed for {did}: {e}"))?; trace!(did = %did, "signature verified"); } - repo_state.rev = Some((&commit.rev).into()); - repo_state.data = Some(repo_commit.data); + repo_state.root = Some(root_commit.into()); repo_state.touch(); batch.insert(&db.repos, keys::repo_key(did), ser_repo_state(&repo_state)?); diff --git a/src/types.rs b/src/types.rs index 94aaf3d..929e7c3 100644 --- a/src/types.rs +++ b/src/types.rs @@ -12,15 +12,56 @@ use smol_str::{SmolStr, ToSmolStr}; use crate::db::types::{DbAction, DbRkey, DbTid, DidKey, TrimmedDid}; use crate::resolver::MiniDoc; -#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] -pub enum RepoStatus { - Backfilling, - Synced, - Error(SmolStr), - Deactivated, - Takendown, - Suspended, +pub(crate) mod v2 { + use super::*; + + #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] + pub enum RepoStatus { + Backfilling, + Synced, + Error(SmolStr), + Deactivated, + Takendown, + Suspended, + } + + #[derive(Debug, Clone, Serialize, Deserialize)] + pub(crate) struct Commit { + pub version: i64, + pub rev: DbTid, + pub data: IpldCid, + pub prev: Option, + #[serde(with = "jacquard_common::serde_bytes_helper")] + pub sig: Bytes, + } + + #[derive(Debug, Clone, Serialize, Deserialize)] + #[serde(bound(deserialize = "'i: 'de"))] + pub(crate) struct RepoState<'i> { + pub status: RepoStatus, + pub root: Option, + // todo: is this actually valid? the spec says this is informal and intermadiate + // services may change it. we should probably document it. if we cant use this + // then how do we dedup account / identity ops? + /// ms since epoch of the last firehose message we processed for this repo. + /// used to deduplicate identity / account events that can arrive from multiple relays at + /// different wall-clock times but represent the same underlying PDS event. + pub last_message_time: Option, + /// this is when we *ingested* any last updates + pub last_updated_at: i64, // unix timestamp + /// whether we are ingesting events for this repo + pub tracked: bool, + /// index id in pending keyspace + pub index_id: u64, + #[serde(borrow)] + pub signing_key: Option>, + #[serde(borrow)] + pub pds: Option>, + #[serde(borrow)] + pub handle: Option>, + } } +pub(crate) use v2::*; impl Display for RepoStatus { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { @@ -35,40 +76,23 @@ impl Display for RepoStatus { } } -#[derive(Debug, Clone, Serialize, Deserialize)] -#[serde(bound(deserialize = "'i: 'de"))] -pub(crate) struct RepoState<'i> { - pub status: RepoStatus, - pub rev: Option, - pub data: Option, - // todo: is this actually valid? the spec says this is informal and intermadiate - // services may change it. we should probably document it. if we cant use this - // then how do we dedup account / identity ops? - /// ms since epoch of the last firehose message we processed for this repo. - /// used to deduplicate identity / account events that can arrive from multiple relays at - /// different wall-clock times but represent the same underlying PDS event. - #[serde(default)] - pub last_message_time: Option, - /// this is when we *ingested* any last updates - pub last_updated_at: i64, // unix timestamp - /// whether we are ingesting events for this repo - pub tracked: bool, - /// index id in pending keyspace - pub index_id: u64, - #[serde(borrow)] - pub signing_key: Option>, - #[serde(borrow)] - pub pds: Option>, - #[serde(borrow)] - pub handle: Option>, +impl<'c> From> for Commit { + fn from(value: jacquard_repo::commit::Commit<'c>) -> Self { + Self { + data: value.data, + prev: value.prev, + rev: DbTid::from(&value.rev), + sig: value.sig, + version: value.version, + } + } } impl<'i> RepoState<'i> { pub fn backfilling(index_id: u64) -> Self { Self { status: RepoStatus::Backfilling, - rev: None, - data: None, + root: None, last_updated_at: chrono::Utc::now().timestamp(), index_id, tracked: true, @@ -115,8 +139,7 @@ impl<'i> IntoStatic for RepoState<'i> { fn into_static(self) -> Self::Output { RepoState { status: self.status, - rev: self.rev, - data: self.data, + root: self.root, last_updated_at: self.last_updated_at, index_id: self.index_id, tracked: self.tracked, @@ -247,7 +270,7 @@ use bytes::Bytes; pub(crate) enum StoredData { Nothing, Ptr(IpldCid), - #[serde(with = "serde_bytes_squared")] + #[serde(with = "jacquard_common::serde_bytes_helper")] Block(Bytes), } @@ -290,19 +313,6 @@ pub(crate) struct StoredEvent<'i> { pub data: StoredData, } -mod serde_bytes_squared { - use bytes::Bytes; - use serde::{Deserialize, Deserializer, Serializer}; - - pub fn serialize(v: impl AsRef<[u8]>, s: S) -> Result { - s.serialize_bytes(serde_bytes::Bytes::new(v.as_ref())) - } - - pub fn deserialize<'de, D: Deserializer<'de>>(d: D) -> Result { - serde_bytes::ByteBuf::deserialize(d).map(|b| b.into_vec().into()) - } -} - #[derive(Debug, PartialEq, Eq, Clone, Copy)] pub(crate) enum GaugeState { Synced,