use std::collections::BTreeMap;
#[derive(Debug, Clone)]
pub enum TriggerEvent {
Object {
bucket: String,
key: String,
size: u64,
last_modified: String, },
ObjectBatch {
bucket: String,
count: usize,
watermark: String,
},
Webhook {
method: String,
body: String,
headers: BTreeMap<String, String>, query: BTreeMap<String, String>,
idem: String,
},
QueueDepth {
queue: String,
depth: u64,
edge: u64,
},
}
impl TriggerEvent {
pub fn type_label(&self) -> &'static str {
match self {
TriggerEvent::Object { .. } | TriggerEvent::ObjectBatch { .. } => "object_arrival",
TriggerEvent::Webhook { .. } => "webhook",
TriggerEvent::QueueDepth { .. } => "queue_depth",
}
}
fn lookup(&self, token: &str, name: &str, fired_at: &str) -> Option<String> {
match token {
"name" => return Some(name.to_string()),
"type" => return Some(self.type_label().to_string()),
"fired_at" => return Some(fired_at.to_string()),
_ => {}
}
match self {
TriggerEvent::Object {
bucket,
key,
size,
last_modified,
} => match token {
"object_key" => Some(key.clone()),
"bucket" => Some(bucket.clone()),
"size" => Some(size.to_string()),
"last_modified" => Some(last_modified.clone()),
_ => None,
},
TriggerEvent::ObjectBatch { bucket, count, .. } => match token {
"bucket" => Some(bucket.clone()),
"object_count" => Some(count.to_string()),
_ => None,
},
TriggerEvent::Webhook {
method,
body,
headers,
query,
..
} => {
if let Some(h) = token.strip_prefix("header.") {
return headers.get(&h.to_ascii_lowercase()).cloned();
}
if let Some(q) = token.strip_prefix("query.") {
return query.get(q).cloned();
}
match token {
"method" => Some(method.clone()),
"body" => Some(body.clone()),
_ => None,
}
}
TriggerEvent::QueueDepth { queue, depth, .. } => match token {
"queue" => Some(queue.clone()),
"depth" => Some(depth.to_string()),
_ => None,
},
}
}
}
pub fn substitute(
text: &str,
event: &TriggerEvent,
name: &str,
fired_at: &str,
) -> Result<String, String> {
let mut out = String::with_capacity(text.len());
let mut rest = text;
while let Some(start) = rest.find("${trigger.") {
out.push_str(&rest[..start]);
let after = &rest[start + 2..]; let Some(end) = after.find('}') else {
return Err("unterminated `${trigger.…}` token".into());
};
let token = &after[8..end]; match event.lookup(token, name, fired_at) {
Some(v) => out.push_str(&yaml_escape(&v)),
None => {
return Err(format!(
"unknown `${{trigger.{token}}}` token for {} trigger '{name}'",
event.type_label()
));
}
}
rest = &after[end + 1..];
}
out.push_str(rest);
Ok(out)
}
fn yaml_escape(v: &str) -> String {
format!(
"\"{}\"",
v.replace('\\', "\\\\")
.replace('"', "\\\"")
.replace('\n', "\\n")
.replace('\r', "\\r")
.replace('\t', "\\t")
)
}
pub fn idempotency_key(name: &str, event: &TriggerEvent) -> String {
match event {
TriggerEvent::Object {
bucket,
key,
last_modified,
..
} => {
format!("trig:{name}:{bucket}:{key}:{last_modified}")
}
TriggerEvent::ObjectBatch { watermark, .. } => format!("trig:{name}:{watermark}"),
TriggerEvent::Webhook { idem, .. } => idem.clone(),
TriggerEvent::QueueDepth { edge, .. } => format!("trig:{name}:edge:{edge}"),
}
}
pub fn labels(name: &str, event: &TriggerEvent) -> BTreeMap<String, String> {
let mut m = BTreeMap::new();
m.insert("faucet.trigger.name".into(), name.to_string());
m.insert("faucet.trigger.type".into(), event.type_label().to_string());
match event {
TriggerEvent::Object { bucket, key, .. } => {
m.insert("faucet.trigger.bucket".into(), bucket.clone());
m.insert("faucet.trigger.object_key".into(), key.clone());
}
TriggerEvent::ObjectBatch { bucket, count, .. } => {
m.insert("faucet.trigger.bucket".into(), bucket.clone());
m.insert("faucet.trigger.object_count".into(), count.to_string());
}
TriggerEvent::QueueDepth { queue, depth, .. } => {
m.insert("faucet.trigger.queue".into(), queue.clone());
m.insert("faucet.trigger.depth".into(), depth.to_string());
}
TriggerEvent::Webhook { method, .. } => {
m.insert("faucet.trigger.method".into(), method.clone());
}
}
m
}
pub fn render_name(template: &str, event: &TriggerEvent, name: &str, fired_at: &str) -> String {
let mut out = String::with_capacity(template.len());
let mut rest = template;
while let Some(start) = rest.find('{') {
out.push_str(&rest[..start]);
let after = &rest[start + 1..];
if let Some(end) = after.find('}') {
let token = &after[..end];
out.push_str(&event.lookup(token, name, fired_at).unwrap_or_default());
rest = &after[end + 1..];
} else {
out.push('{');
rest = after;
}
}
out.push_str(rest);
out
}
#[cfg(test)]
mod tests {
use super::*;
fn obj() -> TriggerEvent {
TriggerEvent::Object {
bucket: "b".into(),
key: "incoming/2026/data:set.json".into(),
size: 42,
last_modified: "2026-06-12T10:00:00Z".into(),
}
}
#[test]
fn substitutes_object_key_with_yaml_escaping() {
let out = substitute(
"key: ${trigger.object_key}",
&obj(),
"t",
"2026-06-12T10:00:01Z",
)
.unwrap();
assert_eq!(out, "key: \"incoming/2026/data:set.json\"");
}
#[test]
fn unknown_token_errors() {
let err = substitute("x: ${trigger.nope}", &obj(), "t", "now").unwrap_err();
assert!(err.contains("trigger.nope"), "{err}");
}
#[test]
fn webhook_header_and_query_tokens() {
let mut headers = BTreeMap::new();
headers.insert("x-tenant".into(), "acme".into());
let mut query = BTreeMap::new();
query.insert("mode".into(), "full".into());
let e = TriggerEvent::Webhook {
method: "POST".into(),
body: "{}".into(),
headers,
query,
idem: "k1".into(),
};
let out = substitute(
"t: ${trigger.header.X-Tenant} m: ${trigger.query.mode}",
&e,
"h",
"now",
)
.unwrap();
assert_eq!(out, "t: \"acme\" m: \"full\"");
}
#[test]
fn idempotency_keys_are_deterministic() {
assert_eq!(
idempotency_key("t", &obj()),
"trig:t:b:incoming/2026/data:set.json:2026-06-12T10:00:00Z"
);
let q = TriggerEvent::QueueDepth {
queue: "jobs".into(),
depth: 9,
edge: 3,
};
assert_eq!(idempotency_key("d", &q), "trig:d:edge:3");
}
#[test]
fn renders_name_template() {
let n = render_name("{name}:{object_key}", &obj(), "t", "now");
assert_eq!(n, "t:incoming/2026/data:set.json");
}
#[test]
fn substitutes_multiline_value_escapes_newline() {
let mut headers = BTreeMap::new();
let mut query = BTreeMap::new();
headers.insert("x-h".into(), "v".into());
query.insert("q".into(), "v".into());
let e = TriggerEvent::Webhook {
method: "POST".into(),
body: "line1\nline2".into(),
headers,
query,
idem: "k".into(),
};
let out = substitute("b: ${trigger.body}", &e, "h", "now").unwrap();
assert_eq!(out, r#"b: "line1\nline2""#);
assert!(!out.contains('\n'), "raw newline must not appear in output");
}
}