Skip to main content

vantage_log_writer/
table_source.rs

1use async_trait::async_trait;
2use indexmap::IndexMap;
3use serde_json::Value;
4use vantage_core::{Result, error};
5use vantage_expressions::Expression;
6use vantage_expressions::traits::associated_expressions::AssociatedExpression;
7use vantage_expressions::traits::datasource::{DataSource, ExprDataSource};
8use vantage_expressions::traits::expressive::{DeferredFn, ExpressiveEnum};
9use vantage_table::column::core::{Column, ColumnType};
10use vantage_table::table::Table;
11use vantage_table::traits::table_source::TableSource;
12use vantage_types::{Entity, Record};
13
14use crate::log_writer::LogWriter;
15use crate::type_system::AnyJsonType;
16use crate::writer_task::WriteOp;
17
18impl DataSource for LogWriter {}
19
20impl ExprDataSource<Value> for LogWriter {
21    async fn execute(&self, _expr: &Expression<Value>) -> Result<Value> {
22        Err(unsupported("execute"))
23    }
24
25    fn defer(&self, _expr: Expression<Value>) -> DeferredFn<Value>
26    where
27        Value: Clone + Send + Sync + 'static,
28    {
29        DeferredFn::new(move || Box::pin(async move { Err(unsupported("defer")) }))
30    }
31}
32
33fn unsupported(method: &'static str) -> vantage_core::VantageError {
34    error!("log-writer is insert-only", method = method)
35        .mark_unsupported()
36        .traced()
37}
38
39#[async_trait]
40impl TableSource for LogWriter {
41    type Column<Type>
42        = Column<Type>
43    where
44        Type: ColumnType;
45    type AnyType = AnyJsonType;
46    type Value = Value;
47    type Id = String;
48    type Condition = Expression<Self::Value>;
49    type Source = String;
50
51    fn create_column<Type: ColumnType>(&self, name: &str) -> Self::Column<Type> {
52        Column::new(name)
53    }
54
55    fn to_any_column<Type: ColumnType>(
56        &self,
57        column: Self::Column<Type>,
58    ) -> Self::Column<Self::AnyType> {
59        Column::from_column(column)
60    }
61
62    fn convert_any_column<Type: ColumnType>(
63        &self,
64        any_column: Self::Column<Self::AnyType>,
65    ) -> Option<Self::Column<Type>> {
66        Some(Column::from_column(any_column))
67    }
68
69    fn expr(
70        &self,
71        template: impl Into<String>,
72        parameters: Vec<ExpressiveEnum<Self::Value>>,
73    ) -> Expression<Self::Value> {
74        Expression::new(template, parameters)
75    }
76
77    fn search_table_condition<E>(
78        &self,
79        _table: &Table<Self, E>,
80        _search_value: &str,
81    ) -> Self::Condition
82    where
83        E: Entity<Self::Value>,
84    {
85        Expression::new("", vec![])
86    }
87
88    async fn list_table_values<E>(
89        &self,
90        _table: &Table<Self, E>,
91    ) -> Result<IndexMap<Self::Id, Record<Self::Value>>>
92    where
93        E: Entity<Self::Value>,
94        Self: Sized,
95    {
96        Err(unsupported("list_table_values"))
97    }
98
99    async fn get_table_value<E>(
100        &self,
101        _table: &Table<Self, E>,
102        _id: &Self::Id,
103    ) -> Result<Option<Record<Self::Value>>>
104    where
105        E: Entity<Self::Value>,
106        Self: Sized,
107    {
108        Err(unsupported("get_table_value"))
109    }
110
111    async fn get_table_some_value<E>(
112        &self,
113        _table: &Table<Self, E>,
114    ) -> Result<Option<(Self::Id, Record<Self::Value>)>>
115    where
116        E: Entity<Self::Value>,
117        Self: Sized,
118    {
119        Err(unsupported("get_table_some_value"))
120    }
121
122    async fn get_table_count<E>(&self, _table: &Table<Self, E>) -> Result<i64>
123    where
124        E: Entity<Self::Value>,
125        Self: Sized,
126    {
127        Err(unsupported("get_table_count"))
128    }
129
130    async fn get_table_sum<E>(
131        &self,
132        _table: &Table<Self, E>,
133        _column: &Self::Column<Self::AnyType>,
134    ) -> Result<Self::Value>
135    where
136        E: Entity<Self::Value>,
137        Self: Sized,
138    {
139        Err(unsupported("get_table_sum"))
140    }
141
142    async fn get_table_max<E>(
143        &self,
144        _table: &Table<Self, E>,
145        _column: &Self::Column<Self::AnyType>,
146    ) -> Result<Self::Value>
147    where
148        E: Entity<Self::Value>,
149        Self: Sized,
150    {
151        Err(unsupported("get_table_max"))
152    }
153
154    async fn get_table_min<E>(
155        &self,
156        _table: &Table<Self, E>,
157        _column: &Self::Column<Self::AnyType>,
158    ) -> Result<Self::Value>
159    where
160        E: Entity<Self::Value>,
161        Self: Sized,
162    {
163        Err(unsupported("get_table_min"))
164    }
165
166    async fn insert_table_value<E>(
167        &self,
168        table: &Table<Self, E>,
169        id: &Self::Id,
170        record: &Record<Self::Value>,
171    ) -> Result<Record<Self::Value>>
172    where
173        E: Entity<Self::Value>,
174        Self: Sized,
175    {
176        let projected = project_record(table.columns().keys(), record, self.id_column(), id);
177        let line = serialize_line(&projected)?;
178        let path = self.file_path(table.table_name());
179        self.sender()
180            .send(WriteOp::Append { path, line })
181            .await
182            .map_err(|e| error!("log writer channel closed", detail = e.to_string()))?;
183        Ok(projected)
184    }
185
186    async fn replace_table_value<E>(
187        &self,
188        _table: &Table<Self, E>,
189        _id: &Self::Id,
190        _record: &Record<Self::Value>,
191    ) -> Result<Record<Self::Value>>
192    where
193        E: Entity<Self::Value>,
194        Self: Sized,
195    {
196        Err(unsupported("replace_table_value"))
197    }
198
199    async fn patch_table_value<E>(
200        &self,
201        _table: &Table<Self, E>,
202        _id: &Self::Id,
203        _partial: &Record<Self::Value>,
204    ) -> Result<Record<Self::Value>>
205    where
206        E: Entity<Self::Value>,
207        Self: Sized,
208    {
209        Err(unsupported("patch_table_value"))
210    }
211
212    async fn delete_table_value<E>(&self, _table: &Table<Self, E>, _id: &Self::Id) -> Result<()>
213    where
214        E: Entity<Self::Value>,
215        Self: Sized,
216    {
217        Err(unsupported("delete_table_value"))
218    }
219
220    async fn delete_table_all_values<E>(&self, _table: &Table<Self, E>) -> Result<()>
221    where
222        E: Entity<Self::Value>,
223        Self: Sized,
224    {
225        Err(unsupported("delete_table_all_values"))
226    }
227
228    async fn insert_table_return_id_value<E>(
229        &self,
230        table: &Table<Self, E>,
231        record: &Record<Self::Value>,
232    ) -> Result<Self::Id>
233    where
234        E: Entity<Self::Value>,
235        Self: Sized,
236    {
237        let id = extract_or_generate_id(record, self.id_column());
238        let projected = project_record(table.columns().keys(), record, self.id_column(), &id);
239        let line = serialize_line(&projected)?;
240        let path = self.file_path(table.table_name());
241        self.sender()
242            .send(WriteOp::Append { path, line })
243            .await
244            .map_err(|e| error!("log writer channel closed", detail = e.to_string()))?;
245        Ok(id)
246    }
247
248    fn related_in_condition<SourceE: Entity<Self::Value> + 'static>(
249        &self,
250        _target_field: &str,
251        _source_table: &Table<Self, SourceE>,
252        _source_column: &str,
253    ) -> Self::Condition
254    where
255        Self: Sized,
256    {
257        Expression::new("", vec![])
258    }
259
260    fn column_table_values_expr<'a, E, Type: ColumnType>(
261        &'a self,
262        _table: &Table<Self, E>,
263        _column: &Self::Column<Type>,
264    ) -> AssociatedExpression<'a, Self, Self::Value, Vec<Type>>
265    where
266        E: Entity<Self::Value> + 'static,
267        Self: Sized,
268    {
269        unimplemented!("log-writer is insert-only; column_table_values_expr is unreachable")
270    }
271}
272
273/// Project a record onto the table's declared column set, then attach the id.
274///
275/// Entity fields not declared as columns are dropped — this is the contract
276/// the user pinned down: "entity values with non-existent columns would be
277/// dropped".
278fn project_record<'a, I>(
279    column_names: I,
280    record: &Record<Value>,
281    id_column: &str,
282    id: &str,
283) -> Record<Value>
284where
285    I: IntoIterator<Item = &'a String>,
286{
287    let mut out = Record::new();
288    let mut wrote_id = false;
289    for col in column_names {
290        if col == id_column {
291            out.insert(col.clone(), Value::String(id.to_string()));
292            wrote_id = true;
293        } else if let Some(v) = record.get(col) {
294            out.insert(col.clone(), v.clone());
295        }
296    }
297    if !wrote_id {
298        out.insert(id_column.to_string(), Value::String(id.to_string()));
299    }
300    out
301}
302
303fn serialize_line(record: &Record<Value>) -> Result<String> {
304    let map: serde_json::Map<String, Value> = record
305        .as_inner()
306        .iter()
307        .map(|(k, v)| (k.clone(), v.clone()))
308        .collect();
309    let mut s = serde_json::to_string(&Value::Object(map))
310        .map_err(|e| error!("failed to serialize record to JSON", detail = e.to_string()))?;
311    s.push('\n');
312    Ok(s)
313}
314
315fn extract_or_generate_id(record: &Record<Value>, id_column: &str) -> String {
316    if let Some(v) = record.get(id_column) {
317        match v {
318            Value::String(s) if !s.is_empty() => return s.clone(),
319            Value::Number(n) => return n.to_string(),
320            _ => {}
321        }
322    }
323    ulid::Ulid::new().to_string()
324}