1use crate::operation::CsvOperation;
2use async_trait::async_trait;
3use indexmap::IndexMap;
4use vantage_core::error;
5use vantage_dataset::traits::Result;
6use vantage_expressions::Expression;
7use vantage_expressions::Expressive;
8use vantage_expressions::traits::associated_expressions::AssociatedExpression;
9use vantage_expressions::traits::datasource::DataSource;
10use vantage_expressions::traits::expressive::{DeferredFn, ExpressiveEnum};
11use vantage_table::column::core::{Column, ColumnType};
12use vantage_table::table::Table;
13use vantage_table::traits::table_source::TableSource;
14use vantage_types::{Entity, Record};
15
16use crate::Csv;
17use crate::condition::apply_condition;
18use crate::type_system::AnyCsvType;
19
20impl DataSource for Csv {}
21
22#[async_trait]
23impl TableSource for Csv {
24 type Column<Type>
25 = Column<Type>
26 where
27 Type: ColumnType;
28 type AnyType = AnyCsvType;
29 type Value = AnyCsvType;
30 type Id = String;
31 type Condition = vantage_expressions::Expression<Self::Value>;
32 type Source = String;
33
34 fn eq_value_condition(&self, field: &str, value: Self::Value) -> Result<Self::Condition> {
35 let column: Column<AnyCsvType> = Column::new(field);
36 Ok(CsvOperation::eq(&column, value))
37 }
38
39 fn create_column<Type: ColumnType>(&self, name: &str) -> Self::Column<Type> {
40 Column::new(name)
41 }
42
43 fn to_any_column<Type: ColumnType>(
44 &self,
45 column: Self::Column<Type>,
46 ) -> Self::Column<Self::AnyType> {
47 Column::from_column(column)
48 }
49
50 fn convert_any_column<Type: ColumnType>(
51 &self,
52 any_column: Self::Column<Self::AnyType>,
53 ) -> Option<Self::Column<Type>> {
54 Some(Column::from_column(any_column))
55 }
56
57 fn expr(
58 &self,
59 template: impl Into<String>,
60 parameters: Vec<ExpressiveEnum<Self::Value>>,
61 ) -> Expression<Self::Value> {
62 Expression::new(template, parameters)
63 }
64
65 fn search_table_condition<E>(
66 &self,
67 _table: &Table<Self, E>,
68 _search_value: &str,
69 ) -> Expression<Self::Value>
70 where
71 E: Entity<Self::Value>,
72 {
73 Expression::new(crate::operation::OP_SEARCH, vec![])
76 }
77
78 async fn list_table_values<E>(
79 &self,
80 table: &Table<Self, E>,
81 ) -> Result<IndexMap<Self::Id, Record<Self::Value>>>
82 where
83 E: Entity<Self::Value>,
84 Self: Sized,
85 {
86 let mut records = self.read_csv(table.table_name(), table.columns())?;
87
88 for condition in table.conditions() {
89 records = apply_condition(records, condition).await?;
90 }
91
92 Ok(records)
93 }
94
95 async fn get_table_value<E>(
96 &self,
97 table: &Table<Self, E>,
98 id: &Self::Id,
99 ) -> Result<Option<Record<Self::Value>>>
100 where
101 E: Entity<Self::Value>,
102 Self: Sized,
103 {
104 let records = self.read_csv(table.table_name(), table.columns())?;
105 Ok(records.get(id).cloned())
106 }
107
108 async fn get_table_some_value<E>(
109 &self,
110 table: &Table<Self, E>,
111 ) -> Result<Option<(Self::Id, Record<Self::Value>)>>
112 where
113 E: Entity<Self::Value>,
114 Self: Sized,
115 {
116 let records = self.read_csv(table.table_name(), table.columns())?;
117 Ok(records.into_iter().next())
118 }
119
120 async fn get_table_count<E>(&self, table: &Table<Self, E>) -> Result<i64>
121 where
122 E: Entity<Self::Value>,
123 Self: Sized,
124 {
125 let records = self.read_csv(table.table_name(), table.columns())?;
126 Ok(records.len() as i64)
127 }
128
129 async fn get_table_sum<E>(
130 &self,
131 _table: &Table<Self, E>,
132 _column: &Self::Column<Self::AnyType>,
133 ) -> Result<Self::Value>
134 where
135 E: Entity<Self::Value>,
136 Self: Sized,
137 {
138 Err(error!("Sum not implemented for CSV backend"))
139 }
140
141 async fn get_table_max<E>(
142 &self,
143 _table: &Table<Self, E>,
144 _column: &Self::Column<Self::AnyType>,
145 ) -> Result<Self::Value>
146 where
147 E: Entity<Self::Value>,
148 Self: Sized,
149 {
150 Err(error!("Max not implemented for CSV backend"))
151 }
152
153 async fn get_table_min<E>(
154 &self,
155 _table: &Table<Self, E>,
156 _column: &Self::Column<Self::AnyType>,
157 ) -> Result<Self::Value>
158 where
159 E: Entity<Self::Value>,
160 Self: Sized,
161 {
162 Err(error!("Min not implemented for CSV backend"))
163 }
164
165 async fn insert_table_value<E>(
166 &self,
167 _table: &Table<Self, E>,
168 _id: &Self::Id,
169 _record: &Record<Self::Value>,
170 ) -> Result<Record<Self::Value>>
171 where
172 E: Entity<Self::Value>,
173 Self: Sized,
174 {
175 Err(error!("CSV is a read-only data source"))
176 }
177
178 async fn replace_table_value<E>(
179 &self,
180 _table: &Table<Self, E>,
181 _id: &Self::Id,
182 _record: &Record<Self::Value>,
183 ) -> Result<Record<Self::Value>>
184 where
185 E: Entity<Self::Value>,
186 Self: Sized,
187 {
188 Err(error!("CSV is a read-only data source"))
189 }
190
191 async fn patch_table_value<E>(
192 &self,
193 _table: &Table<Self, E>,
194 _id: &Self::Id,
195 _partial: &Record<Self::Value>,
196 ) -> Result<Record<Self::Value>>
197 where
198 E: Entity<Self::Value>,
199 Self: Sized,
200 {
201 Err(error!("CSV is a read-only data source"))
202 }
203
204 async fn delete_table_value<E>(&self, _table: &Table<Self, E>, _id: &Self::Id) -> Result<()>
205 where
206 E: Entity<Self::Value>,
207 Self: Sized,
208 {
209 Err(error!("CSV is a read-only data source"))
210 }
211
212 async fn delete_table_all_values<E>(&self, _table: &Table<Self, E>) -> Result<()>
213 where
214 E: Entity<Self::Value>,
215 Self: Sized,
216 {
217 Err(error!("CSV is a read-only data source"))
218 }
219
220 async fn insert_table_return_id_value<E>(
221 &self,
222 _table: &Table<Self, E>,
223 _record: &Record<Self::Value>,
224 ) -> Result<Self::Id>
225 where
226 E: Entity<Self::Value>,
227 Self: Sized,
228 {
229 Err(error!("CSV is a read-only data source"))
230 }
231
232 fn related_in_condition<SourceE: Entity<Self::Value> + 'static>(
233 &self,
234 target_field: &str,
235 source_table: &Table<Self, SourceE>,
236 source_column: &str,
237 ) -> Self::Condition
238 where
239 Self: Sized,
240 {
241 let src_col = self.create_column::<Self::AnyType>(source_column);
242 let fk_values = self.column_table_values_expr(source_table, &src_col);
243 let tgt_col = self.create_column::<Self::AnyType>(target_field);
244 tgt_col.in_(fk_values.expr())
245 }
246
247 fn column_table_values_expr<'a, E, Type: ColumnType>(
248 &'a self,
249 table: &Table<Self, E>,
250 column: &Self::Column<Type>,
251 ) -> AssociatedExpression<'a, Self, Self::Value, Vec<Type>>
252 where
253 E: Entity<Self::Value> + 'static,
254 Self: Sized,
255 {
256 use vantage_expressions::{
257 expr_any,
258 traits::{associated_expressions::AssociatedExpression, datasource::ExprDataSource},
259 };
260
261 let table_clone = table.clone();
262 let col = column.name().to_string();
263 let csv = self.clone();
264
265 let inner = expr_any!("{}", {
266 DeferredFn::new(move || {
267 let csv = csv.clone();
268 let table = table_clone.clone();
269 let col = col.clone();
270 Box::pin(async move {
271 let records = csv.list_table_values(&table).await?;
272 let values: Vec<AnyCsvType> = records
273 .values()
274 .filter_map(|r| r.get(&col).cloned())
275 .collect();
276 Ok(ExpressiveEnum::Scalar(AnyCsvType::new(values)))
277 })
278 })
279 });
280
281 let expr = expr_any!("{}", { self.defer(inner) });
282 AssociatedExpression::new(expr, self)
283 }
284}
285
286#[cfg(test)]
287mod tests {
288 use super::*;
289 use crate::type_system::CsvTypeVariants;
290 use vantage_dataset::prelude::{ReadableValueSet, WritableValueSet};
291 use vantage_types::EmptyEntity;
292
293 fn test_csv() -> Csv {
294 Csv::new(format!("{}/data", env!("CARGO_MANIFEST_DIR")))
295 }
296
297 #[tokio::test]
298 async fn test_list_bakery() {
299 let csv = test_csv();
300 let table = Table::<Csv, EmptyEntity>::new("bakery", csv)
301 .with_column_of::<String>("name")
302 .with_column_of::<i64>("profit_margin");
303
304 let values = table.list_values().await.unwrap();
305 assert_eq!(values.len(), 1);
306 assert!(values.contains_key("hill_valley"));
307
308 let bakery = &values["hill_valley"];
309 let name = bakery["name"].try_get::<String>().unwrap();
310 assert_eq!(name, "Hill Valley Bakery");
311
312 let profit = bakery["profit_margin"].try_get::<i64>().unwrap();
313 assert_eq!(profit, 15);
314 }
315
316 #[tokio::test]
317 async fn test_list_clients() {
318 let csv = test_csv();
319 let table = Table::<Csv, EmptyEntity>::new("client", csv)
320 .with_column_of::<String>("name")
321 .with_column_of::<String>("email")
322 .with_column_of::<bool>("is_paying_client")
323 .with_column_of::<serde_json::Value>("metadata");
324
325 let values = table.list_values().await.unwrap();
326 assert_eq!(values.len(), 3);
327
328 let marty = &values["marty"];
329 assert_eq!(marty["name"].try_get::<String>().unwrap(), "Marty McFly");
330 assert!(marty["is_paying_client"].try_get::<bool>().unwrap());
331
332 let biff = &values["biff"];
333 assert!(!biff["is_paying_client"].try_get::<bool>().unwrap());
334 assert_eq!(biff["metadata"].type_variant(), Some(CsvTypeVariants::Json));
335 }
336
337 #[tokio::test]
338 async fn test_list_products_typed() {
339 let csv = test_csv();
340 let table = Table::<Csv, EmptyEntity>::new("product", csv)
341 .with_column_of::<String>("name")
342 .with_column_of::<i64>("calories")
343 .with_column_of::<i64>("price")
344 .with_column_of::<bool>("is_deleted")
345 .with_column_of::<serde_json::Value>("inventory");
346
347 let values = table.list_values().await.unwrap();
348 assert_eq!(values.len(), 5);
349
350 let cupcake = &values["flux_cupcake"];
351 assert_eq!(
352 cupcake["name"].try_get::<String>().unwrap(),
353 "Flux Capacitor Cupcake"
354 );
355 assert_eq!(cupcake["calories"].try_get::<i64>().unwrap(), 300);
356 assert_eq!(cupcake["price"].try_get::<i64>().unwrap(), 120);
357 assert!(!cupcake["is_deleted"].try_get::<bool>().unwrap());
358
359 let inv = cupcake["inventory"].try_get::<serde_json::Value>().unwrap();
360 assert_eq!(inv["stock"], serde_json::json!(50));
361 }
362
363 #[tokio::test]
364 async fn test_untyped_columns_stay_string() {
365 let csv = test_csv();
366 let table = Table::<Csv, EmptyEntity>::new("product", csv);
367
368 let values = table.list_values().await.unwrap();
369 let cupcake = &values["flux_cupcake"];
370 assert_eq!(
371 cupcake["calories"].type_variant(),
372 Some(CsvTypeVariants::String)
373 );
374 assert_eq!(cupcake["calories"].try_get::<String>().unwrap(), "300");
375 }
376
377 #[tokio::test]
378 async fn test_get_value_by_id() {
379 let csv = test_csv();
380 let table = Table::<Csv, EmptyEntity>::new("client", csv)
381 .with_column_of::<String>("name")
382 .with_column_of::<String>("email");
383
384 let record = table.get_value("doc").await.unwrap().expect("doc exists");
385 assert_eq!(record["name"].try_get::<String>().unwrap(), "Doc Brown");
386 assert_eq!(
387 record["email"].try_get::<String>().unwrap(),
388 "doc@brown.com"
389 );
390 }
391
392 #[tokio::test]
393 async fn test_get_value_not_found() {
394 let csv = test_csv();
395 let table = Table::<Csv, EmptyEntity>::new("client", csv);
396
397 let result = table.get_value("nonexistent").await.unwrap();
398 assert!(result.is_none());
399 }
400
401 #[tokio::test]
402 async fn test_get_some_value() {
403 let csv = test_csv();
404 let table = Table::<Csv, EmptyEntity>::new("bakery", csv).with_column_of::<String>("name");
405
406 let result = table.get_some_value().await.unwrap();
407 assert!(result.is_some());
408 let (id, record) = result.unwrap();
409 assert_eq!(id, "hill_valley");
410 assert_eq!(
411 record["name"].try_get::<String>().unwrap(),
412 "Hill Valley Bakery"
413 );
414 }
415
416 #[tokio::test]
417 async fn test_get_count() {
418 let csv = test_csv();
419 let table = Table::<Csv, EmptyEntity>::new("product", csv);
420
421 let count = table.data_source().get_table_count(&table).await.unwrap();
422 assert_eq!(count, 5);
423 }
424
425 #[tokio::test]
426 async fn test_write_operations_fail() {
427 let csv = test_csv();
428 let table = Table::<Csv, EmptyEntity>::new("bakery", csv);
429
430 let record = Record::new();
431 assert!(
432 WritableValueSet::insert_value(&table, "test", &record)
433 .await
434 .is_err()
435 );
436 assert!(WritableValueSet::delete(&table, "test").await.is_err());
437 assert!(WritableValueSet::delete_all(&table).await.is_err());
438 }
439
440 #[tokio::test]
441 async fn test_missing_file() {
442 let csv = test_csv();
443 let table = Table::<Csv, EmptyEntity>::new("nonexistent", csv);
444
445 let result = table.list_values().await;
446 assert!(result.is_err());
447 }
448
449 #[tokio::test]
454 async fn search_is_unsupported_not_silent_match_all() {
455 let csv = test_csv();
456 let mut table =
457 Table::<Csv, EmptyEntity>::new("client", csv).with_column_of::<String>("name");
458
459 table.add_search("marty");
460
461 let err = table
462 .list_values()
463 .await
464 .expect_err("CSV search must error, not silently return all rows");
465 assert!(
466 err.is_unsupported(),
467 "CSV search error must be classified Unsupported, got: {err:?}"
468 );
469 }
470}