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)
This commit is contained in:
Edison Jwa
2026-06-10 07:24:08 +09:00
parent 8ffd33b3cc
commit 82ea1a117f
2 changed files with 361 additions and 70 deletions
+256
View File
@@ -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<oneshot::Sender<Result<Option<Vec<u8>>, 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<Mutex<Option<ProtocolClient>>>,
semaphore: Arc<Semaphore>,
in_flight: Arc<Mutex<HashMap<String, InFlightWaiters>>>,
negative_cache: Arc<Mutex<HashMap<String, Instant>>>,
}
impl FileTransferService {
pub fn new(cache: BlobCache, protocol: Arc<Mutex<Option<ProtocolClient>>>) -> 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<ProtocolClient>) {
*self.protocol.lock().await = client;
}
pub async fn get_avatar(
&self,
avatar_hash: &str,
client_uid: &str,
) -> Result<Option<Vec<u8>>, 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<u64, FileTransferError> {
Ok(self.cache.total_size().await?)
}
async fn do_download_avatar(
&self,
avatar_hash: &str,
client_uid: &str,
) -> Result<Option<Vec<u8>>, 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<Option<Vec<u8>>, 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);
}
}
+105 -70
View File
@@ -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<chanora_audio::AudioEngine>,
/// 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<Mutex<Option<BookmarkRepository>>>,
blob_cache: Arc<Mutex<Option<chanora_cache::BlobCache>>>,
protocol: Arc<Mutex<Option<chanora_protocol::ProtocolClient>>>,
file_transfer: Arc<Mutex<Option<Arc<file_transfer::FileTransferService>>>>,
/// 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<chanora_protocol::ProtocolClient>) {
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<chanora_protocol::ProtocolClient> {
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<Option<Vec<u8>>, 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<u64, CoreError> {
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<ServerSnapshot, CoreError> {
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<ClientProfile, CoreError> {
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<Mutex<Option<ConnectedState>>>,
protocol: Arc<Mutex<Option<chanora_protocol::ProtocolClient>>>,
events_tx: broadcast::Sender<SessionEvent>,
initial_cfg: ConnectConfig,
initial_lost_rx: oneshot::Receiver<chanora_protocol::DisconnectReason>,
@@ -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,
);