1use std::{marker::PhantomData, sync::Arc};
2
3use async_trait::async_trait;
4use serde::{Deserialize, Serialize};
5use uuid::Uuid;
6
7use crate::{
8 Result, Store, Transaction,
9 backend::RecordMutations,
10 crdt::{
11 Doc,
12 doc::{Value, path::normalize_path},
13 },
14 store::{
15 ProjectionDescriptor, RecordProjection, Registered, StoreStateModel, errors::StoreError,
16 },
17};
18
19const DEFAULT_SCAN_PAGE_SIZE: usize = 128;
20
21struct TableProjection;
22
23fn project_doc_delta(delta: &Doc, prefix: &str, out: &mut RecordMutations) {
24 for (key, value) in delta.iter_all() {
25 let key = if prefix.is_empty() {
26 key.clone()
27 } else {
28 format!("{prefix}.{key}")
29 };
30 match value {
31 Value::Doc(doc) => project_doc_delta(doc, &key, out),
32 Value::Text(value) => {
33 insert_projected_record(out, key.into_bytes(), Some(value.as_bytes().to_vec()))
34 }
35 Value::Deleted => insert_projected_record(out, key.into_bytes(), None),
36 _ => {}
37 }
38 }
39}
40
41fn insert_projected_record(out: &mut RecordMutations, key: Vec<u8>, value: Option<Vec<u8>>) {
42 out.retain(|existing_key, _| !TableProjection.staged_keys_conflict(existing_key, &key));
43 out.insert(key, value);
44}
45
46pub(crate) fn encode_entry_delta(mutations: &RecordMutations) -> Result<Doc> {
47 TableProjection.encode_entry_delta(mutations)
48}
49
50impl RecordProjection<Doc> for TableProjection {
51 fn descriptor(&self) -> ProjectionDescriptor {
52 ProjectionDescriptor {
53 name: "eidetica/table/rows".to_string(),
54 version: 0,
55 }
56 }
57
58 fn project_delta(&self, delta: &Doc, out: &mut RecordMutations) -> Result<()> {
59 project_doc_delta(delta, "", out);
60 Ok(())
61 }
62
63 fn encode_entry_delta(&self, mutations: &RecordMutations) -> Result<Doc> {
64 let mut delta = Doc::new();
65 for (key, value) in mutations {
66 let key =
67 std::str::from_utf8(key).map_err(|error| StoreError::SerializationFailed {
68 store: "Table".to_string(),
69 reason: error.to_string(),
70 })?;
71 let _ = match value {
72 Some(value) => delta.set(
73 key,
74 std::str::from_utf8(value).map_err(|error| {
75 StoreError::SerializationFailed {
76 store: "Table".to_string(),
77 reason: error.to_string(),
78 }
79 })?,
80 ),
81 None => delta.remove(key),
82 };
83 }
84 Ok(delta)
85 }
86
87 fn normalize_record_key(&self, key: &[u8]) -> Result<Option<Vec<u8>>> {
88 let key = std::str::from_utf8(key).map_err(|error| StoreError::SerializationFailed {
89 store: "Table".to_string(),
90 reason: error.to_string(),
91 })?;
92 let normalized = normalize_path(key);
93 Ok((key.is_empty() || !normalized.is_empty()).then_some(normalized.into_bytes()))
94 }
95
96 fn staged_keys_conflict(&self, left: &[u8], right: &[u8]) -> bool {
97 left == right
98 || left
99 .strip_prefix(right)
100 .is_some_and(|suffix| suffix.starts_with(b"."))
101 || right
102 .strip_prefix(left)
103 .is_some_and(|suffix| suffix.starts_with(b"."))
104 }
105
106 fn staged_key_descends_from(&self, staged_key: &[u8], key: &[u8]) -> bool {
107 staged_key
108 .strip_prefix(key)
109 .is_some_and(|suffix| suffix.starts_with(b"."))
110 }
111}
112
113#[derive(Debug, Clone, PartialEq, Eq)]
115pub struct TableCursor(Vec<u8>);
116
117#[derive(Debug, Clone, PartialEq, Eq)]
119pub struct TablePage<T> {
120 pub rows: Vec<(String, T)>,
121 pub next: Option<TableCursor>,
122}
123
124pub struct Table<T>
143where
144 T: Serialize + for<'de> Deserialize<'de> + Clone,
145{
146 name: String,
147 txn: Transaction,
148 phantom: PhantomData<T>,
149}
150
151impl<T> Registered for Table<T>
152where
153 T: Serialize + for<'de> Deserialize<'de> + Clone,
154{
155 fn type_id() -> &'static str {
156 "table:v0"
157 }
158}
159
160#[async_trait]
161impl<T> Store for Table<T>
162where
163 T: Serialize + for<'de> Deserialize<'de> + Clone + Send + Sync,
164{
165 type Data = Doc;
166
167 fn state_model() -> StoreStateModel<Self::Data> {
168 StoreStateModel::Records(Arc::new(TableProjection))
169 }
170
171 async fn load(txn: &Transaction, subtree_name: String) -> Result<Self> {
172 Ok(Self {
173 name: subtree_name,
174 txn: txn.clone(),
175 phantom: PhantomData,
176 })
177 }
178
179 fn name(&self) -> &str {
180 &self.name
181 }
182
183 fn transaction(&self) -> &Transaction {
184 &self.txn
185 }
186}
187
188impl<T> Table<T>
189where
190 T: Serialize + for<'de> Deserialize<'de> + Clone + Send + Sync,
191{
192 pub async fn get(&self, key: impl AsRef<str>) -> Result<T> {
209 let key = key.as_ref();
210
211 let projection = TableProjection;
212 match self
213 .txn
214 .record_get(&self.name, &projection, key.as_bytes())
215 .await?
216 {
217 Some(value) => serde_json::from_slice(&value).map_err(|e| {
218 StoreError::DeserializationFailed {
219 store: self.name.clone(),
220 reason: format!("Failed to deserialize record for key '{key}': {e}"),
221 }
222 .into()
223 }),
224 None => Err(StoreError::KeyNotFound {
225 store: self.name.clone(),
226 key: key.to_string(),
227 }
228 .into()),
229 }
230 }
231
232 pub async fn insert(&self, row: T) -> Result<String> {
248 let primary_key = Uuid::new_v4().to_string();
250
251 let serialized_row =
252 serde_json::to_vec(&row).map_err(|e| StoreError::SerializationFailed {
253 store: self.name.clone(),
254 reason: format!("Failed to serialize record: {e}"),
255 })?;
256 self.txn.stage_record(
257 &self.name,
258 &TableProjection,
259 primary_key.as_bytes().to_vec(),
260 Some(serialized_row),
261 )?;
262
263 Ok(primary_key)
265 }
266
267 pub async fn set(&self, key: impl AsRef<str>, row: T) -> Result<()> {
282 let key_str = key.as_ref();
283 let serialized_row =
284 serde_json::to_vec(&row).map_err(|e| StoreError::SerializationFailed {
285 store: self.name.clone(),
286 reason: format!("Failed to serialize record for key '{key_str}': {e}"),
287 })?;
288 self.txn.stage_record(
289 &self.name,
290 &TableProjection,
291 key_str.as_bytes().to_vec(),
292 Some(serialized_row),
293 )
294 }
295
296 pub async fn delete(&self, key: impl AsRef<str>) -> Result<bool> {
311 let key_str = key.as_ref();
312
313 let exists = self.get(key_str).await.is_ok()
315 || self.txn.record_has_staged_descendant(
316 &self.name,
317 &TableProjection,
318 key_str.as_bytes(),
319 )?;
320
321 if !exists {
323 return Ok(false);
324 }
325
326 self.txn.stage_record(
327 &self.name,
328 &TableProjection,
329 key_str.as_bytes().to_vec(),
330 None,
331 )?;
332
333 Ok(true)
335 }
336
337 pub async fn search(&self, query: impl Fn(&T) -> bool) -> Result<Vec<(String, T)>> {
348 let mut result = Vec::new();
349 let mut cursor = None;
350 loop {
351 let page = self
352 .scan_page(cursor.as_ref(), DEFAULT_SCAN_PAGE_SIZE)
353 .await?;
354 for (key, row) in page.rows {
355 if query(&row) {
356 result.push((key, row));
357 }
358 }
359 cursor = page.next;
360 if cursor.is_none() {
361 break;
362 }
363 }
364 Ok(result)
365 }
366
367 pub async fn scan_page(
373 &self,
374 cursor: Option<&TableCursor>,
375 limit: usize,
376 ) -> Result<TablePage<T>> {
377 let projection = TableProjection;
378 let page = self
379 .txn
380 .record_scan(
381 &self.name,
382 &projection,
383 cursor.map(|cursor| cursor.0.as_slice()),
384 limit,
385 )
386 .await?;
387 let mut rows = Vec::with_capacity(page.records.len());
388 for (key, value) in page.records {
389 let key =
390 String::from_utf8(key).map_err(|error| StoreError::DeserializationFailed {
391 store: self.name.clone(),
392 reason: error.to_string(),
393 })?;
394 let row = serde_json::from_slice(&value).map_err(|error| {
395 StoreError::DeserializationFailed {
396 store: self.name.clone(),
397 reason: format!("Failed to deserialize record for key '{key}': {error}"),
398 }
399 })?;
400 rows.push((key, row));
401 }
402 Ok(TablePage {
403 rows,
404 next: page.next.map(TableCursor),
405 })
406 }
407}