use std::time::SystemTime;
use bytes::Bytes;
use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "lowercase")]
pub enum StreamKind {
Stdout,
Stderr,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(tag = "type", rename_all = "camelCase")]
#[non_exhaustive]
pub enum OutputEvent {
Chunk(OutputChunk),
#[serde(rename_all = "camelCase")]
RunStarted {
generation: u64,
attempt: u32,
#[serde(with = "crate::resource::metadata::time_serde")]
#[cfg_attr(
feature = "schema",
schemars(schema_with = "crate::schema::unix_millis")
)]
started_at: SystemTime,
},
#[serde(rename_all = "camelCase")]
RunFinished {
generation: u64,
attempt: u32,
#[serde(skip_serializing_if = "Option::is_none")]
exit_code: Option<i32>,
#[serde(with = "crate::resource::metadata::time_serde")]
#[cfg_attr(
feature = "schema",
schemars(schema_with = "crate::schema::unix_millis")
)]
finished_at: SystemTime,
},
Lagged {
skipped: u64,
},
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "camelCase")]
pub struct OutputChunk {
pub generation: u64,
pub attempt: u32,
pub stream: StreamKind,
pub seq: u64,
#[serde(with = "crate::resource::metadata::time_serde")]
#[cfg_attr(
feature = "schema",
schemars(schema_with = "crate::schema::unix_millis")
)]
pub ts: SystemTime,
#[serde(with = "bytes_as_base64")]
#[cfg_attr(
feature = "schema",
schemars(schema_with = "crate::schema::base64_bytes")
)]
pub line: Bytes,
}
mod bytes_as_base64 {
use base64::{Engine as _, engine::general_purpose::STANDARD};
use bytes::Bytes;
use serde::{Deserialize, Deserializer, Serializer};
pub(super) fn serialize<S>(b: &Bytes, s: S) -> Result<S::Ok, S::Error>
where
S: Serializer,
{
s.serialize_str(&STANDARD.encode(b))
}
pub(super) fn deserialize<'de, D>(d: D) -> Result<Bytes, D::Error>
where
D: Deserializer<'de>,
{
let s = String::deserialize(d)?;
STANDARD
.decode(s)
.map(Bytes::from)
.map_err(serde::de::Error::custom)
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::time::{Duration, UNIX_EPOCH};
#[test]
fn wire_shape_is_pinned_for_every_event() {
let cases = [
(
OutputEvent::Chunk(OutputChunk {
generation: 2,
attempt: 1,
stream: StreamKind::Stdout,
seq: 0,
ts: UNIX_EPOCH + Duration::from_millis(1_700),
line: Bytes::from_static(b"hi"),
}),
r#"{"type":"chunk","generation":2,"attempt":1,"stream":"stdout","seq":0,"ts":1700,"line":"aGk="}"#,
),
(
OutputEvent::RunStarted {
generation: 4,
attempt: 2,
started_at: UNIX_EPOCH + Duration::from_millis(1_234),
},
r#"{"type":"runStarted","generation":4,"attempt":2,"startedAt":1234}"#,
),
(
OutputEvent::RunFinished {
generation: 4,
attempt: 2,
exit_code: Some(0),
finished_at: UNIX_EPOCH + Duration::from_millis(2_222),
},
r#"{"type":"runFinished","generation":4,"attempt":2,"exitCode":0,"finishedAt":2222}"#,
),
(
OutputEvent::Lagged { skipped: 42 },
r#"{"type":"lagged","skipped":42}"#,
),
];
for (event, expected) in cases {
assert_eq!(serde_json::to_string(&event).unwrap(), expected);
}
}
#[test]
fn every_event_roundtrips_through_json() {
let cases = [
OutputEvent::Chunk(OutputChunk {
generation: 2,
attempt: 1,
stream: StreamKind::Stderr,
seq: 0,
ts: UNIX_EPOCH + Duration::from_millis(1_700_000_000_000),
line: Bytes::from_static(b"warning"),
}),
OutputEvent::RunStarted {
generation: 2,
attempt: 1,
started_at: UNIX_EPOCH + Duration::from_millis(1_700_000_000_000),
},
OutputEvent::RunFinished {
generation: 2,
attempt: 1,
exit_code: Some(42),
finished_at: UNIX_EPOCH + Duration::from_millis(1_700_000_001_000),
},
OutputEvent::Lagged { skipped: 7 },
];
for original in cases {
let json = serde_json::to_string(&original).unwrap();
let back: OutputEvent = serde_json::from_str(&json).unwrap();
assert_eq!(back, original, "roundtrip failed for {json}");
}
}
#[test]
fn binary_chunk_roundtrips_exactly_as_base64() {
let chunk = OutputChunk {
generation: 1,
attempt: 1,
stream: StreamKind::Stdout,
seq: 0,
ts: UNIX_EPOCH,
line: Bytes::from_static(&[b'h', b'i', 0xFF, 0xFE]),
};
let json = serde_json::to_string(&chunk).unwrap();
assert!(json.contains(r#""line":"aGn//g==""#), "{json}");
let decoded: OutputChunk = serde_json::from_str(&json).unwrap();
assert_eq!(decoded, chunk);
}
}