faucet_cli/serve/triggers/
spec.rs1use schemars::JsonSchema;
5use serde::{Deserialize, Serialize};
6use std::collections::BTreeMap;
7
8#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
10#[serde(deny_unknown_fields)]
11pub struct TriggersFile {
12 pub version: u32,
14 pub triggers: Vec<TriggerSpec>,
15}
16
17#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
21pub struct TriggerSpec {
22 pub name: String,
24 #[serde(default = "default_true")]
26 pub enabled: bool,
27 pub config: PipelineRef,
29 #[serde(default)]
31 pub run: RunTemplate,
32 #[serde(flatten)]
34 pub kind: TriggerKind,
35}
36
37#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
40#[serde(untagged)]
41pub enum PipelineRef {
42 Path(String),
43 Inline(serde_json::Value),
44}
45
46#[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 #[serde(default)]
65 dedupe_header: Option<String>,
66 #[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#[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#[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#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonSchema)]
145#[serde(deny_unknown_fields)]
146pub struct RunTemplate {
147 #[serde(default)]
150 pub name: Option<String>,
151 #[serde(default)]
153 pub labels: BTreeMap<String, String>,
154 #[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}