1#[cfg(any(feature = "native-sqlite", feature = "_has-encryption"))]
2use crate::errors::{DynoxideError, Result};
3#[cfg(any(feature = "native-sqlite", feature = "_has-encryption"))]
4use crate::storage_backend::clock::{Clock, SystemClock};
5#[cfg(any(feature = "native-sqlite", feature = "_has-encryption"))]
6use crate::storage_backend::sql_builders::{self, escape_table_name};
7use crate::types::AttributeValue;
8#[cfg(any(feature = "native-sqlite", feature = "_has-encryption"))]
9use rusqlite::{Connection, params};
10#[cfg(any(feature = "native-sqlite", feature = "_has-encryption"))]
11use std::{cell::RefCell, collections::HashMap, sync::Arc};
12
13#[cfg(any(feature = "native-sqlite", feature = "_has-encryption"))]
15const SCHEMA_VERSION: &str = "8";
16
17pub(crate) const HASH_BUCKETS: u32 = 4096;
20
21pub fn compute_hash_prefix(pk_value: &AttributeValue) -> String {
35 let key_bytes = match pk_value {
36 AttributeValue::S(s) => s.as_bytes().to_vec(),
37 AttributeValue::N(n) => num_to_buffer(n),
38 AttributeValue::B(b) => b.clone(),
39 _ => vec![], };
41
42 let digest = md5::compute([b"Outliers" as &[u8], &key_bytes].concat());
43 format!("{:032x}", digest)[..6].to_string()
44}
45
46pub fn hash_bucket(hash_prefix: &str) -> u32 {
48 let prefix_3 = &hash_prefix[..3.min(hash_prefix.len())];
49 u32::from_str_radix(prefix_3, 16).unwrap_or(0)
50}
51
52fn num_to_buffer(num_str: &str) -> Vec<u8> {
58 let trimmed = num_str.trim();
59 if trimmed.is_empty() {
60 return vec![0x80];
61 }
62
63 use bigdecimal::BigDecimal;
64 use std::str::FromStr;
65
66 let bd = match BigDecimal::from_str(trimmed) {
67 Ok(v) => v,
68 Err(_) => return vec![0x80],
69 };
70
71 if bd.sign() == bigdecimal::num_bigint::Sign::NoSign {
72 return vec![0x80];
73 }
74
75 let is_negative = bd.sign() == bigdecimal::num_bigint::Sign::Minus;
76 let bd_abs = if is_negative { -&bd } else { bd.clone() };
77
78 let (mantissa, exponent) = extract_mantissa_and_exponent(&bd_abs);
79 if mantissa.is_empty() {
80 return vec![0x80];
81 }
82
83 let append_zero: i64 = if exponent % 2 != 0 { 1 } else { 0 };
86 let byte_len_no_exp = ((mantissa.len() as i64 + append_zero + 1) / 2) as usize;
87
88 let mut byte_array: Vec<u8>;
89 if byte_len_no_exp < 20 && is_negative {
90 byte_array = vec![0u8; byte_len_no_exp + 2];
91 byte_array[byte_len_no_exp + 1] = 102;
92 } else {
93 byte_array = vec![0u8; byte_len_no_exp + 1];
94 }
95
96 let exp_sum = exponent + append_zero;
100 let exp_byte_val = floor_div(exp_sum, 2) - 64;
101 if is_negative {
102 byte_array[0] = (exp_byte_val ^ !0i64) as u8;
105 } else {
106 byte_array[0] = exp_byte_val as u8;
107 }
108
109 let mut mi: i64 = 0; let mlen = mantissa.len() as i64;
113 let mut appended_zero = false;
114
115 while mi < mlen {
116 let bai = ((mi + append_zero) / 2 + 1) as usize; if append_zero != 0 && mi == 0 && !appended_zero {
118 byte_array[bai] = 0;
119 appended_zero = true;
120 mi -= 1; } else if (mi + append_zero) % 2 == 0 {
122 byte_array[bai] = mantissa[mi as usize] * 10;
123 } else {
124 byte_array[bai] += mantissa[mi as usize];
125 }
126
127 if ((mi + append_zero) % 2 != 0) || (mi == mlen - 1) {
129 if is_negative {
130 byte_array[bai] = 101u8.wrapping_sub(byte_array[bai]);
131 } else {
132 byte_array[bai] = byte_array[bai].wrapping_add(1);
133 }
134 }
135
136 mi += 1; }
138
139 byte_array
140}
141
142fn floor_div(a: i64, b: i64) -> i64 {
144 let d = a / b;
145 let r = a % b;
146 if (r != 0) && ((r ^ b) < 0) { d - 1 } else { d }
147}
148
149fn extract_mantissa_and_exponent(bd: &bigdecimal::BigDecimal) -> (Vec<u8>, i64) {
156 let normalized = bd.normalized();
158
159 let (bigint, scale) = normalized.as_bigint_and_exponent();
162 let digits_str = bigint.to_string();
163 let digits_str = digits_str.trim_start_matches('-');
164
165 let digits: Vec<u8> = digits_str
166 .chars()
167 .map(|c| c.to_digit(10).unwrap() as u8)
168 .collect();
169
170 let exponent = digits.len() as i64 - scale;
173
174 (digits, exponent)
175}
176
177pub fn hash_in_segment(hash_prefix: &str, segment: u32, total_segments: u32) -> bool {
184 let bucket = hash_bucket(hash_prefix);
185 let start = ceiling_div(HASH_BUCKETS * segment, total_segments);
186 let end = ceiling_div(HASH_BUCKETS * (segment + 1), total_segments) - 1;
187 bucket >= start && bucket <= end
188}
189
190pub(crate) fn ceiling_div(a: u32, b: u32) -> u32 {
191 a.div_ceil(b)
192}
193
194#[derive(Debug, Default)]
196pub struct ScanParams<'a> {
197 pub limit: Option<usize>,
198 pub exclusive_start_pk: Option<&'a str>,
199 pub exclusive_start_sk: Option<&'a str>,
200 pub segment: Option<u32>,
201 pub total_segments: Option<u32>,
202 pub exclusive_start_base_pk: Option<&'a str>,
204 pub exclusive_start_base_sk: Option<&'a str>,
206}
207
208#[derive(Debug, Default)]
210pub struct CreateTableMetadata<'a> {
211 pub table_name: &'a str,
212 pub key_schema: &'a str,
213 pub attribute_definitions: &'a str,
214 pub gsi_definitions: Option<&'a str>,
215 pub lsi_definitions: Option<&'a str>,
216 pub provisioned_throughput: Option<&'a str>,
217 pub created_at: i64,
218 pub sse_specification: Option<&'a str>,
219 pub table_class: Option<&'a str>,
220 pub deletion_protection_enabled: bool,
221 pub billing_mode: Option<&'a str>,
222 pub on_demand_throughput: Option<&'a str>,
223}
224
225#[derive(Debug, Default)]
227pub struct QueryParams<'a> {
228 pub sk_condition: Option<&'a str>,
229 pub sk_params: &'a [&'a str],
230 pub forward: bool,
231 pub limit: Option<usize>,
232 pub exclusive_start_sk: Option<&'a str>,
233 pub exclusive_start_base_pk: Option<&'a str>,
235 pub exclusive_start_base_sk: Option<&'a str>,
237}
238
239#[cfg(any(feature = "native-sqlite", feature = "_has-encryption"))]
247pub struct Storage {
248 conn: Connection,
249 metadata_cache: RefCell<HashMap<String, TableMetadata>>,
252 clock: Arc<dyn Clock>,
254}
255
256#[cfg(any(feature = "native-sqlite", feature = "_has-encryption"))]
257impl Storage {
258 pub fn new(path: &str) -> Result<Self> {
260 let conn = Connection::open(path)?;
261 let mut storage = Self {
262 conn,
263 metadata_cache: RefCell::new(HashMap::new()),
264 clock: Arc::new(SystemClock),
265 };
266 storage.initialize().map_err(Self::maybe_encrypted_error)?;
267 Ok(storage)
268 }
269
270 pub fn with_clock(mut self, clock: Arc<dyn Clock>) -> Self {
276 self.clock = clock;
277 self
278 }
279
280 pub(crate) fn clock(&self) -> &dyn Clock {
282 self.clock.as_ref()
283 }
284
285 fn maybe_encrypted_error(err: DynoxideError) -> DynoxideError {
288 if let DynoxideError::SqliteError(ref sqlite_err) = err {
289 if let Some(rusqlite::ErrorCode::NotADatabase) = sqlite_err.sqlite_error_code() {
290 return DynoxideError::InternalServerError(
291 "Database file is encrypted or not a valid SQLite database. \
292 If encrypted, enable the `encryption` or `encryption-cc` feature \
293 and use Database::new_encrypted() with the correct key."
294 .to_string(),
295 );
296 }
297 }
298 err
299 }
300
301 #[cfg(feature = "_has-encryption")]
307 pub fn new_encrypted(path: &str, key: &str) -> Result<Self> {
308 use zeroize::Zeroize;
309
310 let conn = Connection::open(path)?;
311 let mut pragma_val = format!("x'{key}'");
316 conn.pragma_update(None, "key", &pragma_val)?;
317 pragma_val.zeroize();
318 conn.execute_batch("SELECT count(*) FROM sqlite_master;")?;
320 let mut storage = Self {
321 conn,
322 metadata_cache: RefCell::new(HashMap::new()),
323 clock: Arc::new(SystemClock),
324 };
325 storage.initialize()?;
326 Ok(storage)
327 }
328
329 pub fn memory() -> Result<Self> {
331 let conn = Connection::open_in_memory()?;
332 let mut storage = Self {
333 conn,
334 metadata_cache: RefCell::new(HashMap::new()),
335 clock: Arc::new(SystemClock),
336 };
337 storage.initialize()?;
338 Ok(storage)
339 }
340
341 fn initialize(&mut self) -> Result<()> {
343 self.conn.pragma_update(None, "journal_mode", "WAL")?;
345
346 self.conn.create_scalar_function(
349 "fnv1a_hash",
350 1,
351 rusqlite::functions::FunctionFlags::SQLITE_DETERMINISTIC
352 | rusqlite::functions::FunctionFlags::SQLITE_UTF8,
353 |ctx: &rusqlite::functions::Context| -> rusqlite::Result<i64> {
354 let pk_ref = ctx.get_raw(0);
355 let pk_bytes = match pk_ref {
356 rusqlite::types::ValueRef::Text(bytes) => bytes,
357 _ => {
358 return Err(rusqlite::Error::InvalidFunctionParameterType(
359 0,
360 rusqlite::types::Type::Text,
361 ));
362 }
363 };
364 let mut hash: u32 = 2166136261;
365 for &byte in pk_bytes {
366 hash ^= byte as u32;
367 hash = hash.wrapping_mul(16777619);
368 }
369 Ok(hash as i64)
370 },
371 )?;
372
373 self.conn.execute_batch(sql_builders::INIT_SCHEMA)?;
375
376 let _ = self
378 .conn
379 .execute_batch("ALTER TABLE _stream_records ADD COLUMN user_identity TEXT");
380
381 self.conn.execute(
383 "INSERT OR IGNORE INTO _config (key, value) VALUES ('schema_version', ?1)",
384 params![SCHEMA_VERSION],
385 )?;
386
387 let version: i32 = self
389 .conn
390 .query_row(
391 "SELECT value FROM _config WHERE key = 'schema_version'",
392 [],
393 |r| r.get::<_, String>(0),
394 )
395 .unwrap_or_else(|_| "1".to_string())
396 .parse()
397 .unwrap_or(1);
398
399 if version < 2 {
400 self.migrate_v1_to_v2()?;
401 }
402 if version < 3 {
403 self.migrate_v2_to_v3()?;
404 }
405 if version < 4 {
406 self.migrate_v3_to_v4()?;
407 }
408 if version < 5 {
409 self.migrate_v4_to_v5()?;
410 }
411 if version < 6 {
412 self.migrate_v5_to_v6()?;
413 }
414 if version < 7 {
415 self.migrate_v6_to_v7()?;
416 }
417 if version < 8 {
418 self.migrate_v7_to_v8()?;
419 }
420
421 Ok(())
422 }
423
424 fn migrate_v1_to_v2(&self) -> Result<()> {
426 let mut stmt = self.conn.prepare("SELECT table_name FROM _tables")?;
427 let table_names: Vec<String> = stmt
428 .query_map([], |row| row.get(0))?
429 .collect::<std::result::Result<Vec<_>, _>>()?;
430
431 for table_name in &table_names {
432 let escaped = format!("\"{}\"", table_name.replace('"', "\"\""));
433 let _ = self.conn.execute(
434 &format!("ALTER TABLE {escaped} ADD COLUMN cached_at REAL"),
435 [],
436 );
437 }
438
439 self.conn.execute(
440 "INSERT OR REPLACE INTO _config (key, value) VALUES ('schema_version', '2')",
441 [],
442 )?;
443
444 Ok(())
445 }
446
447 fn migrate_v2_to_v3(&self) -> Result<()> {
449 let _ = self
450 .conn
451 .execute("ALTER TABLE _tables ADD COLUMN tags TEXT", []);
452
453 self.conn.execute(
454 "INSERT OR REPLACE INTO _config (key, value) VALUES ('schema_version', '3')",
455 [],
456 )?;
457
458 Ok(())
459 }
460
461 fn migrate_v3_to_v4(&self) -> Result<()> {
463 let _ = self
464 .conn
465 .execute("ALTER TABLE _tables ADD COLUMN sse_specification TEXT", []);
466 let _ = self
467 .conn
468 .execute("ALTER TABLE _tables ADD COLUMN table_class TEXT", []);
469 let _ = self.conn.execute(
470 "ALTER TABLE _tables ADD COLUMN deletion_protection_enabled INTEGER DEFAULT 0",
471 [],
472 );
473
474 self.conn.execute(
475 "INSERT OR REPLACE INTO _config (key, value) VALUES ('schema_version', '4')",
476 [],
477 )?;
478
479 Ok(())
480 }
481
482 fn migrate_v4_to_v5(&self) -> Result<()> {
484 let mut stmt = self
485 .conn
486 .prepare("SELECT table_name, gsi_definitions, lsi_definitions FROM _tables")?;
487 let tables: Vec<(String, Option<String>, Option<String>)> = stmt
488 .query_map([], |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)))?
489 .collect::<std::result::Result<Vec<_>, _>>()?;
490
491 for (table_name, gsi_json, lsi_json) in &tables {
492 if let Some(json) = gsi_json {
493 if let Ok(gsis) = serde_json::from_str::<Vec<serde_json::Value>>(json) {
494 for gsi in &gsis {
495 if let Some(idx) = gsi.get("IndexName").and_then(|v| v.as_str()) {
496 let gsi_table = escape_table_name(&format!("{table_name}::gsi::{idx}"));
497 let idx_name =
498 escape_table_name(&format!("{table_name}::gsi::{idx}::base_key"));
499 let _ = self.conn.execute_batch(&format!(
500 "CREATE INDEX IF NOT EXISTS \"{idx_name}\" ON \"{gsi_table}\" (table_pk, table_sk)"
501 ));
502 }
503 }
504 }
505 }
506 if let Some(json) = lsi_json {
507 if let Ok(lsis) = serde_json::from_str::<Vec<serde_json::Value>>(json) {
508 for lsi in &lsis {
509 if let Some(idx) = lsi.get("IndexName").and_then(|v| v.as_str()) {
510 let lsi_table = escape_table_name(&format!("{table_name}::lsi::{idx}"));
511 let idx_name =
512 escape_table_name(&format!("{table_name}::lsi::{idx}::base_key"));
513 let _ = self.conn.execute_batch(&format!(
514 "CREATE INDEX IF NOT EXISTS \"{idx_name}\" ON \"{lsi_table}\" (base_pk, base_sk)"
515 ));
516 }
517 }
518 }
519 }
520 }
521
522 self.conn.execute(
523 "INSERT OR REPLACE INTO _config (key, value) VALUES ('schema_version', '5')",
524 [],
525 )?;
526
527 Ok(())
528 }
529
530 fn migrate_v5_to_v6(&self) -> Result<()> {
532 let mut stmt = self.conn.prepare("SELECT table_name FROM _tables")?;
533 let table_names: Vec<String> = stmt
534 .query_map([], |row| row.get(0))?
535 .collect::<std::result::Result<Vec<_>, _>>()?;
536
537 for table_name in &table_names {
538 let escaped = escape_table_name(table_name);
539 let _ = self.conn.execute(
540 &format!(
541 "ALTER TABLE \"{escaped}\" ADD COLUMN hash_prefix TEXT NOT NULL DEFAULT ''"
542 ),
543 [],
544 );
545 }
546
547 self.conn.execute(
548 "INSERT OR REPLACE INTO _config (key, value) VALUES ('schema_version', '6')",
549 [],
550 )?;
551
552 Ok(())
553 }
554
555 fn migrate_v6_to_v7(&self) -> Result<()> {
557 let _ = self.conn.execute(
558 "ALTER TABLE _tables ADD COLUMN on_demand_throughput TEXT",
559 [],
560 );
561
562 self.conn.execute(
563 "INSERT OR REPLACE INTO _config (key, value) VALUES ('schema_version', '7')",
564 [],
565 )?;
566
567 Ok(())
568 }
569
570 fn migrate_v7_to_v8(&self) -> Result<()> {
581 let _ = self
582 .conn
583 .execute("ALTER TABLE _tables ADD COLUMN table_id TEXT", []);
584
585 let names: Vec<String> = {
586 let mut stmt = self
587 .conn
588 .prepare("SELECT table_name FROM _tables WHERE table_id IS NULL")?;
589 let rows = stmt.query_map([], |row| row.get::<_, String>(0))?;
590 rows.collect::<rusqlite::Result<Vec<String>>>()?
591 };
592 for name in names {
593 let id = uuid::Uuid::new_v4().to_string();
594 self.conn.execute(
595 "UPDATE _tables SET table_id = ?1 WHERE table_name = ?2",
596 params![id, name],
597 )?;
598 }
599
600 self.conn.execute(
601 "INSERT OR REPLACE INTO _config (key, value) VALUES ('schema_version', '8')",
602 [],
603 )?;
604
605 Ok(())
606 }
607
608 pub fn conn(&self) -> &Connection {
610 &self.conn
611 }
612
613 pub fn conn_mut(&mut self) -> &mut Connection {
615 &mut self.conn
616 }
617
618 pub fn insert_table_metadata(&self, m: &CreateTableMetadata) -> Result<()> {
624 let table_name = m.table_name;
625 let (sql, params) = sql_builders::insert_table_metadata(m);
626 self.conn
627 .execute(&sql, rusqlite::params_from_iter(params.iter()))?;
628 self.metadata_cache.borrow_mut().remove(table_name);
629 Ok(())
630 }
631
632 pub fn get_table_metadata(&self, table_name: &str) -> Result<Option<TableMetadata>> {
638 if let Some(cached) = self.metadata_cache.borrow().get(table_name) {
640 return Ok(Some(cached.clone()));
641 }
642
643 let (sql, params) = sql_builders::get_table_metadata(table_name);
644 let mut stmt = self.conn.prepare(&sql)?;
645
646 let result = stmt.query_row(rusqlite::params_from_iter(params.iter()), row_to_metadata);
647
648 match result {
649 Ok(meta) => {
650 self.metadata_cache
651 .borrow_mut()
652 .insert(table_name.to_string(), meta.clone());
653 Ok(Some(meta))
654 }
655 Err(rusqlite::Error::QueryReturnedNoRows) => Ok(None),
656 Err(e) => Err(DynoxideError::from(e)),
657 }
658 }
659
660 pub fn delete_table_metadata(&self, table_name: &str) -> Result<bool> {
662 let (sql, params) = sql_builders::delete_table_metadata(table_name);
663 let affected = self
664 .conn
665 .execute(&sql, rusqlite::params_from_iter(params.iter()))?;
666 self.metadata_cache.borrow_mut().remove(table_name);
667 Ok(affected > 0)
668 }
669
670 pub fn update_table_metadata(
672 &self,
673 table_name: &str,
674 attribute_definitions: &str,
675 gsi_definitions: Option<&str>,
676 ) -> Result<()> {
677 let (sql, params) =
678 sql_builders::update_table_metadata(table_name, attribute_definitions, gsi_definitions);
679 self.conn
680 .execute(&sql, rusqlite::params_from_iter(params.iter()))?;
681 self.metadata_cache.borrow_mut().remove(table_name);
682 Ok(())
683 }
684
685 pub fn update_provisioned_throughput(
687 &self,
688 table_name: &str,
689 provisioned_throughput: &str,
690 ) -> Result<()> {
691 let (sql, params) =
692 sql_builders::update_provisioned_throughput(table_name, provisioned_throughput);
693 self.conn
694 .execute(&sql, rusqlite::params_from_iter(params.iter()))?;
695 self.metadata_cache.borrow_mut().remove(table_name);
696 Ok(())
697 }
698
699 pub fn clear_provisioned_throughput(&self, table_name: &str) -> Result<()> {
701 let (sql, params) = sql_builders::clear_provisioned_throughput(table_name);
702 self.conn
703 .execute(&sql, rusqlite::params_from_iter(params.iter()))?;
704 self.metadata_cache.borrow_mut().remove(table_name);
705 Ok(())
706 }
707
708 pub fn clear_on_demand_throughput(&self, table_name: &str) -> Result<()> {
710 let (sql, params) = sql_builders::clear_on_demand_throughput(table_name);
711 self.conn
712 .execute(&sql, rusqlite::params_from_iter(params.iter()))?;
713 self.metadata_cache.borrow_mut().remove(table_name);
714 Ok(())
715 }
716
717 pub fn update_billing_mode(&self, table_name: &str, billing_mode: &str) -> Result<()> {
719 let (sql, params) = sql_builders::update_billing_mode(table_name, billing_mode);
720 self.conn
721 .execute(&sql, rusqlite::params_from_iter(params.iter()))?;
722 self.metadata_cache.borrow_mut().remove(table_name);
723 Ok(())
724 }
725
726 pub fn update_table_class(&self, table_name: &str, table_class: &str) -> Result<()> {
728 let (sql, params) = sql_builders::update_table_class(table_name, table_class);
729 self.conn
730 .execute(&sql, rusqlite::params_from_iter(params.iter()))?;
731 self.metadata_cache.borrow_mut().remove(table_name);
732 Ok(())
733 }
734
735 pub fn update_on_demand_throughput(
737 &self,
738 table_name: &str,
739 on_demand_throughput: &str,
740 ) -> Result<()> {
741 let (sql, params) =
742 sql_builders::update_on_demand_throughput(table_name, on_demand_throughput);
743 self.conn
744 .execute(&sql, rusqlite::params_from_iter(params.iter()))?;
745 self.metadata_cache.borrow_mut().remove(table_name);
746 Ok(())
747 }
748
749 pub fn get_tags(&self, table_name: &str) -> Result<Vec<crate::types::Tag>> {
755 let tags_json: Option<String> = self.conn.query_row(
756 "SELECT tags FROM _tables WHERE table_name = ?1",
757 params![table_name],
758 |row| row.get(0),
759 )?;
760
761 match tags_json {
762 Some(json) => serde_json::from_str(&json)
763 .map_err(|e| DynoxideError::InternalServerError(format!("Bad tags JSON: {e}"))),
764 None => Ok(Vec::new()),
765 }
766 }
767
768 pub fn set_tags(&self, table_name: &str, new_tags: &[crate::types::Tag]) -> Result<()> {
770 use std::collections::BTreeMap;
771
772 let existing = self.get_tags(table_name)?;
773 let mut tag_map: BTreeMap<String, String> =
774 existing.into_iter().map(|t| (t.key, t.value)).collect();
775
776 for tag in new_tags {
777 tag_map.insert(tag.key.clone(), tag.value.clone());
778 }
779
780 if tag_map.len() > 50 {
781 return Err(DynoxideError::ValidationException(
782 "One or more parameter values were invalid: \
783 Too many tags: tag limit is 50"
784 .to_string(),
785 ));
786 }
787
788 let merged: Vec<crate::types::Tag> = tag_map
789 .into_iter()
790 .map(|(k, v)| crate::types::Tag { key: k, value: v })
791 .collect();
792
793 let json = serde_json::to_string(&merged)
794 .map_err(|e| DynoxideError::InternalServerError(e.to_string()))?;
795
796 self.conn.execute(
797 "UPDATE _tables SET tags = ?1 WHERE table_name = ?2",
798 params![json, table_name],
799 )?;
800 Ok(())
801 }
802
803 pub fn update_deletion_protection(&self, table_name: &str, enabled: bool) -> Result<()> {
805 let (sql, params) = sql_builders::update_deletion_protection(table_name, enabled);
806 self.conn
807 .execute(&sql, rusqlite::params_from_iter(params.iter()))?;
808 self.metadata_cache.borrow_mut().remove(table_name);
809 Ok(())
810 }
811
812 pub fn remove_tags(&self, table_name: &str, keys: &[String]) -> Result<()> {
814 let mut tags = self.get_tags(table_name)?;
815 tags.retain(|t| !keys.contains(&t.key));
816
817 let json = if tags.is_empty() {
818 None
819 } else {
820 Some(
821 serde_json::to_string(&tags)
822 .map_err(|e| DynoxideError::InternalServerError(e.to_string()))?,
823 )
824 };
825
826 self.conn.execute(
827 "UPDATE _tables SET tags = ?1 WHERE table_name = ?2",
828 params![json, table_name],
829 )?;
830 Ok(())
831 }
832
833 pub fn list_table_names(&self) -> Result<Vec<String>> {
835 let (sql, params) = sql_builders::list_table_names();
836 let mut stmt = self.conn.prepare(&sql)?;
837 let names = stmt
838 .query_map(rusqlite::params_from_iter(params.iter()), |row| row.get(0))?
839 .collect::<std::result::Result<Vec<String>, _>>()?;
840 Ok(names)
841 }
842
843 pub fn table_exists(&self, table_name: &str) -> Result<bool> {
845 let (sql, params) = sql_builders::table_exists(table_name);
846 let count: i32 =
847 self.conn
848 .query_row(&sql, rusqlite::params_from_iter(params.iter()), |row| {
849 row.get(0)
850 })?;
851 Ok(count > 0)
852 }
853
854 #[allow(dead_code)]
856 pub(crate) fn invalidate_metadata_cache(&self, table_name: &str) {
857 self.metadata_cache.borrow_mut().remove(table_name);
858 }
859
860 pub fn create_data_table(&self, table_name: &str) -> Result<()> {
866 let (sql, params) = sql_builders::create_data_table(table_name);
867 self.conn
868 .execute(&sql, rusqlite::params_from_iter(params.iter()))?;
869 Ok(())
870 }
871
872 pub fn drop_data_table(&self, table_name: &str) -> Result<()> {
874 let (sql, params) = sql_builders::drop_data_table(table_name);
875 self.conn
876 .execute(&sql, rusqlite::params_from_iter(params.iter()))?;
877 Ok(())
878 }
879
880 pub fn create_gsi_table(&self, table_name: &str, index_name: &str) -> Result<()> {
882 let (sql, _) = sql_builders::create_gsi_table(table_name, index_name);
883 self.conn.execute_batch(&sql)?;
884 Ok(())
885 }
886
887 pub fn drop_gsi_table(&self, table_name: &str, index_name: &str) -> Result<()> {
889 let (sql, params) = sql_builders::drop_gsi_table(table_name, index_name);
890 self.conn
891 .execute(&sql, rusqlite::params_from_iter(params.iter()))?;
892 Ok(())
893 }
894
895 #[allow(clippy::too_many_arguments)]
901 pub fn insert_gsi_item(
902 &self,
903 table_name: &str,
904 index_name: &str,
905 gsi_pk: &str,
906 gsi_sk: &str,
907 table_pk: &str,
908 table_sk: &str,
909 item_json: &str,
910 ) -> Result<()> {
911 let sql = sql_builders::gsi_insert_sql(table_name, index_name);
912 let params = sql_builders::gsi_insert_params(gsi_pk, gsi_sk, table_pk, table_sk, item_json);
913 self.conn
914 .prepare_cached(&sql)?
915 .execute(rusqlite::params_from_iter(params.iter()))?;
916 Ok(())
917 }
918
919 pub fn insert_gsi_items(
923 &self,
924 table_name: &str,
925 index_name: &str,
926 rows: &[crate::storage_backend::GsiItemRow],
927 ) -> Result<()> {
928 let sql = sql_builders::gsi_insert_sql(table_name, index_name);
929 let mut stmt = self.conn.prepare_cached(&sql)?;
930 for row in rows {
931 let params = sql_builders::gsi_insert_params(
932 &row.gsi_pk,
933 &row.gsi_sk,
934 &row.table_pk,
935 &row.table_sk,
936 &row.item_json,
937 );
938 stmt.execute(rusqlite::params_from_iter(params.iter()))?;
939 }
940 Ok(())
941 }
942
943 pub fn delete_gsi_item(
945 &self,
946 table_name: &str,
947 index_name: &str,
948 table_pk: &str,
949 table_sk: &str,
950 ) -> Result<()> {
951 let (sql, params) =
952 sql_builders::delete_gsi_item(table_name, index_name, table_pk, table_sk);
953 self.conn
954 .prepare_cached(&sql)?
955 .execute(rusqlite::params_from_iter(params.iter()))?;
956 Ok(())
957 }
958
959 pub fn query_gsi_items(
961 &self,
962 table_name: &str,
963 index_name: &str,
964 gsi_pk: &str,
965 params: &QueryParams,
966 ) -> Result<Vec<(String, String, String)>> {
967 let (sql, params_vec) =
968 sql_builders::query_gsi_items(table_name, index_name, gsi_pk, params);
969 let mut stmt = self.conn.prepare(&sql)?;
970 let rows = stmt
971 .query_map(rusqlite::params_from_iter(params_vec.iter()), |row| {
972 Ok((row.get(0)?, row.get(1)?, row.get(2)?))
973 })?
974 .collect::<std::result::Result<Vec<_>, _>>()?;
975
976 Ok(rows)
977 }
978
979 pub fn scan_gsi_items(
981 &self,
982 table_name: &str,
983 index_name: &str,
984 params: &ScanParams,
985 ) -> Result<Vec<(String, String, String)>> {
986 let (sql, params_vec) = sql_builders::scan_gsi_items(table_name, index_name, params);
987 let mut stmt = self.conn.prepare(&sql)?;
988 let rows = stmt
989 .query_map(rusqlite::params_from_iter(params_vec.iter()), |row| {
990 Ok((row.get(0)?, row.get(1)?, row.get(2)?))
991 })?
992 .collect::<std::result::Result<Vec<_>, _>>()?;
993
994 Ok(rows)
995 }
996
997 pub fn create_lsi_table(&self, table_name: &str, index_name: &str) -> Result<()> {
1003 let (sql, _) = sql_builders::create_lsi_table(table_name, index_name);
1004 self.conn.execute_batch(&sql)?;
1005 Ok(())
1006 }
1007
1008 pub fn drop_lsi_table(&self, table_name: &str, index_name: &str) -> Result<()> {
1010 let (sql, params) = sql_builders::drop_lsi_table(table_name, index_name);
1011 self.conn
1012 .execute(&sql, rusqlite::params_from_iter(params.iter()))?;
1013 Ok(())
1014 }
1015
1016 #[allow(clippy::too_many_arguments)]
1022 pub fn insert_lsi_item(
1023 &self,
1024 table_name: &str,
1025 index_name: &str,
1026 pk: &str,
1027 sk: &str,
1028 base_pk: &str,
1029 base_sk: &str,
1030 item_json: &str,
1031 ) -> Result<()> {
1032 let sql = sql_builders::lsi_insert_sql(table_name, index_name);
1033 let params = sql_builders::lsi_insert_params(pk, sk, base_pk, base_sk, item_json);
1034 self.conn
1035 .prepare_cached(&sql)?
1036 .execute(rusqlite::params_from_iter(params.iter()))?;
1037 Ok(())
1038 }
1039
1040 pub fn delete_lsi_item(
1042 &self,
1043 table_name: &str,
1044 index_name: &str,
1045 base_pk: &str,
1046 base_sk: &str,
1047 ) -> Result<()> {
1048 let (sql, params) = sql_builders::delete_lsi_item(table_name, index_name, base_pk, base_sk);
1049 self.conn
1050 .prepare_cached(&sql)?
1051 .execute(rusqlite::params_from_iter(params.iter()))?;
1052 Ok(())
1053 }
1054
1055 pub fn query_lsi_items(
1057 &self,
1058 table_name: &str,
1059 index_name: &str,
1060 pk: &str,
1061 params: &QueryParams,
1062 ) -> Result<Vec<(String, String, String)>> {
1063 let (sql, params_vec) = sql_builders::query_lsi_items(table_name, index_name, pk, params);
1064 let mut stmt = self.conn.prepare(&sql)?;
1065 let rows = stmt
1066 .query_map(rusqlite::params_from_iter(params_vec.iter()), |row| {
1067 Ok((row.get(0)?, row.get(1)?, row.get(2)?))
1068 })?
1069 .collect::<std::result::Result<Vec<_>, _>>()?;
1070
1071 Ok(rows)
1072 }
1073
1074 pub fn scan_lsi_items(
1076 &self,
1077 table_name: &str,
1078 index_name: &str,
1079 params: &ScanParams,
1080 ) -> Result<Vec<(String, String, String)>> {
1081 let (sql, params_vec) = sql_builders::scan_lsi_items(table_name, index_name, params);
1082 let mut stmt = self.conn.prepare(&sql)?;
1083 let rows = stmt
1084 .query_map(rusqlite::params_from_iter(params_vec.iter()), |row| {
1085 Ok((row.get(0)?, row.get(1)?, row.get(2)?))
1086 })?
1087 .collect::<std::result::Result<Vec<_>, _>>()?;
1088
1089 Ok(rows)
1090 }
1091
1092 pub fn begin_transaction(&self) -> Result<()> {
1098 self.conn.execute_batch(sql_builders::BEGIN)?;
1099 Ok(())
1100 }
1101
1102 pub fn commit(&self) -> Result<()> {
1104 self.conn.execute_batch(sql_builders::COMMIT)?;
1105 Ok(())
1106 }
1107
1108 pub fn rollback(&self) -> Result<()> {
1110 self.conn.execute_batch(sql_builders::ROLLBACK)?;
1111 Ok(())
1112 }
1113
1114 pub fn enable_bulk_loading(&self) -> Result<()> {
1120 self.conn.execute_batch(
1121 "PRAGMA synchronous = OFF;
1122 PRAGMA cache_size = -64000;
1123 PRAGMA temp_store = MEMORY;
1124 PRAGMA mmap_size = 268435456;",
1125 )?;
1126 Ok(())
1127 }
1128
1129 pub fn disable_bulk_loading(&self) -> Result<()> {
1131 self.conn.execute_batch(
1132 "PRAGMA synchronous = NORMAL;
1133 PRAGMA cache_size = -2000;
1134 PRAGMA temp_store = DEFAULT;
1135 PRAGMA mmap_size = 0;",
1136 )?;
1137 Ok(())
1138 }
1139
1140 pub fn put_item(
1146 &self,
1147 table_name: &str,
1148 pk: &str,
1149 sk: &str,
1150 item_json: &str,
1151 item_size: usize,
1152 ) -> Result<Option<String>> {
1153 self.put_item_with_hash(table_name, pk, sk, item_json, item_size, "")
1154 }
1155
1156 pub fn put_item_with_hash(
1158 &self,
1159 table_name: &str,
1160 pk: &str,
1161 sk: &str,
1162 item_json: &str,
1163 item_size: usize,
1164 hash_prefix: &str,
1165 ) -> Result<Option<String>> {
1166 let old_item = self.get_item(table_name, pk, sk)?;
1168
1169 let (sql, params) =
1170 sql_builders::put_item_with_hash(table_name, pk, sk, item_json, item_size, hash_prefix);
1171 self.conn
1172 .execute(&sql, rusqlite::params_from_iter(params.iter()))?;
1173
1174 Ok(old_item)
1175 }
1176
1177 pub fn put_base_items(
1182 &self,
1183 table_name: &str,
1184 rows: &[crate::storage_backend::BaseItemRow],
1185 ) -> Result<()> {
1186 let escaped = escape_table_name(table_name);
1187 let sql = format!(
1188 "INSERT OR REPLACE INTO \"{escaped}\" (pk, sk, item_json, item_size, cached_at, hash_prefix) \
1189 VALUES (?1, ?2, ?3, ?4, ?5, ?6)"
1190 );
1191 let mut stmt = self.conn.prepare_cached(&sql)?;
1192 for row in rows {
1193 stmt.execute(params![
1194 row.pk,
1195 row.sk,
1196 row.item_json,
1197 row.item_size as i64,
1198 row.cached_at,
1199 row.hash_prefix
1200 ])?;
1201 }
1202 Ok(())
1203 }
1204
1205 pub fn get_item(&self, table_name: &str, pk: &str, sk: &str) -> Result<Option<String>> {
1207 let (sql, params) = sql_builders::get_item(table_name, pk, sk);
1208 let result = self
1209 .conn
1210 .query_row(&sql, rusqlite::params_from_iter(params.iter()), |row| {
1211 row.get(0)
1212 });
1213
1214 match result {
1215 Ok(json) => Ok(Some(json)),
1216 Err(rusqlite::Error::QueryReturnedNoRows) => Ok(None),
1217 Err(e) => Err(DynoxideError::from(e)),
1218 }
1219 }
1220
1221 pub fn get_partition_size(&self, table_name: &str, pk: &str) -> Result<i64> {
1223 let (sql, params) = sql_builders::get_partition_size(table_name, pk);
1224 let size: i64 =
1225 self.conn
1226 .query_row(&sql, rusqlite::params_from_iter(params.iter()), |row| {
1227 row.get(0)
1228 })?;
1229 Ok(size)
1230 }
1231
1232 pub fn get_lsi_partition_size(
1235 &self,
1236 table_name: &str,
1237 index_name: &str,
1238 pk: &str,
1239 ) -> Result<i64> {
1240 let (sql, params) = sql_builders::get_lsi_partition_size(table_name, index_name, pk);
1241 let size: i64 =
1242 self.conn
1243 .query_row(&sql, rusqlite::params_from_iter(params.iter()), |row| {
1244 row.get(0)
1245 })?;
1246 Ok(size)
1247 }
1248
1249 pub fn delete_item(&self, table_name: &str, pk: &str, sk: &str) -> Result<Option<String>> {
1251 let old_item = self.get_item(table_name, pk, sk)?;
1252
1253 let (sql, params) = sql_builders::delete_item(table_name, pk, sk);
1254 self.conn
1255 .execute(&sql, rusqlite::params_from_iter(params.iter()))?;
1256
1257 Ok(old_item)
1258 }
1259
1260 pub fn query_items(
1265 &self,
1266 table_name: &str,
1267 pk: &str,
1268 params: &QueryParams,
1269 ) -> Result<Vec<(String, String, String)>> {
1270 let (sql, params_vec) = sql_builders::query_items(table_name, pk, params);
1271 let mut stmt = self.conn.prepare(&sql)?;
1272 let rows = stmt
1273 .query_map(rusqlite::params_from_iter(params_vec.iter()), |row| {
1274 Ok((row.get(0)?, row.get(1)?, row.get(2)?))
1275 })?
1276 .collect::<std::result::Result<Vec<_>, _>>()?;
1277
1278 Ok(rows)
1279 }
1280
1281 pub fn scan_items(
1286 &self,
1287 table_name: &str,
1288 params: &ScanParams,
1289 ) -> Result<Vec<(String, String, String)>> {
1290 let (sql, params_vec) = sql_builders::scan_items(table_name, params);
1291 let mut stmt = self.conn.prepare(&sql)?;
1292 let rows = stmt
1293 .query_map(rusqlite::params_from_iter(params_vec.iter()), |row| {
1294 Ok((row.get(0)?, row.get(1)?, row.get(2)?))
1295 })?
1296 .collect::<std::result::Result<Vec<_>, _>>()?;
1297
1298 Ok(rows)
1299 }
1300
1301 pub fn count_items(&self, table_name: &str) -> Result<i64> {
1303 let (sql, params) = sql_builders::count_items(table_name);
1304 let count: i64 =
1305 self.conn
1306 .query_row(&sql, rusqlite::params_from_iter(params.iter()), |row| {
1307 row.get(0)
1308 })?;
1309 Ok(count)
1310 }
1311
1312 pub fn db_path(&self) -> Option<String> {
1318 self.conn
1319 .path()
1320 .filter(|p| !p.is_empty())
1321 .map(|p| p.to_owned())
1322 }
1323
1324 pub fn db_size_bytes(&self) -> Result<u64> {
1326 let size: i64 = self.conn.query_row(
1327 "SELECT page_count * page_size FROM pragma_page_count(), pragma_page_size()",
1328 [],
1329 |row| row.get(0),
1330 )?;
1331 Ok(size as u64)
1332 }
1333
1334 pub fn table_count(&self) -> Result<usize> {
1336 let count: i64 = self
1337 .conn
1338 .query_row("SELECT COUNT(*) FROM _tables", [], |row| row.get(0))?;
1339 Ok(count as usize)
1340 }
1341
1342 pub fn table_stats(&self) -> Result<Vec<TableStats>> {
1346 let table_names = self.list_table_names()?;
1347 let mut stats = Vec::with_capacity(table_names.len());
1348 for name in table_names {
1349 let sql = format!(
1350 "SELECT COUNT(*), COALESCE(SUM(item_size), 0) FROM \"{}\"",
1351 escape_table_name(&name)
1352 );
1353 let (item_count, size_bytes): (i64, i64) = self
1354 .conn
1355 .query_row(&sql, [], |row| Ok((row.get(0)?, row.get(1)?)))?;
1356 stats.push(TableStats {
1357 table_name: name,
1358 item_count,
1359 size_bytes: size_bytes as u64,
1360 });
1361 }
1362 Ok(stats)
1363 }
1364
1365 pub fn database_info(&self) -> Result<DatabaseInfo> {
1370 let path = self.db_path();
1371 let size_bytes = self.db_size_bytes()?;
1372 let table_count = self.table_count()?;
1373 let stats = self.table_stats()?;
1374
1375 let mut table_details = Vec::with_capacity(stats.len());
1376 for s in stats {
1377 let metadata = self.get_table_metadata(&s.table_name)?;
1378 table_details.push(TableInfoEntry { stats: s, metadata });
1379 }
1380
1381 Ok(DatabaseInfo {
1382 path,
1383 size_bytes,
1384 table_count,
1385 tables: table_details,
1386 })
1387 }
1388
1389 pub fn vacuum_into(&self, path: &str) -> Result<()> {
1396 if path.contains('\0') {
1397 return Err(DynoxideError::ValidationException(
1398 "path contains null byte".to_string(),
1399 ));
1400 }
1401 self.conn
1402 .execute_batch(&format!("VACUUM INTO '{}'", path.replace('\'', "''")))?;
1403 Ok(())
1404 }
1405
1406 pub fn vacuum(&self) -> Result<()> {
1408 self.conn.execute_batch("VACUUM")?;
1409 Ok(())
1410 }
1411
1412 pub fn restore_from(&mut self, path: &str) -> Result<()> {
1416 let source = Connection::open(path)?;
1417 self.restore_from_connection(&source)
1418 }
1419
1420 pub fn backup_to_memory(&self) -> Result<Connection> {
1425 let mut dest = Connection::open_in_memory()?;
1426 {
1427 let backup = rusqlite::backup::Backup::new(&self.conn, &mut dest)?;
1428 backup.run_to_completion(100, std::time::Duration::from_millis(0), None)?;
1429 }
1430 Ok(dest)
1431 }
1432
1433 pub fn restore_from_connection(&mut self, source: &Connection) -> Result<()> {
1438 let backup = rusqlite::backup::Backup::new(source, &mut self.conn)?;
1439 backup.run_to_completion(100, std::time::Duration::from_millis(0), None)?;
1440 self.metadata_cache.borrow_mut().clear();
1441 Ok(())
1442 }
1443
1444 pub fn connection_size_bytes(conn: &Connection) -> Result<u64> {
1446 let size: i64 = conn.query_row(
1447 "SELECT page_count * page_size FROM pragma_page_count(), pragma_page_size()",
1448 [],
1449 |row| row.get(0),
1450 )?;
1451 Ok(size as u64)
1452 }
1453
1454 pub fn enable_stream(&self, table_name: &str, view_type: &str, label: &str) -> Result<()> {
1460 self.conn.execute(
1461 "UPDATE _tables SET stream_enabled = 1, stream_view_type = ?1, stream_label = ?2 WHERE table_name = ?3",
1462 params![view_type, label, table_name],
1463 )?;
1464 self.metadata_cache.borrow_mut().remove(table_name);
1465 Ok(())
1466 }
1467
1468 pub fn disable_stream(&self, table_name: &str) -> Result<()> {
1470 self.conn.execute(
1471 "UPDATE _tables SET stream_enabled = 0 WHERE table_name = ?1",
1472 params![table_name],
1473 )?;
1474 self.metadata_cache.borrow_mut().remove(table_name);
1475 Ok(())
1476 }
1477
1478 #[allow(clippy::too_many_arguments)]
1480 pub fn insert_stream_record(
1481 &self,
1482 table_name: &str,
1483 event_name: &str,
1484 keys_json: &str,
1485 new_image: Option<&str>,
1486 old_image: Option<&str>,
1487 sequence_number: &str,
1488 shard_id: &str,
1489 created_at: i64,
1490 ) -> Result<()> {
1491 self.insert_stream_record_with_identity(
1492 table_name,
1493 event_name,
1494 keys_json,
1495 new_image,
1496 old_image,
1497 sequence_number,
1498 shard_id,
1499 created_at,
1500 None,
1501 )
1502 }
1503
1504 #[allow(clippy::too_many_arguments)]
1506 pub fn insert_stream_record_with_identity(
1507 &self,
1508 table_name: &str,
1509 event_name: &str,
1510 keys_json: &str,
1511 new_image: Option<&str>,
1512 old_image: Option<&str>,
1513 sequence_number: &str,
1514 shard_id: &str,
1515 created_at: i64,
1516 user_identity: Option<&str>,
1517 ) -> Result<()> {
1518 self.conn.execute(
1519 "INSERT INTO _stream_records (table_name, event_name, keys_json, new_image, old_image, sequence_number, shard_id, created_at, user_identity)
1520 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)",
1521 params![table_name, event_name, keys_json, new_image, old_image, sequence_number, shard_id, created_at, user_identity],
1522 )?;
1523 Ok(())
1524 }
1525
1526 pub fn next_stream_sequence_number(&self, table_name: &str) -> Result<i64> {
1528 let result: std::result::Result<i64, _> = self.conn.query_row(
1529 "SELECT COALESCE(MAX(CAST(sequence_number AS INTEGER)), 0) + 1 FROM _stream_records WHERE table_name = ?1",
1530 params![table_name],
1531 |row| row.get(0),
1532 );
1533 match result {
1534 Ok(n) => Ok(n),
1535 Err(_) => Ok(1),
1536 }
1537 }
1538
1539 pub fn get_stream_records(
1541 &self,
1542 table_name: &str,
1543 shard_id: &str,
1544 after_sequence: i64,
1545 limit: usize,
1546 ) -> Result<Vec<StreamRecord>> {
1547 let mut stmt = self.conn.prepare(
1548 "SELECT event_name, keys_json, new_image, old_image, sequence_number, created_at, user_identity
1549 FROM _stream_records
1550 WHERE table_name = ?1 AND shard_id = ?2 AND CAST(sequence_number AS INTEGER) > ?3
1551 ORDER BY CAST(sequence_number AS INTEGER) ASC
1552 LIMIT ?4",
1553 )?;
1554 let rows = stmt
1555 .query_map(
1556 params![table_name, shard_id, after_sequence, limit as i64],
1557 |row| {
1558 Ok(StreamRecord {
1559 event_name: row.get(0)?,
1560 keys_json: row.get(1)?,
1561 new_image: row.get(2)?,
1562 old_image: row.get(3)?,
1563 sequence_number: row.get(4)?,
1564 created_at: row.get(5)?,
1565 user_identity: row.get(6)?,
1566 })
1567 },
1568 )?
1569 .collect::<std::result::Result<Vec<_>, _>>()?;
1570 Ok(rows)
1571 }
1572
1573 pub fn list_stream_enabled_tables(&self) -> Result<Vec<TableMetadata>> {
1575 let sql = format!(
1576 "SELECT {} FROM _tables WHERE stream_enabled = 1 ORDER BY table_name",
1577 sql_builders::TABLE_METADATA_COLUMNS
1578 );
1579 let mut stmt = self.conn.prepare(&sql)?;
1580 let rows = stmt
1581 .query_map([], row_to_metadata)?
1582 .collect::<std::result::Result<Vec<_>, _>>()?;
1583 Ok(rows)
1584 }
1585
1586 pub fn update_ttl_config(
1592 &self,
1593 table_name: &str,
1594 attribute_name: Option<&str>,
1595 enabled: bool,
1596 ) -> Result<()> {
1597 self.conn.execute(
1598 "UPDATE _tables SET ttl_attribute = ?1, ttl_enabled = ?2 WHERE table_name = ?3",
1599 params![attribute_name, enabled as i32, table_name],
1600 )?;
1601 self.metadata_cache.borrow_mut().remove(table_name);
1602 Ok(())
1603 }
1604
1605 pub fn list_ttl_enabled_tables(&self) -> Result<Vec<TableMetadata>> {
1607 let sql = format!(
1608 "SELECT {} FROM _tables WHERE ttl_enabled = 1 ORDER BY table_name",
1609 sql_builders::TABLE_METADATA_COLUMNS
1610 );
1611 let mut stmt = self.conn.prepare(&sql)?;
1612 let rows = stmt
1613 .query_map([], row_to_metadata)?
1614 .collect::<std::result::Result<Vec<_>, _>>()?;
1615 Ok(rows)
1616 }
1617
1618 pub fn get_shard_sequence_range(
1620 &self,
1621 table_name: &str,
1622 shard_id: &str,
1623 ) -> Result<(Option<String>, Option<String>)> {
1624 let result: std::result::Result<(Option<String>, Option<String>), _> = self.conn.query_row(
1625 "SELECT MIN(sequence_number), MAX(sequence_number) FROM _stream_records WHERE table_name = ?1 AND shard_id = ?2",
1626 params![table_name, shard_id],
1627 |row| Ok((row.get(0)?, row.get(1)?)),
1628 );
1629 match result {
1630 Ok(range) => Ok(range),
1631 Err(_) => Ok((None, None)),
1632 }
1633 }
1634
1635 pub fn touch_cached_at(
1641 &self,
1642 table_name: &str,
1643 pk: &str,
1644 sk: &str,
1645 timestamp: f64,
1646 ) -> Result<()> {
1647 let sql = format!(
1648 "UPDATE \"{}\" SET cached_at = ?1 WHERE pk = ?2 AND sk = ?3",
1649 escape_table_name(table_name)
1650 );
1651 self.conn.execute(&sql, params![timestamp, pk, sk])?;
1652 Ok(())
1653 }
1654
1655 pub fn get_lru_items(
1660 &self,
1661 table_name: &str,
1662 limit: usize,
1663 ) -> Result<Vec<(String, String, i64)>> {
1664 let sql = format!(
1665 "SELECT pk, sk, item_size FROM \"{}\" WHERE cached_at IS NOT NULL ORDER BY cached_at ASC LIMIT ?1",
1666 escape_table_name(table_name)
1667 );
1668 let mut stmt = self.conn.prepare(&sql)?;
1669 let rows = stmt
1670 .query_map(params![limit as i64], |row| {
1671 Ok((row.get(0)?, row.get(1)?, row.get(2)?))
1672 })?
1673 .collect::<std::result::Result<Vec<_>, _>>()?;
1674 Ok(rows)
1675 }
1676}
1677
1678#[derive(Debug, Clone)]
1680pub struct StreamRecord {
1681 pub event_name: String,
1682 pub keys_json: String,
1683 pub new_image: Option<String>,
1684 pub old_image: Option<String>,
1685 pub sequence_number: String,
1686 pub created_at: i64,
1687 pub user_identity: Option<String>,
1688}
1689
1690#[derive(Debug, Clone)]
1692pub struct TableStats {
1693 pub table_name: String,
1694 pub item_count: i64,
1695 pub size_bytes: u64,
1696}
1697
1698#[derive(Debug, Clone)]
1700pub struct DatabaseInfo {
1701 pub path: Option<String>,
1702 pub size_bytes: u64,
1703 pub table_count: usize,
1704 pub tables: Vec<TableInfoEntry>,
1705}
1706
1707#[derive(Debug, Clone)]
1709pub struct TableInfoEntry {
1710 pub stats: TableStats,
1711 pub metadata: Option<TableMetadata>,
1712}
1713
1714#[derive(Debug, Clone)]
1720pub struct TableMetadata {
1721 pub table_name: String,
1722 pub key_schema: String,
1723 pub attribute_definitions: String,
1724 pub gsi_definitions: Option<String>,
1725 pub lsi_definitions: Option<String>,
1726 pub stream_enabled: bool,
1727 pub stream_view_type: Option<String>,
1728 pub stream_label: Option<String>,
1729 pub ttl_attribute: Option<String>,
1730 pub ttl_enabled: bool,
1731 pub created_at: i64,
1732 pub table_status: String,
1733 pub billing_mode: Option<String>,
1734 pub provisioned_throughput: Option<String>,
1735 pub sse_specification: Option<String>,
1736 pub table_class: Option<String>,
1737 pub deletion_protection_enabled: bool,
1738 pub on_demand_throughput: Option<String>,
1739 pub table_id: Option<String>,
1742}
1743
1744#[cfg(any(feature = "native-sqlite", feature = "_has-encryption"))]
1746fn row_to_metadata(row: &rusqlite::Row) -> rusqlite::Result<TableMetadata> {
1747 Ok(TableMetadata {
1748 table_name: row.get(0)?,
1749 key_schema: row.get(1)?,
1750 attribute_definitions: row.get(2)?,
1751 gsi_definitions: row.get(3)?,
1752 lsi_definitions: row.get(4)?,
1753 stream_enabled: row.get::<_, i32>(5)? != 0,
1754 stream_view_type: row.get(6)?,
1755 stream_label: row.get(7)?,
1756 ttl_attribute: row.get(8)?,
1757 ttl_enabled: row.get::<_, i32>(9)? != 0,
1758 created_at: row.get(10)?,
1759 table_status: row.get(11)?,
1760 billing_mode: row.get(12)?,
1761 provisioned_throughput: row.get(13)?,
1762 sse_specification: row.get(14)?,
1763 table_class: row.get(15)?,
1764 deletion_protection_enabled: row.get::<_, i32>(16).unwrap_or(0) != 0,
1765 on_demand_throughput: row.get(17)?,
1766 table_id: row.get(18)?,
1767 })
1768}
1769
1770#[cfg(all(test, any(feature = "native-sqlite", feature = "_has-encryption")))]
1771mod tests {
1772 use super::*;
1773
1774 fn test_storage() -> Storage {
1775 Storage::memory().expect("Failed to create in-memory storage")
1776 }
1777
1778 #[test]
1779 fn test_initialize_creates_metadata_tables() {
1780 let storage = test_storage();
1781 let version: String = storage
1783 .conn()
1784 .query_row(
1785 "SELECT value FROM _config WHERE key = 'schema_version'",
1786 [],
1787 |row| row.get(0),
1788 )
1789 .unwrap();
1790 assert_eq!(version, SCHEMA_VERSION);
1791 }
1792
1793 #[test]
1794 fn test_migrate_v6_to_v7_adds_on_demand_throughput_column() {
1795 let tmp = tempfile::NamedTempFile::new().unwrap();
1799 let path = tmp.path().to_str().unwrap().to_string();
1800
1801 {
1804 let conn = Connection::open(&path).unwrap();
1805 conn.execute_batch(
1806 "CREATE TABLE _config (key TEXT PRIMARY KEY, value TEXT NOT NULL);
1807 CREATE TABLE _tables (
1808 table_name TEXT PRIMARY KEY,
1809 key_schema TEXT NOT NULL,
1810 attribute_definitions TEXT NOT NULL,
1811 gsi_definitions TEXT,
1812 lsi_definitions TEXT,
1813 stream_enabled INTEGER DEFAULT 0,
1814 stream_view_type TEXT,
1815 stream_label TEXT,
1816 ttl_attribute TEXT,
1817 ttl_enabled INTEGER DEFAULT 0,
1818 created_at INTEGER NOT NULL,
1819 table_status TEXT NOT NULL DEFAULT 'ACTIVE',
1820 billing_mode TEXT DEFAULT 'PAY_PER_REQUEST',
1821 provisioned_throughput TEXT,
1822 tags TEXT,
1823 sse_specification TEXT,
1824 table_class TEXT,
1825 deletion_protection_enabled INTEGER DEFAULT 0
1826 );",
1827 )
1828 .unwrap();
1829 conn.execute(
1830 "INSERT INTO _config (key, value) VALUES ('schema_version', '6')",
1831 [],
1832 )
1833 .unwrap();
1834 conn.execute(
1835 "INSERT INTO _tables (table_name, key_schema, attribute_definitions, created_at) \
1836 VALUES ('LegacyTable', ?1, ?2, 0)",
1837 params![
1838 r#"[{"AttributeName":"pk","KeyType":"HASH"}]"#,
1839 r#"[{"AttributeName":"pk","AttributeType":"S"}]"#,
1840 ],
1841 )
1842 .unwrap();
1843 }
1844
1845 let storage = Storage::new(&path).unwrap();
1848 let version: String = storage
1849 .conn()
1850 .query_row(
1851 "SELECT value FROM _config WHERE key = 'schema_version'",
1852 [],
1853 |r| r.get(0),
1854 )
1855 .unwrap();
1856 assert_eq!(version, SCHEMA_VERSION);
1857
1858 let meta = storage.get_table_metadata("LegacyTable").unwrap().unwrap();
1861 assert_eq!(meta.table_name, "LegacyTable");
1862 assert!(meta.on_demand_throughput.is_none());
1863
1864 let col: Option<String> = storage
1866 .conn()
1867 .query_row(
1868 "SELECT on_demand_throughput FROM _tables WHERE table_name = 'LegacyTable'",
1869 [],
1870 |r| r.get(0),
1871 )
1872 .unwrap();
1873 assert!(col.is_none());
1874 }
1875
1876 #[test]
1880 fn test_migrate_v7_to_v8_backfills_table_id() {
1881 let tmp = tempfile::NamedTempFile::new().unwrap();
1882 let path = tmp.path().to_str().unwrap().to_string();
1883
1884 {
1887 let conn = Connection::open(&path).unwrap();
1888 conn.execute_batch(
1889 "CREATE TABLE _config (key TEXT PRIMARY KEY, value TEXT NOT NULL);
1890 CREATE TABLE _tables (
1891 table_name TEXT PRIMARY KEY,
1892 key_schema TEXT NOT NULL,
1893 attribute_definitions TEXT NOT NULL,
1894 gsi_definitions TEXT,
1895 lsi_definitions TEXT,
1896 stream_enabled INTEGER DEFAULT 0,
1897 stream_view_type TEXT,
1898 stream_label TEXT,
1899 ttl_attribute TEXT,
1900 ttl_enabled INTEGER DEFAULT 0,
1901 created_at INTEGER NOT NULL,
1902 table_status TEXT NOT NULL DEFAULT 'ACTIVE',
1903 billing_mode TEXT DEFAULT 'PAY_PER_REQUEST',
1904 provisioned_throughput TEXT,
1905 tags TEXT,
1906 sse_specification TEXT,
1907 table_class TEXT,
1908 deletion_protection_enabled INTEGER DEFAULT 0,
1909 on_demand_throughput TEXT
1910 );",
1911 )
1912 .unwrap();
1913 conn.execute(
1914 "INSERT INTO _config (key, value) VALUES ('schema_version', '7')",
1915 [],
1916 )
1917 .unwrap();
1918 conn.execute(
1919 "INSERT INTO _tables (table_name, key_schema, attribute_definitions, created_at) \
1920 VALUES ('LegacyTable', ?1, ?2, 0)",
1921 params![
1922 r#"[{"AttributeName":"pk","KeyType":"HASH"}]"#,
1923 r#"[{"AttributeName":"pk","AttributeType":"S"}]"#,
1924 ],
1925 )
1926 .unwrap();
1927 }
1928
1929 let storage = Storage::new(&path).unwrap();
1931 let version: String = storage
1932 .conn()
1933 .query_row(
1934 "SELECT value FROM _config WHERE key = 'schema_version'",
1935 [],
1936 |r| r.get(0),
1937 )
1938 .unwrap();
1939 assert_eq!(version, SCHEMA_VERSION);
1940
1941 let meta = storage.get_table_metadata("LegacyTable").unwrap().unwrap();
1944 let id = meta.table_id.expect("legacy table should be backfilled");
1945 assert!(!id.is_empty());
1946
1947 drop(storage);
1949 let storage2 = Storage::new(&path).unwrap();
1950 let meta2 = storage2.get_table_metadata("LegacyTable").unwrap().unwrap();
1951 assert_eq!(meta2.table_id.as_deref(), Some(id.as_str()));
1952 }
1953
1954 #[test]
1955 fn test_wal_mode_enabled() {
1956 let storage = test_storage();
1957 let mode: String = storage
1958 .conn()
1959 .query_row("PRAGMA journal_mode", [], |row| row.get(0))
1960 .unwrap();
1961 assert!(mode == "wal" || mode == "memory", "Got mode: {mode}");
1963 }
1964
1965 #[test]
1971 fn fnv1a_hash_matches_known_vectors() {
1972 let storage = test_storage();
1973 let cases: [(&str, i64); 6] = [
1974 ("", 2166136261),
1975 ("a", 3826002220),
1976 ("u#1", 2199603432),
1977 ("artist#42", 2385694177),
1978 ("café", 2821410889),
1979 ("tenant#9007199254740993", 2022216178),
1980 ];
1981 for (input, expected) in cases {
1982 let got: i64 = storage
1983 .conn()
1984 .query_row("SELECT fnv1a_hash(?1)", [input], |row| row.get(0))
1985 .unwrap();
1986 assert_eq!(got, expected, "fnv1a_hash({input:?})");
1987 }
1988 }
1989
1990 #[test]
1991 fn test_table_metadata_crud() {
1992 let storage = test_storage();
1993
1994 assert!(!storage.table_exists("TestTable").unwrap());
1996 assert!(storage.list_table_names().unwrap().is_empty());
1997
1998 storage
2000 .insert_table_metadata(&CreateTableMetadata {
2001 table_name: "TestTable",
2002 key_schema: r#"[{"AttributeName":"pk","KeyType":"HASH"}]"#,
2003 attribute_definitions: r#"[{"AttributeName":"pk","AttributeType":"S"}]"#,
2004 created_at: 1000000,
2005 ..Default::default()
2006 })
2007 .unwrap();
2008
2009 assert!(storage.table_exists("TestTable").unwrap());
2010 assert_eq!(storage.list_table_names().unwrap(), vec!["TestTable"]);
2011
2012 let meta = storage.get_table_metadata("TestTable").unwrap().unwrap();
2014 assert_eq!(meta.table_name, "TestTable");
2015 assert_eq!(meta.table_status, "ACTIVE");
2016 assert_eq!(meta.created_at, 1000000);
2017
2018 assert!(storage.delete_table_metadata("TestTable").unwrap());
2020 assert!(!storage.table_exists("TestTable").unwrap());
2021 }
2022
2023 #[test]
2024 fn test_create_and_drop_data_table() {
2025 let storage = test_storage();
2026 storage.create_data_table("MyTable").unwrap();
2027
2028 storage
2030 .put_item("MyTable", "pk1", "", r#"{"pk":{"S":"pk1"}}"#, 10)
2031 .unwrap();
2032
2033 let item = storage.get_item("MyTable", "pk1", "").unwrap();
2034 assert!(item.is_some());
2035
2036 storage.drop_data_table("MyTable").unwrap();
2037 }
2038
2039 #[test]
2040 fn test_item_crud() {
2041 let storage = test_storage();
2042 storage.create_data_table("Items").unwrap();
2043
2044 let old = storage
2046 .put_item(
2047 "Items",
2048 "user#1",
2049 "profile",
2050 r#"{"name":{"S":"Alice"}}"#,
2051 20,
2052 )
2053 .unwrap();
2054 assert!(old.is_none()); let item = storage.get_item("Items", "user#1", "profile").unwrap();
2058 assert_eq!(item.unwrap(), r#"{"name":{"S":"Alice"}}"#);
2059
2060 let old = storage
2062 .put_item("Items", "user#1", "profile", r#"{"name":{"S":"Bob"}}"#, 18)
2063 .unwrap();
2064 assert_eq!(old.unwrap(), r#"{"name":{"S":"Alice"}}"#);
2065
2066 let deleted = storage.delete_item("Items", "user#1", "profile").unwrap();
2068 assert_eq!(deleted.unwrap(), r#"{"name":{"S":"Bob"}}"#);
2069
2070 assert!(
2072 storage
2073 .get_item("Items", "user#1", "profile")
2074 .unwrap()
2075 .is_none()
2076 );
2077 }
2078
2079 #[test]
2080 fn test_query_items() {
2081 let storage = test_storage();
2082 storage.create_data_table("Orders").unwrap();
2083
2084 for i in 1..=5 {
2086 let sk = format!("order#{i:03}");
2087 let json = format!(r#"{{"id":{{"N":"{i}"}}}}"#);
2088 storage
2089 .put_item("Orders", "user#1", &sk, &json, 10)
2090 .unwrap();
2091 }
2092
2093 let results = storage
2095 .query_items(
2096 "Orders",
2097 "user#1",
2098 &QueryParams {
2099 forward: true,
2100 ..Default::default()
2101 },
2102 )
2103 .unwrap();
2104 assert_eq!(results.len(), 5);
2105 assert_eq!(results[0].1, "order#001"); let results = storage
2109 .query_items(
2110 "Orders",
2111 "user#1",
2112 &QueryParams {
2113 forward: true,
2114 limit: Some(2),
2115 ..Default::default()
2116 },
2117 )
2118 .unwrap();
2119 assert_eq!(results.len(), 2);
2120
2121 let results = storage
2123 .query_items(
2124 "Orders",
2125 "user#1",
2126 &QueryParams {
2127 forward: false,
2128 limit: Some(2),
2129 ..Default::default()
2130 },
2131 )
2132 .unwrap();
2133 assert_eq!(results.len(), 2);
2134 assert_eq!(results[0].1, "order#005"); }
2136
2137 #[test]
2138 fn test_scan_items() {
2139 let storage = test_storage();
2140 storage.create_data_table("ScanTest").unwrap();
2141
2142 storage.put_item("ScanTest", "a", "1", r#"{}"#, 2).unwrap();
2143 storage.put_item("ScanTest", "b", "2", r#"{}"#, 2).unwrap();
2144 storage.put_item("ScanTest", "c", "3", r#"{}"#, 2).unwrap();
2145
2146 let results = storage.scan_items("ScanTest", &Default::default()).unwrap();
2147 assert_eq!(results.len(), 3);
2148
2149 let results = storage
2151 .scan_items(
2152 "ScanTest",
2153 &ScanParams {
2154 limit: Some(2),
2155 ..Default::default()
2156 },
2157 )
2158 .unwrap();
2159 assert_eq!(results.len(), 2);
2160
2161 let results = storage
2163 .scan_items(
2164 "ScanTest",
2165 &ScanParams {
2166 limit: Some(2),
2167 exclusive_start_pk: Some("a"),
2168 exclusive_start_sk: Some("1"),
2169 ..Default::default()
2170 },
2171 )
2172 .unwrap();
2173 assert_eq!(results.len(), 2);
2174 assert_eq!(results[0].0, "b"); }
2176
2177 #[test]
2178 fn test_count_items() {
2179 let storage = test_storage();
2180 storage.create_data_table("CountTest").unwrap();
2181
2182 assert_eq!(storage.count_items("CountTest").unwrap(), 0);
2183
2184 storage.put_item("CountTest", "a", "", r#"{}"#, 2).unwrap();
2185 storage.put_item("CountTest", "b", "", r#"{}"#, 2).unwrap();
2186
2187 assert_eq!(storage.count_items("CountTest").unwrap(), 2);
2188 }
2189
2190 #[test]
2191 fn test_gsi_table_lifecycle() {
2192 let storage = test_storage();
2193 storage.create_gsi_table("Orders", "ByDate").unwrap();
2194
2195 let gsi_name = "Orders::gsi::ByDate";
2197 let sql = format!(
2198 "INSERT INTO \"{}\" (gsi_pk, gsi_sk, table_pk, table_sk, item_json) VALUES (?1, ?2, ?3, ?4, ?5)",
2199 gsi_name.replace('"', "\"\"")
2200 );
2201 storage
2202 .conn()
2203 .execute(
2204 &sql,
2205 params!["2024-01-01", "001", "user#1", "order#001", r#"{}"#],
2206 )
2207 .unwrap();
2208
2209 storage.drop_gsi_table("Orders", "ByDate").unwrap();
2210 }
2211
2212 #[test]
2213 fn test_nonexistent_table_metadata() {
2214 let storage = test_storage();
2215 assert!(storage.get_table_metadata("Nonexistent").unwrap().is_none());
2216 assert!(!storage.delete_table_metadata("Nonexistent").unwrap());
2217 }
2218
2219 #[test]
2220 fn test_metadata_cache_hit() {
2221 let storage = test_storage();
2222 storage
2223 .insert_table_metadata(&CreateTableMetadata {
2224 table_name: "CacheTest",
2225 key_schema: r#"[{"AttributeName":"pk","KeyType":"HASH"}]"#,
2226 attribute_definitions: r#"[{"AttributeName":"pk","AttributeType":"S"}]"#,
2227 created_at: 1000000,
2228 ..Default::default()
2229 })
2230 .unwrap();
2231
2232 let meta1 = storage.get_table_metadata("CacheTest").unwrap().unwrap();
2234 assert_eq!(meta1.table_name, "CacheTest");
2235
2236 let meta2 = storage.get_table_metadata("CacheTest").unwrap().unwrap();
2238 assert_eq!(meta2.table_name, "CacheTest");
2239 assert_eq!(meta1.created_at, meta2.created_at);
2240
2241 assert!(storage.metadata_cache.borrow().contains_key("CacheTest"));
2243 }
2244
2245 #[test]
2246 fn test_metadata_cache_invalidated_on_delete() {
2247 let storage = test_storage();
2248 storage
2249 .insert_table_metadata(&CreateTableMetadata {
2250 table_name: "DelCache",
2251 key_schema: r#"[{"AttributeName":"pk","KeyType":"HASH"}]"#,
2252 attribute_definitions: r#"[{"AttributeName":"pk","AttributeType":"S"}]"#,
2253 created_at: 1000000,
2254 ..Default::default()
2255 })
2256 .unwrap();
2257
2258 storage.get_table_metadata("DelCache").unwrap();
2260 assert!(storage.metadata_cache.borrow().contains_key("DelCache"));
2261
2262 storage.delete_table_metadata("DelCache").unwrap();
2264 assert!(!storage.metadata_cache.borrow().contains_key("DelCache"));
2265 }
2266
2267 #[test]
2268 fn test_metadata_cache_invalidated_on_stream_enable() {
2269 let storage = test_storage();
2270 storage
2271 .insert_table_metadata(&CreateTableMetadata {
2272 table_name: "StreamCache",
2273 key_schema: r#"[{"AttributeName":"pk","KeyType":"HASH"}]"#,
2274 attribute_definitions: r#"[{"AttributeName":"pk","AttributeType":"S"}]"#,
2275 created_at: 1000000,
2276 ..Default::default()
2277 })
2278 .unwrap();
2279
2280 let meta = storage.get_table_metadata("StreamCache").unwrap().unwrap();
2282 assert!(!meta.stream_enabled);
2283
2284 storage
2286 .enable_stream("StreamCache", "NEW_AND_OLD_IMAGES", "2024-01-01T00:00:00")
2287 .unwrap();
2288 assert!(!storage.metadata_cache.borrow().contains_key("StreamCache"));
2289
2290 let meta = storage.get_table_metadata("StreamCache").unwrap().unwrap();
2292 assert!(meta.stream_enabled);
2293 }
2294
2295 #[test]
2296 fn test_metadata_cache_invalidated_on_ttl_update() {
2297 let storage = test_storage();
2298 storage
2299 .insert_table_metadata(&CreateTableMetadata {
2300 table_name: "TtlCache",
2301 key_schema: r#"[{"AttributeName":"pk","KeyType":"HASH"}]"#,
2302 attribute_definitions: r#"[{"AttributeName":"pk","AttributeType":"S"}]"#,
2303 created_at: 1000000,
2304 ..Default::default()
2305 })
2306 .unwrap();
2307
2308 let meta = storage.get_table_metadata("TtlCache").unwrap().unwrap();
2310 assert!(!meta.ttl_enabled);
2311
2312 storage
2314 .update_ttl_config("TtlCache", Some("expires_at"), true)
2315 .unwrap();
2316 assert!(!storage.metadata_cache.borrow().contains_key("TtlCache"));
2317
2318 let meta = storage.get_table_metadata("TtlCache").unwrap().unwrap();
2320 assert!(meta.ttl_enabled);
2321 assert_eq!(meta.ttl_attribute, Some("expires_at".to_string()));
2322 }
2323
2324 #[test]
2325 fn test_num_to_buffer_zero() {
2326 assert_eq!(num_to_buffer("0"), vec![0x80]);
2328 assert_eq!(num_to_buffer("-0"), vec![0x80]);
2329 }
2330
2331 #[test]
2332 fn test_hash_prefix_string_keys() {
2333 let h1 = compute_hash_prefix(&AttributeValue::S("3635".into()));
2336 let h2 = compute_hash_prefix(&AttributeValue::S("228".into()));
2337 let h3 = compute_hash_prefix(&AttributeValue::S("1668".into()));
2338 let h4 = compute_hash_prefix(&AttributeValue::S("3435".into()));
2339
2340 assert_eq!(
2343 hash_bucket(&h1),
2344 0,
2345 "3635 should be bucket 0, got hash {h1}"
2346 );
2347 assert_eq!(hash_bucket(&h2), 0, "228 should be bucket 0, got hash {h2}");
2348
2349 assert_eq!(
2351 hash_bucket(&h3),
2352 1,
2353 "1668 should be bucket 1, got hash {h3}"
2354 );
2355
2356 assert_eq!(
2358 hash_bucket(&h4),
2359 4,
2360 "3435 should be bucket 4, got hash {h4}"
2361 );
2362 }
2363
2364 #[test]
2365 fn test_hash_prefix_number_keys() {
2366 let h1 = compute_hash_prefix(&AttributeValue::N("251".into()));
2369 assert_eq!(hash_bucket(&h1), 1, "251 should be bucket 1, got hash {h1}");
2370
2371 let h2 = compute_hash_prefix(&AttributeValue::N("2388".into()));
2373 assert_eq!(
2374 hash_bucket(&h2),
2375 4095,
2376 "2388 should be bucket 4095, got hash {h2}"
2377 );
2378 }
2379
2380 #[test]
2381 fn test_hash_in_segment() {
2382 assert!(hash_in_segment("000000", 0, 4096));
2384 assert!(!hash_in_segment("000000", 1, 4096));
2385
2386 assert!(hash_in_segment("001000", 1, 4096));
2388 assert!(!hash_in_segment("001000", 0, 4096));
2389
2390 assert!(hash_in_segment("fff000", 4095, 4096));
2392 assert!(!hash_in_segment("fff000", 0, 4096));
2393
2394 assert!(hash_in_segment("000000", 0, 2));
2396 assert!(hash_in_segment("7ff000", 0, 2));
2397 assert!(hash_in_segment("800000", 1, 2));
2398 assert!(hash_in_segment("fff000", 1, 2));
2399 }
2400}