1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
//! A set of information representing the terminal state of a node after a graceful shutdown.
//!
//! This state must include all the metrics used by the node (if any exist).
use otel_arrow_dfe_telemetry::metrics::MetricSetSnapshot;
use std::ops::Add;
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
/// Pipeline-wide deadline shared by every terminal metrics handoff.
///
/// The runtime control manager records the real shutdown deadline as soon as
/// it accepts shutdown. Error paths that terminate without a shutdown message
/// lazily establish one finite fallback. Every final reporter then uses the
/// same absolute deadline instead of receiving a fresh timeout.
#[derive(Clone, Debug, Default)]
pub(crate) struct TerminalMetricsDeadline {
deadline: Arc<Mutex<Option<Instant>>>,
}
impl TerminalMetricsDeadline {
const FALLBACK: Duration = Duration::from_secs(5);
/// Records a terminal metrics deadline, preserving the earliest deadline observed.
pub(crate) fn record(&self, deadline: Instant) {
let mut current = self
.deadline
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
*current = Some(current.map_or(deadline, |current| current.min(deadline)));
}
/// Returns the shared deadline, installing a finite fallback if necessary.
pub(crate) fn get(&self) -> Instant {
let mut deadline = self
.deadline
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
*deadline.get_or_insert_with(|| Instant::now() + Self::FALLBACK)
}
}
/// Captures the last metric snapshots produced by a node when it terminates gracefully.
pub struct TerminalState {
deadline: Instant,
metrics: Vec<MetricSetSnapshot>,
}
impl TerminalState {
/// Create a new terminal state with the provided metrics.
pub fn new<MI>(deadline: Instant, metrics: MI) -> Self
where
MI: IntoIterator,
MI::Item: Into<MetricSetSnapshot>,
{
Self {
deadline,
metrics: metrics.into_iter().map(Into::into).collect(),
}
}
/// Returns the deadline by which the node must terminate.
#[must_use]
pub const fn deadline(&self) -> Instant {
self.deadline
}
/// Returns a slice of the metric snapshots captured in this terminal state.
#[must_use]
pub fn metrics(&self) -> &[MetricSetSnapshot] {
&self.metrics
}
/// Consumes the terminal state and returns the contained metric snapshots.
#[must_use]
pub fn into_metrics(self) -> Vec<MetricSetSnapshot> {
self.metrics
}
/// Returns `true` when no metrics were captured.
#[must_use]
pub const fn is_empty(&self) -> bool {
self.metrics.is_empty()
}
}
impl Default for TerminalState {
fn default() -> Self {
Self {
deadline: Instant::now().add(Duration::from_secs(1)),
metrics: Vec::new(),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
/// Scenario: Multiple terminal metric producers record different absolute deadlines.
/// Guarantees: Every producer observes the earliest recorded deadline.
#[test]
fn terminal_metrics_deadline_preserves_the_earliest_recorded_deadline() {
let deadline = TerminalMetricsDeadline::default();
let now = Instant::now();
deadline.record(now + Duration::from_secs(2));
deadline.record(now + Duration::from_secs(1));
deadline.record(now + Duration::from_secs(3));
assert_eq!(deadline.get(), now + Duration::from_secs(1));
assert_eq!(deadline.clone().get(), now + Duration::from_secs(1));
}
/// Scenario: Terminal metrics are requested without an explicit pipeline shutdown deadline.
/// Guarantees: All terminal reporters share one finite fallback deadline.
#[test]
fn terminal_metrics_deadline_installs_only_one_fallback() {
let deadline = TerminalMetricsDeadline::default();
let fallback = deadline.get();
assert_eq!(deadline.get(), fallback);
assert_eq!(deadline.clone().get(), fallback);
}
}