diff --git a/hubble-sync/src/identity/crypto.rs b/hubble-sync/src/identity/crypto.rs index 85291b0..cb9b59e 100644 --- a/hubble-sync/src/identity/crypto.rs +++ b/hubble-sync/src/identity/crypto.rs @@ -19,6 +19,17 @@ pub enum SignatureError { VerificationFailed, } +impl SignatureError { + pub fn could_retry_with_another_key(&self) -> bool { + match self { + Self::BadKey => true, + Self::BadSignature => true, // signature construction depends on key type + Self::HighS => false, // problem with the actual signature + Self::VerificationFailed => true, + } + } +} + /// verify the signature over a message with a secp256k1 key /// /// msg is the encoded bytes -- k256 does the sha256 hashing for us diff --git a/hubble-sync/src/repo_actor/mod.rs b/hubble-sync/src/repo_actor/mod.rs index 0b9e056..c6a806e 100644 --- a/hubble-sync/src/repo_actor/mod.rs +++ b/hubble-sync/src/repo_actor/mod.rs @@ -9,7 +9,7 @@ mod task_processor; pub use actor::RepoMessage; pub use repo_registry::{RepoRegistry, RepoSendError}; -pub use task_processor::{RepoContext, ResyncContext, ValidateCommitError}; +pub use task_processor::{RepoContext, ResyncContext}; use actor::Task; use evictable_state::EvictableState; diff --git a/hubble-sync/src/repo_actor/task_processor/event_validation.rs b/hubble-sync/src/repo_actor/task_processor/event_validation.rs index 495256d..9daa9c9 100644 --- a/hubble-sync/src/repo_actor/task_processor/event_validation.rs +++ b/hubble-sync/src/repo_actor/task_processor/event_validation.rs @@ -5,22 +5,6 @@ use serde::Serialize; use crate::commit_object::CommitObject; use crate::identity::{SignatureError, SigningKey}; -#[derive(Debug, thiserror::Error)] -pub enum ValidateCommitError { - #[error("signature error: {0}")] - BadSignature(#[from] SignatureError), -} - -impl ValidateCommitError { - pub fn could_retry_with_another_key(&self) -> bool { - matches!( - self, - Self::BadSignature(SignatureError::BadKey) - | Self::BadSignature(SignatureError::VerificationFailed) - ) - } -} - /// the commit representation that gets signed /// /// `sig` is absent, gets dag-cbor (drisl)-encoded @@ -55,28 +39,9 @@ impl<'a> From<&'a CommitObject> for UnsignedCommit<'a> { pub fn validate_commit_signature( key: &SigningKey, commit: &CommitObject, -) -> Result<(), ValidateCommitError> { +) -> Result<(), SignatureError> { // 4. verify the commit signature let unsigned: UnsignedCommit = commit.into(); let signed_bytes = drisl::to_vec(&unsigned).expect("infallible encode"); - key.verify_bytes(&signed_bytes, &commit.sig)?; - // TODO: match error here ^^ to separate refresh-and-retry failures - - // // let key = repo.info().identity.signing_key.clone(); - // match asdf { - // Ok(_v) => { - // tracing::info!(%did, %new_rev, "#commit verified commit"); - // unimplemented!("the rest"); - // } - // Err(err) if err.could_retry_with_another_key() => { - // tracing::info!(%did, %new_rev, %err, "#commit verification failed but will refresh & retry"); - // unimplemented!("refresh identity and retry") - // } - // Err(err) => { - // debug!(%did, %new_rev, %err, "#commit dropped (signature check fail)"); - // return Ok(()); - // } - // } - - Ok(()) + key.verify_bytes(&signed_bytes, &commit.sig) } diff --git a/hubble-sync/src/repo_actor/task_processor/mod.rs b/hubble-sync/src/repo_actor/task_processor/mod.rs index d5e2d05..5ed9334 100644 --- a/hubble-sync/src/repo_actor/task_processor/mod.rs +++ b/hubble-sync/src/repo_actor/task_processor/mod.rs @@ -1,7 +1,5 @@ mod event_validation; -pub use event_validation::ValidateCommitError; - use std::path::Path; use std::sync::Arc; use std::time::{Duration, Instant, SystemTime}; @@ -10,15 +8,17 @@ use crate::firehose::{FirehoseCommit, FirehoseSync}; use tokio::sync::{Semaphore, oneshot}; use tokio::task::spawn_blocking; use tokio_util::sync::CancellationToken; -use tracing::{error, trace, warn}; +use tracing::{debug, error, trace, warn}; use super::Task; use crate::firehose::{FirehoseEvent, FirehosePayload}; -use crate::identity::{Resolve, ResolvedIdentity}; +use crate::identity::{Resolve, ResolvedIdentity, Validity}; use crate::resync_data::{ResyncData, ResyncError, Resyncable, TransientResyncError, load_repo}; use crate::storage::engine::{StorageBatch, StorageEngine}; use crate::storage::repo::{AccountStatus, DesyncReason, Repo, RepoIdentity}; -use crate::{CancelExt, ConsumerAppError, Did, Host, StorageError, SyncConsumer, Tid}; +use crate::{ + CancelExt, CommitObject, ConsumerAppError, Did, Host, StorageError, SyncConsumer, Tid, +}; use event_validation::validate_commit_signature; const STREAM_CAR_TIMEOUT: Duration = Duration::from_secs(300); @@ -141,8 +141,8 @@ impl, R: Resolve> TaskProcessor PEResult { - let repo = self.repo.as_ref().expect("repo for task"); - let did = repo.did(); + let mut repo = self.repo.as_ref().expect("repo for task"); + let did = repo.did().clone(); trace!(%did, "started processing #commit event"); // See [`crate::firehose::event_validation::FirehoseCommit::prevalidate`] @@ -153,12 +153,18 @@ impl, R: Resolve> TaskProcessor, R: Resolve> TaskProcessor, R: Resolve> TaskProcessor, R: Resolve> TaskProcessor PEResult { - let repo = self.repo.as_ref().expect("repo for task"); + let mut repo = self.repo.as_ref().expect("repo for task"); let did = repo.did().clone(); trace!(%did, late_by = ?was_due_at.elapsed(), "started processing scheduled resync"); let Some(desync) = repo.current_desync() else { warn!( - ?did, + %did, ?was_due_at, "spurious resync task (repo not desynchronized), ignoring" ); return Ok(()); }; if !desync.is_at(was_due_at) { - warn!(?did, expected_due = ?was_due_at, found_due = ?desync.due_at(), + warn!(%did, expected_due = ?was_due_at, found_due = ?desync.due_at(), "spurious resync task (wrong `due_at`), ignoring"); return Ok(()); } @@ -359,6 +365,23 @@ impl, R: Resolve> TaskProcessor, R: Resolve> TaskProcessor Result { let repo = self.repo.as_ref().expect("repo for task"); let did = repo.did(); @@ -425,25 +454,16 @@ impl, R: Resolve> TaskProcessor, R: Resolve> TaskProcessor PEResult>)> { + let mut rr = None; // Some(repo) if modified + + let ident = &repo.info().identity; + + match validate_commit_signature(&ident.signing_key, commit) { + Ok(()) => return Ok((true, rr)), + Err(e) if !e.could_retry_with_another_key() => return Ok((false, rr)), + _ => {} // worth a retry + }; + + // but only if we haven't refreshed too recently + let did = repo.did().clone(); + if !Validity::can_refresh(&did, ident.age(now)) { + trace!(%did, "not refreshing identity to retry signature verification (too soon)"); + return Ok((false, rr)); + } + + // all clear to (try to) re-resolve + let new_identity = match self.resolver.resolve(&did, now).await { + Ok(id) => id, + Err(err) => { + debug!(%did, %err, "failed identity refresh during signature verification retry"); + return Ok((false, rr)); + } + }; + let new_key = new_identity.signing_key.clone(); + let unchanged = new_key == ident.signing_key; + + // sweet, save it! + // (must still save to update the age, even if unchanged) + let updated = self.commit_refreshed(repo.clone(), new_identity).await?; + rr = Some(Box::new(updated)); + + // if the key hasn't changed, there's nothing more for us to try + if unchanged { + return Ok((false, rr)); + } + + // otherwise, one last shot: + let validated = validate_commit_signature(&new_key, commit).is_ok(); + Ok((validated, rr)) + } + + /// persist an identity change + /// + /// *returns* a [`Repo`] with the changes (does not mutate self) async fn commit_refreshed( - &mut self, - mut repo: Repo, + &self, // for storage only (pass in?) + mut updating_repo: Repo, identity: RepoIdentity, ) -> PEResult { let storage = self.storage.clone(); let updated_repo = spawn_blocking(move || -> PEResult { let mut batch = storage.batch(); - repo.set_identity(identity, &mut batch); + updating_repo.set_identity(identity, &mut batch); batch.commit().map_err(ProcessError::Storage)?; - Ok(repo) + Ok(updating_repo) }) .await .expect("storage not to panic")?; diff --git a/hubble-sync/src/resync_data.rs b/hubble-sync/src/resync_data.rs index 9ff1c90..d3534b3 100644 --- a/hubble-sync/src/resync_data.rs +++ b/hubble-sync/src/resync_data.rs @@ -8,7 +8,6 @@ use tokio::sync::{OwnedSemaphorePermit, Semaphore}; use crate::Did; use crate::commit_object::{CommitConvertError, CommitObject}; -use crate::repo_actor::ValidateCommitError; const MEM_LIMIT_SMALL_MB: usize = 2; const MEM_LIMIT_LARGE_MB: usize = 150; @@ -41,8 +40,6 @@ pub enum TransientResyncError { Spill(String), #[error("wrong DID, expected {expected}, found {got:?}")] WrongDid { expected: Did, got: Did }, - #[error("failed signature verification: {0}")] - ValidateCommit(ValidateCommitError), } #[derive(Debug, thiserror::Error)] diff --git a/hubble-sync/src/storage/repo/info.rs b/hubble-sync/src/storage/repo/info.rs index 0df1bc8..4506f11 100644 --- a/hubble-sync/src/storage/repo/info.rs +++ b/hubble-sync/src/storage/repo/info.rs @@ -162,6 +162,13 @@ impl fmt::Debug for RepoIdentity { } } +impl RepoIdentity { + pub fn age(&self, now: SystemTime) -> Duration { + now.duration_since(self.resolved_at) + .unwrap_or(Duration::ZERO) + } +} + #[derive(Debug, Clone, Serialize, Deserialize)] pub enum AccountStatus { Active, @@ -235,17 +242,22 @@ impl Desynchronized { pub fn backoff(reason: &DesyncReason, attempts: u32, method: DidMethod) -> Option { match reason { + #[allow( + clippy::match_overlapping_arm, + reason = "visual alignment is actually clearer here i think" + )] DesyncReason::UnresolvableIdentity => match method { DidMethod::Plc => match attempts { ..=1 => Some(60), - 2 => Some(300), + ..=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 (!) + ..=4 => Some(300), + ..=8 => Some(3600), + ..=48 => Some(MAX_BACKOFF.as_secs()), + _ => None, // eventuallyyyyyyyyyyyyy }, }, DesyncReason::FirstSeen => Some(Self::backoff_from_zero(1.5, attempts)), @@ -257,7 +269,7 @@ impl Desynchronized { Err(_) => Self::backoff_from(30, 1.5, attempts), }) } - DesyncReason::FirehoseSync => Some(Self::backoff_from_zero(2., attempts)), + DesyncReason::FirehoseSync { .. } => Some(Self::backoff_from_zero(2., attempts)), DesyncReason::Sync11Lax => { // up to twice a day lax catch-ups Some(12 * 3600) @@ -286,7 +298,10 @@ pub enum DesyncReason { /// 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, + FirehoseSync { + #[serde(with = "crate::tid::tid_raw_u64")] + rev: Tid, + }, /// sync1.1 proof validation failure on a #commit /// /// commits dropped otherwise (signature fail, old `rev`,) don't resync.