Skip to main content

dynoxide/storage_backend/
rusqlite_impl.rs

1//! Native [`StorageBackend`] implementation backed by the rusqlite-typed
2//! [`Storage`].
3//!
4//! Each impl method is a thin async wrapper over the existing sync method on
5//! `Storage`: the body invokes the sync method synchronously and returns a
6//! ready future. No `spawn_blocking` is used, so today's behaviour of running
7//! rusqlite calls on the active executor thread is preserved exactly.
8//!
9//! # Load-bearing invariant: native futures never suspend
10//!
11//! Every method here resolves on the first poll (it runs synchronous work and
12//! returns `Ready`). The synchronous `Database` facade depends on this: it
13//! drives these futures with `block_on`, which is only safe to call from inside
14//! the tokio-based HTTP and MCP servers because an always-ready future never
15//! parks the worker thread. Do not introduce a real `.await` on external async
16//! I/O (a socket, a timer, a task spawn, a channel) into any method here
17//! without first moving the facade off `block_on`; otherwise the server worker
18//! threads will stall.
19
20use crate::errors::DynoxideError;
21use crate::storage::{
22    CreateTableMetadata, DatabaseInfo, QueryParams, ScanParams, Storage, StreamRecord,
23    TableMetadata, TableStats,
24};
25use crate::storage_backend::{BackendError, BaseItemRow, Clock, GsiItemRow, StorageBackend, error};
26use crate::types::Tag;
27
28/// Convert a [`DynoxideError`] into a [`BackendError`].
29///
30/// Storage's sync surface returns `Result<T, DynoxideError>`; the trait surface
31/// returns `Result<T, BackendError>`. The conversion preserves rusqlite error
32/// codes via [`error::from_rusqlite`]. A client-facing `ValidationException`
33/// (a backend method such as `set_tags` enforces the tag-count limit) is
34/// preserved as [`BackendError::Validation`] so the reverse conversion can
35/// restore its 400 envelope; any other variant falls through to
36/// [`BackendError::Other`] carrying the original `Display` output.
37fn dyno_to_backend(err: DynoxideError) -> BackendError {
38    match err {
39        DynoxideError::SqliteError(e) => error::from_rusqlite(e),
40        DynoxideError::ValidationException(msg) => BackendError::Validation(msg),
41        other => BackendError::Other(other.to_string()),
42    }
43}
44
45impl StorageBackend for Storage {
46    fn clock(&self) -> &dyn Clock {
47        Storage::clock(self)
48    }
49
50    async fn insert_table_metadata(&self, m: &CreateTableMetadata<'_>) -> Result<(), BackendError> {
51        Storage::insert_table_metadata(self, m).map_err(dyno_to_backend)
52    }
53
54    async fn get_table_metadata(
55        &self,
56        table_name: &str,
57    ) -> Result<Option<TableMetadata>, BackendError> {
58        // Counted at the trait boundary rather than inside `Storage`, because
59        // the figure that matters is how many times a caller asks. This backend
60        // answers most of them from its own cache; the wasm one crosses a bridge
61        // for every single call and caches nothing.
62        crate::bench_counters::record(&crate::bench_counters::METADATA_READS);
63        Storage::get_table_metadata(self, table_name).map_err(dyno_to_backend)
64    }
65
66    async fn delete_table_metadata(&self, table_name: &str) -> Result<bool, BackendError> {
67        Storage::delete_table_metadata(self, table_name).map_err(dyno_to_backend)
68    }
69
70    async fn update_table_metadata(
71        &self,
72        table_name: &str,
73        attribute_definitions: &str,
74        gsi_definitions: Option<&str>,
75    ) -> Result<(), BackendError> {
76        Storage::update_table_metadata(self, table_name, attribute_definitions, gsi_definitions)
77            .map_err(dyno_to_backend)
78    }
79
80    async fn update_vector_index_definitions(
81        &self,
82        table_name: &str,
83        vector_index_definitions: Option<&str>,
84    ) -> Result<(), BackendError> {
85        Storage::update_vector_index_definitions(self, table_name, vector_index_definitions)
86            .map_err(dyno_to_backend)
87    }
88
89    async fn update_provisioned_throughput(
90        &self,
91        table_name: &str,
92        provisioned_throughput: &str,
93    ) -> Result<(), BackendError> {
94        Storage::update_provisioned_throughput(self, table_name, provisioned_throughput)
95            .map_err(dyno_to_backend)
96    }
97
98    async fn clear_provisioned_throughput(&self, table_name: &str) -> Result<(), BackendError> {
99        Storage::clear_provisioned_throughput(self, table_name).map_err(dyno_to_backend)
100    }
101
102    async fn update_billing_mode(
103        &self,
104        table_name: &str,
105        billing_mode: &str,
106    ) -> Result<(), BackendError> {
107        Storage::update_billing_mode(self, table_name, billing_mode).map_err(dyno_to_backend)
108    }
109
110    async fn update_table_class(
111        &self,
112        table_name: &str,
113        table_class: &str,
114    ) -> Result<(), BackendError> {
115        Storage::update_table_class(self, table_name, table_class).map_err(dyno_to_backend)
116    }
117
118    async fn update_on_demand_throughput(
119        &self,
120        table_name: &str,
121        on_demand_throughput: &str,
122    ) -> Result<(), BackendError> {
123        Storage::update_on_demand_throughput(self, table_name, on_demand_throughput)
124            .map_err(dyno_to_backend)
125    }
126
127    async fn clear_on_demand_throughput(&self, table_name: &str) -> Result<(), BackendError> {
128        Storage::clear_on_demand_throughput(self, table_name).map_err(dyno_to_backend)
129    }
130
131    async fn get_tags(&self, table_name: &str) -> Result<Vec<Tag>, BackendError> {
132        Storage::get_tags(self, table_name).map_err(dyno_to_backend)
133    }
134
135    async fn set_tags(&self, table_name: &str, new_tags: &[Tag]) -> Result<(), BackendError> {
136        Storage::set_tags(self, table_name, new_tags).map_err(dyno_to_backend)
137    }
138
139    async fn update_deletion_protection(
140        &self,
141        table_name: &str,
142        enabled: bool,
143    ) -> Result<(), BackendError> {
144        Storage::update_deletion_protection(self, table_name, enabled).map_err(dyno_to_backend)
145    }
146
147    async fn remove_tags(&self, table_name: &str, keys: &[String]) -> Result<(), BackendError> {
148        Storage::remove_tags(self, table_name, keys).map_err(dyno_to_backend)
149    }
150
151    async fn list_table_names(&self) -> Result<Vec<String>, BackendError> {
152        Storage::list_table_names(self).map_err(dyno_to_backend)
153    }
154
155    async fn table_exists(&self, table_name: &str) -> Result<bool, BackendError> {
156        Storage::table_exists(self, table_name).map_err(dyno_to_backend)
157    }
158
159    async fn create_data_table(&self, table_name: &str) -> Result<(), BackendError> {
160        Storage::create_data_table(self, table_name).map_err(dyno_to_backend)
161    }
162
163    async fn drop_data_table(&self, table_name: &str) -> Result<(), BackendError> {
164        Storage::drop_data_table(self, table_name).map_err(dyno_to_backend)
165    }
166
167    async fn create_gsi_table(
168        &self,
169        table_name: &str,
170        index_name: &str,
171    ) -> Result<(), BackendError> {
172        Storage::create_gsi_table(self, table_name, index_name).map_err(dyno_to_backend)
173    }
174
175    async fn drop_gsi_table(&self, table_name: &str, index_name: &str) -> Result<(), BackendError> {
176        Storage::drop_gsi_table(self, table_name, index_name).map_err(dyno_to_backend)
177    }
178
179    async fn create_lsi_table(
180        &self,
181        table_name: &str,
182        index_name: &str,
183    ) -> Result<(), BackendError> {
184        Storage::create_lsi_table(self, table_name, index_name).map_err(dyno_to_backend)
185    }
186
187    async fn drop_lsi_table(&self, table_name: &str, index_name: &str) -> Result<(), BackendError> {
188        Storage::drop_lsi_table(self, table_name, index_name).map_err(dyno_to_backend)
189    }
190
191    async fn create_vector_table(
192        &self,
193        table_name: &str,
194        index_name: &str,
195    ) -> Result<(), BackendError> {
196        Storage::create_vector_table(self, table_name, index_name).map_err(dyno_to_backend)
197    }
198
199    async fn drop_vector_table(
200        &self,
201        table_name: &str,
202        index_name: &str,
203    ) -> Result<(), BackendError> {
204        Storage::drop_vector_table(self, table_name, index_name).map_err(dyno_to_backend)
205    }
206
207    async fn insert_vector_items(
208        &self,
209        table_name: &str,
210        index_name: &str,
211        rows: &[crate::storage_backend::VectorItemRow],
212    ) -> Result<(), BackendError> {
213        Storage::insert_vector_items(self, table_name, index_name, rows).map_err(dyno_to_backend)
214    }
215
216    async fn delete_vector_item(
217        &self,
218        table_name: &str,
219        index_name: &str,
220        table_pk: &str,
221        table_sk: &str,
222    ) -> Result<(), BackendError> {
223        Storage::delete_vector_item(self, table_name, index_name, table_pk, table_sk)
224            .map_err(dyno_to_backend)
225    }
226
227    async fn query_vector_candidates(
228        &self,
229        table_name: &str,
230        index_name: &str,
231        hash_value: Option<&str>,
232    ) -> Result<Vec<crate::storage_backend::VectorCandidateRow>, BackendError> {
233        Storage::query_vector_candidates(self, table_name, index_name, hash_value)
234            .map_err(dyno_to_backend)
235    }
236
237    async fn vector_items_for_keys(
238        &self,
239        table_name: &str,
240        index_name: &str,
241        keys: &[(String, String)],
242    ) -> Result<Vec<(String, String, String)>, BackendError> {
243        Storage::vector_items_for_keys(self, table_name, index_name, keys).map_err(dyno_to_backend)
244    }
245
246    async fn insert_gsi_item(
247        &self,
248        table_name: &str,
249        index_name: &str,
250        gsi_pk: &str,
251        gsi_sk: &str,
252        table_pk: &str,
253        table_sk: &str,
254        item_json: &str,
255    ) -> Result<(), BackendError> {
256        #[cfg(test)]
257        self.check_gsi_insert_fault().map_err(dyno_to_backend)?;
258        Storage::insert_gsi_item(
259            self, table_name, index_name, gsi_pk, gsi_sk, table_pk, table_sk, item_json,
260        )
261        .map_err(dyno_to_backend)
262    }
263
264    async fn insert_gsi_items(
265        &self,
266        table_name: &str,
267        index_name: &str,
268        rows: &[GsiItemRow],
269    ) -> Result<(), BackendError> {
270        #[cfg(test)]
271        self.check_gsi_insert_fault().map_err(dyno_to_backend)?;
272        Storage::insert_gsi_items(self, table_name, index_name, rows).map_err(dyno_to_backend)
273    }
274
275    async fn delete_gsi_item(
276        &self,
277        table_name: &str,
278        index_name: &str,
279        table_pk: &str,
280        table_sk: &str,
281    ) -> Result<(), BackendError> {
282        Storage::delete_gsi_item(self, table_name, index_name, table_pk, table_sk)
283            .map_err(dyno_to_backend)
284    }
285
286    async fn query_gsi_items(
287        &self,
288        table_name: &str,
289        index_name: &str,
290        gsi_pk: &str,
291        params: &QueryParams<'_>,
292    ) -> Result<Vec<(String, String, String)>, BackendError> {
293        Storage::query_gsi_items(self, table_name, index_name, gsi_pk, params)
294            .map_err(dyno_to_backend)
295    }
296
297    async fn scan_gsi_items(
298        &self,
299        table_name: &str,
300        index_name: &str,
301        params: &ScanParams<'_>,
302    ) -> Result<Vec<(String, String, String)>, BackendError> {
303        Storage::scan_gsi_items(self, table_name, index_name, params).map_err(dyno_to_backend)
304    }
305
306    async fn insert_lsi_item(
307        &self,
308        table_name: &str,
309        index_name: &str,
310        pk: &str,
311        sk: &str,
312        base_pk: &str,
313        base_sk: &str,
314        item_json: &str,
315    ) -> Result<(), BackendError> {
316        Storage::insert_lsi_item(
317            self, table_name, index_name, pk, sk, base_pk, base_sk, item_json,
318        )
319        .map_err(dyno_to_backend)
320    }
321
322    async fn delete_lsi_item(
323        &self,
324        table_name: &str,
325        index_name: &str,
326        base_pk: &str,
327        base_sk: &str,
328    ) -> Result<(), BackendError> {
329        Storage::delete_lsi_item(self, table_name, index_name, base_pk, base_sk)
330            .map_err(dyno_to_backend)
331    }
332
333    async fn query_lsi_items(
334        &self,
335        table_name: &str,
336        index_name: &str,
337        pk: &str,
338        params: &QueryParams<'_>,
339    ) -> Result<Vec<(String, String, String)>, BackendError> {
340        Storage::query_lsi_items(self, table_name, index_name, pk, params).map_err(dyno_to_backend)
341    }
342
343    async fn scan_lsi_items(
344        &self,
345        table_name: &str,
346        index_name: &str,
347        params: &ScanParams<'_>,
348    ) -> Result<Vec<(String, String, String)>, BackendError> {
349        Storage::scan_lsi_items(self, table_name, index_name, params).map_err(dyno_to_backend)
350    }
351
352    async fn begin_transaction(&self) -> Result<(), BackendError> {
353        Storage::begin_transaction(self).map_err(dyno_to_backend)
354    }
355
356    async fn commit(&self) -> Result<(), BackendError> {
357        Storage::commit(self).map_err(dyno_to_backend)
358    }
359
360    async fn rollback(&self) -> Result<(), BackendError> {
361        Storage::rollback(self).map_err(dyno_to_backend)
362    }
363
364    async fn recorded_schema_version(&self) -> Result<Option<u32>, BackendError> {
365        Storage::recorded_schema_version(self).map_err(dyno_to_backend)
366    }
367
368    async fn record_schema_version(&self, version: u32) -> Result<(), BackendError> {
369        Storage::record_schema_version(self, version).map_err(dyno_to_backend)
370    }
371
372    async fn enable_bulk_loading(&self) -> Result<(), BackendError> {
373        Storage::enable_bulk_loading(self).map_err(dyno_to_backend)
374    }
375
376    async fn disable_bulk_loading(&self) -> Result<(), BackendError> {
377        Storage::disable_bulk_loading(self).map_err(dyno_to_backend)
378    }
379
380    async fn put_item(
381        &self,
382        table_name: &str,
383        pk: &str,
384        sk: &str,
385        item_json: &str,
386        item_size: usize,
387    ) -> Result<Option<String>, BackendError> {
388        Storage::put_item(self, table_name, pk, sk, item_json, item_size).map_err(dyno_to_backend)
389    }
390
391    async fn put_item_with_hash(
392        &self,
393        table_name: &str,
394        pk: &str,
395        sk: &str,
396        item_json: &str,
397        item_size: usize,
398        hash_prefix: &str,
399    ) -> Result<Option<String>, BackendError> {
400        Storage::put_item_with_hash(self, table_name, pk, sk, item_json, item_size, hash_prefix)
401            .map_err(dyno_to_backend)
402    }
403
404    async fn put_base_items(
405        &self,
406        table_name: &str,
407        rows: &[BaseItemRow],
408    ) -> Result<(), BackendError> {
409        Storage::put_base_items(self, table_name, rows).map_err(dyno_to_backend)
410    }
411
412    async fn get_item(
413        &self,
414        table_name: &str,
415        pk: &str,
416        sk: &str,
417    ) -> Result<Option<String>, BackendError> {
418        Storage::get_item(self, table_name, pk, sk).map_err(dyno_to_backend)
419    }
420
421    async fn get_partition_size(&self, table_name: &str, pk: &str) -> Result<i64, BackendError> {
422        Storage::get_partition_size(self, table_name, pk).map_err(dyno_to_backend)
423    }
424
425    async fn get_lsi_partition_size(
426        &self,
427        table_name: &str,
428        index_name: &str,
429        pk: &str,
430    ) -> Result<i64, BackendError> {
431        Storage::get_lsi_partition_size(self, table_name, index_name, pk).map_err(dyno_to_backend)
432    }
433
434    async fn delete_item(
435        &self,
436        table_name: &str,
437        pk: &str,
438        sk: &str,
439    ) -> Result<Option<String>, BackendError> {
440        Storage::delete_item(self, table_name, pk, sk).map_err(dyno_to_backend)
441    }
442
443    async fn query_items(
444        &self,
445        table_name: &str,
446        pk: &str,
447        params: &QueryParams<'_>,
448    ) -> Result<Vec<(String, String, String)>, BackendError> {
449        Storage::query_items(self, table_name, pk, params).map_err(dyno_to_backend)
450    }
451
452    async fn scan_items(
453        &self,
454        table_name: &str,
455        params: &ScanParams<'_>,
456    ) -> Result<Vec<(String, String, String)>, BackendError> {
457        Storage::scan_items(self, table_name, params).map_err(dyno_to_backend)
458    }
459
460    async fn count_items(&self, table_name: &str) -> Result<i64, BackendError> {
461        Storage::count_items(self, table_name).map_err(dyno_to_backend)
462    }
463
464    async fn item_count_and_size(&self, table_name: &str) -> Result<(i64, i64), BackendError> {
465        Storage::item_count_and_size(self, table_name).map_err(dyno_to_backend)
466    }
467
468    async fn gsi_item_count_and_size(
469        &self,
470        table_name: &str,
471        index_name: &str,
472        whole_items: bool,
473    ) -> Result<(i64, i64), BackendError> {
474        Storage::gsi_item_count_and_size(self, table_name, index_name, whole_items)
475            .map_err(dyno_to_backend)
476    }
477
478    async fn lsi_item_count_and_size(
479        &self,
480        table_name: &str,
481        index_name: &str,
482        whole_items: bool,
483    ) -> Result<(i64, i64), BackendError> {
484        Storage::lsi_item_count_and_size(self, table_name, index_name, whole_items)
485            .map_err(dyno_to_backend)
486    }
487
488    async fn vector_item_count_and_size(
489        &self,
490        table_name: &str,
491        index_name: &str,
492    ) -> Result<(i64, i64), BackendError> {
493        Storage::vector_item_count_and_size(self, table_name, index_name).map_err(dyno_to_backend)
494    }
495
496    async fn set_item_sizes(
497        &self,
498        table_name: &str,
499        sizes: &[(String, String, usize)],
500    ) -> Result<(), BackendError> {
501        Storage::set_item_sizes(self, table_name, sizes).map_err(dyno_to_backend)
502    }
503
504    async fn db_size_bytes(&self) -> Result<u64, BackendError> {
505        Storage::db_size_bytes(self).map_err(dyno_to_backend)
506    }
507
508    async fn table_count(&self) -> Result<usize, BackendError> {
509        Storage::table_count(self).map_err(dyno_to_backend)
510    }
511
512    async fn table_stats(&self) -> Result<Vec<TableStats>, BackendError> {
513        Storage::table_stats(self).map_err(dyno_to_backend)
514    }
515
516    async fn database_info(&self) -> Result<DatabaseInfo, BackendError> {
517        Storage::database_info(self).map_err(dyno_to_backend)
518    }
519
520    async fn vacuum(&self) -> Result<(), BackendError> {
521        Storage::vacuum(self).map_err(dyno_to_backend)
522    }
523
524    async fn enable_stream(
525        &self,
526        table_name: &str,
527        view_type: &str,
528        label: &str,
529    ) -> Result<(), BackendError> {
530        Storage::enable_stream(self, table_name, view_type, label).map_err(dyno_to_backend)
531    }
532
533    async fn disable_stream(&self, table_name: &str) -> Result<(), BackendError> {
534        Storage::disable_stream(self, table_name).map_err(dyno_to_backend)
535    }
536
537    async fn insert_stream_record(
538        &self,
539        table_name: &str,
540        event_name: &str,
541        keys_json: &str,
542        new_image: Option<&str>,
543        old_image: Option<&str>,
544        sequence_number: &str,
545        shard_id: &str,
546        created_at: i64,
547    ) -> Result<(), BackendError> {
548        Storage::insert_stream_record(
549            self,
550            table_name,
551            event_name,
552            keys_json,
553            new_image,
554            old_image,
555            sequence_number,
556            shard_id,
557            created_at,
558        )
559        .map_err(dyno_to_backend)
560    }
561
562    async fn insert_stream_record_with_identity(
563        &self,
564        table_name: &str,
565        event_name: &str,
566        keys_json: &str,
567        new_image: Option<&str>,
568        old_image: Option<&str>,
569        sequence_number: &str,
570        shard_id: &str,
571        created_at: i64,
572        user_identity: Option<&str>,
573    ) -> Result<(), BackendError> {
574        Storage::insert_stream_record_with_identity(
575            self,
576            table_name,
577            event_name,
578            keys_json,
579            new_image,
580            old_image,
581            sequence_number,
582            shard_id,
583            created_at,
584            user_identity,
585        )
586        .map_err(dyno_to_backend)
587    }
588
589    async fn next_stream_sequence_number(&self, table_name: &str) -> Result<i64, BackendError> {
590        Storage::next_stream_sequence_number(self, table_name).map_err(dyno_to_backend)
591    }
592
593    async fn get_stream_records(
594        &self,
595        table_name: &str,
596        shard_id: &str,
597        after_sequence: i64,
598        limit: usize,
599    ) -> Result<Vec<StreamRecord>, BackendError> {
600        Storage::get_stream_records(self, table_name, shard_id, after_sequence, limit)
601            .map_err(dyno_to_backend)
602    }
603
604    async fn list_stream_enabled_tables(&self) -> Result<Vec<TableMetadata>, BackendError> {
605        Storage::list_stream_enabled_tables(self).map_err(dyno_to_backend)
606    }
607
608    async fn update_ttl_config(
609        &self,
610        table_name: &str,
611        attribute_name: Option<&str>,
612        enabled: bool,
613    ) -> Result<(), BackendError> {
614        Storage::update_ttl_config(self, table_name, attribute_name, enabled)
615            .map_err(dyno_to_backend)
616    }
617
618    async fn list_ttl_enabled_tables(&self) -> Result<Vec<TableMetadata>, BackendError> {
619        Storage::list_ttl_enabled_tables(self).map_err(dyno_to_backend)
620    }
621
622    async fn get_shard_sequence_range(
623        &self,
624        table_name: &str,
625        shard_id: &str,
626    ) -> Result<(Option<String>, Option<String>), BackendError> {
627        Storage::get_shard_sequence_range(self, table_name, shard_id).map_err(dyno_to_backend)
628    }
629
630    async fn touch_cached_at(
631        &self,
632        table_name: &str,
633        pk: &str,
634        sk: &str,
635        timestamp: f64,
636    ) -> Result<(), BackendError> {
637        Storage::touch_cached_at(self, table_name, pk, sk, timestamp).map_err(dyno_to_backend)
638    }
639
640    async fn get_lru_items(
641        &self,
642        table_name: &str,
643        limit: usize,
644    ) -> Result<Vec<(String, String, i64)>, BackendError> {
645        Storage::get_lru_items(self, table_name, limit).map_err(dyno_to_backend)
646    }
647}