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        Storage::get_table_metadata(self, table_name).map_err(dyno_to_backend)
59    }
60
61    async fn delete_table_metadata(&self, table_name: &str) -> Result<bool, BackendError> {
62        Storage::delete_table_metadata(self, table_name).map_err(dyno_to_backend)
63    }
64
65    async fn update_table_metadata(
66        &self,
67        table_name: &str,
68        attribute_definitions: &str,
69        gsi_definitions: Option<&str>,
70    ) -> Result<(), BackendError> {
71        Storage::update_table_metadata(self, table_name, attribute_definitions, gsi_definitions)
72            .map_err(dyno_to_backend)
73    }
74
75    async fn update_provisioned_throughput(
76        &self,
77        table_name: &str,
78        provisioned_throughput: &str,
79    ) -> Result<(), BackendError> {
80        Storage::update_provisioned_throughput(self, table_name, provisioned_throughput)
81            .map_err(dyno_to_backend)
82    }
83
84    async fn clear_provisioned_throughput(&self, table_name: &str) -> Result<(), BackendError> {
85        Storage::clear_provisioned_throughput(self, table_name).map_err(dyno_to_backend)
86    }
87
88    async fn update_billing_mode(
89        &self,
90        table_name: &str,
91        billing_mode: &str,
92    ) -> Result<(), BackendError> {
93        Storage::update_billing_mode(self, table_name, billing_mode).map_err(dyno_to_backend)
94    }
95
96    async fn update_table_class(
97        &self,
98        table_name: &str,
99        table_class: &str,
100    ) -> Result<(), BackendError> {
101        Storage::update_table_class(self, table_name, table_class).map_err(dyno_to_backend)
102    }
103
104    async fn update_on_demand_throughput(
105        &self,
106        table_name: &str,
107        on_demand_throughput: &str,
108    ) -> Result<(), BackendError> {
109        Storage::update_on_demand_throughput(self, table_name, on_demand_throughput)
110            .map_err(dyno_to_backend)
111    }
112
113    async fn clear_on_demand_throughput(&self, table_name: &str) -> Result<(), BackendError> {
114        Storage::clear_on_demand_throughput(self, table_name).map_err(dyno_to_backend)
115    }
116
117    async fn get_tags(&self, table_name: &str) -> Result<Vec<Tag>, BackendError> {
118        Storage::get_tags(self, table_name).map_err(dyno_to_backend)
119    }
120
121    async fn set_tags(&self, table_name: &str, new_tags: &[Tag]) -> Result<(), BackendError> {
122        Storage::set_tags(self, table_name, new_tags).map_err(dyno_to_backend)
123    }
124
125    async fn update_deletion_protection(
126        &self,
127        table_name: &str,
128        enabled: bool,
129    ) -> Result<(), BackendError> {
130        Storage::update_deletion_protection(self, table_name, enabled).map_err(dyno_to_backend)
131    }
132
133    async fn remove_tags(&self, table_name: &str, keys: &[String]) -> Result<(), BackendError> {
134        Storage::remove_tags(self, table_name, keys).map_err(dyno_to_backend)
135    }
136
137    async fn list_table_names(&self) -> Result<Vec<String>, BackendError> {
138        Storage::list_table_names(self).map_err(dyno_to_backend)
139    }
140
141    async fn table_exists(&self, table_name: &str) -> Result<bool, BackendError> {
142        Storage::table_exists(self, table_name).map_err(dyno_to_backend)
143    }
144
145    async fn create_data_table(&self, table_name: &str) -> Result<(), BackendError> {
146        Storage::create_data_table(self, table_name).map_err(dyno_to_backend)
147    }
148
149    async fn drop_data_table(&self, table_name: &str) -> Result<(), BackendError> {
150        Storage::drop_data_table(self, table_name).map_err(dyno_to_backend)
151    }
152
153    async fn create_gsi_table(
154        &self,
155        table_name: &str,
156        index_name: &str,
157    ) -> Result<(), BackendError> {
158        Storage::create_gsi_table(self, table_name, index_name).map_err(dyno_to_backend)
159    }
160
161    async fn drop_gsi_table(&self, table_name: &str, index_name: &str) -> Result<(), BackendError> {
162        Storage::drop_gsi_table(self, table_name, index_name).map_err(dyno_to_backend)
163    }
164
165    async fn create_lsi_table(
166        &self,
167        table_name: &str,
168        index_name: &str,
169    ) -> Result<(), BackendError> {
170        Storage::create_lsi_table(self, table_name, index_name).map_err(dyno_to_backend)
171    }
172
173    async fn drop_lsi_table(&self, table_name: &str, index_name: &str) -> Result<(), BackendError> {
174        Storage::drop_lsi_table(self, table_name, index_name).map_err(dyno_to_backend)
175    }
176
177    async fn insert_gsi_item(
178        &self,
179        table_name: &str,
180        index_name: &str,
181        gsi_pk: &str,
182        gsi_sk: &str,
183        table_pk: &str,
184        table_sk: &str,
185        item_json: &str,
186    ) -> Result<(), BackendError> {
187        Storage::insert_gsi_item(
188            self, table_name, index_name, gsi_pk, gsi_sk, table_pk, table_sk, item_json,
189        )
190        .map_err(dyno_to_backend)
191    }
192
193    async fn insert_gsi_items(
194        &self,
195        table_name: &str,
196        index_name: &str,
197        rows: &[GsiItemRow],
198    ) -> Result<(), BackendError> {
199        Storage::insert_gsi_items(self, table_name, index_name, rows).map_err(dyno_to_backend)
200    }
201
202    async fn delete_gsi_item(
203        &self,
204        table_name: &str,
205        index_name: &str,
206        table_pk: &str,
207        table_sk: &str,
208    ) -> Result<(), BackendError> {
209        Storage::delete_gsi_item(self, table_name, index_name, table_pk, table_sk)
210            .map_err(dyno_to_backend)
211    }
212
213    async fn query_gsi_items(
214        &self,
215        table_name: &str,
216        index_name: &str,
217        gsi_pk: &str,
218        params: &QueryParams<'_>,
219    ) -> Result<Vec<(String, String, String)>, BackendError> {
220        Storage::query_gsi_items(self, table_name, index_name, gsi_pk, params)
221            .map_err(dyno_to_backend)
222    }
223
224    async fn scan_gsi_items(
225        &self,
226        table_name: &str,
227        index_name: &str,
228        params: &ScanParams<'_>,
229    ) -> Result<Vec<(String, String, String)>, BackendError> {
230        Storage::scan_gsi_items(self, table_name, index_name, params).map_err(dyno_to_backend)
231    }
232
233    async fn insert_lsi_item(
234        &self,
235        table_name: &str,
236        index_name: &str,
237        pk: &str,
238        sk: &str,
239        base_pk: &str,
240        base_sk: &str,
241        item_json: &str,
242    ) -> Result<(), BackendError> {
243        Storage::insert_lsi_item(
244            self, table_name, index_name, pk, sk, base_pk, base_sk, item_json,
245        )
246        .map_err(dyno_to_backend)
247    }
248
249    async fn delete_lsi_item(
250        &self,
251        table_name: &str,
252        index_name: &str,
253        base_pk: &str,
254        base_sk: &str,
255    ) -> Result<(), BackendError> {
256        Storage::delete_lsi_item(self, table_name, index_name, base_pk, base_sk)
257            .map_err(dyno_to_backend)
258    }
259
260    async fn query_lsi_items(
261        &self,
262        table_name: &str,
263        index_name: &str,
264        pk: &str,
265        params: &QueryParams<'_>,
266    ) -> Result<Vec<(String, String, String)>, BackendError> {
267        Storage::query_lsi_items(self, table_name, index_name, pk, params).map_err(dyno_to_backend)
268    }
269
270    async fn scan_lsi_items(
271        &self,
272        table_name: &str,
273        index_name: &str,
274        params: &ScanParams<'_>,
275    ) -> Result<Vec<(String, String, String)>, BackendError> {
276        Storage::scan_lsi_items(self, table_name, index_name, params).map_err(dyno_to_backend)
277    }
278
279    async fn begin_transaction(&self) -> Result<(), BackendError> {
280        Storage::begin_transaction(self).map_err(dyno_to_backend)
281    }
282
283    async fn commit(&self) -> Result<(), BackendError> {
284        Storage::commit(self).map_err(dyno_to_backend)
285    }
286
287    async fn rollback(&self) -> Result<(), BackendError> {
288        Storage::rollback(self).map_err(dyno_to_backend)
289    }
290
291    async fn enable_bulk_loading(&self) -> Result<(), BackendError> {
292        Storage::enable_bulk_loading(self).map_err(dyno_to_backend)
293    }
294
295    async fn disable_bulk_loading(&self) -> Result<(), BackendError> {
296        Storage::disable_bulk_loading(self).map_err(dyno_to_backend)
297    }
298
299    async fn put_item(
300        &self,
301        table_name: &str,
302        pk: &str,
303        sk: &str,
304        item_json: &str,
305        item_size: usize,
306    ) -> Result<Option<String>, BackendError> {
307        Storage::put_item(self, table_name, pk, sk, item_json, item_size).map_err(dyno_to_backend)
308    }
309
310    async fn put_item_with_hash(
311        &self,
312        table_name: &str,
313        pk: &str,
314        sk: &str,
315        item_json: &str,
316        item_size: usize,
317        hash_prefix: &str,
318    ) -> Result<Option<String>, BackendError> {
319        Storage::put_item_with_hash(self, table_name, pk, sk, item_json, item_size, hash_prefix)
320            .map_err(dyno_to_backend)
321    }
322
323    async fn put_base_items(
324        &self,
325        table_name: &str,
326        rows: &[BaseItemRow],
327    ) -> Result<(), BackendError> {
328        Storage::put_base_items(self, table_name, rows).map_err(dyno_to_backend)
329    }
330
331    async fn get_item(
332        &self,
333        table_name: &str,
334        pk: &str,
335        sk: &str,
336    ) -> Result<Option<String>, BackendError> {
337        Storage::get_item(self, table_name, pk, sk).map_err(dyno_to_backend)
338    }
339
340    async fn get_partition_size(&self, table_name: &str, pk: &str) -> Result<i64, BackendError> {
341        Storage::get_partition_size(self, table_name, pk).map_err(dyno_to_backend)
342    }
343
344    async fn get_lsi_partition_size(
345        &self,
346        table_name: &str,
347        index_name: &str,
348        pk: &str,
349    ) -> Result<i64, BackendError> {
350        Storage::get_lsi_partition_size(self, table_name, index_name, pk).map_err(dyno_to_backend)
351    }
352
353    async fn delete_item(
354        &self,
355        table_name: &str,
356        pk: &str,
357        sk: &str,
358    ) -> Result<Option<String>, BackendError> {
359        Storage::delete_item(self, table_name, pk, sk).map_err(dyno_to_backend)
360    }
361
362    async fn query_items(
363        &self,
364        table_name: &str,
365        pk: &str,
366        params: &QueryParams<'_>,
367    ) -> Result<Vec<(String, String, String)>, BackendError> {
368        Storage::query_items(self, table_name, pk, params).map_err(dyno_to_backend)
369    }
370
371    async fn scan_items(
372        &self,
373        table_name: &str,
374        params: &ScanParams<'_>,
375    ) -> Result<Vec<(String, String, String)>, BackendError> {
376        Storage::scan_items(self, table_name, params).map_err(dyno_to_backend)
377    }
378
379    async fn count_items(&self, table_name: &str) -> Result<i64, BackendError> {
380        Storage::count_items(self, table_name).map_err(dyno_to_backend)
381    }
382
383    async fn db_size_bytes(&self) -> Result<u64, BackendError> {
384        Storage::db_size_bytes(self).map_err(dyno_to_backend)
385    }
386
387    async fn table_count(&self) -> Result<usize, BackendError> {
388        Storage::table_count(self).map_err(dyno_to_backend)
389    }
390
391    async fn table_stats(&self) -> Result<Vec<TableStats>, BackendError> {
392        Storage::table_stats(self).map_err(dyno_to_backend)
393    }
394
395    async fn database_info(&self) -> Result<DatabaseInfo, BackendError> {
396        Storage::database_info(self).map_err(dyno_to_backend)
397    }
398
399    async fn vacuum(&self) -> Result<(), BackendError> {
400        Storage::vacuum(self).map_err(dyno_to_backend)
401    }
402
403    async fn enable_stream(
404        &self,
405        table_name: &str,
406        view_type: &str,
407        label: &str,
408    ) -> Result<(), BackendError> {
409        Storage::enable_stream(self, table_name, view_type, label).map_err(dyno_to_backend)
410    }
411
412    async fn disable_stream(&self, table_name: &str) -> Result<(), BackendError> {
413        Storage::disable_stream(self, table_name).map_err(dyno_to_backend)
414    }
415
416    async fn insert_stream_record(
417        &self,
418        table_name: &str,
419        event_name: &str,
420        keys_json: &str,
421        new_image: Option<&str>,
422        old_image: Option<&str>,
423        sequence_number: &str,
424        shard_id: &str,
425        created_at: i64,
426    ) -> Result<(), BackendError> {
427        Storage::insert_stream_record(
428            self,
429            table_name,
430            event_name,
431            keys_json,
432            new_image,
433            old_image,
434            sequence_number,
435            shard_id,
436            created_at,
437        )
438        .map_err(dyno_to_backend)
439    }
440
441    async fn insert_stream_record_with_identity(
442        &self,
443        table_name: &str,
444        event_name: &str,
445        keys_json: &str,
446        new_image: Option<&str>,
447        old_image: Option<&str>,
448        sequence_number: &str,
449        shard_id: &str,
450        created_at: i64,
451        user_identity: Option<&str>,
452    ) -> Result<(), BackendError> {
453        Storage::insert_stream_record_with_identity(
454            self,
455            table_name,
456            event_name,
457            keys_json,
458            new_image,
459            old_image,
460            sequence_number,
461            shard_id,
462            created_at,
463            user_identity,
464        )
465        .map_err(dyno_to_backend)
466    }
467
468    async fn next_stream_sequence_number(&self, table_name: &str) -> Result<i64, BackendError> {
469        Storage::next_stream_sequence_number(self, table_name).map_err(dyno_to_backend)
470    }
471
472    async fn get_stream_records(
473        &self,
474        table_name: &str,
475        shard_id: &str,
476        after_sequence: i64,
477        limit: usize,
478    ) -> Result<Vec<StreamRecord>, BackendError> {
479        Storage::get_stream_records(self, table_name, shard_id, after_sequence, limit)
480            .map_err(dyno_to_backend)
481    }
482
483    async fn list_stream_enabled_tables(&self) -> Result<Vec<TableMetadata>, BackendError> {
484        Storage::list_stream_enabled_tables(self).map_err(dyno_to_backend)
485    }
486
487    async fn update_ttl_config(
488        &self,
489        table_name: &str,
490        attribute_name: Option<&str>,
491        enabled: bool,
492    ) -> Result<(), BackendError> {
493        Storage::update_ttl_config(self, table_name, attribute_name, enabled)
494            .map_err(dyno_to_backend)
495    }
496
497    async fn list_ttl_enabled_tables(&self) -> Result<Vec<TableMetadata>, BackendError> {
498        Storage::list_ttl_enabled_tables(self).map_err(dyno_to_backend)
499    }
500
501    async fn get_shard_sequence_range(
502        &self,
503        table_name: &str,
504        shard_id: &str,
505    ) -> Result<(Option<String>, Option<String>), BackendError> {
506        Storage::get_shard_sequence_range(self, table_name, shard_id).map_err(dyno_to_backend)
507    }
508
509    async fn touch_cached_at(
510        &self,
511        table_name: &str,
512        pk: &str,
513        sk: &str,
514        timestamp: f64,
515    ) -> Result<(), BackendError> {
516        Storage::touch_cached_at(self, table_name, pk, sk, timestamp).map_err(dyno_to_backend)
517    }
518
519    async fn get_lru_items(
520        &self,
521        table_name: &str,
522        limit: usize,
523    ) -> Result<Vec<(String, String, i64)>, BackendError> {
524        Storage::get_lru_items(self, table_name, limit).map_err(dyno_to_backend)
525    }
526}