From 82ea1a117fa8394e59fdbf6ba871aa308e364fd8 Mon Sep 17 00:00:00 2001 From: Edison Jwa Date: Wed, 10 Jun 2026 07:24:08 +0900 Subject: [PATCH] feat(core): FileTransferService with coalescing, throttling, negative cache - New file_transfer module with FileTransferService struct - Semaphore(2) throttles concurrent downloads - In-flight HashMap coalesces duplicate avatar requests - 5-min negative cache short-circuits ServerRejected misses - ChanoraSession delegates get_avatar through the service - connect/disconnect update shared protocol handle - clear_cache/cache_size delegate to service - 2 new unit tests (cached hit, negative cache) --- core/chanora_core/src/file_transfer.rs | 256 +++++++++++++++++++++++++ core/chanora_core/src/lib.rs | 175 ++++++++++------- 2 files changed, 361 insertions(+), 70 deletions(-) create mode 100644 core/chanora_core/src/file_transfer.rs diff --git a/core/chanora_core/src/file_transfer.rs b/core/chanora_core/src/file_transfer.rs new file mode 100644 index 0000000..28d4540 --- /dev/null +++ b/core/chanora_core/src/file_transfer.rs @@ -0,0 +1,256 @@ +use std::collections::HashMap; +use std::sync::Arc; +use std::time::{Duration, Instant}; + +use chanora_cache::{BlobCache, BlobCacheError, PREFIX_AVATAR}; +use chanora_protocol::{ProtocolClient, ProtocolError}; +use tokio::sync::{Mutex, Semaphore, oneshot}; +use tracing::warn; + +const MAX_CONCURRENT_DOWNLOADS: usize = 2; +const NEGATIVE_CACHE_TTL: Duration = Duration::from_secs(5 * 60); + +type InFlightWaiters = Vec>, FileTransferError>>>; + +/// Errors raised while resolving protocol-owned file assets. +#[derive(Debug, thiserror::Error)] +pub enum FileTransferError { + /// No live protocol client is available for a download. + #[error("not connected")] + NotConnected, + /// The protocol layer failed while downloading the asset. + #[error("protocol error: {0}")] + Protocol(#[from] ProtocolError), + /// The blob cache failed while reading or writing the asset. + #[error("cache error: {0}")] + Cache(#[from] BlobCacheError), +} + +impl Clone for FileTransferError { + fn clone(&self) -> Self { + match self { + Self::NotConnected => Self::NotConnected, + Self::Protocol(error) => Self::Protocol(clone_protocol_error(error)), + Self::Cache(error) => Self::Cache(clone_blob_cache_error(error)), + } + } +} + +pub struct FileTransferService { + cache: BlobCache, + protocol: Arc>>, + semaphore: Arc, + in_flight: Arc>>, + negative_cache: Arc>>, +} + +impl FileTransferService { + pub fn new(cache: BlobCache, protocol: Arc>>) -> Self { + Self { + cache, + protocol, + semaphore: Arc::new(Semaphore::new(MAX_CONCURRENT_DOWNLOADS)), + in_flight: Arc::new(Mutex::new(HashMap::new())), + negative_cache: Arc::new(Mutex::new(HashMap::new())), + } + } + + pub async fn set_protocol(&self, client: Option) { + *self.protocol.lock().await = client; + } + + pub async fn get_avatar( + &self, + avatar_hash: &str, + client_uid: &str, + ) -> Result>, FileTransferError> { + if let Some(bytes) = self.cache.get(PREFIX_AVATAR, avatar_hash).await? { + return Ok(Some(bytes)); + } + + if self.is_negative_cache_hit(avatar_hash).await { + return Ok(None); + } + + let rx = { + let mut in_flight = self.in_flight.lock().await; + if let Some(waiters) = in_flight.get_mut(avatar_hash) { + let (tx, rx) = oneshot::channel(); + waiters.push(tx); + Some(rx) + } else { + in_flight.insert(avatar_hash.to_string(), Vec::new()); + None + } + }; + + if let Some(rx) = rx { + return rx.await.unwrap_or_else(|_| { + Err(FileTransferError::Protocol(ProtocolError::Lost( + "coalesced avatar download waiter dropped".to_string(), + ))) + }); + } + + let _permit = self + .semaphore + .acquire() + .await + .expect("file transfer semaphore should stay open"); + + let result = self.do_download_avatar(avatar_hash, client_uid).await; + self.finish_in_flight(avatar_hash, &result).await; + result + } + + pub async fn clear_cache(&self) -> Result<(), FileTransferError> { + self.cache.clear().await?; + self.negative_cache.lock().await.clear(); + Ok(()) + } + + pub async fn cache_size(&self) -> Result { + Ok(self.cache.total_size().await?) + } + + async fn do_download_avatar( + &self, + avatar_hash: &str, + client_uid: &str, + ) -> Result>, FileTransferError> { + if let Some(bytes) = self.cache.get(PREFIX_AVATAR, avatar_hash).await? { + return Ok(Some(bytes)); + } + + if self.is_negative_cache_hit(avatar_hash).await { + return Ok(None); + } + + let protocol = self.protocol.lock().await; + let client = protocol.as_ref().ok_or(FileTransferError::NotConnected)?; + match client.download_avatar(client_uid).await { + Ok(bytes) => { + self.cache.put(PREFIX_AVATAR, avatar_hash, &bytes).await?; + self.negative_cache.lock().await.remove(avatar_hash); + Ok(Some(bytes)) + } + Err(ProtocolError::ServerRejected { .. }) => { + self.negative_cache + .lock() + .await + .insert(avatar_hash.to_string(), Instant::now() + NEGATIVE_CACHE_TTL); + Ok(None) + } + Err(error) => Err(FileTransferError::Protocol(error)), + } + } + + async fn finish_in_flight( + &self, + avatar_hash: &str, + result: &Result>, FileTransferError>, + ) { + let waiters = self.in_flight.lock().await.remove(avatar_hash).unwrap_or_default(); + for waiter in waiters { + if waiter.send(result.clone()).is_err() { + warn!(target: "chanora_core", avatar_hash, "avatar download waiter dropped"); + } + } + } + + async fn is_negative_cache_hit(&self, avatar_hash: &str) -> bool { + let now = Instant::now(); + let mut negative_cache = self.negative_cache.lock().await; + match negative_cache.get(avatar_hash).copied() { + Some(expires_at) if expires_at > now => true, + Some(_) => { + negative_cache.remove(avatar_hash); + false + } + None => false, + } + } +} + +fn clone_protocol_error(error: &ProtocolError) -> ProtocolError { + match error { + ProtocolError::Invalid(message) => ProtocolError::Invalid(message.clone()), + ProtocolError::DnsFailed { host, reason } => ProtocolError::DnsFailed { + host: host.clone(), + reason: reason.clone(), + }, + ProtocolError::Connect(message) => ProtocolError::Connect(message.clone()), + ProtocolError::DisconnectedEarly(message) => { + ProtocolError::DisconnectedEarly(message.clone()) + } + ProtocolError::Lost(message) => ProtocolError::Lost(message.clone()), + ProtocolError::Identity(message) => ProtocolError::Identity(message.clone()), + ProtocolError::Timeout => ProtocolError::Timeout, + ProtocolError::ServerRejected { code, message } => ProtocolError::ServerRejected { + code: *code, + message: message.clone(), + }, + ProtocolError::Backend(message) => ProtocolError::Backend(message.clone()), + ProtocolError::FileTransfer(message) => ProtocolError::FileTransfer(message.clone()), + } +} + +fn clone_blob_cache_error(error: &BlobCacheError) -> BlobCacheError { + match error { + BlobCacheError::Io(message) => BlobCacheError::Io(message.clone()), + BlobCacheError::InvalidKey(message) => BlobCacheError::InvalidKey(message.clone()), + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn test_cache_dir(name: &str) -> std::path::PathBuf { + let mut path = std::env::temp_dir(); + let nanos = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap() + .as_nanos(); + path.push(format!("chanora-core-file-transfer-{name}-{nanos}")); + path + } + + #[tokio::test] + async fn returns_cached_avatar_without_connection() { + let cache_dir = test_cache_dir("cache-hit"); + let cache = BlobCache::new(&cache_dir, 1024).unwrap(); + cache.put(PREFIX_AVATAR, "a1b2c3d4e5f6a7b8c9d0e1f2a3b4c5d6", b"avatar") + .await + .unwrap(); + let service = FileTransferService::new(cache, Arc::new(Mutex::new(None))); + + let avatar = service + .get_avatar("a1b2c3d4e5f6a7b8c9d0e1f2a3b4c5d6", "client") + .await + .unwrap(); + + assert_eq!(avatar, Some(b"avatar".to_vec())); + let _ = std::fs::remove_dir_all(cache_dir); + } + + #[tokio::test] + async fn negative_cache_short_circuits_not_connected() { + let cache_dir = test_cache_dir("negative-cache"); + let cache = BlobCache::new(&cache_dir, 1024).unwrap(); + let service = FileTransferService::new(cache, Arc::new(Mutex::new(None))); + service + .negative_cache + .lock() + .await + .insert("a1b2c3d4e5f6a7b8c9d0e1f2a3b4c5d6".to_string(), Instant::now() + NEGATIVE_CACHE_TTL); + + let avatar = service + .get_avatar("a1b2c3d4e5f6a7b8c9d0e1f2a3b4c5d6", "client") + .await + .unwrap(); + + assert_eq!(avatar, None); + let _ = std::fs::remove_dir_all(cache_dir); + } +} diff --git a/core/chanora_core/src/lib.rs b/core/chanora_core/src/lib.rs index 374d213..9b4120c 100644 --- a/core/chanora_core/src/lib.rs +++ b/core/chanora_core/src/lib.rs @@ -53,6 +53,7 @@ use chanora_state::channel_join::{ }; mod events; +mod file_transfer; mod network_diagnostics; pub mod ptt; @@ -75,6 +76,7 @@ pub use events::{ NetworkState, PersistedPttBinding, PttDescriptorSnapshot, SessionEvent, VoiceJoinErrorCode, VoiceJoinSyncState, }; +pub use file_transfer::FileTransferError; use network_diagnostics::NetworkDiagnostics; /// Errors that can arise during top-level orchestration. @@ -95,6 +97,9 @@ pub enum CoreError { /// Blob-cache failure. #[error("cache: {0}")] Cache(#[from] chanora_cache::BlobCacheError), + /// File-transfer failure. + #[error("file transfer: {0}")] + FileTransfer(#[from] FileTransferError), /// Diagnostics error. #[error("diagnostics: {0}")] Diagnostics(#[from] chanora_diagnostics::DiagnosticsError), @@ -133,7 +138,6 @@ struct SupervisorInner { } struct ConnectedState { - protocol: chanora_protocol::ProtocolClient, audio: Option, /// Active PTT controller (SDD-088). Owns the platform input /// backend, the active binding, and the capability watch @@ -200,7 +204,8 @@ pub struct ChanoraSession { /// extension). Lives alongside the identity file. Wired by /// [`Self::init_storage`]. bookmark_store: Arc>>, - blob_cache: Arc>>, + protocol: Arc>>, + file_transfer: Arc>>>, /// Invisible server-address prefetch cache. Warmed by Flutter typing /// but validated by Rust before Connect can reuse it. server_prefetch: ServerPrefetcher, @@ -256,7 +261,8 @@ impl ChanoraSession { network_tx, identity_store: Arc::new(Mutex::new(None)), bookmark_store: Arc::new(Mutex::new(None)), - blob_cache: Arc::new(Mutex::new(None)), + protocol: Arc::new(Mutex::new(None)), + file_transfer: Arc::new(Mutex::new(None)), server_prefetch: ServerPrefetcher::new(), voice_selector: selector, release_tail, @@ -278,6 +284,19 @@ impl ChanoraSession { ConnectionEpoch(epoch) } + async fn store_protocol(&self, client: Option) { + let service = { self.file_transfer.lock().await.clone() }; + if let Some(service) = service { + service.set_protocol(client).await; + } else { + *self.protocol.lock().await = client; + } + } + + async fn take_protocol(&self) -> Option { + self.protocol.lock().await.take() + } + /// Wire a directory-backed identity store. Called by the bridge /// during `bridge_init` once Flutter has resolved the platform /// app-private storage directory. Subsequent [`Self::connect`] @@ -346,8 +365,12 @@ impl ChanoraSession { pub async fn init_cache(&self, dir: &str) -> Result<(), CoreError> { let cache = chanora_cache::BlobCache::new(dir, 100 * 1024 * 1024)?; cache.evict().await?; - let mut guard = self.blob_cache.lock().await; - *guard = Some(cache); + let service = Arc::new(file_transfer::FileTransferService::new( + cache, + self.protocol.clone(), + )); + let mut guard = self.file_transfer.lock().await; + *guard = Some(service); Ok(()) } @@ -357,47 +380,31 @@ impl ChanoraSession { avatar_hash: &str, client_uid: &str, ) -> Result>, CoreError> { - { - let guard = self.blob_cache.lock().await; - if let Some(ref cache) = *guard { - if let Some(bytes) = cache.get(chanora_cache::PREFIX_AVATAR, avatar_hash).await? { - return Ok(Some(bytes)); - } - } + let service = { self.file_transfer.lock().await.clone() }; + if let Some(service) = service { + return Ok(service.get_avatar(avatar_hash, client_uid).await?); } - let guard = self.inner.lock().await; - let state = guard.as_ref().ok_or(CoreError::NotConnected)?; - let bytes = state.protocol.download_avatar(client_uid).await?; - - { - let cache_guard = self.blob_cache.lock().await; - if let Some(ref cache) = *cache_guard { - let _ = cache - .put(chanora_cache::PREFIX_AVATAR, avatar_hash, &bytes) - .await; - } - } - - Ok(Some(bytes)) + let protocol = self.protocol.lock().await; + let client = protocol.as_ref().ok_or(CoreError::NotConnected)?; + Ok(Some(client.download_avatar(client_uid).await?)) } /// Purge cached protocol-owned assets. pub async fn clear_cache(&self) -> Result<(), CoreError> { - let guard = self.blob_cache.lock().await; - if let Some(ref cache) = *guard { - cache.clear().await?; + let service = { self.file_transfer.lock().await.clone() }; + if let Some(service) = service { + service.clear_cache().await?; } Ok(()) } /// Report the configured blob-cache size. pub async fn cache_size(&self) -> Result { - let guard = self.blob_cache.lock().await; - if let Some(ref cache) = *guard { - Ok(cache.total_size().await?) - } else { - Ok(0) + let service = { self.file_transfer.lock().await.clone() }; + match service { + Some(service) => Ok(service.cache_size().await?), + None => Ok(0), } } @@ -578,6 +585,7 @@ impl ChanoraSession { let supervisor = tokio::spawn(supervisor_loop(SupervisorContext { state_arc: self.inner.clone(), + protocol: self.protocol.clone(), events_tx: self.events_tx.clone(), initial_cfg: cfg.clone(), initial_lost_rx: lost_rx, @@ -632,9 +640,9 @@ impl ChanoraSession { } spawn_event_forwarders(&client, &self.events_tx); + self.store_protocol(Some(client)).await; *guard = Some(ConnectedState { - protocol: client, audio: None, ptt_controller: None, cancel_tx: Some(cancel_tx), @@ -695,9 +703,13 @@ impl ChanoraSession { /// Return a fresh snapshot of the current server state. pub async fn snapshot(&self) -> Result { + let snap = { + let protocol = self.protocol.lock().await; + let client = protocol.as_ref().ok_or(CoreError::NotConnected)?; + client.snapshot().await? + }; let mut guard = self.inner.lock().await; let state = guard.as_mut().ok_or(CoreError::NotConnected)?; - let snap = state.protocol.snapshot().await?; let current_channel = self .find_own_in(&snap) .await @@ -725,9 +737,9 @@ impl ChanoraSession { /// Fetch richer profile and live connection details for one online client. pub async fn client_profile(&self, client_id: u64) -> Result { - let guard = self.inner.lock().await; - let state = guard.as_ref().ok_or(CoreError::NotConnected)?; - Ok(state.protocol.client_profile(client_id).await?) + let protocol = self.protocol.lock().await; + let client = protocol.as_ref().ok_or(CoreError::NotConnected)?; + Ok(client.client_profile(client_id).await?) } /// True if a connection is currently active. @@ -744,9 +756,9 @@ impl ChanoraSession { if !should_dispatch_text_message(&message, &target) { return Ok(()); } - let guard = self.inner.lock().await; - let state = guard.as_ref().ok_or(CoreError::NotConnected)?; - state.protocol.send_text_message(message, target).await?; + let protocol = self.protocol.lock().await; + let client = protocol.as_ref().ok_or(CoreError::NotConnected)?; + client.send_text_message(message, target).await?; Ok(()) } @@ -805,11 +817,16 @@ impl ChanoraSession { // session permanently unable to restart audio without a // reconnect (the user saw "voice_in already taken" on the // second channel switch). - let voice_out = state.protocol.voice_out(); - let voice_in = state - .protocol - .take_voice_in() - .ok_or(CoreError::Invariant("voice_in already taken"))?; + let (voice_out, voice_in) = { + let protocol = self.protocol.lock().await; + let client = protocol.as_ref().ok_or(CoreError::NotConnected)?; + ( + client.voice_out(), + client + .take_voice_in() + .ok_or(CoreError::Invariant("voice_in already taken"))?, + ) + }; let gate = AudioTransmitGate::new(cfg.ptt_initial); cfg.voice_activity_selector = Some(self.voice_selector.clone()); let new_engine = match chanora_audio::AudioEngine::start_with_gate( @@ -1040,10 +1057,11 @@ impl ChanoraSession { let password_to_send = requested_password .clone() .or_else(|| state.channel_passwords.get(&channel_id).cloned()); - state - .protocol - .move_to_channel(channel_id, password_to_send) - .await?; + { + let protocol = self.protocol.lock().await; + let client = protocol.as_ref().ok_or(CoreError::NotConnected)?; + client.move_to_channel(channel_id, password_to_send).await?; + } if let Some(pw) = requested_password { state.channel_passwords.insert(channel_id, pw); } @@ -1068,7 +1086,11 @@ impl ChanoraSession { if let Some(muted) = output { state.local_output_muted = muted; } - state.protocol.set_muted(input, output).await?; + { + let protocol = self.protocol.lock().await; + let client = protocol.as_ref().ok_or(CoreError::NotConnected)?; + client.set_muted(input, output).await?; + } if let Some(muted) = output { if let Some(audio) = state.audio.as_ref() { audio.set_output_muted(muted); @@ -1391,11 +1413,12 @@ impl ChanoraSession { let password_to_send = requested_password .clone() .or_else(|| state.channel_passwords.get(&channel_id).cloned()); - if let Err(e) = state - .protocol - .queue_move_to_channel(channel_id, password_to_send) - .await - { + let move_result = { + let protocol = self.protocol.lock().await; + let client = protocol.as_ref().ok_or(CoreError::NotConnected)?; + client.queue_move_to_channel(channel_id, password_to_send).await + }; + if let Err(e) = move_result { // TS3 error 0x0302 = `channel_already_in`: we're already // in the target channel, so this is a no-op success. // Rolling `in_channel` back to false would break PTT @@ -1657,7 +1680,9 @@ impl ChanoraSession { audio.stop(); let _ = self.events_tx.send(SessionEvent::AudioStopped); } - state.protocol.disconnect().await; + if let Some(protocol) = self.take_protocol().await { + protocol.disconnect().await; + } // Wait for the supervisor to wind down so we don't race // a redial against the explicit disconnect. if let Some(handle) = state.supervisor.take() { @@ -1719,6 +1744,7 @@ async fn await_supervisor_shutdown(mut handle: JoinHandle<()>, timeout_duration: struct SupervisorContext { state_arc: Arc>>, + protocol: Arc>>, events_tx: broadcast::Sender, initial_cfg: ConnectConfig, initial_lost_rx: oneshot::Receiver, @@ -1853,6 +1879,7 @@ fn spawn_event_forwarders( async fn supervisor_loop(ctx: SupervisorContext) { let SupervisorContext { state_arc, + protocol, events_tx, initial_cfg, initial_lost_rx, @@ -2122,20 +2149,21 @@ async fn supervisor_loop(ctx: SupervisorContext) { // Reattach into the session state. let restart_audio = { + let guard = state_arc.lock().await; + if guard.is_none() { + // Session was disposed mid-reconnect. + return; + } + drop(guard); + let old = protocol.lock().await.replace(new_client); + drop(old); let mut guard = state_arc.lock().await; let state = match guard.as_mut() { Some(s) => s, None => { - // Session was disposed mid-reconnect. return; } }; - // Replace the dead protocol client with the new one. - // The old client's background task either already - // exited (loss notifier fired) or will exit when - // its request channel drops (watchdog path). - let old = std::mem::replace(&mut state.protocol, new_client); - drop(old); let _ = channel_join::reduce( &mut state.join_state, @@ -2165,9 +2193,9 @@ async fn supervisor_loop(ctx: SupervisorContext) { }); { - let guard = state_arc.lock().await; - if let Some(state) = guard.as_ref() { - spawn_event_forwarders(&state.protocol, &events_tx); + let protocol = protocol.lock().await; + if let Some(client) = protocol.as_ref() { + spawn_event_forwarders(client, &events_tx); } } @@ -2179,8 +2207,15 @@ async fn supervisor_loop(ctx: SupervisorContext) { }; let mut guard = state_arc.lock().await; if let Some(state) = guard.as_mut() { - let voice_out = state.protocol.voice_out(); - if let Some(voice_in) = state.protocol.take_voice_in() { + let (voice_out, voice_in) = { + let protocol = protocol.lock().await; + let client = match protocol.as_ref() { + Some(client) => client, + None => return, + }; + (client.voice_out(), client.take_voice_in()) + }; + if let Some(voice_in) = voice_in { let gate = chanora_audio::AudioTransmitGate::new( audio_cfg.ptt_initial, );