Skip to main content

taquba_workflow/bulk/
progress.rs

1//! Batch-level progress and the final run report.
2
3use std::time::{Duration, Instant};
4
5use serde::{Deserialize, Serialize};
6
7use crate::bulk::cost::CostReport;
8
9/// The durable marker of one terminated item, stored under
10/// `workflow/bulk/batches/{batch_id}/items/{key}` in the settlement that
11/// commits the item's terminal outcome. Only a success and a failure
12/// settle through the item's own step, so those are the recorded
13/// outcomes; a cancelled item leaves no marker.
14#[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/// The durable state of a batch, read from its manifest and item markers
42/// by [`Batch::status`](crate::bulk::Batch::status). An item with no
43/// recorded outcome is neither succeeded nor failed: it has not run to
44/// a settlement of its own, or it was cancelled.
45#[derive(Debug, Clone)]
46pub struct BatchStatus {
47    /// The batch id.
48    pub batch_id: String,
49    /// Number of items in the batch's manifest.
50    pub total: usize,
51    /// Items whose last recorded outcome is a success.
52    pub succeeded: usize,
53    /// Items whose last recorded outcome is a failure.
54    pub failed: usize,
55    /// Cost counters rolled up across the recorded items.
56    pub cost: CostReport,
57    /// Keys of the items whose last recorded outcome is a failure.
58    pub failed_keys: Vec<String>,
59}
60
61/// A point-in-time view of a batch run's progress. Returned by
62/// [`Batch::progress`](crate::bulk::Batch::progress) and suitable for a
63/// status line or a polling UI.
64#[derive(Debug, Clone)]
65pub struct ProgressSnapshot {
66    /// Number of items in the batch.
67    pub total: usize,
68    /// Items that have reached any terminal state.
69    pub completed: usize,
70    /// Items that terminated successfully.
71    pub succeeded: usize,
72    /// Items that terminated failed.
73    pub failed: usize,
74    /// Items that terminated cancelled.
75    pub cancelled: usize,
76    /// Wall-clock time since the run started.
77    pub elapsed: Duration,
78    /// Completed items per second over the elapsed window.
79    pub rate_per_sec: f64,
80    /// Estimated time to finish the remaining items at the current rate, or
81    /// `None` when the total is unknown or the rate is zero.
82    pub time_remaining: Option<Duration>,
83    /// Cost counters rolled up across completed items.
84    pub cost: CostReport,
85}
86
87/// The outcome of a finished (or drained) bulk run, returned by
88/// [`Batch::run`](crate::bulk::Batch::run).
89#[derive(Debug, Clone)]
90pub struct BulkReport {
91    /// The batch this report describes.
92    pub batch_id: String,
93    /// Number of items that were expected to complete.
94    pub total: usize,
95    /// Items that terminated successfully.
96    pub succeeded: usize,
97    /// Items that terminated failed.
98    pub failed: usize,
99    /// Items that terminated cancelled.
100    pub cancelled: usize,
101    /// Wall-clock duration of the run.
102    pub elapsed: Duration,
103    /// Cost counters rolled up across all completed items.
104    pub cost: CostReport,
105    /// Keys of the items that failed. A later run of the same batch runs
106    /// them again and skips the items that succeeded.
107    pub failed_keys: Vec<String>,
108}
109
110/// Internal, mutex-guarded counters of a batch being run, updated as its
111/// items terminate and read by [`ProgressSnapshot`] / [`BulkReport`].
112#[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}