use std::fmt;
use std::sync::Arc;
use serde::{Deserialize, Serialize};
use crate::step::{Status, StepRef};
pub const EVENT_SCHEMA_VERSION: u32 = 1;
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(tag = "event", rename_all = "snake_case")]
pub enum Event {
RunStarted {
schema: u32,
run_id: Arc<str>,
},
ScenarioStarted {
scenario: Arc<str>,
file: Arc<str>,
#[serde(default, skip_serializing_if = "Option::is_none")]
timestamp_ms: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
worker: Option<u64>,
},
BatchStarted {
scenario: Arc<str>,
engine: Arc<str>,
steps: usize,
},
EntryRunning {
scenario: Arc<str>,
engine: Arc<str>,
entry: usize,
retry: u32,
},
StepFinished {
scenario: Arc<str>,
engine: Arc<str>,
step: StepRef,
status: Status,
attempts: u32,
duration_ms: u64,
captures: Vec<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
detail: Option<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
attempt_details: Vec<String>,
},
ScenarioFinished {
scenario: Arc<str>,
#[serde(default = "unknown_file")]
file: Arc<str>,
status: Status,
#[serde(default, skip_serializing_if = "Option::is_none")]
timestamp_ms: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
worker: Option<u64>,
},
RunFinished {
passed: usize,
failed: usize,
skipped: usize,
#[serde(default, skip_serializing_if = "std::ops::Not::not")]
cancelled: bool,
},
}
fn unknown_file() -> Arc<str> {
Arc::from("")
}
#[derive(Clone)]
pub struct EventSink(Arc<dyn Fn(&Event) + Send + Sync>);
impl EventSink {
pub fn new(f: impl Fn(&Event) + Send + Sync + 'static) -> Self {
Self(Arc::new(f))
}
pub fn null() -> Self {
Self(Arc::new(|_| {}))
}
pub fn emit(&self, event: &Event) {
(self.0)(event);
}
}
impl fmt::Debug for EventSink {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str("EventSink(..)")
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn run_started_wire_shape_is_stable() {
let event = Event::RunStarted {
schema: EVENT_SCHEMA_VERSION,
run_id: Arc::from("run-0001"),
};
let json = serde_json::to_string(&event).unwrap_or_default();
assert_eq!(
json,
r#"{"event":"run_started","schema":1,"run_id":"run-0001"}"#
);
}
#[test]
fn events_round_trip_through_jsonl() {
let event = Event::StepFinished {
scenario: Arc::from("Search finds a record"),
engine: Arc::from("http"),
step: StepRef {
file: Arc::from("tests/features/501_search.feature"),
line: 12,
text: Arc::from("the admin searches for \"Jansen\""),
},
status: Status::Passed,
attempts: 2,
duration_ms: 42,
captures: vec!["recordId".to_owned()],
detail: None,
attempt_details: vec!["attempt 1: HTTP 404 (retried)".to_owned()],
};
let json = serde_json::to_string(&event).unwrap_or_default();
let back: Event = serde_json::from_str(&json).unwrap_or(Event::RunFinished {
passed: 0,
failed: 0,
skipped: 0,
cancelled: false,
});
assert_eq!(back, event);
}
#[test]
fn sink_fans_out_borrowed_events() {
use std::sync::atomic::{AtomicUsize, Ordering};
let seen = Arc::new(AtomicUsize::new(0));
let counter = Arc::clone(&seen);
let sink = EventSink::new(move |_| {
counter.fetch_add(1, Ordering::SeqCst);
});
let clone = sink.clone();
clone.emit(&Event::RunFinished {
passed: 1,
failed: 0,
skipped: 0,
cancelled: false,
});
sink.emit(&Event::RunFinished {
passed: 1,
failed: 0,
skipped: 0,
cancelled: false,
});
assert_eq!(seen.load(Ordering::SeqCst), 2);
}
}