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}