fix(desktop): align workflow read/save commands to the frontend contract (#820)

Signed-off-by: Wes <wesbillman@users.noreply.github.com>
Co-authored-by: Brain <21994759fc7a6fa6b965551d35cfd7897d262f2495467f2d78694ddcfa6a5c7e@sprout-oss.stage.blox.sqprod.co>
This commit is contained in:
Wes
2026-06-02 12:52:56 -07:00
committed by GitHub
co-authored by Brain
parent 5b572d6f5e
commit b72eee365f
4 changed files with 405 additions and 123 deletions
+1
View File
@@ -8669,6 +8669,7 @@ dependencies = [
"rubato", "rubato",
"serde", "serde",
"serde_json", "serde_json",
"serde_yaml",
"sha2 0.11.0", "sha2 0.11.0",
"sherpa-onnx", "sherpa-onnx",
"sprout-core", "sprout-core",
+1
View File
@@ -49,6 +49,7 @@ opus = "0.3"
neteq = { version = "0.8", default-features = false } neteq = { version = "0.8", default-features = false }
serde = { version = "1", features = ["derive"] } serde = { version = "1", features = ["derive"] }
serde_json = "1" serde_json = "1"
serde_yaml = "0.9"
nostr = { version = "0.44", features = ["nip44"] } nostr = { version = "0.44", features = ["nip44"] }
zeroize = "1" zeroize = "1"
reqwest = { version = "0.13", features = ["json", "query", "stream"] } reqwest = { version = "0.13", features = ["json", "query", "stream"] }
+194 -123
View File
@@ -1,3 +1,4 @@
use serde::Serialize;
use serde_json::Value; use serde_json::Value;
use tauri::State; use tauri::State;
@@ -7,13 +8,52 @@ use crate::{
relay::{parse_command_response, query_relay, submit_event}, relay::{parse_command_response, query_relay, submit_event},
}; };
// ── Reads ─────────────────────────────────────────────────────────────────── // ── Wire shapes (snake_case, consumed by tauriWorkflows.ts) ──────────────────
/// A workflow definition as the desktop frontend expects it. Mirrors the
/// `RawWorkflow` type in `desktop/src/shared/api/tauriWorkflows.ts`.
///
/// The relay stores a workflow as a single kind:30620 event whose content is
/// the raw YAML. Everything the UI needs is derived from that event:
/// - `id` / `channel_id` from the `d` / `h` tags,
/// - `definition` from parsing the YAML body into a free-form object,
/// - `name` from `definition.name`,
/// - `owner_pubkey` / timestamps from the event itself.
///
/// `status` is always `"active"` here: the relay's disable/archive lifecycle is
/// not reflected back into the kind:30620 event, and the UI derives a
/// "disabled" display state from `definition.enabled` on its own
/// (`getWorkflowDisplayStatus`).
#[derive(Debug, Clone, Serialize, PartialEq)]
pub struct WorkflowWire {
pub id: String,
pub name: String,
pub owner_pubkey: String,
pub channel_id: Option<String>,
pub definition: Value,
pub status: String,
pub created_at: i64,
pub updated_at: i64,
}
/// Response shape for create/update. Mirrors `RawWorkflowSaveResponse` in the
/// frontend: a full workflow record plus an optional webhook secret (only
/// present for webhook-triggered workflows on creation).
#[derive(Debug, Clone, Serialize, PartialEq)]
pub struct WorkflowSaveWire {
#[serde(flatten)]
pub workflow: WorkflowWire,
#[serde(skip_serializing_if = "Option::is_none")]
pub webhook_secret: Option<String>,
}
// ── Reads ────────────────────────────────────────────────────────────────────
#[tauri::command] #[tauri::command]
pub async fn get_channel_workflows( pub async fn get_channel_workflows(
channel_id: String, channel_id: String,
state: State<'_, AppState>, state: State<'_, AppState>,
) -> Result<Value, String> { ) -> Result<Vec<WorkflowWire>, String> {
let events = query_relay( let events = query_relay(
&state, &state,
&[serde_json::json!({ &[serde_json::json!({
@@ -23,15 +63,14 @@ pub async fn get_channel_workflows(
) )
.await?; .await?;
let workflows: Vec<Value> = events.iter().map(workflow_from_event).collect(); Ok(events.iter().map(workflow_from_event).collect())
Ok(serde_json::json!({ "workflows": workflows }))
} }
#[tauri::command] #[tauri::command]
pub async fn get_workflow( pub async fn get_workflow(
workflow_id: String, workflow_id: String,
state: State<'_, AppState>, state: State<'_, AppState>,
) -> Result<Value, String> { ) -> Result<WorkflowWire, String> {
let events = query_relay( let events = query_relay(
&state, &state,
&[serde_json::json!({ &[serde_json::json!({
@@ -52,58 +91,66 @@ pub async fn get_workflow(
pub async fn get_workflow_runs( pub async fn get_workflow_runs(
workflow_id: String, workflow_id: String,
limit: Option<u32>, limit: Option<u32>,
state: State<'_, AppState>, _state: State<'_, AppState>,
) -> Result<Value, String> { ) -> Result<Vec<Value>, String> {
let cap = limit.unwrap_or(50).min(200); // TODO(workflow-runs): Run reconstruction is a clearly-scoped follow-up.
let events = query_relay( // The authoritative run record the frontend's `WorkflowRun` shape needs
&state, // (status / current_step / execution_trace / error_message) lives in the
&[serde_json::json!({ // relay DB and is not exposed to the desktop client as a single queryable
"kinds": [46001, 46002, 46003, 46004, 46005, 46006, 46007, 46010, 46011, 46012], // record. If the relay starts emitting lifecycle events (46001–46007, …),
"#d": [workflow_id], // folding that stream into `WorkflowRun` would be another viable design.
"limit": cap, // The important bit for this command is that raw lifecycle events are not
})], // the `RawWorkflowRun` contract.
) //
.await?; // Until then we return a bare empty array — NOT a raw-event wrapper. The
// frontend wrapper (`getWorkflowRuns`) does `raw.map(fromRawWorkflowRun)`,
let runs: Vec<Value> = events // so it must receive an array; the wrapped `{ runs: [...] }` shape would
.iter() // make `.map()` throw and crash the detail panel (the same TypeError class
.map(|ev| { // as the original page bug). Raw lifecycle events also don't carry the
serde_json::json!({ // `id`/`workflow_id`/`status`/… fields `RawWorkflowRun` expects, so an
"event_id": ev.id.to_hex(), // empty list is the honest, safe placeholder.
"kind": ev.kind.as_u16(), let _ = (workflow_id, limit);
"pubkey": ev.pubkey.to_hex(), Ok(Vec::new())
"created_at": ev.created_at.as_secs(),
"content": ev.content,
"tags": ev.tags.iter().map(|t| t.as_slice().to_vec()).collect::<Vec<_>>(),
})
})
.collect();
Ok(serde_json::json!({ "runs": runs }))
} }
// ── Writes ────────────────────────────────────────────────────────────────── // ── Writes ───────────────────────────────────────────────────────────────────
#[tauri::command] #[tauri::command]
pub async fn create_workflow( pub async fn create_workflow(
channel_id: String, channel_id: String,
yaml_definition: String, yaml_definition: String,
state: State<'_, AppState>, state: State<'_, AppState>,
) -> Result<Value, String> { ) -> Result<WorkflowSaveWire, String> {
let workflow_id = uuid::Uuid::new_v4().to_string(); let workflow_id = uuid::Uuid::new_v4().to_string();
let builder = events::build_workflow_definition(&workflow_id, &channel_id, &yaml_definition)?; let builder = events::build_workflow_definition(&workflow_id, &channel_id, &yaml_definition)?;
let result = submit_event(builder, &state).await?; let result = submit_event(builder, &state).await?;
// The relay returns webhook_secret in the OK response message for new workflows. // The relay returns `webhook_secret` in the OK response message for
let mut response = serde_json::json!({ // webhook-triggered workflows. Everything else in the save record is built
"workflow_id": workflow_id, // locally from the inputs we already hold — the relay's create response
"event_id": result.event_id, // only carries `{ workflow_id, webhook_secret? }`.
}); let webhook_secret = parse_command_response::<Value>(&result.message)
if let Ok(cmd_resp) = parse_command_response::<Value>(&result.message) { .ok()
if let Some(secret) = cmd_resp.get("webhook_secret") { .and_then(|v| {
response["webhook_secret"] = secret.clone(); v.get("webhook_secret")
} .and_then(Value::as_str)
} .map(str::to_string)
Ok(response) });
let now = now_secs();
let workflow = workflow_record(
workflow_id,
Some(channel_id),
current_pubkey_hex(&state)?,
&yaml_definition,
now,
now,
);
Ok(WorkflowSaveWire {
workflow,
webhook_secret,
})
} }
#[tauri::command] #[tauri::command]
@@ -111,9 +158,10 @@ pub async fn update_workflow(
workflow_id: String, workflow_id: String,
yaml_definition: String, yaml_definition: String,
state: State<'_, AppState>, state: State<'_, AppState>,
) -> Result<Value, String> { ) -> Result<WorkflowSaveWire, String> {
// Find the channel id from the existing workflow event so the new event // Find the channel id (and creation time) from the existing workflow event
// carries the same `h` tag — kind:30620 is replaceable by (pubkey, d-tag). // so the new event carries the same `h` tag — kind:30620 is replaceable by
// (pubkey, d-tag).
let prior = query_relay( let prior = query_relay(
&state, &state,
&[serde_json::json!({ &[serde_json::json!({
@@ -124,26 +172,30 @@ pub async fn update_workflow(
) )
.await?; .await?;
let channel_id = prior let prior_event = prior
.first() .first()
.and_then(|ev| {
ev.tags.iter().find_map(|t| {
let s = t.as_slice();
if s.len() >= 2 && s[0] == "h" {
Some(s[1].clone())
} else {
None
}
})
})
.ok_or_else(|| "workflow not found".to_string())?; .ok_or_else(|| "workflow not found".to_string())?;
let channel_id = tag_value(prior_event, "h").ok_or_else(|| "workflow not found".to_string())?;
let created_at = prior_event.created_at.as_secs() as i64;
let builder = events::build_workflow_definition(&workflow_id, &channel_id, &yaml_definition)?; let builder = events::build_workflow_definition(&workflow_id, &channel_id, &yaml_definition)?;
let result = submit_event(builder, &state).await?; submit_event(builder, &state).await?;
Ok(serde_json::json!({
"workflow_id": workflow_id, let updated_at = now_secs();
"event_id": result.event_id, let workflow = workflow_record(
})) workflow_id,
Some(channel_id),
current_pubkey_hex(&state)?,
&yaml_definition,
created_at,
updated_at,
);
Ok(WorkflowSaveWire {
workflow,
// Updates never rotate the webhook secret.
webhook_secret: None,
})
} }
#[tauri::command] #[tauri::command]
@@ -166,38 +218,21 @@ pub async fn trigger_workflow(
Ok(serde_json::json!({ "event_id": result.event_id })) Ok(serde_json::json!({ "event_id": result.event_id }))
} }
// ── Approvals ─────────────────────────────────────────────────────────────── // ── Approvals ────────────────────────────────────────────────────────────────
#[tauri::command] #[tauri::command]
pub async fn get_run_approvals( pub async fn get_run_approvals(
workflow_id: String, workflow_id: String,
run_id: String, run_id: String,
state: State<'_, AppState>, _state: State<'_, AppState>,
) -> Result<Value, String> { ) -> Result<Vec<Value>, String> {
let _ = run_id; // TODO(workflow-runs): Like runs (see `get_workflow_runs`), reconstructing
// Approval-request events for a workflow are kinds 46010/46011/46012. // approvals into the frontend's `WorkflowApproval` shape from lifecycle
let events = query_relay( // events (46010/46011/46012) is a clearly-scoped follow-up tracked under
&state, // TODO(workflow-runs). Return a bare empty array so the frontend's
&[serde_json::json!({ // `getRunApprovals` (`raw.map(fromRawApproval)`) is safe.
"kinds": [46010, 46011, 46012], let _ = (workflow_id, run_id);
"#d": [workflow_id], Ok(Vec::new())
})],
)
.await?;
let approvals: Vec<Value> = events
.iter()
.map(|ev| {
serde_json::json!({
"event_id": ev.id.to_hex(),
"kind": ev.kind.as_u16(),
"pubkey": ev.pubkey.to_hex(),
"created_at": ev.created_at.as_secs(),
"content": ev.content,
"tags": ev.tags.iter().map(|t| t.as_slice().to_vec()).collect::<Vec<_>>(),
})
})
.collect();
Ok(serde_json::json!({ "approvals": approvals }))
} }
#[tauri::command] #[tauri::command]
@@ -222,42 +257,78 @@ pub async fn deny_approval(
Ok(serde_json::json!({ "event_id": result.event_id })) Ok(serde_json::json!({ "event_id": result.event_id }))
} }
// ── Helpers (pure, unit-tested in workflows_tests.rs) ─────────────────────────
fn current_pubkey_hex(state: &AppState) -> Result<String, String> { fn current_pubkey_hex(state: &AppState) -> Result<String, String> {
let keys = state.keys.lock().map_err(|e| e.to_string())?; let keys = state.keys.lock().map_err(|e| e.to_string())?;
Ok(keys.public_key().to_hex()) Ok(keys.public_key().to_hex())
} }
fn workflow_from_event(ev: &nostr::Event) -> Value { fn now_secs() -> i64 {
let workflow_id = ev std::time::SystemTime::now()
.tags .duration_since(std::time::UNIX_EPOCH)
.iter() .map(|d| d.as_secs() as i64)
.find_map(|t| { .unwrap_or_default()
let s = t.as_slice(); }
if s.len() >= 2 && s[0] == "d" {
Some(s[1].clone()) /// First value of the tag whose name matches `name` (e.g. `d`, `h`).
} else { fn tag_value(ev: &nostr::Event, name: &str) -> Option<String> {
None ev.tags.iter().find_map(|t| {
} let s = t.as_slice();
}) (s.len() >= 2 && s[0] == name).then(|| s[1].clone())
.unwrap_or_default();
let channel_id = ev
.tags
.iter()
.find_map(|t| {
let s = t.as_slice();
if s.len() >= 2 && s[0] == "h" {
Some(s[1].clone())
} else {
None
}
})
.unwrap_or_default();
serde_json::json!({
"workflow_id": workflow_id,
"channel_id": channel_id,
"yaml_definition": ev.content,
"event_id": ev.id.to_hex(),
"pubkey": ev.pubkey.to_hex(),
"created_at": ev.created_at.as_secs(),
}) })
} }
/// Parse a workflow's YAML body into a free-form JSON object. The frontend
/// consumes `definition` as `Record<string, unknown>`, so we preserve the full
/// document. On parse failure (or a non-object document) we fall back to an
/// empty object rather than failing the whole list query — a single malformed
/// workflow must not break the page.
fn parse_definition(yaml: &str) -> Value {
match serde_yaml::from_str::<Value>(yaml) {
Ok(v @ Value::Object(_)) => v,
_ => Value::Object(serde_json::Map::new()),
}
}
/// Build a [`WorkflowWire`] record from its parts. Shared by the read path
/// (from a relay event) and the write path (from local inputs).
fn workflow_record(
id: String,
channel_id: Option<String>,
owner_pubkey: String,
yaml_definition: &str,
created_at: i64,
updated_at: i64,
) -> WorkflowWire {
let definition = parse_definition(yaml_definition);
let name = definition
.get("name")
.and_then(Value::as_str)
.filter(|s| !s.trim().is_empty())
.map(str::to_string)
.unwrap_or_else(|| id.clone());
WorkflowWire {
id,
name,
owner_pubkey,
channel_id,
definition,
status: "active".to_string(),
created_at,
updated_at,
}
}
/// Convert a kind:30620 workflow definition event into a [`WorkflowWire`].
fn workflow_from_event(ev: &nostr::Event) -> WorkflowWire {
let id = tag_value(ev, "d").unwrap_or_default();
let channel_id = tag_value(ev, "h");
let ts = ev.created_at.as_secs() as i64;
workflow_record(id, channel_id, ev.pubkey.to_hex(), &ev.content, ts, ts)
}
#[cfg(test)]
#[path = "workflows_tests.rs"]
mod tests;
@@ -0,0 +1,209 @@
// Tests for commands/workflows.rs — split into a sibling file to keep
// workflows.rs focused. These exercise the pure helpers (no relay): event →
// wire conversion, YAML definition parsing, name derivation, and the
// create/update record shaping.
use super::*;
use nostr::{EventBuilder, Keys, Kind, Tag};
/// Build a signed kind:30620 workflow definition event with the given YAML
/// content and d/h tags.
fn wf_event(d: &str, h: &str, yaml: &str) -> nostr::Event {
let keys = Keys::generate();
let tags: Vec<Tag> = [vec!["d", d], vec!["h", h]]
.into_iter()
.map(|t| Tag::parse(t).expect("parse tag"))
.collect();
EventBuilder::new(Kind::Custom(30620), yaml)
.tags(tags)
.sign_with_keys(&keys)
.expect("sign")
}
const CHAN: &str = "11111111-1111-1111-1111-111111111111";
const WF: &str = "22222222-2222-2222-2222-222222222222";
const YAML: &str = "\
name: Greet on join
description: Says hi
enabled: true
trigger:
on: message_posted
filter: hello
steps:
- id: reply
action: post_message
";
#[test]
fn workflow_from_event_maps_all_fields() {
let ev = wf_event(WF, CHAN, YAML);
let wf = workflow_from_event(&ev);
assert_eq!(wf.id, WF);
assert_eq!(wf.channel_id.as_deref(), Some(CHAN));
assert_eq!(wf.owner_pubkey, ev.pubkey.to_hex());
assert_eq!(wf.name, "Greet on join");
assert_eq!(wf.status, "active");
assert_eq!(wf.created_at, ev.created_at.as_secs() as i64);
assert_eq!(wf.updated_at, ev.created_at.as_secs() as i64);
}
#[test]
fn definition_is_parsed_into_object_with_nested_fields() {
let ev = wf_event(WF, CHAN, YAML);
let wf = workflow_from_event(&ev);
// The whole YAML document is preserved as a free-form object.
let def = wf.definition.as_object().expect("definition is an object");
assert_eq!(
def.get("description").and_then(Value::as_str),
Some("Says hi")
);
assert_eq!(def.get("enabled").and_then(Value::as_bool), Some(true));
assert_eq!(
wf.definition.pointer("/trigger/on").and_then(Value::as_str),
Some("message_posted")
);
assert_eq!(
wf.definition
.pointer("/steps/0/action")
.and_then(Value::as_str),
Some("post_message")
);
}
#[test]
fn name_falls_back_to_id_when_missing() {
let yaml = "trigger:\n on: schedule\n cron: '* * * * *'\n";
let ev = wf_event(WF, CHAN, yaml);
let wf = workflow_from_event(&ev);
assert_eq!(wf.name, WF);
}
#[test]
fn name_falls_back_to_id_when_blank() {
let yaml = "name: ' '\ntrigger:\n on: schedule\n";
let ev = wf_event(WF, CHAN, yaml);
let wf = workflow_from_event(&ev);
assert_eq!(wf.name, WF);
}
#[test]
fn malformed_yaml_yields_empty_object_not_error() {
// A broken workflow must not break the whole list — definition falls back
// to an empty object and the name falls back to the id. (YAML is permissive,
// so this uses an unterminated flow mapping that genuinely fails to parse.)
let ev = wf_event(WF, CHAN, "{ name: oops, unterminated: [1, 2");
let wf = workflow_from_event(&ev);
assert_eq!(wf.definition, Value::Object(serde_json::Map::new()));
assert_eq!(wf.name, WF);
}
#[test]
fn scalar_yaml_document_yields_empty_object() {
// A bare scalar parses as valid YAML but isn't an object; treat as empty.
let ev = wf_event(WF, CHAN, "just a string");
let wf = workflow_from_event(&ev);
assert_eq!(wf.definition, Value::Object(serde_json::Map::new()));
}
#[test]
fn tag_value_reads_d_and_h_and_misses_absent() {
let ev = wf_event(WF, CHAN, YAML);
assert_eq!(tag_value(&ev, "d").as_deref(), Some(WF));
assert_eq!(tag_value(&ev, "h").as_deref(), Some(CHAN));
assert_eq!(tag_value(&ev, "z"), None);
}
#[test]
fn workflow_record_shapes_save_inputs() {
let wf = workflow_record(
WF.to_string(),
Some(CHAN.to_string()),
"deadbeef".to_string(),
YAML,
100,
200,
);
assert_eq!(wf.id, WF);
assert_eq!(wf.name, "Greet on join");
assert_eq!(wf.owner_pubkey, "deadbeef");
assert_eq!(wf.channel_id.as_deref(), Some(CHAN));
assert_eq!(wf.created_at, 100);
assert_eq!(wf.updated_at, 200);
assert_eq!(wf.status, "active");
}
#[test]
fn save_wire_serializes_flat_with_optional_secret() {
let workflow = workflow_record(
WF.to_string(),
Some(CHAN.to_string()),
"deadbeef".to_string(),
YAML,
1,
1,
);
// With a secret: present, flattened alongside the workflow fields.
let with = WorkflowSaveWire {
workflow: workflow.clone(),
webhook_secret: Some("s3cr3t".to_string()),
};
let v = serde_json::to_value(&with).expect("serialize");
assert_eq!(v.get("id").and_then(Value::as_str), Some(WF));
assert_eq!(v.get("name").and_then(Value::as_str), Some("Greet on join"));
assert_eq!(
v.get("webhook_secret").and_then(Value::as_str),
Some("s3cr3t")
);
// Without a secret: the key is omitted entirely (frontend treats as null).
let without = WorkflowSaveWire {
workflow,
webhook_secret: None,
};
let v = serde_json::to_value(&without).expect("serialize");
assert!(v.get("webhook_secret").is_none());
assert_eq!(v.get("id").and_then(Value::as_str), Some(WF));
}
#[test]
fn workflow_wire_serializes_with_snake_case_keys() {
// Guard the wire contract the frontend's RawWorkflow depends on.
let ev = wf_event(WF, CHAN, YAML);
let v = serde_json::to_value(workflow_from_event(&ev)).expect("serialize");
for key in [
"id",
"name",
"owner_pubkey",
"channel_id",
"definition",
"status",
"created_at",
"updated_at",
] {
assert!(v.get(key).is_some(), "missing wire key: {key}");
}
}
#[test]
fn runs_and_approvals_serialize_to_bare_empty_array() {
// Regression guard for the crash class this fix closed. The frontend
// wrappers `getWorkflowRuns` / `getRunApprovals` do `raw.map(...)`, so the
// Rust side MUST return a bare JSON array. A wrapped `{ runs: [...] }` /
// `{ approvals: [...] }` shape would make `.map()` throw and crash the
// detail panel — the same TypeError class as the original page bug.
//
// The commands take `State<AppState>`, so we can't invoke them directly in
// a unit test; instead we pin the exact value they return (`Vec::new()` of
// their `Vec<Value>` element type) and assert its serialized shape.
let runs: Vec<Value> = Vec::new();
let approvals: Vec<Value> = Vec::new();
assert_eq!(serde_json::to_string(&runs).expect("serialize runs"), "[]");
assert_eq!(
serde_json::to_string(&approvals).expect("serialize approvals"),
"[]"
);
}