Skip to main content

faucet_source_sqs/
config.rs

1//! Configuration for the SQS source. No I/O here.
2
3use faucet_common_sqs::SqsCredentials;
4use schemars::JsonSchema;
5use serde::{Deserialize, Serialize};
6
7/// SQS long-poll ceiling: `ReceiveMessage` accepts at most a 20-second wait.
8pub const MAX_WAIT_TIME_SECONDS: i32 = 20;
9/// SQS hard limit: at most 10 messages per `ReceiveMessage` call.
10pub const MAX_RECEIVE_BATCH: i32 = 10;
11
12/// Configuration for [`SqsSource`](crate::SqsSource).
13#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
14pub struct SqsSourceConfig {
15    /// SQS queue URL (e.g. `https://sqs.us-east-1.amazonaws.com/1234/my-q`).
16    pub queue_url: String,
17    /// AWS region. `None` uses the SDK default chain (env, profile, IMDS).
18    #[serde(default, skip_serializing_if = "Option::is_none")]
19    pub region: Option<String>,
20    /// Custom endpoint URL for LocalStack / VPC endpoints.
21    #[serde(default, skip_serializing_if = "Option::is_none")]
22    pub endpoint_url: Option<String>,
23    /// AWS credentials. Defaults to the SDK default provider chain.
24    #[serde(default)]
25    pub credentials: SqsCredentials,
26
27    /// Stop after this many seconds without a new message. At least one of
28    /// `idle_timeout_secs` and `max_messages` must be set (mirrors the Kafka /
29    /// Kinesis sources) — a batch run must terminate.
30    #[serde(default, skip_serializing_if = "Option::is_none")]
31    pub idle_timeout_secs: Option<u64>,
32    /// Stop after this many messages in total.
33    #[serde(default, skip_serializing_if = "Option::is_none")]
34    pub max_messages: Option<usize>,
35
36    /// Long-poll wait per `ReceiveMessage` call, in seconds (0–20). Default 10.
37    #[serde(default = "default_wait_time_seconds")]
38    pub wait_time_seconds: i32,
39
40    /// Records per emitted [`StreamPage`](faucet_core::StreamPage). `0` is the
41    /// "no batching" sentinel: one page per assembled drain buffer. Default
42    /// 1000.
43    #[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    /// Minimal config with defaults for everything but the queue URL.
56    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    /// Fail-fast validation, called from `SqsSource::new`.
70    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}