Skip to main content

ic_testkit/artifacts/
transaction_batch.rs

1use crate::batch::{BatchLabelError, validate_labels};
2
3use std::time::{Duration, Instant};
4
5use super::transaction::{
6    ArtifactBuildTransaction, ArtifactCacheError, ArtifactCacheOutcome, ArtifactCachePreparation,
7    ArtifactCacheSpec, ArtifactCacheTimings, prepare_artifact_cache,
8};
9
10/// Caller-labeled specification for one generic artifact batch entry.
11///
12/// The label is report and callback identity, not artifact-cache identity. The
13/// caller must keep it stable anywhere reports are composed across stages.
14#[derive(Clone, Debug, Eq, PartialEq)]
15pub struct LabeledArtifactCacheSpec {
16    label: String,
17    spec: ArtifactCacheSpec,
18}
19
20/// Ordered labeled entries from a collect-all artifact transaction batch.
21#[derive(Debug)]
22pub struct ArtifactCacheBatchReport<E> {
23    entries: Vec<ArtifactCacheBatchEntry<E>>,
24    total: Duration,
25}
26
27/// One ordered labeled result from a generic artifact batch.
28#[derive(Debug)]
29pub struct ArtifactCacheBatchEntry<E> {
30    index: usize,
31    label: String,
32    result: Result<ArtifactCacheOutcome, ArtifactCacheBatchFailure<E>>,
33    entry_elapsed: Duration,
34}
35
36/// One successful generic artifact batch entry.
37#[derive(Clone, Copy, Debug)]
38pub struct ArtifactCacheBatchOutcomeEntry<'a> {
39    index: usize,
40    label: &'a str,
41    outcome: &'a ArtifactCacheOutcome,
42    entry_elapsed: Duration,
43}
44
45/// One failed generic artifact batch entry.
46#[derive(Debug)]
47pub struct ArtifactCacheBatchFailedEntry<'a, E> {
48    index: usize,
49    label: &'a str,
50    failure: &'a ArtifactCacheBatchFailure<E>,
51    entry_elapsed: Duration,
52}
53
54/// Structural error that prevents a labeled artifact batch from starting.
55#[non_exhaustive]
56#[derive(Clone, Debug, Eq, PartialEq)]
57pub enum ArtifactCacheBatchContractError {
58    /// An entry label was empty.
59    EmptyLabel {
60        /// Zero-based position of the invalid entry.
61        index: usize,
62    },
63    /// Two entries used the same label.
64    DuplicateLabel {
65        /// Duplicated caller label.
66        label: String,
67        /// Position where the label first appeared.
68        first_index: usize,
69        /// Position where the label was repeated.
70        duplicate_index: usize,
71    },
72}
73
74/// Primary phase in which a generic artifact batch entry failed.
75#[derive(Clone, Copy, Debug, Eq, PartialEq)]
76pub enum ArtifactCacheBatchFailurePhase {
77    /// Cache preparation failed before the population callback ran.
78    Preparation,
79    /// The caller's population callback failed.
80    Callback,
81    /// Transaction commit failed after successful population.
82    Commit,
83}
84
85/// Partial phase timings retained for one failed artifact batch entry.
86#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
87pub struct ArtifactCacheBatchFailureTimings {
88    preparation: Duration,
89    callback: Option<Duration>,
90    cleanup: Option<Duration>,
91    commit: Option<Duration>,
92    total: Duration,
93}
94
95/// Failure from one entry in a collect-all artifact transaction batch.
96#[derive(Debug)]
97pub enum ArtifactCacheBatchFailure<E> {
98    /// Cache preparation or commit failed.
99    Cache {
100        /// Failed cache phase.
101        phase: ArtifactCacheBatchFailurePhase,
102        /// Cache preparation or commit failure.
103        source: Box<ArtifactCacheError>,
104        /// Phase timings completed before the failure returned.
105        timings: ArtifactCacheBatchFailureTimings,
106    },
107    /// The caller's population callback failed.
108    Build {
109        /// Caller population failure.
110        source: Box<E>,
111        /// Failure from synchronously aborting the active transaction.
112        cleanup_error: Option<Box<ArtifactCacheError>>,
113        /// Preparation, callback, and explicit cleanup timings.
114        timings: ArtifactCacheBatchFailureTimings,
115    },
116}
117
118/// Aggregate counters and successful-acquisition timings for an artifact batch.
119#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
120pub struct ArtifactCacheBatchMetrics {
121    entries: usize,
122    succeeded: usize,
123    failed: usize,
124    built: usize,
125    reused: usize,
126    successful_timings: ArtifactCacheTimings,
127    total: Duration,
128}
129
130impl LabeledArtifactCacheSpec {
131    /// Attach a caller-owned stable label to one artifact specification.
132    #[must_use]
133    pub fn new(label: impl Into<String>, spec: ArtifactCacheSpec) -> Self {
134        Self {
135            label: label.into(),
136            spec,
137        }
138    }
139
140    /// Caller-owned report and callback label.
141    #[must_use]
142    pub fn label(&self) -> &str {
143        &self.label
144    }
145
146    /// Underlying exact artifact specification.
147    #[must_use]
148    pub const fn spec(&self) -> &ArtifactCacheSpec {
149        &self.spec
150    }
151
152    /// Consume the entry into its label and artifact specification.
153    #[must_use]
154    pub fn into_parts(self) -> (String, ArtifactCacheSpec) {
155        (self.label, self.spec)
156    }
157}
158
159impl<E> ArtifactCacheBatchReport<E> {
160    /// Ordered labeled entries.
161    #[must_use]
162    pub fn entries(&self) -> &[ArtifactCacheBatchEntry<E>] {
163        &self.entries
164    }
165
166    /// Consume the report into its ordered labeled entries.
167    #[must_use]
168    pub fn into_entries(self) -> Vec<ArtifactCacheBatchEntry<E>> {
169        self.entries
170    }
171
172    /// Structured successful entries with labels and wall-clock times.
173    pub fn outcomes(&self) -> impl Iterator<Item = ArtifactCacheBatchOutcomeEntry<'_>> {
174        self.entries.iter().filter_map(|entry| {
175            entry
176                .outcome()
177                .map(|outcome| ArtifactCacheBatchOutcomeEntry {
178                    index: entry.index,
179                    label: &entry.label,
180                    outcome,
181                    entry_elapsed: entry.entry_elapsed,
182                })
183        })
184    }
185
186    /// Structured failed entries with labels and partial phase timings.
187    pub fn failures(&self) -> impl Iterator<Item = ArtifactCacheBatchFailedEntry<'_, E>> {
188        self.entries.iter().filter_map(|entry| {
189            entry
190                .failure()
191                .map(|failure| ArtifactCacheBatchFailedEntry {
192                    index: entry.index,
193                    label: &entry.label,
194                    failure,
195                    entry_elapsed: entry.entry_elapsed,
196                })
197        })
198    }
199
200    /// Complete wall-clock time for the sequential collect-all batch.
201    #[must_use]
202    pub const fn total(&self) -> Duration {
203        self.total
204    }
205
206    /// Whether every specification completed successfully.
207    #[must_use]
208    pub fn is_success(&self) -> bool {
209        self.entries.iter().all(ArtifactCacheBatchEntry::is_success)
210    }
211
212    /// Aggregate outcome counters and successful-acquisition timings.
213    #[must_use]
214    pub fn metrics(&self) -> ArtifactCacheBatchMetrics {
215        let mut metrics = ArtifactCacheBatchMetrics {
216            entries: self.entries.len(),
217            total: self.total,
218            ..ArtifactCacheBatchMetrics::default()
219        };
220        for entry in &self.entries {
221            match &entry.result {
222                Ok(outcome) => {
223                    metrics.succeeded += 1;
224                    if outcome.is_reused() {
225                        metrics.reused += 1;
226                    } else {
227                        metrics.built += 1;
228                    }
229                    metrics.successful_timings = metrics
230                        .successful_timings
231                        .saturating_add(outcome.record().timings());
232                }
233                Err(_) => metrics.failed += 1,
234            }
235        }
236        metrics
237    }
238}
239
240impl<E> ArtifactCacheBatchEntry<E> {
241    /// Zero-based position in the supplied specification slice.
242    #[must_use]
243    pub const fn index(&self) -> usize {
244        self.index
245    }
246
247    /// Caller-owned stable label.
248    #[must_use]
249    pub fn label(&self) -> &str {
250        &self.label
251    }
252
253    /// Structured success or failure for this entry.
254    pub const fn result(&self) -> Result<&ArtifactCacheOutcome, &ArtifactCacheBatchFailure<E>> {
255        self.result.as_ref()
256    }
257
258    /// Successful artifact outcome, when this entry succeeded.
259    #[must_use]
260    pub fn outcome(&self) -> Option<&ArtifactCacheOutcome> {
261        self.result.as_ref().ok()
262    }
263
264    /// Structured batch failure, when this entry failed.
265    #[must_use]
266    pub fn failure(&self) -> Option<&ArtifactCacheBatchFailure<E>> {
267        self.result.as_ref().err()
268    }
269
270    /// Complete wall-clock time retained for this entry.
271    #[must_use]
272    pub const fn entry_elapsed(&self) -> Duration {
273        self.entry_elapsed
274    }
275
276    /// Whether this entry completed successfully.
277    #[must_use]
278    pub const fn is_success(&self) -> bool {
279        self.result.is_ok()
280    }
281
282    /// Consume the entry into its ordered identity, result, and wall time.
283    pub fn into_parts(
284        self,
285    ) -> (
286        usize,
287        String,
288        Result<ArtifactCacheOutcome, ArtifactCacheBatchFailure<E>>,
289        Duration,
290    ) {
291        (self.index, self.label, self.result, self.entry_elapsed)
292    }
293}
294
295impl<'a> ArtifactCacheBatchOutcomeEntry<'a> {
296    /// Zero-based position in the supplied specification slice.
297    #[must_use]
298    pub const fn index(self) -> usize {
299        self.index
300    }
301
302    /// Caller-owned stable label.
303    #[must_use]
304    pub const fn label(self) -> &'a str {
305        self.label
306    }
307
308    /// Successful artifact outcome.
309    #[must_use]
310    pub const fn outcome(self) -> &'a ArtifactCacheOutcome {
311        self.outcome
312    }
313
314    /// Complete wall-clock time retained for this successful entry.
315    #[must_use]
316    pub const fn entry_elapsed(self) -> Duration {
317        self.entry_elapsed
318    }
319}
320
321impl<'a, E> ArtifactCacheBatchFailedEntry<'a, E> {
322    /// Zero-based position in the supplied specification slice.
323    #[must_use]
324    pub const fn index(&self) -> usize {
325        self.index
326    }
327
328    /// Caller-owned stable label.
329    #[must_use]
330    pub const fn label(&self) -> &'a str {
331        self.label
332    }
333
334    /// Structured cache or caller-build failure.
335    #[must_use]
336    pub const fn failure(&self) -> &'a ArtifactCacheBatchFailure<E> {
337        self.failure
338    }
339
340    /// Partial phase timings retained with the failure.
341    #[must_use]
342    pub const fn timings(&self) -> ArtifactCacheBatchFailureTimings {
343        self.failure.timings()
344    }
345
346    /// Complete wall-clock time retained for this failed entry.
347    #[must_use]
348    pub const fn entry_elapsed(&self) -> Duration {
349        self.entry_elapsed
350    }
351}
352
353impl<E> ArtifactCacheBatchFailure<E> {
354    /// Primary phase in which this entry failed.
355    #[must_use]
356    pub const fn phase(&self) -> ArtifactCacheBatchFailurePhase {
357        match self {
358            Self::Cache { phase, .. } => *phase,
359            Self::Build { .. } => ArtifactCacheBatchFailurePhase::Callback,
360        }
361    }
362
363    /// Partial phase timings completed before the failure returned.
364    #[must_use]
365    pub const fn timings(&self) -> ArtifactCacheBatchFailureTimings {
366        match self {
367            Self::Cache { timings, .. } | Self::Build { timings, .. } => *timings,
368        }
369    }
370
371    /// Cleanup failure after a caller population error, when one occurred.
372    #[must_use]
373    pub fn cleanup_error(&self) -> Option<&ArtifactCacheError> {
374        match self {
375            Self::Build { cleanup_error, .. } => cleanup_error.as_deref(),
376            Self::Cache { .. } => None,
377        }
378    }
379}
380
381impl ArtifactCacheBatchFailureTimings {
382    /// Time spent in cache preparation before the failure path continued.
383    #[must_use]
384    pub const fn preparation(self) -> Duration {
385        self.preparation
386    }
387
388    /// Time spent in the population callback, when it ran.
389    #[must_use]
390    pub const fn callback(self) -> Option<Duration> {
391        self.callback
392    }
393
394    /// Time spent explicitly aborting after a callback failure, when attempted.
395    #[must_use]
396    pub const fn cleanup(self) -> Option<Duration> {
397        self.cleanup
398    }
399
400    /// Time spent committing, including commit-owned failure cleanup, when attempted.
401    #[must_use]
402    pub const fn commit(self) -> Option<Duration> {
403        self.commit
404    }
405
406    /// Complete wall-clock time retained for the failed entry.
407    #[must_use]
408    pub const fn total(self) -> Duration {
409        self.total
410    }
411}
412
413impl ArtifactCacheBatchMetrics {
414    /// Number of supplied specifications.
415    #[must_use]
416    pub const fn entries(self) -> usize {
417        self.entries
418    }
419
420    /// Number of successful specifications.
421    #[must_use]
422    pub const fn succeeded(self) -> usize {
423        self.succeeded
424    }
425
426    /// Number of failed specifications.
427    #[must_use]
428    pub const fn failed(self) -> usize {
429        self.failed
430    }
431
432    /// Number of newly built artifact sets.
433    #[must_use]
434    pub const fn built(self) -> usize {
435        self.built
436    }
437
438    /// Number of artifact sets reused from the exact cache.
439    #[must_use]
440    pub const fn reused(self) -> usize {
441        self.reused
442    }
443
444    /// Sum of timings from successful acquisitions.
445    #[must_use]
446    pub const fn successful_timings(self) -> ArtifactCacheTimings {
447        self.successful_timings
448    }
449
450    /// Complete wall-clock time for the sequential batch.
451    #[must_use]
452    pub const fn total(self) -> Duration {
453        self.total
454    }
455}
456
457impl<E> std::fmt::Display for ArtifactCacheBatchReport<E> {
458    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
459        let metrics = self.metrics();
460        write!(
461            formatter,
462            "entries={} succeeded={} failed={} built={} reused={} successful_timings=({}) total={:?}",
463            metrics.entries(),
464            metrics.succeeded(),
465            metrics.failed(),
466            metrics.built(),
467            metrics.reused(),
468            metrics.successful_timings(),
469            metrics.total(),
470        )
471    }
472}
473
474/// Build every independent caller-labeled artifact specification.
475///
476/// Labels must be nonempty and unique; the complete batch is rejected before
477/// work begins otherwise. The callback runs only for misses and receives the
478/// stable label plus one live transaction at a time. Sequential acquisition
479/// avoids self-deadlock when entries share a coordination scope. Every
480/// independent result is retained, and a callback error aborts and
481/// synchronously cleans that transaction before the next entry begins.
482///
483/// This operation is not atomic across specifications. A caller requiring one
484/// all-or-nothing multi-output publication should declare those outputs on one
485/// [`ArtifactCacheSpec`] instead.
486pub fn build_artifact_caches_batch<E, F>(
487    specs: &[LabeledArtifactCacheSpec],
488    mut populate: F,
489) -> Result<ArtifactCacheBatchReport<E>, ArtifactCacheBatchContractError>
490where
491    F: FnMut(&str, &ArtifactBuildTransaction) -> Result<(), E>,
492{
493    validate_batch_labels(specs)?;
494    let started = Instant::now();
495    let mut entries = Vec::with_capacity(specs.len());
496    for (index, labeled) in specs.iter().enumerate() {
497        let entry_started = Instant::now();
498        let preparation_started = Instant::now();
499        let preparation = prepare_artifact_cache(&labeled.spec);
500        let preparation_elapsed = preparation_started.elapsed();
501        let (result, entry_elapsed) = match preparation {
502            Ok(ArtifactCachePreparation::Reused(record)) => (
503                Ok(ArtifactCacheOutcome::Reused(record)),
504                entry_started.elapsed(),
505            ),
506            Ok(ArtifactCachePreparation::Build(transaction)) => {
507                let callback_started = Instant::now();
508                let callback_result = populate(&labeled.label, &transaction);
509                let callback_elapsed = callback_started.elapsed();
510                if let Err(source) = callback_result {
511                    let cleanup_started = Instant::now();
512                    let cleanup_error = transaction.abort().err().map(Box::new);
513                    let cleanup_elapsed = cleanup_started.elapsed();
514                    let entry_elapsed = entry_started.elapsed();
515                    (
516                        Err(ArtifactCacheBatchFailure::Build {
517                            source: Box::new(source),
518                            cleanup_error,
519                            timings: ArtifactCacheBatchFailureTimings {
520                                preparation: preparation_elapsed,
521                                callback: Some(callback_elapsed),
522                                cleanup: Some(cleanup_elapsed),
523                                commit: None,
524                                total: entry_elapsed,
525                            },
526                        }),
527                        entry_elapsed,
528                    )
529                } else {
530                    let commit_started = Instant::now();
531                    match transaction.commit() {
532                        Ok(outcome) => (Ok(outcome), entry_started.elapsed()),
533                        Err(source) => {
534                            let commit_elapsed = commit_started.elapsed();
535                            let entry_elapsed = entry_started.elapsed();
536                            (
537                                Err(ArtifactCacheBatchFailure::Cache {
538                                    phase: ArtifactCacheBatchFailurePhase::Commit,
539                                    source: Box::new(source),
540                                    timings: ArtifactCacheBatchFailureTimings {
541                                        preparation: preparation_elapsed,
542                                        callback: Some(callback_elapsed),
543                                        cleanup: None,
544                                        commit: Some(commit_elapsed),
545                                        total: entry_elapsed,
546                                    },
547                                }),
548                                entry_elapsed,
549                            )
550                        }
551                    }
552                }
553            }
554            Err(source) => {
555                let entry_elapsed = entry_started.elapsed();
556                (
557                    Err(ArtifactCacheBatchFailure::Cache {
558                        phase: ArtifactCacheBatchFailurePhase::Preparation,
559                        source: Box::new(source),
560                        timings: ArtifactCacheBatchFailureTimings {
561                            preparation: preparation_elapsed,
562                            callback: None,
563                            cleanup: None,
564                            commit: None,
565                            total: entry_elapsed,
566                        },
567                    }),
568                    entry_elapsed,
569                )
570            }
571        };
572        entries.push(ArtifactCacheBatchEntry {
573            index,
574            label: labeled.label.clone(),
575            result,
576            entry_elapsed,
577        });
578    }
579    Ok(ArtifactCacheBatchReport {
580        entries,
581        total: started.elapsed(),
582    })
583}
584
585fn validate_batch_labels(
586    specs: &[LabeledArtifactCacheSpec],
587) -> Result<(), ArtifactCacheBatchContractError> {
588    validate_labels(specs.iter().map(|labeled| labeled.label.as_str())).map_err(|error| match error
589    {
590        BatchLabelError::Empty { index } => ArtifactCacheBatchContractError::EmptyLabel { index },
591        BatchLabelError::Duplicate {
592            label,
593            first_index,
594            duplicate_index,
595        } => ArtifactCacheBatchContractError::DuplicateLabel {
596            label,
597            first_index,
598            duplicate_index,
599        },
600    })
601}
602
603impl std::fmt::Display for ArtifactCacheBatchContractError {
604    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
605        match self {
606            Self::EmptyLabel { index } => {
607                write!(formatter, "artifact batch label at index {index} is empty")
608            }
609            Self::DuplicateLabel {
610                label,
611                first_index,
612                duplicate_index,
613            } => write!(
614                formatter,
615                "artifact batch label {label:?} at index {duplicate_index} duplicates index {first_index}",
616            ),
617        }
618    }
619}
620
621impl std::error::Error for ArtifactCacheBatchContractError {}
622
623impl std::fmt::Display for ArtifactCacheBatchFailurePhase {
624    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
625        formatter.write_str(match self {
626            Self::Preparation => "preparation",
627            Self::Callback => "callback",
628            Self::Commit => "commit",
629        })
630    }
631}
632
633impl std::fmt::Display for ArtifactCacheBatchFailureTimings {
634    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
635        write!(
636            formatter,
637            "total={:?} preparation={:?} callback={:?} cleanup={:?} commit={:?}",
638            self.total, self.preparation, self.callback, self.cleanup, self.commit,
639        )
640    }
641}
642
643impl<E: std::fmt::Display> std::fmt::Display for ArtifactCacheBatchFailure<E> {
644    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
645        match self {
646            Self::Cache {
647                phase,
648                source,
649                timings,
650            } => write!(
651                formatter,
652                "artifact cache {phase} failed: {source}; timings=({timings})",
653            ),
654            Self::Build {
655                source,
656                cleanup_error,
657                timings,
658            } => {
659                write!(
660                    formatter,
661                    "artifact callback failed: {source}; timings=({timings})"
662                )?;
663                if let Some(cleanup_error) = cleanup_error {
664                    write!(formatter, "; cleanup also failed: {cleanup_error}")?;
665                }
666                Ok(())
667            }
668        }
669    }
670}
671
672impl<E> std::error::Error for ArtifactCacheBatchFailure<E>
673where
674    E: std::error::Error + 'static,
675{
676    fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
677        match self {
678            Self::Cache { source, .. } => Some(source.as_ref()),
679            Self::Build { source, .. } => Some(source.as_ref()),
680        }
681    }
682}
683
684#[cfg(test)]
685mod tests;