A materialized copy of an M5 group stateful template (cron/queue/email) now reads as `materialized = true` — inherited, read-only — distinct from a hand-authored trigger. Threads a derived `materialized` bool (= `materialized_from IS NOT NULL`) row → domain → API → CLI, mirroring the `sealed`/`shared` columns; `list_for_app` SELECTs `materialized_from`. No migration. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
291 lines
8.6 KiB
Rust
291 lines
8.6 KiB
Rust
//! `pic triggers` subcommands.
|
|
//!
|
|
//! Three per-kind create helpers cover the most common trigger shapes
|
|
//! (KV, cron, dead_letter). For docs/files/pubsub/email/queue the
|
|
//! generic `create-from-json` accepts a raw JSON body so headless
|
|
//! operators don't get blocked waiting for per-kind wrappers.
|
|
|
|
use anyhow::{anyhow, Context, Result};
|
|
use serde_json::json;
|
|
|
|
use crate::client::Client;
|
|
use crate::config;
|
|
use crate::output::{KvBlock, OutputMode, Table};
|
|
|
|
pub async fn ls(app: &str, mode: OutputMode) -> Result<()> {
|
|
let creds = config::resolve()?;
|
|
let client = Client::from_creds(&creds)?;
|
|
let resp = client.triggers_list(app).await?;
|
|
let mut table = Table::new([
|
|
"id",
|
|
"kind",
|
|
"script_id",
|
|
"enabled",
|
|
"dispatch",
|
|
"retry_max",
|
|
"materialized",
|
|
"created_at",
|
|
]);
|
|
for t in resp.triggers {
|
|
table.row([
|
|
t.id,
|
|
t.kind,
|
|
t.script_id.to_string(),
|
|
t.enabled.to_string(),
|
|
t.dispatch_mode,
|
|
t.retry_max_attempts.to_string(),
|
|
t.materialized.to_string(),
|
|
t.created_at.to_rfc3339(),
|
|
]);
|
|
}
|
|
table.print(mode);
|
|
Ok(())
|
|
}
|
|
|
|
/// `pic triggers ls --group <g>` — read-only view of a group's §11 tail trigger
|
|
/// TEMPLATES (event kinds that fan out live to descendant apps). Authoring is
|
|
/// declarative via the `[group]` manifest `[[triggers.*]]`.
|
|
pub async fn ls_group(group: &str, mode: OutputMode) -> Result<()> {
|
|
let creds = config::resolve()?;
|
|
let client = Client::from_creds(&creds)?;
|
|
let rows = client.group_triggers_list(group).await?;
|
|
let mut table = Table::new(["kind", "target", "script", "enabled", "sealed", "shared"]);
|
|
for t in rows {
|
|
table.row([
|
|
t.kind,
|
|
t.target,
|
|
t.script,
|
|
t.enabled.to_string(),
|
|
t.sealed.to_string(),
|
|
t.shared.to_string(),
|
|
]);
|
|
}
|
|
table.print(mode);
|
|
Ok(())
|
|
}
|
|
|
|
pub async fn rm(app: &str, trigger_id: &str) -> Result<()> {
|
|
let creds = config::resolve()?;
|
|
let client = Client::from_creds(&creds)?;
|
|
client.triggers_delete(app, trigger_id).await?;
|
|
println!("Deleted trigger {trigger_id}");
|
|
Ok(())
|
|
}
|
|
|
|
/// Common create body shape — every kind shares `script_id`,
|
|
/// `dispatch_mode`, and the retry knobs. Per-kind fields layer on
|
|
/// top via `serde_json::json!`.
|
|
fn base_body(script_id: &str, dispatch: &str) -> serde_json::Value {
|
|
json!({
|
|
"script_id": script_id,
|
|
"dispatch_mode": dispatch,
|
|
})
|
|
}
|
|
|
|
pub async fn create_kv(
|
|
app: &str,
|
|
script_id: &str,
|
|
collection_glob: &str,
|
|
ops: &[String],
|
|
dispatch: &str,
|
|
mode: OutputMode,
|
|
) -> Result<()> {
|
|
let creds = config::resolve()?;
|
|
let client = Client::from_creds(&creds)?;
|
|
let mut body = base_body(script_id, dispatch);
|
|
body["collection_glob"] = json!(collection_glob);
|
|
if !ops.is_empty() {
|
|
body["ops"] = json!(ops);
|
|
}
|
|
let created = client.triggers_create(app, "kv", &body).await?;
|
|
print_created(&created, mode);
|
|
Ok(())
|
|
}
|
|
|
|
pub async fn create_cron(
|
|
app: &str,
|
|
script_id: &str,
|
|
schedule: &str,
|
|
timezone: Option<&str>,
|
|
dispatch: &str,
|
|
mode: OutputMode,
|
|
) -> Result<()> {
|
|
let creds = config::resolve()?;
|
|
let client = Client::from_creds(&creds)?;
|
|
let mut body = base_body(script_id, dispatch);
|
|
body["schedule"] = json!(schedule);
|
|
if let Some(tz) = timezone {
|
|
body["timezone"] = json!(tz);
|
|
}
|
|
let created = client.triggers_create(app, "cron", &body).await?;
|
|
print_created(&created, mode);
|
|
Ok(())
|
|
}
|
|
|
|
pub async fn create_dead_letter(
|
|
app: &str,
|
|
script_id: &str,
|
|
source_filter: Option<&str>,
|
|
trigger_id_filter: Option<&str>,
|
|
script_id_filter: Option<&str>,
|
|
dispatch: &str,
|
|
mode: OutputMode,
|
|
) -> Result<()> {
|
|
let creds = config::resolve()?;
|
|
let client = Client::from_creds(&creds)?;
|
|
let mut body = base_body(script_id, dispatch);
|
|
if let Some(s) = source_filter {
|
|
body["source_filter"] = json!(s);
|
|
}
|
|
if let Some(t) = trigger_id_filter {
|
|
body["trigger_id_filter"] = json!(t);
|
|
}
|
|
if let Some(s) = script_id_filter {
|
|
body["script_id_filter"] = json!(s);
|
|
}
|
|
let created = client.triggers_create(app, "dead_letter", &body).await?;
|
|
print_created(&created, mode);
|
|
Ok(())
|
|
}
|
|
|
|
/// docs/files share KV's `{collection_glob, ops?}` shape.
|
|
async fn create_collection_trigger(
|
|
kind: &str,
|
|
app: &str,
|
|
script_id: &str,
|
|
collection_glob: &str,
|
|
ops: &[String],
|
|
dispatch: &str,
|
|
mode: OutputMode,
|
|
) -> Result<()> {
|
|
let creds = config::resolve()?;
|
|
let client = Client::from_creds(&creds)?;
|
|
let mut body = base_body(script_id, dispatch);
|
|
body["collection_glob"] = json!(collection_glob);
|
|
if !ops.is_empty() {
|
|
body["ops"] = json!(ops);
|
|
}
|
|
let created = client.triggers_create(app, kind, &body).await?;
|
|
print_created(&created, mode);
|
|
Ok(())
|
|
}
|
|
|
|
pub async fn create_docs(
|
|
app: &str,
|
|
script_id: &str,
|
|
collection_glob: &str,
|
|
ops: &[String],
|
|
dispatch: &str,
|
|
mode: OutputMode,
|
|
) -> Result<()> {
|
|
create_collection_trigger("docs", app, script_id, collection_glob, ops, dispatch, mode).await
|
|
}
|
|
|
|
pub async fn create_files(
|
|
app: &str,
|
|
script_id: &str,
|
|
collection_glob: &str,
|
|
ops: &[String],
|
|
dispatch: &str,
|
|
mode: OutputMode,
|
|
) -> Result<()> {
|
|
create_collection_trigger(
|
|
"files",
|
|
app,
|
|
script_id,
|
|
collection_glob,
|
|
ops,
|
|
dispatch,
|
|
mode,
|
|
)
|
|
.await
|
|
}
|
|
|
|
pub async fn create_pubsub(
|
|
app: &str,
|
|
script_id: &str,
|
|
topic_pattern: &str,
|
|
dispatch: &str,
|
|
mode: OutputMode,
|
|
) -> Result<()> {
|
|
let creds = config::resolve()?;
|
|
let client = Client::from_creds(&creds)?;
|
|
let mut body = base_body(script_id, dispatch);
|
|
body["topic_pattern"] = json!(topic_pattern);
|
|
let created = client.triggers_create(app, "pubsub", &body).await?;
|
|
print_created(&created, mode);
|
|
Ok(())
|
|
}
|
|
|
|
pub async fn create_queue(
|
|
app: &str,
|
|
script_id: &str,
|
|
queue_name: &str,
|
|
visibility_timeout_secs: Option<u32>,
|
|
dispatch: &str,
|
|
mode: OutputMode,
|
|
) -> Result<()> {
|
|
let creds = config::resolve()?;
|
|
let client = Client::from_creds(&creds)?;
|
|
let mut body = base_body(script_id, dispatch);
|
|
body["queue_name"] = json!(queue_name);
|
|
if let Some(v) = visibility_timeout_secs {
|
|
body["visibility_timeout_secs"] = json!(v);
|
|
}
|
|
let created = client.triggers_create(app, "queue", &body).await?;
|
|
print_created(&created, mode);
|
|
Ok(())
|
|
}
|
|
|
|
pub async fn create_email(
|
|
app: &str,
|
|
script_id: &str,
|
|
inbound_secret: Option<&str>,
|
|
mode: OutputMode,
|
|
) -> Result<()> {
|
|
let creds = config::resolve()?;
|
|
let client = Client::from_creds(&creds)?;
|
|
// Email triggers take no dispatch_mode (inbound webhook only) and an
|
|
// optional shared HMAC secret the provider signs POSTs with.
|
|
let mut body = json!({ "script_id": script_id });
|
|
if let Some(secret) = inbound_secret {
|
|
body["inbound_secret"] = json!(secret);
|
|
}
|
|
let created = client.triggers_create(app, "email", &body).await?;
|
|
print_created(&created, mode);
|
|
Ok(())
|
|
}
|
|
|
|
pub async fn create_from_json(app: &str, kind: &str, body: &str, mode: OutputMode) -> Result<()> {
|
|
let creds = config::resolve()?;
|
|
let client = Client::from_creds(&creds)?;
|
|
let body_json: serde_json::Value = if body == "-" {
|
|
let mut s = String::new();
|
|
std::io::Read::read_to_string(&mut std::io::stdin(), &mut s)
|
|
.context("read body JSON from stdin")?;
|
|
serde_json::from_str(&s).map_err(|e| anyhow!("parse JSON from stdin: {e}"))?
|
|
} else if let Some(path) = body.strip_prefix('@') {
|
|
let raw =
|
|
std::fs::read_to_string(path).with_context(|| format!("read JSON body from {path}"))?;
|
|
serde_json::from_str(&raw).map_err(|e| anyhow!("parse JSON body in {path}: {e}"))?
|
|
} else {
|
|
serde_json::from_str(body).map_err(|e| anyhow!("parse inline JSON body: {e}"))?
|
|
};
|
|
let created = client.triggers_create(app, kind, &body_json).await?;
|
|
print_created(&created, mode);
|
|
Ok(())
|
|
}
|
|
|
|
fn print_created(t: &crate::client::TriggerDto, mode: OutputMode) {
|
|
let mut block = KvBlock::new();
|
|
block
|
|
.field("id", t.id.clone())
|
|
.field("kind", t.kind.clone())
|
|
.field("script_id", t.script_id.to_string())
|
|
.field("enabled", t.enabled.to_string())
|
|
.field("dispatch_mode", t.dispatch_mode.clone())
|
|
.field("retry_max_attempts", t.retry_max_attempts.to_string())
|
|
.field("created_at", t.created_at.to_rfc3339());
|
|
block.print(mode);
|
|
}
|