faucet_cli/serve/triggers/
context.rs1use std::collections::BTreeMap;
6
7#[derive(Debug, Clone)]
9pub enum TriggerEvent {
10 Object {
11 bucket: String,
12 key: String,
13 size: u64,
14 last_modified: String, },
16 ObjectBatch {
18 bucket: String,
19 count: usize,
20 watermark: String,
22 },
23 Webhook {
24 method: String,
25 body: String,
26 headers: BTreeMap<String, String>, query: BTreeMap<String, String>,
28 idem: String,
30 },
31 QueueDepth {
32 queue: String,
33 depth: u64,
34 edge: u64,
36 },
37}
38
39impl TriggerEvent {
40 pub fn type_label(&self) -> &'static str {
41 match self {
42 TriggerEvent::Object { .. } | TriggerEvent::ObjectBatch { .. } => "object_arrival",
43 TriggerEvent::Webhook { .. } => "webhook",
44 TriggerEvent::QueueDepth { .. } => "queue_depth",
45 }
46 }
47
48 fn lookup(&self, token: &str, name: &str, fired_at: &str) -> Option<String> {
51 match token {
52 "name" => return Some(name.to_string()),
53 "type" => return Some(self.type_label().to_string()),
54 "fired_at" => return Some(fired_at.to_string()),
55 _ => {}
56 }
57 match self {
58 TriggerEvent::Object {
59 bucket,
60 key,
61 size,
62 last_modified,
63 } => match token {
64 "object_key" => Some(key.clone()),
65 "bucket" => Some(bucket.clone()),
66 "size" => Some(size.to_string()),
67 "last_modified" => Some(last_modified.clone()),
68 _ => None,
69 },
70 TriggerEvent::ObjectBatch { bucket, count, .. } => match token {
71 "bucket" => Some(bucket.clone()),
72 "object_count" => Some(count.to_string()),
73 _ => None,
74 },
75 TriggerEvent::Webhook {
76 method,
77 body,
78 headers,
79 query,
80 ..
81 } => {
82 if let Some(h) = token.strip_prefix("header.") {
85 return headers.get(&h.to_ascii_lowercase()).cloned();
86 }
87 if let Some(q) = token.strip_prefix("query.") {
88 return query.get(q).cloned();
89 }
90 match token {
91 "method" => Some(method.clone()),
92 "body" => Some(body.clone()),
93 _ => None,
94 }
95 }
96 TriggerEvent::QueueDepth { queue, depth, .. } => match token {
97 "queue" => Some(queue.clone()),
98 "depth" => Some(depth.to_string()),
99 _ => None,
100 },
101 }
102 }
103}
104
105pub fn substitute(
110 text: &str,
111 event: &TriggerEvent,
112 name: &str,
113 fired_at: &str,
114) -> Result<String, String> {
115 let mut out = String::with_capacity(text.len());
116 let mut rest = text;
117 while let Some(start) = rest.find("${trigger.") {
118 out.push_str(&rest[..start]);
119 let after = &rest[start + 2..]; let Some(end) = after.find('}') else {
121 return Err("unterminated `${trigger.…}` token".into());
122 };
123 let token = &after[8..end]; match event.lookup(token, name, fired_at) {
125 Some(v) => out.push_str(&yaml_escape(&v)),
126 None => {
127 return Err(format!(
128 "unknown `${{trigger.{token}}}` token for {} trigger '{name}'",
129 event.type_label()
130 ));
131 }
132 }
133 rest = &after[end + 1..];
134 }
135 out.push_str(rest);
136 Ok(out)
137}
138
139fn yaml_escape(v: &str) -> String {
141 format!(
142 "\"{}\"",
143 v.replace('\\', "\\\\")
144 .replace('"', "\\\"")
145 .replace('\n', "\\n")
146 .replace('\r', "\\r")
147 .replace('\t', "\\t")
148 )
149}
150
151pub fn idempotency_key(name: &str, event: &TriggerEvent) -> String {
153 match event {
154 TriggerEvent::Object {
155 bucket,
156 key,
157 last_modified,
158 ..
159 } => {
160 format!("trig:{name}:{bucket}:{key}:{last_modified}")
161 }
162 TriggerEvent::ObjectBatch { watermark, .. } => format!("trig:{name}:{watermark}"),
163 TriggerEvent::Webhook { idem, .. } => idem.clone(),
164 TriggerEvent::QueueDepth { edge, .. } => format!("trig:{name}:edge:{edge}"),
165 }
166}
167
168pub fn labels(name: &str, event: &TriggerEvent) -> BTreeMap<String, String> {
170 let mut m = BTreeMap::new();
171 m.insert("faucet.trigger.name".into(), name.to_string());
172 m.insert("faucet.trigger.type".into(), event.type_label().to_string());
173 match event {
174 TriggerEvent::Object { bucket, key, .. } => {
175 m.insert("faucet.trigger.bucket".into(), bucket.clone());
176 m.insert("faucet.trigger.object_key".into(), key.clone());
177 }
178 TriggerEvent::ObjectBatch { bucket, count, .. } => {
179 m.insert("faucet.trigger.bucket".into(), bucket.clone());
180 m.insert("faucet.trigger.object_count".into(), count.to_string());
181 }
182 TriggerEvent::QueueDepth { queue, depth, .. } => {
183 m.insert("faucet.trigger.queue".into(), queue.clone());
184 m.insert("faucet.trigger.depth".into(), depth.to_string());
185 }
186 TriggerEvent::Webhook { method, .. } => {
187 m.insert("faucet.trigger.method".into(), method.clone());
188 }
189 }
190 m
191}
192
193pub fn render_name(template: &str, event: &TriggerEvent, name: &str, fired_at: &str) -> String {
195 let mut out = String::with_capacity(template.len());
196 let mut rest = template;
197 while let Some(start) = rest.find('{') {
198 out.push_str(&rest[..start]);
199 let after = &rest[start + 1..];
200 if let Some(end) = after.find('}') {
201 let token = &after[..end];
202 out.push_str(&event.lookup(token, name, fired_at).unwrap_or_default());
203 rest = &after[end + 1..];
204 } else {
205 out.push('{');
206 rest = after;
207 }
208 }
209 out.push_str(rest);
210 out
211}
212
213#[cfg(test)]
214mod tests {
215 use super::*;
216
217 fn obj() -> TriggerEvent {
218 TriggerEvent::Object {
219 bucket: "b".into(),
220 key: "incoming/2026/data:set.json".into(),
221 size: 42,
222 last_modified: "2026-06-12T10:00:00Z".into(),
223 }
224 }
225
226 #[test]
227 fn substitutes_object_key_with_yaml_escaping() {
228 let out = substitute(
229 "key: ${trigger.object_key}",
230 &obj(),
231 "t",
232 "2026-06-12T10:00:01Z",
233 )
234 .unwrap();
235 assert_eq!(out, "key: \"incoming/2026/data:set.json\"");
237 }
238
239 #[test]
240 fn unknown_token_errors() {
241 let err = substitute("x: ${trigger.nope}", &obj(), "t", "now").unwrap_err();
242 assert!(err.contains("trigger.nope"), "{err}");
243 }
244
245 #[test]
246 fn webhook_header_and_query_tokens() {
247 let mut headers = BTreeMap::new();
248 headers.insert("x-tenant".into(), "acme".into());
249 let mut query = BTreeMap::new();
250 query.insert("mode".into(), "full".into());
251 let e = TriggerEvent::Webhook {
252 method: "POST".into(),
253 body: "{}".into(),
254 headers,
255 query,
256 idem: "k1".into(),
257 };
258 let out = substitute(
259 "t: ${trigger.header.X-Tenant} m: ${trigger.query.mode}",
260 &e,
261 "h",
262 "now",
263 )
264 .unwrap();
265 assert_eq!(out, "t: \"acme\" m: \"full\"");
266 }
267
268 #[test]
269 fn idempotency_keys_are_deterministic() {
270 assert_eq!(
271 idempotency_key("t", &obj()),
272 "trig:t:b:incoming/2026/data:set.json:2026-06-12T10:00:00Z"
273 );
274 let q = TriggerEvent::QueueDepth {
275 queue: "jobs".into(),
276 depth: 9,
277 edge: 3,
278 };
279 assert_eq!(idempotency_key("d", &q), "trig:d:edge:3");
280 }
281
282 #[test]
283 fn renders_name_template() {
284 let n = render_name("{name}:{object_key}", &obj(), "t", "now");
285 assert_eq!(n, "t:incoming/2026/data:set.json");
286 }
287
288 #[test]
289 fn substitutes_multiline_value_escapes_newline() {
290 let mut headers = BTreeMap::new();
291 let mut query = BTreeMap::new();
292 headers.insert("x-h".into(), "v".into());
293 query.insert("q".into(), "v".into());
294 let e = TriggerEvent::Webhook {
295 method: "POST".into(),
296 body: "line1\nline2".into(),
297 headers,
298 query,
299 idem: "k".into(),
300 };
301 let out = substitute("b: ${trigger.body}", &e, "h", "now").unwrap();
302 assert_eq!(out, r#"b: "line1\nline2""#);
304 assert!(!out.contains('\n'), "raw newline must not appear in output");
305 }
306}