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 built: usize,
123 reused: usize,
124 successful_timings: ArtifactCacheTimings,
125 total: Duration,
126}
127
128impl LabeledArtifactCacheSpec {
129 #[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 #[must_use]
140 pub fn label(&self) -> &str {
141 &self.label
142 }
143
144 #[must_use]
146 pub const fn spec(&self) -> &ArtifactCacheSpec {
147 &self.spec
148 }
149
150 #[must_use]
152 pub fn into_parts(self) -> (String, ArtifactCacheSpec) {
153 (self.label, self.spec)
154 }
155}
156
157impl<E> ArtifactCacheBatchReport<E> {
158 #[must_use]
160 pub fn entries(&self) -> &[ArtifactCacheBatchEntry<E>] {
161 &self.entries
162 }
163
164 #[must_use]
166 pub fn into_entries(self) -> Vec<ArtifactCacheBatchEntry<E>> {
167 self.entries
168 }
169
170 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 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 #[must_use]
200 pub const fn total(&self) -> Duration {
201 self.total
202 }
203
204 #[must_use]
206 pub fn is_success(&self) -> bool {
207 self.entries.iter().all(ArtifactCacheBatchEntry::is_success)
208 }
209
210 #[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 #[must_use]
236 pub const fn index(&self) -> usize {
237 self.index
238 }
239
240 #[must_use]
242 pub fn label(&self) -> &str {
243 &self.label
244 }
245
246 pub const fn result(&self) -> Result<&ArtifactCacheOutcome, &ArtifactCacheBatchFailure<E>> {
248 self.result.as_ref()
249 }
250
251 #[must_use]
253 pub fn outcome(&self) -> Option<&ArtifactCacheOutcome> {
254 self.result.as_ref().ok()
255 }
256
257 #[must_use]
259 pub fn failure(&self) -> Option<&ArtifactCacheBatchFailure<E>> {
260 self.result.as_ref().err()
261 }
262
263 #[must_use]
265 pub const fn entry_elapsed(&self) -> Duration {
266 self.entry_elapsed
267 }
268
269 #[must_use]
271 pub const fn is_success(&self) -> bool {
272 self.result.is_ok()
273 }
274
275 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 #[must_use]
291 pub const fn index(self) -> usize {
292 self.index
293 }
294
295 #[must_use]
297 pub const fn label(self) -> &'a str {
298 self.label
299 }
300
301 #[must_use]
303 pub const fn outcome(self) -> &'a ArtifactCacheOutcome {
304 self.outcome
305 }
306
307 #[must_use]
309 pub const fn entry_elapsed(self) -> Duration {
310 self.entry_elapsed
311 }
312}
313
314impl<'a, E> ArtifactCacheBatchFailedEntry<'a, E> {
315 #[must_use]
317 pub const fn index(&self) -> usize {
318 self.index
319 }
320
321 #[must_use]
323 pub const fn label(&self) -> &'a str {
324 self.label
325 }
326
327 #[must_use]
329 pub const fn failure(&self) -> &'a ArtifactCacheBatchFailure<E> {
330 self.failure
331 }
332
333 #[must_use]
335 pub const fn timings(&self) -> ArtifactCacheBatchFailureTimings {
336 self.failure.timings()
337 }
338
339 #[must_use]
341 pub const fn entry_elapsed(&self) -> Duration {
342 self.entry_elapsed
343 }
344}
345
346impl<E> ArtifactCacheBatchFailure<E> {
347 #[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 #[must_use]
358 pub const fn timings(&self) -> ArtifactCacheBatchFailureTimings {
359 match self {
360 Self::Cache { timings, .. } | Self::Build { timings, .. } => *timings,
361 }
362 }
363
364 #[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 #[must_use]
377 pub const fn preparation(self) -> Duration {
378 self.preparation
379 }
380
381 #[must_use]
383 pub const fn callback(self) -> Option<Duration> {
384 self.callback
385 }
386
387 #[must_use]
389 pub const fn cleanup(self) -> Option<Duration> {
390 self.cleanup
391 }
392
393 #[must_use]
395 pub const fn commit(self) -> Option<Duration> {
396 self.commit
397 }
398
399 #[must_use]
401 pub const fn total(self) -> Duration {
402 self.total
403 }
404}
405
406impl ArtifactCacheBatchMetrics {
407 #[must_use]
409 pub const fn entries(self) -> usize {
410 self.entries
411 }
412
413 #[must_use]
415 pub const fn succeeded(self) -> usize {
416 self.built + self.reused
417 }
418
419 #[must_use]
421 pub const fn failed(self) -> usize {
422 self.entries - self.succeeded()
423 }
424
425 #[must_use]
427 pub const fn built(self) -> usize {
428 self.built
429 }
430
431 #[must_use]
433 pub const fn reused(self) -> usize {
434 self.reused
435 }
436
437 #[must_use]
439 pub const fn successful_timings(self) -> ArtifactCacheTimings {
440 self.successful_timings
441 }
442
443 #[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
467pub 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;