use serde_json::Value;
use sha2::{Digest, Sha256};
use crate::errors::OrionError;
use crate::storage::models::{Channel, Connector, Workflow};
use crate::storage::repositories::channels::CreateChannelRequest;
use crate::storage::repositories::connectors::CreateConnectorRequest;
use crate::storage::repositories::workflows::CreateWorkflowRequest;
pub fn workflow_content(w: &Workflow) -> Result<Value, OrionError> {
Ok(orion_api::content::workflow_content(&serde_json::json!({
"name": w.name,
"description": w.description,
"priority": w.priority,
"condition": serde_json::from_str::<Value>(&w.condition_json)?,
"tasks": serde_json::from_str::<Value>(&w.tasks_json)?,
"tags": serde_json::from_str::<Value>(&w.tags_json)?,
"loop": w.loop_json.as_deref().map(serde_json::from_str::<Value>).transpose()?,
"continue_on_error": w.continue_on_error,
})))
}
pub fn workflow_request_content(r: &CreateWorkflowRequest) -> Value {
orion_api::content::workflow_content(&serde_json::json!({
"name": r.name,
"description": r.description,
"priority": r.priority,
"condition": r.condition,
"tasks": r.tasks,
"tags": r.tags,
"loop": r.loop_config,
"continue_on_error": r.continue_on_error,
}))
}
pub fn plugin_content(p: &crate::storage::models::Plugin) -> Result<Value, OrionError> {
Ok(orion_api::content::plugin_content(&serde_json::json!({
"manifest": serde_json::from_str::<Value>(&p.manifest_json)?,
"digest": p.digest,
"tags": serde_json::from_str::<Value>(&p.tags_json)?,
})))
}
pub fn plugin_request_content(manifest: &Value, digest: &str, tags: &[String]) -> Value {
orion_api::content::plugin_content(&serde_json::json!({
"manifest": manifest,
"digest": digest,
"tags": tags,
}))
}
pub fn model_content(m: &crate::storage::models::Model) -> Result<Value, OrionError> {
Ok(orion_api::content::model_content(&serde_json::json!({
"manifest": serde_json::from_str::<Value>(&m.manifest_json)?,
"artifact": serde_json::from_str::<Value>(&m.artifact_json)?,
"tags": serde_json::from_str::<Value>(&m.tags_json)?,
})))
}
pub fn model_request_content(manifest: &Value, artifact: &Value, tags: &[String]) -> Value {
orion_api::content::model_content(&serde_json::json!({
"manifest": manifest,
"artifact": artifact,
"tags": tags,
}))
}
pub fn channel_content(c: &Channel) -> Result<Value, OrionError> {
let methods = c
.methods_json
.as_deref()
.map(serde_json::from_str::<Value>)
.transpose()?;
Ok(orion_api::content::channel_content(&serde_json::json!({
"name": c.name,
"description": c.description,
"channel_type": c.channel_type,
"protocol": c.protocol,
"methods": methods,
"route_pattern": c.route_pattern,
"topic": c.topic,
"consumer_group": c.consumer_group,
"transport_config": serde_json::from_str::<Value>(&c.transport_config_json)?,
"workflow_id": c.workflow_id,
"config": serde_json::from_str::<Value>(&c.config_json)?,
"priority": c.priority,
"tags": serde_json::from_str::<Value>(&c.tags_json)?,
})))
}
pub fn channel_request_content(r: &CreateChannelRequest) -> Value {
orion_api::content::channel_content(&serde_json::json!({
"name": r.name,
"description": r.description,
"channel_type": r.channel_type.as_str(),
"protocol": r.protocol.as_str(),
"methods": r.methods,
"route_pattern": r.route_pattern,
"topic": r.topic,
"consumer_group": r.consumer_group,
"transport_config": r.transport_config,
"workflow_id": r.workflow_id,
"config": r.config,
"priority": r.priority,
"tags": r.tags,
}))
}
pub fn connector_content(c: &Connector) -> Result<Value, OrionError> {
Ok(orion_api::content::connector_content(&serde_json::json!({
"name": c.name,
"connector_type": c.connector_type,
"config": serde_json::from_str::<Value>(&c.config_json)?,
"enabled": c.enabled,
"tags": serde_json::from_str::<Value>(&c.tags_json)?,
})))
}
pub fn connector_request_content(r: &CreateConnectorRequest) -> Value {
orion_api::content::connector_content(&serde_json::json!({
"name": r.name,
"connector_type": r.connector_type.as_str(),
"config": r.config,
"enabled": r.enabled,
"tags": r.tags,
}))
}
pub fn canonical_json(value: &Value) -> String {
fn write(value: &Value, out: &mut String) {
match value {
Value::Object(map) => {
let mut keys: Vec<&String> = map.keys().collect();
keys.sort_unstable();
out.push('{');
for (i, key) in keys.iter().enumerate() {
if i > 0 {
out.push(',');
}
out.push_str(&serde_json::to_string(key).expect("string serializes"));
out.push(':');
write(&map[key.as_str()], out);
}
out.push('}');
}
Value::Array(items) => {
out.push('[');
for (i, item) in items.iter().enumerate() {
if i > 0 {
out.push(',');
}
write(item, out);
}
out.push(']');
}
leaf => out.push_str(&serde_json::to_string(leaf).expect("scalar serializes")),
}
}
let mut out = String::new();
write(value, &mut out);
out
}
pub fn content_hash(value: &Value) -> String {
let mut hasher = Sha256::new();
hasher.update(canonical_json(value).as_bytes());
format!("sha256:{}", hex::encode(hasher.finalize()))
}
pub const CONTENT_VERSION_KEYWORD: &str = "content";
pub const CONTENT_VERSION_HEX_LEN: usize = 12;
pub fn content_version(content_hash: &str, prefix: Option<&str>) -> Result<String, String> {
if !crate::crypto::is_sha256_digest(content_hash) {
return Err(format!(
"'{content_hash}' is not a content hash (sha256:<64 lowercase hex>)"
));
}
let prefix = prefix.unwrap_or(CONTENT_VERSION_KEYWORD);
check_version_prefix(prefix)?;
let hex = &content_hash["sha256:".len()..][..CONTENT_VERSION_HEX_LEN];
Ok(format!("{prefix}-{hex}"))
}
pub fn check_version_prefix(prefix: &str) -> Result<(), String> {
crate::validation::package_key(
"--version-prefix",
prefix,
crate::validation::MAX_PACKAGE_VERSION_LEN - 1 - CONTENT_VERSION_HEX_LEN,
)
}
pub fn content_version_hex(version: &str) -> Option<&str> {
version
.strip_prefix(CONTENT_VERSION_KEYWORD)
.and_then(|rest| rest.strip_prefix('-'))
.filter(|hex| {
hex.len() == CONTENT_VERSION_HEX_LEN
&& hex
.bytes()
.all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
})
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
const HASH: &str = "sha256:3f9c2a1b7d04e5f60718293a4b5c6d7e8f90a1b2c3d4e5f60718293a4b5c6d7e";
#[test]
fn a_content_version_is_the_hash_prefix() {
assert_eq!(
content_version(HASH, None).expect("keyword form"),
"content-3f9c2a1b7d04"
);
assert_eq!(
content_version(HASH, Some("1.4.0")).expect("prefixed"),
"1.4.0-3f9c2a1b7d04"
);
assert_eq!(
content_version_hex("content-3f9c2a1b7d04"),
Some("3f9c2a1b7d04")
);
assert_eq!(content_version_hex("1.4.0-3f9c2a1b7d04"), None);
assert_eq!(content_version_hex("content-3F9C2A1B7D04"), None);
assert_eq!(content_version_hex("content-3f9c"), None);
}
#[test]
fn a_content_version_refuses_what_the_receipt_would() {
assert!(content_version("sha256:abc", None).is_err());
assert!(content_version(&HASH.to_uppercase(), None).is_err());
for bad in ["", "1.4/0", "1.4.0+x"] {
assert!(content_version(HASH, Some(bad)).is_err(), "{bad:?}");
}
let longest = "p".repeat(51);
let version = content_version(HASH, Some(&longest)).expect("51 fits");
assert_eq!(version.len(), 64);
crate::validation::package_key("version", &version, 64).expect("a receipt key");
assert!(content_version(HASH, Some(&"p".repeat(52))).is_err());
}
#[test]
fn content_hash_spelling_is_pinned_to_plain_sha256_lowercase_hex() {
let v: Value =
serde_json::from_str(r#"{"b": [1, {"y": 2, "x": 3}], "a": null}"#).expect("test");
assert_eq!(canonical_json(&v), r#"{"a":null,"b":[1,{"x":3,"y":2}]}"#);
assert_eq!(
content_hash(&v),
"sha256:fd5905a59ba4aec9fd37e5214d395b1f9ac0db9d3a7addf85ee5c31e89e8a5bc"
);
}
#[test]
fn canonical_hash_is_order_insensitive_and_value_sensitive() {
let a: Value =
serde_json::from_str(r#"{"b": [1, {"y": 2, "x": 3}], "a": null}"#).expect("test");
let b: Value =
serde_json::from_str(r#"{"a": null, "b": [1, {"x": 3, "y": 2}]}"#).expect("test");
assert_eq!(canonical_json(&a), r#"{"a":null,"b":[1,{"x":3,"y":2}]}"#);
assert_eq!(content_hash(&a), content_hash(&b));
let c = json!({"a": null, "b": [1, {"x": 3, "y": 999}]});
assert_ne!(content_hash(&a), content_hash(&c));
}
#[test]
fn channel_and_connector_projections_agree() {
let now = chrono::NaiveDateTime::default();
let req: crate::storage::repositories::channels::CreateChannelRequest =
serde_json::from_value(json!({
"channel_id": "ch-1",
"name": "Hash Ch",
"channel_type": "sync",
"protocol": "rest",
"methods": ["POST"],
"route_pattern": "/hash",
"workflow_id": "wf-1",
"config": {"a": 1},
"priority": 2,
"tags": ["t"],
}))
.expect("channel request");
let row = Channel {
channel_id: "ch-1".to_string(),
version: 4, name: "Hash Ch".to_string(),
description: None,
channel_type: "sync".to_string(),
protocol: "rest".to_string(),
methods_json: Some(r#"["POST"]"#.to_string()),
route_pattern: Some("/hash".to_string()),
topic: None,
consumer_group: None,
transport_config_json: "{}".to_string(),
workflow_id: Some("wf-1".to_string()),
config_json: r#"{"a":1}"#.to_string(),
status: "active".to_string(),
priority: 2,
tags_json: r#"["t"]"#.to_string(),
created_at: now,
updated_at: now,
};
assert_eq!(
channel_content(&row).expect("channel content"),
channel_request_content(&req)
);
assert_eq!(
canonical_json(&channel_request_content(&req)),
r#"{"channel_type":"sync","config":{"a":1},"consumer_group":null,"description":null,"methods":["POST"],"name":"Hash Ch","priority":2,"protocol":"rest","route_pattern":"/hash","tags":["t"],"topic":null,"transport_config":{},"workflow_id":"wf-1"}"#
);
let req: crate::storage::repositories::connectors::CreateConnectorRequest =
serde_json::from_value(json!({
"name": "hash-conn",
"connector_type": "http",
"config": {"url": "https://example.com"},
"enabled": false,
"tags": ["t"],
}))
.expect("connector request");
let row = Connector {
id: "generated".to_string(), name: "hash-conn".to_string(),
connector_type: "http".to_string(),
config_json: r#"{"url":"https://example.com"}"#.to_string(),
enabled: false,
tags_json: r#"["t"]"#.to_string(),
created_at: now,
updated_at: now,
};
assert_eq!(
connector_content(&row).expect("connector content"),
connector_request_content(&req)
);
assert_eq!(
canonical_json(&connector_request_content(&req)),
r#"{"config":{"url":"https://example.com"},"connector_type":"http","enabled":false,"name":"hash-conn","tags":["t"]}"#
);
}
#[test]
fn model_row_and_request_projections_agree() {
let now = chrono::NaiveDateTime::default();
let manifest = json!({"abi": "1", "name": "m", "version": "1.0.0", "format": "onnx",
"inputs": [{"name": "x"}], "outputs": [{"name": "y"}]});
let artifact = json!({"connector": "models", "key": "m/1.onnx",
"digest": "sha256:d", "size": 10});
let row = crate::storage::models::Model {
model_id: "m".to_string(),
version: 3, status: "active".to_string(),
digest: "sha256:d".to_string(),
manifest_json: serde_json::to_string(&manifest).expect("test"),
artifact_json: r#"{"connector":"models","key":"m/1.onnx","digest":"sha256:d"}"#
.to_string(),
admission_json: r#"{"state":"passed","node":"n1"}"#.to_string(),
stats_json: Some(r#"{"parameters":5}"#.to_string()),
tags_json: r#"["t"]"#.to_string(),
signature: Some("sig".to_string()),
created_at: now,
updated_at: now,
};
let row_content = model_content(&row).expect("model content");
assert_eq!(
row_content,
model_request_content(&manifest, &artifact, &["t".to_string()])
);
assert_eq!(
canonical_json(&row_content),
r#"{"artifact":{"connector":"models","digest":"sha256:d","key":"m/1.onnx"},"manifest":{"abi":"1","format":"onnx","inputs":[{"name":"x"}],"name":"m","outputs":[{"name":"y"}],"version":"1.0.0"},"tags":["t"]}"#
);
}
#[test]
fn row_and_request_projections_agree() {
let req: CreateWorkflowRequest = serde_json::from_value(json!({
"workflow_id": "wf-1",
"name": "Hash Me",
"priority": 5,
"tags": ["a"],
"tasks": [{"id": "t", "name": "T",
"function": {"name": "log", "input": {"message": "x"}}}],
}))
.expect("request");
let now = chrono::NaiveDateTime::default();
let row = Workflow {
workflow_id: "wf-1".to_string(),
version: 3, name: "Hash Me".to_string(),
description: None,
priority: 5,
status: "active".to_string(),
rollout_percentage: 50,
condition_json: "true".to_string(),
tasks_json: serde_json::to_string(&req.tasks).expect("test"),
tags_json: r#"["a"]"#.to_string(),
loop_json: None,
continue_on_error: false,
created_at: now,
updated_at: now,
};
let row_content = workflow_content(&row).expect("content");
assert_eq!(
canonical_json(&row_content),
r#"{"condition":true,"continue_on_error":false,"description":null,"name":"Hash Me","priority":5,"tags":["a"],"tasks":[{"function":{"input":{"message":"x"},"name":"log"},"id":"t","name":"T"}]}"#
);
assert_eq!(row_content, workflow_request_content(&req));
assert_eq!(
content_hash(&row_content),
content_hash(&workflow_request_content(&req))
);
}
}