Skip to main content

faucet_cli/serve/triggers/
spec.rs

1//! Serde config types for the `--triggers` file. Pure data + `JsonSchema`; no IO.
2//! Validation lives in `compiled.rs`.
3
4use schemars::JsonSchema;
5use serde::{Deserialize, Serialize};
6use std::collections::BTreeMap;
7
8/// Top-level `--triggers` document.
9#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
10#[serde(deny_unknown_fields)]
11pub struct TriggersFile {
12    /// Schema version; must be `1`.
13    pub version: u32,
14    pub triggers: Vec<TriggerSpec>,
15}
16
17/// One configured trigger.
18// Note: `deny_unknown_fields` is intentionally absent here — serde does not
19// support it on structs that contain a `#[serde(flatten)]` field.
20#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
21pub struct TriggerSpec {
22    /// Unique name; used in metrics, idempotency keys, and the webhook path.
23    pub name: String,
24    /// Spawn this trigger? Default `true`.
25    #[serde(default = "default_true")]
26    pub enabled: bool,
27    /// The pipeline to enqueue when the trigger fires (path string or inline doc).
28    pub config: PipelineRef,
29    /// Optional run-shaping.
30    #[serde(default)]
31    pub run: RunTemplate,
32    /// Type-specific settings.
33    #[serde(flatten)]
34    pub kind: TriggerKind,
35}
36
37/// A pipeline reference: a path to a config file, OR an inline config document.
38/// Untagged so YAML `config: ./x.yaml` and `config: { pipeline: … }` both parse.
39#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
40#[serde(untagged)]
41pub enum PipelineRef {
42    Path(String),
43    Inline(serde_json::Value),
44}
45
46/// Trigger type + its settings. Internally-tagged on `type` — variant fields
47/// sit flat alongside `type` in the serialized form.
48#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
49#[serde(tag = "type", rename_all = "snake_case")]
50pub enum TriggerKind {
51    ObjectArrival {
52        store: StoreSpec,
53        #[serde(default = "default_poll_secs")]
54        poll_interval_secs: u64,
55        #[serde(default)]
56        mode: ArrivalMode,
57        #[serde(default)]
58        start_at: StartAt,
59    },
60    Webhook {
61        #[serde(default = "default_webhook_methods")]
62        methods: Vec<String>,
63        /// Header whose value is used as the idempotency key (else per-request UUID).
64        #[serde(default)]
65        dedupe_header: Option<String>,
66        /// Leading-edge debounce window in seconds: coalesce fires that arrive
67        /// within this many seconds of the last accepted fire. Default 0 (off).
68        #[serde(default)]
69        debounce_secs: u64,
70    },
71    QueueDepth {
72        queue: QueueSpec,
73        #[serde(default = "default_threshold")]
74        threshold: u64,
75        #[serde(default = "default_poll_secs")]
76        poll_interval_secs: u64,
77    },
78}
79
80#[derive(Debug, Clone, Copy, Default, Serialize, Deserialize, JsonSchema)]
81#[serde(rename_all = "snake_case")]
82pub enum ArrivalMode {
83    #[default]
84    PerObject,
85    Batch,
86}
87
88#[derive(Debug, Clone, Copy, Default, Serialize, Deserialize, JsonSchema)]
89#[serde(rename_all = "snake_case")]
90pub enum StartAt {
91    #[default]
92    Now,
93    Beginning,
94}
95
96/// Object-store connection for `object_arrival`. Internally-tagged on `type`;
97/// variant fields sit flat alongside `type`.
98#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
99#[serde(tag = "type", rename_all = "snake_case")]
100pub enum StoreSpec {
101    S3 {
102        bucket: String,
103        #[serde(default)]
104        prefix: Option<String>,
105        #[serde(default)]
106        region: Option<String>,
107        #[serde(default)]
108        endpoint: Option<String>,
109    },
110    Gcs {
111        bucket: String,
112        #[serde(default)]
113        prefix: Option<String>,
114    },
115}
116
117/// Queue connection for `queue_depth`. Internally-tagged on `type`;
118/// variant fields sit flat alongside `type`.
119#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
120#[serde(tag = "type", rename_all = "snake_case")]
121pub enum QueueSpec {
122    Redis {
123        url: String,
124        key: String,
125        #[serde(default)]
126        kind: RedisQueueKind,
127    },
128    Kafka {
129        brokers: String,
130        topic: String,
131        group: String,
132    },
133}
134
135#[derive(Debug, Clone, Copy, Default, Serialize, Deserialize, JsonSchema)]
136#[serde(rename_all = "snake_case")]
137pub enum RedisQueueKind {
138    #[default]
139    List,
140    Stream,
141}
142
143/// Optional run-shaping applied to the enqueued run.
144#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonSchema)]
145#[serde(deny_unknown_fields)]
146pub struct RunTemplate {
147    /// Run name template; `{field}` tokens resolve from trigger fields
148    /// (`name`, `type`, `object_key`, `bucket`, `queue`, `depth`).
149    #[serde(default)]
150    pub name: Option<String>,
151    /// Static labels merged with the auto-derived trigger labels.
152    #[serde(default)]
153    pub labels: BTreeMap<String, String>,
154    /// Per-run timeout in seconds.
155    #[serde(default)]
156    pub timeout_secs: Option<u64>,
157}
158
159fn default_true() -> bool {
160    true
161}
162fn default_poll_secs() -> u64 {
163    30
164}
165fn default_threshold() -> u64 {
166    1
167}
168fn default_webhook_methods() -> Vec<String> {
169    vec!["POST".to_string()]
170}
171
172#[cfg(test)]
173mod tests {
174    use super::*;
175
176    #[test]
177    fn parses_object_arrival_with_defaults() {
178        let yaml = r#"
179version: 1
180triggers:
181  - name: drop
182    type: object_arrival
183    config: ./pipelines/load.yaml
184    store: { type: s3, bucket: b, prefix: incoming/ }
185"#;
186        let f: TriggersFile = serde_yaml::from_str(yaml).unwrap();
187        assert_eq!(f.version, 1);
188        assert_eq!(f.triggers.len(), 1);
189        let t = &f.triggers[0];
190        assert_eq!(t.name, "drop");
191        assert!(t.enabled);
192        assert!(matches!(t.config, PipelineRef::Path(ref p) if p == "./pipelines/load.yaml"));
193        match &t.kind {
194            TriggerKind::ObjectArrival {
195                poll_interval_secs,
196                mode,
197                start_at,
198                ..
199            } => {
200                assert_eq!(*poll_interval_secs, 30);
201                assert!(matches!(mode, ArrivalMode::PerObject));
202                assert!(matches!(start_at, StartAt::Now));
203            }
204            _ => panic!("wrong kind"),
205        }
206    }
207
208    #[test]
209    fn parses_inline_pipeline_and_webhook_and_queue() {
210        let yaml = r#"
211version: 1
212triggers:
213  - name: hook
214    type: webhook
215    config: { pipeline: { sources: {}, sinks: {} } }
216    dedupe_header: Idempotency-Key
217    debounce_secs: 30
218  - name: drain
219    type: queue_depth
220    config: ./drain.yaml
221    queue: { type: redis, url: "redis://x", key: jobs, kind: stream }
222    threshold: 5
223"#;
224        let f: TriggersFile = serde_yaml::from_str(yaml).unwrap();
225        assert!(matches!(f.triggers[0].config, PipelineRef::Inline(_)));
226        match &f.triggers[0].kind {
227            TriggerKind::Webhook {
228                methods,
229                dedupe_header,
230                debounce_secs,
231            } => {
232                assert_eq!(methods, &vec!["POST".to_string()]);
233                assert_eq!(dedupe_header.as_deref(), Some("Idempotency-Key"));
234                assert_eq!(*debounce_secs, 30);
235            }
236            _ => panic!("wrong kind"),
237        }
238        match &f.triggers[1].kind {
239            TriggerKind::QueueDepth {
240                threshold, queue, ..
241            } => {
242                assert_eq!(*threshold, 5);
243                assert!(matches!(
244                    queue,
245                    QueueSpec::Redis {
246                        kind: RedisQueueKind::Stream,
247                        ..
248                    }
249                ));
250            }
251            _ => panic!("wrong kind"),
252        }
253    }
254
255    #[test]
256    fn rejects_unknown_top_level_field() {
257        let yaml = "version: 1\ntriggers: []\nbogus: 1\n";
258        assert!(serde_yaml::from_str::<TriggersFile>(yaml).is_err());
259    }
260}