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