Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions src/icloud/cloudkit.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1373,6 +1373,7 @@ pub struct CloudKitOpenContainer<'t, T: AnisetteProvider> {
container: &'t CloudKitContainer<'t>,
pub user_id: String,
pub client: Arc<CloudKitClient<T>>,
// Always lock PasswordManager state first, then these CloudKit container keys, then Keychain state.
pub keys: DebugMutex<HashMap<String, PCSZoneConfig>>,
pub database_type: cloudkit_proto::request_operation::header::Database,
}
Expand Down Expand Up @@ -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<ShareInfo>)], client: &KeychainClient<T>, pcs_service: &PCSService<'_>, sync_keychain: bool) -> Result<Vec<Result<PCSZoneConfig, PushError>>, 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();
Expand Down
52 changes: 32 additions & 20 deletions src/icloud/keychain.rs
Original file line number Diff line number Diff line change
Expand Up @@ -661,6 +661,7 @@ pub const SECURITYD_CONTAINER: CloudKitContainer = CloudKitContainer {
pub struct KeychainClient<P: AnisetteProvider> {
pub anisette: ArcAnisetteClient<P>,
pub token_provider: Arc<TokenProvider<P>>,
// 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<KeychainClientState>,
pub config: Arc<dyn OSConfig>,
pub update_state: Box<dyn Fn(&KeychainClientState) + Send + Sync>,
Expand Down Expand Up @@ -1336,13 +1337,18 @@ impl<P: AnisetteProvider> KeychainClient<P> {
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?;
Expand Down Expand Up @@ -1788,13 +1794,13 @@ impl<P: AnisetteProvider> KeychainClient<P> {
(self.update_state)(&state);
}

pub async fn fetch_shares_for(&self, user: &KeychainUserIdentity<impl KeystoreDeriveKey>) -> Result<Vec<CuttlefishSerializedKey>, 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<impl KeystoreDeriveKey>, peers: &HashMap<String, EncodedPeer>) -> Result<Vec<CuttlefishSerializedKey>, 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 {
Expand All @@ -1803,8 +1809,8 @@ impl<P: AnisetteProvider> KeychainClient<P> {
};
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::<Vec<_>>());
let Some(sending_peer) = peers.get(&item.sender) else {
warn!("missing sender {} in state! {:?}", item.sender, peers.keys().collect::<Vec<_>>());
continue
};
sending_peer.verify_signature_dig(MessageDigest::sha256(), &item.data_for_signing(), &base64_decode(&item.signature))?;
Expand Down Expand Up @@ -1911,10 +1917,12 @@ impl<P: AnisetteProvider> KeychainClient<P> {

info!("Self vouching as {} {:?}", other_identity.identifier, state.state.keys().collect::<Vec<_>>());
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)
}
Expand Down Expand Up @@ -1986,10 +1994,11 @@ impl<P: AnisetteProvider> KeychainClient<P> {
(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);
Expand Down Expand Up @@ -2111,14 +2120,12 @@ impl<P: AnisetteProvider> KeychainClient<P> {
}

async fn get_escrow_headers(&self) -> Result<HeaderMap, PushError> {
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();

Expand All @@ -2133,9 +2140,14 @@ impl<P: AnisetteProvider> KeychainClient<P> {
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)?)
Expand Down
81 changes: 67 additions & 14 deletions src/ids/identity_manager.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<KeyCache>,
pub users: DebugRwLock<Vec<IDSUser>>,
pub identity: IDSNGMIdentity,
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -454,13 +454,20 @@ impl IdentityResource {
}

pub async fn get_possible_handles(&self) -> Result<HashSet<String>, 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 {
Expand Down Expand Up @@ -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<IDSUser>, 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<String, HashMap<String, CachedHandle>>, 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) {
Expand Down Expand Up @@ -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 {
Expand All @@ -684,6 +691,7 @@ impl IdentityResource {
}

pub async fn get_sms_targets(&self, handle: &str, refresh: bool) -> Result<Vec<PrivateDeviceInfo>, 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;
Expand All @@ -697,6 +705,7 @@ impl IdentityResource {
}

pub async fn token_to_uuid(&self, handle: &str, token: &[u8]) -> Result<String, 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;
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) {
Expand Down Expand Up @@ -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 {
Expand All @@ -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;
Expand Down Expand Up @@ -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"
);
}
}
2 changes: 2 additions & 0 deletions src/passwords.rs
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,7 @@ pub struct PasswordManager<P: AnisetteProvider> {
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<PasswordState>,
update_state: Box<dyn Fn(&PasswordState) + Send + Sync>,
notif_watcher: CloudKitNotifWatcher,
Expand Down Expand Up @@ -1293,6 +1294,7 @@ impl<P: AnisetteProvider + Send + Sync + 'static> PasswordManager<P> {
}

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| {
Expand Down
32 changes: 17 additions & 15 deletions src/sharedstreams.rs
Original file line number Diff line number Diff line change
Expand Up @@ -222,9 +222,8 @@ pub struct SharedStreamClient<P: AnisetteProvider> {
}

impl<P: AnisetteProvider> SharedStreamClient<P> {
async fn get_headers(&self) -> Result<HeaderMap, PushError> {
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<HeaderMap, PushError> {
let mme_token = self.token_provider.get_mme_token("mmeAuthToken").await?;

let mut map = HeaderMap::new();
Expand All @@ -235,9 +234,7 @@ impl<P: AnisetteProvider> SharedStreamClient<P> {
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();

Expand All @@ -249,12 +246,15 @@ impl<P: AnisetteProvider> SharedStreamClient<P> {
}

pub async fn get_album<T: DeserializeOwned>(&self, album: &str, url: &str, enter: impl Serialize) -> Result<T, PushError> {
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?;
Expand All @@ -269,15 +269,17 @@ impl<P: AnisetteProvider> SharedStreamClient<P> {
}

pub async fn request_me(&self, url: &str, enter: impl Serialize) -> Result<Vec<u8>, 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();
Expand Down