Skip to main content

mocra_core/common/model/
config.rs

1use crate::common::policy::PolicyConfig;
2use mocra_proxy::ProxyConfig;
3use serde::{Deserialize, Serialize};
4use std::fmt;
5
6/// API Configuration
7#[derive(Serialize, Deserialize, Clone)]
8pub struct Api {
9    /// Port number for the API server
10    pub port: u16,
11    /// Optional API key for authentication
12    pub api_key: Option<String>,
13    /// Optional Rate Limit (requests per second)
14    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/// Database Configuration
28#[derive(Serialize, Deserialize, Clone)]
29pub struct DatabaseConfig {
30    /// Connection URL
31    pub url: Option<String>,
32    /// Database schema
33    pub database_schema: Option<String>,
34    /// Connection pool size
35    pub pool_size: Option<u32>,
36    /// Enable TLS (sslmode=require)
37    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/// Downloader Configuration
52#[derive(Serialize, Deserialize, Debug, Clone)]
53pub struct DownloadConfig {
54    /// Downloader expiration time in seconds
55    pub downloader_expire: u64,
56    /// Request timeout in seconds
57    pub timeout: u32,
58    /// Rate limit in requests per second
59    pub rate_limit: f32,
60    /// Enable session cache/sync behavior
61    pub enable_session: bool,
62    /// Enable distributed locking
63    pub enable_locker: bool,
64    /// Enable rate limiting
65    pub enable_rate_limit: bool,
66    /// Cache TTL in seconds
67    pub cache_ttl: u64,
68    /// WebSocket timeout in seconds
69    pub wss_timeout: u32,
70    /// Connection pool size for HTTP client (default: 200)
71    pub pool_size: Option<usize>,
72    /// Maximum response size in bytes (default: 10MB)
73    pub max_response_size: Option<usize>,
74}
75
76/// Synchronization Configuration
77#[derive(Serialize, Deserialize, Debug, Clone)]
78pub struct SyncConfig {
79    /// Kafka configuration for synchronization
80    pub kafka: Option<KafkaConfig>,
81    /// Allow rollback to older versions (default: true)
82    #[serde(default = "default_sync_allow_rollback")]
83    pub allow_rollback: bool,
84    /// Enable versioned envelope for sync payloads (default: false)
85    #[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/// Cache Configuration
108#[derive(Serialize, Deserialize, Debug, Clone)]
109pub struct CacheConfig {
110    /// Default TTL for cache items
111    pub ttl: u64,
112    /// Compression threshold in bytes (payloads larger than this will be compressed)
113    pub compression_threshold: Option<usize>,
114    /// Enable L1 local cache layer (default: false)
115    pub enable_l1: Option<bool>,
116    /// L1 cache TTL in seconds (default: 30)
117    pub l1_ttl_secs: Option<u64>,
118    /// L1 cache max entries before eviction (default: 10000)
119    pub l1_max_entries: Option<usize>,
120}
121
122/// Crawler Configuration
123#[derive(Serialize, Deserialize, Debug, Clone)]
124pub struct CrawlerConfig {
125    /// Maximum retries for failed requests
126    pub request_max_retries: usize,
127    /// Maximum allowed errors per task
128    pub task_max_errors: usize,
129    /// Maximum allowed errors per module
130    pub module_max_errors: usize,
131    /// TTL for module locks
132    pub module_locker_ttl: u64,
133    /// Optional node ID (useful for stable node identity across restarts)
134    pub node_id: Option<String>,
135    /// Concurrency for task processor (JS execution)
136    pub task_concurrency: Option<usize>,
137    /// Concurrency for request publishing
138    pub publish_concurrency: Option<usize>,
139    /// Concurrency for parser task processor
140    pub parser_concurrency: Option<usize>,
141    /// Concurrency for error task processor
142    pub error_task_concurrency: Option<usize>,
143    /// Delay in milliseconds before retrying when backpressure is detected
144    pub backpressure_retry_delay_ms: Option<u64>,
145    /// Request deduplication TTL in seconds (default: 3600)
146    pub dedup_ttl_secs: Option<u64>,
147    /// Idle stop timeout in seconds (stop engine if local queues have no data for this duration)
148    pub idle_stop_secs: Option<u64>,
149}
150
151/// Scheduler Configuration
152#[derive(Serialize, Deserialize, Debug, Clone)]
153pub struct SchedulerConfig {
154    /// Maximum tolerance in seconds for missed cron jobs (default: 300)
155    pub misfire_tolerance_secs: Option<i64>,
156    /// Max concurrency for processing scheduled contexts (default: 100)
157    pub concurrency: Option<usize>,
158    /// Refresh interval in seconds for scheduler cache (default: 60)
159    pub refresh_interval_secs: Option<u64>,
160    /// Max allowed staleness in seconds before forcing refresh (default: 120)
161    pub max_staleness_secs: Option<u64>,
162}
163
164/// Kafka Configuration
165#[derive(Serialize, Deserialize, Clone)]
166pub struct KafkaConfig {
167    /// Comma-separated list of brokers
168    pub brokers: String,
169    /// SASL username
170    pub username: Option<String>,
171    /// SASL password
172    pub password: Option<String>,
173    /// Enable TLS
174    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/// NATS (JetStream) Configuration
192#[derive(Serialize, Deserialize, Clone)]
193pub struct NatsConfig {
194    /// Server address (e.g. `nats://127.0.0.1:4222`; comma-separated for multiple).
195    pub url: String,
196    /// Optional username.
197    pub username: Option<String>,
198    /// Optional password.
199    pub password: Option<String>,
200    /// Optional token authentication.
201    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/// Blob Storage Configuration
219#[derive(Serialize, Deserialize, Debug, Clone)]
220pub struct BlobStorageConfig {
221    /// Local file system path for blob storage
222    pub path: Option<String>,
223}
224
225/// Channel (Queue) Configuration
226#[derive(Serialize, Deserialize, Debug, Clone)]
227pub struct ChannelConfig {
228    /// Blob storage configuration
229    pub blob_storage: Option<BlobStorageConfig>,
230    /// Kafka configuration for queue
231    pub kafka: Option<KafkaConfig>,
232    /// NATS (JetStream) configuration for queue
233    #[serde(default)]
234    pub nats: Option<NatsConfig>,
235    /// Minimum ID time (snowflake/uuid related)
236    pub minid_time: u64,
237    /// Channel capacity
238    pub capacity: usize,
239    /// Queue codec: json | msgpack
240    pub queue_codec: Option<String>,
241    /// Concurrency limit for batch flushing (default: 10)
242    pub batch_concurrency: Option<usize>,
243    /// Compression threshold in bytes (payloads larger than this will be compressed)
244    pub compression_threshold: Option<usize>,
245    /// Max retries for NACK before sending to DLQ (default: 0)
246    pub nack_max_retries: Option<u32>,
247    /// Backoff in milliseconds before retrying NACK (default: 0)
248    pub nack_backoff_ms: Option<u64>,
249}
250
251use super::logger_config::LoggerConfig;
252
253/// Event Bus Configuration
254#[derive(Serialize, Deserialize, Debug, Clone)]
255pub struct EventBusConfig {
256    /// Channel capacity for events (default: 1024)
257    pub capacity: usize,
258    /// Concurrency limit for event handlers (default: 64)
259    pub concurrency: usize,
260}
261
262/// Main Configuration
263#[derive(Serialize, Deserialize, Debug, Clone)]
264pub struct Config {
265    /// Application instance name
266    pub name: String,
267    /// Database configuration
268    pub db: DatabaseConfig,
269    /// Download configuration
270    pub download_config: DownloadConfig,
271    /// Cache configuration
272    pub cache: CacheConfig,
273    /// Crawler behavior configuration
274    pub crawler: CrawlerConfig,
275    /// Scheduler configuration
276    pub scheduler: Option<SchedulerConfig>,
277    /// Synchronization configuration
278    pub sync: Option<SyncConfig>,
279    /// Message channel configuration
280    pub channel_config: ChannelConfig,
281    /// Proxy configuration
282    pub proxy: Option<ProxyConfig>,
283    /// API server configuration
284    pub api: Option<Api>,
285    /// Event Bus configuration
286    pub event_bus: Option<EventBusConfig>,
287    /// Logger configuration
288    pub logger: Option<LoggerConfig>,
289    /// Policy configuration (override default error strategies)
290    pub policy: Option<PolicyConfig>,
291}
292impl Config {
293    /// Loads configuration from a TOML file
294    pub fn load(path: &str) -> Result<Self, String> {
295        // Read `config.toml` file content.
296        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}