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#[derive(Clone, Debug, Eq, PartialEq)]
15pub struct LabeledArtifactCacheSpec {
16 label: String,
17 spec: ArtifactCacheSpec,
18}
19
20#[derive(Debug)]
22pub struct ArtifactCacheBatchReport<E> {
23 entries: Vec<ArtifactCacheBatchEntry<E>>,
24 total: Duration,
25}
26
27#[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#[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#[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#[non_exhaustive]
56#[derive(Clone, Debug, Eq, PartialEq)]
57pub enum ArtifactCacheBatchContractError {
58 EmptyLabel {
60 index: usize,
62 },
63 DuplicateLabel {
65 label: String,
67 first_index: usize,
69 duplicate_index: usize,
71 },
72}
73
74#[derive(Clone, Copy, Debug, Eq, PartialEq)]
76pub enum ArtifactCacheBatchFailurePhase {
77 Preparation,
79 Callback,
81 Commit,
83}
84
85#[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#[derive(Debug)]
97pub enum ArtifactCacheBatchFailure<E> {
98 Cache {
100 phase: ArtifactCacheBatchFailurePhase,
102 source: Box<ArtifactCacheError>,
104 timings: ArtifactCacheBatchFailureTimings,
106 },
107 Build {
109 source: Box<E>,
111 cleanup_error: Option<Box<ArtifactCacheError>>,
113 timings: ArtifactCacheBatchFailureTimings,
115 },
116}
117
118#[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 #[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 #[must_use]
142 pub fn label(&self) -> &str {
143 &self.label
144 }
145
146 #[must_use]
148 pub const fn spec(&self) -> &ArtifactCacheSpec {
149 &self.spec
150 }
151
152 #[must_use]
154 pub fn into_parts(self) -> (String, ArtifactCacheSpec) {
155 (self.label, self.spec)
156 }
157}
158
159impl<E> ArtifactCacheBatchReport<E> {
160 #[must_use]
162 pub fn entries(&self) -> &[ArtifactCacheBatchEntry<E>] {
163 &self.entries
164 }
165
166 #[must_use]
168 pub fn into_entries(self) -> Vec<ArtifactCacheBatchEntry<E>> {
169 self.entries
170 }
171
172 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 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 #[must_use]
202 pub const fn total(&self) -> Duration {
203 self.total
204 }
205
206 #[must_use]
208 pub fn is_success(&self) -> bool {
209 self.entries.iter().all(ArtifactCacheBatchEntry::is_success)
210 }
211
212 #[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 #[must_use]
243 pub const fn index(&self) -> usize {
244 self.index
245 }
246
247 #[must_use]
249 pub fn label(&self) -> &str {
250 &self.label
251 }
252
253 pub const fn result(&self) -> Result<&ArtifactCacheOutcome, &ArtifactCacheBatchFailure<E>> {
255 self.result.as_ref()
256 }
257
258 #[must_use]
260 pub fn outcome(&self) -> Option<&ArtifactCacheOutcome> {
261 self.result.as_ref().ok()
262 }
263
264 #[must_use]
266 pub fn failure(&self) -> Option<&ArtifactCacheBatchFailure<E>> {
267 self.result.as_ref().err()
268 }
269
270 #[must_use]
272 pub const fn entry_elapsed(&self) -> Duration {
273 self.entry_elapsed
274 }
275
276 #[must_use]
278 pub const fn is_success(&self) -> bool {
279 self.result.is_ok()
280 }
281
282 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 #[must_use]
298 pub const fn index(self) -> usize {
299 self.index
300 }
301
302 #[must_use]
304 pub const fn label(self) -> &'a str {
305 self.label
306 }
307
308 #[must_use]
310 pub const fn outcome(self) -> &'a ArtifactCacheOutcome {
311 self.outcome
312 }
313
314 #[must_use]
316 pub const fn entry_elapsed(self) -> Duration {
317 self.entry_elapsed
318 }
319}
320
321impl<'a, E> ArtifactCacheBatchFailedEntry<'a, E> {
322 #[must_use]
324 pub const fn index(&self) -> usize {
325 self.index
326 }
327
328 #[must_use]
330 pub const fn label(&self) -> &'a str {
331 self.label
332 }
333
334 #[must_use]
336 pub const fn failure(&self) -> &'a ArtifactCacheBatchFailure<E> {
337 self.failure
338 }
339
340 #[must_use]
342 pub const fn timings(&self) -> ArtifactCacheBatchFailureTimings {
343 self.failure.timings()
344 }
345
346 #[must_use]
348 pub const fn entry_elapsed(&self) -> Duration {
349 self.entry_elapsed
350 }
351}
352
353impl<E> ArtifactCacheBatchFailure<E> {
354 #[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 #[must_use]
365 pub const fn timings(&self) -> ArtifactCacheBatchFailureTimings {
366 match self {
367 Self::Cache { timings, .. } | Self::Build { timings, .. } => *timings,
368 }
369 }
370
371 #[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 #[must_use]
384 pub const fn preparation(self) -> Duration {
385 self.preparation
386 }
387
388 #[must_use]
390 pub const fn callback(self) -> Option<Duration> {
391 self.callback
392 }
393
394 #[must_use]
396 pub const fn cleanup(self) -> Option<Duration> {
397 self.cleanup
398 }
399
400 #[must_use]
402 pub const fn commit(self) -> Option<Duration> {
403 self.commit
404 }
405
406 #[must_use]
408 pub const fn total(self) -> Duration {
409 self.total
410 }
411}
412
413impl ArtifactCacheBatchMetrics {
414 #[must_use]
416 pub const fn entries(self) -> usize {
417 self.entries
418 }
419
420 #[must_use]
422 pub const fn succeeded(self) -> usize {
423 self.succeeded
424 }
425
426 #[must_use]
428 pub const fn failed(self) -> usize {
429 self.failed
430 }
431
432 #[must_use]
434 pub const fn built(self) -> usize {
435 self.built
436 }
437
438 #[must_use]
440 pub const fn reused(self) -> usize {
441 self.reused
442 }
443
444 #[must_use]
446 pub const fn successful_timings(self) -> ArtifactCacheTimings {
447 self.successful_timings
448 }
449
450 #[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
474pub 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;