use serde::{Deserialize, Serialize};
use std::sync::Arc;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum Disposition {
Running,
Completed,
Interrupted,
Skipped,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum Validity {
Succeeded,
Failed,
}
pub trait Payload: Send + Sync + std::fmt::Debug + 'static {}
impl<T: ?Sized + crate::adapter::ResultBody + 'static> Payload for T {}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Outcome {
pub disposition: Disposition,
pub validity: Validity,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub reason: Option<String>,
#[serde(skip)]
pub payload: Option<Arc<dyn Payload>>,
}
impl Outcome {
pub const fn new(disposition: Disposition, validity: Validity) -> Self {
Self {
disposition,
validity,
reason: None,
payload: None,
}
}
pub const fn completed() -> Self {
Self::new(Disposition::Completed, Validity::Succeeded)
}
pub const fn failed() -> Self {
Self::new(Disposition::Interrupted, Validity::Failed)
}
pub const fn completed_failed() -> Self {
Self::new(Disposition::Completed, Validity::Failed)
}
pub const fn interrupted() -> Self {
Self::new(Disposition::Interrupted, Validity::Succeeded)
}
pub const fn skipped() -> Self {
Self::new(Disposition::Skipped, Validity::Succeeded)
}
pub fn with_reason(mut self, reason: impl Into<String>) -> Self {
self.reason = Some(reason.into());
self
}
pub fn is_failure(&self) -> bool {
matches!(self.validity, Validity::Failed)
}
pub fn glyph(&self) -> char {
match (self.disposition, self.validity) {
(Disposition::Running, _) => '⋯',
(Disposition::Skipped, _) => '~',
(_, Validity::Failed) => '✗',
(Disposition::Completed, Validity::Succeeded) => '✓',
(Disposition::Interrupted, Validity::Succeeded) => '…',
}
}
pub fn label(&self) -> &'static str {
match (self.disposition, self.validity) {
(Disposition::Running, _) => "running",
(Disposition::Skipped, _) => "skipped",
(Disposition::Completed, Validity::Succeeded) => "completed",
(Disposition::Completed, Validity::Failed) => "completed_failed",
(Disposition::Interrupted, Validity::Succeeded) => "interrupted",
(Disposition::Interrupted, Validity::Failed) => "failed",
}
}
pub fn with_payload(mut self, payload: Arc<dyn Payload>) -> Self {
self.payload = Some(payload);
self
}
}
impl PartialEq for Outcome {
fn eq(&self, other: &Self) -> bool {
self.disposition == other.disposition
&& self.validity == other.validity
&& self.reason == other.reason
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum SessionDisposition {
Success,
Failure,
}
impl SessionDisposition {
pub fn exit_code(&self) -> i32 {
match self {
SessionDisposition::Success => 0,
SessionDisposition::Failure => 1,
}
}
pub fn label(&self) -> &'static str {
match self {
SessionDisposition::Success => "SUCCESS",
SessionDisposition::Failure => "FAILURE",
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
pub struct PhaseIdentity {
pub name: String,
pub labels: String,
}
impl PhaseIdentity {
pub fn new(name: impl Into<String>, labels: impl Into<String>) -> Self {
Self {
name: name.into(),
labels: labels.into(),
}
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct PhaseErrorDetail {
pub class: String,
pub message: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub op_name: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub cycle: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub op_template: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub op_resolved: Option<String>,
pub at_nanos: u64,
#[serde(default)]
pub retryable: bool,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, Default)]
pub struct ResumeCursor {
#[serde(default)]
pub opaque: Vec<u8>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ReasonClass {
Timeout,
StopCondition,
Error,
Panic,
Operator,
}
impl ReasonClass {
pub fn as_str(&self) -> &'static str {
match self {
ReasonClass::Timeout => "timeout",
ReasonClass::StopCondition => "stop_condition",
ReasonClass::Error => "error",
ReasonClass::Panic => "panic",
ReasonClass::Operator => "operator",
}
}
pub fn from_error_class(class: &str) -> Self {
match class {
"timeout" | "poll_timeout" => ReasonClass::Timeout,
"panic" => ReasonClass::Panic,
"operator" => ReasonClass::Operator,
"error_rate_exceeded" => ReasonClass::StopCondition,
c if c.starts_with("stop_condition") => ReasonClass::StopCondition,
_ => ReasonClass::Error,
}
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(from = "PhaseOutcomeWire")]
pub struct PhaseOutcome {
pub phase_id: PhaseIdentity,
pub disposition: Disposition,
pub validity: Validity,
pub duration_secs: f64,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub errors: Vec<PhaseErrorDetail>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub resume_cursor: Option<ResumeCursor>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub phase_hash: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub params_consumed: Option<String>,
}
#[derive(Deserialize)]
struct PhaseOutcomeWire {
phase_id: PhaseIdentity,
disposition: Option<Disposition>,
validity: Option<Validity>,
status: Option<LegacyPhaseStatus>,
duration_secs: f64,
#[serde(default)]
errors: Vec<PhaseErrorDetail>,
#[serde(default)]
resume_cursor: Option<ResumeCursor>,
#[serde(default)]
phase_hash: Option<String>,
#[serde(default)]
params_consumed: Option<String>,
}
#[derive(Deserialize)]
#[serde(rename_all = "snake_case")]
enum LegacyPhaseStatus {
Completed,
Failed,
Skipped,
CursorSuspended,
}
impl From<PhaseOutcomeWire> for PhaseOutcome {
fn from(w: PhaseOutcomeWire) -> Self {
let (disposition, validity) = match (w.disposition, w.validity, w.status) {
(Some(d), Some(v), _) => (d, v),
(_, _, Some(LegacyPhaseStatus::Completed)) => {
(Disposition::Completed, Validity::Succeeded)
}
(_, _, Some(LegacyPhaseStatus::Failed)) => (Disposition::Interrupted, Validity::Failed),
(_, _, Some(LegacyPhaseStatus::Skipped)) => (Disposition::Skipped, Validity::Succeeded),
(_, _, Some(LegacyPhaseStatus::CursorSuspended)) => {
(Disposition::Interrupted, Validity::Succeeded)
}
_ => (Disposition::Completed, Validity::Succeeded),
};
Self {
phase_id: w.phase_id,
disposition,
validity,
duration_secs: w.duration_secs,
errors: w.errors,
resume_cursor: w.resume_cursor,
phase_hash: w.phase_hash,
params_consumed: w.params_consumed,
}
}
}
impl PhaseOutcome {
pub fn reason_class(&self) -> Option<ReasonClass> {
match self.validity {
Validity::Succeeded => None,
Validity::Failed => Some(ReasonClass::from_error_class(
self.errors
.first()
.map(|e| e.class.as_str())
.unwrap_or("error"),
)),
}
}
pub fn protocol_class(&self) -> &'static str {
if matches!(self.disposition, Disposition::Skipped) {
return "SKIPPED";
}
match self.reason_class() {
None => "COMPLETED",
Some(ReasonClass::Timeout) => "OUT-OF-RANGE",
Some(_) => "FAILED",
}
}
pub fn completed(phase_id: PhaseIdentity, duration_secs: f64) -> Self {
Self {
phase_id,
disposition: Disposition::Completed,
validity: Validity::Succeeded,
duration_secs,
errors: Vec::new(),
resume_cursor: None,
phase_hash: None,
params_consumed: None,
}
}
pub fn failed(
phase_id: PhaseIdentity,
duration_secs: f64,
errors: Vec<PhaseErrorDetail>,
) -> Self {
assert!(
!errors.is_empty(),
"PhaseOutcome::failed requires at least one error"
);
Self {
phase_id,
disposition: Disposition::Interrupted,
validity: Validity::Failed,
duration_secs,
errors,
resume_cursor: None,
phase_hash: None,
params_consumed: None,
}
}
pub fn skipped(phase_id: PhaseIdentity) -> Self {
Self {
phase_id,
disposition: Disposition::Skipped,
validity: Validity::Succeeded,
duration_secs: 0.0,
errors: Vec::new(),
resume_cursor: None,
phase_hash: None,
params_consumed: None,
}
}
pub fn interrupted(
phase_id: PhaseIdentity,
duration_secs: f64,
resume_cursor: Option<ResumeCursor>,
) -> Self {
Self {
phase_id,
disposition: Disposition::Interrupted,
validity: Validity::Succeeded,
duration_secs,
errors: Vec::new(),
resume_cursor,
phase_hash: None,
params_consumed: None,
}
}
pub fn with_phase_hash(mut self, hex_hash: String) -> Self {
self.phase_hash = Some(hex_hash);
self
}
pub fn with_params_consumed(mut self, json: Option<String>) -> Self {
self.params_consumed = json;
self
}
pub fn first_error_message(&self) -> Option<&str> {
self.errors.first().map(|e| e.message.as_str())
}
pub fn outcome(&self) -> Outcome {
Outcome::new(self.disposition, self.validity)
}
pub fn is_failure(&self) -> bool {
self.outcome().is_failure()
}
pub fn glyph(&self) -> char {
self.outcome().glyph()
}
pub fn label(&self) -> &'static str {
self.outcome().label()
}
pub fn to_sqlite_row(
&self,
session: &str,
exec_id: u64,
started_at_nanos: i64,
) -> nmbrs_metrics::reporters::sqlite::PhaseOutcomeRow {
let ended_at_nanos: i64 = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_nanos() as i64)
.unwrap_or(started_at_nanos);
nmbrs_metrics::reporters::sqlite::PhaseOutcomeRow {
session: session.to_string(),
exec_id,
phase_name: self.phase_id.name.clone(),
phase_labels: self.phase_id.labels.clone(),
status: self.label().to_string(),
duration_secs: self.duration_secs,
started_at_nanos,
ended_at_nanos,
reason_class: self.reason_class().map(|c| c.as_str().to_string()),
phase_hash: self.phase_hash.clone(),
params_consumed: self.params_consumed.clone(),
errors: self
.errors
.iter()
.map(|e| nmbrs_metrics::reporters::sqlite::PhaseErrorRow {
class: e.class.clone(),
message: e.message.clone(),
op_name: e.op_name.clone(),
cycle: e.cycle,
op_template: e.op_template.clone(),
op_resolved: e.op_resolved.clone(),
at_nanos: e.at_nanos as i64,
retryable: e.retryable,
})
.collect(),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn is_failure_keys_on_validity_alone() {
assert!(!Outcome::completed().is_failure());
assert!(Outcome::failed().is_failure());
assert!(Outcome::completed_failed().is_failure());
assert!(!Outcome::interrupted().is_failure());
assert!(!Outcome::skipped().is_failure());
}
#[test]
fn glyphs_and_labels_cover_the_axis_pairs() {
assert_eq!(Outcome::completed().glyph(), '✓');
assert_eq!(Outcome::failed().glyph(), '✗');
assert_eq!(Outcome::completed_failed().glyph(), '✗');
assert_eq!(Outcome::interrupted().glyph(), '…');
assert_eq!(Outcome::skipped().glyph(), '~');
assert_eq!(Outcome::completed().label(), "completed");
assert_eq!(Outcome::failed().label(), "failed");
assert_eq!(Outcome::completed_failed().label(), "completed_failed");
assert_eq!(Outcome::interrupted().label(), "interrupted");
assert_eq!(Outcome::skipped().label(), "skipped");
}
#[test]
fn session_disposition_exit_codes() {
assert_eq!(SessionDisposition::Success.exit_code(), 0);
assert_ne!(SessionDisposition::Failure.exit_code(), 0);
}
#[test]
fn outcome_completed_has_no_errors() {
let o = PhaseOutcome::completed(PhaseIdentity::new("p", "x=1"), 1.5);
assert_eq!(o.disposition, Disposition::Completed);
assert_eq!(o.validity, Validity::Succeeded);
assert!(o.errors.is_empty());
assert!(o.first_error_message().is_none());
}
#[test]
fn outcome_failed_carries_errors() {
let errors = vec![PhaseErrorDetail {
class: "Timeout".into(),
message: "connection timed out".into(),
op_name: Some("read_state".into()),
cycle: Some(0),
op_template: Some("SELECT ...".into()),
op_resolved: Some("SELECT * FROM ks.t".into()),
at_nanos: 1_000_000_000,
retryable: true,
}];
let o = PhaseOutcome::failed(PhaseIdentity::new("p", "x=1"), 30.0, errors.clone());
assert_eq!(o.disposition, Disposition::Interrupted);
assert_eq!(o.validity, Validity::Failed);
assert_eq!(o.errors, errors);
assert_eq!(o.first_error_message(), Some("connection timed out"));
}
#[test]
#[should_panic(expected = "at least one error")]
fn outcome_failed_requires_non_empty_errors() {
let _ = PhaseOutcome::failed(PhaseIdentity::new("p", ""), 1.0, Vec::new());
}
#[test]
fn outcome_skipped_has_zero_duration_and_no_errors() {
let o = PhaseOutcome::skipped(PhaseIdentity::new("p", ""));
assert_eq!(o.disposition, Disposition::Skipped);
assert_eq!(o.duration_secs, 0.0);
assert!(o.errors.is_empty());
}
#[test]
fn outcome_round_trips_through_json() {
let original = PhaseOutcome {
phase_id: PhaseIdentity::new("ensure_compacted", "(k=10)"),
disposition: Disposition::Interrupted,
validity: Validity::Failed,
duration_secs: 14400.0,
errors: vec![PhaseErrorDetail {
class: "poll_timeout".into(),
message: "deadline reached".into(),
op_name: None,
cycle: None,
op_template: None,
op_resolved: None,
at_nanos: 0,
retryable: false,
}],
resume_cursor: None,
phase_hash: None,
params_consumed: None,
};
let json = serde_json::to_string(&original).expect("serialise");
let parsed: PhaseOutcome = serde_json::from_str(&json).expect("deserialise");
assert_eq!(parsed, original);
}
#[test]
fn two_axis_outcome_projects_to_and_from_status() {
assert_eq!(
Outcome::completed(),
Outcome::new(Disposition::Completed, Validity::Succeeded)
);
assert_eq!(
Outcome::failed(),
Outcome::new(Disposition::Interrupted, Validity::Failed)
);
assert_eq!(
Outcome::interrupted(),
Outcome::new(Disposition::Interrupted, Validity::Succeeded)
);
assert!(Outcome::failed().is_failure());
assert!(!Outcome::interrupted().is_failure());
assert!(!Outcome::completed().is_failure());
let oc = PhaseOutcome::completed(PhaseIdentity::new("p", ""), 1.0);
assert_eq!(oc.outcome(), Outcome::completed());
assert!(!oc.outcome().is_failure());
}
#[test]
fn legacy_status_records_still_deserialize() {
let cases = [
("completed", Disposition::Completed, Validity::Succeeded),
("failed", Disposition::Interrupted, Validity::Failed),
("skipped", Disposition::Skipped, Validity::Succeeded),
(
"cursor_suspended",
Disposition::Interrupted,
Validity::Succeeded,
),
];
for (status, d, v) in cases {
let json = format!(
r#"{{"phase_id":{{"name":"p","labels":""}},"status":"{status}","duration_secs":1.0}}"#
);
let parsed: PhaseOutcome = serde_json::from_str(&json)
.unwrap_or_else(|e| panic!("legacy '{status}' must parse: {e}"));
assert_eq!(parsed.disposition, d, "disposition for '{status}'");
assert_eq!(parsed.validity, v, "validity for '{status}'");
}
}
}