use crate::engine::error::{DataflowError, ErrorInfo};
use chrono::{DateTime, Utc};
use datavalue::OwnedDataValue;
use serde::{Deserialize, Serialize};
use serde_json::Value as JsonValue;
use std::sync::Arc;
use uuid::Uuid;
#[derive(Clone)]
pub(crate) enum MessageId {
Uuid([u8; 36]),
Custom(String),
}
impl MessageId {
fn new_uuid_v7() -> Self {
let mut buf = [0u8; 36];
Uuid::now_v7().hyphenated().encode_lower(&mut buf);
Self::Uuid(buf)
}
fn as_str(&self) -> &str {
match self {
Self::Uuid(buf) => std::str::from_utf8(buf).expect("uuid ids are ascii"),
Self::Custom(s) => s,
}
}
}
impl std::fmt::Debug for MessageId {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
std::fmt::Debug::fmt(self.as_str(), f)
}
}
#[derive(Debug, Clone)]
pub struct Message {
pub(crate) id: MessageId,
pub(crate) payload: Arc<OwnedDataValue>,
pub context: OwnedDataValue,
pub(crate) audit_trail: Vec<AuditTrail>,
pub(crate) errors: Vec<ErrorInfo>,
pub(crate) capture_changes: bool,
pub(crate) routing_bucket: Option<u8>,
}
impl Serialize for Message {
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
where
S: serde::Serializer,
{
use serde::ser::SerializeStruct;
let mut state = serializer.serialize_struct("Message", 5)?;
state.serialize_field("id", &self.id.as_str())?;
state.serialize_field("payload", &self.payload)?;
state.serialize_field("context", &self.context)?;
state.serialize_field("audit_trail", &self.audit_trail)?;
state.serialize_field("errors", &self.errors)?;
state.end()
}
}
impl<'de> Deserialize<'de> for Message {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: serde::Deserializer<'de>,
{
#[derive(Deserialize)]
struct MessageData {
id: String,
payload: Arc<OwnedDataValue>,
context: OwnedDataValue,
audit_trail: Vec<AuditTrail>,
errors: Vec<ErrorInfo>,
}
let data = MessageData::deserialize(deserializer)?;
Ok(Message {
id: MessageId::Custom(data.id),
payload: data.payload,
context: data.context,
audit_trail: data.audit_trail,
errors: data.errors,
capture_changes: true,
routing_bucket: None,
})
}
}
impl Message {
pub fn builder() -> MessageBuilder {
MessageBuilder::new()
}
pub fn new(payload: Arc<OwnedDataValue>) -> Self {
Self {
id: MessageId::new_uuid_v7(),
payload,
context: empty_context(),
audit_trail: vec![],
errors: vec![],
capture_changes: true,
routing_bucket: None,
}
}
pub fn from_value(payload: &JsonValue) -> Self {
Self::new(Arc::new(OwnedDataValue::from(payload)))
}
pub fn from_json_str(payload: &str) -> crate::engine::error::Result<Self> {
let value: JsonValue = serde_json::from_str(payload).map_err(DataflowError::from_serde)?;
Ok(Self::from_value(&value))
}
pub fn add_error(&mut self, error: ErrorInfo) {
self.errors.push(error);
}
pub fn has_errors(&self) -> bool {
!self.errors.is_empty()
}
#[inline]
pub fn id(&self) -> &str {
self.id.as_str()
}
#[inline]
pub fn payload(&self) -> &OwnedDataValue {
&self.payload
}
#[inline]
pub fn payload_arc(&self) -> &Arc<OwnedDataValue> {
&self.payload
}
#[inline]
pub fn audit_trail(&self) -> &[AuditTrail] {
&self.audit_trail
}
#[inline]
pub fn errors(&self) -> &[ErrorInfo] {
&self.errors
}
#[inline]
pub fn capture_changes(&self) -> bool {
self.capture_changes
}
#[inline]
pub fn routing_bucket(&self) -> Option<u8> {
self.routing_bucket
}
pub fn data(&self) -> &OwnedDataValue {
&self.context["data"]
}
pub fn metadata(&self) -> &OwnedDataValue {
&self.context["metadata"]
}
pub fn temp_data(&self) -> &OwnedDataValue {
&self.context["temp_data"]
}
}
#[must_use = "MessageBuilder must be `.build()` to produce a Message"]
#[derive(Default)]
pub struct MessageBuilder {
id: Option<String>,
payload: Option<Arc<OwnedDataValue>>,
capture_changes: Option<bool>,
data: Option<OwnedDataValue>,
metadata: Option<OwnedDataValue>,
temp_data: Option<OwnedDataValue>,
routing_bucket: Option<u8>,
}
impl MessageBuilder {
pub fn new() -> Self {
Self::default()
}
pub fn id(mut self, id: impl Into<String>) -> Self {
self.id = Some(id.into());
self
}
pub fn payload(mut self, payload: Arc<OwnedDataValue>) -> Self {
self.payload = Some(payload);
self
}
pub fn payload_json(mut self, payload: &JsonValue) -> Self {
self.payload = Some(Arc::new(OwnedDataValue::from(payload)));
self
}
pub fn data(mut self, data: OwnedDataValue) -> Self {
if data.is_object() {
self.data = Some(data);
}
self
}
pub fn data_json(self, data: &JsonValue) -> Self {
self.data(OwnedDataValue::from(data))
}
pub fn metadata(mut self, metadata: OwnedDataValue) -> Self {
if metadata.is_object() {
self.metadata = Some(metadata);
}
self
}
pub fn metadata_json(self, metadata: &JsonValue) -> Self {
self.metadata(OwnedDataValue::from(metadata))
}
pub fn temp_data(mut self, temp_data: OwnedDataValue) -> Self {
if temp_data.is_object() {
self.temp_data = Some(temp_data);
}
self
}
pub fn temp_data_json(self, temp_data: &JsonValue) -> Self {
self.temp_data(OwnedDataValue::from(temp_data))
}
pub fn routing_bucket(mut self, bucket: u8) -> Self {
self.routing_bucket = Some(bucket.min(99));
self
}
pub fn capture_changes(mut self, on: bool) -> Self {
self.capture_changes = Some(on);
self
}
pub fn build(self) -> Message {
Message {
id: self
.id
.map(MessageId::Custom)
.unwrap_or_else(MessageId::new_uuid_v7),
payload: self
.payload
.unwrap_or_else(|| Arc::new(OwnedDataValue::Null)),
context: context_from(self.data, self.metadata, self.temp_data),
audit_trail: vec![],
errors: vec![],
capture_changes: self.capture_changes.unwrap_or(true),
routing_bucket: self.routing_bucket,
}
}
}
fn context_from(
data: Option<OwnedDataValue>,
metadata: Option<OwnedDataValue>,
temp_data: Option<OwnedDataValue>,
) -> OwnedDataValue {
fn slot(v: Option<OwnedDataValue>) -> OwnedDataValue {
v.unwrap_or_else(|| OwnedDataValue::Object(Vec::new()))
}
OwnedDataValue::Object(vec![
("data".to_string(), slot(data)),
("metadata".to_string(), slot(metadata)),
("temp_data".to_string(), slot(temp_data)),
])
}
fn empty_context() -> OwnedDataValue {
context_from(None, None, None)
}
#[derive(Serialize, Deserialize, Debug, Clone)]
pub struct AuditTrail {
pub workflow_id: Arc<str>,
pub task_id: Arc<str>,
pub timestamp: DateTime<Utc>,
pub changes: Vec<Change>,
pub status: usize,
}
#[derive(Serialize, Deserialize, Debug, Clone)]
pub struct Change {
pub path: Arc<str>,
pub old_value: OwnedDataValue,
pub new_value: OwnedDataValue,
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn from_json_str_parses_valid_payload() {
let msg =
Message::from_json_str(r#"{"order": {"total": 42}}"#).expect("valid JSON should parse");
let payload_json = serde_json::to_value(msg.payload()).unwrap();
assert_eq!(payload_json, serde_json::json!({"order": {"total": 42}}));
}
#[test]
fn from_json_str_rejects_malformed_payload() {
let err = Message::from_json_str("{ not json").expect_err("malformed input should fail");
assert!(matches!(err, DataflowError::Deserialization(_)));
}
#[test]
fn builder_with_no_seed_matches_the_historical_empty_context() {
let m = Message::builder().build();
let v = serde_json::to_value(&m).unwrap();
let ctx = v["context"].as_object().unwrap();
assert_eq!(
ctx.keys().collect::<Vec<_>>(),
vec!["data", "metadata", "temp_data"]
);
assert_eq!(
v["context"],
serde_json::json!({
"data": {}, "metadata": {}, "temp_data": {}
})
);
}
#[test]
fn each_setter_lands_in_its_own_root_field() {
let m = Message::builder()
.data_json(&serde_json::json!({"d": 1}))
.build();
assert_eq!(
serde_json::Value::from(m.data()),
serde_json::json!({"d": 1})
);
assert_eq!(serde_json::Value::from(m.metadata()), serde_json::json!({}));
assert_eq!(
serde_json::Value::from(m.temp_data()),
serde_json::json!({})
);
let m = Message::builder()
.metadata_json(&serde_json::json!({"m": 1}))
.build();
assert_eq!(
serde_json::Value::from(m.metadata()),
serde_json::json!({"m": 1})
);
assert_eq!(serde_json::Value::from(m.data()), serde_json::json!({}));
let m = Message::builder()
.temp_data_json(&serde_json::json!({"t": 1}))
.build();
assert_eq!(
serde_json::Value::from(m.temp_data()),
serde_json::json!({"t": 1})
);
assert_eq!(serde_json::Value::from(m.data()), serde_json::json!({}));
}
#[test]
fn the_owned_and_json_setter_forms_agree() {
let v = serde_json::json!({"a": {"b": [1, 2]}});
let via_json = Message::builder().data_json(&v).build();
let via_owned = Message::builder().data(OwnedDataValue::from(&v)).build();
assert_eq!(via_json.context, via_owned.context);
}
#[test]
fn seeding_records_no_audit_entry_or_change() {
for capture in [true, false] {
let m = Message::builder()
.capture_changes(capture)
.data_json(&serde_json::json!({"d": 1}))
.build();
assert!(
m.audit_trail().is_empty(),
"seeding is initial state, not a mutation"
);
assert_eq!(
serde_json::Value::from(m.data()),
serde_json::json!({"d": 1})
);
}
}
#[test]
fn calling_a_setter_twice_keeps_the_last_value() {
let m = Message::builder()
.data_json(&serde_json::json!({"first": 1}))
.data_json(&serde_json::json!({"second": 2}))
.build();
assert_eq!(
serde_json::Value::from(m.data()),
serde_json::json!({"second": 2})
);
}
#[test]
fn an_empty_object_seed_is_indistinguishable_from_no_seed() {
let seeded = Message::builder().data_json(&serde_json::json!({})).build();
let bare = Message::builder().build();
assert_eq!(seeded.context, bare.context);
}
#[test]
fn seed_keys_are_literal_not_paths() {
use crate::engine::utils::get_nested_value;
let m = Message::builder()
.metadata_json(&serde_json::json!({"a.b": 1}))
.build();
assert_eq!(
serde_json::Value::from(m.metadata()),
serde_json::json!({"a.b": 1})
);
assert!(
get_nested_value(&m.context, "metadata.a.b").is_none(),
"a dotted key must not become a nested path"
);
let m = Message::builder()
.metadata_json(&serde_json::json!({"#20": 1}))
.build();
assert_eq!(
serde_json::Value::from(m.metadata()),
serde_json::json!({"#20": 1})
);
let m = Message::builder()
.data_json(&serde_json::json!({"0": 1}))
.build();
assert!(matches!(m.data(), OwnedDataValue::Object(_)));
assert_eq!(
serde_json::Value::from(m.data()),
serde_json::json!({"0": 1})
);
}
#[test]
fn non_ascii_keys_and_values_round_trip() {
let v = serde_json::json!({"régión": "東京"});
let m = Message::builder().data_json(&v).build();
assert_eq!(serde_json::Value::from(m.data()), v);
assert_eq!(serde_json::to_value(&m).unwrap()["context"]["data"], v);
}
#[test]
fn nested_seeds_resolve_through_the_path_api() {
use crate::engine::utils::get_nested_value;
let m = Message::builder()
.data_json(&serde_json::json!({"order": {"items": [1, 2]}}))
.build();
assert_eq!(
get_nested_value(&m.context, "data.order.items.1"),
Some(&OwnedDataValue::from(&serde_json::json!(2)))
);
}
#[test]
fn non_object_seeds_are_ignored_to_preserve_the_context_invariant() {
for bad in [
serde_json::json!("scalar"),
serde_json::json!([1, 2]),
serde_json::json!(null),
serde_json::json!(7),
] {
let m = Message::builder()
.data_json(&bad)
.metadata_json(&bad)
.temp_data_json(&bad)
.build();
assert_eq!(
serde_json::Value::from(&m.context),
serde_json::json!({"data": {}, "metadata": {}, "temp_data": {}}),
"non-object seed {bad} must be ignored"
);
}
}
}