faucet_source_sqs/
config.rs1use faucet_common_sqs::SqsCredentials;
4use schemars::JsonSchema;
5use serde::{Deserialize, Serialize};
6
7pub const MAX_WAIT_TIME_SECONDS: i32 = 20;
9pub const MAX_RECEIVE_BATCH: i32 = 10;
11
12#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
14pub struct SqsSourceConfig {
15 pub queue_url: String,
17 #[serde(default, skip_serializing_if = "Option::is_none")]
19 pub region: Option<String>,
20 #[serde(default, skip_serializing_if = "Option::is_none")]
22 pub endpoint_url: Option<String>,
23 #[serde(default)]
25 pub credentials: SqsCredentials,
26
27 #[serde(default, skip_serializing_if = "Option::is_none")]
31 pub idle_timeout_secs: Option<u64>,
32 #[serde(default, skip_serializing_if = "Option::is_none")]
34 pub max_messages: Option<usize>,
35
36 #[serde(default = "default_wait_time_seconds")]
38 pub wait_time_seconds: i32,
39
40 #[serde(default = "default_batch_size")]
44 pub batch_size: usize,
45}
46
47fn default_wait_time_seconds() -> i32 {
48 10
49}
50fn default_batch_size() -> usize {
51 faucet_core::DEFAULT_BATCH_SIZE
52}
53
54impl SqsSourceConfig {
55 pub fn new(queue_url: impl Into<String>) -> Self {
57 Self {
58 queue_url: queue_url.into(),
59 region: None,
60 endpoint_url: None,
61 credentials: SqsCredentials::default(),
62 idle_timeout_secs: None,
63 max_messages: None,
64 wait_time_seconds: default_wait_time_seconds(),
65 batch_size: default_batch_size(),
66 }
67 }
68
69 pub fn validate(&self) -> Result<(), faucet_core::FaucetError> {
71 use faucet_core::FaucetError;
72 if self.queue_url.trim().is_empty() {
73 return Err(FaucetError::Config(
74 "sqs source: queue_url must not be empty".into(),
75 ));
76 }
77 faucet_core::validate_batch_size(self.batch_size)?;
78 if self.wait_time_seconds < 0 || self.wait_time_seconds > MAX_WAIT_TIME_SECONDS {
79 return Err(FaucetError::Config(format!(
80 "sqs source: wait_time_seconds must be 0..={MAX_WAIT_TIME_SECONDS} (got {})",
81 self.wait_time_seconds
82 )));
83 }
84 if self.idle_timeout_secs.is_none() && self.max_messages.is_none() {
85 return Err(FaucetError::Config(
86 "sqs source: set at least one of idle_timeout_secs / max_messages so a run can \
87 terminate (mirrors the kafka/kinesis sources)"
88 .into(),
89 ));
90 }
91 if self.idle_timeout_secs == Some(0) {
92 return Err(FaucetError::Config(
93 "sqs source: idle_timeout_secs must be at least 1".into(),
94 ));
95 }
96 if self.max_messages == Some(0) {
97 return Err(FaucetError::Config(
98 "sqs source: max_messages must be at least 1".into(),
99 ));
100 }
101 Ok(())
102 }
103}
104
105#[cfg(test)]
106mod tests {
107 use super::*;
108
109 fn valid() -> SqsSourceConfig {
110 let mut c = SqsSourceConfig::new("https://sqs.us-east-1.amazonaws.com/1/events");
111 c.max_messages = Some(100);
112 c
113 }
114
115 #[test]
116 fn defaults_are_sensible() {
117 let c = SqsSourceConfig::new("https://q");
118 assert_eq!(c.wait_time_seconds, 10);
119 assert_eq!(c.batch_size, faucet_core::DEFAULT_BATCH_SIZE);
120 assert!(c.idle_timeout_secs.is_none());
121 assert!(c.max_messages.is_none());
122 }
123
124 #[test]
125 fn validation_bounds() {
126 valid().validate().unwrap();
127
128 let mut c = valid();
129 c.queue_url = " ".into();
130 assert!(c.validate().is_err());
131
132 let mut c = valid();
133 c.wait_time_seconds = 21;
134 assert!(c.validate().is_err());
135 c.wait_time_seconds = -1;
136 assert!(c.validate().is_err());
137
138 let mut c = valid();
139 c.max_messages = None;
140 c.idle_timeout_secs = None;
141 let err = c.validate().unwrap_err();
142 assert!(err.to_string().contains("idle_timeout_secs"), "{err}");
143
144 let mut c = valid();
145 c.idle_timeout_secs = Some(0);
146 assert!(c.validate().is_err());
147 let mut c = valid();
148 c.max_messages = Some(0);
149 assert!(c.validate().is_err());
150
151 let mut c = valid();
152 c.batch_size = faucet_core::MAX_BATCH_SIZE + 1;
153 assert!(c.validate().is_err());
154 }
155
156 #[test]
157 fn full_config_parses_from_yaml() {
158 let yaml = r#"
159queue_url: https://sqs.us-east-1.amazonaws.com/1/events
160region: us-east-1
161endpoint_url: http://127.0.0.1:4566
162credentials: { type: access_key, config: { access_key_id: a, secret_access_key: b } }
163idle_timeout_secs: 5
164max_messages: 250
165wait_time_seconds: 5
166batch_size: 100
167"#;
168 let c: SqsSourceConfig = serde_yaml::from_str(yaml).unwrap();
169 c.validate().unwrap();
170 assert_eq!(c.wait_time_seconds, 5);
171 assert_eq!(c.batch_size, 100);
172 assert_eq!(c.idle_timeout_secs, Some(5));
173 assert_eq!(c.max_messages, Some(250));
174 }
175}