akar_common/
extension_utils.rs1use crate::data_chunk::DataChunk;
9use crate::types::PhysicalTypeID;
10use arrow::array::{ArrayRef, StringArray};
11use std::collections::VecDeque;
12use std::path::PathBuf;
13use std::sync::{Arc, Mutex};
14
15pub fn fill_chunk_with_strings(chunk: &mut DataChunk, field_name: &str, values: &[String]) {
21 let array: ArrayRef = Arc::new(StringArray::from_iter_values(values.iter().map(String::as_str)));
22 chunk.fields.clear();
23 chunk.field_types.clear();
24 chunk.field_names.clear();
25 chunk.fields.push(array);
26 chunk.field_types.push(PhysicalTypeID::String);
27 chunk.field_names.push(field_name.to_string());
28 chunk.size = values.len();
29}
30
31pub fn quote_sql_identifier(name: &str) -> String {
33 format!("\"{}\"", name.replace('"', "\"\""))
34}
35
36pub fn quote_sql_table_name(name: &str) -> String {
42 name.split('.').map(quote_sql_identifier).collect::<Vec<_>>().join(".")
43}
44
45const MAX_RETAINED_TEMP_FILES: usize = 64;
47
48static RETAINED_TEMP_FILES: Mutex<VecDeque<PathBuf>> = Mutex::new(VecDeque::new());
51
52pub fn retain_temp_file(path: PathBuf) -> String {
59 let mut queue = RETAINED_TEMP_FILES.lock().unwrap_or_else(|p| p.into_inner());
60 queue.push_back(path.clone());
61 while queue.len() > MAX_RETAINED_TEMP_FILES {
62 if let Some(oldest) = queue.pop_front() {
63 let _ = std::fs::remove_file(&oldest);
64 }
65 }
66 path.to_string_lossy().to_string()
67}
68
69pub fn clear_retained_temp_files() {
71 let mut queue = RETAINED_TEMP_FILES.lock().unwrap_or_else(|p| p.into_inner());
72 for path in queue.drain(..) {
73 let _ = std::fs::remove_file(&path);
74 }
75}
76
77#[cfg(test)]
78mod tests {
79 use super::*;
80 use crate::vector::DataChunk;
81
82 #[test]
83 fn test_fill_chunk_with_strings() {
84 let mut chunk = DataChunk::new(Vec::new(), Vec::new());
85 fill_chunk_with_strings(&mut chunk, "path", &["a.parquet".into(), "b.parquet".into()]);
86 assert_eq!(chunk.size, 2);
87 assert_eq!(chunk.num_fields(), 1);
88 assert_eq!(chunk.field_names, vec!["path".to_string()]);
89 assert_eq!(chunk.field_types, vec![PhysicalTypeID::String]);
90 }
91
92 #[test]
93 fn test_fill_chunk_with_strings_replaces_schema() {
94 let mut chunk = DataChunk::new(Vec::new(), Vec::new());
95 fill_chunk_with_strings(&mut chunk, "first", &["a".into()]);
96 fill_chunk_with_strings(&mut chunk, "second", &["x".into(), "y".into(), "z".into()]);
97 assert_eq!(chunk.size, 3);
98 assert_eq!(chunk.num_fields(), 1);
99 assert_eq!(chunk.field_names, vec!["second".to_string()]);
100 }
101
102 #[test]
103 fn test_quote_sql_identifier() {
104 assert_eq!(quote_sql_identifier("tbl"), "\"tbl\"");
105 assert_eq!(quote_sql_identifier("we\"ird"), "\"we\"\"ird\"");
106 }
107
108 #[test]
109 fn test_quote_sql_table_name() {
110 assert_eq!(
111 quote_sql_table_name("main.default.people"),
112 "\"main\".\"default\".\"people\""
113 );
114 assert_eq!(
115 quote_sql_table_name("x.\"y\"; DROP TABLE t"),
116 "\"x\".\"\"\"y\"\"; DROP TABLE t\""
117 );
118 }
119
120 #[test]
121 fn test_retain_temp_file_evicts_oldest() {
122 clear_retained_temp_files();
123 let dir = std::env::temp_dir();
124 for i in 0..(MAX_RETAINED_TEMP_FILES + 5) {
125 let p = dir.join(format!("akar_retain_test_{i}.tmp"));
126 std::fs::write(&p, b"x").unwrap();
127 retain_temp_file(p);
128 }
129 for i in 0..5 {
131 let p = dir.join(format!("akar_retain_test_{i}.tmp"));
132 assert!(!p.exists(), "file {i} should have been evicted");
133 }
134 for i in 5..(MAX_RETAINED_TEMP_FILES + 5) {
136 let p = dir.join(format!("akar_retain_test_{i}.tmp"));
137 assert!(p.exists(), "file {i} should still be retained");
138 }
139 clear_retained_temp_files();
140 }
141}