use std::collections::BTreeMap;
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct Metadata {
metadata_version: u16,
timestamp: uhlc::Timestamp,
pub parameters: MetadataParameters,
}
impl Metadata {
pub const CURRENT_VERSION: u16 = 2;
pub fn new(timestamp: uhlc::Timestamp) -> Self {
Self::from_parameters(timestamp, Default::default())
}
pub fn startup_marker(timestamp: uhlc::Timestamp) -> Self {
Self::from_parameters(
timestamp,
BTreeMap::from([(STARTUP_MARKER_PARAM.to_owned(), Parameter::Bool(true))]),
)
}
pub fn is_startup_marker(&self) -> bool {
get_bool_param(&self.parameters, STARTUP_MARKER_PARAM).unwrap_or(false)
}
pub fn startup_ack(timestamp: uhlc::Timestamp, consumer_node: &str, input_id: &str) -> Self {
Self::from_parameters(
timestamp,
BTreeMap::from([
(STARTUP_ACK_PARAM.to_owned(), Parameter::Bool(true)),
(
STARTUP_ACK_CONSUMER_PARAM.to_owned(),
Parameter::String(consumer_node.to_owned()),
),
(
STARTUP_ACK_INPUT_PARAM.to_owned(),
Parameter::String(input_id.to_owned()),
),
]),
)
}
pub fn startup_ack_identity(&self) -> Option<(&str, &str)> {
if !get_bool_param(&self.parameters, STARTUP_ACK_PARAM).unwrap_or(false) {
return None;
}
let consumer = get_string_param(&self.parameters, STARTUP_ACK_CONSUMER_PARAM)?;
let input = get_string_param(&self.parameters, STARTUP_ACK_INPUT_PARAM)?;
Some((consumer, input))
}
pub fn from_parameters(timestamp: uhlc::Timestamp, parameters: MetadataParameters) -> Self {
Self {
metadata_version: Self::CURRENT_VERSION,
timestamp,
parameters,
}
}
pub fn metadata_version(&self) -> u16 {
self.metadata_version
}
pub fn timestamp(&self) -> uhlc::Timestamp {
self.timestamp
}
pub fn open_telemetry_context(&self) -> String {
get_string_param(&self.parameters, OPEN_TELEMETRY_CONTEXT)
.unwrap_or("")
.to_string()
}
}
pub const STARTUP_MARKER_PARAM: &str = "__dora_startup_marker";
pub const STARTUP_ACK_PARAM: &str = "__dora_startup_ack";
pub const STARTUP_ACK_CONSUMER_PARAM: &str = "__dora_startup_ack_consumer";
pub const STARTUP_ACK_INPUT_PARAM: &str = "__dora_startup_ack_input";
pub type MetadataParameters = BTreeMap<String, Parameter>;
#[derive(Debug, PartialEq, Clone, Serialize, Deserialize)]
pub enum Parameter {
Bool(bool),
Integer(i64),
String(String),
ListInt(Vec<i64>),
Float(f64),
ListFloat(Vec<f64>),
ListString(Vec<String>),
Timestamp(DateTime<Utc>),
}
pub fn get_string_param<'a>(params: &'a MetadataParameters, key: &str) -> Option<&'a str> {
params.get(key).and_then(|p| match p {
Parameter::String(s) => Some(s.as_str()),
_ => None,
})
}
pub fn get_integer_param(params: &MetadataParameters, key: &str) -> Option<i64> {
params.get(key).and_then(|p| match p {
Parameter::Integer(n) => Some(*n),
_ => None,
})
}
pub fn get_bool_param(params: &MetadataParameters, key: &str) -> Option<bool> {
params.get(key).and_then(|p| match p {
Parameter::Bool(b) => Some(*b),
_ => None,
})
}
pub const REQUEST_ID: &str = "request_id";
pub const GOAL_ID: &str = "goal_id";
pub const GOAL_STATUS: &str = "goal_status";
pub const GOAL_STATUS_SUCCEEDED: &str = "succeeded";
pub const GOAL_STATUS_ABORTED: &str = "aborted";
pub const GOAL_STATUS_CANCELED: &str = "canceled";
pub const OPEN_TELEMETRY_CONTEXT: &str = "open_telemetry_context";
pub const SESSION_ID: &str = "session_id";
pub const SEGMENT_ID: &str = "segment_id";
pub const SEQ: &str = "seq";
pub const FIN: &str = "fin";
pub const FLUSH: &str = "flush";
pub const FRAMING: &str = "_framing";
pub const FRAMING_ARROW_IPC: &str = "arrow-ipc";
pub const SCHEMA_HASH: &str = "_schema_hash";
pub const WIRE_SIZE: &str = "_wire_size";
pub fn debug_frame_wire_size(params: &MetadataParameters, data: Option<&[u8]>) -> usize {
get_integer_param(params, WIRE_SIZE)
.and_then(|n| usize::try_from(n).ok())
.or_else(|| data.map(|d| d.len()))
.unwrap_or(0)
}
pub fn carries_pattern_correlation(params: &MetadataParameters) -> bool {
params.contains_key(REQUEST_ID)
|| params.contains_key(GOAL_ID)
|| params.contains_key(GOAL_STATUS)
}
pub fn strip_internal_parameters(params: &mut MetadataParameters) {
params.remove(SCHEMA_HASH);
params.remove(FRAMING);
}
pub fn fnv1a(bytes: &[u8]) -> u64 {
let mut hash: u64 = 0xcbf29ce484222325;
for b in bytes {
hash ^= *b as u64;
hash = hash.wrapping_mul(0x100000001b3);
}
hash
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn fnv1a_matches_standard_vectors() {
assert_eq!(fnv1a(b""), 0xcbf29ce484222325);
assert_eq!(fnv1a(b"a"), 0xaf63dc4c8601ec8c);
}
fn test_timestamp() -> uhlc::Timestamp {
uhlc::HLC::default().new_timestamp()
}
#[test]
fn startup_marker_is_detected_and_survives_the_wire() {
let marker = Metadata::startup_marker(test_timestamp());
assert!(marker.is_startup_marker());
let bytes = crate::encode(&marker).expect("serialize");
let decoded: Metadata = crate::decode(&bytes).expect("deserialize");
assert!(decoded.is_startup_marker());
}
#[test]
fn ordinary_metadata_is_not_a_startup_marker() {
assert!(!Metadata::new(test_timestamp()).is_startup_marker());
let wrong_type = Metadata::from_parameters(
test_timestamp(),
BTreeMap::from([(
STARTUP_MARKER_PARAM.to_owned(),
Parameter::String("true".into()),
)]),
);
assert!(!wrong_type.is_startup_marker());
let explicit_false = Metadata::from_parameters(
test_timestamp(),
BTreeMap::from([(STARTUP_MARKER_PARAM.to_owned(), Parameter::Bool(false))]),
);
assert!(!explicit_false.is_startup_marker());
}
#[test]
fn startup_marker_key_is_reserved_and_distinct() {
assert_eq!(STARTUP_MARKER_PARAM, "__dora_startup_marker");
assert!(STARTUP_MARKER_PARAM.starts_with("__dora_"));
for key in [
REQUEST_ID,
GOAL_ID,
GOAL_STATUS,
SESSION_ID,
SEGMENT_ID,
SEQ,
FIN,
FLUSH,
] {
assert_ne!(STARTUP_MARKER_PARAM, key);
}
}
#[test]
fn startup_ack_round_trips_and_extracts_identity() {
let ack = Metadata::startup_ack(test_timestamp(), "camera-consumer", "image/depth");
assert_eq!(
ack.startup_ack_identity(),
Some(("camera-consumer", "image/depth"))
);
assert!(!ack.is_startup_marker());
let bytes = crate::encode(&ack).expect("serialize");
let decoded: Metadata = crate::decode(&bytes).expect("deserialize");
assert_eq!(
decoded.startup_ack_identity(),
Some(("camera-consumer", "image/depth"))
);
}
#[test]
fn malformed_startup_acks_are_rejected() {
assert_eq!(Metadata::new(test_timestamp()).startup_ack_identity(), None);
assert_eq!(
Metadata::startup_marker(test_timestamp()).startup_ack_identity(),
None
);
let flag_only = Metadata::from_parameters(
test_timestamp(),
BTreeMap::from([(STARTUP_ACK_PARAM.to_owned(), Parameter::Bool(true))]),
);
assert_eq!(flag_only.startup_ack_identity(), None);
let wrong_types = Metadata::from_parameters(
test_timestamp(),
BTreeMap::from([
(STARTUP_ACK_PARAM.to_owned(), Parameter::Bool(true)),
(STARTUP_ACK_CONSUMER_PARAM.to_owned(), Parameter::Integer(1)),
(STARTUP_ACK_INPUT_PARAM.to_owned(), Parameter::Integer(2)),
]),
);
assert_eq!(wrong_types.startup_ack_identity(), None);
let explicit_false = Metadata::from_parameters(
test_timestamp(),
BTreeMap::from([
(STARTUP_ACK_PARAM.to_owned(), Parameter::Bool(false)),
(
STARTUP_ACK_CONSUMER_PARAM.to_owned(),
Parameter::String("c".into()),
),
(
STARTUP_ACK_INPUT_PARAM.to_owned(),
Parameter::String("i".into()),
),
]),
);
assert_eq!(explicit_false.startup_ack_identity(), None);
}
#[test]
fn startup_ack_keys_are_reserved_and_distinct() {
let ack_keys = [
STARTUP_ACK_PARAM,
STARTUP_ACK_CONSUMER_PARAM,
STARTUP_ACK_INPUT_PARAM,
];
for key in ack_keys {
assert!(key.starts_with("__dora_"));
assert_ne!(key, STARTUP_MARKER_PARAM);
}
for (i, a) in ack_keys.iter().enumerate() {
for b in &ack_keys[i + 1..] {
assert_ne!(a, b);
}
}
}
#[test]
fn well_known_keys_have_stable_values() {
assert_eq!(REQUEST_ID, "request_id");
assert_eq!(GOAL_ID, "goal_id");
assert_eq!(GOAL_STATUS, "goal_status");
assert_eq!(GOAL_STATUS_SUCCEEDED, "succeeded");
assert_eq!(GOAL_STATUS_ABORTED, "aborted");
assert_eq!(GOAL_STATUS_CANCELED, "canceled");
assert_eq!(SESSION_ID, "session_id");
assert_eq!(SEGMENT_ID, "segment_id");
assert_eq!(SEQ, "seq");
assert_eq!(FIN, "fin");
assert_eq!(FLUSH, "flush");
}
#[test]
fn well_known_keys_are_distinct() {
let keys = [
REQUEST_ID,
GOAL_ID,
GOAL_STATUS,
SESSION_ID,
SEGMENT_ID,
SEQ,
FIN,
FLUSH,
];
for (i, a) in keys.iter().enumerate() {
for b in &keys[i + 1..] {
assert_ne!(a, b);
}
}
}
#[test]
fn goal_status_values_are_distinct() {
let vals = [
GOAL_STATUS_SUCCEEDED,
GOAL_STATUS_ABORTED,
GOAL_STATUS_CANCELED,
];
for (i, a) in vals.iter().enumerate() {
for b in &vals[i + 1..] {
assert_ne!(a, b);
}
}
}
#[test]
fn outgoing_metadata_is_stamped_with_current_version() {
assert_eq!(Metadata::CURRENT_VERSION, 2);
let ts = uhlc::HLC::default().new_timestamp();
assert_eq!(
Metadata::new(ts).metadata_version(),
Metadata::CURRENT_VERSION
);
assert_eq!(
Metadata::from_parameters(ts, Default::default()).metadata_version(),
Metadata::CURRENT_VERSION
);
}
#[test]
fn get_string_param_extracts_string() {
let mut params = MetadataParameters::default();
params.insert("key".to_string(), Parameter::String("value".to_string()));
assert_eq!(get_string_param(¶ms, "key"), Some("value"));
assert_eq!(get_string_param(¶ms, "missing"), None);
}
#[test]
fn get_string_param_returns_none_for_non_string() {
let mut params = MetadataParameters::default();
params.insert("num".to_string(), Parameter::Integer(42));
assert_eq!(get_string_param(¶ms, "num"), None);
}
#[test]
fn get_integer_param_extracts_integer() {
let mut params = MetadataParameters::default();
params.insert("key".to_string(), Parameter::Integer(42));
assert_eq!(get_integer_param(¶ms, "key"), Some(42));
assert_eq!(get_integer_param(¶ms, "missing"), None);
}
#[test]
fn get_integer_param_returns_none_for_non_integer() {
let mut params = MetadataParameters::default();
params.insert("s".to_string(), Parameter::String("hello".to_string()));
assert_eq!(get_integer_param(¶ms, "s"), None);
}
#[test]
fn get_bool_param_extracts_bool() {
let mut params = MetadataParameters::default();
params.insert("key".to_string(), Parameter::Bool(true));
assert_eq!(get_bool_param(¶ms, "key"), Some(true));
assert_eq!(get_bool_param(¶ms, "missing"), None);
}
#[test]
fn get_bool_param_returns_none_for_non_bool() {
let mut params = MetadataParameters::default();
params.insert("n".to_string(), Parameter::Integer(1));
assert_eq!(get_bool_param(¶ms, "n"), None);
}
#[test]
fn debug_frame_wire_size_prefers_stamped_value() {
let mut params = MetadataParameters::default();
params.insert(WIRE_SIZE.to_string(), Parameter::Integer(17));
assert_eq!(debug_frame_wire_size(¶ms, Some(&[0u8; 42])), 17);
}
#[test]
fn debug_frame_wire_size_falls_back_to_buffer_len() {
let params = MetadataParameters::default();
assert_eq!(debug_frame_wire_size(¶ms, Some(&[0u8; 42])), 42);
}
#[test]
fn debug_frame_wire_size_falls_back_on_out_of_range_stamp() {
let mut params = MetadataParameters::default();
params.insert(WIRE_SIZE.to_string(), Parameter::Integer(-1));
assert_eq!(debug_frame_wire_size(¶ms, Some(&[0u8; 42])), 42);
}
#[test]
fn debug_frame_wire_size_zero_without_stamp_or_buffer() {
assert_eq!(
debug_frame_wire_size(&MetadataParameters::default(), None),
0
);
}
}