taquba_workflow/bulk/
progress.rs1use std::time::{Duration, Instant};
4
5use serde::{Deserialize, Serialize};
6
7use crate::bulk::cost::CostReport;
8
9#[derive(Debug, Clone, Serialize, Deserialize)]
15pub(crate) struct ItemMarker {
16 pub(crate) status: MarkerStatus,
17 pub(crate) error: Option<String>,
18 pub(crate) cost: CostReport,
19}
20
21#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
22pub(crate) enum MarkerStatus {
23 Succeeded,
24 Failed,
25}
26
27impl ItemMarker {
28 pub(crate) fn new(status: MarkerStatus, error: Option<String>, cost: CostReport) -> Self {
29 Self {
30 status,
31 error,
32 cost,
33 }
34 }
35
36 pub(crate) fn encode(&self) -> Result<Vec<u8>, crate::Error> {
37 Ok(rmp_serde::to_vec_named(self)?)
38 }
39}
40
41#[derive(Debug, Clone)]
46pub struct BatchStatus {
47 pub batch_id: String,
49 pub total: usize,
51 pub succeeded: usize,
53 pub failed: usize,
55 pub cost: CostReport,
57 pub failed_keys: Vec<String>,
59}
60
61#[derive(Debug, Clone)]
65pub struct ProgressSnapshot {
66 pub total: usize,
68 pub completed: usize,
70 pub succeeded: usize,
72 pub failed: usize,
74 pub cancelled: usize,
76 pub elapsed: Duration,
78 pub rate_per_sec: f64,
80 pub time_remaining: Option<Duration>,
83 pub cost: CostReport,
85}
86
87#[derive(Debug, Clone)]
90pub struct BulkReport {
91 pub batch_id: String,
93 pub total: usize,
95 pub succeeded: usize,
97 pub failed: usize,
99 pub cancelled: usize,
101 pub elapsed: Duration,
103 pub cost: CostReport,
105 pub failed_keys: Vec<String>,
108}
109
110#[derive(Debug)]
113pub(crate) struct ProgressState {
114 pub total: usize,
115 pub succeeded: usize,
116 pub failed: usize,
117 pub cancelled: usize,
118 pub cost: CostReport,
119 pub failed_keys: Vec<String>,
120 started_at: Instant,
121}
122
123impl ProgressState {
124 pub(crate) fn new(total: usize) -> Self {
125 Self {
126 total,
127 succeeded: 0,
128 failed: 0,
129 cancelled: 0,
130 cost: CostReport::new(),
131 failed_keys: Vec::new(),
132 started_at: Instant::now(),
133 }
134 }
135
136 pub(crate) fn completed(&self) -> usize {
137 self.succeeded + self.failed + self.cancelled
138 }
139
140 pub(crate) fn snapshot(&self) -> ProgressSnapshot {
141 let elapsed = self.started_at.elapsed();
142 let completed = self.completed();
143 let secs = elapsed.as_secs_f64();
144 let rate_per_sec = if secs > 0.0 {
145 completed as f64 / secs
146 } else {
147 0.0
148 };
149 let remaining = self.total.saturating_sub(completed);
150 let time_remaining = if rate_per_sec > 0.0 && remaining > 0 {
151 Some(Duration::from_secs_f64(remaining as f64 / rate_per_sec))
152 } else {
153 None
154 };
155 ProgressSnapshot {
156 total: self.total,
157 completed,
158 succeeded: self.succeeded,
159 failed: self.failed,
160 cancelled: self.cancelled,
161 elapsed,
162 rate_per_sec,
163 time_remaining,
164 cost: self.cost.clone(),
165 }
166 }
167
168 pub(crate) fn to_report(&self, batch_id: &str) -> BulkReport {
169 BulkReport {
170 batch_id: batch_id.to_string(),
171 total: self.total,
172 succeeded: self.succeeded,
173 failed: self.failed,
174 cancelled: self.cancelled,
175 elapsed: self.started_at.elapsed(),
176 cost: self.cost.clone(),
177 failed_keys: self.failed_keys.clone(),
178 }
179 }
180}
181
182#[cfg(test)]
183mod tests {
184 use super::*;
185
186 #[test]
187 fn completed_sums_terminal_buckets() {
188 let mut st = ProgressState::new(6);
189 st.succeeded = 3;
190 st.failed = 2;
191 st.cancelled = 1;
192 assert_eq!(st.completed(), 6);
193 }
194
195 #[test]
196 fn snapshot_reports_no_time_remaining_before_progress() {
197 let st = ProgressState::new(10);
198 let snap = st.snapshot();
199 assert_eq!(snap.total, 10);
200 assert_eq!(snap.completed, 0);
201 assert!(snap.time_remaining.is_none());
202 }
203}