Skip to main content

faucet_cli/serve/triggers/
context.rs

1//! The fired event (`TriggerEvent`) and pure derivations from it:
2//! `${trigger.*}` text substitution, the deterministic idempotency key, and the
3//! auto run-labels. No IO.
4
5use std::collections::BTreeMap;
6
7/// A concrete event that fired a trigger.
8#[derive(Debug, Clone)]
9pub enum TriggerEvent {
10    Object {
11        bucket: String,
12        key: String,
13        size: u64,
14        last_modified: String, // RFC3339
15    },
16    /// One poll-cycle's worth of new objects (batch mode).
17    ObjectBatch {
18        bucket: String,
19        count: usize,
20        /// Max last_modified seen this batch (for the idempotency key).
21        watermark: String,
22    },
23    Webhook {
24        method: String,
25        body: String,
26        headers: BTreeMap<String, String>, // lowercased header names
27        query: BTreeMap<String, String>,
28        /// Pre-resolved idempotency key (dedupe header value or a UUID).
29        idem: String,
30    },
31    QueueDepth {
32        queue: String,
33        depth: u64,
34        /// Rising-edge ordinal (so re-arm + re-cross → distinct run).
35        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    /// Map `${trigger.<token>}` → value. `fired_at` is supplied by the caller
49    /// (so this stays pure / clock-free).
50    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                // HTTP header names are case-insensitive (looked up lowercased); URI query
83                // keys are case-sensitive per the URI spec and matched verbatim.
84                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
105/// Substitute every `${trigger.<token>}` in `text`. Substituted values are
106/// YAML-escaped (wrapped + quotes doubled) so a value containing `:`/quotes can
107/// land in a scalar position without breaking the document. An unknown token is
108/// an error (never silently passed through). Returns the substituted text.
109pub 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..]; // skip "${"
120        let Some(end) = after.find('}') else {
121            return Err("unterminated `${trigger.…}` token".into());
122        };
123        let token = &after[8..end]; // after "trigger."
124        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
139/// Quote a value for safe scalar substitution into YAML.
140fn 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
151/// Deterministic idempotency key for the event (see spec §9).
152pub 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
168/// Auto labels attached to the enqueued run (low-cardinality by construction).
169pub 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
193/// Render a `{field}` run-name template against trigger fields.
194pub 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        // The ':' in the key must be quoted so YAML stays valid.
236        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        // Must contain literal backslash-n, not a raw newline.
303        assert_eq!(out, r#"b: "line1\nline2""#);
304        assert!(!out.contains('\n'), "raw newline must not appear in output");
305    }
306}