1use crate::common::policy::PolicyConfig;
2use mocra_proxy::ProxyConfig;
3use serde::{Deserialize, Serialize};
4use std::fmt;
5
6#[derive(Serialize, Deserialize, Clone)]
8pub struct Api {
9 pub port: u16,
11 pub api_key: Option<String>,
13 pub rate_limit: Option<f64>,
15}
16
17impl fmt::Debug for Api {
18 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
19 f.debug_struct("Api")
20 .field("port", &self.port)
21 .field("api_key", &self.api_key.as_ref().map(|_| "***REDACTED***"))
22 .field("rate_limit", &self.rate_limit)
23 .finish()
24 }
25}
26
27#[derive(Serialize, Deserialize, Clone)]
29pub struct DatabaseConfig {
30 pub url: Option<String>,
32 pub database_schema: Option<String>,
34 pub pool_size: Option<u32>,
36 pub tls: Option<bool>,
38}
39
40impl fmt::Debug for DatabaseConfig {
41 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
42 f.debug_struct("DatabaseConfig")
43 .field("url", &self.url.as_ref().map(|_| "***REDACTED***"))
44 .field("database_schema", &self.database_schema)
45 .field("pool_size", &self.pool_size)
46 .field("tls", &self.tls)
47 .finish()
48 }
49}
50
51#[derive(Serialize, Deserialize, Debug, Clone)]
53pub struct DownloadConfig {
54 pub downloader_expire: u64,
56 pub timeout: u32,
58 pub rate_limit: f32,
60 pub enable_session: bool,
62 pub enable_locker: bool,
64 pub enable_rate_limit: bool,
66 pub cache_ttl: u64,
68 pub wss_timeout: u32,
70 pub pool_size: Option<usize>,
72 pub max_response_size: Option<usize>,
74}
75
76#[derive(Serialize, Deserialize, Debug, Clone)]
78pub struct SyncConfig {
79 pub kafka: Option<KafkaConfig>,
81 #[serde(default = "default_sync_allow_rollback")]
83 pub allow_rollback: bool,
84 #[serde(default = "default_sync_envelope_enabled")]
86 pub envelope_enabled: bool,
87}
88
89fn default_sync_allow_rollback() -> bool {
90 true
91}
92
93fn default_sync_envelope_enabled() -> bool {
94 false
95}
96
97impl Default for SyncConfig {
98 fn default() -> Self {
99 Self {
100 kafka: None,
101 allow_rollback: true,
102 envelope_enabled: false,
103 }
104 }
105}
106
107#[derive(Serialize, Deserialize, Debug, Clone)]
109pub struct CacheConfig {
110 pub ttl: u64,
112 pub compression_threshold: Option<usize>,
114 pub enable_l1: Option<bool>,
116 pub l1_ttl_secs: Option<u64>,
118 pub l1_max_entries: Option<usize>,
120}
121
122#[derive(Serialize, Deserialize, Debug, Clone)]
124pub struct CrawlerConfig {
125 pub request_max_retries: usize,
127 pub task_max_errors: usize,
129 pub module_max_errors: usize,
131 pub module_locker_ttl: u64,
133 pub node_id: Option<String>,
135 pub task_concurrency: Option<usize>,
137 pub publish_concurrency: Option<usize>,
139 pub parser_concurrency: Option<usize>,
141 pub error_task_concurrency: Option<usize>,
143 pub backpressure_retry_delay_ms: Option<u64>,
145 pub dedup_ttl_secs: Option<u64>,
147 pub idle_stop_secs: Option<u64>,
149}
150
151#[derive(Serialize, Deserialize, Debug, Clone)]
153pub struct SchedulerConfig {
154 pub misfire_tolerance_secs: Option<i64>,
156 pub concurrency: Option<usize>,
158 pub refresh_interval_secs: Option<u64>,
160 pub max_staleness_secs: Option<u64>,
162}
163
164#[derive(Serialize, Deserialize, Clone)]
166pub struct KafkaConfig {
167 pub brokers: String,
169 pub username: Option<String>,
171 pub password: Option<String>,
173 pub tls: Option<bool>,
175}
176
177impl fmt::Debug for KafkaConfig {
178 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
179 f.debug_struct("KafkaConfig")
180 .field("brokers", &self.brokers)
181 .field("username", &self.username)
182 .field(
183 "password",
184 &self.password.as_ref().map(|_| "***REDACTED***"),
185 )
186 .field("tls", &self.tls)
187 .finish()
188 }
189}
190
191#[derive(Serialize, Deserialize, Clone)]
193pub struct NatsConfig {
194 pub url: String,
196 pub username: Option<String>,
198 pub password: Option<String>,
200 pub token: Option<String>,
202}
203
204impl fmt::Debug for NatsConfig {
205 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
206 f.debug_struct("NatsConfig")
207 .field("url", &self.url)
208 .field("username", &self.username)
209 .field(
210 "password",
211 &self.password.as_ref().map(|_| "***REDACTED***"),
212 )
213 .field("token", &self.token.as_ref().map(|_| "***REDACTED***"))
214 .finish()
215 }
216}
217
218#[derive(Serialize, Deserialize, Debug, Clone)]
220pub struct BlobStorageConfig {
221 pub path: Option<String>,
223}
224
225#[derive(Serialize, Deserialize, Debug, Clone)]
227pub struct ChannelConfig {
228 pub blob_storage: Option<BlobStorageConfig>,
230 pub kafka: Option<KafkaConfig>,
232 #[serde(default)]
234 pub nats: Option<NatsConfig>,
235 pub minid_time: u64,
237 pub capacity: usize,
239 pub queue_codec: Option<String>,
241 pub batch_concurrency: Option<usize>,
243 pub compression_threshold: Option<usize>,
245 pub nack_max_retries: Option<u32>,
247 pub nack_backoff_ms: Option<u64>,
249}
250
251use super::logger_config::LoggerConfig;
252
253#[derive(Serialize, Deserialize, Debug, Clone)]
255pub struct EventBusConfig {
256 pub capacity: usize,
258 pub concurrency: usize,
260}
261
262#[derive(Serialize, Deserialize, Debug, Clone)]
264pub struct Config {
265 pub name: String,
267 pub db: DatabaseConfig,
269 pub download_config: DownloadConfig,
271 pub cache: CacheConfig,
273 pub crawler: CrawlerConfig,
275 pub scheduler: Option<SchedulerConfig>,
277 pub sync: Option<SyncConfig>,
279 pub channel_config: ChannelConfig,
281 pub proxy: Option<ProxyConfig>,
283 pub api: Option<Api>,
285 pub event_bus: Option<EventBusConfig>,
287 pub logger: Option<LoggerConfig>,
289 pub policy: Option<PolicyConfig>,
291}
292impl Config {
293 pub fn load(path: &str) -> Result<Self, String> {
295 let config_str = std::fs::read_to_string(path).map_err(|e| e.to_string())?;
297 let config: Config = toml::from_str(&config_str).map_err(|e| e.to_string())?;
298 Ok(config)
299 }
300}
301
302#[cfg(test)]
303mod tests {
304 use super::*;
305 #[allow(unused_imports)]
306 use std::io::Write;
307 #[allow(unused_imports)]
308 use tempfile::NamedTempFile;
309
310 #[test]
311 fn test_config_deserialization() {
312 let toml_str = r#"
313 name = "test_app"
314
315 [db]
316 url = "postgres://user:password@localhost:5432/db"
317 database_schema = "public"
318
319 [download_config]
320 downloader_expire = 3600
321 timeout = 30
322 rate_limit = 10.0
323 enable_session = true
324 enable_locker = false
325 enable_rate_limit = true
326 cache_ttl = 600
327 wss_timeout = 60
328
329 [cache]
330 ttl = 3600
331
332 [crawler]
333 request_max_retries = 3
334 task_max_errors = 5
335 module_max_errors = 10
336 module_locker_ttl = 60
337
338 [channel_config]
339 minid_time = 0
340 capacity = 1000
341 "#;
342
343 let config: Config = toml::from_str(toml_str).unwrap();
344 assert_eq!(config.name, "test_app");
345 assert_eq!(
346 config.db.url.as_deref(),
347 Some("postgres://user:password@localhost:5432/db")
348 );
349 assert_eq!(config.crawler.request_max_retries, 3);
350 assert!(config.proxy.is_none());
351 }
352}