Skip to content

Commit 9e60fba

Browse files
fix(sql): preserve typed storage keys (#227)
1 parent 047da93 commit 9e60fba

2 files changed

Lines changed: 164 additions & 13 deletions

File tree

src/sql.rs

Lines changed: 116 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -60,6 +60,11 @@ use gluesql_core::{
6060
#[cfg(feature = "sql")]
6161
use crate::git::versioned_store::ThreadSafeGitVersionedKvStore;
6262

63+
#[cfg(feature = "sql")]
64+
const LEGACY_SCHEMA_ID: &str = "__schema__";
65+
#[cfg(feature = "sql")]
66+
const ROW_KEY_PREFIX: &[u8] = b":row:";
67+
6368
/// GlueSQL storage backend using ProllyTree.
6469
///
6570
/// Wraps a [`ThreadSafeGitVersionedKvStore`] and implements GlueSQL's async
@@ -109,21 +114,57 @@ impl<const D: usize> ProllyStorage<D> {
109114

110115
/// Convert table name and row key to storage key
111116
fn make_storage_key(table_name: &str, key: &Key) -> Vec<u8> {
117+
let mut storage_key = Self::row_prefix(table_name);
118+
let encoded_key = serde_json::to_vec(key).expect("GlueSQL key serialization cannot fail");
119+
storage_key.extend_from_slice(&encoded_key);
120+
storage_key
121+
}
122+
123+
/// Convert table name and row key to the legacy untyped storage key.
124+
fn legacy_storage_key(table_name: &str, key: &Key) -> Vec<u8> {
112125
match key {
113126
Key::I64(id) => format!("{table_name}:{id}").into_bytes(),
114127
Key::Str(id) => format!("{table_name}:{id}").into_bytes(),
115-
Key::None => format!("{table_name}:__schema__").into_bytes(),
128+
Key::None => format!("{table_name}:{LEGACY_SCHEMA_ID}").into_bytes(),
116129
_ => format!("{table_name}:{key:?}").into_bytes(),
117130
}
118131
}
119132

120133
/// Get schema key for a table
121134
fn schema_key(table_name: &str) -> Vec<u8> {
122-
Self::make_storage_key(table_name, &Key::None)
135+
format!("{table_name}:{LEGACY_SCHEMA_ID}").into_bytes()
136+
}
137+
138+
/// Get the typed row key prefix for a table.
139+
fn row_prefix(table_name: &str) -> Vec<u8> {
140+
let mut prefix = table_name.as_bytes().to_vec();
141+
prefix.extend_from_slice(ROW_KEY_PREFIX);
142+
prefix
143+
}
144+
145+
/// Get the legacy table prefix for untyped row keys.
146+
fn legacy_table_prefix(table_name: &str) -> Vec<u8> {
147+
format!("{table_name}:").into_bytes()
123148
}
124149

125150
/// Parse key from storage key string
126-
fn parse_key_from_storage_key(storage_key: &[u8], table_prefix: &str) -> Key {
151+
fn parse_key_from_storage_key(storage_key: &[u8], table_prefix: &str) -> Result<Key> {
152+
let row_prefix = Self::row_prefix(table_prefix);
153+
if let Some(encoded_key) = storage_key.strip_prefix(row_prefix.as_slice()) {
154+
let key = serde_json::from_slice(encoded_key).map_err(|e| {
155+
Error::StorageMsg(format!("Failed to deserialize typed row key: {e}"))
156+
})?;
157+
return Ok(key);
158+
}
159+
160+
Ok(Self::parse_legacy_key_from_storage_key(
161+
storage_key,
162+
table_prefix,
163+
))
164+
}
165+
166+
/// Parse a key from the legacy untyped storage key string.
167+
fn parse_legacy_key_from_storage_key(storage_key: &[u8], table_prefix: &str) -> Key {
127168
let key_str = String::from_utf8_lossy(storage_key);
128169
let key_part = key_str
129170
.strip_prefix(&format!("{table_prefix}:"))
@@ -136,6 +177,11 @@ impl<const D: usize> ProllyStorage<D> {
136177
}
137178
}
138179

180+
/// Check whether a key is the legacy schema sentinel for this table.
181+
fn is_schema_storage_key(table_name: &str, storage_key: &[u8]) -> bool {
182+
storage_key == Self::schema_key(table_name).as_slice()
183+
}
184+
139185
/// Commit with a custom message.
140186
///
141187
/// Uses `spawn_blocking` to offload the synchronous git commit to a
@@ -218,11 +264,22 @@ impl<const D: usize> Store for ProllyStorage<D> {
218264
async fn fetch_data(&self, table_name: &str, key: &Key) -> Result<Option<DataRow>> {
219265
let store = self.store.clone();
220266
let storage_key = Self::make_storage_key(table_name, key);
267+
let legacy_storage_key = Self::legacy_storage_key(table_name, key);
268+
let schema_key = Self::schema_key(table_name);
221269
tokio::task::spawn_blocking(move || {
222270
if let Some(row_data) = store.get(&storage_key) {
223271
let row: DataRow = serde_json::from_slice(&row_data)
224272
.map_err(|e| Error::StorageMsg(format!("Failed to deserialize row: {e}")))?;
225273
Ok(Some(row))
274+
} else if legacy_storage_key != schema_key {
275+
if let Some(row_data) = store.get(&legacy_storage_key) {
276+
let row: DataRow = serde_json::from_slice(&row_data).map_err(|e| {
277+
Error::StorageMsg(format!("Failed to deserialize row: {e}"))
278+
})?;
279+
Ok(Some(row))
280+
} else {
281+
Ok(None)
282+
}
226283
} else {
227284
Ok(None)
228285
}
@@ -235,20 +292,35 @@ impl<const D: usize> Store for ProllyStorage<D> {
235292
let store = self.store.clone();
236293
let table_name = table_name.to_string();
237294
tokio::task::spawn_blocking(move || {
238-
let prefix = format!("{table_name}:");
239-
let prefix_bytes = prefix.as_bytes();
295+
let row_prefix = Self::row_prefix(&table_name);
296+
let legacy_prefix = Self::legacy_table_prefix(&table_name);
240297

241298
let all_keys = store
242299
.list_keys()
243300
.map_err(|e| Error::StorageMsg(format!("Failed to list keys: {e}")))?;
244-
let mut rows = Vec::new();
301+
let mut rows_by_key = HashMap::new();
302+
303+
for storage_key in all_keys.iter() {
304+
if storage_key.starts_with(&legacy_prefix)
305+
&& !storage_key.starts_with(&row_prefix)
306+
&& !Self::is_schema_storage_key(&table_name, storage_key)
307+
{
308+
if let Some(row_data) = store.get(storage_key) {
309+
let row: DataRow = serde_json::from_slice(&row_data).map_err(|e| {
310+
Error::StorageMsg(format!("Failed to deserialize row: {e}"))
311+
})?;
245312

246-
for storage_key in all_keys {
247-
if storage_key.starts_with(prefix_bytes) {
248-
if storage_key.ends_with(b":__schema__") {
249-
continue;
313+
let key = ProllyStorage::<D>::parse_key_from_storage_key(
314+
storage_key,
315+
&table_name,
316+
)?;
317+
rows_by_key.entry(key).or_insert(row);
250318
}
319+
}
320+
}
251321

322+
for storage_key in all_keys {
323+
if storage_key.starts_with(&row_prefix) {
252324
if let Some(row_data) = store.get(&storage_key) {
253325
let row: DataRow = serde_json::from_slice(&row_data).map_err(|e| {
254326
Error::StorageMsg(format!("Failed to deserialize row: {e}"))
@@ -257,12 +329,17 @@ impl<const D: usize> Store for ProllyStorage<D> {
257329
let key = ProllyStorage::<D>::parse_key_from_storage_key(
258330
&storage_key,
259331
&table_name,
260-
);
261-
rows.push(Ok((key, row)));
332+
)?;
333+
rows_by_key.insert(key, row);
262334
}
263335
}
264336
}
265337

338+
let mut rows: Vec<_> = rows_by_key
339+
.into_iter()
340+
.map(|(key, row)| Ok((key, row)))
341+
.collect();
342+
266343
rows.sort_by(|a, b| match (a, b) {
267344
(Ok((key_a, _)), Ok((key_b, _))) => key_a.cmp(key_b),
268345
_ => std::cmp::Ordering::Equal,
@@ -371,10 +448,17 @@ impl<const D: usize> StoreMut for ProllyStorage<D> {
371448
tokio::task::spawn_blocking(move || {
372449
for key in keys {
373450
let storage_key = ProllyStorage::<D>::make_storage_key(&table_name, &key);
451+
let legacy_storage_key = ProllyStorage::<D>::legacy_storage_key(&table_name, &key);
452+
let schema_key = ProllyStorage::<D>::schema_key(&table_name);
374453

375454
store
376455
.delete(&storage_key)
377456
.map_err(|e| Error::StorageMsg(format!("Failed to delete row: {e}")))?;
457+
if legacy_storage_key != schema_key {
458+
store.delete(&legacy_storage_key).map_err(|e| {
459+
Error::StorageMsg(format!("Failed to delete legacy row: {e}"))
460+
})?;
461+
}
378462
}
379463
Ok(())
380464
})
@@ -493,4 +577,24 @@ mod tests {
493577
let first = iter.next().await.unwrap().unwrap();
494578
assert_eq!(first.0, key);
495579
}
580+
581+
#[test]
582+
fn test_text_key_round_trip_preserves_key_variant() {
583+
let storage_key = ProllyStorage::<32>::make_storage_key("users", &Key::Str("42".into()));
584+
let parsed =
585+
ProllyStorage::<32>::parse_key_from_storage_key(&storage_key, "users").unwrap();
586+
587+
assert_eq!(parsed, Key::Str("42".into()));
588+
}
589+
590+
#[test]
591+
fn test_non_string_key_round_trip_preserves_key_variant() {
592+
for key in [Key::I64(42), Key::Bool(true), Key::None] {
593+
let storage_key = ProllyStorage::<32>::make_storage_key("users", &key);
594+
let parsed =
595+
ProllyStorage::<32>::parse_key_from_storage_key(&storage_key, "users").unwrap();
596+
597+
assert_eq!(parsed, key);
598+
}
599+
}
496600
}

tests/sql_integration.rs

Lines changed: 48 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,12 @@ limitations under the License.
1818

1919
mod common;
2020

21-
use gluesql_core::prelude::Glue;
21+
use futures::StreamExt;
22+
use gluesql_core::{
23+
data::{Key, Value},
24+
prelude::Glue,
25+
store::{DataRow, Store, StoreMut},
26+
};
2227
use prollytree::sql::ProllyStorage;
2328

2429
// ---------------------------------------------------------------------------
@@ -104,6 +109,48 @@ async fn test_sql_multiple_tables() {
104109
assert!(!r2.is_empty());
105110
}
106111

112+
#[tokio::test]
113+
async fn test_sql_storage_scan_preserves_text_primary_key_order() {
114+
let (_temp, dataset) = common::setup_repo_and_dataset();
115+
let mut storage = ProllyStorage::<32>::init(&dataset).expect("ProllyStorage init");
116+
117+
storage
118+
.insert_data(
119+
"text_keys",
120+
vec![
121+
(
122+
Key::Str("2".into()),
123+
DataRow::Vec(vec![Value::Str("two".into())]),
124+
),
125+
(
126+
Key::Str("10".into()),
127+
DataRow::Vec(vec![Value::Str("ten".into())]),
128+
),
129+
(
130+
Key::Str("42".into()),
131+
DataRow::Vec(vec![Value::Str("forty-two".into())]),
132+
),
133+
],
134+
)
135+
.await
136+
.unwrap();
137+
138+
let mut rows = storage.scan_data("text_keys").await.unwrap();
139+
let mut keys = Vec::new();
140+
while let Some(row) = rows.next().await {
141+
keys.push(row.unwrap().0);
142+
}
143+
144+
assert_eq!(
145+
keys,
146+
[
147+
Key::Str("10".into()),
148+
Key::Str("2".into()),
149+
Key::Str("42".into())
150+
]
151+
);
152+
}
153+
107154
// ---------------------------------------------------------------------------
108155
// Drop table
109156
// ---------------------------------------------------------------------------

0 commit comments

Comments
 (0)