use std::fs::read_to_string;
use anyhow::{Context, Result, bail};
use clap::{Args, Subcommand};
use comfy_table::{ContentArrangement, Table};
use ironflow_sdk::IronflowClient;
use ironflow_sdk::client::ListSignalsFilter;
use ironflow_sdk::types::{SendSignalRequest, SignalDeliveryResponse, SignalResponse};
use serde_json::{Value, from_str, from_value, json};
use crate::output;
#[derive(Debug, Args)]
pub struct SignalArgs {
#[command(subcommand)]
pub command: SignalCommands,
}
#[derive(Debug, Subcommand)]
pub enum SignalCommands {
Send {
name: String,
#[arg(long)]
key: String,
#[arg(long)]
payload: Option<String>,
#[arg(long)]
idempotency_id: Option<String>,
},
List {
#[arg(long)]
name: Option<String>,
#[arg(long)]
key: Option<String>,
#[arg(long)]
page: Option<u32>,
#[arg(long)]
per_page: Option<u32>,
},
}
fn read_payload_file(path: &str) -> Result<String> {
read_to_string(path).with_context(|| format!("cannot read payload file '{path}'"))
}
fn parse_payload(raw: Option<&str>) -> Result<Value> {
let Some(raw) = raw else {
return Ok(json!({}));
};
let text = match raw.strip_prefix('@') {
Some(path) => read_payload_file(path)?,
None => raw.to_string(),
};
let payload: Value = from_str(&text).context("--payload must be valid JSON")?;
if !payload.is_object() {
bail!("--payload must be a JSON object");
}
Ok(payload)
}
fn delivery_table(delivery: &SignalDeliveryResponse) -> Table {
let mut table = Table::new();
table.set_content_arrangement(ContentArrangement::Dynamic);
table.set_header(vec!["SIGNAL ID", "DUPLICATE", "RESUMED", "REJECTED"]);
table.add_row(vec![
delivery.signal_id.to_string(),
delivery.duplicate.to_string(),
delivery.resumed.len().to_string(),
delivery.rejected.len().to_string(),
]);
table
}
fn signals_table(items: &[SignalResponse]) -> Table {
let mut table = Table::new();
table.set_content_arrangement(ContentArrangement::Dynamic);
table.set_header(vec!["ID", "NAME", "KEY", "RECEIVED AT"]);
for s in items {
table.add_row(vec![
s.id.to_string(),
s.name.clone(),
s.key.clone(),
s.received_at.to_string(),
]);
}
table
}
pub async fn execute(client: &IronflowClient, args: &SignalArgs, json_mode: bool) -> Result<()> {
match &args.command {
SignalCommands::Send {
name,
key,
payload,
idempotency_id,
} => {
let payload = parse_payload(payload.as_deref())?;
let request: SendSignalRequest = from_value(json!({
"name": name,
"key": key,
"payload": payload,
"idempotency_id": idempotency_id,
}))
.context("cannot build the signal request")?;
let response = client.send_signal(&request).await?;
output::print_output(json_mode, &response, || delivery_table(&response.data))
}
SignalCommands::List {
name,
key,
page,
per_page,
} => {
let filter = ListSignalsFilter {
name: name.clone(),
key: key.clone(),
page: *page,
per_page: *per_page,
};
let response = client.list_signals(&filter).await?;
output::print_output(json_mode, &response, || signals_table(&response.data))
}
}
}
#[cfg(test)]
mod tests {
use std::fs::write;
use tempfile::tempdir;
use super::*;
#[test]
fn parse_payload_defaults_to_an_empty_object() {
assert_eq!(parse_payload(None).unwrap(), json!({}));
}
#[test]
fn parse_payload_accepts_inline_json() {
let payload = parse_payload(Some(r#"{"status":"success"}"#)).unwrap();
assert_eq!(payload, json!({"status": "success"}));
}
#[test]
fn parse_payload_reads_a_file() {
let dir = tempdir().unwrap();
let path = dir.path().join("payload.json");
write(&path, r#"{"sha":"abc"}"#).unwrap();
let payload = parse_payload(Some(&format!("@{}", path.display()))).unwrap();
assert_eq!(payload, json!({"sha": "abc"}));
}
#[test]
fn parse_payload_rejects_invalid_json_and_non_objects() {
assert!(parse_payload(Some("{not json")).is_err());
assert!(parse_payload(Some("[1, 2]")).is_err());
assert!(parse_payload(Some("@/nonexistent/payload.json")).is_err());
}
}