core+apps: file transfer progress
Add a Transfer record {name, received, total, incoming, done} and an engine-side
transfers map updated by the RTC/LAN transports — fine-grained on receive (per
chunk in apply_file_frame), per-chunk on send. Engine.transfers() exposes
in-flight transfers (completed linger ~6s, then prune). Both apps poll it and
show a progress row with a bar + percent above the feed.
Verified with examples/progress_demo.rs: a 5MB file's transfer reports
received/total advancing to done.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
44
core/examples/progress_demo.rs
Normal file
44
core/examples/progress_demo.rs
Normal file
@@ -0,0 +1,44 @@
|
||||
//! Sanity-check transfer progress: A sends a 5 MB file; poll B.transfers() and
|
||||
//! A.transfers() to see received/total advance to done.
|
||||
//! cargo run --example progress_demo
|
||||
|
||||
use std::time::Duration;
|
||||
|
||||
use tethercore::{Engine, Message, MessageHandler};
|
||||
|
||||
struct Noop;
|
||||
impl MessageHandler for Noop {
|
||||
fn on_message(&self, _m: Message) {}
|
||||
fn on_status(&self, _c: bool) {}
|
||||
fn on_file(&self, name: String, _mime: String, data: Vec<u8>) {
|
||||
eprintln!(" received {name} ({} bytes)", data.len());
|
||||
}
|
||||
}
|
||||
|
||||
fn main() {
|
||||
let dead = "http://127.0.0.1:1".to_string(); // LAN only
|
||||
let room = "progress-room".to_string();
|
||||
let payload = vec![9u8; 5_000_000];
|
||||
let a = Engine::new(dead.clone(), room.clone(), "dev-a".into(), "macos".into(), "A".into(), "".into());
|
||||
let b = Engine::new(dead, room, "dev-b".into(), "ios".into(), "B".into(), "".into());
|
||||
a.start(Box::new(Noop));
|
||||
b.start(Box::new(Noop));
|
||||
|
||||
eprintln!("waiting for LAN link…");
|
||||
std::thread::sleep(Duration::from_secs(5));
|
||||
a.send_file("big.bin".into(), "application/octet-stream".into(), payload);
|
||||
|
||||
for _ in 0..30 {
|
||||
std::thread::sleep(Duration::from_millis(150));
|
||||
let rx = b.transfers();
|
||||
let tx = a.transfers();
|
||||
if let Some(t) = rx.first() {
|
||||
println!("↓ recv {}/{} {}", t.received, t.total, if t.done { "✅ done" } else { "…" });
|
||||
if t.done { break; }
|
||||
} else if let Some(t) = tx.first() {
|
||||
println!("↑ send {}/{} {}", t.received, t.total, if t.done { "done" } else { "…" });
|
||||
}
|
||||
}
|
||||
a.stop();
|
||||
b.stop();
|
||||
}
|
||||
@@ -38,10 +38,16 @@ pub(crate) struct LanTransport {
|
||||
incoming: Arc<Mutex<HashMap<String, crate::IncomingFile>>>, // in-flight inbound files
|
||||
discovered: Discovered, // nearby devices (any room), for the engine's nearby()
|
||||
direct: DirectSet, // LAN-connected peers (a direct path; suppresses relay)
|
||||
transfers: crate::Transfers, // file-transfer progress
|
||||
}
|
||||
|
||||
impl LanTransport {
|
||||
pub(crate) fn new(cfg: Config, discovered: Discovered, direct: DirectSet) -> Self {
|
||||
pub(crate) fn new(
|
||||
cfg: Config,
|
||||
discovered: Discovered,
|
||||
direct: DirectSet,
|
||||
transfers: crate::Transfers,
|
||||
) -> Self {
|
||||
Self {
|
||||
cfg,
|
||||
peers: Arc::new(Mutex::new(HashMap::new())),
|
||||
@@ -50,6 +56,7 @@ impl LanTransport {
|
||||
incoming: Arc::new(Mutex::new(HashMap::new())),
|
||||
discovered,
|
||||
direct,
|
||||
transfers,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -123,20 +130,29 @@ impl Transport for LanTransport {
|
||||
|
||||
fn send_file(&self, rt: &Handle, id: String, name: String, mime: String, data: Vec<u8>) {
|
||||
let peers = self.peers.clone();
|
||||
let transfers = self.transfers.clone();
|
||||
rt.spawn(async move {
|
||||
let writers: Vec<Writer> = peers.lock().await.values().cloned().collect();
|
||||
if writers.is_empty() {
|
||||
return;
|
||||
}
|
||||
let total = data.len() as u64;
|
||||
crate::transfer_start(&transfers, &id, &name, total, false);
|
||||
let frames = crate::file_frames(&id, &name, &mime, &data);
|
||||
for w in writers {
|
||||
let mut g = w.lock().await;
|
||||
let mut sent = 0u64;
|
||||
for (tag, payload) in &frames {
|
||||
if write_frame(&mut g, *tag, payload).await.is_err() {
|
||||
break;
|
||||
}
|
||||
if *tag == b'C' {
|
||||
sent += (payload.len() - 16) as u64;
|
||||
crate::transfer_progress(&transfers, &id, sent);
|
||||
}
|
||||
}
|
||||
}
|
||||
crate::transfer_done(&transfers, &id, total);
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -307,7 +323,7 @@ impl LanTransport {
|
||||
b'M' | b'C' | b'E' => {
|
||||
let done = {
|
||||
let mut map = self.incoming.lock().await;
|
||||
crate::apply_file_frame(&mut map, tag, &payload)
|
||||
crate::apply_file_frame(&mut map, &self.transfers, tag, &payload)
|
||||
};
|
||||
if let Some((name, mime, bytes)) = done {
|
||||
eprintln!("tethercore[lan]: received file {name:?} ({} bytes)", bytes.len());
|
||||
|
||||
@@ -155,6 +155,19 @@ pub(crate) struct IncomingFile {
|
||||
pub data: Vec<u8>,
|
||||
}
|
||||
|
||||
/// Progress of an in-flight file transfer, for the UI.
|
||||
#[derive(Clone, Debug, uniffi::Record)]
|
||||
pub struct Transfer {
|
||||
pub name: String,
|
||||
pub received: u64, // bytes transferred so far
|
||||
pub total: u64, // total bytes
|
||||
pub incoming: bool, // true = receiving, false = sending
|
||||
pub done: bool,
|
||||
}
|
||||
|
||||
/// Transfer table, keyed by file id, with last-update time for pruning.
|
||||
pub(crate) type Transfers = Arc<Mutex<HashMap<String, (Transfer, Instant)>>>;
|
||||
|
||||
/// Build the frames for a file: M<json meta>, then C<16-byte id><chunk>…, then
|
||||
/// E<16-byte id>. Each is (tag, payload); the transport adds its own envelope
|
||||
/// (RTC: tag-prefixed binary message; LAN: length-prefixed).
|
||||
@@ -173,10 +186,32 @@ pub(crate) fn file_frames(id: &str, name: &str, mime: &str, data: &[u8]) -> Vec<
|
||||
out
|
||||
}
|
||||
|
||||
/// Apply one inbound file frame to the in-progress map; returns the completed
|
||||
/// (name, mime, data) on the terminating 'E' frame.
|
||||
/// Send-progress helpers (used by the RTC/LAN transports).
|
||||
pub(crate) fn transfer_start(transfers: &Transfers, id: &str, name: &str, total: u64, incoming: bool) {
|
||||
transfers.lock().unwrap().insert(
|
||||
id.to_string(),
|
||||
(Transfer { name: name.to_string(), received: 0, total, incoming, done: false }, Instant::now()),
|
||||
);
|
||||
}
|
||||
pub(crate) fn transfer_progress(transfers: &Transfers, id: &str, received: u64) {
|
||||
if let Some((t, ts)) = transfers.lock().unwrap().get_mut(id) {
|
||||
t.received = received;
|
||||
*ts = Instant::now();
|
||||
}
|
||||
}
|
||||
pub(crate) fn transfer_done(transfers: &Transfers, id: &str, total: u64) {
|
||||
if let Some((t, ts)) = transfers.lock().unwrap().get_mut(id) {
|
||||
t.received = total;
|
||||
t.done = true;
|
||||
*ts = Instant::now();
|
||||
}
|
||||
}
|
||||
|
||||
/// Apply one inbound file frame to the in-progress map (and update transfer
|
||||
/// progress); returns the completed (name, mime, data) on the 'E' frame.
|
||||
pub(crate) fn apply_file_frame(
|
||||
map: &mut HashMap<String, IncomingFile>,
|
||||
transfers: &Transfers,
|
||||
tag: u8,
|
||||
rest: &[u8],
|
||||
) -> Option<(String, String, Vec<u8>)> {
|
||||
@@ -185,14 +220,20 @@ pub(crate) fn apply_file_frame(
|
||||
if let Ok(v) = serde_json::from_slice::<serde_json::Value>(rest) {
|
||||
let id = v["id"].as_str().unwrap_or("").to_string();
|
||||
if !id.is_empty() {
|
||||
let name = v["name"].as_str().unwrap_or("file").to_string();
|
||||
let total = v["size"].as_u64().unwrap_or(0);
|
||||
map.insert(
|
||||
id,
|
||||
id.clone(),
|
||||
IncomingFile {
|
||||
name: v["name"].as_str().unwrap_or("file").to_string(),
|
||||
name: name.clone(),
|
||||
mime: v["mime"].as_str().unwrap_or("application/octet-stream").to_string(),
|
||||
data: Vec::new(),
|
||||
},
|
||||
);
|
||||
transfers.lock().unwrap().insert(
|
||||
id,
|
||||
(Transfer { name, received: 0, total, incoming: true, done: false }, Instant::now()),
|
||||
);
|
||||
}
|
||||
}
|
||||
None
|
||||
@@ -202,10 +243,19 @@ pub(crate) fn apply_file_frame(
|
||||
if let Some(f) = map.get_mut(&id) {
|
||||
f.data.extend_from_slice(&rest[16..]);
|
||||
}
|
||||
if let Some((t, ts)) = transfers.lock().unwrap().get_mut(&id) {
|
||||
t.received += (rest.len() - 16) as u64;
|
||||
*ts = Instant::now();
|
||||
}
|
||||
None
|
||||
}
|
||||
b'E' if rest.len() >= 16 => {
|
||||
let id = String::from_utf8_lossy(&rest[..16]).to_string();
|
||||
if let Some((t, ts)) = transfers.lock().unwrap().get_mut(&id) {
|
||||
t.done = true;
|
||||
t.received = t.total;
|
||||
*ts = Instant::now();
|
||||
}
|
||||
map.remove(&id).map(|f| (f.name, f.mime, f.data))
|
||||
}
|
||||
_ => None,
|
||||
@@ -246,6 +296,7 @@ pub struct Engine {
|
||||
direct: DirectSet,
|
||||
present: PresentSet,
|
||||
file_counter: AtomicU64,
|
||||
transfers: Transfers,
|
||||
}
|
||||
|
||||
#[uniffi::export]
|
||||
@@ -274,13 +325,24 @@ impl Engine {
|
||||
let discovered: Discovered = Arc::new(Mutex::new(HashMap::new()));
|
||||
let direct: DirectSet = Arc::new(Mutex::new(HashSet::new()));
|
||||
let present: PresentSet = Arc::new(Mutex::new(HashMap::new()));
|
||||
let transfers: Transfers = Arc::new(Mutex::new(HashMap::new()));
|
||||
let transports: Vec<Arc<dyn Transport>> = vec![
|
||||
Arc::new(SseTransport {
|
||||
cfg: cfg.clone(),
|
||||
http: reqwest::Client::new(),
|
||||
}),
|
||||
Arc::new(rtc::RtcTransport::new(cfg.clone(), direct.clone(), present.clone())),
|
||||
Arc::new(lan::LanTransport::new(cfg.clone(), discovered.clone(), direct.clone())),
|
||||
Arc::new(rtc::RtcTransport::new(
|
||||
cfg.clone(),
|
||||
direct.clone(),
|
||||
present.clone(),
|
||||
transfers.clone(),
|
||||
)),
|
||||
Arc::new(lan::LanTransport::new(
|
||||
cfg.clone(),
|
||||
discovered.clone(),
|
||||
direct.clone(),
|
||||
transfers.clone(),
|
||||
)),
|
||||
];
|
||||
Arc::new(Self {
|
||||
cfg,
|
||||
@@ -293,9 +355,19 @@ impl Engine {
|
||||
direct,
|
||||
present,
|
||||
file_counter: AtomicU64::new(0),
|
||||
transfers,
|
||||
})
|
||||
}
|
||||
|
||||
/// In-flight file transfers (send + receive) for a progress UI. Completed
|
||||
/// ones linger ~6s so the UI can show "done", then are pruned.
|
||||
pub fn transfers(&self) -> Vec<Transfer> {
|
||||
let now = Instant::now();
|
||||
let mut map = self.transfers.lock().unwrap();
|
||||
map.retain(|_, (t, ts)| !(t.done && now.duration_since(*ts) > Duration::from_secs(6)));
|
||||
map.values().map(|(t, _)| t.clone()).collect()
|
||||
}
|
||||
|
||||
/// Send a file/image to the room. Direct (RTC) only — files are never
|
||||
/// relayed. Chunked + reassembled on the receiver, surfaced via on_file.
|
||||
pub fn send_file(&self, name: String, mime: String, data: Vec<u8>) {
|
||||
|
||||
@@ -51,6 +51,7 @@ pub(crate) struct RtcTransport {
|
||||
ice: StdMutex<Vec<RTCIceServer>>, // refreshed from /api/turn-cred at start
|
||||
direct: DirectSet, // peers with an open data channel
|
||||
present: PresentSet, // peers seen via presence chirps
|
||||
transfers: crate::Transfers, // file-transfer progress
|
||||
}
|
||||
|
||||
fn default_ice() -> Vec<RTCIceServer> {
|
||||
@@ -61,7 +62,12 @@ fn default_ice() -> Vec<RTCIceServer> {
|
||||
}
|
||||
|
||||
impl RtcTransport {
|
||||
pub(crate) fn new(cfg: Config, direct: DirectSet, present: PresentSet) -> Self {
|
||||
pub(crate) fn new(
|
||||
cfg: Config,
|
||||
direct: DirectSet,
|
||||
present: PresentSet,
|
||||
transfers: crate::Transfers,
|
||||
) -> Self {
|
||||
// Data-channel-only: a bare media engine, no codecs/interceptors needed.
|
||||
let api = APIBuilder::new()
|
||||
.with_media_engine(MediaEngine::default())
|
||||
@@ -77,6 +83,7 @@ impl RtcTransport {
|
||||
ice: StdMutex::new(default_ice()),
|
||||
direct,
|
||||
present,
|
||||
transfers,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -264,7 +271,7 @@ impl RtcTransport {
|
||||
};
|
||||
let done = {
|
||||
let mut map = self.incoming.lock().await;
|
||||
crate::apply_file_frame(&mut map, tag, rest)
|
||||
crate::apply_file_frame(&mut map, &self.transfers, tag, rest)
|
||||
};
|
||||
if let Some((name, mime, bytes)) = done {
|
||||
eprintln!("tethercore[rtc]: received file {name:?} ({} bytes)", bytes.len());
|
||||
@@ -583,6 +590,7 @@ impl Transport for RtcTransport {
|
||||
/// The channel is reliable+ordered, so chunks reassemble in arrival order.
|
||||
fn send_file(&self, rt: &Handle, id: String, name: String, mime: String, data: Vec<u8>) {
|
||||
let peers = self.peers.clone();
|
||||
let transfers = self.transfers.clone();
|
||||
rt.spawn(async move {
|
||||
let dcs: Vec<Arc<RTCDataChannel>> = {
|
||||
let map = peers.lock().await;
|
||||
@@ -597,9 +605,12 @@ impl Transport for RtcTransport {
|
||||
if dcs.is_empty() {
|
||||
return;
|
||||
}
|
||||
let total = data.len() as u64;
|
||||
crate::transfer_start(&transfers, &id, &name, total, false);
|
||||
// Each frame is sent as one binary message: tag byte + payload.
|
||||
let frames = crate::file_frames(&id, &name, &mime, &data);
|
||||
for dc in dcs {
|
||||
let mut sent = 0u64;
|
||||
for (tag, payload) in &frames {
|
||||
let mut f = Vec::with_capacity(1 + payload.len());
|
||||
f.push(*tag);
|
||||
@@ -607,8 +618,13 @@ impl Transport for RtcTransport {
|
||||
if dc.send(&Bytes::from(f)).await.is_err() {
|
||||
break;
|
||||
}
|
||||
if *tag == b'C' {
|
||||
sent += (payload.len() - 16) as u64;
|
||||
crate::transfer_progress(&transfers, &id, sent);
|
||||
}
|
||||
}
|
||||
}
|
||||
crate::transfer_done(&transfers, &id, total);
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
@@ -16,6 +16,10 @@ struct ContentView: View {
|
||||
nearbyStrip
|
||||
Divider()
|
||||
}
|
||||
if !store.transfers.isEmpty {
|
||||
transfersRow
|
||||
Divider()
|
||||
}
|
||||
feedArea
|
||||
Divider()
|
||||
if !store.receipts.isEmpty {
|
||||
@@ -42,6 +46,26 @@ struct ContentView: View {
|
||||
.padding(.horizontal, 14).padding(.vertical, 5)
|
||||
}
|
||||
|
||||
// MARK: - File transfers (progress)
|
||||
|
||||
private var transfersRow: some View {
|
||||
VStack(spacing: 4) {
|
||||
ForEach(store.transfers, id: \.name) { t in
|
||||
let frac = t.total > 0 ? Double(t.received) / Double(t.total) : 0
|
||||
HStack(spacing: 8) {
|
||||
Image(systemName: t.incoming ? "arrow.down.circle" : "arrow.up.circle")
|
||||
.foregroundStyle(.secondary)
|
||||
Text(t.name).font(.caption).lineLimit(1).frame(maxWidth: 140, alignment: .leading)
|
||||
ProgressView(value: frac).frame(width: 120)
|
||||
Text(t.done ? "done" : "\(Int(frac * 100))%")
|
||||
.font(.caption2).foregroundStyle(.secondary)
|
||||
Spacer()
|
||||
}
|
||||
}
|
||||
}
|
||||
.padding(.horizontal, 14).padding(.vertical, 5)
|
||||
}
|
||||
|
||||
// MARK: - Presence ("who's here")
|
||||
|
||||
private var presenceRow: some View {
|
||||
|
||||
@@ -31,6 +31,7 @@ final class MacTetherStore {
|
||||
var nearby: [Peer] = [] // tether devices on the LAN, for the picker
|
||||
var present: [Peer] = [] // devices online in the room ("who's here")
|
||||
var receipts: [Receipt] = [] // delivery acks — "seen by <device>"
|
||||
var transfers: [Transfer] = [] // in-flight file transfers
|
||||
var connected = false
|
||||
var lastError: String?
|
||||
var hasAttemptedConnect = false
|
||||
@@ -75,6 +76,7 @@ final class MacTetherStore {
|
||||
self.nearby = engine.nearby().sorted { $0.name < $1.name }
|
||||
self.present = engine.present().sorted { $0.name < $1.name }
|
||||
self.receipts = engine.receipts()
|
||||
self.transfers = engine.transfers()
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -19,6 +19,10 @@ struct RoomView: View {
|
||||
nearbyStrip
|
||||
Divider()
|
||||
}
|
||||
if !store.transfers.isEmpty {
|
||||
transfersRow
|
||||
Divider()
|
||||
}
|
||||
feedList
|
||||
Divider()
|
||||
if !store.receipts.isEmpty {
|
||||
@@ -113,6 +117,23 @@ struct RoomView: View {
|
||||
}
|
||||
}
|
||||
|
||||
private var transfersRow: some View {
|
||||
VStack(spacing: 4) {
|
||||
ForEach(store.transfers, id: \.name) { t in
|
||||
let frac = t.total > 0 ? Double(t.received) / Double(t.total) : 0
|
||||
HStack(spacing: 8) {
|
||||
Image(systemName: t.incoming ? "arrow.down.circle" : "arrow.up.circle")
|
||||
.foregroundStyle(.secondary)
|
||||
Text(t.name).font(.caption).lineLimit(1)
|
||||
ProgressView(value: frac).frame(maxWidth: 120)
|
||||
Text(t.done ? "done" : "\(Int(frac * 100))%")
|
||||
.font(.caption2).foregroundStyle(.secondary)
|
||||
}
|
||||
}
|
||||
}
|
||||
.padding(.horizontal, 14).padding(.vertical, 5)
|
||||
}
|
||||
|
||||
private var seenByBar: some View {
|
||||
let names = Array(Set(store.receipts.map(\.name))).sorted()
|
||||
return HStack(spacing: 6) {
|
||||
|
||||
@@ -36,6 +36,7 @@ final class TetherStore {
|
||||
var nearby: [Peer] = [] // tether devices on the LAN (empty on iOS w/o multicast entitlement)
|
||||
var present: [Peer] = [] // devices online in the room ("who's here")
|
||||
var receipts: [Receipt] = [] // delivery acks — "seen by <device>"
|
||||
var transfers: [Transfer] = [] // in-flight file transfers
|
||||
var connected = false
|
||||
var lastError: String?
|
||||
|
||||
@@ -74,6 +75,7 @@ final class TetherStore {
|
||||
self.nearby = engine.nearby().sorted { $0.name < $1.name }
|
||||
self.present = engine.present().sorted { $0.name < $1.name }
|
||||
self.receipts = engine.receipts()
|
||||
self.transfers = engine.transfers()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user