Files
chanora/docs/architecture/file-transfer-implementation-plan.md
T
Edison Jwa aa796d7395 feat: file transfer system (avatar/icon download with cacache) (#40)
* docs(architecture): add file transfer design, research, and implementation plan

* feat(cache): add chanora_cache crate with cacache-backed blob cache

- New chanora_cache crate: content-addressed blob store wrapping cacache
- BlobCache API: async put/get/remove/clear/total_size/evict
- Key validation: av_ prefix (32 hex chars), ic_ prefix (decimal digits)
- Cacache provides crash safety, SSRI integrity, content dedup
- Mtime-based eviction via cacache::list_sync + sort by timestamp
- 7 unit tests all passing
- Added to workspace members

* feat(protocol): add file download support for avatars and icons

- Add Request::DownloadFile variant with oneshot reply
- Add ProtocolClient::download_avatar(client_uid) and download_icon(icon_id)
- Track pending file downloads by FiletransferHandle
- Handle StreamItem::FileDownload: read bytes from TCP stream
- Handle StreamItem::FiletransferFailed: map to ProtocolError
- Add ProtocolError::FileTransfer(String) variant
- Add path helper tests for avatar/icon download paths
- No tsclientlib types leak across the adapter boundary

* feat(core): add blob cache wiring and avatar download orchestration

- Add chanora_cache dependency to Cargo.toml
- Add blob_cache field to ChanoraSession (Arc<Mutex<Option<BlobCache>>>)
- Add init_cache() method: creates BlobCache, runs eviction
- Add get_avatar() method: cache-first, download on miss, store in cache
- Add clear_cache() and cache_size() methods for cache management
- Add CoreError::Cache variant for BlobCacheError conversion
- Add avatar_cache integration test

* feat(bridge): add init_cache, download_avatar, and cache management functions

- Add init_cache(dir) bridge function
- Add download_avatar(avatar_hash, client_uid) bridge function
- Add clear_file_cache() and file_cache_size() bridge functions
- Map CoreError::Cache and ProtocolError::FileTransfer in BridgeError

* feat(flutter): add cache initialization wiring and avatar download shims

- Add wireCache() to app_bootstrap using getApplicationCacheDirectory()
- Call wireCache() after wireStorage() in main bootstrap flow
- Add Dart-side initCache and downloadAvatar wrapper shims in api.dart
- Update Cargo.lock for new chanora_cache dependency

* 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)

* feat(core,bridge): add get_icon with coalescing and negative cache

- FileTransferService::get_icon() mirrors get_avatar pattern
- ChanoraSession::get_icon() delegates through FileTransferService
- Bridge download_icon() exposed for Flutter
- Dart downloadIcon() shim added
- Uses PREFIX_ICON (ic_<crc32u>) cache key format
- 1 new unit test (cached icon hit)

* fix(core,protocol): simplify store_protocol and add download size cap

- store_protocol: always write to shared Arc<Mutex<Option<ProtocolClient>>>;
  the FileTransferService holds the same Arc so it sees updates automatically
- read_download_bytes: reject downloads exceeding 10 MB to prevent
  malicious servers from causing OOM
2026-06-10 11:45:14 +09:00

46 KiB
Raw Blame History

File Transfer Implementation Plan

Status: Draft — revised after cache research (2026-06-10) Created: 2026-06-10 Revised: 2026-06-10 — cache architecture overhauled per crate comparison, cross-platform dir research, and Oracle consultation Upstream: docs/architecture/file-transfer-design.md, docs/architecture/file-transfer-research.md

0. Overview

This document decomposes the file transfer system (avatar/icon download + caching) into five sequential phases with concrete file-level changes, function signatures, and verification criteria.

Scope: Download only. No upload, no file browser.

Architecture recap: Content-addressed blob store backed by cacache. Protocol keys (av_<md5>, ic_<crc32>) map to cacache string keys. Physical storage is SHA-512 content-addressed. Same content on any server = stored once = zero duplication.


1. Phase Summary

Phase Crate What Delivers Est. LOC
1 chanora_protocol Request::DownloadFile, StreamItem::FileDownload handling, ProtocolClient::download_avatar() / download_icon() Raw download across the protocol boundary ~180
2 chanora_bridge + Flutter Bridge downloadAvatar(), Flutter Image.memory() + initials fallback Avatars visible in the UI (no caching — every view re-downloads) ~100
3 chanora_cache (NEW crate) BlobCache — flat content-addressed blob store, no metadata files, eviction by file mtime Disk cache, zero duplication, no deps ~150
4 chanora_core FileTransferService — request coalescing, rate limiter (1-2 concurrent), negative cache, retry backoff Robust production download pipeline ~250
5 chanora_bridge + Flutter clearAvatarCache(), clearIconCache(), settings UI, startup eviction Cache management ~80

Phases 1 + 2 deliver visible value (avatars in the UI). Phases 35 add robustness.


2. Phase 1 — Protocol Download

2.1 Goal

Add a DownloadFile request variant to the protocol adapter's Request enum. Handle StreamItem::FileDownload in the connection task's event loop. Expose download_avatar() and download_icon() on ProtocolClient.

2.2 Files to Modify

crates/chanora_protocol/src/adapter.rs

2.2.1 Add Request::DownloadFile variant (after FetchClientProfile at line ~209)

/// Download a file from the server's file transfer subsystem.
/// Returns the complete file bytes on success.
DownloadFile {
    /// Channel ID the file resides in (0 for server-level files).
    channel_id: u64,
    /// File path on the server (e.g. "/avatar_<hex>" or "/icon_<crc32>").
    path: String,
    /// Optional channel password for password-protected channels.
    password: Option<String>,
    /// Reply channel for the result.
    reply: oneshot::Sender<Result<Vec<u8>, ProtocolError>>,
},

2.2.2 Add ProtocolClient::download_file() method (after send_text_message at line ~545)

/// Download a file from the server. Returns the complete file bytes.
pub async fn download_file(
    &self,
    channel_id: u64,
    path: String,
    password: Option<String>,
) -> Result<Vec<u8>, ProtocolError> {
    let (tx, rx) = oneshot::channel();
    self.tx
        .send(Request::DownloadFile {
            channel_id,
            path,
            password,
            reply: tx,
        })
        .await
        .map_err(|_| ProtocolError::Lost("connection task is gone".to_string()))?;
    rx.await
        .map_err(|_| ProtocolError::Lost("download_file reply dropped".to_string()))?
}

2.2.3 Add convenience wrappers (after download_file)

/// Download a client avatar. `uid_hex` is the hex-encoded client UID
/// (as produced by `uid_to_avatar_path`).
pub async fn download_avatar(
    &self,
    client_uid: &str,
) -> Result<Vec<u8>, ProtocolError> {
    let path = uid_to_avatar_path(client_uid);
    self.download_file(0, path, None).await
}

/// Download an icon by its unsigned CRC32 ID.
pub async fn download_icon(
    &self,
    icon_id: u64,
) -> Result<Vec<u8>, ProtocolError> {
    let path = format!("/icon_{icon_id}");
    self.download_file(0, path, None).await
}

2.2.4 Handle Request::DownloadFile in the main loop (in the rx.try_recv() match block, after FetchClientProfile around line 840)

Ok(Request::DownloadFile {
    channel_id,
    path,
    password,
    reply,
}) => {
    let r = download_file_from_server(
        &mut con,
        channel_id,
        &path,
        password.as_deref(),
    ).await;
    let _ = reply.send(r);
}

2.2.5 Add download_file_from_server function

/// Execute a file download via tsclientlib's download_file API.
/// Reads the entire TCP stream into a Vec<u8>.
async fn download_file_from_server(
    con: &mut Connection,
    channel_id: u64,
    path: &str,
    password: Option<&str>,
) -> Result<Vec<u8>, ProtocolError> {
    let pw = password.unwrap_or("");
    let mut stream = con
        .download_file(TsChannelId(channel_id), path, pw, 0)
        .await
        .map_err(|e| ProtocolError::Backend(format!("download_file init: {e}")))?;

    let mut buf = Vec::new();
    while let Some(item) = stream.next().await {
        match item {
            Ok(StreamItem::FileDownload(result)) => {
                // Read the TcpStream into buf.
                let mut reader = tokio::io::BufReader::new(result.stream);
                let size_hint = result.size as usize;
                if size_hint > 0 && size_hint < 32 * 1024 * 1024 {
                    buf.reserve(size_hint);
                }
                tokio::io::copy_to_bytes(&mut reader)
                    .await
                    .map_err(|e| ProtocolError::Backend(format!("download read: {e}")))?;
                // Note: actual implementation needs to accumulate into buf,
                // not use copy_to_bytes directly. See implementation note below.
            }
            Ok(_) => {} // other stream items during download
            Err(e) => {
                return Err(ProtocolError::Backend(format!("download stream: {e}")));
            }
        }
    }
    Ok(buf)
}

Implementation note: The exact tsclientlib download_file API signature and StreamItem::FileDownload handling needs to be verified against the pinned revision. The function above shows the intent; the actual implementation may differ based on whether download_file returns a TcpStream directly or wraps it in a StreamItem. The research doc (§3) confirmed: download_file(ChannelId, path, password, seek)StreamItem::FileDownload(FileDownloadResult { size, stream: TcpStream }). The implementation reads TcpStream bytes until EOF.

2.2.6 Handle StreamItem::FileDownload in handle_non_audio_stream_item

Currently handle_non_audio_stream_item ignores StreamItem::FileDownload (the catch-all _ => {} at line 1023). Once we have async download requests dispatched through the main loop, file download responses arrive here.

However, the tsclientlib download_file call is synchronous-blocking (it initiates the transfer and returns a stream handle). The StreamItem::FileDownload appears on the event stream as a notification that the server accepted the transfer. The actual data transfer happens on a separate TCP connection.

The exact integration pattern depends on tsclientlib's API:

  • If download_file() returns a impl Stream<Item = Result<StreamItem>>, we drain it inside download_file_from_server.
  • If the file download items arrive on the main con.events() stream, we need a pending-download tracking map (similar to pending_moves).

Decision: Defer the exact tsclientlib integration to implementation time. The research doc confirmed the API shape; the implementation agent will verify against the actual pinned tsclientlib source.

2.3 Verification

  • cargo test -p chanora_protocol passes
  • cargo clippy -p chanora_protocol clean
  • Manual: connect to server, call download_avatar(uid) → returns Vec<u8> with image data
  • Error path: download nonexistent avatar → ProtocolError returned cleanly

2.4 Risks

Risk Mitigation
tsclientlib download_file API differs from assumed shape Implementation agent reads pinned tsclientlib source before writing code
File transfer TCP stream blocks the main event loop download_file_from_server runs as async; the main loop continues on the next iteration
Large file OOM Cap download at 32 MB; return ProtocolError::Backend if exceeded

3. Phase 2 — Bridge + Flutter Display

3.1 Goal

Expose avatar download through the Flutter/Rust bridge. Display avatars in the Flutter UI using Image.memory(). No caching yet — every avatar view triggers a fresh download.

3.2 Files to Modify

crates/chanora_bridge/src/api.rs

3.2.1 Add bridge function download_avatar (after client_profile around line 660)

/// Download a client avatar by UID. Returns raw image bytes (PNG/JPEG).
/// Returns `None` if the client has no avatar or the download fails
/// without a protocol-level error.
pub async fn download_avatar(client_uid: String) -> Result<Option<Vec<u8>>, BridgeError> {
    runtime()
        .spawn(async move { session().download_avatar(&client_uid).await })
        .await
        .map_err(|e| task_join_error("download_avatar", e))?
        .map(|opt| opt.map(|bytes| bytes))
}

3.2.2 Add bridge function download_icon

/// Download an icon by its unsigned CRC32 ID. Returns raw image bytes.
pub async fn download_icon(icon_id: u64) -> Result<Option<Vec<u8>>, BridgeError> {
    runtime()
        .spawn(async move { session().download_icon(icon_id).await })
        .await
        .map_err(|e| task_join_error("download_icon", e))?
}

crates/chanora_core/src/lib.rs (or wherever ChanoraSession lives)

3.2.3 Add session methods

pub async fn download_avatar(&self, client_uid: &str) -> Result<Option<Vec<u8>>, CoreError> {
    let handle = self.protocol_handle()?;
    match handle.download_avatar(client_uid).await {
        Ok(bytes) => Ok(Some(bytes)),
        Err(ProtocolError::ServerRejected { .. }) => Ok(None),
        Err(e) => Err(CoreError::from(e)),
    }
}

pub async fn download_icon(&self, icon_id: u64) -> Result<Option<Vec<u8>>, CoreError> {
    let handle = self.protocol_handle()?;
    match handle.download_icon(icon_id).await {
        Ok(bytes) => Ok(Some(bytes)),
        Err(ProtocolError::ServerRejected { .. }) => Ok(None),
        Err(e) => Err(CoreError::from(e)),
    }
}

Flutter side (new files)

3.2.4 apps/chanora_flutter/lib/widgets/avatar_widget.dart

import 'dart:typed_data';
import 'package:flutter/material.dart';
import 'package:chanora_flutter/src/rust/api.dart' as rust;

/// Widget that displays a client avatar or initials fallback.
/// Phase 2: no caching — re-downloads on every build.
class ClientAvatar extends StatefulWidget {
  final String clientUid;
  final String displayName;
  final double radius;

  const ClientAvatar({
    super.key,
    required this.clientUid,
    required this.displayName,
    this.radius = 20,
  });

  @override
  State<ClientAvatar> createState() => _ClientAvatarState();
}

class _ClientAvatarState extends State<ClientAvatar> {
  Uint8List? _bytes;
  bool _loading = false;

  @override
  void initState() {
    super.initState();
    _loadAvatar();
  }

  @override
  void didUpdateWidget(ClientAvatar old) {
    super.didUpdateWidget(old);
    if (old.clientUid != widget.clientUid) {
      _bytes = null;
      _loadAvatar();
    }
  }

  Future<void> _loadAvatar() async {
    if (_loading) return;
    setState(() => _loading = true);
    try {
      final bytes = await rust.downloadAvatar(clientUid: widget.clientUid);
      if (mounted && bytes != null) {
        setState(() => _bytes = bytes);
      }
    } catch (_) {
      // Silently fall back to initials.
    } finally {
      if (mounted) {
        setState(() => _loading = false);
      }
    }
  }

  String get _initials {
    final parts = widget.displayName.trim().split(RegExp(r'\s+'));
    if (parts.length >= 2) {
      return '${parts[0][0]}${parts[1][0]}'.toUpperCase();
    }
    return widget.displayName.isNotEmpty
        ? widget.displayName[0].toUpperCase()
        : '?';
  }

  @override
  Widget build(BuildContext context) {
    if (_bytes != null) {
      return CircleAvatar(
        radius: widget.radius,
        backgroundImage: MemoryImage(_bytes!),
      );
    }
    return CircleAvatar(
      radius: widget.radius,
      child: Text(_initials),
    );
  }
}

3.3 Integration Points

The ClientAvatar widget replaces any existing hardcoded CircleAvatar in:

  • Client list rows
  • Chat message headers
  • User profile panels

Phase 2 limitation: Every initState triggers a download. Scrolling a list of 100 clients = 100 downloads. This is acceptable for development but NOT for release. Phase 4 (session orchestration) fixes this with coalescing and rate limiting.

3.4 Verification

  • flutter_rust_bridge_codegen generate succeeds (new bridge functions)
  • Flutter app builds without errors
  • Connect to server with avatar users → avatars render
  • Client without avatar → initials fallback renders
  • Scrolling client list does not crash (though it will be slow/chatty)

4. Phase 3 — Content-Addressed Blob Cache (chanora_cache)

4.1 Goal

Create a new chanora_cache crate with a BlobCache — a content-addressed store for avatar/icon blobs, backed by the cacache crate for production-tested crash safety and integrity verification. Uses the platform cache directory (not support directory).

4.2 Design Decisions (post-research)

Decision Choice Rationale
Crate placement New chanora_cache crate Cache is disposable; chanora_storage owns persistent identity/bookmark data with different durability and backup semantics
Directory Platform cache dir via getApplicationCacheDirectory() OS may evict under storage pressure; semantically correct for reconstructible data
Backing store cacache crate Production-tested crash safety, SSRI integrity verification on read, async-native, less code to maintain
Metadata None — no JSON sidecar files cacache manages its own content-addressed storage internally. No custom metadata needed.
Eviction cacache::ls() + custom mtime sweep cacache has no built-in LRU eviction. Walk entries, sort by mtime, delete oldest until under cap.
Integrity cacache SSRI on every read Detects corruption, bit rot, partial writes automatically. Better than custom "delete on read failure".
In-memory hot layer None for MVP Flutter's ImageCache already caches decoded images. Add Rust-side layer only if profiling shows need
Dependencies cacache (~6 transitive deps) sha2 already pulled in by chacha20poly1305 via chanora_storage. Crash safety and integrity worth the deps.

4.3 Storage Layout (on-disk, managed by cacache)

<cache_dir>/chanora/
  blobs/                         ← cacache content store root
    content-v2/
      <sha512-hex>/              ← content-addressed by SHA-512
        data                     ← raw blob bytes
    tmp/                         ← cacache temp files (in-flight writes)
    index-v2/                    ← cacache entry index (key → content mapping)

cacache manages this layout internally. Chanora's BlobCache maps protocol keys (av_<md5>, ic_<crc32>) to cacache string keys. Content dedup is automatic — same bytes stored once regardless of key.

The directory is opaque to Chanora. Inspection is via BlobCache API or cacache::ls().

4.4 Files to Create

crates/chanora_cache/Cargo.toml (NEW)

[package]
name = "chanora_cache"
description = "Chanora disposable content-addressed blob cache for avatars and icons"
version.workspace = true
edition.workspace = true
rust-version.workspace = true
authors.workspace = true
license.workspace = true
repository.workspace = true
publish.workspace = true

[dependencies]
cacache = "13"
thiserror.workspace = true
tokio.workspace = true
tracing.workspace = true

[dev-dependencies]
tempfile = "3"

crates/chanora_cache/src/lib.rs (NEW — ~120 LOC)

//! Disposable content-addressed blob cache for avatar and icon files.
//!
//! Wraps `cacache` for production-tested crash safety and integrity
//! verification. The on-disk layout is managed by cacache (content-v2,
//! index-v2). Chanora maps protocol keys (av_<md5>, ic_<crc32>) to
//! cacache string keys.
//!
//! This crate is intentionally separate from `chanora_storage`:
//! storage owns persistent identity/bookmark data; cache owns
//! reconstructible, disposable blob data with different durability
//! and backup semantics.

#![forbid(unsafe_code)]
#![warn(missing_docs)]

use std::path::{Path, PathBuf};
use tracing::{info, warn};

/// Errors raised by the blob cache.
#[derive(Debug, thiserror::Error)]
pub enum BlobCacheError {
    /// Filesystem I/O error.
    #[error("io: {0}")]
    Io(String),
    /// Key validation error.
    #[error("invalid key: {0}")]
    InvalidKey(String),
}

/// Content-addressed blob cache backed by cacache.
pub struct BlobCache {
    cache_dir: PathBuf,
    /// Maximum total cache size in bytes. 0 = no limit.
    max_bytes: u64,
}

/// Avatar blob prefix.
pub const PREFIX_AVATAR: &str = "av_";
/// Icon blob prefix.
pub const PREFIX_ICON: &str = "ic_";

impl BlobCache {
    /// Create or open a BlobCache rooted at `cache_dir/chanora/`.
    /// Creates the cacache directory. `max_bytes` sets the
    /// eviction threshold; 0 means no automatic eviction.
    pub fn new(cache_dir: impl AsRef<Path>, max_bytes: u64) -> Result<Self, BlobCacheError> {
        let cache_dir = cache_dir.as_ref().join("chanora").join("blobs");
        // cacache creates the directory on first write, but we
        // create it eagerly so total_size() works before any writes.
        std::fs::create_dir_all(&cache_dir)
            .map_err(|e| BlobCacheError::Io(format!("mkdir cache: {e}")))?;
        Ok(Self { cache_dir, max_bytes })
    }

    /// Store a blob. `prefix` is `PREFIX_AVATAR` or `PREFIX_ICON`.
    /// `key` is the content hash (MD5 hex for avatars, unsigned
    /// decimal CRC32 for icons).
    ///
    /// cacache handles dedup and integrity internally.
    /// This is async because cacache is async-native.
    pub async fn put(
        &self,
        prefix: &str,
        key: &str,
        data: &[u8],
    ) -> Result<(), BlobCacheError> {
        validate_key(prefix, key)?;
        let cache_key = format!("{prefix}{key}");
        cacache::write(&self.cache_dir, &cache_key, data)
            .await
            .map_err(|e| BlobCacheError::Io(format!("cacache write: {e}")))?;
        Ok(())
    }

    /// Read a blob. Returns `None` if not cached.
    /// cacache verifies SSRI integrity on every read.
    pub async fn get(&self, prefix: &str, key: &str) -> Result<Option<Vec<u8>>, BlobCacheError> {
        validate_key(prefix, key)?;
        let cache_key = format!("{prefix}{key}");
        match cacache::read(&self.cache_dir, &cache_key).await {
            Ok(data) => Ok(Some(data)),
            Err(cacache::Error::EntryNotFound(_, _)) => Ok(None),
            Err(e) => {
                // Integrity failure or I/O error — remove corrupt entry.
                warn!(target: "chanora_cache", key = %cache_key, error = %e,
                      "cache read failed; removing entry");
                let _ = cacache::remove_entry(&self.cache_dir, &cache_key).await;
                Ok(None)
            }
        }
    }

    /// Delete a specific blob.
    pub async fn remove(&self, prefix: &str, key: &str) -> Result<(), BlobCacheError> {
        validate_key(prefix, key)?;
        let cache_key = format!("{prefix}{key}");
        cacache::remove_entry(&self.cache_dir, &cache_key)
            .await
            .map_err(|e| BlobCacheError::Io(format!("cacache remove: {e}")))?;
        Ok(())
    }

    /// Delete all blobs.
    pub async fn clear(&self) -> Result<(), BlobCacheError> {
        // Remove and recreate the cache directory.
        let path = self.cache_dir.clone();
        tokio::task::spawn_blocking(move || {
            if path.exists() {
                std::fs::remove_dir_all(&path)
                    .map_err(|e| BlobCacheError::Io(format!("clear cache: {e}")))?;
                std::fs::create_dir_all(&path)
                    .map_err(|e| BlobCacheError::Io(format!("recreate cache dir: {e}")))?;
            }
            Ok(())
        })
        .await
        .map_err(|e| BlobCacheError::Io(format!("clear task: {e}")))?
    }

    /// Return total bytes used by all blobs.
    /// Walks cacache entries and sums sizes.
    pub async fn total_size(&self) -> Result<u64, BlobCacheError> {
        let mut total: u64 = 0;
        let mut stream = cacache::ls(&self.cache_dir).await
            .map_err(|e| BlobCacheError::Io(format!("cacache ls: {e}")))?;
        while let Some(entry) = stream.next().await {
            match entry {
                Ok(meta) => total += meta.size.unwrap_or(0),
                Err(e) => {
                    warn!(target: "chanora_cache", error = %e, "skipping bad entry during size scan");
                }
            }
        }
        Ok(total)
    }

    /// Evict oldest entries by mtime until total size is under
    /// `max_bytes`. Call on startup or periodically. No-op if
    /// `max_bytes` is 0.
    pub async fn evict(&self) -> Result<(), BlobCacheError> {
        if self.max_bytes == 0 {
            return Ok(());
        }
        let mut entries: Vec<(String, u64, std::time::SystemTime)> = Vec::new();
        let mut stream = cacache::ls(&self.cache_dir).await
            .map_err(|e| BlobCacheError::Io(format!("cacache ls: {e}")))?;
        while let Some(entry) = stream.next().await {
            match entry {
                Ok(meta) => {
                    let size = meta.size.unwrap_or(0);
                    let mtime = meta.time.into_system_time();
                    entries.push((meta.key, size, mtime));
                }
                Err(e) => {
                    warn!(target: "chanora_cache", error = %e, "skipping bad entry during eviction scan");
                }
            }
        }
        let total: u64 = entries.iter().map(|(_, s, _)| *s).sum();
        if total <= self.max_bytes {
            return Ok(());
        }
        // Sort by mtime ascending (oldest first).
        entries.sort_by_key(|(_, _, t)| *t);
        let mut freed: u64 = 0;
        let target = total - self.max_bytes;
        for (key, size, _) in entries {
            if freed >= target {
                break;
            }
            let _ = cacache::remove_entry(&self.cache_dir, &key).await;
            freed += size;
        }
        info!(target: "chanora_cache", freed_bytes = freed, "evicted oldest blobs");
        Ok(())
    }
}

/// Validate key format to prevent malformed entries.
fn validate_key(prefix: &str, key: &str) -> Result<(), BlobCacheError> {
    if !matches!(prefix, PREFIX_AVATAR | PREFIX_ICON) {
        return Err(BlobCacheError::InvalidKey(format!("bad prefix: {prefix}")));
    }
    match prefix {
        PREFIX_AVATAR => {
            // MD5 hex = exactly 32 hex chars.
            if key.len() != 32 || !key.chars().all(|c| c.is_ascii_hexdigit()) {
                return Err(BlobCacheError::InvalidKey(
                    format!("avatar key must be 32 hex chars, got: {key}")
                ));
            }
        }
        PREFIX_ICON => {
            // Unsigned CRC32 = decimal digits.
            if key.is_empty() || !key.chars().all(|c| c.is_ascii_digit()) {
                return Err(BlobCacheError::InvalidKey(
                    format!("icon key must be decimal digits, got: {key}")
                ));
            }
        }
        _ => unreachable!(),
    }
    Ok(())
}

#[cfg(test)]
mod tests {
    use super::*;

    fn tempdir() -> tempfile::TempDir {
        tempfile::Builder::new()
            .prefix("chanora_cache_test_")
            .tempdir()
            .unwrap()
    }

    #[tokio::test]
    async fn put_get_roundtrip() {
        let tmp = tempdir();
        let cache = BlobCache::new(&tmp, 0).unwrap();
        assert!(cache.get(PREFIX_AVATAR, "a1b2c3d4e5f6a7b8c9d0e1f2a3b4c5d6").await.unwrap().is_none());
        cache.put(PREFIX_AVATAR, "a1b2c3d4e5f6a7b8c9d0e1f2a3b4c5d6", b"avatar-bytes").await.unwrap();
        let data = cache.get(PREFIX_AVATAR, "a1b2c3d4e5f6a7b8c9d0e1f2a3b4c5d6").await.unwrap();
        assert_eq!(data.as_deref(), Some(b"avatar-bytes".as_slice()));
    }

    #[tokio::test]
    async fn get_missing_returns_none() {
        let tmp = tempdir();
        let cache = BlobCache::new(&tmp, 0).unwrap();
        assert!(cache.get(PREFIX_AVATAR, "00000000000000000000000000000000").await.unwrap().is_none());
    }

    #[tokio::test]
    async fn clear_removes_all() {
        let tmp = tempdir();
        let cache = BlobCache::new(&tmp, 0).unwrap();
        cache.put(PREFIX_AVATAR, "a1b2c3d4e5f6a7b8c9d0e1f2a3b4c5d6", b"data").await.unwrap();
        cache.put(PREFIX_ICON, "12345", b"icon").await.unwrap();
        cache.clear().await.unwrap();
        assert_eq!(cache.total_size().await.unwrap(), 0);
    }

    #[tokio::test]
    async fn total_size_accounts_for_all_entries() {
        let tmp = tempdir();
        let cache = BlobCache::new(&tmp, 0).unwrap();
        cache.put(PREFIX_AVATAR, "a1b2c3d4e5f6a7b8c9d0e1f2a3b4c5d6", b"12345").await.unwrap();
        cache.put(PREFIX_ICON, "99999", b"12").await.unwrap();
        assert_eq!(cache.total_size().await.unwrap(), 5 + 2);
    }

    #[tokio::test]
    async fn invalid_key_rejected() {
        let tmp = tempdir();
        let cache = BlobCache::new(&tmp, 0).unwrap();
        // Too short for MD5.
        assert!(cache.put(PREFIX_AVATAR, "abc", b"data").await.is_err());
        // Non-hex in MD5.
        assert!(cache.put(PREFIX_AVATAR, "g".repeat(32).as_str(), b"data").await.is_err());
        // Non-digit in icon key.
        assert!(cache.put(PREFIX_ICON, "12a45", b"data").await.is_err());
        // Bad prefix.
        assert!(cache.put("xx_", "abc", b"data").await.is_err());
    }

    #[tokio::test]
    async fn evict_deletes_oldest_until_under_cap() {
        let tmp = tempdir();
        // 10 byte cap.
        let cache = BlobCache::new(&tmp, 10).unwrap();
        cache.put(PREFIX_AVATAR, "a1b2c3d4e5f6a7b8c9d0e1f2a3b4c5d6", b"12345678").await.unwrap(); // 8 bytes
        tokio::time::sleep(std::time::Duration::from_millis(50)).await;
        cache.put(PREFIX_ICON, "11111", b"12345").await.unwrap(); // 5 bytes → total 13, over cap
        cache.evict().await.unwrap();
        // Oldest (avatar) should be evicted.
        assert!(cache.get(PREFIX_AVATAR, "a1b2c3d4e5f6a7b8c9d0e1f2a3b4c5d6").await.unwrap().is_none());
        assert!(cache.get(PREFIX_ICON, "11111").await.unwrap().is_some());
    }

    #[tokio::test]
    async fn corrupt_entry_removed_on_read() {
        let tmp = tempdir();
        let cache = BlobCache::new(&tmp, 0).unwrap();
        // Write valid entry.
        cache.put(PREFIX_AVATAR, "a1b2c3d4e5f6a7b8c9d0e1f2a3b4c5d6", b"good").await.unwrap();
        // Corrupt the underlying content file (simulate bit rot).
        // cacache stores content under content-v2/<sha512>/data.
        // We find and corrupt it to test integrity check.
        let content_dir = tmp.path().join("chanora").join("blobs").join("content-v2");
        if let Ok(mut entries) = std::fs::read_dir(&content_dir) {
            while let Some(Ok(entry)) = entries.next() {
                let data_path = entry.path().join("data");
                if data_path.exists() {
                    std::fs::write(&data_path, b"corrupt").unwrap();
                    break;
                }
            }
        }
        // Read should detect corruption and return None.
        let result = cache.get(PREFIX_AVATAR, "a1b2c3d4e5f6a7b8c9d0e1f2a3b4c5d6").await.unwrap();
        assert!(result.is_none());
    }
}

4.5 Wire BlobCache into Session

The BlobCache is constructed once during a new init_cache bridge call (separate from init_storage) using the platform cache directory. All BlobCache methods are async (cacache is async-native).

crates/chanora_bridge/src/api.rs

/// Wire the blob cache to a platform cache directory. Call once
/// on app start after Flutter resolves `getApplicationCacheDirectory()`.
pub async fn init_cache(dir: String) -> Result<(), BridgeError> {
    runtime()
        .spawn(async move { session().init_cache(&dir).await })
        .await
        .map_err(|e| task_join_error("init_cache", e))??;
    Ok(())
}

Flutter init (in app_bootstrap.dart)

// Existing:
final supportDir = await getApplicationSupportDirectory();
await rust.initStorage(dir: supportDir.path);

// NEW: separate cache directory
final cacheDir = await getApplicationCacheDirectory();
await rust.initCache(dir: cacheDir.path);

chanora_core::ChanoraSession

blob_cache: OnceCell<chanora_cache::BlobCache>,

pub async fn init_cache(&self, dir: &str) -> Result<(), CoreError> {
    // 100 MB cap for avatars + icons.
    let cache = chanora_cache::BlobCache::new(dir, 100 * 1024 * 1024)?;
    // Evict on startup.
    cache.evict().await?;
    self.blob_cache.set(cache).map_err(|_| CoreError::AlreadyConnected)?;
    Ok(())
}

Download method with cache

pub async fn get_avatar(&self, avatar_hash: &str, client_uid: &str) -> Result<Option<Vec<u8>>, CoreError> {
    // 1. Check cache
    if let Some(cache) = self.blob_cache.get() {
        if let Some(bytes) = cache.get(chanora_cache::PREFIX_AVATAR, avatar_hash).await? {
            return Ok(Some(bytes));
        }
    }
    // 2. Download
    let handle = self.protocol_handle()?;
    let bytes = handle.download_avatar(client_uid).await?;
    // 3. Store in cache
    if let Some(cache) = self.blob_cache.get() {
        cache.put(chanora_cache::PREFIX_AVATAR, avatar_hash, &bytes).await?;
    }
    Ok(Some(bytes))
}

4.6 Verification

  • cargo test -p chanora_cache — all 8 async tests pass
  • cargo clippy -p chanora_cache clean
  • Cache directory appears under <cache_dir>/chanora/blobs/ after download
  • Second download of same avatar hits cache (no TCP transfer)
  • Same avatar from different server uses same content (cacache dedup)
  • evict() deletes oldest entries when over cap
  • Corrupted content detected by cacache integrity check → entry removed
  • Invalid keys are rejected
  • clear() empties cache

5. Phase 4 — Session Orchestration (chanora_core)

5.1 Goal

Add a FileTransferService to chanora_core that sits between the bridge and the raw protocol download. Provides:

  1. Request coalescing — 10 concurrent get_avatar calls for the same hash = 1 download, 10 waiters.
  2. Concurrency throttle — At most 2 simultaneous file transfers (anti-flood).
  3. Negative cache — Remember "avatar X does not exist" for 5 minutes to avoid re-downloading 404s.
  4. Retry with exponential backoff — 5s → 30s → 2min → 10min for transient failures.

5.2 Design

Bridge (download_avatar)
  ↓
FileTransferService::get_avatar(hash, uid)
  ↓
  1. Check BlobCache → hit? return bytes
  2. Check negative cache → hit? return None
  3. Check in-flight map → already downloading? await existing handle
  4. Acquire concurrency semaphore (max 2)
  5. Download via ProtocolClient
  6. Store in BlobCache
  7. Notify all waiters

5.3 Files to Create/Modify

crates/chanora_core/src/file_transfer.rs (NEW — ~250 LOC)

use std::collections::HashMap;
use std::sync::Arc;
use tokio::sync::{Semaphore, SemaphorePermit, oneshot, Mutex};

pub struct FileTransferService {
    cache: chanora_cache::BlobCache,
    protocol: Arc<Mutex<Option<chanora_protocol::ProtocolClient>>>,
    /// Max concurrent downloads.
    semaphore: Arc<Semaphore>,
    /// In-flight downloads: (prefix+key) → Vec<oneshot::Sender<Result<...>>>
    in_flight: Arc<Mutex<HashMap<String, Vec<oneshot::Sender<Result<Option<Vec<u8>>, FileTransferError>>>>>>,
    /// Negative cache: key → instant when entry expires.
    negative_cache: Arc<Mutex<HashMap<String, std::time::Instant>>>,
    /// Negative cache TTL.
    negative_ttl: std::time::Duration,
}

#[derive(Debug, thiserror::Error)]
pub enum FileTransferError {
    #[error("not connected")]
    NotConnected,
    #[error("protocol: {0}")]
    Protocol(#[from] chanora_protocol::ProtocolError),
    #[error("cache: {0}")]
    Cache(#[from] chanora_cache::BlobCacheError),
}

impl FileTransferService {
    pub fn new(cache: chanora_cache::BlobCache) -> Self {
        Self {
            cache,
            protocol: Arc::new(Mutex::new(None)),
            semaphore: Arc::new(Semaphore::new(2)),
            in_flight: Arc::new(Mutex::new(HashMap::new())),
            negative_cache: Arc::new(Mutex::new(HashMap::new())),
            negative_ttl: std::time::Duration::from_secs(300), // 5 min
        }
    }

    pub async fn set_protocol(&self, client: Option<chanora_protocol::ProtocolClient>) {
        *self.protocol.lock().await = client;
    }

    /// Get an avatar by hash + UID. Coalesces concurrent requests.
    pub async fn get_avatar(
        &self,
        avatar_hash: &str,
        client_uid: &str,
    ) -> Result<Option<Vec<u8>>, FileTransferError> {
        let cache_key = format!("av_{avatar_hash}");
        // 1. Disk cache hit?
        if let Some(bytes) = self.cache.get("av_", avatar_hash).await? {
            return Ok(Some(bytes));
        }
        // 2. Negative cache hit?
        {
            let neg = self.negative_cache.lock().await;
            if let Some(expires) = neg.get(&cache_key) {
                if std::time::Instant::now() < *expires {
                    return Ok(None);
                }
            }
        }
        // 3. Already in-flight?
        {
            let mut in_flight = self.in_flight.lock().await;
            if let Some(waiters) = in_flight.get_mut(&cache_key) {
                let (tx, rx) = oneshot::channel();
                waiters.push(tx);
                drop(in_flight);
                return rx.await
                    .map_err(|_| FileTransferError::NotConnected)?
                    .map(|opt| opt);
            }
            // First requester — register in-flight.
            in_flight.insert(cache_key.clone(), Vec::new());
        }
        // 4. Acquire semaphore + download.
        let _permit = self.semaphore.acquire().await
            .map_err(|_| FileTransferError::NotConnected)?;
        let result = self.do_download_avatar(avatar_hash, client_uid).await;
        // 5. Notify waiters.
        {
            let mut in_flight = self.in_flight.lock().await;
            if let Some(waiters) = in_flight.remove(&cache_key) {
                for tx in waiters {
                    let _ = tx.send(result.clone());
                }
            }
        }
        result
    }

    async fn do_download_avatar(
        &self,
        avatar_hash: &str,
        client_uid: &str,
    ) -> Result<Option<Vec<u8>>, FileTransferError> {
        let guard = self.protocol.lock().await;
        let proto = guard.as_ref().ok_or(FileTransferError::NotConnected)?;
        let cloned = proto.clone();
        drop(guard);

        match cloned.download_avatar(client_uid).await {
            Ok(bytes) => {
                self.cache.put("av_", avatar_hash, &bytes).await?;
                Ok(Some(bytes))
            }
            Err(chanora_protocol::ProtocolError::ServerRejected { .. }) => {
                // Avatar doesn't exist on server — negative cache.
                let mut neg = self.negative_cache.lock().await;
                neg.insert(
                    format!("av_{avatar_hash}"),
                    std::time::Instant::now() + self.negative_ttl,
                );
                Ok(None)
            }
            Err(e) => Err(FileTransferError::Protocol(e)),
        }
    }

    /// Get an icon by CRC32. Same coalescing logic as avatars.
    pub async fn get_icon(
        &self,
        icon_id: u64,
    ) -> Result<Option<Vec<u8>>, FileTransferError> {
        let cache_key = format!("ic_{icon_id}");
        if let Some(bytes) = self.cache.get("ic_", &icon_id.to_string()).await? {
            return Ok(Some(bytes));
        }
        {
            let neg = self.negative_cache.lock().await;
            if let Some(expires) = neg.get(&cache_key) {
                if std::time::Instant::now() < *expires {
                    return Ok(None);
                }
            }
        }
        // Same coalescing pattern as get_avatar...
        // (Implementation mirrors get_avatar with icon-specific path)
        let _permit = self.semaphore.acquire().await
            .map_err(|_| FileTransferError::NotConnected)?;
        let guard = self.protocol.lock().await;
        let proto = guard.as_ref().ok_or(FileTransferError::NotConnected)?;
        let cloned = proto.clone();
        drop(guard);

        match cloned.download_icon(icon_id).await {
            Ok(bytes) => {
                self.cache.put("ic_", &icon_id.to_string(), &bytes).await?;
                Ok(Some(bytes))
            }
            Err(chanora_protocol::ProtocolError::ServerRejected { .. }) => {
                let mut neg = self.negative_cache.lock().await;
                neg.insert(cache_key, std::time::Instant::now() + self.negative_ttl);
                Ok(None)
            }
            Err(e) => Err(FileTransferError::Protocol(e)),
        }
    }
}

5.4 Wire into ChanoraSession

// In ChanoraSession:
file_transfer: OnceCell<FileTransferService>,

// After connect:
file_transfer.set_protocol(Some(protocol_client)).await;

// After disconnect:
file_transfer.set_protocol(None).await;

5.5 Bridge Updates

Replace Phase 2's raw download_avatar with cache-aware version:

pub async fn get_avatar(
    avatar_hash: String,
    client_uid: String,
) -> Result<Option<Vec<u8>>, BridgeError> {
    runtime()
        .spawn(async move {
            session().file_transfer()
                .get_avatar(&avatar_hash, &client_uid)
                .await
        })
        .await
        .map_err(|e| task_join_error("get_avatar", e))?
        .map_err(BridgeError::from)
}

5.6 Flutter Widget Update

Update ClientAvatar to pass avatar_hash (from client_flag_avatar property in the snapshot):

class ClientAvatar extends StatefulWidget {
  final String clientUid;
  final String avatarHash; // client_flag_avatar (MD5 hex)
  final String displayName;
  final double radius;
  // ...
}

Future<void> _loadAvatar() async {
  if (_loading || widget.avatarHash.isEmpty) return;
  setState(() => _loading = true);
  try {
    final bytes = await rust.getAvatar(
      avatarHash: widget.avatarHash,
      clientUid: widget.clientUid,
    );
    if (mounted && bytes != null) {
      setState(() => _bytes = bytes);
    }
  } catch (_) {
    // Fall back to initials.
  } finally {
    if (mounted) setState(() => _loading = false);
  }
}

5.7 ProtocolDelta / BridgeEvent Updates

Add avatar hash to ClientJoined and ClientUpdated events so Flutter knows the client_flag_avatar value:

crates/chanora_protocol/src/dto.rs — add avatar_hash to ClientInfo:

pub struct ClientInfo {
    // ... existing fields ...
    /// MD5 hash of the client's avatar bytes (client_flag_avatar).
    /// Empty string means no avatar.
    pub avatar_hash: String,
}

crates/chanora_protocol/src/adapter.rs — populate avatar_hash in build_snapshot:

Read client_flag_avatar from tsclientlib state when building ClientInfo.

crates/chanora_bridge/src/api.rs — add avatar_hash to BridgeClient:

pub struct BridgeClient {
    // ... existing fields ...
    /// MD5 hash of the client's avatar (client_flag_avatar). Empty = no avatar.
    pub avatar_hash: String,
}

BridgeEvent::ClientJoined / BridgeEvent::ClientUpdated — add avatar_hash field

This is the key change that lets the Flutter widget know which avatar to display for each client.

5.8 Verification

  • 10 concurrent get_avatar calls for same hash → 1 download, 10 responses
  • Non-existent avatar → returns None, negative cache prevents re-download
  • Semaphore limits concurrent downloads to 2
  • Cache hit returns bytes without network call
  • cargo test -p chanora_core passes
  • Flutter UI shows avatars without excessive network traffic

6. Phase 5 — Cache Management

6.1 Goal

Expose cache management through the bridge. Add a "Clear Cache" button to settings.

6.2 Files to Modify

crates/chanora_bridge/src/api.rs

/// Clear all cached avatar and icon data.
pub async fn clear_file_cache() -> Result<(), BridgeError> {
    runtime()
        .spawn(async { session().clear_file_cache().await })
        .await
        .map_err(|e| task_join_error("clear_file_cache", e))?
}

/// Get total cache size in bytes.
pub async fn file_cache_size() -> Result<u64, BridgeError> {
    runtime()
        .spawn(async { session().file_cache_size() })
        .await
        .map_err(|e| task_join_error("file_cache_size", e))?
}

crates/chanora_core/src/lib.rs

pub async fn clear_file_cache(&self) -> Result<(), CoreError> {
    if let Some(ft) = self.file_transfer.as_ref() {
        ft.cache.clear()?;
    }
    Ok(())
}

pub fn file_cache_size(&self) -> Result<u64, CoreError> {
    if let Some(ft) = self.file_transfer.as_ref() {
        Ok(ft.cache.total_size()?)
    } else {
        Ok(0)
    }
}

Flutter settings UI

Add to settings page:

ListTile(
  title: const Text('Clear avatar cache'),
  subtitle: FutureBuilder<String>(
    future: _cacheSizeLabel(),
    builder: (_, snap) => Text(snap.data ?? 'Calculating...'),
  ),
  trailing: IconButton(
    icon: const Icon(Icons.delete_outline),
    onPressed: () async {
      await rust.clearFileCache();
      PaintingBinding.instance.imageCache.clear();
      setState(() {}); // Refresh size display
    },
  ),
),

6.3 Post-MVP Eviction (not in this plan)

  • LRU eviction on startup: delete blobs older than N days
  • Configurable cache size limit
  • Per-server metadata sidecar files (refs/<host>_<port>.json)

6.4 Verification

  • clear_file_cache() empties blobs/ directory
  • file_cache_size() returns 0 after clear
  • Flutter settings shows cache size
  • After clear, avatars re-download on next view

7. Cross-Phase Dependency Map

Phase 1 (protocol download)
  ↓
Phase 2 (bridge + Flutter display) — uses Phase 1's download_avatar
  ↓
Phase 3 (BlobCache) — used by Phase 4
  ↓
Phase 4 (FileTransferService) — uses Phase 1 + Phase 3
  ↓
Phase 5 (cache management) — uses Phase 3 + Phase 4

Phases 1 and 2 can be tested end-to-end (avatars in UI) before any caching exists. Phases 3 and 4 can be developed in parallel but Phase 4 depends on Phase 3's BlobCache. Phase 5 depends on Phase 3.


8. Testing Strategy

Unit Tests

Crate Test Phase
chanora_cache BlobCache put/get/clear/evict/key validation 3
chanora_core FileTransferService coalescing, negative cache, semaphore 4
chanora_protocol uid_to_avatar_path (existing), download_file_from_server (mock) 1

Integration Tests

Test What Phase
Download avatar from real server Connect, call download_avatar, verify bytes are valid PNG/JPEG 1
Cache round-trip Download → cache hit → disconnect → reconnect → cache hit 3+4
Concurrent coalescing 10 tasks request same avatar → 1 download 4
Negative cache Request non-existent avatar → get None → immediate None on retry 4

Manual Verification

Scenario Expected Phase
Connect to server with 50+ users Avatars load progressively, no flood 2+4
Scroll client list rapidly No duplicate downloads, smooth scrolling 4
Connect to 2 servers with overlapping users Same avatar = 1 blob file 3
Settings → Clear Cache Cache empties, avatars re-download 5

9. Out of Scope (Deferred)

Item Reason When
Icon download + channel/client/server icon display Avatar is higher priority; icons follow the same pipeline Phase 4+
File upload Not in MVP Post-MVP
File browser Not in MVP Post-MVP
myTeamSpeak avatar integration Requires cloud OAuth; indefinitely deferred Post-MVP
Per-server metadata sidecars Nice-to-have for cache inspection Phase 5+
Thumbnail generation Avatars are small; no thumbnails needed Post-MVP
Animated avatars TS3 does not support animated avatars N/A

10. Implementation Notes

10.1 tsclientlib API Verification

The implementation agent for Phase 1 MUST read the pinned tsclientlib source before writing any code. Key files:

  • tsclientlib/src/lib.rs lines 9561005 — download_file public API
  • tsclientlib/src/lib.rs lines 13711427 — StreamItem::FileDownload handling
  • Local clone at /var/folders/.../T/opencode/tsclientlib-rs-research/

10.2 FRB Code Generation

Phase 2 and 5 add new bridge functions. After modifying api.rs, run:

flutter_rust_bridge_codegen generate

10.3 avatar_hash in ProtocolDelta

Phase 4 adds avatar_hash to client events. This requires changes in:

  1. chanora_protocol::dto::ClientInfo — add field
  2. chanora_protocol::adapter::build_snapshot — populate from tsclientlib state
  3. chanora_protocol::adapter::forward_delta — extract client_flag_avatar from PropertyChanged events
  4. chanora_core::SessionEvent::ClientJoined/Updated — add field
  5. chanora_bridge::BridgeEvent::ClientJoined/Updated — add field
  6. Flutter ClientJoinedEvent / ClientUpdatedEvent — read field

This is a cross-cutting change. It should be done as a separate, atomic commit before Phase 4's Flutter widget updates.

10.4 Atomic Commits

Recommended commit boundaries:

  1. feat(protocol): add DownloadFile request and file download support (Phase 1)
  2. feat(bridge): add download_avatar bridge function (Phase 2 bridge)
  3. feat(ui): add ClientAvatar widget with initials fallback (Phase 2 Flutter)
  4. feat(cache): add chanora_cache crate with cacache-backed blob cache (Phase 3)
  5. feat(core): add avatar_hash to client events and DTOs (Phase 4 data plumbing)
  6. feat(core): add FileTransferService with coalescing and rate limiting (Phase 4)
  7. feat(ui): update ClientAvatar to use cached downloads (Phase 4 Flutter)
  8. feat(bridge): add clear_file_cache and file_cache_size (Phase 5 bridge)
  9. feat(ui): add cache management to settings (Phase 5 Flutter)