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()))
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
#[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))
);
}
}