Skip to main content

vantage_csv/
table_source.rs

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        // CSV cannot search server-side; this sentinel makes `apply_condition`
74        // surface an Unsupported error instead of silently matching every row.
75        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    // CSV has no query engine, so it cannot perform full-table search. Asking it
450    // to must surface an `Unsupported` error — never the old silent match-all,
451    // which masqueraded a "no filter applied" as a successful search and leaked
452    // every row. In-memory search belongs to the Lens/Diorama layer.
453    #[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}