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
273fn 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}