use serde::{Deserialize, Serialize};
use cloacina_api_types::InputSlot;
pub const RESERVED_FIRE_KEYS: &[&str] = &[
"scheduled_time",
"schedule_id",
"schedule_timezone",
"schedule_expression",
"trigger_name",
"triggered_at",
];
#[derive(Debug, thiserror::Error)]
pub enum WorkflowInstanceError {
#[error("unknown param '{0}' — the workflow does not declare it")]
UnknownParam(String),
#[error("required param '{0}' was not supplied and has no default")]
MissingParam(String),
#[error("param '{0}' conflicts with the reserved scheduler key of the same name")]
ReservedParam(String),
#[error("param '{name}' failed schema validation: {message}")]
InvalidParam { name: String, message: String },
#[error(
"param '{0}' is declared as an encrypted secret and must be bound with a \
{{\"$secret\": \"name\"}} reference, not a literal value"
)]
SecretRequiresRef(String),
#[error(
"param '{0}' is not declared as a secret but was bound with a \
{{\"$secret\": \"name\"}} reference"
)]
UnexpectedSecretRef(String),
#[error("param '{name}' has a malformed secret reference: {message}")]
MalformedSecretRef { name: String, message: String },
#[error("serialization: {0}")]
Serialization(String),
}
pub const SECRET_REF_MARKER: &str = "$secret";
pub fn secret_ref_target(value: &serde_json::Value) -> Result<Option<String>, String> {
let serde_json::Value::Object(map) = value else {
return Ok(None);
};
if !map.contains_key(SECRET_REF_MARKER) {
return Ok(None);
}
if map.len() != 1 {
return Err(format!(
"a '{}' reference object must contain only the '{}' key",
SECRET_REF_MARKER, SECRET_REF_MARKER
));
}
match map.get(SECRET_REF_MARKER) {
Some(serde_json::Value::String(name)) if !name.is_empty() => Ok(Some(name.clone())),
_ => Err(format!(
"'{}' must reference a non-empty secret name string",
SECRET_REF_MARKER
)),
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WorkflowInstance {
pub workflow_name: String,
pub params: serde_json::Map<String, serde_json::Value>,
}
impl WorkflowInstance {
pub fn builder(workflow_name: impl Into<String>) -> WorkflowInstanceBuilder {
WorkflowInstanceBuilder {
workflow_name: workflow_name.into(),
supplied: serde_json::Map::new(),
}
}
pub fn from_resolved(
workflow_name: impl Into<String>,
params: serde_json::Map<String, serde_json::Value>,
) -> Self {
Self {
workflow_name: workflow_name.into(),
params,
}
}
pub fn params_json(&self) -> Result<String, WorkflowInstanceError> {
serde_json::to_string(&self.params)
.map_err(|e| WorkflowInstanceError::Serialization(e.to_string()))
}
}
pub struct WorkflowInstanceBuilder {
workflow_name: String,
supplied: serde_json::Map<String, serde_json::Value>,
}
impl WorkflowInstanceBuilder {
pub fn param(
mut self,
name: impl Into<String>,
value: impl Serialize,
) -> Result<Self, WorkflowInstanceError> {
let name = name.into();
let value = serde_json::to_value(value)
.map_err(|e| WorkflowInstanceError::Serialization(e.to_string()))?;
self.supplied.insert(name, value);
Ok(self)
}
pub fn build(self, declared: &[InputSlot]) -> Result<WorkflowInstance, WorkflowInstanceError> {
for name in self.supplied.keys() {
if !declared.iter().any(|s| &s.name == name) {
return Err(WorkflowInstanceError::UnknownParam(name.clone()));
}
if RESERVED_FIRE_KEYS.contains(&name.as_str()) {
return Err(WorkflowInstanceError::ReservedParam(name.clone()));
}
}
let mut resolved = serde_json::Map::new();
for slot in declared {
if RESERVED_FIRE_KEYS.contains(&slot.name.as_str()) {
return Err(WorkflowInstanceError::ReservedParam(slot.name.clone()));
}
match self.supplied.get(&slot.name) {
Some(v) => {
let is_ref = secret_ref_target(v)
.map_err(|message| WorkflowInstanceError::MalformedSecretRef {
name: slot.name.clone(),
message,
})?
.is_some();
if slot.encrypted && !is_ref {
return Err(WorkflowInstanceError::SecretRequiresRef(slot.name.clone()));
}
if !slot.encrypted && is_ref {
return Err(WorkflowInstanceError::UnexpectedSecretRef(
slot.name.clone(),
));
}
resolved.insert(slot.name.clone(), v.clone());
}
None => match &slot.default {
Some(d) => {
resolved.insert(slot.name.clone(), d.clone());
}
None if slot.required => {
return Err(WorkflowInstanceError::MissingParam(slot.name.clone()));
}
None => {}
},
}
}
Ok(WorkflowInstance {
workflow_name: self.workflow_name,
params: resolved,
})
}
}
pub fn merge_instance_params(
context: &mut crate::Context<serde_json::Value>,
params_json: &str,
) -> Result<(), String> {
let params: serde_json::Map<String, serde_json::Value> = serde_json::from_str(params_json)
.map_err(|e| format!("instance params JSON parse: {}", e))?;
let mut secret_refs = serde_json::Map::new();
for (k, v) in params {
if RESERVED_FIRE_KEYS.contains(&k.as_str()) {
continue;
}
if k == cloacina_workflow::secret::SECRET_REFS_KEY {
return Err(format!(
"instance param '{}' collides with the reserved secret-reference key",
k
));
}
match secret_ref_target(&v).map_err(|m| {
format!(
"instance param '{}' has a malformed secret reference: {}",
k, m
)
})? {
Some(secret_name) => {
secret_refs.insert(k, serde_json::Value::String(secret_name));
}
None => {
if context.update(&k, v.clone()).is_err() {
context
.insert(k.as_str(), v)
.map_err(|e| format!("instance param '{}' insert: {}", k, e))?;
}
}
}
}
if !secret_refs.is_empty() {
let alias_map = serde_json::Value::Object(secret_refs);
let key = cloacina_workflow::secret::SECRET_REFS_KEY;
if context.update(key, alias_map.clone()).is_err() {
context
.insert(key, alias_map)
.map_err(|e| format!("secret-reference map insert: {}", e))?;
}
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
fn slot(name: &str, required: bool, default: Option<serde_json::Value>) -> InputSlot {
InputSlot {
name: name.to_string(),
required,
default,
schema: serde_json::json!({"type": "string"}),
encrypted: false,
}
}
#[test]
fn build_resolves_defaults_and_validates() {
let declared = vec![
slot("source", true, None),
slot("mode", false, Some(serde_json::json!("copy"))),
];
let err = WorkflowInstance::builder("sync")
.build(&declared)
.unwrap_err();
assert!(matches!(err, WorkflowInstanceError::MissingParam(_)));
let err = WorkflowInstance::builder("sync")
.param("nope", "x")
.unwrap()
.build(&declared)
.unwrap_err();
assert!(matches!(err, WorkflowInstanceError::UnknownParam(_)));
let inst = WorkflowInstance::builder("sync")
.param("source", "/a")
.unwrap()
.build(&declared)
.unwrap();
assert_eq!(inst.params["source"], serde_json::json!("/a"));
assert_eq!(inst.params["mode"], serde_json::json!("copy"));
let json = serde_json::to_string(&inst).unwrap();
let back: WorkflowInstance = serde_json::from_str(&json).unwrap();
assert_eq!(back.params, inst.params);
}
#[test]
fn merge_skips_reserved_and_overrides_payload() {
let mut ctx: crate::Context<serde_json::Value> = crate::Context::new();
ctx.insert("from_trigger", serde_json::json!("payload"))
.unwrap();
let params = serde_json::json!({
"source": "/a",
"from_trigger": "bound-wins",
"scheduled_time": "spoof-attempt"
})
.to_string();
merge_instance_params(&mut ctx, ¶ms).unwrap();
assert_eq!(ctx.get("source").cloned().unwrap(), serde_json::json!("/a"));
assert_eq!(
ctx.get("from_trigger").cloned().unwrap(),
serde_json::json!("bound-wins")
);
assert!(ctx.get("scheduled_time").is_none());
}
fn secret_slot(name: &str) -> InputSlot {
InputSlot::secret(name)
}
#[test]
fn secret_ref_target_classifies_values() {
assert_eq!(
secret_ref_target(&serde_json::json!({"$secret": "db_prod"})).unwrap(),
Some("db_prod".to_string())
);
assert_eq!(
secret_ref_target(&serde_json::json!("plain")).unwrap(),
None
);
assert_eq!(
secret_ref_target(&serde_json::json!({"host": "x"})).unwrap(),
None
);
assert!(secret_ref_target(&serde_json::json!({"$secret": "a", "x": 1})).is_err());
assert!(secret_ref_target(&serde_json::json!({"$secret": 5})).is_err());
assert!(secret_ref_target(&serde_json::json!({"$secret": ""})).is_err());
}
#[test]
fn merge_routes_secret_refs_away_from_plaintext_context() {
let mut ctx: crate::Context<serde_json::Value> = crate::Context::new();
let params = serde_json::json!({
"region": "us-east-1",
"dst_credentials": {"$secret": "s3_prod"}
})
.to_string();
merge_instance_params(&mut ctx, ¶ms).unwrap();
assert_eq!(
ctx.get("region").cloned().unwrap(),
serde_json::json!("us-east-1")
);
assert!(ctx.get("dst_credentials").is_none());
let refs = ctx
.get(cloacina_workflow::secret::SECRET_REFS_KEY)
.cloned()
.unwrap();
assert_eq!(refs, serde_json::json!({"dst_credentials": "s3_prod"}));
let json = ctx.to_json().unwrap();
assert!(!json.contains("$secret"));
assert!(json.contains("s3_prod")); }
#[test]
fn merge_rejects_malformed_secret_ref_and_reserved_key() {
let mut ctx: crate::Context<serde_json::Value> = crate::Context::new();
let bad = serde_json::json!({"cred": {"$secret": 3}}).to_string();
assert!(merge_instance_params(&mut ctx, &bad).is_err());
let mut ctx2: crate::Context<serde_json::Value> = crate::Context::new();
let reserved = serde_json::json!({
cloacina_workflow::secret::SECRET_REFS_KEY: {"x": "y"}
})
.to_string();
assert!(merge_instance_params(&mut ctx2, &reserved).is_err());
}
#[test]
fn build_requires_secret_slot_bound_with_ref() {
let declared = vec![slot("source", true, None), secret_slot("db")];
let err = WorkflowInstance::builder("sync")
.param("source", "/a")
.unwrap()
.build(&declared)
.unwrap_err();
assert!(matches!(err, WorkflowInstanceError::MissingParam(_)));
let err = WorkflowInstance::builder("sync")
.param("source", "/a")
.unwrap()
.param("db", "plaintext-password")
.unwrap()
.build(&declared)
.unwrap_err();
assert!(matches!(err, WorkflowInstanceError::SecretRequiresRef(_)));
let inst = WorkflowInstance::builder("sync")
.param("source", "/a")
.unwrap()
.param("db", serde_json::json!({"$secret": "db_prod"}))
.unwrap()
.build(&declared)
.unwrap();
assert_eq!(inst.params["db"], serde_json::json!({"$secret": "db_prod"}));
}
#[test]
fn build_rejects_secret_ref_on_plaintext_slot() {
let declared = vec![slot("mode", true, None)];
let err = WorkflowInstance::builder("sync")
.param("mode", serde_json::json!({"$secret": "leak"}))
.unwrap()
.build(&declared)
.unwrap_err();
assert!(matches!(err, WorkflowInstanceError::UnexpectedSecretRef(_)));
}
}