use liminal::protocol::SchemaId as ProtocolSchemaId;
use serde_json::Value;
use crate::config::types::ChannelDef;
pub(super) struct ChannelSchema {
pub document: Value,
pub protocol_id: ProtocolSchemaId,
}
const EMPTY_SCHEMA_BYTES: &[u8] = b"{}";
pub(super) fn resolve_channel_schema(channel: &ChannelDef) -> ChannelSchema {
channel.loaded_schema.as_ref().map_or_else(
|| ChannelSchema {
document: Value::Object(serde_json::Map::new()),
protocol_id: schema_id_from_bytes(EMPTY_SCHEMA_BYTES),
},
|loaded| ChannelSchema {
document: loaded.document.clone(),
protocol_id: schema_id_from_bytes(&loaded.bytes),
},
)
}
pub(super) fn resolve_schema_bytes(
schema_bytes: Option<&[u8]>,
) -> Result<ChannelSchema, serde_json::Error> {
let Some(bytes) = schema_bytes else {
return Ok(ChannelSchema {
document: Value::Object(serde_json::Map::new()),
protocol_id: schema_id_from_bytes(EMPTY_SCHEMA_BYTES),
});
};
Ok(ChannelSchema {
document: serde_json::from_slice(bytes)?,
protocol_id: schema_id_from_bytes(bytes),
})
}
fn schema_id_from_bytes(schema_bytes: &[u8]) -> ProtocolSchemaId {
let mut id = [0_u8; ProtocolSchemaId::WIRE_LEN];
let mut hash = fnv1a(schema_bytes).to_be_bytes();
for (index, slot) in id.iter_mut().enumerate() {
*slot = hash[index % hash.len()];
if index % hash.len() == hash.len() - 1 {
hash = fnv1a(&hash).to_be_bytes();
}
}
ProtocolSchemaId::new(id)
}
fn fnv1a(bytes: &[u8]) -> u64 {
const OFFSET_BASIS: u64 = 0xcbf2_9ce4_8422_2325;
const PRIME: u64 = 0x0000_0100_0000_01b3;
let mut hash = OFFSET_BASIS;
for byte in bytes {
hash ^= u64::from(*byte);
hash = hash.wrapping_mul(PRIME);
}
hash
}
#[cfg(test)]
mod tests {
use super::{resolve_channel_schema, resolve_schema_bytes, schema_id_from_bytes};
use crate::config::types::{ChannelDef, LoadedSchema};
fn channel(schema_ref: Option<&str>, loaded: Option<LoadedSchema>) -> ChannelDef {
ChannelDef {
name: "orders".to_owned(),
schema_ref: schema_ref.map(Into::into),
durable: false,
loaded_schema: loaded,
}
}
#[test]
fn schema_ids_are_deterministic_and_content_addressed() {
assert_eq!(schema_id_from_bytes(b"a"), schema_id_from_bytes(b"a"));
assert_ne!(schema_id_from_bytes(b"a"), schema_id_from_bytes(b"b"));
}
#[test]
fn loaded_schema_id_derives_from_loaded_bytes() {
let bytes = br#"{"type":"object"}"#.to_vec();
let document = serde_json::json!({"type": "object"});
let resolved = resolve_channel_schema(&channel(
Some("orders.json"),
Some(LoadedSchema {
bytes: bytes.clone(),
document: document.clone(),
}),
));
assert_eq!(resolved.document, document);
assert_eq!(resolved.protocol_id, schema_id_from_bytes(&bytes));
}
#[test]
fn schema_less_channel_is_permissive_and_uses_empty_schema_id() {
let resolved = resolve_channel_schema(&channel(None, None));
assert_eq!(resolved.document, serde_json::json!({}));
assert_eq!(resolved.protocol_id, schema_id_from_bytes(b"{}"));
}
#[test]
fn raw_bytes_resolution_matches_the_channel_def_form() -> Result<(), serde_json::Error> {
let empty = resolve_schema_bytes(None)?;
let from_def = resolve_channel_schema(&channel(None, None));
assert_eq!(empty.document, from_def.document);
assert_eq!(empty.protocol_id, from_def.protocol_id);
let bytes = br#"{"type":"object"}"#.to_vec();
let loaded = resolve_schema_bytes(Some(&bytes))?;
let loaded_from_def = resolve_channel_schema(&channel(
Some("orders.json"),
Some(LoadedSchema {
bytes,
document: serde_json::json!({"type": "object"}),
}),
));
assert_eq!(loaded.document, loaded_from_def.document);
assert_eq!(loaded.protocol_id, loaded_from_def.protocol_id);
Ok(())
}
#[test]
fn raw_bytes_resolution_refuses_bytes_that_are_not_json() {
assert!(resolve_schema_bytes(Some(b"not json")).is_err());
}
}