use std::{
sync::Arc,
time::Duration,
};
use crate::{
MetricSnapshot,
OperationAttributes,
Stage,
};
#[cfg(feature = "serde")]
use crate::{
Metric,
validation::{
validate_attributes,
validate_metrics,
},
};
#[cfg_attr(feature = "serde", derive(serde::Deserialize, serde::Serialize))]
#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum Phase {
Started,
Running,
Succeeded,
Failed,
Cancelled,
}
impl Phase {
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
Self::Started => "started",
Self::Running => "running",
Self::Succeeded => "succeeded",
Self::Failed => "failed",
Self::Cancelled => "cancelled",
}
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct Event {
operation_id: u64,
sequence: u64,
phase: Phase,
stage: Option<Stage>,
attributes: Arc<OperationAttributes>,
metrics: Vec<MetricSnapshot>,
elapsed: Duration,
}
impl Event {
pub(crate) fn new(
operation_id: u64,
sequence: u64,
phase: Phase,
stage: Option<Stage>,
attributes: Arc<OperationAttributes>,
metrics: Vec<MetricSnapshot>,
elapsed: Duration,
) -> Self {
Self {
operation_id,
sequence,
phase,
stage,
attributes,
metrics,
elapsed,
}
}
#[must_use]
pub const fn operation_id(&self) -> u64 {
self.operation_id
}
#[must_use]
pub const fn sequence(&self) -> u64 {
self.sequence
}
#[must_use]
pub const fn phase(&self) -> Phase {
self.phase
}
#[must_use]
pub const fn stage(&self) -> Option<&Stage> {
self.stage.as_ref()
}
#[must_use]
pub fn attributes(&self) -> &OperationAttributes {
&self.attributes
}
#[must_use]
pub fn attribute(&self, key: &str) -> Option<&str> {
self.attributes.get(key)
}
#[must_use]
pub fn metrics(&self) -> &[MetricSnapshot] {
&self.metrics
}
#[must_use]
pub fn metric(&self, metric_id: &str) -> Option<&MetricSnapshot> {
self.metrics.iter().find(|metric| metric.id() == metric_id)
}
#[must_use]
pub const fn elapsed(&self) -> Duration {
self.elapsed
}
}
#[cfg(feature = "serde")]
impl serde::Serialize for Event {
#[cfg_attr(coverage, inline(never))]
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
where
S: serde::Serializer,
{
serde::Serialize::serialize(
&EventWireRef {
operation_id: self.operation_id,
sequence: self.sequence,
phase: self.phase,
stage: self.stage.as_ref(),
attributes: self.attributes.as_ref(),
metrics: &self.metrics,
elapsed: format_duration(self.elapsed),
},
serializer,
)
}
}
#[cfg(feature = "serde")]
impl<'de> serde::Deserialize<'de> for Event {
#[cfg_attr(coverage, inline(never))]
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: serde::Deserializer<'de>,
{
let wire =
<EventWire as serde::Deserialize>::deserialize(deserializer)?;
let elapsed =
parse_duration(&wire.elapsed).map_err(serde::de::Error::custom)?;
validate_wire_event(&wire, elapsed)
.map_err(serde::de::Error::custom)?;
Ok(Self::new(
wire.operation_id,
wire.sequence,
wire.phase,
wire.stage,
Arc::new(wire.attributes),
wire.metrics,
elapsed,
))
}
}
#[cfg(feature = "serde")]
#[derive(serde::Serialize)]
struct EventWireRef<'a> {
operation_id: u64,
sequence: u64,
phase: Phase,
stage: Option<&'a Stage>,
#[serde(skip_serializing_if = "OperationAttributes::is_empty")]
attributes: &'a OperationAttributes,
metrics: &'a [MetricSnapshot],
elapsed: String,
}
#[cfg(feature = "serde")]
#[derive(serde::Deserialize)]
struct EventWire {
operation_id: u64,
sequence: u64,
phase: Phase,
stage: Option<Stage>,
#[serde(default)]
attributes: OperationAttributes,
metrics: Vec<MetricSnapshot>,
elapsed: String,
}
#[cfg(feature = "serde")]
fn format_duration(duration: Duration) -> String {
if duration.is_zero() {
return "0s".into();
}
let nanoseconds = duration.as_nanos();
for (unit, suffix) in [
(3_600_000_000_000_u128, "h"),
(60_000_000_000_u128, "m"),
(1_000_000_000_u128, "s"),
(1_000_000_u128, "ms"),
(1_000_u128, "us"),
(1_u128, "ns"),
] {
if nanoseconds.is_multiple_of(unit) {
return format!("{}{}", nanoseconds / unit, suffix);
}
}
unreachable!("nanosecond duration is always divisible by one nanosecond")
}
#[cfg(feature = "serde")]
fn parse_duration(text: &str) -> Result<Duration, String> {
let (amount, unit) = ["ms", "us", "ns", "h", "m", "s"]
.into_iter()
.find_map(|unit| text.strip_suffix(unit).map(|amount| (amount, unit)))
.ok_or_else(|| {
"elapsed must end in h, m, s, ms, us, or ns".to_owned()
})?;
if amount.is_empty() || !amount.bytes().all(|byte| byte.is_ascii_digit()) {
return Err("elapsed amount must be an unsigned integer".into());
}
let multiplier = match unit {
"h" => 3_600_000_000_000_u128,
"m" => 60_000_000_000_u128,
"s" => 1_000_000_000_u128,
"ms" => 1_000_000_u128,
"us" => 1_000_u128,
"ns" => 1_u128,
_ => unreachable!("unit list is exhaustive"),
};
let nanoseconds = amount
.parse::<u128>()
.map_err(|_| {
"elapsed amount is outside the supported range".to_owned()
})?
.checked_mul(multiplier)
.ok_or_else(|| "elapsed duration overflows".to_owned())?;
let seconds = nanoseconds / 1_000_000_000;
if seconds > u128::from(u64::MAX) {
return Err("elapsed duration overflows".into());
}
Ok(Duration::new(
seconds as u64,
(nanoseconds % 1_000_000_000) as u32,
))
}
#[cfg(feature = "serde")]
#[cfg_attr(coverage, inline(never))]
fn validate_wire_event(
wire: &EventWire,
elapsed: Duration,
) -> Result<(), String> {
if wire.operation_id == 0 {
return Err("operation_id must be nonzero".into());
}
let definitions = wire
.metrics
.iter()
.map(metric_definition)
.collect::<Vec<_>>();
validate_attributes(&wire.attributes).map_err(|error| error.to_string())?;
validate_metrics(&definitions).map_err(|error| error.to_string())?;
match wire.phase {
Phase::Started => {
if wire.sequence != 0
|| !elapsed.is_zero()
|| wire.metrics.iter().any(has_dynamic_counts)
{
return Err(
"started event must have sequence 0, zero elapsed, and zero counts".into(),
);
}
}
Phase::Running
| Phase::Succeeded
| Phase::Failed
| Phase::Cancelled
if wire.sequence == 0 =>
{
return Err(
"non-started event must have a positive sequence".into()
);
}
Phase::Running
| Phase::Succeeded
| Phase::Failed
| Phase::Cancelled => {}
}
Ok(())
}
#[cfg(feature = "serde")]
fn metric_definition(snapshot: &MetricSnapshot) -> Metric {
let metric = Metric::new(snapshot.id(), snapshot.name());
match snapshot.total() {
Some(total) => metric.total(total),
None => metric,
}
}
#[cfg(feature = "serde")]
const fn has_dynamic_counts(snapshot: &MetricSnapshot) -> bool {
snapshot.completed() != 0
|| snapshot.active() != 0
|| snapshot.succeeded() != 0
|| snapshot.failed() != 0
|| snapshot.cancelled() != 0
}
#[cfg(all(feature = "json-lines", coverage))]
#[doc(hidden)]
pub fn __coverage_event_serde() {
let value = serde_json::json!({
"operation_id": 1,
"sequence": 0,
"phase": "started",
"stage": null,
"metrics": [{
"id": "tasks",
"name": "Tasks",
"total": null,
"completed": 0,
"active": 0,
"succeeded": 0,
"failed": 0,
"cancelled": 0
}],
"elapsed": "0s"
});
let text =
serde_json::to_string(&value).expect("coverage JSON must serialize");
let mut deserializer = serde_json::Deserializer::from_str(&text);
let event = <Event as serde::Deserialize>::deserialize(&mut deserializer)
.expect("coverage event must deserialize");
let wire = EventWire {
operation_id: event.operation_id(),
sequence: event.sequence(),
phase: event.phase(),
stage: event.stage().cloned(),
attributes: event.attributes().clone(),
metrics: event.metrics().to_vec(),
elapsed: "0s".into(),
};
validate_wire_event(&wire, event.elapsed())
.expect("coverage event validation must succeed");
let mut invalid_wire = wire;
invalid_wire.attributes.insert(" ", "invalid");
assert!(validate_wire_event(&invalid_wire, event.elapsed()).is_err());
let mut output = Vec::new();
let mut serializer = serde_json::Serializer::new(&mut output);
coverage_serialize_event(&event, &mut serializer)
.expect("coverage event must serialize");
}
#[cfg(all(feature = "json-lines", coverage))]
#[inline(never)]
fn coverage_serialize_event(
event: &Event,
serializer: &mut serde_json::Serializer<&mut Vec<u8>>,
) -> Result<(), serde_json::Error> {
<Event as serde::Serialize>::serialize(event, serializer)
}