Skip to main content

dynoxide/storage_backend/
mod.rs

1//! Storage backend abstraction.
2//!
3//! Defines the [`StorageBackend`] trait that decouples Dynoxide's data layer
4//! from a specific SQLite binding. The native [`rusqlite`]-backed
5//! [`Storage`](crate::storage::Storage) implements the trait, and the
6//! `wasm-sqlite` build adds `wasm_backend::WasmBridgeBackend`, which runs the
7//! same SQL against a JS SQLite database over a wasm-bindgen bridge. Both
8//! backends issue identical SQL because they share the builders in
9//! [`sql_builders`].
10//!
11//! The native build consumes the trait monomorphically through `Storage`; the
12//! wasm build consumes it through `WasmBridgeBackend`. The escape hatches
13//! `Storage::conn()` and `Storage::conn_mut()` are not exposed by the trait
14//! and remain native-only.
15//!
16//! # No `Send + Sync` super-trait
17//!
18//! [`Storage`](crate::storage::Storage) carries a `RefCell<HashMap<...>>` for
19//! its metadata cache, so `Storage: !Sync`. A `Send + Sync` super-trait would
20//! refuse the impl on `Storage`. With no dynamic dispatch site in scope,
21//! auto-trait propagation across `.await` is decided per-callsite anyway, so
22//! adding `Send` to the super-trait would not earn any compile-time
23//! guarantee on the futures returned by trait methods.
24
25pub mod clock;
26pub mod error;
27// Internal: the shared SQL contract between the rusqlite and wasm backends.
28// `pub` only because both backend modules consume it across the cfg split; it
29// is not a stable API, hence `#[doc(hidden)]`.
30#[doc(hidden)]
31pub mod sql_builders;
32// Schema versions and the migrations a recorded version calls for on open.
33pub(crate) mod schema;
34// The native rusqlite-backed `Storage` exists whenever either SQLite backend
35// feature is on (the crate refuses to build with neither), and the handlers
36// now consume `StorageBackend` through it, so the impl must track the same
37// condition rather than `native-sqlite` alone.
38#[cfg(any(feature = "native-sqlite", feature = "_has-encryption"))]
39pub mod rusqlite_impl;
40#[cfg(feature = "wasm-sqlite")]
41pub mod wasm_backend;
42// Compile-only proof that the trait is satisfiable by a fresh impl. In-crate
43// rather than under `tests/`, because a sealed trait cannot be implemented
44// from outside.
45#[cfg(test)]
46mod trait_compile;
47
48use crate::storage::{
49    CreateTableMetadata, DatabaseInfo, QueryParams, ScanParams, StreamRecord, TableMetadata,
50    TableStats,
51};
52use crate::types::Tag;
53
54pub use clock::{Clock, ManualClock, SystemClock};
55pub use error::BackendError;
56#[cfg(any(feature = "native-sqlite", feature = "_has-encryption"))]
57pub use error::from_rusqlite;
58#[doc(hidden)]
59pub use sql_builders::SqlParam;
60#[cfg(feature = "wasm-sqlite")]
61pub use wasm_backend::WasmBridgeBackend;
62
63/// One base-table row for a bulk insert via [`StorageBackend::put_base_items`].
64///
65/// Unlike [`StorageBackend::put_item_with_hash`], which preserves any existing
66/// `cached_at` value, the bulk path writes `cached_at` verbatim: this mirrors
67/// the import flow, which sets the timestamp explicitly (or clears it) for
68/// every row it loads.
69#[derive(Debug, Clone)]
70pub struct BaseItemRow {
71    /// Partition key string.
72    pub pk: String,
73    /// Sort key string (empty for tables without a sort key).
74    pub sk: String,
75    /// Serialised item JSON.
76    pub item_json: String,
77    /// Item size in bytes.
78    pub item_size: usize,
79    /// Cache timestamp written verbatim; `None` clears the column.
80    pub cached_at: Option<f64>,
81    /// Hash prefix used for parallel-scan ordering.
82    pub hash_prefix: String,
83}
84
85/// One GSI-table row for a bulk insert via [`StorageBackend::insert_gsi_items`].
86///
87/// The fields mirror the argument order of the single-row
88/// [`StorageBackend::insert_gsi_item`].
89#[derive(Debug, Clone)]
90pub struct GsiItemRow {
91    /// GSI partition key string.
92    pub gsi_pk: String,
93    /// GSI sort key string (empty when the index has no sort key).
94    pub gsi_sk: String,
95    /// Base-table partition key string.
96    pub table_pk: String,
97    /// Base-table sort key string.
98    pub table_sk: String,
99    /// Projected item JSON.
100    pub item_json: String,
101}
102
103/// The DynamoDB size of one stored index entry, read from its projected JSON.
104/// Both backends size a non-ALL index this way, so they agree on it.
105pub(crate) fn projected_entry_size(item_json: &str) -> Result<i64, BackendError> {
106    let item: crate::types::Item = serde_json::from_str(item_json)
107        .map_err(|e| BackendError::Other(format!("unreadable index entry: {e}")))?;
108    Ok(crate::types::item_size(&item) as i64)
109}
110
111/// One vector shadow-table row for a bulk insert via
112/// [`StorageBackend::insert_vector_items`].
113///
114/// The fields mirror the shadow table's column order.
115#[derive(Debug, Clone)]
116pub struct VectorItemRow {
117    /// Base-table partition key string.
118    pub table_pk: String,
119    /// Base-table sort key string.
120    pub table_sk: String,
121    /// SearchSchema HASH attribute key string (empty when the schema declares
122    /// no HASH element).
123    pub hash_value: String,
124    /// The f32-truncated vector as a JSON array.
125    pub vector_json: String,
126    /// INLINE_FILTER attribute values as a JSON object.
127    pub filter_json: String,
128    /// Projected item JSON.
129    pub item_json: String,
130    /// The entry's billable size in bytes, by the captured write formula.
131    /// Stored so a search sums a column instead of rebuilding every entry it
132    /// scanned.
133    pub entry_bytes: i64,
134}
135
136/// One candidate row loaded from a vector shadow table for scoring by
137/// SearchVectors, via [`StorageBackend::query_vector_candidates`].
138///
139/// Carries no `hash_value`: scoping by HASH value happens in the query
140/// itself, so every returned row already matches the scope.
141#[derive(Debug, Clone)]
142pub struct VectorCandidateRow {
143    /// Base-table partition key string.
144    pub table_pk: String,
145    /// Base-table sort key string.
146    pub table_sk: String,
147    /// The stored f32-truncated vector as a JSON array.
148    pub vector_json: String,
149    /// INLINE_FILTER attribute values as a wire-shaped JSON object.
150    pub filter_json: String,
151    /// The entry's billable size, read back rather than recomputed. The
152    /// projected item is deliberately absent: a search scores from the vector
153    /// and the filter, and only the rows it returns need the item, so carrying
154    /// it here made a scan of a large index materialise every entry to throw
155    /// almost all of them away.
156    pub entry_bytes: i64,
157}
158
159/// One index-table write operation, backend-neutral.
160///
161/// The per-write and per-delete GSI/LSI fan-out builds an ordered list of these
162/// and hands it to [`StorageBackend::apply_index_writes`] in a single call. The
163/// default impl replays each op through the matching per-item method, identical
164/// to the per-op loop it replaces; the wasm backend overrides it to collapse the
165/// list into one bridge crossing. Each variant's fields mirror the argument
166/// order of the per-item method it stands in for.
167#[derive(Debug, Clone)]
168pub enum IndexWriteOp {
169    /// Remove this base key's entry from a GSI table.
170    DeleteGsi {
171        table_name: String,
172        index_name: String,
173        table_pk: String,
174        table_sk: String,
175    },
176    /// Insert (or replace) this item's projected entry into a GSI table.
177    InsertGsi {
178        table_name: String,
179        index_name: String,
180        gsi_pk: String,
181        gsi_sk: String,
182        table_pk: String,
183        table_sk: String,
184        item_json: String,
185    },
186    /// Remove this base key's entry from an LSI table.
187    DeleteLsi {
188        table_name: String,
189        index_name: String,
190        base_pk: String,
191        base_sk: String,
192    },
193    /// Insert (or replace) this item's projected entry into an LSI table.
194    InsertLsi {
195        table_name: String,
196        index_name: String,
197        pk: String,
198        sk: String,
199        base_pk: String,
200        base_sk: String,
201        item_json: String,
202    },
203    /// Remove this base key's row from a vector index shadow table.
204    DeleteVector {
205        table_name: String,
206        index_name: String,
207        table_pk: String,
208        table_sk: String,
209    },
210    /// Insert (or replace) this item's derived row into a vector index shadow
211    /// table. Carries the whole [`VectorItemRow`] rather than flattened fields
212    /// because the row struct already mirrors the shadow table's column order.
213    /// Boxed so the six-string row does not widen the whole enum: every op in
214    /// a fan-out pays for the largest variant, and the size should stay
215    /// pinned by `InsertGsi` (see the size test below).
216    InsertVector {
217        table_name: String,
218        index_name: String,
219        row: Box<VectorItemRow>,
220    },
221}
222
223/// Restricts [`StorageBackend`] to this crate's own backends.
224///
225/// The trait's method set tracks engine features rather than a stable
226/// interface: vector index support added seven methods to it in a single
227/// release. Sealing keeps that growth out of the versioning contract, because
228/// a trait nothing downstream can implement cannot be broken by a new method.
229/// The trait stays nameable, so [`Database`](crate::Database) and the wasm
230/// dispatch bounds are unaffected.
231///
232/// Unsealing later is a minor, if third-party backends ever earn their place.
233/// Sealing after 1.0.0 would not be, which is why it happens before.
234mod private {
235    pub trait Sealed {}
236
237    #[cfg(any(feature = "native-sqlite", feature = "_has-encryption"))]
238    impl Sealed for crate::storage::Storage {}
239
240    #[cfg(feature = "wasm-sqlite")]
241    impl Sealed for super::wasm_backend::WasmBridgeBackend {}
242}
243
244/// Backend-neutral storage interface.
245///
246/// Sealed: implemented by this crate's own backends only. It is public because
247/// [`Database`](crate::Database) is generic over it and the wasm dispatch
248/// functions take it as a bound, so it can be named and used as a bound from
249/// anywhere. Its method set tracks engine features rather than a stable
250/// interface, so closing it to outside impls keeps adding one out of the
251/// versioning contract. The [versioning policy] carries the reasoning.
252///
253/// [versioning policy]: https://github.com/nubo-db/dynoxide/blob/main/docs/versioning.md
254///
255/// Method signatures mirror [`Storage`](crate::storage::Storage)'s public
256/// surface 1:1, with three mechanical transformations:
257///
258/// 1. `Result<T, DynoxideError>` becomes `Result<T, BackendError>`.
259/// 2. `fn` becomes `async fn`.
260/// 3. Filesystem-typed and rusqlite-typed methods are excluded; they remain
261///    on the native [`Storage`](crate::storage::Storage) only.
262///
263/// The trait is not consumed dynamically today. The native
264/// [`Storage`](crate::storage::Storage) and the wasm
265/// `WasmBridgeBackend` each implement it
266/// monomorphically.
267///
268/// The `#[allow(async_fn_in_trait)]` reflects the monomorphic-only consumption
269/// model. The lint can be revisited if and when `dyn StorageBackend` becomes
270/// a real callsite.
271#[allow(async_fn_in_trait)]
272pub trait StorageBackend: private::Sealed {
273    // -----------------------------------------------------------------------
274    // Capabilities
275    // -----------------------------------------------------------------------
276
277    /// Wall-clock access for the stream and TTL paths.
278    ///
279    /// Sync because reading the clock is not I/O. The native backend returns
280    /// its injected [`Clock`]; the wasm SQLite backend supplies its own.
281    fn clock(&self) -> &dyn Clock;
282
283    /// Whether this backend can record and serve DynamoDB Streams.
284    ///
285    /// A backend that answers `false` still gets its stream methods called
286    /// nowhere on the happy path: the actions consult this before mutating, so
287    /// a request that needs streams is refused before anything is created,
288    /// rather than half-applied and then failed. Defaults to `true`; the wasm
289    /// backend answers `false` until a delivery mechanism exists.
290    fn supports_streams(&self) -> bool {
291        true
292    }
293
294    /// Whether this backend stores resource tags. Same contract as
295    /// [`supports_streams`](Self::supports_streams): consulted before a
296    /// mutation that would need `set_tags`, so a tagged `CreateTable` on a
297    /// backend without tags is refused whole rather than creating the table
298    /// and then failing.
299    fn supports_tags(&self) -> bool {
300        true
301    }
302
303    // -----------------------------------------------------------------------
304    // Table metadata
305    // -----------------------------------------------------------------------
306
307    async fn insert_table_metadata(&self, m: &CreateTableMetadata<'_>) -> Result<(), BackendError>;
308
309    async fn get_table_metadata(
310        &self,
311        table_name: &str,
312    ) -> Result<Option<TableMetadata>, BackendError>;
313
314    async fn delete_table_metadata(&self, table_name: &str) -> Result<bool, BackendError>;
315
316    async fn update_table_metadata(
317        &self,
318        table_name: &str,
319        attribute_definitions: &str,
320        gsi_definitions: Option<&str>,
321    ) -> Result<(), BackendError>;
322
323    /// Update a table's vector index definitions. The write the
324    /// add/delete-vector path of `UpdateTable` issues; `None` clears the
325    /// column when no vector indexes remain.
326    async fn update_vector_index_definitions(
327        &self,
328        table_name: &str,
329        vector_index_definitions: Option<&str>,
330    ) -> Result<(), BackendError>;
331
332    async fn update_provisioned_throughput(
333        &self,
334        table_name: &str,
335        provisioned_throughput: &str,
336    ) -> Result<(), BackendError>;
337
338    async fn clear_provisioned_throughput(&self, table_name: &str) -> Result<(), BackendError>;
339
340    async fn update_billing_mode(
341        &self,
342        table_name: &str,
343        billing_mode: &str,
344    ) -> Result<(), BackendError>;
345
346    async fn update_table_class(
347        &self,
348        table_name: &str,
349        table_class: &str,
350    ) -> Result<(), BackendError>;
351
352    /// Store the serialised on-demand throughput ceilings. Implementations
353    /// must store the payload opaquely: the `clear_on_demand_throughput`
354    /// default passes a JSON `null` through this method, which every reader
355    /// treats as absent, so parsing or rejecting the payload here would break
356    /// that default.
357    async fn update_on_demand_throughput(
358        &self,
359        table_name: &str,
360        on_demand_throughput: &str,
361    ) -> Result<(), BackendError>;
362
363    /// Remove any stored on-demand throughput ceilings. The default stores a
364    /// JSON `null`, which every reader treats as absent, so existing backend
365    /// implementations keep working without changes; the in-tree backends
366    /// override this to clear the underlying value properly.
367    async fn clear_on_demand_throughput(&self, table_name: &str) -> Result<(), BackendError> {
368        self.update_on_demand_throughput(table_name, "null").await
369    }
370
371    async fn get_tags(&self, table_name: &str) -> Result<Vec<Tag>, BackendError>;
372
373    async fn set_tags(&self, table_name: &str, new_tags: &[Tag]) -> Result<(), BackendError>;
374
375    async fn update_deletion_protection(
376        &self,
377        table_name: &str,
378        enabled: bool,
379    ) -> Result<(), BackendError>;
380
381    async fn remove_tags(&self, table_name: &str, keys: &[String]) -> Result<(), BackendError>;
382
383    async fn list_table_names(&self) -> Result<Vec<String>, BackendError>;
384
385    async fn table_exists(&self, table_name: &str) -> Result<bool, BackendError>;
386
387    // -----------------------------------------------------------------------
388    // Dynamic data tables (DDL)
389    // -----------------------------------------------------------------------
390
391    async fn create_data_table(&self, table_name: &str) -> Result<(), BackendError>;
392
393    async fn drop_data_table(&self, table_name: &str) -> Result<(), BackendError>;
394
395    async fn create_gsi_table(
396        &self,
397        table_name: &str,
398        index_name: &str,
399    ) -> Result<(), BackendError>;
400
401    async fn drop_gsi_table(&self, table_name: &str, index_name: &str) -> Result<(), BackendError>;
402
403    async fn create_lsi_table(
404        &self,
405        table_name: &str,
406        index_name: &str,
407    ) -> Result<(), BackendError>;
408
409    async fn drop_lsi_table(&self, table_name: &str, index_name: &str) -> Result<(), BackendError>;
410
411    /// Create the physical shadow table for a vector index.
412    async fn create_vector_table(
413        &self,
414        table_name: &str,
415        index_name: &str,
416    ) -> Result<(), BackendError>;
417
418    /// Drop the physical shadow table for a vector index, if it exists.
419    async fn drop_vector_table(
420        &self,
421        table_name: &str,
422        index_name: &str,
423    ) -> Result<(), BackendError>;
424
425    /// Bulk-insert many rows into one vector shadow table.
426    ///
427    /// Batch-shaped so a backend can amortise per-row round-trips, mirroring
428    /// [`insert_gsi_items`](Self::insert_gsi_items). Used by the vector index
429    /// backfill path, and by the default
430    /// [`apply_index_writes`](Self::apply_index_writes) with a one-row slice
431    /// for live-write maintenance.
432    async fn insert_vector_items(
433        &self,
434        table_name: &str,
435        index_name: &str,
436        rows: &[VectorItemRow],
437    ) -> Result<(), BackendError>;
438
439    /// Delete a vector shadow-table row by base-table primary key. Used by the
440    /// live-write maintenance fan-out, mirroring
441    /// [`delete_gsi_item`](Self::delete_gsi_item).
442    async fn delete_vector_item(
443        &self,
444        table_name: &str,
445        index_name: &str,
446        table_pk: &str,
447        table_sk: &str,
448    ) -> Result<(), BackendError>;
449
450    /// Load the rows SearchVectors scores, optionally scoped to one
451    /// SearchSchema HASH value (key-string encoded, matching the write
452    /// path's `hash_value` column). Rows come back ordered by base-table
453    /// primary key so equal scores tie-break deterministically.
454    async fn query_vector_candidates(
455        &self,
456        table_name: &str,
457        index_name: &str,
458        hash_value: Option<&str>,
459    ) -> Result<Vec<VectorCandidateRow>, BackendError>;
460
461    /// Load the projected entries for the base-table keys a search returns.
462    ///
463    /// Paired with [`query_vector_candidates`](Self::query_vector_candidates),
464    /// which leaves the projected item behind so a scan of a large index does
465    /// not materialise every entry only to score it and drop it. `keys` are the
466    /// post-`TopK` survivors, so this reads at most `TopK` rows. Returns
467    /// `(table_pk, table_sk, item_json)` per row found; a key with no row is
468    /// omitted rather than erroring.
469    async fn vector_items_for_keys(
470        &self,
471        table_name: &str,
472        index_name: &str,
473        keys: &[(String, String)],
474    ) -> Result<Vec<(String, String, String)>, BackendError>;
475
476    // -----------------------------------------------------------------------
477    // GSI item operations
478    // -----------------------------------------------------------------------
479
480    #[allow(clippy::too_many_arguments)]
481    async fn insert_gsi_item(
482        &self,
483        table_name: &str,
484        index_name: &str,
485        gsi_pk: &str,
486        gsi_sk: &str,
487        table_pk: &str,
488        table_sk: &str,
489        item_json: &str,
490    ) -> Result<(), BackendError>;
491
492    /// Bulk-insert many rows into one GSI table.
493    ///
494    /// Batch-shaped so a backend can amortise per-row round-trips (the native
495    /// backend reuses a single cached prepared statement). Used by the GSI
496    /// backfill path; the per-row [`insert_gsi_item`](Self::insert_gsi_item)
497    /// covers single writes during normal fan-out.
498    async fn insert_gsi_items(
499        &self,
500        table_name: &str,
501        index_name: &str,
502        rows: &[GsiItemRow],
503    ) -> Result<(), BackendError>;
504
505    async fn delete_gsi_item(
506        &self,
507        table_name: &str,
508        index_name: &str,
509        table_pk: &str,
510        table_sk: &str,
511    ) -> Result<(), BackendError>;
512
513    async fn query_gsi_items(
514        &self,
515        table_name: &str,
516        index_name: &str,
517        gsi_pk: &str,
518        params: &QueryParams<'_>,
519    ) -> Result<Vec<(String, String, String)>, BackendError>;
520
521    async fn scan_gsi_items(
522        &self,
523        table_name: &str,
524        index_name: &str,
525        params: &ScanParams<'_>,
526    ) -> Result<Vec<(String, String, String)>, BackendError>;
527
528    // -----------------------------------------------------------------------
529    // LSI item operations
530    // -----------------------------------------------------------------------
531
532    #[allow(clippy::too_many_arguments)]
533    async fn insert_lsi_item(
534        &self,
535        table_name: &str,
536        index_name: &str,
537        pk: &str,
538        sk: &str,
539        base_pk: &str,
540        base_sk: &str,
541        item_json: &str,
542    ) -> Result<(), BackendError>;
543
544    async fn delete_lsi_item(
545        &self,
546        table_name: &str,
547        index_name: &str,
548        base_pk: &str,
549        base_sk: &str,
550    ) -> Result<(), BackendError>;
551
552    async fn query_lsi_items(
553        &self,
554        table_name: &str,
555        index_name: &str,
556        pk: &str,
557        params: &QueryParams<'_>,
558    ) -> Result<Vec<(String, String, String)>, BackendError>;
559
560    async fn scan_lsi_items(
561        &self,
562        table_name: &str,
563        index_name: &str,
564        params: &ScanParams<'_>,
565    ) -> Result<Vec<(String, String, String)>, BackendError>;
566
567    // -----------------------------------------------------------------------
568    // Index write fan-out
569    // -----------------------------------------------------------------------
570
571    /// Apply an ordered batch of GSI/LSI write operations.
572    ///
573    /// The GSI/LSI maintenance helpers build the list and call this once per
574    /// fan-out instead of invoking the per-item methods one at a time. The
575    /// default impl replays each op through the matching per-item method in
576    /// order, so a backend that does not override it behaves exactly as the
577    /// per-op loop did. The wasm backend overrides this to issue the whole list
578    /// in a single bridge crossing.
579    ///
580    /// Owns no transaction: the caller's open transaction supplies atomicity, so
581    /// a mid-batch failure is rolled back by that caller. An empty list does no
582    /// work.
583    async fn apply_index_writes(&self, ops: &[IndexWriteOp]) -> Result<(), BackendError> {
584        for op in ops {
585            match op {
586                IndexWriteOp::DeleteGsi {
587                    table_name,
588                    index_name,
589                    table_pk,
590                    table_sk,
591                } => {
592                    self.delete_gsi_item(table_name, index_name, table_pk, table_sk)
593                        .await?;
594                }
595                IndexWriteOp::InsertGsi {
596                    table_name,
597                    index_name,
598                    gsi_pk,
599                    gsi_sk,
600                    table_pk,
601                    table_sk,
602                    item_json,
603                } => {
604                    self.insert_gsi_item(
605                        table_name, index_name, gsi_pk, gsi_sk, table_pk, table_sk, item_json,
606                    )
607                    .await?;
608                }
609                IndexWriteOp::DeleteLsi {
610                    table_name,
611                    index_name,
612                    base_pk,
613                    base_sk,
614                } => {
615                    self.delete_lsi_item(table_name, index_name, base_pk, base_sk)
616                        .await?;
617                }
618                IndexWriteOp::InsertLsi {
619                    table_name,
620                    index_name,
621                    pk,
622                    sk,
623                    base_pk,
624                    base_sk,
625                    item_json,
626                } => {
627                    self.insert_lsi_item(
628                        table_name, index_name, pk, sk, base_pk, base_sk, item_json,
629                    )
630                    .await?;
631                }
632                IndexWriteOp::DeleteVector {
633                    table_name,
634                    index_name,
635                    table_pk,
636                    table_sk,
637                } => {
638                    self.delete_vector_item(table_name, index_name, table_pk, table_sk)
639                        .await?;
640                }
641                IndexWriteOp::InsertVector {
642                    table_name,
643                    index_name,
644                    row,
645                } => {
646                    // The batch-shaped insert with a one-row slice: no separate
647                    // per-item method exists because the row struct already
648                    // carries the full column set.
649                    self.insert_vector_items(
650                        table_name,
651                        index_name,
652                        std::slice::from_ref(row.as_ref()),
653                    )
654                    .await?;
655                }
656            }
657        }
658        Ok(())
659    }
660
661    // -----------------------------------------------------------------------
662    // Transactions
663    // -----------------------------------------------------------------------
664
665    async fn begin_transaction(&self) -> Result<(), BackendError>;
666    async fn commit(&self) -> Result<(), BackendError>;
667    async fn rollback(&self) -> Result<(), BackendError>;
668
669    // -----------------------------------------------------------------------
670    // Schema version
671    // -----------------------------------------------------------------------
672
673    /// The schema version the database records, or `None` when it records
674    /// none, or one that isn't a number.
675    async fn recorded_schema_version(&self) -> Result<Option<u32>, BackendError>;
676
677    /// Record `version` as the database's schema version, replacing any
678    /// recorded one.
679    async fn record_schema_version(&self, version: u32) -> Result<(), BackendError>;
680
681    // -----------------------------------------------------------------------
682    // Bulk-loading PRAGMAs
683    // -----------------------------------------------------------------------
684
685    async fn enable_bulk_loading(&self) -> Result<(), BackendError>;
686    async fn disable_bulk_loading(&self) -> Result<(), BackendError>;
687
688    // -----------------------------------------------------------------------
689    // Item CRUD
690    // -----------------------------------------------------------------------
691
692    async fn put_item(
693        &self,
694        table_name: &str,
695        pk: &str,
696        sk: &str,
697        item_json: &str,
698        item_size: usize,
699    ) -> Result<Option<String>, BackendError>;
700
701    #[allow(clippy::too_many_arguments)]
702    async fn put_item_with_hash(
703        &self,
704        table_name: &str,
705        pk: &str,
706        sk: &str,
707        item_json: &str,
708        item_size: usize,
709        hash_prefix: &str,
710    ) -> Result<Option<String>, BackendError>;
711
712    /// Bulk-insert many base-table rows in one call (`INSERT OR REPLACE`).
713    ///
714    /// Batch-shaped so a backend can amortise per-row round-trips (the native
715    /// backend reuses a single cached prepared statement). Used by the import
716    /// path. Writes `cached_at` verbatim from each [`BaseItemRow`]; see the
717    /// note there for how this differs from
718    /// [`put_item_with_hash`](Self::put_item_with_hash).
719    async fn put_base_items(
720        &self,
721        table_name: &str,
722        rows: &[BaseItemRow],
723    ) -> Result<(), BackendError>;
724
725    async fn get_item(
726        &self,
727        table_name: &str,
728        pk: &str,
729        sk: &str,
730    ) -> Result<Option<String>, BackendError>;
731
732    async fn get_partition_size(&self, table_name: &str, pk: &str) -> Result<i64, BackendError>;
733
734    async fn get_lsi_partition_size(
735        &self,
736        table_name: &str,
737        index_name: &str,
738        pk: &str,
739    ) -> Result<i64, BackendError>;
740
741    async fn delete_item(
742        &self,
743        table_name: &str,
744        pk: &str,
745        sk: &str,
746    ) -> Result<Option<String>, BackendError>;
747
748    async fn query_items(
749        &self,
750        table_name: &str,
751        pk: &str,
752        params: &QueryParams<'_>,
753    ) -> Result<Vec<(String, String, String)>, BackendError>;
754
755    async fn scan_items(
756        &self,
757        table_name: &str,
758        params: &ScanParams<'_>,
759    ) -> Result<Vec<(String, String, String)>, BackendError>;
760
761    async fn count_items(&self, table_name: &str) -> Result<i64, BackendError>;
762
763    /// The number of items in a data table and the sum of their stored
764    /// `item_size`, which is the DynamoDB item size each write computed.
765    async fn item_count_and_size(&self, table_name: &str) -> Result<(i64, i64), BackendError>;
766
767    /// The number of entries in a GSI and the total of their projected item
768    /// sizes: its `ItemCount` and `IndexSizeBytes`. `whole_items` is true for
769    /// an ALL projection, where each entry is the base item and the base
770    /// table's stored size is its size; any other projection sizes each
771    /// entry's stored JSON.
772    async fn gsi_item_count_and_size(
773        &self,
774        table_name: &str,
775        index_name: &str,
776        whole_items: bool,
777    ) -> Result<(i64, i64), BackendError>;
778
779    /// [`gsi_item_count_and_size`](Self::gsi_item_count_and_size) for an LSI.
780    async fn lsi_item_count_and_size(
781        &self,
782        table_name: &str,
783        index_name: &str,
784        whole_items: bool,
785    ) -> Result<(i64, i64), BackendError>;
786
787    /// The number of entries in a vector index and the total of their stored
788    /// billable sizes: its `ItemCount` and `IndexSizeBytes`.
789    async fn vector_item_count_and_size(
790        &self,
791        table_name: &str,
792        index_name: &str,
793    ) -> Result<(i64, i64), BackendError>;
794
795    /// Rewrite the stored `item_size` of the given rows, each a
796    /// `(pk, sk, item_size)`, leaving the rows otherwise unchanged.
797    async fn set_item_sizes(
798        &self,
799        table_name: &str,
800        sizes: &[(String, String, usize)],
801    ) -> Result<(), BackendError>;
802
803    // -----------------------------------------------------------------------
804    // Introspection
805    // -----------------------------------------------------------------------
806
807    async fn db_size_bytes(&self) -> Result<u64, BackendError>;
808    async fn table_count(&self) -> Result<usize, BackendError>;
809    async fn table_stats(&self) -> Result<Vec<TableStats>, BackendError>;
810    async fn database_info(&self) -> Result<DatabaseInfo, BackendError>;
811    async fn vacuum(&self) -> Result<(), BackendError>;
812
813    // -----------------------------------------------------------------------
814    // Streams
815    // -----------------------------------------------------------------------
816
817    async fn enable_stream(
818        &self,
819        table_name: &str,
820        view_type: &str,
821        label: &str,
822    ) -> Result<(), BackendError>;
823
824    async fn disable_stream(&self, table_name: &str) -> Result<(), BackendError>;
825
826    #[allow(clippy::too_many_arguments)]
827    async fn insert_stream_record(
828        &self,
829        table_name: &str,
830        event_name: &str,
831        keys_json: &str,
832        new_image: Option<&str>,
833        old_image: Option<&str>,
834        sequence_number: &str,
835        shard_id: &str,
836        created_at: i64,
837    ) -> Result<(), BackendError>;
838
839    #[allow(clippy::too_many_arguments)]
840    async fn insert_stream_record_with_identity(
841        &self,
842        table_name: &str,
843        event_name: &str,
844        keys_json: &str,
845        new_image: Option<&str>,
846        old_image: Option<&str>,
847        sequence_number: &str,
848        shard_id: &str,
849        created_at: i64,
850        user_identity: Option<&str>,
851    ) -> Result<(), BackendError>;
852
853    async fn next_stream_sequence_number(&self, table_name: &str) -> Result<i64, BackendError>;
854
855    async fn get_stream_records(
856        &self,
857        table_name: &str,
858        shard_id: &str,
859        after_sequence: i64,
860        limit: usize,
861    ) -> Result<Vec<StreamRecord>, BackendError>;
862
863    async fn list_stream_enabled_tables(&self) -> Result<Vec<TableMetadata>, BackendError>;
864
865    // -----------------------------------------------------------------------
866    // TTL operations
867    // -----------------------------------------------------------------------
868
869    async fn update_ttl_config(
870        &self,
871        table_name: &str,
872        attribute_name: Option<&str>,
873        enabled: bool,
874    ) -> Result<(), BackendError>;
875
876    async fn list_ttl_enabled_tables(&self) -> Result<Vec<TableMetadata>, BackendError>;
877
878    async fn get_shard_sequence_range(
879        &self,
880        table_name: &str,
881        shard_id: &str,
882    ) -> Result<(Option<String>, Option<String>), BackendError>;
883
884    // -----------------------------------------------------------------------
885    // Cache tracking
886    // -----------------------------------------------------------------------
887
888    async fn touch_cached_at(
889        &self,
890        table_name: &str,
891        pk: &str,
892        sk: &str,
893        timestamp: f64,
894    ) -> Result<(), BackendError>;
895
896    async fn get_lru_items(
897        &self,
898        table_name: &str,
899        limit: usize,
900    ) -> Result<Vec<(String, String, i64)>, BackendError>;
901}
902
903#[cfg(test)]
904mod tests {
905    use super::IndexWriteOp;
906
907    /// `InsertVector` carries its row boxed so the vector family does not
908    /// widen the whole enum: every op in a fan-out list pays for the largest
909    /// variant, and before the vector variants arrived that was `InsertGsi`
910    /// at seven `String`s plus the discriminant word. Pin the enum to that
911    /// size so a future field on the vector row cannot regrow it unnoticed.
912    #[test]
913    fn index_write_op_stays_at_the_insert_gsi_driven_size() {
914        let insert_gsi_driven = 7 * std::mem::size_of::<String>() + std::mem::size_of::<usize>();
915        assert!(
916            std::mem::size_of::<IndexWriteOp>() <= insert_gsi_driven,
917            "IndexWriteOp grew past the InsertGsi-driven size: {} > {}",
918            std::mem::size_of::<IndexWriteOp>(),
919            insert_gsi_driven
920        );
921    }
922}