use std::collections::BTreeMap;
use serde::{Deserialize, Serialize, ser::SerializeMap};
use serde_json::Value;
use crate::core::{InboundEvent, Timestamp};
pub const SPEC_VERSION: &str = "1.0";
pub const CONTENT_TYPE: &str = "application/cloudevents+json; charset=UTF-8";
pub const HEADER_PREFIX: &str = "ce-";
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
#[non_exhaustive]
pub enum CloudEventError {
#[error("not a JSON object: {0}")]
Malformed(String),
#[error(
"specversion '{0}' is not '{expected}' — this plane does not guess at an \
envelope it has not been written against",
expected = SPEC_VERSION
)]
UnknownSpecVersion(String),
#[error("CloudEvents requires a non-empty '{0}'")]
MissingAttribute(&'static str),
#[error(
"extension attribute '{0}' is not a CloudEvents attribute name: names are \
lowercase letters and digits, so an event that travels over a binding \
with case-insensitive keys cannot arrive as a different event"
)]
BadExtensionName(String),
#[error(
"extension attribute '{0}' names a core attribute — an extension that \
shadows 'id' or 'source' would let the serializer emit an envelope \
whose identity is the extension's value"
)]
ReservedExtensionName(String),
#[error(
"extension attribute '{0}' carries a JSON {1}, which the CloudEvents \
type system has no extension type for — send a string, number, or \
boolean, or put structure in 'data'"
)]
BadExtensionValue(String, &'static str),
#[error(
"'{0}' contains a control character — the deduplication identity is \
'source' and 'id' joined by one, so a value that embeds it could \
spell another producer's pair"
)]
ControlCharacter(&'static str),
#[error(
"header '{0}' appears more than once — two values for one attribute \
are two different events wearing one envelope"
)]
DuplicateHeader(String),
#[error("'time' is not an RFC 3339 timestamp: {0}")]
BadTime(String),
#[error(
"'data_base64' carries bytes, and an event payload here is a JSON value — \
send the data as JSON or address the bytes through a blob reference"
)]
BinaryData,
#[error(
"datacontenttype '{0}' is not JSON, and the body of a binary-mode event \
becomes a JSON payload — wrapping other bytes in a string would retype \
them silently"
)]
UnsupportedDataContentType(String),
#[error("a percent-encoded header value is not valid UTF-8: {0}")]
BadHeaderEncoding(String),
}
#[derive(Debug, Clone, PartialEq, Deserialize)]
#[serde(try_from = "WireCloudEvent")]
pub struct CloudEvent {
id: String,
source: String,
event_type: String,
subject: Option<String>,
time: Option<Timestamp>,
datacontenttype: Option<String>,
dataschema: Option<String>,
extensions: BTreeMap<String, Value>,
data: Option<Value>,
}
impl CloudEvent {
pub fn new(
source: impl Into<String>,
id: impl Into<String>,
event_type: impl Into<String>,
) -> Result<Self, CloudEventError> {
let event = Self {
id: id.into(),
source: source.into(),
event_type: event_type.into(),
subject: None,
time: None,
datacontenttype: None,
dataschema: None,
extensions: BTreeMap::new(),
data: None,
};
event.checked()
}
fn checked(self) -> Result<Self, CloudEventError> {
if self.id.is_empty() {
return Err(CloudEventError::MissingAttribute("id"));
}
if self.source.is_empty() {
return Err(CloudEventError::MissingAttribute("source"));
}
if self.event_type.is_empty() {
return Err(CloudEventError::MissingAttribute("type"));
}
for (name, value) in [
("id", &self.id),
("source", &self.source),
("type", &self.event_type),
] {
if value.chars().any(char::is_control) {
return Err(CloudEventError::ControlCharacter(name));
}
}
for (name, value) in &self.extensions {
if name.is_empty()
|| !name
.bytes()
.all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit())
{
return Err(CloudEventError::BadExtensionName(name.clone()));
}
if matches!(
name.as_str(),
"specversion"
| "id"
| "source"
| "type"
| "subject"
| "time"
| "datacontenttype"
| "dataschema"
| "data"
| "data_base64"
) {
return Err(CloudEventError::ReservedExtensionName(name.clone()));
}
match value {
Value::String(_) | Value::Bool(_) | Value::Number(_) => {}
Value::Null => {
return Err(CloudEventError::BadExtensionValue(name.clone(), "null"));
}
Value::Array(_) => {
return Err(CloudEventError::BadExtensionValue(name.clone(), "array"));
}
Value::Object(_) => {
return Err(CloudEventError::BadExtensionValue(name.clone(), "object"));
}
}
}
Ok(self)
}
#[must_use]
pub fn with_data(mut self, data: Value) -> Self {
self.data = Some(data);
self.datacontenttype = Some("application/json".to_owned());
self
}
#[must_use]
pub fn with_subject(mut self, subject: impl Into<String>) -> Self {
self.subject = Some(subject.into());
self
}
pub fn with_extension(
mut self,
name: impl Into<String>,
value: Value,
) -> Result<Self, CloudEventError> {
self.extensions.insert(name.into(), value);
self.checked()
}
#[must_use]
pub fn id(&self) -> &str {
&self.id
}
#[must_use]
pub fn source(&self) -> &str {
&self.source
}
#[must_use]
pub fn event_type(&self) -> &str {
&self.event_type
}
#[must_use]
pub fn subject(&self) -> Option<&str> {
self.subject.as_deref()
}
#[must_use]
pub const fn time(&self) -> Option<Timestamp> {
self.time
}
#[must_use]
pub fn data(&self) -> Option<&Value> {
self.data.as_ref()
}
#[must_use]
pub fn extension(&self, name: &str) -> Option<&Value> {
self.extensions.get(name)
}
#[must_use]
pub fn origin_id(&self) -> String {
crate::core::origin_key(&self.source, &self.id)
}
#[must_use]
pub fn into_inbound(self, transport_source: impl Into<String>) -> InboundEvent {
let mut event = InboundEvent::new(
transport_source,
self.origin_id(),
self.event_type.clone(),
self.data.unwrap_or(Value::Null),
);
if let Some(subject) = self.subject {
event = event.correlate(crate::core::CorrelationKey::new("subject", subject));
}
event
}
pub fn from_json(body: &[u8]) -> Result<Self, CloudEventError> {
let wire: WireCloudEvent =
serde_json::from_slice(body).map_err(|e| CloudEventError::Malformed(e.to_string()))?;
Self::try_from(wire)
}
pub fn from_http<'a, I>(headers: I, body: &[u8]) -> Result<Self, CloudEventError>
where
I: IntoIterator<Item = (&'a str, &'a str)>,
{
let mut content_type = None;
let mut attributes: BTreeMap<String, String> = BTreeMap::new();
for (name, value) in headers {
let name = name.to_ascii_lowercase();
if name == "content-type" {
content_type = Some(value.to_owned());
} else if let Some(attribute) = name.strip_prefix(HEADER_PREFIX) {
if attributes
.insert(attribute.to_owned(), percent_decode(value)?)
.is_some()
{
return Err(CloudEventError::DuplicateHeader(name));
}
}
}
if content_type
.as_deref()
.is_some_and(is_structured_media_type)
{
return Self::from_json(body);
}
Self::from_binary(&attributes, content_type.as_deref(), body)
}
fn from_binary(
attributes: &BTreeMap<String, String>,
content_type: Option<&str>,
body: &[u8],
) -> Result<Self, CloudEventError> {
let take = |name: &str| attributes.get(name).cloned();
let specversion = take("specversion").unwrap_or_default();
if specversion != SPEC_VERSION {
return Err(CloudEventError::UnknownSpecVersion(specversion));
}
let data = if body.is_empty() {
None
} else {
let media = content_type.map_or_else(|| "application/json".to_owned(), media_type_of);
if !is_json_media_type(&media) {
return Err(CloudEventError::UnsupportedDataContentType(media));
}
Some(
serde_json::from_slice(body)
.map_err(|e| CloudEventError::Malformed(e.to_string()))?,
)
};
let time = take("time").map(|raw| parse_time(&raw)).transpose()?;
let known = [
"specversion",
"id",
"source",
"type",
"subject",
"time",
"dataschema",
"datacontenttype",
];
let extensions = attributes
.iter()
.filter(|(name, _)| !known.contains(&name.as_str()))
.map(|(name, value)| (name.clone(), Value::String(value.clone())))
.collect();
Self {
id: take("id").unwrap_or_default(),
source: take("source").unwrap_or_default(),
event_type: take("type").unwrap_or_default(),
subject: take("subject"),
time,
datacontenttype: content_type.map(media_type_of),
dataschema: take("dataschema"),
extensions,
data,
}
.checked()
}
#[must_use]
pub fn to_bytes(&self) -> Vec<u8> {
crate::core::canon::value_bytes(&self.to_value())
}
#[must_use]
pub fn to_value(&self) -> Value {
let mut map = serde_json::Map::new();
map.insert("specversion".to_owned(), Value::String(SPEC_VERSION.into()));
map.insert("id".to_owned(), Value::String(self.id.clone()));
map.insert("source".to_owned(), Value::String(self.source.clone()));
map.insert("type".to_owned(), Value::String(self.event_type.clone()));
for (name, value) in [
("subject", self.subject.as_ref()),
("datacontenttype", self.datacontenttype.as_ref()),
("dataschema", self.dataschema.as_ref()),
] {
if let Some(value) = value {
map.insert(name.to_owned(), Value::String(value.clone()));
}
}
if let Some(time) = self.time {
map.insert(
"time".to_owned(),
Value::String(crate::core::format_timestamp(time)),
);
}
for (name, value) in &self.extensions {
map.insert(name.clone(), value.clone());
}
if let Some(data) = &self.data {
map.insert("data".to_owned(), data.clone());
}
Value::Object(map)
}
}
impl Serialize for CloudEvent {
fn serialize<S: serde::Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
let value = self.to_value();
let object = value.as_object().expect("to_value builds an object");
let mut map = serializer.serialize_map(Some(object.len()))?;
for (key, value) in object {
map.serialize_entry(key, value)?;
}
map.end()
}
}
#[derive(Deserialize)]
struct WireCloudEvent {
specversion: String,
#[serde(default)]
id: String,
#[serde(default)]
source: String,
#[serde(default, rename = "type")]
event_type: String,
#[serde(default)]
subject: Option<String>,
#[serde(default)]
time: Option<String>,
#[serde(default)]
datacontenttype: Option<String>,
#[serde(default)]
dataschema: Option<String>,
#[serde(default)]
data: Option<Value>,
#[serde(default)]
data_base64: Option<String>,
#[serde(flatten)]
extensions: BTreeMap<String, Value>,
}
impl TryFrom<WireCloudEvent> for CloudEvent {
type Error = CloudEventError;
fn try_from(wire: WireCloudEvent) -> Result<Self, Self::Error> {
if wire.specversion != SPEC_VERSION {
return Err(CloudEventError::UnknownSpecVersion(wire.specversion));
}
if wire.data_base64.is_some() {
return Err(CloudEventError::BinaryData);
}
let time = wire.time.map(|raw| parse_time(&raw)).transpose()?;
Self {
id: wire.id,
source: wire.source,
event_type: wire.event_type,
subject: wire.subject,
time,
datacontenttype: wire.datacontenttype,
dataschema: wire.dataschema,
extensions: wire.extensions,
data: wire.data,
}
.checked()
}
}
fn parse_time(raw: &str) -> Result<Timestamp, CloudEventError> {
Timestamp::parse(raw, &time::format_description::well_known::Rfc3339)
.map_err(|error| CloudEventError::BadTime(format!("{raw}: {error}")))
}
fn media_type_of(header: &str) -> String {
header
.split(';')
.next()
.unwrap_or(header)
.trim()
.to_ascii_lowercase()
}
#[must_use]
pub fn is_structured_media_type(header: &str) -> bool {
media_type_of(header) == "application/cloudevents+json"
}
fn is_json_media_type(media: &str) -> bool {
media.eq_ignore_ascii_case("application/json")
|| media.eq_ignore_ascii_case("text/json")
|| media.to_ascii_lowercase().ends_with("+json")
}
fn percent_decode(value: &str) -> Result<String, CloudEventError> {
if !value.contains('%') {
return Ok(value.to_owned());
}
let bytes = value.as_bytes();
let mut out = Vec::with_capacity(bytes.len());
let mut index = 0;
while index < bytes.len() {
if bytes[index] == b'%' && index + 2 < bytes.len() {
let hex = std::str::from_utf8(&bytes[index + 1..index + 3])
.ok()
.and_then(|pair| u8::from_str_radix(pair, 16).ok());
if let Some(byte) = hex {
out.push(byte);
index += 3;
continue;
}
}
out.push(bytes[index]);
index += 1;
}
String::from_utf8(out).map_err(|error| CloudEventError::BadHeaderEncoding(error.to_string()))
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_control_character_cannot_forge_another_producers_pair() {
for (source, id) in [("a\u{1f}b", "c"), ("a", "b\u{1f}c"), ("a\nb", "c")] {
let error = CloudEvent::new(id, source, "t").expect_err("a forgeable pair");
assert!(
matches!(error, CloudEventError::ControlCharacter(_)),
"{source:?}/{id:?} was accepted: {error}"
);
}
}
#[test]
fn an_extension_cannot_shadow_a_core_attribute() {
let event = CloudEvent::new("1", "urn:a", "t").expect("event");
let error = event
.with_extension("id", serde_json::Value::String("other".to_owned()))
.expect_err("a shadowing extension");
assert!(matches!(error, CloudEventError::ReservedExtensionName(_)));
}
#[test]
fn a_structured_extension_value_must_be_a_scalar() {
let error = structured(&json!({
"specversion": "1.0", "id": "1", "source": "urn:a", "type": "t",
"tenantid": { "nested": true },
}))
.expect_err("an object extension");
assert!(
matches!(error, CloudEventError::BadExtensionValue(..)),
"{error}"
);
}
#[test]
fn structured_mode_is_chosen_case_insensitively() {
let body = json!({ "specversion": "1.0", "id": "1", "source": "urn:a", "type": "t" });
let event = CloudEvent::from_http(
[(
"content-type",
"Application/CloudEvents+JSON; charset=utf-8",
)],
body.to_string().as_bytes(),
)
.expect("structured mode despite capitalization");
assert_eq!(event.id(), "1");
assert!(is_structured_media_type("APPLICATION/CLOUDEVENTS+JSON"));
assert!(!is_structured_media_type("application/cloudevents+jsonx"));
}
#[test]
fn a_duplicated_attribute_header_is_refused() {
let error = CloudEvent::from_http(
[
("ce-specversion", "1.0"),
("ce-id", "1"),
("ce-id", "2"),
("ce-source", "urn:a"),
("ce-type", "t"),
],
b"",
)
.expect_err("two ids");
assert!(
matches!(error, CloudEventError::DuplicateHeader(_)),
"{error}"
);
}
#[test]
fn the_subject_becomes_the_correlation_key_and_nothing_else_does() {
let with = CloudEvent::new("1", "urn:a", "t")
.expect("event")
.with_subject("order-9")
.into_inbound("peer:bus");
assert_eq!(
with.correlation,
vec![crate::core::CorrelationKey::new("subject", "order-9")]
);
let without = CloudEvent::new("1", "urn:a", "t")
.expect("event")
.into_inbound("peer:bus");
assert!(without.correlation.is_empty());
}
use serde_json::json;
fn structured(body: &serde_json::Value) -> Result<CloudEvent, CloudEventError> {
CloudEvent::from_json(body.to_string().as_bytes())
}
#[test]
fn a_structured_event_parses_into_its_attributes() {
let event = structured(&json!({
"specversion": "1.0",
"type": "de.messwert.reading.direct.stored",
"source": "/edmd",
"id": "1",
"subject": "malo/42",
"time": "2026-08-19T09:00:00Z",
"datacontenttype": "application/json",
"tenantid": "acme",
"data": {"malo": "42"},
}))
.expect("a conformant event");
assert_eq!(event.event_type(), "de.messwert.reading.direct.stored");
assert_eq!(event.source(), "/edmd");
assert_eq!(event.id(), "1");
assert_eq!(event.subject(), Some("malo/42"));
assert_eq!(event.data(), Some(&json!({"malo": "42"})));
assert_eq!(
event.extension("tenantid"),
Some(&json!("acme")),
"an unknown attribute is an extension, not a discard: a deployment \
binds tenants on one"
);
assert!(
event.time().is_some(),
"the producer's clock was dropped rather than carried"
);
}
#[test]
fn two_producers_numbering_from_one_are_not_the_same_event() {
let first = CloudEvent::new("/edmd", "1", "reading.stored").unwrap();
let second = CloudEvent::new("/erp", "1", "reading.stored").unwrap();
assert_ne!(first.origin_id(), second.origin_id());
let sneaky = CloudEvent::new("/edmd\u{1f}1", "", "x");
assert!(
sneaky.is_err(),
"an empty id is not refused, so a source can be padded to spell \
another pair"
);
}
#[test]
fn a_binary_mode_event_parses_from_its_headers() {
let event = CloudEvent::from_http(
[
("Content-Type", "application/json"),
("ce-specversion", "1.0"),
("ce-id", "42"),
("ce-source", "/edmd"),
("ce-type", "reading.stored"),
("ce-tenantid", "acme"),
],
br#"{"malo":"7"}"#,
)
.expect("a conformant binary-mode event");
assert_eq!(event.id(), "42");
assert_eq!(event.event_type(), "reading.stored");
assert_eq!(event.data(), Some(&json!({"malo": "7"})));
assert_eq!(event.extension("tenantid"), Some(&json!("acme")));
}
#[test]
fn the_content_type_decides_which_mode_a_message_is_in() {
let event = CloudEvent::from_http(
[
(
"content-type",
"application/cloudevents+json; charset=UTF-8",
),
("ce-id", "from-the-headers"),
],
br#"{"specversion":"1.0","id":"from-the-body","source":"/x","type":"t"}"#,
)
.expect("a structured message");
assert_eq!(event.id(), "from-the-body");
}
#[test]
fn a_percent_encoded_header_value_is_decoded() {
let event = CloudEvent::from_http(
[
("ce-specversion", "1.0"),
("ce-id", "1"),
("ce-source", "/z%C3%A4hler"),
("ce-type", "t"),
],
b"",
)
.expect("a conformant event");
assert_eq!(event.source(), "/zähler");
assert_eq!(event.data(), None, "an empty body is no data, not null");
}
#[test]
fn what_is_refused_and_why() {
assert!(matches!(
structured(&json!({"specversion": "0.3", "id": "1", "source": "/x", "type": "t"})),
Err(CloudEventError::UnknownSpecVersion(_)),
));
assert!(matches!(
structured(&json!({"specversion": "1.0", "id": "", "source": "/x", "type": "t"})),
Err(CloudEventError::MissingAttribute("id"))
));
assert!(matches!(
structured(&json!({"specversion": "1.0", "id": "1", "source": "", "type": "t"})),
Err(CloudEventError::MissingAttribute("source"))
));
assert!(matches!(
structured(&json!({"specversion": "1.0", "id": "1", "source": "/x", "type": ""})),
Err(CloudEventError::MissingAttribute("type"))
));
assert!(
matches!(
structured(&json!({
"specversion": "1.0", "id": "1", "source": "/x", "type": "t",
"data_base64": "aGk=",
})),
Err(CloudEventError::BinaryData)
),
"bytes were decoded into a payload a run cannot address"
);
assert!(
matches!(
structured(&json!({
"specversion": "1.0", "id": "1", "source": "/x", "type": "t",
"tenantId": "acme",
})),
Err(CloudEventError::BadExtensionName(_))
),
"an extension name that is not lowercase alphanumeric arrives as a \
different name over a binding with case-insensitive keys"
);
assert!(matches!(
structured(&json!({
"specversion": "1.0", "id": "1", "source": "/x", "type": "t",
"time": "yesterday",
})),
Err(CloudEventError::BadTime(_))
));
assert!(
matches!(
CloudEvent::from_http(
[
("content-type", "application/octet-stream"),
("ce-specversion", "1.0"),
("ce-id", "1"),
("ce-source", "/x"),
("ce-type", "t"),
],
b"\x00\x01",
),
Err(CloudEventError::UnsupportedDataContentType(_))
),
"arbitrary bytes were retyped as a JSON payload"
);
}
#[test]
fn an_emitted_event_parses_back_to_itself() {
let event = CloudEvent::new("urn:mako:agentd", "run-1", "io.agentplane.run.completed")
.expect("the three required attributes")
.with_subject("run-1")
.with_extension("tenantid", json!("acme"))
.expect("a conformant extension name")
.with_data(json!({"outcome": "success"}));
let parsed = CloudEvent::from_json(&event.to_bytes()).expect("our own bytes");
assert_eq!(parsed, event);
assert_eq!(
parsed.to_bytes(),
event.to_bytes(),
"the encoding is not a function of the event, so a signature over \
it cannot be reproduced"
);
}
#[test]
fn an_inbound_event_keeps_the_producers_pair_without_trusting_it() {
let event = CloudEvent::new("/edmd", "7", "reading.stored")
.unwrap()
.with_data(json!({"malo": "42"}));
let inbound = event.clone().into_inbound("peer:gateway");
assert_eq!(
inbound.source, "peer:gateway",
"a self-asserted source lets a caller pick the namespace it \
deduplicates in"
);
assert_eq!(inbound.id, event.origin_id());
assert_eq!(inbound.kind, "reading.stored");
assert_eq!(inbound.payload, json!({"malo": "42"}));
}
}