laminar_connectors/files/
mod.rs1use std::sync::Arc;
4
5use crate::config::ConnectorInfo;
6use crate::registry::ConnectorRegistry;
7
8pub mod arrow_ipc_codec;
9pub mod config;
10pub mod discovery;
11pub mod manifest;
12pub mod sink;
13pub mod source;
14pub mod text_decoder;
15
16pub use config::{FileFormat, FileSinkConfig, FileSourceConfig};
17pub use manifest::FileIngestionManifest;
18pub use sink::FileSink;
19pub use source::FileSource;
20pub use text_decoder::TextLineDecoder;
21
22pub fn register_file_source(
32 registry: &ConnectorRegistry,
33) -> Result<(), crate::error::ConnectorError> {
34 use crate::config::ConfigKeySpec;
35 let info = ConnectorInfo {
36 name: "files".to_string(),
37 display_name: "File Source (AutoLoader)".to_string(),
38 version: env!("CARGO_PKG_VERSION").to_string(),
39 is_source: true,
40 is_sink: false,
41 config_keys: vec![
42 ConfigKeySpec::required("path", "Directory path, glob pattern, or cloud storage URL"),
43 ConfigKeySpec::optional(
44 "format",
45 "Data format (csv, tsv, json, jsonl, text, txt, parquet, arrow)",
46 "auto-detect",
47 ),
48 ConfigKeySpec::optional(
49 "glob_pattern",
50 "Optional glob pattern to filter files by name",
51 "*",
52 ),
53 ],
54 };
55 registry.register_source(
56 "files",
57 info,
58 Arc::new(|registry: Option<&Arc<prometheus::Registry>>| {
59 Ok(Box::new(FileSource::with_registry(
60 registry.map(Arc::as_ref),
61 )))
62 }),
63 )
64}
65
66pub fn register_file_sink(
74 registry: &ConnectorRegistry,
75) -> Result<(), crate::error::ConnectorError> {
76 use crate::config::ConfigKeySpec;
77 let info = ConnectorInfo {
78 name: "files".to_string(),
79 display_name: "File Sink".to_string(),
80 version: env!("CARGO_PKG_VERSION").to_string(),
81 is_source: false,
82 is_sink: true,
83 config_keys: vec![
84 ConfigKeySpec::required("path", "Output directory path"),
85 ConfigKeySpec::required("format", "Output format (csv, json, text, parquet, arrow)"),
86 ConfigKeySpec::optional("prefix", "Immutable output file name prefix", "part"),
87 ],
88 };
89 registry.register_sink(
90 "files",
91 info,
92 Arc::new(|_config, registry: Option<&Arc<prometheus::Registry>>| {
93 Ok(Box::new(FileSink::with_registry(registry.map(Arc::as_ref))))
94 }),
95 )
96}
97
98#[cfg(test)]
99mod tests {
100 use super::*;
101
102 #[test]
103 fn test_register_file_source() {
104 let registry = ConnectorRegistry::new();
105 register_file_source(®istry).unwrap();
106
107 let sources = registry.list_sources();
108 assert!(sources.contains(&"files".to_string()));
109
110 let info = registry.source_info("files").unwrap();
111 assert_eq!(info.name, "files");
112 assert!(info.is_source);
113 assert!(!info.is_sink);
114 }
115
116 #[test]
117 fn test_register_file_sink() {
118 let registry = ConnectorRegistry::new();
119 register_file_sink(®istry).unwrap();
120
121 let sinks = registry.list_sinks();
122 assert!(sinks.contains(&"files".to_string()));
123
124 let info = registry.sink_info("files").unwrap();
125 assert_eq!(info.name, "files");
126 assert!(!info.is_source);
127 assert!(info.is_sink);
128 }
129
130 #[test]
131 fn test_create_source_from_registry() {
132 let registry = ConnectorRegistry::new();
133 register_file_source(®istry).unwrap();
134
135 let config = crate::config::ConnectorConfig::new("files");
136 let source = registry.create_source(&config, None);
137 assert!(source.is_ok());
138 }
139
140 #[test]
141 fn test_create_sink_from_registry() {
142 let registry = ConnectorRegistry::new();
143 register_file_sink(®istry).unwrap();
144
145 let config = crate::config::ConnectorConfig::new("files");
146 let sink = registry.create_sink(&config, None);
147 assert!(sink.is_ok());
148 }
149}