#![allow(clippy::doc_markdown)]
use serde::{Deserialize, Serialize};
use time::OffsetDateTime;
use uuid::Uuid;
#[derive(Debug, Clone, Deserialize)]
pub struct InboundMakoEvent {
pub specversion: String,
pub id: String,
pub source: String,
#[serde(rename = "type")]
pub ce_type: String,
#[serde(with = "time::serde::rfc3339")]
pub time: OffsetDateTime,
pub subject: Option<String>,
pub dataschema: Option<String>,
pub datacontenttype: Option<String>,
#[serde(default)]
pub makopid: Option<u32>,
#[serde(default)]
pub makoconvid: Option<String>,
#[serde(default)]
pub makocausationid: Option<String>,
#[serde(default)]
pub makofailreason: Option<String>,
#[serde(default)]
pub makoworkflow: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub traceparent: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub tracestate: Option<String>,
pub data: serde_json::Value,
}
impl InboundMakoEvent {
#[must_use]
pub fn process_id(&self) -> Option<Uuid> {
self.subject.as_deref().and_then(|s| s.parse().ok())
}
#[must_use]
pub fn conv_id(&self) -> Option<Uuid> {
self.makoconvid.as_deref().and_then(|s| s.parse().ok())
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct MarktEvent {
pub specversion: String,
pub id: String,
pub source: String,
#[serde(rename = "type")]
pub ce_type: String,
#[serde(with = "time::serde::rfc3339")]
pub time: OffsetDateTime,
pub subject: String,
pub datacontenttype: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub marktmaloid: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub marktmeloid: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub marktcontractid: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub marktrole: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub markterpref: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub makoconvid: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub makopid: Option<u32>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub makoworkflow: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub makoerc: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub traceparent: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub tracestate: Option<String>,
pub data: serde_json::Value,
}
impl MarktEvent {
#[must_use]
pub fn new(
tenant: &str,
ce_type: impl Into<String>,
subject: impl Into<String>,
data: serde_json::Value,
) -> Self {
Self {
specversion: "1.0".into(),
id: Uuid::new_v4().to_string(),
source: format!("urn:markt:tenant:{tenant}"),
ce_type: ce_type.into(),
time: OffsetDateTime::now_utc(),
subject: subject.into(),
datacontenttype: "application/json".into(),
marktmaloid: None,
marktmeloid: None,
marktcontractid: None,
marktrole: None,
markterpref: None,
makoconvid: None,
makopid: None,
makoworkflow: None,
makoerc: None,
traceparent: None,
tracestate: None,
data,
}
}
#[must_use]
pub fn with_extensions(mut self, ext: EventExtensions) -> Self {
self.marktmaloid = ext.marktmaloid;
self.marktmeloid = ext.marktmeloid;
self.marktcontractid = ext.marktcontractid;
self.marktrole = ext.marktrole;
self.markterpref = ext.markterpref;
self.makoconvid = ext.makoconvid;
self.makopid = ext.makopid;
self.makoworkflow = ext.makoworkflow;
self.makoerc = ext.makoerc;
self.traceparent = ext.traceparent;
self.tracestate = ext.tracestate;
self
}
}
#[derive(Debug, Default, Clone)]
pub struct EventExtensions {
pub marktmaloid: Option<String>,
pub marktmeloid: Option<String>,
pub marktcontractid: Option<String>,
pub marktrole: Option<String>,
pub markterpref: Option<String>,
pub makoconvid: Option<String>,
pub makopid: Option<u32>,
pub makoworkflow: Option<String>,
pub makoerc: Option<String>,
pub traceparent: Option<String>,
pub tracestate: Option<String>,
}
#[must_use]
pub fn compute_signature(secret: &[u8], body: &[u8]) -> String {
use hmac::{Hmac, Mac};
use sha2::Sha256;
let mut mac = <Hmac<Sha256>>::new_from_slice(secret).expect("HMAC accepts any key length");
mac.update(body);
let result = mac.finalize().into_bytes();
hex::bytes_to_hex_str(&result)
}
#[must_use]
pub fn verify_signature(secret: &[u8], body: &[u8], provided_hex: &str) -> bool {
use subtle::ConstantTimeEq;
let expected = compute_signature(secret, body);
expected.as_bytes().ct_eq(provided_hex.as_bytes()).into()
}
mod hex {
const HEX: &[u8; 16] = b"0123456789abcdef";
pub(super) fn bytes_to_hex_str(bytes: &[u8]) -> String {
let mut s = String::with_capacity(bytes.len() * 2);
for &b in bytes {
s.push(HEX[(b >> 4) as usize] as char);
s.push(HEX[(b & 0xf) as usize] as char);
}
s
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn hmac_roundtrip() {
let secret = b"test-secret";
let body = b"{\"type\":\"de.mako.process.completed\"}";
let sig = compute_signature(secret, body);
assert!(verify_signature(secret, body, &sig));
assert!(!verify_signature(secret, body, "deadbeef"));
}
#[test]
fn makoworkflow_deserializes_from_cloudevent_json() {
let json = r#"{
"specversion": "1.0",
"id": "evt-001",
"source": "urn:mako:tenant:9900000000001",
"type": "de.mako.gpke.lieferbeginn.completed",
"time": "2025-10-01T10:00:00Z",
"makopid": 55003,
"makoconvid": "a0000000-0000-0000-0000-000000000001",
"makocausationid": "b0000000-0000-0000-0000-000000000002",
"makoworkflow": "gpke-supplier-change",
"data": {}
}"#;
let event: InboundMakoEvent = serde_json::from_str(json).unwrap();
assert_eq!(event.makoworkflow.as_deref(), Some("gpke-supplier-change"));
}
#[test]
fn makoworkflow_absent_deserializes_to_none() {
let json = r#"{
"specversion": "1.0",
"id": "evt-002",
"source": "urn:mako:tenant:9900000000001",
"type": "de.mako.gpke.lieferbeginn.completed",
"time": "2025-10-01T10:00:00Z",
"data": {}
}"#;
let event: InboundMakoEvent = serde_json::from_str(json).unwrap();
assert_eq!(event.makoworkflow, None);
}
#[test]
fn markt_event_source_prefix() {
let event = MarktEvent::new(
"9900000000001",
"de.markt.malo.updated",
"DE000000000001",
serde_json::json!({}),
);
assert!(event.source.starts_with("urn:markt:tenant:"));
}
#[test]
fn markt_event_with_extensions_sets_marktrole() {
let event = MarktEvent::new(
"9900000000001",
"de.markt.malo.updated",
"DE000000000001",
serde_json::json!({}),
)
.with_extensions(EventExtensions {
marktrole: Some("NB".into()),
..Default::default()
});
assert_eq!(event.marktrole.as_deref(), Some("NB"));
let json = serde_json::to_value(&event).unwrap();
assert_eq!(json["marktrole"], "NB");
assert!(
json.get("mdmrole").is_none(),
"must not emit legacy mdmrole field"
);
}
#[test]
fn markt_event_roundtrip_json() {
let orig = MarktEvent::new(
"9900000000001",
"de.markt.versorgung.beliefert",
"51238696780",
serde_json::json!({"lieferstatus": "Beliefert"}),
)
.with_extensions(EventExtensions {
marktmaloid: Some("51238696780".into()),
marktrole: Some("LF".into()),
makopid: Some(55003),
..Default::default()
});
let json = serde_json::to_string(&orig).unwrap();
let back: MarktEvent = serde_json::from_str(&json).unwrap();
assert_eq!(back.marktmaloid, orig.marktmaloid);
assert_eq!(back.marktrole, orig.marktrole);
assert_eq!(back.makopid, orig.makopid);
assert_eq!(back.ce_type, "de.markt.versorgung.beliefert");
}
}