From 1e677bb9f00e70a41964934906ffef9e838d31ee Mon Sep 17 00:00:00 2001 From: mtalexan Date: Mon, 5 Oct 2026 18:59:28 -0600 Subject: [PATCH] Fix mutex locking order consistency and deadlocks. Define a consistent mutex locking order across all the mutexes that are locked together so we can't end up in mutex inversion. Fix some deadlocks with cross-thread tokenio lock queueing. Reduce locking scopes to only the code that needs the lock. Long async actions don't hold locks while watiing to run, they either get a cloned copy or re-acquire the lock. --- src/icloud/cloudkit.rs | 2 + src/icloud/keychain.rs | 52 +++++++++++++++--------- src/ids/identity_manager.rs | 81 ++++++++++++++++++++++++++++++------- src/passwords.rs | 2 + src/sharedstreams.rs | 32 ++++++++------- 5 files changed, 120 insertions(+), 49 deletions(-) diff --git a/src/icloud/cloudkit.rs b/src/icloud/cloudkit.rs index ece3518..3654c2f 100644 --- a/src/icloud/cloudkit.rs +++ b/src/icloud/cloudkit.rs @@ -1373,6 +1373,7 @@ pub struct CloudKitOpenContainer<'t, T: AnisetteProvider> { container: &'t CloudKitContainer<'t>, pub user_id: String, pub client: Arc>, + // Always lock PasswordManager state first, then these CloudKit container keys, then Keychain state. pub keys: DebugMutex>, pub database_type: cloudkit_proto::request_operation::header::Database, } @@ -1438,6 +1439,7 @@ impl<'t, T: AnisetteProvider> CloudKitOpenContainer<'t, T> { } pub async fn get_zone_encryption_config_sev(&self, zone_ids: &[(cloudkit_proto::RecordZoneIdentifier, Option)], client: &KeychainClient, pcs_service: &PCSService<'_>, sync_keychain: bool) -> Result>, PushError> { + // This holds the container keys lock while sync_keychain may take the Keychain state lock. The PasswordManager state lock, when held by the caller, comes before both. let mut cached_keys = self.keys.lock().await; let mut get_needed = zone_ids.iter().filter(|(zone_id, share)| { let zone_name = zone_id.value.as_ref().unwrap().name().to_string(); diff --git a/src/icloud/keychain.rs b/src/icloud/keychain.rs index c8afb85..f41b1ed 100644 --- a/src/icloud/keychain.rs +++ b/src/icloud/keychain.rs @@ -661,6 +661,7 @@ pub const SECURITYD_CONTAINER: CloudKitContainer = CloudKitContainer { pub struct KeychainClient { pub anisette: ArcAnisetteClient

, pub token_provider: Arc>, + // Always lock PasswordManager state first, then the CloudKit container keys, then this Keychain state. Do not take either of the earlier locks while this one is held. pub state: DebugRwLock, pub config: Arc, pub update_state: Box, @@ -1336,13 +1337,18 @@ impl KeychainClient

{ return Err(PushError::NotInClique) } - let state = self.state.read().await; - if state.keystore.0.is_empty() { - let shares = self.fetch_shares_for(state.user_identity.as_ref().unwrap()).await?; - drop(state); + // Clone the identity and peer map so this read lock is released before fetch_shares_for waits on CloudKit. + let fetch = { + let state = self.state.read().await; + if state.keystore.0.is_empty() { + Some((state.user_identity.as_ref().unwrap().clone(), state.state.clone())) + } else { + None + } + }; + if let Some((identity, peers)) = fetch { + let shares = self.fetch_shares_for(&identity, &peers).await?; self.store_keys(&shares).await?; - } else { - drop(state); } let security_container = self.get_security_container().await?; @@ -1788,13 +1794,13 @@ impl KeychainClient

{ (self.update_state)(&state); } - pub async fn fetch_shares_for(&self, user: &KeychainUserIdentity) -> Result, PushError> { + // The identity and peer map are supplied by the caller. This function waits on CloudKit, so it must not take the state lock. + pub async fn fetch_shares_for(&self, user: &KeychainUserIdentity, peers: &HashMap) -> Result, PushError> { let response: CuttlefishFetchRecoverableTlkSharesResponse = self.invoke_cuttlefish("fetchRecoverableTLKShares", CuttlefishFetchRecoverableTlkSharesRequest { for_peer: Some(user.identifier.clone()), }).await?; let mut keys = vec![]; - let state = self.state.read().await; for share in response.shares { info!("Entering on key {}", share.service()); let Some(share_record) = &share.share else { @@ -1803,8 +1809,8 @@ impl KeychainClient

{ }; let item = CuttlefishTlkShare::from_record(&share_record.inner.as_ref().unwrap().record_field); - let Some(sending_peer) = state.state.get(&item.sender) else { - warn!("missing sender {} in state! {:?}", item.sender, state.state.keys().collect::>()); + let Some(sending_peer) = peers.get(&item.sender) else { + warn!("missing sender {} in state! {:?}", item.sender, peers.keys().collect::>()); continue }; sending_peer.verify_signature_dig(MessageDigest::sha256(), &item.data_for_signing(), &base64_decode(&item.signature))?; @@ -1911,10 +1917,12 @@ impl KeychainClient

{ info!("Self vouching as {} {:?}", other_identity.identifier, state.state.keys().collect::>()); let voucher = other_identity.vouch_for(my_identity.identifier.clone())?; + // Clone the peer map so this read lock is released before fetch_shares_for waits on CloudKit. + let peers = state.state.clone(); drop(state); - let shares = self.fetch_shares_for(&other_identity).await?; + let shares = self.fetch_shares_for(&other_identity, &peers).await?; if shares.is_empty() { return Err(PushError::PeerNoShares) } @@ -1986,10 +1994,11 @@ impl KeychainClient

{ (self.update_state)(&state); if with_tlk_shares.is_empty() { - let state = state.downgrade(); - // fetch tlk shares - let shares = self.fetch_shares_for(state.user_identity.as_ref().unwrap()).await?; + // Clone the identity and peer map, then drop this write lock. Downgrading it to a read lock would keep the lock held while fetch_shares_for waits on CloudKit. + let identity = state.user_identity.as_ref().unwrap().clone(); + let peers = state.state.clone(); drop(state); + let shares = self.fetch_shares_for(&identity, &peers).await?; self.store_keys(&shares).await?; } else { drop(state); @@ -2111,14 +2120,12 @@ impl KeychainClient

{ } async fn get_escrow_headers(&self) -> Result { - let state_lock = self.state.read().await; let mut map = HeaderMap::new(); map.insert("User-Agent", self.config.get_normal_ua("com.apple.sbd/638.100.48").parse().unwrap()); map.insert("Accept-Language", "en-US,en;q=0.9".parse().unwrap()); map.insert("x-apple-i-device-type", "1".parse().unwrap()); map.insert("Accept", "*/*".parse().unwrap()); - map.insert("X-Apple-I-Locale", "en_US".parse().unwrap()); - drop(state_lock); + map.insert("X-Apple-I-Locale", "en_US".parse().unwrap()); let mut base_headers = self.anisette.lock().await.get_headers().await?.clone(); @@ -2133,9 +2140,14 @@ impl KeychainClient

{ let auth = self.token_provider.get_gsa_token("com.apple.gs.idms.pet").await.ok_or(PushError::TokenMissing)?; let email = self.token_provider.get_gsa_email().await.expect("no email!"); - let state = self.state.read().await; - let resp = REQWEST.post(format!("{}/escrowproxy/api/{}", state.host, request.command.get_url())) - .headers(self.get_escrow_headers().await?) + // Clone the host so this read lock is released before the header lookup and HTTP request. + let host = { + let state = self.state.read().await; + state.host.clone() + }; + let headers = self.get_escrow_headers().await?; + let resp = REQWEST.post(format!("{}/escrowproxy/api/{}", host, request.command.get_url())) + .headers(headers) .header("Content-Type", "application/x-apple-plst") .basic_auth(&email, Some(&auth)) .body(plist_to_string(&request)?) diff --git a/src/ids/identity_manager.rs b/src/ids/identity_manager.rs index 02392e1..ece4b93 100644 --- a/src/ids/identity_manager.rs +++ b/src/ids/identity_manager.rs @@ -346,6 +346,7 @@ impl KeyCache { } pub struct IdentityResource { + // Always lock the cache first, then the user list, then APS state. Do not take the cache lock while either of the other two is held. pub cache: DebugMutex, pub users: DebugRwLock>, pub identity: IDSNGMIdentity, @@ -383,7 +384,6 @@ impl Resource for IdentityResource { // drop, not downgrade, to process any readers holding cache lock right now drop(users_lock); - let mut cache_lock = self.cache.lock().await; cache_lock.verity(&self.aps, &self.users.read().await, self.services).await; drop(cache_lock); @@ -454,13 +454,20 @@ impl IdentityResource { } pub async fn get_possible_handles(&self) -> Result, PushError> { - let users_locked = self.users.read().await; - let state = self.aps.state.read().await; - let mut possible_handles = HashSet::new(); - for user in &*users_locked { - let data = user.get_handle_data(&*state).await?; + // The cache lock must be taken before the user list lock, so finish this fetch and release both locks before updating the cache. + let fetched = { + let users_locked = self.users.read().await; + let state = self.aps.state.read().await; + let mut fetched = Vec::new(); + for user in &*users_locked { + fetched.push(user.get_handle_data(&*state).await?); + } + fetched + }; - let mut cache_lock = self.cache.lock().await; + let mut possible_handles = HashSet::new(); + let mut cache_lock = self.cache.lock().await; + for data in fetched { for handle in data { for (alias, attributes) in handle.aliases { for (service, _) in attributes.allowed_services { @@ -530,12 +537,11 @@ impl IdentityResource { users.iter().find(|user| user.registration["com.apple.madrid"].handles.contains(&handle.to_string())).ok_or(PushError::HandleNotFound(handle.to_string())) } - pub async fn user_by_handle<'t>(&self, service: &str, users: &'t Vec, mut handle: &str) -> Result<&'t IDSUser, PushError> { - let cache_lock = self.cache.lock().await; - if let Some(real) = cache_lock.cache.get(service).and_then(|service| service.get(handle)).and_then(|s| s.real_handle.as_ref()) { - handle = real.as_str(); - } - Self::user_by_real_handle(users, handle) + fn resolve_real_handle(cache: &HashMap>, topic: &str, handle: &str) -> String { + cache.get(topic) + .and_then(|service| service.get(handle)) + .and_then(|cached| cached.real_handle.clone()) + .unwrap_or_else(|| handle.to_string()) } pub async fn register_pseudonym(&self, services: &[&str], handle: &str, pseud: &str, exp: f64) { @@ -660,6 +666,7 @@ impl IdentityResource { Ok(response) } + // The caller already holds the cache lock, so the user list and APS state locks taken here come after it. pub async fn ensure_private_self(&self, cache_lock: &mut KeyCache, handle: &str, refresh: bool) -> Result<(), PushError> { let my_cache = cache_lock.cache.get_mut("com.apple.madrid").unwrap().get_mut(handle).unwrap(); if my_cache.private_data.len() != 0 && !refresh { @@ -684,6 +691,7 @@ impl IdentityResource { } pub async fn get_sms_targets(&self, handle: &str, refresh: bool) -> Result, PushError> { + // Takes the cache lock first. ensure_private_self then takes the user list and APS state locks while it is held. let mut cache_lock = self.cache.lock().await; self.ensure_private_self(&mut cache_lock, handle, refresh).await?; let private_self = &cache_lock.cache["com.apple.madrid"].get(handle).unwrap().private_data; @@ -697,6 +705,7 @@ impl IdentityResource { } pub async fn token_to_uuid(&self, handle: &str, token: &[u8]) -> Result { + // Takes the cache lock first. ensure_private_self then takes the user list and APS state locks while it is held. let mut cache_lock = self.cache.lock().await; let private_self = &cache_lock.cache["com.apple.madrid"].get(handle).unwrap().private_data; if let Some(found) = private_self.iter().find(|i| i.token == token) { @@ -726,8 +735,13 @@ impl IdentityResource { drop(key_cache); for chunk in fetch.chunks(18) { debug!("Fetching keys for chunk {:?}", chunk); + // Copy the real handle out so the cache lock is released before the user list lock is taken. + let real_handle = { + let cache_lock = self.cache.lock().await; + Self::resolve_real_handle(&cache_lock.cache, topic, handle) + }; let users = self.users.read().await; - let results = match self.user_by_handle(topic, &users, handle).await?.query(&*self.config, &self.aps, topic, self.get_main_service(topic), handle, chunk, meta).await { + let results = match Self::user_by_real_handle(&users, &real_handle)?.query(&*self.config, &self.aps, topic, self.get_main_service(topic), handle, chunk, meta).await { Ok(results) => results, Err(err) => { if let PushError::LookupFailed(IDSError(6005)) = err { @@ -739,6 +753,8 @@ impl IdentityResource { return Err(err) } }; + drop(users); + // Release the user list lock before taking the cache lock again. The cache lock has to come first. debug!("Got keys for {:?}", chunk); let mut key_cache = self.cache.lock().await; @@ -1258,3 +1274,40 @@ impl InnerSendJob { Ok(()) } } + +#[cfg(test)] +mod tests { + use super::*; + + fn cached_handle(real_handle: Option<&str>) -> CachedHandle { + CachedHandle { + keys: HashMap::new(), + env_hash: [0; 20], + private_data: Vec::new(), + real_handle: real_handle.map(str::to_string), + expiry: None, + } + } + + #[test] + fn resolve_real_handle_follows_pseudonym() { + let mut service = HashMap::new(); + service.insert("alias".to_string(), cached_handle(Some("real@example.com"))); + let mut cache = HashMap::new(); + cache.insert("com.apple.madrid".to_string(), service); + + assert_eq!( + IdentityResource::resolve_real_handle(&cache, "com.apple.madrid", "alias"), + "real@example.com" + ); + } + + #[test] + fn resolve_real_handle_keeps_unmapped_handle() { + let cache = HashMap::new(); + assert_eq!( + IdentityResource::resolve_real_handle(&cache, "com.apple.madrid", "tel:+15555550100"), + "tel:+15555550100" + ); + } +} diff --git a/src/passwords.rs b/src/passwords.rs index aeb62ed..6b328a8 100644 --- a/src/passwords.rs +++ b/src/passwords.rs @@ -55,6 +55,7 @@ pub struct PasswordManager { pub conn: APSConnection, pub identity: IdentityManager, _interest_token: APSInterestToken, + // Always lock this PasswordManager state first, then the CloudKit container keys, then Keychain state. pub state: DebugRwLock, update_state: Box, notif_watcher: CloudKitNotifWatcher, @@ -1293,6 +1294,7 @@ impl PasswordManager

{ } async fn sync_zones(&self, container: &CloudKitOpenContainer<'_, P>, zones_to_fetch: &[RecordZoneIdentifier]) -> Result<(), PushError> { + // This holds PasswordManager state while the CloudKit work may take the container keys lock and then the Keychain state lock. let mut groups = self.state.write().await; let zone_records = FetchRecordChangesOperation::do_sync(&container, &zones_to_fetch.iter().map(|identifier| { diff --git a/src/sharedstreams.rs b/src/sharedstreams.rs index 7f4f4bd..ebe4044 100644 --- a/src/sharedstreams.rs +++ b/src/sharedstreams.rs @@ -222,9 +222,8 @@ pub struct SharedStreamClient { } impl SharedStreamClient

{ - async fn get_headers(&self) -> Result { - let state_lock = self.state.read().await; - + // dsid comes from the caller. This function waits on the network, so it must not take the state lock. + async fn get_headers(&self, dsid: &str) -> Result { let mme_token = self.token_provider.get_mme_token("mmeAuthToken").await?; let mut map = HeaderMap::new(); @@ -235,9 +234,7 @@ impl SharedStreamClient

{ map.insert("x-apple-i-device-type", "1".parse().unwrap()); map.insert("Accept", "*/*".parse().unwrap()); map.insert("X-Apple-I-Locale", "en_US".parse().unwrap()); - map.insert("Authorization", format!("X-MobileMe-AuthToken {}", base64_encode(format!("{}:{}", state_lock.dsid, mme_token).as_bytes())).parse().unwrap()); - - drop(state_lock); + map.insert("Authorization", format!("X-MobileMe-AuthToken {}", base64_encode(format!("{}:{}", dsid, mme_token).as_bytes())).parse().unwrap()); let mut base_headers = self.anisette.lock().await.get_headers().await?.clone(); @@ -249,12 +246,15 @@ impl SharedStreamClient

{ } pub async fn get_album(&self, album: &str, url: &str, enter: impl Serialize) -> Result { - let state = self.state.read().await; - let location = state.albums.iter().find(|a| a.albumguid == album).unwrap().albumlocation.clone().expect("Not a confirmed location?"); - drop(state); + // Copy the location and dsid so this read lock is released before get_headers and the HTTP request. + let (location, dsid) = { + let state = self.state.read().await; + let location = state.albums.iter().find(|a| a.albumguid == album).unwrap().albumlocation.clone().expect("Not a confirmed location?"); + (location, state.dsid.clone()) + }; let resp = REQWEST.post(format!("{}{}", location, url)) - .headers(self.get_headers().await?) + .headers(self.get_headers(&dsid).await?) .header("Content-Type", "text/plist") .body(plist_to_string(&enter)?) .send().await?; @@ -269,15 +269,17 @@ impl SharedStreamClient

{ } pub async fn request_me(&self, url: &str, enter: impl Serialize) -> Result, PushError> { - let state_lock = self.state.read().await; - let resp = REQWEST.post(format!("{}/{}/sharedstreams/{}", state_lock.host, state_lock.dsid, url)) - .headers(self.get_headers().await?) + // Clone the host and dsid so this read lock is released before get_headers and the HTTP request. + let (host, dsid) = { + let state_lock = self.state.read().await; + (state_lock.host.clone(), state_lock.dsid.clone()) + }; + let resp = REQWEST.post(format!("{}/{}/sharedstreams/{}", host, dsid, url)) + .headers(self.get_headers(&dsid).await?) .header("Content-Type", "text/plist") .body(plist_to_string(&enter)?) .send().await?; - drop(state_lock); - let mut state_lock = self.state.write().await; if let Some(host) = resp.headers().get("X-Apple-MME-Host") { state_lock.host = host.to_str().unwrap().to_string();