1use 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
55fn 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
64async 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 pub fn data_dir(&self) -> &PathBuf {
86 self.store.data_dir()
87 }
88
89 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 pub async fn is_empty(&self) -> Result<bool> {
113 Ok(self.len().await? == 0)
114 }
115
116 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 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 fn validate_key(key: &str) -> Result<()> {
152 crate::providers::kv::validate_key(key)
153 }
154
155 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 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 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 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 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 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; }
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 {
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 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 #[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 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 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 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 time::sleep(Duration::from_millis(450)).await;
822
823 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 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}