use serde_json::Value;
use super::enums::parse_json_field;
use super::rows::{
AuditLogEntry, Channel, Connector, CronOccurrence, Model, PackageReceipt, Plugin,
TraceDlqEntry, TraceDlqSummary, TraceListRow, Workflow,
};
use crate::errors::OrionError;
pub use orion_api::dto::{
AuditLogEntryResponse, ChannelResponse, ConnectorResponse, CronOccurrenceResponse,
CronOccurrenceSummaryResponse, CronScheduleStatusResponse, ModelAdmission, ModelArtifactRef,
ModelHealth, ModelResponse, ModelStats, PackageReceiptResponse, PluginHealth, PluginResponse,
TraceDlqEntryResponse, TraceDlqSummaryResponse, TraceListItemResponse, WorkflowResponse,
};
impl From<&CronOccurrence> for CronOccurrenceSummaryResponse {
fn from(row: &CronOccurrence) -> Self {
Self {
id: row.id.clone(),
channel_id: row.channel_id.clone(),
channel_name: row.channel_name.clone(),
trigger: row.trigger.clone(),
scheduled_for: row.scheduled_for,
status: row.status.clone(),
attempt: row.attempt,
started_at: row.started_at,
completed_at: row.completed_at,
created_at: row.created_at,
}
}
}
impl From<&CronOccurrence> for CronOccurrenceResponse {
fn from(row: &CronOccurrence) -> Self {
Self {
id: row.id.clone(),
channel_id: row.channel_id.clone(),
channel_name: row.channel_name.clone(),
channel_version: row.channel_version,
executing_version: row.executing_version,
workflow_id: row.workflow_id.clone(),
trigger: row.trigger.clone(),
scheduled_for: row.scheduled_for,
status: row.status.clone(),
attempt: row.attempt,
claimed_by: row.claimed_by.clone(),
claimed_until: row.claimed_until,
singleton_key: row.singleton_key.clone(),
fencing_token: row.fencing_token,
trace_id: row.trace_id.clone(),
error_message: row.error_message.clone(),
started_at: row.started_at,
completed_at: row.completed_at,
created_at: row.created_at,
updated_at: row.updated_at,
}
}
}
impl TryFrom<&Plugin> for PluginResponse {
type Error = OrionError;
fn try_from(plugin: &Plugin) -> Result<Self, Self::Error> {
let id = &plugin.plugin_id;
let manifest: Value =
parse_json_field(&plugin.manifest_json, "plugin", id, "manifest_json")?;
let functions = manifest
.get("functions")
.and_then(Value::as_array)
.map(|fs| {
fs.iter()
.filter_map(|f| f.get("name").and_then(Value::as_str))
.map(str::to_string)
.collect()
})
.unwrap_or_default();
let field = |key: &str| {
manifest
.get(key)
.and_then(Value::as_str)
.unwrap_or_default()
.to_string()
};
Ok(Self {
plugin_id: plugin.plugin_id.clone(),
version: plugin.version,
status: plugin.status.clone(),
digest: plugin.digest.clone(),
abi: field("abi"),
plugin_version: field("version"),
manifest,
functions,
tags: parse_json_field(&plugin.tags_json, "plugin", id, "tags_json")?,
content_hash: crate::storage::content::content_hash(
&crate::storage::content::plugin_content(plugin)?,
),
signature: plugin.signature.clone(),
health: None,
created_at: plugin.created_at,
updated_at: plugin.updated_at,
})
}
}
impl TryFrom<&Model> for ModelResponse {
type Error = OrionError;
fn try_from(model: &Model) -> Result<Self, Self::Error> {
let id = &model.model_id;
let manifest: Value = parse_json_field(&model.manifest_json, "model", id, "manifest_json")?;
let names = |key: &str| -> Vec<String> {
manifest
.get(key)
.and_then(Value::as_array)
.map(|entries| {
entries
.iter()
.filter_map(|entry| entry.get("name").and_then(Value::as_str))
.map(str::to_string)
.collect()
})
.unwrap_or_default()
};
let field = |key: &str| {
manifest
.get(key)
.and_then(Value::as_str)
.unwrap_or_default()
.to_string()
};
let abi = field("abi");
let model_version = field("version");
let format = field("format");
let inputs = names("inputs");
let outputs = names("outputs");
let artifact: ModelArtifactRef =
parse_json_field(&model.artifact_json, "model", id, "artifact_json")?;
let admission: ModelAdmission =
parse_json_field(&model.admission_json, "model", id, "admission_json")?;
let stats: Option<ModelStats> = model
.stats_json
.as_deref()
.map(|json| parse_json_field(json, "model", id, "stats_json"))
.transpose()?;
Ok(Self {
model_id: model.model_id.clone(),
version: model.version,
status: model.status.clone(),
digest: model.digest.clone(),
abi,
model_version,
format,
manifest,
inputs,
outputs,
artifact,
admission,
stats,
tags: parse_json_field(&model.tags_json, "model", id, "tags_json")?,
content_hash: crate::storage::content::content_hash(
&crate::storage::content::model_content(model)?,
),
signature: model.signature.clone(),
health: None,
created_at: model.created_at,
updated_at: model.updated_at,
})
}
}
impl TryFrom<&Workflow> for WorkflowResponse {
type Error = OrionError;
fn try_from(workflow: &Workflow) -> Result<Self, Self::Error> {
let id = &workflow.workflow_id;
Ok(Self {
workflow_id: workflow.workflow_id.clone(),
version: workflow.version,
name: workflow.name.clone(),
description: workflow.description.clone(),
priority: workflow.priority,
status: workflow.status.clone(),
rollout_percentage: workflow.rollout_percentage,
condition: parse_json_field(
&workflow.condition_json,
"workflow",
id,
"condition_json",
)?,
tasks: parse_json_field(&workflow.tasks_json, "workflow", id, "tasks_json")?,
tags: parse_json_field(&workflow.tags_json, "workflow", id, "tags_json")?,
loop_config: workflow
.loop_json
.as_deref()
.map(|json| parse_json_field(json, "workflow", id, "loop_json"))
.transpose()?,
continue_on_error: workflow.continue_on_error,
content_hash: crate::storage::content::content_hash(
&crate::storage::content::workflow_content(workflow)?,
),
created_at: workflow.created_at,
updated_at: workflow.updated_at,
})
}
}
impl TryFrom<&Channel> for ChannelResponse {
type Error = OrionError;
fn try_from(channel: &Channel) -> Result<Self, Self::Error> {
let id = &channel.channel_id;
let methods = channel
.methods_json
.as_ref()
.map(|m| parse_json_field(m, "channel", id, "methods_json"))
.transpose()?;
Ok(Self {
channel_id: channel.channel_id.clone(),
version: channel.version,
name: channel.name.clone(),
description: channel.description.clone(),
channel_type: channel.channel_type.clone(),
protocol: channel.protocol.clone(),
methods,
route_pattern: channel.route_pattern.clone(),
topic: channel.topic.clone(),
consumer_group: channel.consumer_group.clone(),
transport_config: parse_json_field(
&channel.transport_config_json,
"channel",
id,
"transport_config_json",
)?,
workflow_id: channel.workflow_id.clone(),
config: {
let mut config =
parse_json_field(&channel.config_json, "channel", id, "config_json")?;
crate::connector::mask_channel_config(&mut config);
config
},
status: channel.status.clone(),
priority: channel.priority,
tags: parse_json_field(&channel.tags_json, "channel", id, "tags_json")?,
content_hash: crate::storage::content::content_hash(
&crate::storage::content::channel_content(channel)?,
),
created_at: channel.created_at,
updated_at: channel.updated_at,
})
}
}
impl From<&Connector> for ConnectorResponse {
fn from(connector: &Connector) -> Self {
Self {
id: connector.id.clone(),
name: connector.name.clone(),
connector_type: connector.connector_type.clone(),
config_json: connector.config_json.clone(),
config: Value::Null,
enabled: connector.enabled,
tags: serde_json::from_str(&connector.tags_json)
.unwrap_or_else(|_| Value::Array(Vec::new())),
content_hash: crate::storage::content::connector_content(connector)
.map(|v| crate::storage::content::content_hash(&v))
.unwrap_or_default(),
created_at: connector.created_at,
updated_at: connector.updated_at,
}
}
}
impl From<&PackageReceipt> for PackageReceiptResponse {
fn from(receipt: &PackageReceipt) -> Self {
Self {
name: receipt.name.clone(),
version: receipt.version.clone(),
content_hash: receipt.content_hash.clone(),
state: receipt.state.clone(),
principal: receipt.principal.clone(),
created_at: receipt.created_at,
updated_at: receipt.updated_at,
}
}
}
impl From<&TraceDlqEntry> for TraceDlqEntryResponse {
fn from(entry: &TraceDlqEntry) -> Self {
Self {
id: entry.id.clone(),
trace_id: entry.trace_id.clone(),
channel: entry.channel.clone(),
payload_json: entry.payload_json.clone(),
metadata_json: entry.metadata_json.clone(),
error_message: entry.error_message.clone(),
retry_count: entry.retry_count,
max_retries: entry.max_retries,
next_retry_at: entry.next_retry_at,
created_at: entry.created_at,
updated_at: entry.updated_at,
}
}
}
impl From<&TraceDlqSummary> for TraceDlqSummaryResponse {
fn from(entry: &TraceDlqSummary) -> Self {
Self {
id: entry.id.clone(),
trace_id: entry.trace_id.clone(),
channel: entry.channel.clone(),
error_message: entry.error_message.clone(),
retry_count: entry.retry_count,
max_retries: entry.max_retries,
next_retry_at: entry.next_retry_at,
created_at: entry.created_at,
updated_at: entry.updated_at,
}
}
}
impl From<&TraceListRow> for TraceListItemResponse {
fn from(trace: &TraceListRow) -> Self {
Self {
id: trace.id.clone(),
channel: trace.channel.clone(),
channel_id: trace.channel_id.clone(),
mode: trace.mode.clone(),
status: trace.status.clone(),
error_message: trace.error_message.clone(),
duration_ms: trace.duration_ms,
started_at: trace.started_at,
completed_at: trace.completed_at,
created_at: trace.created_at,
updated_at: trace.updated_at,
}
}
}
impl From<&AuditLogEntry> for AuditLogEntryResponse {
fn from(entry: &AuditLogEntry) -> Self {
Self {
id: entry.id.clone(),
principal: entry.principal.clone(),
action: entry.action.clone(),
resource_type: entry.resource_type.clone(),
resource_id: entry.resource_id.clone(),
details: entry.details.clone(),
created_at: entry.created_at,
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::storage::models::enums::{CHANNEL_TYPE_ASYNC, CHANNEL_TYPE_SYNC};
use crate::storage::models::{ChannelProtocol, EntityStatus};
use chrono::{NaiveDate, NaiveDateTime};
#[test]
fn json_columns_carry_the_json_suffix() {
const SOURCE: &str = include_str!("dto.rs");
const CALL: &str = "parse_json_field(";
let mut labels = Vec::new();
let code = SOURCE
.split_once("\n#[cfg(test)]\n")
.map(|(before, _)| before)
.unwrap_or(SOURCE);
let mut rest = code;
while let Some(start) = rest.find(CALL) {
rest = &rest[start + CALL.len()..];
let mut depth = 0usize;
let mut args = vec![String::new()];
let mut end = rest.len();
for (i, ch) in rest.char_indices() {
match ch {
'(' | '[' => depth += 1,
')' | ']' if depth == 0 => {
end = i;
break;
}
')' | ']' => depth -= 1,
',' if depth == 0 => {
args.push(String::new());
continue;
}
_ => {}
}
args.last_mut().expect("at least one arg").push(ch);
}
args.retain(|a| !a.trim().is_empty());
assert_eq!(
args.len(),
4,
"parse_json_field takes (value, entity, id, column); found {} args in \
`{}`",
args.len(),
&rest[..end.min(120)]
);
labels.push(args[3].trim().trim_matches('"').to_string());
rest = &rest[end..];
}
assert!(
labels.len() >= 6,
"the scan found only {} parse_json_field call(s) — it has stopped \
matching the source and is no longer checking anything: {labels:?}",
labels.len()
);
for label in &labels {
assert!(
label.ends_with("_json"),
"`{label}` is decoded as a JSON document but its column is not named \
`*_json` (D26). The suffix is the only signal a reader gets that the \
value has to go through serde_json before it means anything — every \
other column in the schema is used as-is."
);
}
for expected in ["tags_json", "methods_json"] {
assert!(
labels.iter().any(|l| l == expected),
"expected `{expected}` among the decoded columns: {labels:?}"
);
}
}
fn sample_datetime() -> NaiveDateTime {
NaiveDate::from_ymd_opt(2025, 1, 1)
.expect("test")
.and_hms_opt(0, 0, 0)
.expect("test")
}
fn sample_workflow() -> Workflow {
Workflow {
workflow_id: "wf-1".to_string(),
name: "Test Workflow".to_string(),
description: Some("A test workflow".to_string()),
priority: 10,
version: 1,
status: EntityStatus::Active.as_str().to_string(),
rollout_percentage: 100,
condition_json: r#"{"==": [1, 1]}"#.to_string(),
tasks_json: r#"[{"id": "t1", "function": "http_call"}]"#.to_string(),
tags_json: r#"["test"]"#.to_string(),
loop_json: None,
continue_on_error: false,
created_at: sample_datetime(),
updated_at: sample_datetime(),
}
}
fn sample_channel() -> Channel {
Channel {
tags_json: "[]".to_string(),
channel_id: "ch-1".to_string(),
version: 1,
name: "orders".to_string(),
description: Some("Order processing channel".to_string()),
channel_type: CHANNEL_TYPE_SYNC.to_string(),
protocol: ChannelProtocol::Rest.as_str().to_string(),
methods_json: Some(r#"["POST"]"#.to_string()),
route_pattern: Some("/orders".to_string()),
topic: None,
consumer_group: None,
transport_config_json: "{}".to_string(),
workflow_id: Some("wf-1".to_string()),
config_json: r#"{"timeout_ms": 5000}"#.to_string(),
status: EntityStatus::Active.as_str().to_string(),
priority: 0,
created_at: sample_datetime(),
updated_at: sample_datetime(),
}
}
#[test]
fn test_workflow_response_try_from_valid() {
let workflow = sample_workflow();
let response = WorkflowResponse::try_from(&workflow).expect("test");
assert_eq!(response.workflow_id, "wf-1");
assert_eq!(response.name, "Test Workflow");
assert_eq!(response.priority, 10);
assert_eq!(response.version, 1);
assert_eq!(response.status, EntityStatus::Active.as_str());
assert_eq!(response.rollout_percentage, 100);
assert_eq!(response.condition, serde_json::json!({"==": [1, 1]}));
assert_eq!(
response.tasks,
serde_json::json!([{"id": "t1", "function": "http_call"}])
);
assert_eq!(response.tags, serde_json::json!(["test"]));
assert!(!response.continue_on_error);
}
#[test]
fn test_workflow_response_try_from_invalid_condition_json() {
let mut workflow = sample_workflow();
workflow.condition_json = "not valid json {{{".to_string();
let result = WorkflowResponse::try_from(&workflow);
assert!(result.is_err());
}
#[test]
fn test_workflow_response_try_from_invalid_tasks_json() {
let mut workflow = sample_workflow();
workflow.tasks_json = "invalid".to_string();
let result = WorkflowResponse::try_from(&workflow);
assert!(result.is_err());
}
#[test]
fn test_workflow_response_try_from_invalid_tags_json() {
let mut workflow = sample_workflow();
workflow.tags_json = "not json".to_string();
let result = WorkflowResponse::try_from(&workflow);
assert!(result.is_err());
}
#[test]
fn test_workflow_response_try_from_no_description() {
let mut workflow = sample_workflow();
workflow.description = None;
let response = WorkflowResponse::try_from(&workflow).expect("test");
assert!(response.description.is_none());
}
#[test]
fn test_channel_response_try_from_valid() {
let channel = sample_channel();
let response = ChannelResponse::try_from(&channel).expect("test");
assert_eq!(response.channel_id, "ch-1");
assert_eq!(response.name, "orders");
assert_eq!(response.channel_type, CHANNEL_TYPE_SYNC);
assert_eq!(response.protocol, ChannelProtocol::Rest.as_str());
assert_eq!(response.methods, Some(serde_json::json!(["POST"])));
assert_eq!(response.route_pattern, Some("/orders".to_string()));
assert!(response.topic.is_none());
assert_eq!(response.workflow_id, Some("wf-1".to_string()));
assert_eq!(response.config, serde_json::json!({"timeout_ms": 5000}));
}
#[test]
fn test_channel_response_try_from_async() {
let mut channel = sample_channel();
channel.channel_type = CHANNEL_TYPE_ASYNC.to_string();
channel.protocol = ChannelProtocol::Kafka.as_str().to_string();
channel.methods_json = None;
channel.route_pattern = None;
channel.topic = Some("order.placed".to_string());
channel.consumer_group = Some("orion".to_string());
let response = ChannelResponse::try_from(&channel).expect("test");
assert_eq!(response.channel_type, CHANNEL_TYPE_ASYNC);
assert_eq!(response.protocol, ChannelProtocol::Kafka.as_str());
assert!(response.methods.is_none());
assert_eq!(response.topic, Some("order.placed".to_string()));
}
#[test]
fn test_channel_response_try_from_invalid_config_json() {
let mut channel = sample_channel();
channel.config_json = "bad json".to_string();
let result = ChannelResponse::try_from(&channel);
assert!(result.is_err());
}
fn field_names(value: &serde_json::Value) -> Vec<String> {
value
.as_object()
.expect("DTO serializes to an object")
.keys()
.cloned()
.collect()
}
#[test]
fn connector_response_keeps_the_row_wire_shape() {
let row = Connector {
tags_json: "[]".to_string(),
id: "c-1".to_string(),
name: "pg".to_string(),
connector_type: "db".to_string(),
config_json: r#"{"password":"******"}"#.to_string(),
enabled: true,
created_at: sample_datetime(),
updated_at: sample_datetime(),
};
let value = serde_json::to_value(ConnectorResponse::from(&row)).expect("test");
assert_eq!(
field_names(&value),
[
"id",
"name",
"connector_type",
"config_json",
"config",
"enabled",
"tags",
"content_hash",
"created_at",
"updated_at"
]
);
assert_eq!(value["config_json"], r#"{"password":"******"}"#);
assert_eq!(value["config"], serde_json::Value::Null);
assert_eq!(value["tags"], serde_json::json!([]));
assert!(
value["content_hash"]
.as_str()
.is_some_and(|h| h.starts_with("sha256:")),
"{value}"
);
assert_eq!(value["created_at"], "2025-01-01T00:00:00");
}
#[test]
fn audit_log_response_keeps_the_row_wire_shape() {
let row = AuditLogEntry {
id: "a-1".to_string(),
principal: "admin...".to_string(),
action: "activate".to_string(),
resource_type: "workflow".to_string(),
resource_id: "wf-1".to_string(),
details: None,
created_at: sample_datetime(),
};
let value = serde_json::to_value(AuditLogEntryResponse::from(&row)).expect("test");
assert_eq!(
field_names(&value),
[
"id",
"principal",
"action",
"resource_type",
"resource_id",
"details",
"created_at"
]
);
assert!(value["details"].is_null());
}
#[test]
fn trace_dlq_responses_keep_the_row_wire_shapes() {
let entry = TraceDlqEntry {
id: "d-1".to_string(),
trace_id: "t-1".to_string(),
channel: "orders".to_string(),
payload_json: r#"{"a":1}"#.to_string(),
metadata_json: "{}".to_string(),
error_message: "boom".to_string(),
retry_count: 1,
max_retries: 3,
next_retry_at: sample_datetime(),
created_at: sample_datetime(),
updated_at: sample_datetime(),
};
let value = serde_json::to_value(TraceDlqEntryResponse::from(&entry)).expect("test");
assert_eq!(
field_names(&value),
[
"id",
"trace_id",
"channel",
"payload_json",
"metadata_json",
"error_message",
"retry_count",
"max_retries",
"next_retry_at",
"created_at",
"updated_at"
]
);
let summary = TraceDlqSummary {
id: "d-1".to_string(),
trace_id: "t-1".to_string(),
channel: "orders".to_string(),
error_message: "boom".to_string(),
retry_count: 1,
max_retries: 3,
next_retry_at: sample_datetime(),
created_at: sample_datetime(),
updated_at: sample_datetime(),
};
let value = serde_json::to_value(TraceDlqSummaryResponse::from(&summary)).expect("test");
assert_eq!(
field_names(&value),
[
"id",
"trace_id",
"channel",
"error_message",
"retry_count",
"max_retries",
"next_retry_at",
"created_at",
"updated_at"
]
);
assert!(value.get("payload_json").is_none());
assert!(value.get("metadata_json").is_none());
}
fn sample_model() -> Model {
Model {
model_id: "fraud-scorer".to_string(),
version: 2,
status: EntityStatus::Active.as_str().to_string(),
digest: "sha256:abc".to_string(),
manifest_json: r#"{"abi":"1","name":"fraud-scorer","version":"1.4.0","format":"onnx","inputs":[{"name":"features","dtype":"float32","shape":[1,32]}],"outputs":[{"name":"score","dtype":"float32","shape":[1]}]}"#.to_string(),
artifact_json: r#"{"connector":"models","key":"fraud/1.4.0.onnx","digest":"sha256:abc","size":4096}"#.to_string(),
admission_json: r#"{"state":"passed","node":"node-a","at":"2025-01-01T00:00:00"}"#.to_string(),
stats_json: Some(r#"{"parameters":1200,"nodes":17,"artifact_bytes":4096,"probe_ms":12.5,"ir_version":9,"opset":17,"runtime":"ort","device":"cpu"}"#.to_string()),
tags_json: r#"["fraud"]"#.to_string(),
signature: None,
created_at: sample_datetime(),
updated_at: sample_datetime(),
}
}
#[test]
fn model_response_projects_the_manifest_and_pins_the_wire_shape() {
let response = ModelResponse::try_from(&sample_model()).expect("test");
assert_eq!(response.abi, "1");
assert_eq!(response.model_version, "1.4.0");
assert_eq!(response.format, "onnx");
assert_eq!(response.inputs, ["features"]);
assert_eq!(response.outputs, ["score"]);
assert_eq!(response.artifact.connector, "models");
assert_eq!(response.artifact.key, "fraud/1.4.0.onnx");
assert_eq!(response.artifact.size, Some(4096));
assert_eq!(response.admission.state, "passed");
assert_eq!(response.admission.node.as_deref(), Some("node-a"));
assert!(response.admission.at.is_some());
assert_eq!(response.stats.as_ref().map(|s| s.parameters), Some(1200));
assert_eq!(response.stats.as_ref().map(|s| s.opset), Some(17));
assert!(response.health.is_none());
let value = serde_json::to_value(&response).expect("test");
assert_eq!(
field_names(&value),
[
"model_id",
"version",
"status",
"digest",
"abi",
"model_version",
"format",
"manifest",
"inputs",
"outputs",
"artifact",
"admission",
"stats",
"tags",
"content_hash",
"created_at",
"updated_at"
]
);
assert_eq!(value["tags"], serde_json::json!(["fraud"]));
assert!(
value["content_hash"]
.as_str()
.is_some_and(|h| h.starts_with("sha256:")),
"{value}"
);
let mut admission_keys = field_names(&value["admission"]);
admission_keys.sort();
assert_eq!(admission_keys, ["at", "node", "state"]);
}
#[test]
fn model_response_carries_null_stats_until_admission_passes() {
let mut model = sample_model();
model.stats_json = None;
model.admission_json =
crate::storage::repositories::models::ADMISSION_PENDING_JSON.to_string();
let response = ModelResponse::try_from(&model).expect("test");
assert_eq!(
response.admission,
ModelAdmission {
state: "pending".to_string(),
..Default::default()
}
);
assert!(response.stats.is_none());
let value = serde_json::to_value(&response).expect("test");
assert!(value.get("stats").is_some_and(Value::is_null));
}
#[test]
fn model_response_refuses_a_corrupt_json_column() {
for column in [
"manifest_json",
"artifact_json",
"admission_json",
"stats_json",
"tags_json",
] {
let mut model = sample_model();
match column {
"manifest_json" => model.manifest_json = "nope".to_string(),
"artifact_json" => model.artifact_json = "nope".to_string(),
"admission_json" => model.admission_json = "nope".to_string(),
"stats_json" => model.stats_json = Some("nope".to_string()),
_ => model.tags_json = "nope".to_string(),
}
let err = ModelResponse::try_from(&model).expect_err(column);
assert!(err.to_string().contains(column), "{column}: {err}");
}
}
}