Skip to content

Commit 13645b7

Browse files
committed
fix: expose transfer bytes in ItemProgressReport for smooth download progress
ItemBridgeState::compute_diff hardcoded all total_transfer_bytes fields to 0, suppressing the fine-grained network-level progress that xet-core already tracks per HTTP chunk. This caused Python callbacks to only receive coarse bytes_completed updates (per disk write, ~256MB batches). - Add transfer_bytes and transfer_bytes_completed to ItemProgressReport - Update ItemBridgeState::compute_diff to diff transfer fields (matching GroupBridgeState which already does this correctly) - Fire callbacks when transfer progress changes, even if bytes_completed has not changed yet
1 parent 20198a9 commit 13645b7

2 files changed

Lines changed: 112 additions & 6 deletions

File tree

xet_data/src/progress_tracking/progress_types.rs

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -48,6 +48,8 @@ impl ItemProgress {
4848
item_name: self.name.to_string(),
4949
total_bytes,
5050
bytes_completed,
51+
transfer_bytes,
52+
transfer_bytes_completed,
5153
}
5254
}
5355
}
@@ -380,6 +382,8 @@ pub struct ItemProgressReport {
380382
pub item_name: String,
381383
pub total_bytes: u64,
382384
pub bytes_completed: u64,
385+
pub transfer_bytes: u64,
386+
pub transfer_bytes_completed: u64,
383387
}
384388

385389
#[cfg(test)]

xet_pkg/src/legacy/progress_tracking/callback_bridge.rs

Lines changed: 108 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -75,9 +75,11 @@ impl GroupBridgeState {
7575
for (&id, report) in &items {
7676
let prev = self.prev_items.get(&id);
7777
let prev_completed = prev.map_or(0, |p| p.bytes_completed);
78+
let prev_transfer_completed = prev.map_or(0, |p| p.transfer_bytes_completed);
7879
let increment = report.bytes_completed.saturating_sub(prev_completed);
80+
let transfer_increment = report.transfer_bytes_completed.saturating_sub(prev_transfer_completed);
7981

80-
if increment > 0 || prev.is_none() {
82+
if increment > 0 || transfer_increment > 0 || prev.is_none() {
8183
item_updates.push(ItemProgressUpdate {
8284
tracking_id: id,
8385
item_name: Arc::from(report.item_name.as_str()),
@@ -123,11 +125,17 @@ impl ItemBridgeState {
123125
fn compute_diff(&mut self, item_id: UniqueID, report: ItemProgressReport) -> ProgressUpdate {
124126
let prev_completed = self.prev.as_ref().map_or(0, |p| p.bytes_completed);
125127
let prev_total = self.prev.as_ref().map_or(0, |p| p.total_bytes);
128+
let prev_transfer_bytes = self.prev.as_ref().map_or(0, |p| p.transfer_bytes);
129+
let prev_transfer_completed = self.prev.as_ref().map_or(0, |p| p.transfer_bytes_completed);
126130

127131
let bytes_increment = report.bytes_completed.saturating_sub(prev_completed);
128132
let total_increment = report.total_bytes.saturating_sub(prev_total);
133+
let transfer_bytes_increment = report.transfer_bytes.saturating_sub(prev_transfer_bytes);
134+
let transfer_completion_increment = report.transfer_bytes_completed.saturating_sub(prev_transfer_completed);
129135

130-
let item_updates = if bytes_increment > 0 || self.prev.is_none() {
136+
let has_progress = bytes_increment > 0 || transfer_completion_increment > 0 || self.prev.is_none();
137+
138+
let item_updates = if has_progress {
131139
vec![ItemProgressUpdate {
132140
tracking_id: item_id,
133141
item_name: Arc::from(report.item_name.as_str()),
@@ -146,10 +154,10 @@ impl ItemBridgeState {
146154
total_bytes_completed: report.bytes_completed,
147155
total_bytes_completion_increment: bytes_increment,
148156
total_bytes_completion_rate: None,
149-
total_transfer_bytes: 0,
150-
total_transfer_bytes_increment: 0,
151-
total_transfer_bytes_completed: 0,
152-
total_transfer_bytes_completion_increment: 0,
157+
total_transfer_bytes: report.transfer_bytes,
158+
total_transfer_bytes_increment: transfer_bytes_increment,
159+
total_transfer_bytes_completed: report.transfer_bytes_completed,
160+
total_transfer_bytes_completion_increment: transfer_completion_increment,
153161
total_transfer_bytes_completion_rate: None,
154162
};
155163

@@ -347,6 +355,24 @@ mod tests {
347355
item_name: name.to_string(),
348356
total_bytes,
349357
bytes_completed,
358+
transfer_bytes: 0,
359+
transfer_bytes_completed: 0,
360+
}
361+
}
362+
363+
fn make_item_report_with_transfer(
364+
name: &str,
365+
total_bytes: u64,
366+
bytes_completed: u64,
367+
transfer_bytes: u64,
368+
transfer_bytes_completed: u64,
369+
) -> ItemProgressReport {
370+
ItemProgressReport {
371+
item_name: name.to_string(),
372+
total_bytes,
373+
bytes_completed,
374+
transfer_bytes,
375+
transfer_bytes_completed,
350376
}
351377
}
352378

@@ -432,6 +458,27 @@ mod tests {
432458
assert_eq!(update.item_updates[0].bytes_completion_increment, 0);
433459
}
434460

461+
#[test]
462+
fn test_group_bridge_transfer_progress_includes_item() {
463+
let mut state = GroupBridgeState::new();
464+
let id = UniqueID::new();
465+
466+
// Initial state
467+
let group1 = make_group_report(1000, 0, 800, 0);
468+
let items1 = HashMap::from([(id, make_item_report_with_transfer("a.bin", 1000, 0, 800, 0))]);
469+
state.compute_diff(group1, items1);
470+
471+
// Transfer progress only (bytes_completed unchanged)
472+
let group2 = make_group_report(1000, 0, 800, 200);
473+
let items2 = HashMap::from([(id, make_item_report_with_transfer("a.bin", 1000, 0, 800, 200))]);
474+
let update = state.compute_diff(group2, items2);
475+
476+
assert_eq!(update.total_transfer_bytes_completion_increment, 200);
477+
assert_eq!(update.total_bytes_completion_increment, 0);
478+
// Item should be included despite bytes_completed not changing
479+
assert_eq!(update.item_updates.len(), 1);
480+
}
481+
435482
#[test]
436483
fn test_item_bridge_first_diff() {
437484
let mut state = ItemBridgeState::new();
@@ -486,4 +533,59 @@ mod tests {
486533
assert_eq!(update.total_bytes_completion_increment, 0);
487534
assert!(update.item_updates.is_empty());
488535
}
536+
537+
#[test]
538+
fn test_item_bridge_transfer_bytes_reported() {
539+
let mut state = ItemBridgeState::new();
540+
let id = UniqueID::new();
541+
542+
let report = make_item_report_with_transfer("file.bin", 1000, 0, 800, 200);
543+
let update = state.compute_diff(id, report);
544+
545+
assert_eq!(update.total_transfer_bytes, 800);
546+
assert_eq!(update.total_transfer_bytes_completed, 200);
547+
assert_eq!(update.total_transfer_bytes_completion_increment, 200);
548+
}
549+
550+
#[test]
551+
fn test_item_bridge_transfer_bytes_incremental() {
552+
let mut state = ItemBridgeState::new();
553+
let id = UniqueID::new();
554+
555+
state.compute_diff(id, make_item_report_with_transfer("file.bin", 1000, 0, 800, 200));
556+
557+
let update = state.compute_diff(id, make_item_report_with_transfer("file.bin", 1000, 0, 800, 500));
558+
559+
assert_eq!(update.total_transfer_bytes_increment, 0);
560+
assert_eq!(update.total_transfer_bytes_completion_increment, 300);
561+
}
562+
563+
#[test]
564+
fn test_item_bridge_transfer_progress_triggers_callback() {
565+
let mut state = ItemBridgeState::new();
566+
let id = UniqueID::new();
567+
568+
// Initial state: no bytes completed, no transfer
569+
state.compute_diff(id, make_item_report_with_transfer("file.bin", 1000, 0, 800, 0));
570+
571+
// Transfer progress without bytes_completed change should still produce an update
572+
let update = state.compute_diff(id, make_item_report_with_transfer("file.bin", 1000, 0, 800, 100));
573+
574+
assert!(!update.is_empty());
575+
assert_eq!(update.total_transfer_bytes_completion_increment, 100);
576+
assert_eq!(update.total_bytes_completion_increment, 0);
577+
}
578+
579+
#[test]
580+
fn test_item_bridge_no_transfer_no_bytes_is_empty() {
581+
let mut state = ItemBridgeState::new();
582+
let id = UniqueID::new();
583+
584+
state.compute_diff(id, make_item_report_with_transfer("file.bin", 1000, 100, 800, 200));
585+
586+
// Same values again: should be empty
587+
let update = state.compute_diff(id, make_item_report_with_transfer("file.bin", 1000, 100, 800, 200));
588+
589+
assert!(update.is_empty());
590+
}
489591
}

0 commit comments

Comments
 (0)