use std::time::{Duration, Instant};
use serde::{Deserialize, Serialize};
use crate::bulk::cost::CostReport;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub(crate) struct ItemMarker {
pub(crate) status: MarkerStatus,
pub(crate) error: Option<String>,
pub(crate) cost: CostReport,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub(crate) enum MarkerStatus {
Succeeded,
Failed,
}
impl ItemMarker {
pub(crate) fn new(status: MarkerStatus, error: Option<String>, cost: CostReport) -> Self {
Self {
status,
error,
cost,
}
}
pub(crate) fn encode(&self) -> Result<Vec<u8>, crate::Error> {
Ok(rmp_serde::to_vec_named(self)?)
}
}
#[derive(Debug, Clone)]
pub struct BatchStatus {
pub batch_id: String,
pub total: usize,
pub succeeded: usize,
pub failed: usize,
pub cost: CostReport,
pub failed_keys: Vec<String>,
}
#[derive(Debug, Clone)]
pub struct ProgressSnapshot {
pub total: usize,
pub completed: usize,
pub succeeded: usize,
pub failed: usize,
pub cancelled: usize,
pub elapsed: Duration,
pub rate_per_sec: f64,
pub time_remaining: Option<Duration>,
pub cost: CostReport,
}
#[derive(Debug, Clone)]
pub struct BulkReport {
pub batch_id: String,
pub total: usize,
pub succeeded: usize,
pub failed: usize,
pub cancelled: usize,
pub elapsed: Duration,
pub cost: CostReport,
pub failed_keys: Vec<String>,
}
#[derive(Debug)]
pub(crate) struct ProgressState {
pub total: usize,
pub succeeded: usize,
pub failed: usize,
pub cancelled: usize,
pub cost: CostReport,
pub failed_keys: Vec<String>,
started_at: Instant,
}
impl ProgressState {
pub(crate) fn new(total: usize) -> Self {
Self {
total,
succeeded: 0,
failed: 0,
cancelled: 0,
cost: CostReport::new(),
failed_keys: Vec::new(),
started_at: Instant::now(),
}
}
pub(crate) fn completed(&self) -> usize {
self.succeeded + self.failed + self.cancelled
}
pub(crate) fn snapshot(&self) -> ProgressSnapshot {
let elapsed = self.started_at.elapsed();
let completed = self.completed();
let secs = elapsed.as_secs_f64();
let rate_per_sec = if secs > 0.0 {
completed as f64 / secs
} else {
0.0
};
let remaining = self.total.saturating_sub(completed);
let time_remaining = if rate_per_sec > 0.0 && remaining > 0 {
Some(Duration::from_secs_f64(remaining as f64 / rate_per_sec))
} else {
None
};
ProgressSnapshot {
total: self.total,
completed,
succeeded: self.succeeded,
failed: self.failed,
cancelled: self.cancelled,
elapsed,
rate_per_sec,
time_remaining,
cost: self.cost.clone(),
}
}
pub(crate) fn to_report(&self, batch_id: &str) -> BulkReport {
BulkReport {
batch_id: batch_id.to_string(),
total: self.total,
succeeded: self.succeeded,
failed: self.failed,
cancelled: self.cancelled,
elapsed: self.started_at.elapsed(),
cost: self.cost.clone(),
failed_keys: self.failed_keys.clone(),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn completed_sums_terminal_buckets() {
let mut st = ProgressState::new(6);
st.succeeded = 3;
st.failed = 2;
st.cancelled = 1;
assert_eq!(st.completed(), 6);
}
#[test]
fn snapshot_reports_no_time_remaining_before_progress() {
let st = ProgressState::new(10);
let snap = st.snapshot();
assert_eq!(snap.total, 10);
assert_eq!(snap.completed, 0);
assert!(snap.time_remaining.is_none());
}
}