Skip to main content

alien_bindings/providers/kv/
local.rs

1//! Local disk-persisted KV backed by turso (`localkv.v1`), multi-process safe.
2//!
3//! # Connection strategy
4//!
5//! turso is async-native and its `Connection` is `Send + Sync`, so there is no
6//! `spawn_blocking` boundary and no `Mutex<Connection>` anywhere. `LocalKv`
7//! holds one `turso::Database` handle on `<dataDir>/localkv.sqlite`; each
8//! operation opens its **own** short-lived connection from it and drops it
9//! when the operation completes, so no statement state leaks between
10//! operations.
11//!
12//! Correctness under concurrent access (multiple handles on one file, i.e.
13//! multiple processes) comes from turso's multi-process WAL mode — enabled
14//! explicitly, experimental upstream, and gated by the multi-handle tests
15//! below — plus a `busy_timeout` (writers wait for the write lock instead of
16//! failing with `Busy`). The schema is created once in [`LocalKv::new`]. Reads
17//! of live rows never take the write lock; a read that encounters an expired
18//! row escalates to a short delete (see `get` and `exists`). Conditional puts
19//! are a single atomic `INSERT ... ON CONFLICT DO UPDATE ... WHERE` so the
20//! race is resolved by the database, not by application-level locking.
21//!
22//! See `crates/alien-bindings/FORMAT.md` for the on-disk `localkv.v1` contract.
23use crate::error::{ErrorData, Result};
24use crate::providers::local_store::{
25    as_blob, as_i64, as_opt_i64, as_text, opt_i64_value, query_all, LocalStore, StoreSpec,
26};
27use crate::traits::{Binding, Kv, PutOptions, ScanResult};
28use alien_error::{AlienError, Context as _, IntoAlienError as _};
29use async_trait::async_trait;
30use base64::{engine::general_purpose::URL_SAFE_NO_PAD, Engine as _};
31use chrono::Utc;
32use serde::{Deserialize, Serialize};
33use std::path::PathBuf;
34use turso::Connection;
35
36static KV_SPEC: StoreSpec = StoreSpec {
37    db_filename: "localkv.sqlite",
38    format_version: "localkv.v1",
39    binding_type: "local KV",
40    schema_ddl: "CREATE TABLE IF NOT EXISTS kv (key TEXT PRIMARY KEY, value BLOB NOT NULL, expires_at INTEGER);",
41};
42
43#[derive(Debug)]
44pub struct LocalKv {
45    store: LocalStore,
46}
47
48#[derive(Debug, Serialize, Deserialize)]
49struct CursorState {
50    version: u8,
51    prefix: String,
52    last_key: String,
53}
54
55/// Build the standard KV operation error context.
56fn kv_error(operation: &str, key: &str, reason: &str) -> ErrorData {
57    ErrorData::KvOperationFailed {
58        operation: operation.to_string(),
59        key: key.to_string(),
60        reason: reason.to_string(),
61    }
62}
63
64/// Delete a row that is already expired (`expires_at <= now`), used by the
65/// lazy-expiry paths of `get` and `exists`.
66async fn delete_expired(conn: &Connection, operation: &str, key: &str, now: i64) -> Result<()> {
67    conn.execute(
68        "DELETE FROM kv WHERE key = ?1 AND expires_at IS NOT NULL AND expires_at <= ?2",
69        (key, now),
70    )
71    .await
72    .into_alien_error()
73    .context(kv_error(operation, key, "failed to delete expired row"))?;
74    Ok(())
75}
76
77impl LocalKv {
78    pub async fn new(data_dir: PathBuf) -> Result<Self> {
79        Ok(Self {
80            store: LocalStore::open(data_dir, &KV_SPEC).await?,
81        })
82    }
83
84    /// Get the data directory path (the directory that holds `localkv.sqlite`).
85    pub fn data_dir(&self) -> &PathBuf {
86        self.store.data_dir()
87    }
88
89    /// Get the number of items currently stored (including expired items).
90    /// Useful for testing.
91    pub async fn len(&self) -> Result<usize> {
92        self.store
93            .with_conn(|conn| async move {
94                let rows = query_all(&conn, "SELECT COUNT(*) FROM kv", ())
95                    .await
96                    .into_alien_error()
97                    .context(kv_error("len", "*", "failed to count rows"))?;
98                let count = rows
99                    .first()
100                    .and_then(|row| row.first())
101                    .and_then(as_i64)
102                    .ok_or_else(|| {
103                        AlienError::new(kv_error("len", "*", "count query returned no value"))
104                    })?;
105                Ok(count as usize)
106            })
107            .await
108    }
109
110    /// Check if the store is empty (including expired items).
111    /// Useful for testing.
112    pub async fn is_empty(&self) -> Result<bool> {
113        Ok(self.len().await? == 0)
114    }
115
116    /// Clear all data from the store.
117    /// Useful for testing.
118    pub async fn clear(&self) -> Result<()> {
119        self.store
120            .with_conn(|conn| async move {
121                conn.execute("DELETE FROM kv", ())
122                    .await
123                    .into_alien_error()
124                    .context(kv_error("clear", "*", "failed to clear local KV store"))?;
125                Ok(())
126            })
127            .await
128    }
129
130    /// Get all keys currently in the store (including expired ones).
131    /// Useful for testing and debugging.
132    pub async fn keys(&self) -> Result<Vec<String>> {
133        self.store
134            .with_conn(|conn| async move {
135                let rows = query_all(&conn, "SELECT key FROM kv", ())
136                    .await
137                    .into_alien_error()
138                    .context(kv_error("keys", "*", "failed to scan keys"))?;
139                let mut keys = Vec::with_capacity(rows.len());
140                for row in &rows {
141                    keys.push(row.first().and_then(as_text).ok_or_else(|| {
142                        AlienError::new(kv_error("keys", "*", "failed to read key row"))
143                    })?);
144                }
145                Ok(keys)
146            })
147            .await
148    }
149
150    /// Validate key constraints using global KV validation.
151    fn validate_key(key: &str) -> Result<()> {
152        crate::providers::kv::validate_key(key)
153    }
154
155    /// Validate value constraints using global KV validation.
156    fn validate_value(value: &[u8]) -> Result<()> {
157        crate::providers::kv::validate_value(value)
158    }
159
160    fn encode_cursor(state: &CursorState) -> Result<String> {
161        let bytes =
162            serde_json::to_vec(state)
163                .into_alien_error()
164                .context(ErrorData::InvalidInput {
165                    operation_context: "Local KV cursor encoding".to_string(),
166                    details: "Failed to serialize cursor state".to_string(),
167                    field_name: Some("cursor".to_string()),
168                })?;
169        Ok(URL_SAFE_NO_PAD.encode(bytes))
170    }
171
172    fn decode_cursor(prefix: &str, cursor: &str) -> Result<CursorState> {
173        let bytes =
174            URL_SAFE_NO_PAD
175                .decode(cursor)
176                .into_alien_error()
177                .context(ErrorData::InvalidInput {
178                    operation_context: "Local KV cursor decoding".to_string(),
179                    details: "Invalid cursor encoding".to_string(),
180                    field_name: Some("cursor".to_string()),
181                })?;
182        let state: CursorState =
183            serde_json::from_slice(&bytes)
184                .into_alien_error()
185                .context(ErrorData::InvalidInput {
186                    operation_context: "Local KV cursor decoding".to_string(),
187                    details: "Invalid cursor data".to_string(),
188                    field_name: Some("cursor".to_string()),
189                })?;
190        if state.version != 1 || state.prefix != prefix || !state.last_key.starts_with(prefix) {
191            return Err(AlienError::new(ErrorData::InvalidInput {
192                operation_context: "Local KV cursor validation".to_string(),
193                details: "Cursor does not belong to this prefix scan".to_string(),
194                field_name: Some("cursor".to_string()),
195            }));
196        }
197        Ok(state)
198    }
199}
200
201impl Binding for LocalKv {}
202
203#[async_trait]
204impl Kv for LocalKv {
205    async fn get(&self, key: &str) -> Result<Option<Vec<u8>>> {
206        Self::validate_key(key)?;
207
208        self.store
209            .with_conn(|conn| async move {
210                let now = Utc::now().timestamp_millis();
211                let rows = query_all(
212                    &conn,
213                    "SELECT value, expires_at FROM kv WHERE key = ?1",
214                    (key,),
215                )
216                .await
217                .into_alien_error()
218                .context(kv_error("get", key, "failed to read value"))?;
219
220                let Some(row) = rows.first() else {
221                    return Ok(None);
222                };
223                let value = row.first().and_then(as_blob).ok_or_else(|| {
224                    AlienError::new(kv_error("get", key, "stored value is not a blob"))
225                })?;
226                let expires_at = row.get(1).and_then(as_opt_i64).ok_or_else(|| {
227                    AlienError::new(kv_error("get", key, "stored expires_at is not an integer"))
228                })?;
229
230                if matches!(expires_at, Some(exp) if exp <= now) {
231                    // Lazily remove the expired row.
232                    delete_expired(&conn, "get", key, now).await?;
233                    Ok(None)
234                } else {
235                    Ok(Some(value))
236                }
237            })
238            .await
239    }
240
241    async fn put(&self, key: &str, value: Vec<u8>, options: Option<PutOptions>) -> Result<bool> {
242        Self::validate_key(key)?;
243        Self::validate_value(&value)?;
244        let options = options.unwrap_or_default();
245
246        self.store
247            .with_conn(|conn| async move {
248                let now = Utc::now().timestamp_millis();
249                let expires_at: Option<i64> = options
250                    .ttl
251                    .map(|d| now.saturating_add(i64::try_from(d.as_millis()).unwrap_or(i64::MAX)));
252
253                if options.if_not_exists {
254                    // One atomic statement: insert if absent, otherwise overwrite
255                    // ONLY when the existing row is already expired. The changed
256                    // row count (returned by `execute`) is 1 for the winner, 0
257                    // for a loser.
258                    let changed = conn
259                        .execute(
260                            "INSERT INTO kv (key, value, expires_at) VALUES (?1, ?2, ?3) \
261                             ON CONFLICT(key) DO UPDATE SET value = ?2, expires_at = ?3 \
262                             WHERE kv.expires_at IS NOT NULL AND kv.expires_at <= ?4",
263                            (key, value, opt_i64_value(expires_at), now),
264                        )
265                        .await
266                        .into_alien_error()
267                        .context(kv_error("put", key, "failed conditional put"))?;
268                    Ok(changed == 1)
269                } else {
270                    conn.execute(
271                        "INSERT INTO kv (key, value, expires_at) VALUES (?1, ?2, ?3) \
272                         ON CONFLICT(key) DO UPDATE SET value = ?2, expires_at = ?3",
273                        (key, value, opt_i64_value(expires_at)),
274                    )
275                    .await
276                    .into_alien_error()
277                    .context(kv_error("put", key, "failed to upsert value"))?;
278                    Ok(true)
279                }
280            })
281            .await
282    }
283
284    async fn delete(&self, key: &str) -> Result<()> {
285        Self::validate_key(key)?;
286
287        self.store
288            .with_conn(|conn| async move {
289                conn.execute("DELETE FROM kv WHERE key = ?1", (key,))
290                    .await
291                    .into_alien_error()
292                    .context(kv_error("delete", key, "failed to delete key"))?;
293                Ok(())
294            })
295            .await
296    }
297
298    async fn exists(&self, key: &str) -> Result<bool> {
299        Self::validate_key(key)?;
300
301        self.store
302            .with_conn(|conn| async move {
303                let now = Utc::now().timestamp_millis();
304                let rows = query_all(&conn, "SELECT expires_at FROM kv WHERE key = ?1", (key,))
305                    .await
306                    .into_alien_error()
307                    .context(kv_error("exists", key, "failed to check existence"))?;
308
309                let Some(row) = rows.first() else {
310                    return Ok(false);
311                };
312                let expires_at = row.first().and_then(as_opt_i64).ok_or_else(|| {
313                    AlienError::new(kv_error(
314                        "exists",
315                        key,
316                        "stored expires_at is not an integer",
317                    ))
318                })?;
319
320                if matches!(expires_at, Some(exp) if exp <= now) {
321                    // Lazily remove the expired row.
322                    delete_expired(&conn, "exists", key, now).await?;
323                    Ok(false)
324                } else {
325                    Ok(true)
326                }
327            })
328            .await
329    }
330
331    async fn scan_prefix(
332        &self,
333        prefix: &str,
334        limit: Option<usize>,
335        cursor: Option<String>,
336    ) -> Result<ScanResult> {
337        Self::validate_key(prefix)?;
338
339        let limit = limit.unwrap_or(1000);
340        let last_key = cursor
341            .as_deref()
342            .map(|cursor| Self::decode_cursor(prefix, cursor))
343            .transpose()?
344            .map(|state| state.last_key);
345        if limit == 0 {
346            return Ok(ScanResult {
347                items: Vec::new(),
348                next_cursor: cursor,
349            });
350        }
351        // Collect matching, non-expired items after the last key in sorted order.
352        let matching: Vec<(String, Vec<u8>)> =
353            self.store
354                .with_conn(|conn| async move {
355                    let now = Utc::now().timestamp_millis();
356                    let rows = match last_key.as_deref() {
357                        Some(last_key) => {
358                            query_all(
359                                &conn,
360                                "SELECT key, value, expires_at FROM kv WHERE key > ?1 ORDER BY key",
361                                (last_key,),
362                            )
363                            .await
364                        }
365                        None => query_all(
366                            &conn,
367                            "SELECT key, value, expires_at FROM kv WHERE key >= ?1 ORDER BY key",
368                            (prefix,),
369                        )
370                        .await,
371                    }
372                    .into_alien_error()
373                    .context(kv_error(
374                        "scan_prefix",
375                        prefix,
376                        "failed to scan prefix",
377                    ))?;
378
379                    let mut matching = Vec::new();
380                    for row in &rows {
381                        let k = row.first().and_then(as_text).ok_or_else(|| {
382                            AlienError::new(kv_error(
383                                "scan_prefix",
384                                prefix,
385                                "failed to read scan row key",
386                            ))
387                        })?;
388                        // Keys are ordered ascending starting at `prefix`; once a key
389                        // stops matching the prefix, no later key can match either.
390                        if !k.starts_with(prefix) {
391                            break;
392                        }
393                        let v = row.get(1).and_then(as_blob).ok_or_else(|| {
394                            AlienError::new(kv_error(
395                                "scan_prefix",
396                                prefix,
397                                "stored value is not a blob",
398                            ))
399                        })?;
400                        let exp = row.get(2).and_then(as_opt_i64).ok_or_else(|| {
401                            AlienError::new(kv_error(
402                                "scan_prefix",
403                                prefix,
404                                "stored expires_at is not an integer",
405                            ))
406                        })?;
407                        if matches!(exp, Some(e) if e <= now) {
408                            continue; // expired: treat as absent
409                        }
410                        matching.push((k, v));
411                    }
412                    Ok(matching)
413                })
414                .await?;
415
416        let has_more = matching.len() > limit;
417        let items = matching.into_iter().take(limit).collect::<Vec<_>>();
418        let next_cursor = if has_more {
419            items
420                .last()
421                .map(|(last_key, _)| {
422                    Self::encode_cursor(&CursorState {
423                        version: 1,
424                        prefix: prefix.to_string(),
425                        last_key: last_key.clone(),
426                    })
427                })
428                .transpose()?
429        } else {
430            None
431        };
432
433        Ok(ScanResult { items, next_cursor })
434    }
435}
436
437#[cfg(test)]
438mod tests {
439    use super::*;
440    use crate::providers::local_store::open_database;
441    use std::sync::Arc;
442    use std::time::Duration;
443    use tempfile::TempDir;
444    use tokio::time;
445
446    async fn create_test_kv() -> (LocalKv, TempDir) {
447        let temp_dir = tempfile::tempdir().expect("Failed to create temp dir");
448        let kv = LocalKv::new(temp_dir.path().join("kv.db"))
449            .await
450            .expect("Failed to create LocalKv");
451        (kv, temp_dir)
452    }
453
454    #[tokio::test]
455    async fn test_basic_operations() {
456        let (kv, _temp_dir) = create_test_kv().await;
457
458        assert!(kv
459            .put("test_key", b"test_value".to_vec(), None)
460            .await
461            .unwrap());
462        let value = kv.get("test_key").await.unwrap();
463        assert_eq!(value, Some(b"test_value".to_vec()));
464
465        assert!(kv.exists("test_key").await.unwrap());
466        assert!(!kv.exists("nonexistent").await.unwrap());
467
468        kv.delete("test_key").await.unwrap();
469        assert!(!kv.exists("test_key").await.unwrap());
470        assert_eq!(kv.get("test_key").await.unwrap(), None);
471    }
472
473    #[tokio::test]
474    async fn test_conditional_put() {
475        let (kv, _temp_dir) = create_test_kv().await;
476
477        let options = Some(PutOptions {
478            ttl: None,
479            if_not_exists: true,
480        });
481        assert!(kv
482            .put("key", b"value1".to_vec(), options.clone())
483            .await
484            .unwrap());
485
486        assert!(!kv.put("key", b"value2".to_vec(), options).await.unwrap());
487
488        assert_eq!(kv.get("key").await.unwrap(), Some(b"value1".to_vec()));
489
490        assert!(kv.put("key", b"value3".to_vec(), None).await.unwrap());
491        assert_eq!(kv.get("key").await.unwrap(), Some(b"value3".to_vec()));
492    }
493
494    #[tokio::test]
495    async fn test_ttl_expiration() {
496        let (kv, _temp_dir) = create_test_kv().await;
497
498        let options = Some(PutOptions {
499            ttl: Some(Duration::from_millis(500)),
500            if_not_exists: false,
501        });
502
503        kv.put("expiring_key", b"value".to_vec(), options)
504            .await
505            .unwrap();
506
507        assert!(kv.exists("expiring_key").await.unwrap());
508        assert_eq!(
509            kv.get("expiring_key").await.unwrap(),
510            Some(b"value".to_vec())
511        );
512
513        time::sleep(Duration::from_millis(750)).await;
514
515        assert!(!kv.exists("expiring_key").await.unwrap());
516        assert_eq!(kv.get("expiring_key").await.unwrap(), None);
517    }
518
519    #[tokio::test]
520    async fn test_prefix_scanning() {
521        let (kv, _temp_dir) = create_test_kv().await;
522
523        kv.put("prefix:key1", b"value1".to_vec(), None)
524            .await
525            .unwrap();
526        kv.put("prefix:key2", b"value2".to_vec(), None)
527            .await
528            .unwrap();
529        kv.put("prefix:key3", b"value3".to_vec(), None)
530            .await
531            .unwrap();
532        kv.put("other:key", b"other".to_vec(), None).await.unwrap();
533
534        let result = kv.scan_prefix("prefix:", None, None).await.unwrap();
535        assert_eq!(result.items.len(), 3);
536        assert!(result.next_cursor.is_none());
537
538        assert_eq!(result.items[0].0, "prefix:key1");
539        assert_eq!(result.items[1].0, "prefix:key2");
540        assert_eq!(result.items[2].0, "prefix:key3");
541
542        let result = kv.scan_prefix("prefix:", Some(2), None).await.unwrap();
543        assert_eq!(result.items.len(), 2);
544        assert!(result.next_cursor.is_some());
545
546        let cursor = result.next_cursor.unwrap();
547        let result = kv
548            .scan_prefix("prefix:", Some(2), Some(cursor))
549            .await
550            .unwrap();
551        assert_eq!(result.items.len(), 1);
552        assert_eq!(result.items[0].0, "prefix:key3");
553        assert!(result.next_cursor.is_none());
554    }
555
556    #[tokio::test]
557    async fn prefix_cursor_does_not_skip_after_an_earlier_key_is_deleted() {
558        let (kv, _temp_dir) = create_test_kv().await;
559        for key in ["prefix:key1", "prefix:key2", "prefix:key3"] {
560            kv.put(key, key.as_bytes().to_vec(), None).await.unwrap();
561        }
562
563        let first = kv.scan_prefix("prefix:", Some(2), None).await.unwrap();
564        assert_eq!(
565            first
566                .items
567                .iter()
568                .map(|(key, _)| key.as_str())
569                .collect::<Vec<_>>(),
570            ["prefix:key1", "prefix:key2"]
571        );
572
573        kv.delete("prefix:key1").await.unwrap();
574        let second = kv
575            .scan_prefix("prefix:", Some(2), first.next_cursor)
576            .await
577            .unwrap();
578        assert_eq!(second.items[0].0, "prefix:key3");
579        assert!(second.next_cursor.is_none());
580    }
581
582    #[tokio::test]
583    async fn prefix_cursor_is_rejected_for_another_prefix() {
584        let (kv, _temp_dir) = create_test_kv().await;
585        for key in ["first:key1", "first:key2"] {
586            kv.put(key, key.as_bytes().to_vec(), None).await.unwrap();
587        }
588
589        let cursor = kv
590            .scan_prefix("first:", Some(1), None)
591            .await
592            .unwrap()
593            .next_cursor
594            .expect("the first page should have a cursor");
595        assert!(kv
596            .scan_prefix("second:", Some(1), Some(cursor.clone()))
597            .await
598            .is_err());
599        assert!(kv
600            .scan_prefix("second:", Some(0), Some(cursor))
601            .await
602            .is_err());
603        assert!(kv
604            .scan_prefix("first:", Some(0), Some("not-a-cursor".to_string()))
605            .await
606            .is_err());
607    }
608
609    #[tokio::test]
610    async fn test_persistence_across_reopens() {
611        let temp_dir = tempfile::tempdir().expect("Failed to create temp dir");
612        let db_path = temp_dir.path().join("kv.db");
613
614        {
615            let kv = LocalKv::new(db_path.clone())
616                .await
617                .expect("Failed to create LocalKv");
618            kv.put("persistent_key", b"persistent_value".to_vec(), None)
619                .await
620                .unwrap();
621        }
622
623        {
624            let kv = LocalKv::new(db_path)
625                .await
626                .expect("Failed to reopen LocalKv");
627            let value = kv.get("persistent_key").await.unwrap();
628            assert_eq!(value, Some(b"persistent_value".to_vec()));
629        }
630    }
631
632    #[tokio::test]
633    async fn test_key_validation() {
634        let (kv, _temp_dir) = create_test_kv().await;
635
636        assert!(kv.put("", b"value".to_vec(), None).await.is_err());
637        assert!(kv.get("").await.is_err());
638
639        let long_key = "a".repeat(513);
640        assert!(kv.put(&long_key, b"value".to_vec(), None).await.is_err());
641
642        assert!(kv
643            .put("key with spaces", b"value".to_vec(), None)
644            .await
645            .is_err());
646        assert!(kv
647            .put("key\nwith\nnewlines", b"value".to_vec(), None)
648            .await
649            .is_err());
650        assert!(kv
651            .put("key/with/slashes", b"value".to_vec(), None)
652            .await
653            .is_err());
654
655        assert!(kv
656            .put("valid_key-123", b"value".to_vec(), None)
657            .await
658            .is_ok());
659        assert!(kv
660            .put("domain.com:8080", b"value".to_vec(), None)
661            .await
662            .is_ok());
663    }
664
665    #[tokio::test]
666    async fn test_value_validation() {
667        let (kv, _temp_dir) = create_test_kv().await;
668
669        let large_value = vec![0u8; 24_577];
670        assert!(kv.put("key", large_value, None).await.is_err());
671
672        let max_value = vec![0u8; 24_576];
673        assert!(kv.put("key", max_value, None).await.is_ok());
674    }
675
676    #[tokio::test]
677    async fn test_utility_methods() {
678        let (kv, _temp_dir) = create_test_kv().await;
679
680        assert!(kv.is_empty().await.unwrap());
681        assert_eq!(kv.len().await.unwrap(), 0);
682        assert_eq!(kv.keys().await.unwrap(), Vec::<String>::new());
683
684        kv.put("key1", b"value1".to_vec(), None).await.unwrap();
685        kv.put("key2", b"value2".to_vec(), None).await.unwrap();
686
687        assert!(!kv.is_empty().await.unwrap());
688        assert_eq!(kv.len().await.unwrap(), 2);
689
690        let mut keys = kv.keys().await.unwrap();
691        keys.sort();
692        assert_eq!(keys, vec!["key1", "key2"]);
693
694        kv.clear().await.unwrap();
695        assert!(kv.is_empty().await.unwrap());
696        assert_eq!(kv.len().await.unwrap(), 0);
697    }
698
699    #[tokio::test]
700    async fn test_unknown_format_rejected_on_open() {
701        let temp_dir = tempfile::tempdir().expect("Failed to create temp dir");
702        let dir = temp_dir.path().join("kv");
703
704        // Create a valid store, then rewrite its format marker to a future version.
705        {
706            let kv = LocalKv::new(dir.clone()).await.expect("initial open");
707            kv.put("k", b"v".to_vec(), None).await.unwrap();
708        }
709        {
710            let db = open_database(&dir.join("localkv.sqlite"), "test")
711                .await
712                .expect("raw open");
713            let conn = db.connect().expect("raw connect");
714            conn.execute(
715                "UPDATE meta SET value = 'localkv.v2' WHERE key = 'format'",
716                (),
717            )
718            .await
719            .expect("format overwrite");
720        }
721
722        // Reopening must fail fast, naming both the found and expected formats.
723        let err = LocalKv::new(dir)
724            .await
725            .expect_err("unknown format must be rejected");
726        let msg = err.to_string();
727        assert!(
728            msg.contains("localkv.v2"),
729            "error must name the found format, got: {msg}"
730        );
731        assert!(
732            msg.contains("localkv.v1"),
733            "error must name the expected format, got: {msg}"
734        );
735    }
736
737    // ---- Multi-process-safety proofs ----
738
739    #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
740    async fn test_conditional_put_atomicity_across_handles() {
741        let temp_dir = tempfile::tempdir().expect("Failed to create temp dir");
742        let dir = temp_dir.path().join("kv");
743        // Two independent handles on the SAME data_dir == two processes sharing the file.
744        let kv_a = Arc::new(LocalKv::new(dir.clone()).await.expect("open handle a"));
745        let kv_b = Arc::new(LocalKv::new(dir.clone()).await.expect("open handle b"));
746
747        let n = 16;
748        let mut handles = Vec::new();
749        for i in 0..n {
750            let kv = if i % 2 == 0 {
751                kv_a.clone()
752            } else {
753                kv_b.clone()
754            };
755            handles.push(tokio::spawn(async move {
756                let val = format!("val-{i}").into_bytes();
757                let opts = Some(PutOptions {
758                    ttl: None,
759                    if_not_exists: true,
760                });
761                let won = kv.put("race", val.clone(), opts).await.expect("put ok");
762                (won, val)
763            }));
764        }
765
766        let mut winners = Vec::new();
767        for h in handles {
768            let (won, val) = h.await.expect("task join");
769            if won {
770                winners.push(val);
771            }
772        }
773
774        assert_eq!(
775            winners.len(),
776            1,
777            "exactly one conditional put must win across both handles"
778        );
779        let stored = kv_a.get("race").await.unwrap().expect("key present");
780        assert_eq!(stored, winners[0], "stored value must equal the winner");
781        assert_eq!(
782            kv_b.get("race").await.unwrap().expect("key present via b"),
783            winners[0]
784        );
785    }
786
787    #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
788    async fn test_ttl_expiry_takeover_conditional_put() {
789        let temp_dir = tempfile::tempdir().expect("Failed to create temp dir");
790        let dir = temp_dir.path().join("kv");
791        let kv_a = Arc::new(LocalKv::new(dir.clone()).await.expect("open handle a"));
792        let kv_b = Arc::new(LocalKv::new(dir.clone()).await.expect("open handle b"));
793
794        // Seed a short-lived key with a conditional put.
795        assert!(kv_a
796            .put(
797                "k",
798                b"initial".to_vec(),
799                Some(PutOptions {
800                    ttl: Some(Duration::from_millis(300)),
801                    if_not_exists: true,
802                }),
803            )
804            .await
805            .unwrap());
806
807        // While it is still live, a conditional put must lose.
808        assert!(!kv_b
809            .put(
810                "k",
811                b"early".to_vec(),
812                Some(PutOptions {
813                    ttl: None,
814                    if_not_exists: true,
815                }),
816            )
817            .await
818            .unwrap());
819
820        // Wait for the seeded key to expire.
821        time::sleep(Duration::from_millis(450)).await;
822
823        // Race the takeover: many conditional puts against the now-expired key.
824        let n = 12;
825        let mut handles = Vec::new();
826        for i in 0..n {
827            let kv = if i % 2 == 0 {
828                kv_a.clone()
829            } else {
830                kv_b.clone()
831            };
832            handles.push(tokio::spawn(async move {
833                let val = format!("takeover-{i}").into_bytes();
834                let won = kv
835                    .put(
836                        "k",
837                        val.clone(),
838                        Some(PutOptions {
839                            ttl: None,
840                            if_not_exists: true,
841                        }),
842                    )
843                    .await
844                    .expect("put ok");
845                (won, val)
846            }));
847        }
848
849        let mut winners = Vec::new();
850        for h in handles {
851            let (won, val) = h.await.expect("task join");
852            if won {
853                winners.push(val);
854            }
855        }
856
857        assert_eq!(
858            winners.len(),
859            1,
860            "exactly one takeover conditional put must win after expiry"
861        );
862        let stored = kv_b.get("k").await.unwrap().expect("key present");
863        assert_eq!(
864            stored, winners[0],
865            "stored value must equal the takeover winner"
866        );
867    }
868
869    #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
870    async fn test_multi_handle_concurrent_smoke() {
871        let temp_dir = tempfile::tempdir().expect("Failed to create temp dir");
872        let dir = temp_dir.path().join("kv");
873        let kv_a = Arc::new(LocalKv::new(dir.clone()).await.expect("open handle a"));
874        let kv_b = Arc::new(LocalKv::new(dir.clone()).await.expect("open handle b"));
875
876        let mut handles = Vec::new();
877        for i in 0..50 {
878            let kv = if i % 2 == 0 {
879                kv_a.clone()
880            } else {
881                kv_b.clone()
882            };
883            handles.push(tokio::spawn(async move {
884                let key = format!("key_{i}");
885                let val = format!("v{i}").into_bytes();
886                // No busy errors expected under multi-process WAL + busy_timeout.
887                kv.put(&key, val.clone(), None).await.expect("put ok");
888                let got = kv.get(&key).await.expect("get ok");
889                assert_eq!(got, Some(val));
890            }));
891        }
892        for h in handles {
893            h.await.expect("task join");
894        }
895
896        assert_eq!(kv_a.len().await.unwrap(), 50, "handle a sees all keys");
897        assert_eq!(kv_b.len().await.unwrap(), 50, "handle b sees all keys");
898    }
899}