diff --git a/src/api/messages.rs b/src/api/messages.rs index 40cb039..e9f75d4 100644 --- a/src/api/messages.rs +++ b/src/api/messages.rs @@ -487,11 +487,11 @@ mod tests { } #[tokio::test] - async fn send_message_uses_the_state_sms_sender() { + async fn send_message_uses_the_runtime_modem_path() { let sender = RecordingSmsSender::default(); let store = crate::storage::MessageStore::open_in_memory().unwrap(); let modem = crate::modem::ModemService::new(); - modem.set_verified_path(Some("/org/freedesktop/ModemManager1/Modem/0".to_string())); + modem.set_runtime_path(Some("/org/freedesktop/ModemManager1/Modem/1".to_string())); let state = super::super::ApiState { config: std::sync::Arc::new(crate::config::AppConfig::default()), config_path: std::path::PathBuf::from("/tmp/not-used.toml"), @@ -559,7 +559,7 @@ mod tests { assert_eq!( sender.calls.lock().unwrap().as_slice(), [( - "/org/freedesktop/ModemManager1/Modem/0".to_string(), + "/org/freedesktop/ModemManager1/Modem/1".to_string(), "+15551234567".to_string(), "test body".to_string(), )] diff --git a/src/api/mod.rs b/src/api/mod.rs index 2c38ba4..ff0660d 100644 --- a/src/api/mod.rs +++ b/src/api/mod.rs @@ -47,7 +47,7 @@ impl ApiState { self.delivery_wakeup.clone(), self.sms_sender.clone(), ) - .with_verified_modem(self.modem.clone()) + .with_modem_service(self.modem.clone()) } } diff --git a/src/inbound.rs b/src/inbound.rs index abb70df..196498b 100644 --- a/src/inbound.rs +++ b/src/inbound.rs @@ -12,7 +12,7 @@ use crate::dbus::{ InboundEvent, InboundSms, InboundSmsProperties, InboundSubscription, SystemInboundSource, }; use crate::messaging::{Messaging, ReceiveMessage}; -use crate::modem::ModemService; +use crate::modem::{ModemService, ModemTargets}; use crate::persistence::Store; const MAX_INBOUND_TASKS: usize = 16; @@ -22,6 +22,11 @@ const BODY_POLL_INTERVAL: Duration = Duration::from_millis(100); const MAX_BODY_POLLS: usize = 600; const INITIAL_PERSISTENCE_RETRY_DELAY: Duration = Duration::from_millis(100); const MAX_PERSISTENCE_RETRY_DELAY: Duration = Duration::from_secs(30); +const LEGACY_SINGLE_MODEM_FINGERPRINT_SEED: &str = "sms-relayed-single-modem"; +#[cfg(not(test))] +const RUNTIME_IDENTITY_REFRESH_INTERVAL: Duration = Duration::from_secs(60); +#[cfg(test)] +const RUNTIME_IDENTITY_REFRESH_INTERVAL: Duration = Duration::from_millis(10); #[derive(Debug, Clone)] pub struct ReceivedSms { @@ -149,19 +154,17 @@ impl InboundWorker { ) .await { - Ok(path) => { - self.modem_service.set_verified_path(path.clone()); - path - } + Ok(resolved) => publish_resolved_path(&self.modem_service, resolved), Err(error) => { error!("modem resolution failed: {}", error); - self.modem_service.set_verified_path(None); + self.modem_service + .set_modem_targets(ModemTargets::default()); None } }; } let Some(path) = current_path.clone() else { - warn!("no verified modem identity available; retrying resolution"); + warn!("no runtime modem path available; retrying resolution"); tokio::time::sleep(delay).await; delay = next_reconnect_delay(delay); continue; @@ -174,29 +177,28 @@ impl InboundWorker { Err(error) => { error!("D-Bus monitor lost: {}", error); crate::monitoring::capture_failure("dbus", "dbus.monitor_lost"); - let resolved = match resolve_monitor_path( + let resolved_path = match resolve_monitor_path( &self.settings.configured_modem_path, &self.modem_service, &self.store, ) .await { - Ok(resolved) => resolved, + Ok(resolved) => publish_resolved_path(&self.modem_service, resolved), Err(error) => { error!("modem resolution failed: {}", error); - self.modem_service.set_verified_path(None); + self.modem_service + .set_modem_targets(ModemTargets::default()); None } }; - if let Some(new_path) = resolved { + if let Some(new_path) = resolved_path { if new_path != path { info!("modem path changed from {} to {}", path, new_path); } - self.modem_service.set_verified_path(Some(new_path.clone())); current_path = Some(new_path); } else { warn!("modem re-resolution failed; will retry"); - self.modem_service.set_verified_path(None); current_path = None; } } @@ -210,11 +212,20 @@ impl InboundWorker { async fn run_subscription(&self, actual_path: &str, children: &mut JoinSet<()>) -> Result<()> { let mut subscription = self.source.subscribe(actual_path).await?; + let mut identity_refresh = tokio::time::interval(RUNTIME_IDENTITY_REFRESH_INTERVAL); + identity_refresh.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); + identity_refresh.tick().await; info!("SMS monitor ready on {}", actual_path); loop { let sms = tokio::select! { + _ = identity_refresh.tick() => { + if let Err(error) = self.observe_runtime_identity(actual_path).await { + warn!("runtime modem identity refresh failed: {}", error); + } + continue; + } joined = children.join_next(), if !children.is_empty() => { report_child_result(joined); continue; @@ -243,6 +254,78 @@ impl InboundWorker { }); } } + + async fn observe_runtime_identity(&self, actual_path: &str) -> Result<()> { + retry_pending_identity_mismatches(&self.modem_service, &self.store).await?; + let action_path = self.modem_service.verified_path(); + if let Some(action_path) = action_path.as_deref().filter(|path| *path != actual_path) { + if let Some(identity) = self.modem_service.extract_identity(action_path).await { + let fingerprint = ModemService::compute_fingerprint(&identity); + match self.store.modem_fingerprint().await? { + None => { + self.store.set_modem_fingerprint(fingerprint).await?; + self.store.backfill_dedupe_keys().await?; + } + Some(enrolled) if enrolled == fingerprint => {} + Some(enrolled) => { + warn!("configured modem identity changed; revoking verified action target"); + self.modem_service + .set_modem_targets(ModemTargets::runtime_only(actual_path)); + persist_identity_mismatch( + &self.modem_service, + &self.store, + action_path, + &enrolled, + ) + .await?; + } + } + } + } + + let Some(identity) = self.modem_service.extract_identity(actual_path).await else { + return Ok(()); + }; + let fingerprint = ModemService::compute_fingerprint(&identity); + if action_path.as_deref() == Some(actual_path) { + match self.store.modem_fingerprint().await? { + None => { + self.store.set_modem_fingerprint(fingerprint).await?; + self.store.backfill_dedupe_keys().await?; + } + Some(enrolled) if enrolled == fingerprint => {} + Some(enrolled) => { + warn!("configured modem identity changed; revoking verified action target"); + self.modem_service + .set_modem_targets(ModemTargets::runtime_only(actual_path)); + persist_identity_mismatch( + &self.modem_service, + &self.store, + actual_path, + &enrolled, + ) + .await?; + self.store + .ensure_modem_dedupe_namespace_with(fingerprint.clone()) + .await?; + self.store + .set_runtime_modem_fingerprint(fingerprint) + .await?; + } + } + } else if self.store.runtime_modem_fingerprint().await?.as_deref() + != Some(fingerprint.as_str()) + { + self.store + .ensure_modem_dedupe_namespace_with(fingerprint.clone()) + .await?; + self.store + .set_runtime_modem_fingerprint(fingerprint) + .await?; + self.store.backfill_dedupe_keys().await?; + } + Ok(()) + } } fn report_child_result(joined: Option>) { @@ -256,6 +339,48 @@ fn report_child_failure() { crate::monitoring::capture_failure("dbus", "dbus.inbound_processing_failed"); } +fn publish_resolved_path(modem_service: &ModemService, resolved: ModemTargets) -> Option { + let path = resolved.runtime_path().map(ToString::to_string); + modem_service.set_modem_targets(resolved); + path +} + +async fn persist_identity_mismatch( + modem_service: &ModemService, + store: &Store, + modem_path: &str, + enrolled_fingerprint: &str, +) -> Result<()> { + modem_service.remember_pending_identity_mismatch(modem_path, enrolled_fingerprint); + store + .mark_modem_identity_mismatch(modem_path.to_string(), enrolled_fingerprint.to_string()) + .await?; + modem_service.finish_pending_identity_mismatch(modem_path, enrolled_fingerprint); + Ok(()) +} + +async fn retry_pending_identity_mismatches( + modem_service: &ModemService, + store: &Store, +) -> Result<()> { + for (path, fingerprint) in modem_service.pending_identity_mismatches() { + persist_identity_mismatch(modem_service, store, &path, &fingerprint).await?; + } + Ok(()) +} + +async fn clear_identity_mismatch( + modem_service: &ModemService, + store: &Store, + modem_path: &str, +) -> Result<()> { + store + .clear_modem_identity_mismatch(modem_path.to_string()) + .await?; + modem_service.clear_pending_identity_mismatch(modem_path); + Ok(()) +} + fn next_reconnect_delay(delay: Duration) -> Duration { (delay * 2).min(MAX_RECONNECT_DELAY) } @@ -326,45 +451,186 @@ fn should_ignore_storage(storage: u32, filters: &[StorageType]) -> bool { .any(|filter| !matches!(filter, StorageType::All) && filter.should_ignore(storage)) } -/// Resolve the actual modem path for monitoring. -/// First tries the configured path directly. If it fails and a fingerprint is -/// stored, scans all modems and matches by fingerprint exactly once. +/// Resolve independent runtime and action modem targets. +/// +/// A verified identity can serve both roles. When the configured action path +/// has no readable identity but a different runtime modem is matched, the +/// targets remain separate. A modem selected only by an observed runtime +/// fingerprint or because it is the sole available candidate is runtime-only +/// and must never receive control actions. An action-only result preserves an +/// exact configured target when runtime selection is unavailable or ambiguous; +/// the empty default means neither role could be resolved safely. pub(crate) async fn resolve_monitor_path( configured_path: &str, modem_service: &ModemService, store: &Store, -) -> Result> { - let stored_fingerprint = store.modem_fingerprint().await?; +) -> Result { + retry_pending_identity_mismatches(modem_service, store).await?; + let legacy_fingerprint = + ModemService::compute_fingerprint(LEGACY_SINGLE_MODEM_FINGERPRINT_SEED); + let stored_fingerprint = match store.modem_fingerprint().await? { + Some(fingerprint) if fingerprint == legacy_fingerprint => { + if store.migrate_legacy_modem_fingerprint(fingerprint).await? { + None + } else { + store.modem_fingerprint().await? + } + } + fingerprint => fingerprint, + }; + let runtime_fingerprint = store.runtime_modem_fingerprint().await?; + let needs_fallback_dedupe_namespace = + stored_fingerprint.is_none() && runtime_fingerprint.is_none(); let identity = modem_service.extract_identity(configured_path).await; + let quarantined_fingerprint = store + .modem_identity_mismatch_fingerprint(configured_path.to_string()) + .await?; + let pending_quarantined_fingerprint = modem_service.pending_identity_mismatch(configured_path); + let mut configured_identity_mismatch = stored_fingerprint.as_deref().is_some_and(|enrolled| { + quarantined_fingerprint.as_deref() == Some(enrolled) + || pending_quarantined_fingerprint.as_deref() == Some(enrolled) + }); if let Some(identity) = identity { let current_fingerprint = ModemService::compute_fingerprint(&identity); match stored_fingerprint.as_deref() { Some(enrolled_fingerprint) if enrolled_fingerprint == current_fingerprint => { + clear_identity_mismatch(modem_service, store, configured_path).await?; store.backfill_dedupe_keys().await?; - return Ok(Some(configured_path.to_string())); + return Ok(ModemTargets::verified(configured_path)); } Some(enrolled_fingerprint) => { - warn!("configured modem identity changed; refusing path reuse"); - return Ok(modem_service - .scan_and_match_fingerprint(enrolled_fingerprint) - .await); + configured_identity_mismatch = true; + persist_identity_mismatch( + modem_service, + store, + configured_path, + enrolled_fingerprint, + ) + .await?; + warn!("configured modem identity changed; checking available modems"); } None => { + clear_identity_mismatch(modem_service, store, configured_path).await?; store.set_modem_fingerprint(current_fingerprint).await?; // Backfill dedupe keys for legacy modem-inbound messages now // that the fingerprint is available for stable hashing. store.backfill_dedupe_keys().await?; - return Ok(Some(configured_path.to_string())); + return Ok(ModemTargets::verified(configured_path)); } } } - let Some(stored_fingerprint) = stored_fingerprint.as_deref() else { - return Ok(None); + let paths = modem_service.list_all_modem_paths().await; + let mut candidates = Vec::with_capacity(paths.len()); + for path in paths { + let fingerprint = modem_service + .extract_identity(&path) + .await + .map(|identity| ModemService::compute_fingerprint(&identity)); + candidates.push((path, fingerprint)); + } + let configured_action_path = candidates + .iter() + .find(|(path, fingerprint)| { + !configured_identity_mismatch && path == configured_path && fingerprint.is_none() + }) + .map(|(path, _)| path.clone()); + + if let Some(enrolled_fingerprint) = stored_fingerprint.as_deref() { + let mut matches = candidates.iter().filter_map(|(path, fingerprint)| { + (fingerprint.as_deref() == Some(enrolled_fingerprint)).then(|| path.clone()) + }); + let matched = matches.next(); + let ambiguous = matches.next().is_some(); + if ambiguous { + return Ok(configured_action_path + .clone() + .map(ModemTargets::action_only) + .unwrap_or_default()); + } + if let Some(path) = matched { + clear_identity_mismatch(modem_service, store, &path).await?; + store.backfill_dedupe_keys().await?; + return Ok(ModemTargets::verified(path)); + } + } + + if let Some(observed_fingerprint) = runtime_fingerprint.as_deref() { + let mut matches = candidates.iter().filter_map(|(path, fingerprint)| { + (fingerprint.as_deref() == Some(observed_fingerprint)).then(|| path.clone()) + }); + let matched = matches.next(); + let ambiguous = matches.next().is_some(); + if ambiguous { + return Ok(configured_action_path + .clone() + .map(ModemTargets::action_only) + .unwrap_or_default()); + } + if let Some(path) = matched { + store + .ensure_modem_dedupe_namespace_with(observed_fingerprint.to_string()) + .await?; + store.backfill_dedupe_keys().await?; + return Ok(match configured_action_path { + Some(action_path) => ModemTargets::separate(path, action_path), + None => ModemTargets::runtime_only(path), + }); + } + } + + if candidates.len() != 1 { + return Ok(configured_action_path + .map(ModemTargets::action_only) + .unwrap_or_default()); + } + + let (selected_path, selected_fingerprint) = + candidates.into_iter().next().expect("one modem candidate"); + warn!( + "selecting the only available modem as the runtime target at {}", + selected_path + ); + let resolved = if selected_path == configured_path && !configured_identity_mismatch { + match selected_fingerprint { + Some(fingerprint) if stored_fingerprint.is_none() => { + store.set_modem_fingerprint(fingerprint).await?; + ModemTargets::verified(selected_path) + } + Some(fingerprint) => { + store + .ensure_modem_dedupe_namespace_with(fingerprint.clone()) + .await?; + store.set_runtime_modem_fingerprint(fingerprint).await?; + ModemTargets::runtime_only(selected_path) + } + None => { + if needs_fallback_dedupe_namespace { + warn!("modem identity is unavailable; using a local inbound dedupe namespace"); + store.ensure_modem_dedupe_namespace().await?; + } + ModemTargets::verified(selected_path) + } + } + } else { + match selected_fingerprint { + Some(fingerprint) => { + store + .ensure_modem_dedupe_namespace_with(fingerprint.clone()) + .await?; + store.set_runtime_modem_fingerprint(fingerprint).await?; + } + None => { + if needs_fallback_dedupe_namespace { + warn!("modem identity is unavailable; using a local inbound dedupe namespace"); + store.ensure_modem_dedupe_namespace().await?; + } + } + } + ModemTargets::runtime_only(selected_path) }; - Ok(modem_service - .scan_and_match_fingerprint(stored_fingerprint) - .await) + store.backfill_dedupe_keys().await?; + Ok(resolved) } trait InboundSourceAdapter: Send + Sync { @@ -437,7 +703,7 @@ mod tests { use crate::delivery::DeliveryWakeup; use crate::events::EventBus; use crate::message::{MessageDirection, MessageSource, MessageStatus}; - use crate::modem::{MmcliOutput, MmcliRunner, ModemError}; + use crate::modem::{MmcliOutput, MmcliRunner, ModemAction, ModemError}; const MODEM_PATH: &str = "/org/freedesktop/ModemManager1/Modem/0"; const OTHER_MODEM_PATH: &str = "/org/freedesktop/ModemManager1/Modem/1"; @@ -528,6 +794,7 @@ mod tests { struct ScriptedSource { subscription: Mutex>>, subscribed: Arc, + subscribed_paths: Arc>>, } impl ScriptedSource { @@ -535,6 +802,7 @@ mod tests { Self { subscription: Mutex::new(Some(Box::new(subscription))), subscribed: Arc::new(Notify::new()), + subscribed_paths: Arc::new(Mutex::new(Vec::new())), } } } @@ -542,8 +810,12 @@ mod tests { impl InboundSourceAdapter for ScriptedSource { fn subscribe<'a>( &'a self, - _modem_path: &'a str, + modem_path: &'a str, ) -> BoxFuture<'a, Result>> { + self.subscribed_paths + .lock() + .unwrap() + .push(modem_path.to_string()); let subscription = self.subscription.lock().unwrap().take(); self.subscribed.notify_one(); Box::pin(async move { @@ -581,6 +853,187 @@ mod tests { } } + #[derive(Clone)] + struct SingleModemRunner { + candidate_identity: Option<&'static str>, + } + + impl MmcliRunner for SingleModemRunner { + fn run<'a>( + &'a self, + args: &'a [&'a str], + _timeout: Duration, + ) -> Pin> + Send + 'a>> { + Box::pin(async move { + let stdout = match args { + ["-L"] => { + format!("{OTHER_MODEM_PATH} [test] modem\n") + } + ["--modem", OTHER_MODEM_PATH, "--output-json"] => self + .candidate_identity + .map(|identity| { + format!( + r#"{{"modem":{{"generic":{{"equipment-identifier":"{identity}"}}}}}}"# + ) + }) + .unwrap_or_default(), + _ => String::new(), + }; + Ok(MmcliOutput { + status_success: !stdout.is_empty(), + stdout, + stderr: String::new(), + }) + }) + } + } + + #[derive(Clone)] + struct RecoveringIdentityRunner { + candidate_identity: Arc>>, + } + + impl MmcliRunner for RecoveringIdentityRunner { + fn run<'a>( + &'a self, + args: &'a [&'a str], + _timeout: Duration, + ) -> Pin> + Send + 'a>> { + Box::pin(async move { + let stdout = match args { + ["-L"] => format!("{OTHER_MODEM_PATH} [test] modem\n"), + ["--modem", OTHER_MODEM_PATH, "--output-json"] => self + .candidate_identity + .lock() + .unwrap() + .map(|identity| { + format!( + r#"{{"modem":{{"generic":{{"equipment-identifier":"{identity}"}}}}}}"# + ) + }) + .unwrap_or_default(), + _ => String::new(), + }; + Ok(MmcliOutput { + status_success: !stdout.is_empty(), + stdout, + stderr: String::new(), + }) + }) + } + } + + #[derive(Clone)] + struct MultipleModemsRunner; + + impl MmcliRunner for MultipleModemsRunner { + fn run<'a>( + &'a self, + args: &'a [&'a str], + _timeout: Duration, + ) -> Pin> + Send + 'a>> { + Box::pin(async move { + let stdout = match args { + ["-L"] => format!( + "{OTHER_MODEM_PATH} [test] first modem\n/org/freedesktop/ModemManager1/Modem/2 [test] second modem\n" + ), + _ => String::new(), + }; + Ok(MmcliOutput { + status_success: !stdout.is_empty(), + stdout, + stderr: String::new(), + }) + }) + } + } + + #[derive(Clone)] + struct ConfiguredModemWithoutIdentityRunner; + + impl MmcliRunner for ConfiguredModemWithoutIdentityRunner { + fn run<'a>( + &'a self, + args: &'a [&'a str], + _timeout: Duration, + ) -> Pin> + Send + 'a>> { + Box::pin(async move { + let stdout = match args { + ["-L"] => format!("{MODEM_PATH} [test] configured modem\n"), + _ => String::new(), + }; + Ok(MmcliOutput { + status_success: !stdout.is_empty(), + stdout, + stderr: String::new(), + }) + }) + } + } + + #[derive(Clone)] + struct ConfiguredIdentityRunner { + identity: Arc>>, + } + + impl MmcliRunner for ConfiguredIdentityRunner { + fn run<'a>( + &'a self, + args: &'a [&'a str], + _timeout: Duration, + ) -> Pin> + Send + 'a>> { + Box::pin(async move { + let stdout = match args { + ["-L"] => format!("{MODEM_PATH} [test] configured modem\n"), + ["--modem", MODEM_PATH, "--output-json"] => self + .identity + .lock() + .unwrap() + .map(|identity| { + format!( + r#"{{"modem":{{"generic":{{"equipment-identifier":"{identity}"}}}}}}"# + ) + }) + .unwrap_or_default(), + _ => String::new(), + }; + Ok(MmcliOutput { + status_success: !stdout.is_empty(), + stdout, + stderr: String::new(), + }) + }) + } + } + + #[derive(Clone)] + struct RuntimeAndConfiguredModemsRunner; + + impl MmcliRunner for RuntimeAndConfiguredModemsRunner { + fn run<'a>( + &'a self, + args: &'a [&'a str], + _timeout: Duration, + ) -> Pin> + Send + 'a>> { + Box::pin(async move { + let stdout = match args { + ["-L"] => format!( + "{MODEM_PATH} [test] configured modem\n{OTHER_MODEM_PATH} [test] runtime modem\n" + ), + ["--modem", OTHER_MODEM_PATH, "--output-json"] => { + r#"{"modem":{"generic":{"equipment-identifier":"target"}}}"#.to_string() + } + _ => String::new(), + }; + Ok(MmcliOutput { + status_success: !stdout.is_empty(), + stdout, + stderr: String::new(), + }) + }) + } + } + fn properties(body: &str, storage: u32) -> InboundSmsProperties { InboundSmsProperties { phone_number: "+15550000000".to_string(), @@ -909,13 +1362,524 @@ mod tests { .await .unwrap(); - assert_eq!(resolved.as_deref(), Some(OTHER_MODEM_PATH)); + assert_eq!(resolved.runtime_path(), Some(OTHER_MODEM_PATH)); assert_eq!( store.modem_fingerprint().await.unwrap().as_deref(), Some(target.as_str()) ); } + #[tokio::test] + async fn stale_configured_path_selects_the_only_available_modem() { + let store = Store::open_in_memory().unwrap(); + let service = ModemService::new_with_runner(SingleModemRunner { + candidate_identity: Some("target"), + }); + + let resolved = resolve_monitor_path(MODEM_PATH, &service, &store) + .await + .unwrap(); + + assert_eq!(resolved.runtime_path(), Some(OTHER_MODEM_PATH)); + assert_eq!( + store.runtime_modem_fingerprint().await.unwrap().as_deref(), + Some(ModemService::compute_fingerprint("target").as_str()) + ); + assert_eq!(store.modem_fingerprint().await.unwrap(), None); + } + + #[tokio::test] + async fn observed_runtime_identity_never_becomes_a_verified_action_target() { + let store = Store::open_in_memory().unwrap(); + let service = ModemService::new_with_runner(SingleModemRunner { + candidate_identity: Some("target"), + }); + + let first = resolve_monitor_path(MODEM_PATH, &service, &store) + .await + .unwrap(); + publish_resolved_path(&service, first); + assert_eq!(service.verified_path(), None); + + let second = resolve_monitor_path(MODEM_PATH, &service, &store) + .await + .unwrap(); + publish_resolved_path(&service, second); + + assert_eq!(service.verified_path(), None); + assert_eq!(store.modem_fingerprint().await.unwrap(), None); + } + + #[tokio::test] + async fn worker_keeps_the_only_available_modem_out_of_the_verified_action_path() { + let source = Arc::new(ScriptedSource::new(ScriptedSubscription { + messages: VecDeque::new(), + reads: Arc::new(AtomicUsize::new(0)), + dropped: Arc::new(AtomicBool::new(false)), + })); + let subscribed = source.subscribed.clone(); + let subscribed_paths = source.subscribed_paths.clone(); + let store = Store::open_in_memory().unwrap(); + let modem_service = ModemService::new_with_runner(SingleModemRunner { + candidate_identity: Some("target"), + }); + let worker = InboundWorker::new( + store.clone(), + messaging(store), + modem_service.clone(), + settings(Vec::new(), Vec::new()), + ) + .with_source(source); + + let run = tokio::spawn(async move { worker.run().await }); + tokio::time::timeout(Duration::from_secs(1), subscribed.notified()) + .await + .expect("worker subscribes after resolving the modem"); + + assert_eq!( + subscribed_paths.lock().unwrap().as_slice(), + [OTHER_MODEM_PATH] + ); + assert_eq!(modem_service.verified_path(), None); + let action_error = modem_service + .run_action("test-session", ModemAction::Disable) + .await + .unwrap_err(); + assert_eq!(action_error.code(), "modem_path_unresolved"); + + run.abort(); + let _ = run.await; + } + + #[tokio::test] + async fn worker_observes_runtime_identity_when_it_becomes_available() { + let source = Arc::new(ScriptedSource::new(ScriptedSubscription { + messages: VecDeque::new(), + reads: Arc::new(AtomicUsize::new(0)), + dropped: Arc::new(AtomicBool::new(false)), + })); + let subscribed = source.subscribed.clone(); + let candidate_identity = Arc::new(Mutex::new(None)); + let store = Store::open_in_memory().unwrap(); + let modem_service = ModemService::new_with_runner(RecoveringIdentityRunner { + candidate_identity: candidate_identity.clone(), + }); + let worker = InboundWorker::new( + store.clone(), + messaging(store.clone()), + modem_service.clone(), + settings(Vec::new(), Vec::new()), + ) + .with_source(source); + + let run = tokio::spawn(async move { worker.run().await }); + subscribed.notified().await; + assert_eq!(store.runtime_modem_fingerprint().await.unwrap(), None); + + *candidate_identity.lock().unwrap() = Some("target"); + let expected = ModemService::compute_fingerprint("target"); + tokio::time::timeout(Duration::from_secs(1), async { + while store.runtime_modem_fingerprint().await.unwrap().as_deref() + != Some(expected.as_str()) + { + tokio::task::yield_now().await; + } + }) + .await + .expect("the refresh observes the recovered runtime identity"); + + assert_eq!( + store.runtime_modem_fingerprint().await.unwrap().as_deref(), + Some(expected.as_str()) + ); + assert_eq!(modem_service.verified_path(), None); + + run.abort(); + let _ = run.await; + } + + #[tokio::test] + async fn only_available_modem_is_selected_when_identity_is_temporarily_unavailable() { + let store = Store::open_in_memory().unwrap(); + let enrolled = ModemService::compute_fingerprint("target"); + store.set_modem_fingerprint(enrolled.clone()).await.unwrap(); + let service = ModemService::new_with_runner(SingleModemRunner { + candidate_identity: None, + }); + + let resolved = resolve_monitor_path(MODEM_PATH, &service, &store) + .await + .unwrap(); + + assert_eq!(resolved.runtime_path(), Some(OTHER_MODEM_PATH)); + assert_eq!( + store.modem_fingerprint().await.unwrap().as_deref(), + Some(enrolled.as_str()) + ); + } + + #[tokio::test] + async fn temporary_identity_loss_keeps_the_existing_inbound_dedupe_namespace() { + let store = Store::open_in_memory().unwrap(); + let enrolled = ModemService::compute_fingerprint("target"); + store.set_modem_fingerprint(enrolled).await.unwrap(); + process_incoming_sms( + Box::new(ScriptedSms::new( + SMS_PATH, + vec![PropertyAction::Return(Ok(properties( + "same message", + StorageType::Me as u32, + )))], + )), + &[], + messaging(store.clone()), + Vec::new(), + ) + .await + .unwrap(); + + let unavailable = ModemService::new_with_runner(SingleModemRunner { + candidate_identity: None, + }); + resolve_monitor_path(MODEM_PATH, &unavailable, &store) + .await + .unwrap(); + process_incoming_sms( + Box::new(ScriptedSms::new( + "/org/freedesktop/ModemManager1/SMS/99", + vec![PropertyAction::Return(Ok(properties( + "same message", + StorageType::Me as u32, + )))], + )), + &[], + messaging(store.clone()), + Vec::new(), + ) + .await + .unwrap(); + + assert_eq!(store.sqlite().count_messages().unwrap(), 1); + } + + #[tokio::test] + async fn configured_path_remains_a_verified_action_target_when_identity_is_unavailable() { + let store = Store::open_in_memory().unwrap(); + let service = ModemService::new_with_runner(ConfiguredModemWithoutIdentityRunner); + + let resolved = resolve_monitor_path(MODEM_PATH, &service, &store) + .await + .unwrap(); + publish_resolved_path(&service, resolved); + + assert_eq!(service.verified_path().as_deref(), Some(MODEM_PATH)); + assert_eq!(store.modem_fingerprint().await.unwrap(), None); + } + + #[tokio::test] + async fn configured_identity_mismatch_stays_quarantined_when_identity_becomes_unavailable() { + let store = Store::open_in_memory().unwrap(); + store + .set_modem_fingerprint(ModemService::compute_fingerprint("target")) + .await + .unwrap(); + let identity = Arc::new(Mutex::new(Some("other"))); + let service = ModemService::new_with_runner(ConfiguredIdentityRunner { + identity: identity.clone(), + }); + + let mismatched = resolve_monitor_path(MODEM_PATH, &service, &store) + .await + .unwrap(); + assert_eq!(mismatched.runtime_path(), Some(MODEM_PATH)); + assert_eq!(mismatched.action_path(), None); + assert_eq!( + store + .modem_identity_mismatch_fingerprint(MODEM_PATH.to_string()) + .await + .unwrap() + .as_deref(), + Some(ModemService::compute_fingerprint("target").as_str()) + ); + + *identity.lock().unwrap() = None; + let unavailable = resolve_monitor_path(MODEM_PATH, &service, &store) + .await + .unwrap(); + + assert_eq!(unavailable.runtime_path(), Some(MODEM_PATH)); + assert_eq!(unavailable.action_path(), None); + + *identity.lock().unwrap() = Some("target"); + let recovered = resolve_monitor_path(MODEM_PATH, &service, &store) + .await + .unwrap(); + + assert_eq!(recovered.action_path(), Some(MODEM_PATH)); + assert_eq!( + store + .modem_identity_mismatch_fingerprint(MODEM_PATH.to_string()) + .await + .unwrap(), + None + ); + } + + #[tokio::test] + async fn failed_mismatch_persistence_still_blocks_action_during_identity_loss() { + let store = Store::open_in_memory().unwrap(); + store + .set_modem_fingerprint(ModemService::compute_fingerprint("target")) + .await + .unwrap(); + store.fail_next_modem_identity_mismatch_marks(1); + let identity = Arc::new(Mutex::new(Some("other"))); + let modem_service = ModemService::new_with_runner(ConfiguredIdentityRunner { + identity: identity.clone(), + }); + modem_service.set_verified_path(Some(MODEM_PATH.to_string())); + let worker = InboundWorker::new( + store.clone(), + messaging(store.clone()), + modem_service.clone(), + settings(Vec::new(), Vec::new()), + ); + + assert!(worker.observe_runtime_identity(MODEM_PATH).await.is_err()); + assert_eq!(modem_service.verified_path(), None); + *identity.lock().unwrap() = None; + + let resolved = resolve_monitor_path(MODEM_PATH, &modem_service, &store) + .await + .unwrap(); + + assert_eq!(resolved.runtime_path(), Some(MODEM_PATH)); + assert_eq!(resolved.action_path(), None); + assert_eq!( + store + .modem_identity_mismatch_fingerprint(MODEM_PATH.to_string()) + .await + .unwrap() + .as_deref(), + Some(ModemService::compute_fingerprint("target").as_str()) + ); + } + + #[tokio::test] + async fn runtime_rebind_keeps_the_exact_configured_action_target() { + let store = Store::open_in_memory().unwrap(); + store + .set_runtime_modem_fingerprint(ModemService::compute_fingerprint("target")) + .await + .unwrap(); + let service = ModemService::new_with_runner(RuntimeAndConfiguredModemsRunner); + + let resolved = resolve_monitor_path(MODEM_PATH, &service, &store) + .await + .unwrap(); + publish_resolved_path(&service, resolved); + + assert_eq!(service.runtime_path().as_deref(), Some(OTHER_MODEM_PATH)); + assert_eq!(service.verified_path().as_deref(), Some(MODEM_PATH)); + } + + #[tokio::test] + async fn runtime_rebind_dedupe_survives_configured_action_identity_enrollment() { + let store = Store::open_in_memory().unwrap(); + store + .set_runtime_modem_fingerprint(ModemService::compute_fingerprint("target")) + .await + .unwrap(); + let service = ModemService::new_with_runner(RuntimeAndConfiguredModemsRunner); + resolve_monitor_path(MODEM_PATH, &service, &store) + .await + .unwrap(); + process_incoming_sms( + Box::new(ScriptedSms::new( + SMS_PATH, + vec![PropertyAction::Return(Ok(properties( + "same message", + StorageType::Me as u32, + )))], + )), + &[], + messaging(store.clone()), + Vec::new(), + ) + .await + .unwrap(); + + store + .set_modem_fingerprint(ModemService::compute_fingerprint("configured")) + .await + .unwrap(); + process_incoming_sms( + Box::new(ScriptedSms::new( + "/org/freedesktop/ModemManager1/SMS/99", + vec![PropertyAction::Return(Ok(properties( + "same message", + StorageType::Me as u32, + )))], + )), + &[], + messaging(store.clone()), + Vec::new(), + ) + .await + .unwrap(); + + assert_eq!(store.sqlite().count_messages().unwrap(), 1); + } + + #[tokio::test] + async fn first_only_modem_without_identity_receives_without_enrolling_a_fake_identity() { + let store = Store::open_in_memory().unwrap(); + let service = ModemService::new_with_runner(SingleModemRunner { + candidate_identity: None, + }); + + let resolved = resolve_monitor_path(MODEM_PATH, &service, &store) + .await + .unwrap(); + + assert_eq!(resolved.runtime_path(), Some(OTHER_MODEM_PATH)); + assert_eq!(store.modem_fingerprint().await.unwrap(), None); + + process_incoming_sms( + Box::new(ScriptedSms::new( + SMS_PATH, + vec![PropertyAction::Return(Ok(properties( + "identity unavailable", + StorageType::Me as u32, + )))], + )), + &[], + messaging(store.clone()), + Vec::new(), + ) + .await + .unwrap(); + assert_eq!(store.sqlite().count_messages().unwrap(), 1); + } + + #[tokio::test] + async fn recovered_runtime_identity_is_observed_without_changing_inbound_dedupe() { + let store = Store::open_in_memory().unwrap(); + let unavailable = ModemService::new_with_runner(SingleModemRunner { + candidate_identity: None, + }); + resolve_monitor_path(MODEM_PATH, &unavailable, &store) + .await + .unwrap(); + process_incoming_sms( + Box::new(ScriptedSms::new( + SMS_PATH, + vec![PropertyAction::Return(Ok(properties( + "same message", + StorageType::Me as u32, + )))], + )), + &[], + messaging(store.clone()), + Vec::new(), + ) + .await + .unwrap(); + + let recovered = ModemService::new_with_runner(SingleModemRunner { + candidate_identity: Some("target"), + }); + resolve_monitor_path(MODEM_PATH, &recovered, &store) + .await + .unwrap(); + assert_eq!( + store.runtime_modem_fingerprint().await.unwrap().as_deref(), + Some(ModemService::compute_fingerprint("target").as_str()) + ); + assert_eq!(store.modem_fingerprint().await.unwrap(), None); + + process_incoming_sms( + Box::new(ScriptedSms::new( + "/org/freedesktop/ModemManager1/SMS/99", + vec![PropertyAction::Return(Ok(properties( + "same message", + StorageType::Me as u32, + )))], + )), + &[], + messaging(store.clone()), + Vec::new(), + ) + .await + .unwrap(); + assert_eq!(store.sqlite().count_messages().unwrap(), 1); + } + + #[tokio::test] + async fn legacy_single_modem_fingerprint_migrates_without_changing_inbound_dedupe() { + let store = Store::open_in_memory().unwrap(); + let legacy_fingerprint = ModemService::compute_fingerprint("sms-relayed-single-modem"); + store + .set_modem_fingerprint(legacy_fingerprint) + .await + .unwrap(); + process_incoming_sms( + Box::new(ScriptedSms::new( + SMS_PATH, + vec![PropertyAction::Return(Ok(properties( + "same message", + StorageType::Me as u32, + )))], + )), + &[], + messaging(store.clone()), + Vec::new(), + ) + .await + .unwrap(); + + let recovered = ModemService::new_with_runner(SingleModemRunner { + candidate_identity: Some("target"), + }); + resolve_monitor_path(MODEM_PATH, &recovered, &store) + .await + .unwrap(); + + assert_eq!(store.modem_fingerprint().await.unwrap(), None); + assert_eq!( + store.runtime_modem_fingerprint().await.unwrap().as_deref(), + Some(ModemService::compute_fingerprint("target").as_str()) + ); + process_incoming_sms( + Box::new(ScriptedSms::new( + "/org/freedesktop/ModemManager1/SMS/99", + vec![PropertyAction::Return(Ok(properties( + "same message", + StorageType::Me as u32, + )))], + )), + &[], + messaging(store.clone()), + Vec::new(), + ) + .await + .unwrap(); + assert_eq!(store.sqlite().count_messages().unwrap(), 1); + } + + #[tokio::test] + async fn stale_configured_path_does_not_select_from_multiple_unidentified_modems() { + let store = Store::open_in_memory().unwrap(); + let service = ModemService::new_with_runner(MultipleModemsRunner); + + let resolved = resolve_monitor_path(MODEM_PATH, &service, &store) + .await + .unwrap(); + + assert_eq!(resolved, ModemTargets::default()); + assert_eq!(store.modem_fingerprint().await.unwrap(), None); + } + #[tokio::test] async fn matching_enrolled_fingerprint_backfills_before_monitoring() { let store = Store::open_in_memory().unwrap(); @@ -942,7 +1906,7 @@ mod tests { .await .unwrap(); - assert_eq!(resolved.as_deref(), Some(MODEM_PATH)); + assert_eq!(resolved.runtime_path(), Some(MODEM_PATH)); let replay = crate::storage::NewMessage::modem_inbound( "+1", "legacy", diff --git a/src/messaging.rs b/src/messaging.rs index ef2bb4e..9cd1fdc 100644 --- a/src/messaging.rs +++ b/src/messaging.rs @@ -62,7 +62,7 @@ pub struct Messaging { events: EventBus, delivery_wakeup: DeliveryWakeup, sms_sender: Arc, - verified_modem: Option, + modem_service: Option, } impl Messaging { @@ -82,12 +82,12 @@ impl Messaging { events, delivery_wakeup, sms_sender, - verified_modem: None, + modem_service: None, } } - pub fn with_verified_modem(mut self, modem: crate::modem::ModemService) -> Self { - self.verified_modem = Some(modem); + pub fn with_modem_service(mut self, modem: crate::modem::ModemService) -> Self { + self.modem_service = Some(modem); self } @@ -115,10 +115,10 @@ impl Messaging { mut request: SendMessage, owner: String, ) -> anyhow::Result { - if let Some(modem) = self.verified_modem.as_ref() { + if let Some(modem) = self.modem_service.as_ref() { request.modem_path = modem - .verified_path() - .ok_or_else(|| anyhow::anyhow!("verified modem identity is not ready"))?; + .runtime_path() + .ok_or_else(|| anyhow::anyhow!("runtime modem path is not ready"))?; } let SendMessage { ref phone_number, @@ -595,11 +595,11 @@ impl Messaging { owner: &str, modem_sms_path: &str, ) -> anyhow::Result { - let verified_path = match self.verified_modem.as_ref() { + let runtime_path = match self.modem_service.as_ref() { Some(modem) => Some( modem - .verified_path() - .ok_or_else(|| anyhow::anyhow!("verified modem identity is not ready"))?, + .runtime_path() + .ok_or_else(|| anyhow::anyhow!("runtime modem path is not ready"))?, ), None => None, }; @@ -607,7 +607,7 @@ impl Messaging { message_id, owner, self.sms_sender - .sms_snapshot(verified_path.as_deref(), modem_sms_path), + .sms_snapshot(runtime_path.as_deref(), modem_sms_path), ) .await } @@ -1863,7 +1863,7 @@ mod tests { } #[tokio::test] - async fn inbound_without_an_enrolled_fingerprint_is_not_persisted() { + async fn inbound_without_a_dedupe_namespace_is_not_persisted() { let store = Store::open_in_memory().unwrap(); let messaging = Messaging::new( store.clone(), diff --git a/src/modem.rs b/src/modem.rs index 4668047..6afb8c6 100644 --- a/src/modem.rs +++ b/src/modem.rs @@ -537,6 +537,56 @@ pub enum ModemAction { Reset, } +/// The independently resolved paths for SMS traffic and modem control actions. +/// Fields remain private so callers can only publish one complete snapshot. +#[derive(Debug, Clone, Default, PartialEq, Eq)] +pub(crate) struct ModemTargets { + runtime: Option, + action: Option, +} + +impl ModemTargets { + pub(crate) fn runtime_only(path: impl Into) -> Self { + Self { + runtime: Some(path.into()), + action: None, + } + } + + pub(crate) fn verified(path: impl Into) -> Self { + let path = path.into(); + Self { + runtime: Some(path.clone()), + action: Some(path), + } + } + + pub(crate) fn separate( + runtime_path: impl Into, + action_path: impl Into, + ) -> Self { + Self { + runtime: Some(runtime_path.into()), + action: Some(action_path.into()), + } + } + + pub(crate) fn action_only(path: impl Into) -> Self { + Self { + runtime: None, + action: Some(path.into()), + } + } + + pub(crate) fn runtime_path(&self) -> Option<&str> { + self.runtime.as_deref() + } + + pub(crate) fn action_path(&self) -> Option<&str> { + self.action.as_deref() + } +} + #[derive(Debug, Clone, Serialize)] pub struct ActionResponse { pub accepted: bool, @@ -647,7 +697,8 @@ pub struct ModemService { pub(crate) action_lock: Arc>, reset_limits: Arc>>, health_refresh_lock: Arc>, - verified_path: Arc>>, + modem_targets: Arc>, + pending_identity_mismatches: Arc>>, } impl ModemService { @@ -676,16 +727,85 @@ impl ModemService { action_lock: Arc::new(tokio::sync::Mutex::new(())), reset_limits: Arc::new(Mutex::new(HashMap::new())), health_refresh_lock: Arc::new(tokio::sync::Mutex::new(())), - verified_path: Arc::new(Mutex::new(None)), + modem_targets: Arc::new(Mutex::new(ModemTargets::default())), + pending_identity_mismatches: Arc::new(Mutex::new(HashMap::new())), + } + } + + pub(crate) fn set_modem_targets(&self, targets: ModemTargets) { + *self.modem_targets.lock().unwrap() = targets; + } + + pub(crate) fn remember_pending_identity_mismatch( + &self, + modem_path: impl Into, + enrolled_fingerprint: impl Into, + ) { + self.pending_identity_mismatches + .lock() + .unwrap() + .insert(modem_path.into(), enrolled_fingerprint.into()); + } + + pub(crate) fn pending_identity_mismatch(&self, modem_path: &str) -> Option { + self.pending_identity_mismatches + .lock() + .unwrap() + .get(modem_path) + .cloned() + } + + pub(crate) fn pending_identity_mismatches(&self) -> Vec<(String, String)> { + self.pending_identity_mismatches + .lock() + .unwrap() + .iter() + .map(|(path, fingerprint)| (path.clone(), fingerprint.clone())) + .collect() + } + + pub(crate) fn finish_pending_identity_mismatch( + &self, + modem_path: &str, + enrolled_fingerprint: &str, + ) { + let mut pending = self.pending_identity_mismatches.lock().unwrap(); + if pending.get(modem_path).map(String::as_str) == Some(enrolled_fingerprint) { + pending.remove(modem_path); } } + pub(crate) fn clear_pending_identity_mismatch(&self, modem_path: &str) { + self.pending_identity_mismatches + .lock() + .unwrap() + .remove(modem_path); + } + + #[cfg(test)] + pub(crate) fn set_runtime_path(&self, path: Option) { + self.set_modem_targets(path.map(ModemTargets::runtime_only).unwrap_or_default()); + } + + #[cfg(test)] pub(crate) fn set_verified_path(&self, path: Option) { - *self.verified_path.lock().unwrap() = path; + self.set_modem_targets(path.map(ModemTargets::verified).unwrap_or_default()); + } + + pub fn runtime_path(&self) -> Option { + self.modem_targets + .lock() + .unwrap() + .runtime_path() + .map(ToString::to_string) } pub fn verified_path(&self) -> Option { - self.verified_path.lock().unwrap().clone() + self.modem_targets + .lock() + .unwrap() + .action_path() + .map(ToString::to_string) } pub async fn status(&self, configured_path: &str) -> ModemStatus { @@ -1129,24 +1249,6 @@ impl ModemService { } } - pub async fn scan_and_match_fingerprint(&self, target_fingerprint: &str) -> Option { - let paths = self.list_all_modem_paths().await; - let mut matched: Vec = Vec::new(); - for path in &paths { - if let Some(identity) = self.extract_identity(path).await { - let fp = Self::compute_fingerprint(&identity); - if fp == target_fingerprint { - matched.push(path.clone()); - } - } - } - if matched.len() == 1 { - Some(matched.remove(0)) - } else { - None - } - } - #[allow(dead_code)] pub fn runner_ref(&self) -> &Arc { &self.runner diff --git a/src/persistence/mod.rs b/src/persistence/mod.rs index 6c69ab9..007c9b7 100644 --- a/src/persistence/mod.rs +++ b/src/persistence/mod.rs @@ -21,6 +21,9 @@ pub use delivery::{ }; const MODEM_FINGERPRINT_META_KEY: &str = "modem_fingerprint"; +const MODEM_DEDUPE_NAMESPACE_META_KEY: &str = "modem_dedupe_namespace"; +const RUNTIME_MODEM_FINGERPRINT_META_KEY: &str = "runtime_modem_fingerprint"; +const MODEM_IDENTITY_MISMATCH_META_KEY_PREFIX: &str = "modem_identity_mismatch:"; fn outbound_phase_to_str(phase: OutboundPhase) -> &'static str { match phase { @@ -98,6 +101,8 @@ pub struct Store { #[cfg(test)] outbound_finalization_failures: Arc, #[cfg(test)] + modem_identity_mismatch_failures: Arc, + #[cfg(test)] outbound_creation_pause: Arc>>, } @@ -115,6 +120,8 @@ impl From for Store { #[cfg(test)] outbound_finalization_failures: Arc::new(AtomicUsize::new(0)), #[cfg(test)] + modem_identity_mismatch_failures: Arc::new(AtomicUsize::new(0)), + #[cfg(test)] outbound_creation_pause: Arc::new(Mutex::new(None)), } } @@ -155,16 +162,16 @@ impl Store { ) -> Result { let result = self .run(move |sqlite| { - let fingerprint = sqlite - .get_meta(MODEM_FINGERPRINT_META_KEY)? + let dedupe_namespace = sqlite + .inbound_dedupe_namespace()? .filter(|value| !value.is_empty()) - .ok_or_else(|| anyhow::anyhow!("modem fingerprint is not enrolled"))?; + .ok_or_else(|| anyhow::anyhow!("modem dedupe namespace is not enrolled"))?; let message = NewMessage::modem_inbound( &input.phone_number, &input.body, &input.timestamp, &input.modem_sms_path, - &fingerprint, + &dedupe_namespace, ); sqlite.insert_inbound_message_with_deliveries(message, &profile_keys) }) @@ -416,6 +423,73 @@ impl Store { .await } + pub async fn runtime_modem_fingerprint(&self) -> Result> { + self.run(|sqlite| sqlite.get_meta(RUNTIME_MODEM_FINGERPRINT_META_KEY)) + .await + } + + pub async fn set_runtime_modem_fingerprint(&self, fingerprint: String) -> Result<()> { + self.run(move |sqlite| sqlite.set_meta(RUNTIME_MODEM_FINGERPRINT_META_KEY, &fingerprint)) + .await + } + + pub async fn modem_identity_mismatch_fingerprint( + &self, + modem_path: String, + ) -> Result> { + let key = format!("{MODEM_IDENTITY_MISMATCH_META_KEY_PREFIX}{modem_path}"); + self.run(move |sqlite| sqlite.get_meta(&key)).await + } + + pub async fn mark_modem_identity_mismatch( + &self, + modem_path: String, + enrolled_fingerprint: String, + ) -> Result<()> { + #[cfg(test)] + if self + .modem_identity_mismatch_failures + .fetch_update(Ordering::SeqCst, Ordering::SeqCst, |remaining| { + remaining.checked_sub(1) + }) + .is_ok() + { + anyhow::bail!("injected modem identity mismatch persistence failure"); + } + let key = format!("{MODEM_IDENTITY_MISMATCH_META_KEY_PREFIX}{modem_path}"); + self.run(move |sqlite| sqlite.set_meta(&key, &enrolled_fingerprint)) + .await + } + + #[cfg(test)] + pub(crate) fn fail_next_modem_identity_mismatch_marks(&self, count: usize) { + self.modem_identity_mismatch_failures + .store(count, Ordering::SeqCst); + } + + pub async fn clear_modem_identity_mismatch(&self, modem_path: String) -> Result<()> { + let key = format!("{MODEM_IDENTITY_MISMATCH_META_KEY_PREFIX}{modem_path}"); + self.run(move |sqlite| sqlite.delete_meta(&key)).await + } + + pub async fn migrate_legacy_modem_fingerprint( + &self, + legacy_fingerprint: String, + ) -> Result { + self.run(move |sqlite| sqlite.migrate_legacy_modem_fingerprint(&legacy_fingerprint)) + .await + } + + pub async fn ensure_modem_dedupe_namespace(&self) -> Result { + self.ensure_modem_dedupe_namespace_with(uuid::Uuid::new_v4().to_string()) + .await + } + + pub async fn ensure_modem_dedupe_namespace_with(&self, candidate: String) -> Result { + self.run(move |sqlite| sqlite.ensure_meta(MODEM_DEDUPE_NAMESPACE_META_KEY, &candidate)) + .await + } + pub async fn backfill_dedupe_keys(&self) -> Result { self.run(|sqlite| sqlite.backfill_dedupe_keys()).await } diff --git a/src/runtime.rs b/src/runtime.rs index 08f3d89..9e873ab 100644 --- a/src/runtime.rs +++ b/src/runtime.rs @@ -35,7 +35,7 @@ pub async fn run_forwarding(config_path: &Path) -> Result<()> { delivery_wakeup.clone(), sms_sender.clone(), ) - .with_verified_modem(modem_service.clone()); + .with_modem_service(modem_service.clone()); let client = Arc::new(build_http_client(&config.http)); let webhook_client = Arc::new(build_webhook_http_client(&config.http)); @@ -167,13 +167,14 @@ pub async fn send_interactive(config_path: &Path) -> Result<()> { let sms_sender = Arc::new(dbus::SystemSmsSender::connect().await?); let store = Store::open(Path::new(&config.api.database_path)).await?; let modem_service = ModemService::new(); - let verified_path = - inbound::resolve_monitor_path(&config.app.modem_path, &modem_service, &store) - .await? - .ok_or_else(|| anyhow::anyhow!("no verified modem identity available"))?; - modem_service.set_verified_path(Some(verified_path)); + let modem_targets = + inbound::resolve_monitor_path(&config.app.modem_path, &modem_service, &store).await?; + if modem_targets.runtime_path().is_none() { + return Err(anyhow::anyhow!("no runtime modem path available")); + } + modem_service.set_modem_targets(modem_targets); let messaging = Messaging::new(store, EventBus::new(), DeliveryWakeup::new(), sms_sender) - .with_verified_modem(modem_service); + .with_modem_service(modem_service); if messaging.has_pending_outbound().await? && !Confirm::new( "An outbound SMS is unresolved. Sending another may duplicate it. Continue anyway?", diff --git a/src/storage/metadata.rs b/src/storage/metadata.rs index 3fdb5b2..a040096 100644 --- a/src/storage/metadata.rs +++ b/src/storage/metadata.rs @@ -31,17 +31,62 @@ impl MessageStore { )?; Ok(()) } -} -pub(super) fn backfill_dedupe_keys_on(conn: &Connection) -> Result { - let fingerprint: Option = conn - .query_row( - "SELECT value FROM meta WHERE key = 'modem_fingerprint'", - [], + pub fn delete_meta(&self, key: &str) -> Result<()> { + let conn = self.conn.lock().unwrap(); + conn.execute("DELETE FROM meta WHERE key = ?1", params![key])?; + Ok(()) + } + + pub fn ensure_meta(&self, key: &str, value: &str) -> Result { + let conn = self.conn.lock().unwrap(); + conn.execute( + "INSERT OR IGNORE INTO meta (key, value) VALUES (?1, ?2)", + params![key, value], + )?; + conn.query_row( + "SELECT value FROM meta WHERE key = ?1", + params![key], |row| row.get(0), ) - .optional()?; - let Some(fingerprint) = fingerprint else { + .map_err(Into::into) + } + + pub fn inbound_dedupe_namespace(&self) -> Result> { + let conn = self.conn.lock().unwrap(); + inbound_dedupe_namespace_on(&conn) + } + + pub fn migrate_legacy_modem_fingerprint(&self, legacy_fingerprint: &str) -> Result { + let mut conn = self.conn.lock().unwrap(); + let tx = conn.transaction()?; + let current: Option = tx + .query_row( + "SELECT value FROM meta WHERE key = 'modem_fingerprint'", + [], + |row| row.get(0), + ) + .optional()?; + if current.as_deref() != Some(legacy_fingerprint) { + return Ok(false); + } + + tx.execute( + "INSERT OR IGNORE INTO meta (key, value) + VALUES ('modem_dedupe_namespace', ?1)", + params![legacy_fingerprint], + )?; + tx.execute( + "DELETE FROM meta WHERE key = 'modem_fingerprint' AND value = ?1", + params![legacy_fingerprint], + )?; + tx.commit()?; + Ok(true) + } +} + +pub(super) fn backfill_dedupe_keys_on(conn: &Connection) -> Result { + let Some(dedupe_namespace) = inbound_dedupe_namespace_on(conn)? else { return Ok(0); }; let mut statement = conn.prepare( @@ -60,7 +105,7 @@ pub(super) fn backfill_dedupe_keys_on(conn: &Connection) -> Result { conn.prepare("SELECT COUNT(*) > 0 FROM messages WHERE inbound_dedupe_key = ?1")?; let mut count = 0; for (id, phone, body, timestamp) in &rows { - let dedupe_key = compute_inbound_dedupe_key(&fingerprint, timestamp, phone, body); + let dedupe_key = compute_inbound_dedupe_key(&dedupe_namespace, timestamp, phone, body); let exists: bool = exists_statement.query_row(params![dedupe_key], |row| row.get(0))?; if !exists { conn.execute( @@ -73,19 +118,40 @@ pub(super) fn backfill_dedupe_keys_on(conn: &Connection) -> Result { Ok(count) } +fn inbound_dedupe_namespace_on(conn: &Connection) -> Result> { + conn.query_row( + "SELECT value FROM meta + WHERE key IN ( + 'modem_dedupe_namespace', + 'modem_fingerprint', + 'runtime_modem_fingerprint' + ) AND value <> '' + ORDER BY CASE key + WHEN 'modem_dedupe_namespace' THEN 0 + WHEN 'modem_fingerprint' THEN 1 + ELSE 2 + END + LIMIT 1", + [], + |row| row.get(0), + ) + .optional() + .map_err(Into::into) +} + #[cfg(test)] mod tests { use super::*; #[test] - fn backfill_without_modem_fingerprint_is_a_no_op() { + fn backfill_without_a_dedupe_namespace_is_a_no_op() { let store = MessageStore::open_in_memory().unwrap(); assert_eq!(store.backfill_dedupe_keys().unwrap(), 0); } #[test] - fn backfill_propagates_modem_fingerprint_query_errors() { + fn backfill_propagates_dedupe_namespace_query_errors() { let store = MessageStore::open_in_memory().unwrap(); { let conn = store.conn.lock().unwrap(); @@ -96,7 +162,7 @@ mod tests { assert!( error.to_string().contains("no such table: meta"), - "expected the fingerprint query error, got: {error:#}" + "expected the dedupe namespace query error, got: {error:#}" ); } }