use jsonschema::validator_for;
use schemars::JsonSchema;
use serde::de::DeserializeOwned;
use serde::{Deserialize, Serialize};
use serde_json::{Value, json};
use uuid::Uuid;
use ironflow_store::entities::Signal as StoredSignal;
pub trait Signal: DeserializeOwned + Serialize + JsonSchema {
const NAME: &'static str;
}
pub const SIGNAL_SCHEMA_KEY: &str = "schema";
pub const SIGNAL_TIMED_OUT_KEY: &str = "timed_out";
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct SignalDelivery {
pub signal_id: Uuid,
pub duplicate: bool,
pub resumed: Vec<SignalResumed>,
pub rejected: Vec<SignalRejected>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct SignalResumed {
pub run_id: Uuid,
pub step_id: Uuid,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct SignalRejected {
pub run_id: Uuid,
pub step_id: Uuid,
pub error: String,
}
pub(crate) fn received_output(signal: &StoredSignal) -> Value {
json!({
SIGNAL_TIMED_OUT_KEY: false,
"signal_id": signal.id,
"payload": signal.payload,
})
}
pub(crate) fn timed_out_output() -> Value {
json!({ SIGNAL_TIMED_OUT_KEY: true })
}
pub(crate) fn validate_payload(schema: &Value, payload: &Value) -> Result<(), String> {
let validator = validator_for(schema).map_err(|e| format!("invalid payload schema: {e}"))?;
let errors: Vec<String> = validator
.iter_errors(payload)
.map(|e| e.to_string())
.collect();
if errors.is_empty() {
Ok(())
} else {
Err(errors.join("; "))
}
}
pub(crate) fn validate_step_payload(input: Option<&Value>, payload: &Value) -> Result<(), String> {
let schema = input
.and_then(|i| i.get(SIGNAL_SCHEMA_KEY))
.ok_or_else(|| "step has no stored schema".to_string())?;
validate_payload(schema, payload)
}
#[cfg(test)]
mod tests {
use chrono::Utc;
use schemars::schema_for;
use serde_json::to_value;
use super::*;
#[derive(Serialize, Deserialize, JsonSchema)]
struct PipelineFinished {
status: String,
}
fn schema() -> Value {
to_value(schema_for!(PipelineFinished)).unwrap()
}
#[test]
fn validate_payload_accepts_a_matching_payload() {
assert_eq!(
validate_payload(&schema(), &json!({"status": "success"})),
Ok(())
);
}
#[test]
fn validate_payload_rejects_a_mismatching_payload() {
let err = validate_payload(&schema(), &json!({"state": 1})).unwrap_err();
assert!(err.contains("status"), "got {err}");
}
#[test]
fn validate_payload_rejects_an_invalid_schema() {
let err = validate_payload(&json!({"type": 12}), &json!({})).unwrap_err();
assert!(err.contains("invalid payload schema"), "got {err}");
}
#[test]
fn validate_step_payload_requires_a_stored_schema() {
let err = validate_step_payload(None, &json!({})).unwrap_err();
assert_eq!(err, "step has no stored schema");
let input = json!({ SIGNAL_SCHEMA_KEY: schema() });
assert_eq!(
validate_step_payload(Some(&input), &json!({"status": "ok"})),
Ok(())
);
}
#[test]
fn signal_outputs_carry_the_timeout_flag() {
let signal = StoredSignal {
id: Uuid::now_v7(),
name: "ci.done".to_string(),
key: "abc".to_string(),
payload: json!({"status": "success"}),
idempotency_id: None,
received_at: Utc::now(),
};
let received = received_output(&signal);
assert_eq!(received[SIGNAL_TIMED_OUT_KEY], json!(false));
assert_eq!(received["payload"], signal.payload);
assert_eq!(received["signal_id"], json!(signal.id));
assert_eq!(timed_out_output(), json!({"timed_out": true}));
}
}