1use std::collections::HashMap;
6use std::fmt::Write as _;
7use std::path::Path;
8
9use rusqlite::{Connection, OptionalExtension, TransactionBehavior, params};
10
11use crate::model::{ListEventsOptions, ListObjectsOptions};
12use crate::{
13 EventLog, MemoryEvent, MemoryObject, ObjectKey, ObjectStore, QueueClaimOptions, QueueJob,
14 QueueJobStatus, QueueNackOptions, QueueStore, ThingdError, ThingdResult, u64_to_i64,
15 unix_timestamp_millis,
16};
17
18pub const SQLITE_SCHEMA_VERSION: u32 = 5;
20
21pub struct SqliteThingStore {
23 connection: Connection,
24 db_path: Option<String>,
25 event_idempotency_keys: HashMap<(String, String), u64>,
26}
27
28impl SqliteThingStore {
29 pub fn open(path: impl AsRef<Path>) -> ThingdResult<Self> {
35 let path_str = path.as_ref().to_str().map(String::from);
36 let connection = Connection::open(path).map_err(ThingdError::from)?;
37 connection
38 .busy_timeout(std::time::Duration::from_secs(5))
39 .map_err(ThingdError::from)?;
40 let store = Self {
41 connection,
42 db_path: path_str,
43 event_idempotency_keys: HashMap::new(),
44 };
45 store.initialize()?;
46 Ok(store)
47 }
48
49 pub fn open_in_memory() -> ThingdResult<Self> {
55 let connection = Connection::open_in_memory().map_err(ThingdError::from)?;
56 connection
57 .busy_timeout(std::time::Duration::from_secs(5))
58 .map_err(ThingdError::from)?;
59 let store = Self {
60 connection,
61 db_path: None,
62 event_idempotency_keys: HashMap::new(),
63 };
64 store.initialize()?;
65 Ok(store)
66 }
67
68 fn initialize(&self) -> ThingdResult<()> {
69 let current_mode: String = self
70 .connection
71 .query_row("PRAGMA journal_mode;", [], |row| row.get(0))
72 .unwrap_or_else(|_| "delete".to_string());
73
74 if current_mode.to_lowercase() != "wal" {
75 self.connection
76 .query_row("PRAGMA journal_mode = WAL;", [], |_| Ok(()))
77 .map_err(|e| {
78 eprintln!("warning: failed to enable WAL journal mode: {e}");
79 })
80 .ok();
81 }
82
83 self.connection
84 .execute_batch(
85 r"
86 PRAGMA synchronous = NORMAL;
87 PRAGMA foreign_keys = ON;
88 PRAGMA busy_timeout = 5000;
89 ",
90 )
91 .map_err(ThingdError::from)?;
92
93 self.connection
94 .execute_batch(
95 r"
96 CREATE TABLE IF NOT EXISTS thingd_schema_migrations (
97 version INTEGER PRIMARY KEY,
98 name TEXT NOT NULL,
99 applied_at TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ', 'now'))
100 );
101 ",
102 )
103 .map_err(ThingdError::from)?;
104
105 let current_version = self.schema_version()?;
106
107 if current_version > 0 && current_version < SQLITE_SCHEMA_VERSION {
109 #[allow(clippy::collapsible_if)]
110 if let Some(ref path) = self.db_path {
111 let backup_path = format!("{path}.pre-v{current_version}");
112 let escaped = backup_path.replace('\'', "''");
113 let _ = self
114 .connection
115 .execute_batch(&format!("VACUUM INTO '{escaped}'"));
116 }
117 }
118
119 if current_version < 1 {
120 self.apply_schema_v1()?;
121 }
122
123 let current_version = self.schema_version()?;
124
125 if current_version < 2 {
126 self.apply_schema_v2()?;
127 }
128
129 let current_version = self.schema_version()?;
130
131 if current_version < 3 {
132 self.apply_schema_v3()?;
133 }
134
135 let current_version = self.schema_version()?;
136
137 if current_version < 4 {
138 self.apply_schema_v4()?;
139 }
140
141 let current_version = self.schema_version()?;
142
143 if current_version < 5 {
144 self.apply_schema_v5()?;
145 }
146
147 if current_version > SQLITE_SCHEMA_VERSION {
148 return Err(ThingdError::Storage(format!(
149 "database schema version {current_version} is newer than supported version {SQLITE_SCHEMA_VERSION}"
150 )));
151 }
152
153 let ok: String = self
155 .connection
156 .query_row("PRAGMA quick_check", [], |row| row.get(0))
157 .map_err(ThingdError::from)?;
158 if ok != "ok" {
159 return Err(ThingdError::Storage(format!(
160 "database integrity check failed: {ok}"
161 )));
162 }
163
164 Ok(())
165 }
166
167 pub fn schema_version(&self) -> ThingdResult<u32> {
173 let version = self
174 .connection
175 .query_row(
176 "SELECT COALESCE(MAX(version), 0) FROM thingd_schema_migrations",
177 [],
178 |row| row.get::<_, i64>(0),
179 )
180 .map_err(ThingdError::from)?;
181
182 u32::try_from(version).map_err(|error| ThingdError::Storage(error.to_string()))
183 }
184
185 pub fn wal_checkpoint(&self) -> ThingdResult<(i32, i32)> {
194 self.connection
195 .query_row("PRAGMA wal_checkpoint(TRUNCATE)", [], |row| {
196 Ok((row.get::<_, i32>(0)?, row.get::<_, i32>(1)?))
197 })
198 .map_err(ThingdError::from)
199 }
200
201 pub fn optimize_search_index(&self) -> ThingdResult<()> {
210 self.connection
211 .execute_batch("INSERT INTO search_index(search_index) VALUES('optimize')")
212 .map_err(ThingdError::from)
213 }
214
215 pub fn close(&self) -> ThingdResult<()> {
221 self.wal_checkpoint()?;
222 Ok(())
223 }
224
225 pub fn backup_to(&self, path: &str) -> ThingdResult<()> {
235 if Path::new(path)
237 .components()
238 .filter_map(|c| match c {
239 std::path::Component::Normal(s) => Some(s.to_str().unwrap_or("")),
240 std::path::Component::ParentDir => Some(".."),
241 _ => None,
242 })
243 .any(|x| x == "..")
244 {
245 return Err(ThingdError::InvalidInput(
246 "Backup path must not contain '..' (path traversal)".to_string(),
247 ));
248 }
249 let escaped = path.replace('\'', "''");
250 self.connection
251 .execute_batch(&format!("VACUUM INTO '{escaped}'"))
252 .map_err(ThingdError::from)
253 }
254
255 fn apply_schema_v1(&self) -> ThingdResult<()> {
256 self.connection
257 .execute_batch(
258 r"
259 BEGIN;
260
261 CREATE TABLE IF NOT EXISTS objects (
262 collection TEXT NOT NULL,
263 id TEXT NOT NULL,
264 body TEXT NOT NULL,
265 version INTEGER NOT NULL,
266 created_at TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ', 'now')),
267 updated_at TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ', 'now')),
268 PRIMARY KEY (collection, id)
269 );
270
271 CREATE TABLE IF NOT EXISTS events (
272 sequence INTEGER PRIMARY KEY AUTOINCREMENT,
273 stream TEXT NOT NULL,
274 event_type TEXT NOT NULL,
275 body TEXT NOT NULL,
276 created_at TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ', 'now'))
277 );
278
279 CREATE INDEX IF NOT EXISTS idx_events_stream_sequence
280 ON events (stream, sequence);
281
282 CREATE TABLE IF NOT EXISTS queue_jobs (
283 queue TEXT NOT NULL,
284 id TEXT NOT NULL,
285 body TEXT NOT NULL,
286 attempts INTEGER NOT NULL,
287 max_attempts INTEGER NOT NULL,
288 status TEXT NOT NULL,
289 available_at_ms INTEGER NOT NULL,
290 leased_at_ms INTEGER,
291 lease_expires_at_ms INTEGER,
292 completed_at_ms INTEGER,
293 dead_at_ms INTEGER,
294 created_at TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ', 'now')),
295 updated_at TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ', 'now')),
296 PRIMARY KEY (queue, id)
297 );
298
299 CREATE INDEX IF NOT EXISTS idx_queue_jobs_queue_status_created
300 ON queue_jobs (queue, status, created_at);
301
302 CREATE INDEX IF NOT EXISTS idx_queue_jobs_status
303 ON queue_jobs (status);
304
305 INSERT OR IGNORE INTO thingd_schema_migrations (version, name, applied_at)
306 VALUES (1, 'initial_objects_events_queues', strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
307
308 COMMIT;
309 ",
310 )
311 .map_err(ThingdError::from)?;
312
313 Ok(())
314 }
315
316 fn apply_schema_v2(&self) -> ThingdResult<()> {
317 let tx = self
318 .connection
319 .unchecked_transaction()
320 .map_err(ThingdError::from)?;
321
322 tx.execute_batch(
323 "CREATE VIRTUAL TABLE IF NOT EXISTS search_index USING fts5(
324 collection UNINDEXED,
325 id UNINDEXED,
326 kind UNINDEXED,
327 text,
328 tokenize='porter unicode61'
329 );",
330 )
331 .map_err(ThingdError::from)?;
332
333 Self::reindex_all_into(&tx)?;
335
336 tx.execute(
337 "INSERT OR IGNORE INTO thingd_schema_migrations (version, name, applied_at)
338 VALUES (2, 'fts5_search_index', strftime('%Y-%m-%dT%H:%M:%fZ', 'now'))",
339 [],
340 )
341 .map_err(ThingdError::from)?;
342
343 tx.commit().map_err(ThingdError::from)?;
344 Ok(())
345 }
346
347 fn apply_schema_v3(&self) -> ThingdResult<()> {
348 self.connection
349 .execute(
350 "ALTER TABLE queue_jobs ADD COLUMN last_error TEXT NOT NULL DEFAULT ''",
351 [],
352 )
353 .map_err(ThingdError::from)?;
354 self.connection
355 .execute(
356 "INSERT OR IGNORE INTO thingd_schema_migrations (version, name, applied_at)
357 VALUES (3, 'queue_jobs_last_error', strftime('%Y-%m-%dT%H:%M:%fZ', 'now'))",
358 [],
359 )
360 .map_err(ThingdError::from)?;
361 Ok(())
362 }
363
364 fn apply_schema_v4(&self) -> ThingdResult<()> {
365 self.connection
366 .execute_batch(
367 r"
368 CREATE TABLE IF NOT EXISTS links (
369 id TEXT PRIMARY KEY,
370 from_ref TEXT NOT NULL,
371 type TEXT NOT NULL,
372 to_ref TEXT NOT NULL,
373 weight REAL,
374 metadata_json TEXT NOT NULL DEFAULT '{}',
375 created_at TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ', 'now'))
376 );
377
378 CREATE INDEX IF NOT EXISTS idx_links_from_ref ON links (from_ref);
379 CREATE INDEX IF NOT EXISTS idx_links_to_ref ON links (to_ref);
380 CREATE INDEX IF NOT EXISTS idx_links_type ON links (type);
381 ",
382 )
383 .map_err(ThingdError::from)?;
384 self.connection
385 .execute(
386 "INSERT OR IGNORE INTO thingd_schema_migrations (version, name, applied_at)
387 VALUES (4, 'graph_links', strftime('%Y-%m-%dT%H:%M:%fZ', 'now'))",
388 [],
389 )
390 .map_err(ThingdError::from)?;
391 Ok(())
392 }
393
394 fn apply_schema_v5(&self) -> ThingdResult<()> {
395 self.connection
396 .execute_batch(
397 r"
398 CREATE INDEX IF NOT EXISTS idx_objects_collection ON objects (collection);
399 CREATE INDEX IF NOT EXISTS idx_objects_created_at ON objects (created_at);
400 CREATE INDEX IF NOT EXISTS idx_events_stream ON events (stream);
401 ",
402 )
403 .map_err(ThingdError::from)?;
404 self.connection
405 .execute(
406 "INSERT OR IGNORE INTO thingd_schema_migrations (version, name, applied_at)
407 VALUES (5, 'performance_indexes', strftime('%Y-%m-%dT%H:%M:%fZ', 'now'))",
408 [],
409 )
410 .map_err(ThingdError::from)?;
411 Ok(())
412 }
413
414 fn reindex_all_into(tx: &rusqlite::Transaction<'_>) -> ThingdResult<()> {
415 let mut stmt_objects = tx
416 .prepare("SELECT collection, id, body FROM objects")
417 .map_err(ThingdError::from)?;
418 let rows_objects = stmt_objects
419 .query_map([], |row| {
420 Ok((
421 row.get::<_, String>(0)?,
422 row.get::<_, String>(1)?,
423 row.get::<_, String>(2)?,
424 ))
425 })
426 .map_err(ThingdError::from)?;
427
428 let mut stmt_events = tx
429 .prepare("SELECT stream, sequence, body FROM events")
430 .map_err(ThingdError::from)?;
431 let rows_events = stmt_events
432 .query_map([], |row| {
433 Ok((
434 row.get::<_, String>(0)?,
435 row.get::<_, i64>(1)?.to_string(),
436 row.get::<_, String>(2)?,
437 ))
438 })
439 .map_err(ThingdError::from)?;
440
441 tx.execute("DELETE FROM search_index", [])
443 .map_err(ThingdError::from)?;
444
445 for row in rows_objects {
446 let (collection, id, body) = row.map_err(ThingdError::from)?;
447 let text = extract_text_from_json(&body);
448 tx.execute(
449 "INSERT INTO search_index (collection, id, kind, text) VALUES (?1, ?2, 'object', ?3)",
450 params![collection, id, text],
451 )
452 .map_err(ThingdError::from)?;
453 }
454
455 for row in rows_events {
456 let (stream, sequence, body) = row.map_err(ThingdError::from)?;
457 let text = extract_text_from_json(&body);
458 tx.execute(
459 "INSERT INTO search_index (collection, id, kind, text) VALUES (?1, ?2, 'event', ?3)",
460 params![stream, sequence, text],
461 )
462 .map_err(ThingdError::from)?;
463 }
464
465 Ok(())
466 }
467}
468
469impl ObjectStore for SqliteThingStore {
470 fn put_object(&mut self, mut object: MemoryObject) -> ThingdResult<MemoryObject> {
471 let transaction = self.connection.transaction().map_err(ThingdError::from)?;
472
473 let row = transaction
475 .query_row(
476 r"
477 INSERT INTO objects (collection, id, body, version, created_at, updated_at)
478 VALUES (?1, ?2, ?3, 1, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'), strftime('%Y-%m-%dT%H:%M:%fZ', 'now'))
479 ON CONFLICT(collection, id) DO UPDATE SET
480 body = excluded.body,
481 version = version + 1,
482 updated_at = strftime('%Y-%m-%dT%H:%M:%fZ', 'now')
483 RETURNING version, created_at, updated_at
484 ",
485 params![&object.key.collection, &object.key.id, &object.body],
486 |row| {
487 Ok((
488 row.get::<_, i64>(0)?,
489 row.get::<_, String>(1)?,
490 row.get::<_, String>(2)?,
491 ))
492 },
493 )
494 .map_err(ThingdError::from)?;
495 object.version = u64::try_from(row.0).map_err(|e| ThingdError::Storage(e.to_string()))?;
496 object.created_at = row.1;
497 object.updated_at = row.2;
498
499 let text = extract_text_from_json(&object.body);
500 transaction
501 .execute(
502 "DELETE FROM search_index WHERE collection = ?1 AND id = ?2 AND kind = 'object'",
503 params![&object.key.collection, &object.key.id],
504 )
505 .map_err(ThingdError::from)?;
506 transaction
507 .execute(
508 "INSERT INTO search_index (collection, id, kind, text) VALUES (?1, ?2, 'object', ?3)",
509 params![&object.key.collection, &object.key.id, text],
510 )
511 .map_err(ThingdError::from)?;
512
513 transaction.commit().map_err(ThingdError::from)?;
514
515 Ok(object)
516 }
517
518 fn put_objects_batch(&mut self, objects: Vec<MemoryObject>) -> ThingdResult<Vec<MemoryObject>> {
519 if objects.is_empty() {
520 return Ok(Vec::new());
521 }
522
523 let transaction = self.connection.transaction().map_err(ThingdError::from)?;
524
525 let mut fts_updates: Vec<(String, String, String)> = Vec::with_capacity(objects.len());
526 let mut results = Vec::with_capacity(objects.len());
527 let mut param_values: Vec<String> = Vec::new();
528 let mut value_sql = String::new();
529
530 for (i, object) in objects.into_iter().enumerate() {
531 if i > 0 {
532 value_sql.push_str(", ");
533 }
534 let ci = i * 3 + 1;
535 let ii = i * 3 + 2;
536 let bi = i * 3 + 3;
537 let _ = std::fmt::Write::write_fmt(
538 &mut value_sql,
539 format_args!(
540 "(?{ci}, ?{ii}, ?{bi}, 1, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'), strftime('%Y-%m-%dT%H:%M:%fZ', 'now'))"
541 ),
542 );
543 let text = extract_text_from_json(&object.body);
544 fts_updates.push((object.key.collection.clone(), object.key.id.clone(), text));
545
546 param_values.push(object.key.collection.clone());
547 param_values.push(object.key.id.clone());
548 param_values.push(object.body.clone());
549
550 results.push(object);
551 }
552
553 let sql = format!(
554 "INSERT INTO objects (collection, id, body, version, created_at, updated_at) VALUES {value_sql} \
555 ON CONFLICT(collection, id) DO UPDATE SET \
556 body = excluded.body, version = version + 1, updated_at = strftime('%Y-%m-%dT%H:%M:%fZ', 'now') \
557 RETURNING version, created_at, updated_at"
558 );
559
560 let param_slices: Vec<&dyn rusqlite::types::ToSql> = param_values
561 .iter()
562 .map(|s| s as &dyn rusqlite::types::ToSql)
563 .collect();
564
565 let mut statement = transaction.prepare(&sql).map_err(ThingdError::from)?;
566 let rows = statement
567 .query_map(param_slices.as_slice(), |row| {
568 Ok((
569 row.get::<_, i64>(0)?,
570 row.get::<_, String>(1)?,
571 row.get::<_, String>(2)?,
572 ))
573 })
574 .map_err(ThingdError::from)?;
575
576 for (result, row) in results.iter_mut().zip(rows) {
577 let (version, created_at, updated_at) = row.map_err(ThingdError::from)?;
578 result.version =
579 u64::try_from(version).map_err(|e| ThingdError::Storage(e.to_string()))?;
580 result.created_at = created_at;
581 result.updated_at = updated_at;
582 }
583
584 drop(statement);
585
586 for (collection, id, text) in &fts_updates {
588 transaction
589 .execute(
590 "DELETE FROM search_index WHERE collection = ?1 AND id = ?2 AND kind = 'object'",
591 params![collection, id],
592 )
593 .map_err(ThingdError::from)?;
594 transaction
595 .execute(
596 "INSERT INTO search_index (collection, id, kind, text) VALUES (?1, ?2, 'object', ?3)",
597 params![collection, id, text],
598 )
599 .map_err(ThingdError::from)?;
600 }
601
602 transaction.commit().map_err(ThingdError::from)?;
603
604 Ok(results)
605 }
606
607 fn put_object_with_options(
608 &mut self,
609 mut object: MemoryObject,
610 options: crate::PutObjectOptions,
611 ) -> ThingdResult<MemoryObject> {
612 let transaction = self.connection.transaction().map_err(ThingdError::from)?;
613 let current_version = transaction
614 .query_row(
615 "SELECT version FROM objects WHERE collection = ?1 AND id = ?2",
616 params![&object.key.collection, &object.key.id],
617 |row| row.get::<_, i64>(0),
618 )
619 .optional()
620 .map_err(ThingdError::from)?;
621
622 if let Some(expected) = options.expected_version {
624 match current_version {
625 Some(actual) if u64::try_from(actual).unwrap_or(0) != expected => {
626 return Err(ThingdError::Conflict(format!(
627 "Version mismatch for {}/{}: expected {expected}, got {actual}",
628 object.key.collection, object.key.id,
629 )));
630 },
631 None => {
632 return Err(ThingdError::Conflict(format!(
633 "Version mismatch for {}/{}: expected {expected}, object does not exist",
634 object.key.collection, object.key.id,
635 )));
636 },
637 _ => {},
638 }
639 }
640
641 let version = current_version.map_or(Ok::<u64, ThingdError>(1), |existing| {
642 u64::try_from(existing)
643 .map(|existing| existing + 1)
644 .map_err(|error| ThingdError::Storage(error.to_string()))
645 })?;
646
647 object.version = version;
648 let stored_version = i64::try_from(object.version)
649 .map_err(|error| ThingdError::Storage(error.to_string()))?;
650
651 let timestamps = transaction
652 .query_row(
653 r"
654 INSERT INTO objects (collection, id, body, version, created_at, updated_at)
655 VALUES (?1, ?2, ?3, ?4, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'), strftime('%Y-%m-%dT%H:%M:%fZ', 'now'))
656 ON CONFLICT(collection, id) DO UPDATE SET
657 body = excluded.body,
658 version = excluded.version,
659 updated_at = strftime('%Y-%m-%dT%H:%M:%fZ', 'now')
660 RETURNING created_at, updated_at
661 ",
662 params![
663 &object.key.collection,
664 &object.key.id,
665 &object.body,
666 stored_version
667 ],
668 |row| {
669 Ok((
670 row.get::<_, String>(0)?,
671 row.get::<_, String>(1)?,
672 ))
673 },
674 )
675 .map_err(ThingdError::from)?;
676 object.created_at = timestamps.0;
677 object.updated_at = timestamps.1;
678
679 if options.index {
680 let text = extract_text_from_json(&object.body);
681 transaction
682 .execute(
683 "DELETE FROM search_index WHERE collection = ?1 AND id = ?2 AND kind = 'object'",
684 params![&object.key.collection, &object.key.id],
685 )
686 .map_err(ThingdError::from)?;
687 transaction
688 .execute(
689 "INSERT INTO search_index (collection, id, kind, text) VALUES (?1, ?2, 'object', ?3)",
690 params![&object.key.collection, &object.key.id, text],
691 )
692 .map_err(ThingdError::from)?;
693 }
694
695 transaction.commit().map_err(ThingdError::from)?;
696
697 Ok(object)
698 }
699
700 fn get_object(&self, collection: &str, id: &str) -> ThingdResult<Option<MemoryObject>> {
701 self.connection
702 .query_row(
703 "SELECT collection, id, body, version, created_at, updated_at FROM objects WHERE collection = ?1 AND id = ?2",
704 params![collection, id],
705 |row| {
706 let version = row.get::<_, i64>(3)?;
707
708 Ok(MemoryObject {
709 key: ObjectKey::new(row.get::<_, String>(0)?, row.get::<_, String>(1)?),
710 body: row.get(2)?,
711 version: u64::try_from(version).map_err(|error| {
712 rusqlite::Error::FromSqlConversionFailure(
713 3,
714 rusqlite::types::Type::Integer,
715 Box::new(error),
716 )
717 })?,
718 created_at: row.get::<_, String>(4).unwrap_or_default(),
719 updated_at: row.get::<_, String>(5).unwrap_or_default(),
720 })
721 },
722 )
723 .optional()
724 .map_err(ThingdError::from)
725 }
726
727 fn list_objects(
728 &self,
729 collections: Option<&[String]>,
730 options: &ListObjectsOptions,
731 ) -> ThingdResult<Vec<MemoryObject>> {
732 let mut conditions: Vec<String> = Vec::new();
735 let mut bound_values: Vec<Box<dyn rusqlite::types::ToSql>> = Vec::new();
736
737 let has_collections = collections.is_some_and(|c| !c.is_empty());
739 if has_collections {
740 let cols = collections.unwrap_or_default();
741 let placeholders = cols.iter().map(|_| "?").collect::<Vec<_>>().join(", ");
742 conditions.push(format!("collection IN ({placeholders})"));
743 for col in cols {
744 bound_values.push(Box::new(col.clone()));
745 }
746 }
747
748 for (key, value) in &options.filter {
750 if !key
752 .chars()
753 .all(|c| c.is_alphanumeric() || c == '_' || c == '.')
754 {
755 return Err(ThingdError::InvalidInput(format!(
756 "Invalid filter key: '{key}'. Only alphanumeric, underscore, and dot characters are allowed."
757 )));
758 }
759 conditions.push(format!("json_extract(body, '$.{key}') = ?"));
760 let sql_val: Box<dyn rusqlite::types::ToSql> = match value {
761 serde_json::Value::String(s) => Box::new(s.clone()),
762 serde_json::Value::Number(n) => {
763 if let Some(i) = n.as_i64() {
764 Box::new(i)
765 } else {
766 Box::new(n.as_f64().unwrap_or(0.0))
767 }
768 },
769 serde_json::Value::Bool(b) => Box::new(i64::from(*b)),
770 serde_json::Value::Null => Box::new(rusqlite::types::Null),
771 other => Box::new(other.to_string()),
772 };
773 bound_values.push(sql_val);
774 }
775
776 let where_clause = if conditions.is_empty() {
777 String::new()
778 } else {
779 format!("WHERE {}", conditions.join(" AND "))
780 };
781
782 let limit_clause = match (options.limit, options.offset) {
784 (Some(l), Some(o)) => {
785 bound_values.push(Box::new(i64::try_from(l).unwrap_or(i64::MAX)));
786 bound_values.push(Box::new(i64::try_from(o).unwrap_or(i64::MAX)));
787 " LIMIT ? OFFSET ?"
788 },
789 (Some(l), None) => {
790 bound_values.push(Box::new(i64::try_from(l).unwrap_or(i64::MAX)));
791 " LIMIT ?"
792 },
793 (None, Some(o)) => {
794 bound_values.push(Box::new(i64::try_from(o).unwrap_or(i64::MAX)));
795 " LIMIT -1 OFFSET ?"
796 },
797 (None, None) => "",
798 };
799
800 let order_clause = options.sort_by.as_ref().map_or_else(
801 || "ORDER BY collection, id".to_string(),
802 |sort_by| {
803 let col = match sort_by.field.as_str() {
804 "id" => "id",
805 "collection" => "collection",
806 "created_at" => "created_at",
807 "updated_at" => "updated_at",
808 "version" => "version",
809 _ => "collection, id",
810 };
811 let dir = match sort_by.direction {
812 crate::model::SortDirection::Asc => "ASC",
813 crate::model::SortDirection::Desc => "DESC",
814 };
815 format!("ORDER BY {col} {dir}")
816 },
817 );
818
819 let sql = format!(
820 "SELECT collection, id, body, version, created_at, updated_at FROM objects {where_clause} {order_clause} {limit_clause}"
821 );
822
823 let mut statement = self.connection.prepare(&sql).map_err(ThingdError::from)?;
824 let params: Vec<&dyn rusqlite::types::ToSql> =
825 bound_values.iter().map(AsRef::as_ref).collect();
826 let rows = statement
827 .query_map(params.as_slice(), row_to_object)
828 .map_err(ThingdError::from)?;
829
830 let mut objects = Vec::new();
831 for row in rows {
832 objects.push(row.map_err(ThingdError::from)?);
833 }
834 Ok(objects)
835 }
836
837 fn delete_object(&mut self, collection: &str, id: &str) -> ThingdResult<bool> {
838 let transaction = self.connection.transaction().map_err(ThingdError::from)?;
839 let changed = transaction
840 .execute(
841 "DELETE FROM objects WHERE collection = ?1 AND id = ?2",
842 params![collection, id],
843 )
844 .map_err(ThingdError::from)?;
845
846 if changed > 0 {
847 transaction
848 .execute(
849 "DELETE FROM search_index WHERE collection = ?1 AND id = ?2 AND kind = 'object'",
850 params![collection, id],
851 )
852 .map_err(ThingdError::from)?;
853 }
854
855 transaction.commit().map_err(ThingdError::from)?;
856 Ok(changed > 0)
857 }
858
859 fn delete_objects_batch(&mut self, keys: &[(String, String)]) -> ThingdResult<u64> {
860 use std::fmt::Write;
861
862 let transaction = self.connection.transaction().map_err(ThingdError::from)?;
863
864 if keys.is_empty() {
865 return Ok(0);
866 }
867
868 let mut total_deleted = 0u64;
869 for chunk in keys.chunks(500) {
871 let mut sql = String::from("DELETE FROM objects WHERE ");
872 let mut fts_sql = String::from("DELETE FROM search_index WHERE kind = 'object' AND (");
873 let mut param_values: Vec<String> = Vec::with_capacity(chunk.len() * 2);
874 for (i, (collection, id)) in chunk.iter().enumerate() {
875 if i > 0 {
876 sql.push_str(" OR ");
877 fts_sql.push_str(" OR ");
878 }
879 let ci = i * 2 + 1;
880 let ii = i * 2 + 2;
881 let _ = write!(sql, "(collection = ?{ci} AND id = ?{ii})");
882 let _ = write!(fts_sql, "(collection = ?{ci} AND id = ?{ii})");
883 param_values.push(collection.clone());
884 param_values.push(id.clone());
885 }
886 fts_sql.push(')');
887 let param_slices: Vec<&dyn rusqlite::types::ToSql> = param_values
888 .iter()
889 .map(|s| s as &dyn rusqlite::types::ToSql)
890 .collect();
891
892 let deleted = transaction
893 .execute(&sql, param_slices.as_slice())
894 .map_err(ThingdError::from)?;
895 transaction
896 .execute(&fts_sql, param_slices.as_slice())
897 .map_err(ThingdError::from)?;
898 total_deleted += deleted as u64;
899 }
900
901 transaction.commit().map_err(ThingdError::from)?;
902 Ok(total_deleted)
903 }
904
905 fn count_objects(&self) -> ThingdResult<u64> {
906 let count: i64 = self
907 .connection
908 .query_row("SELECT COUNT(*) FROM objects", [], |row| row.get(0))
909 .map_err(ThingdError::from)?;
910 Ok(u64::try_from(count).unwrap_or(0))
911 }
912
913 fn list_collections(&self) -> ThingdResult<Vec<String>> {
914 let mut statement = self
915 .connection
916 .prepare("SELECT DISTINCT collection FROM objects ORDER BY collection")
917 .map_err(ThingdError::from)?;
918 let rows = statement
919 .query_map([], |row| row.get::<_, String>(0))
920 .map_err(ThingdError::from)?;
921
922 let mut collections = Vec::new();
923 for row in rows {
924 collections.push(row.map_err(ThingdError::from)?);
925 }
926 Ok(collections)
927 }
928}
929
930impl EventLog for SqliteThingStore {
931 fn is_protected_stream(&self, stream: &str) -> bool {
932 stream == "__thingd:mcp:audit"
933 }
934 fn append_event(&mut self, mut event: MemoryEvent) -> ThingdResult<MemoryEvent> {
935 if !event.idempotency_key.is_empty()
937 && let Some(&existing_seq) = self
938 .event_idempotency_keys
939 .get(&(event.stream.clone(), event.idempotency_key.clone()))
940 {
941 let existing = self.connection.query_row(
942 "SELECT stream, event_type, body, sequence, created_at FROM events WHERE stream = ?1 AND sequence = ?2",
943 params![&event.stream, existing_seq.cast_signed()],
944 |row| {
945 Ok(MemoryEvent {
946 stream: row.get(0)?,
947 event_type: row.get(1)?,
948 body: row.get(2)?,
949 sequence: row.get::<_, i64>(3)?.cast_unsigned(),
950 created_at: row.get(4)?,
951 idempotency_key: event.idempotency_key.clone(),
952 })
953 },
954 ).map_err(ThingdError::from)?;
955 return Ok(existing);
956 }
957
958 let transaction = self.connection.transaction().map_err(ThingdError::from)?;
959
960 let (sequence, created_at): (i64, String) = transaction
961 .query_row(
962 r"
963 INSERT INTO events (stream, event_type, body, created_at)
964 VALUES (?1, ?2, ?3, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'))
965 RETURNING sequence, created_at
966 ",
967 params![&event.stream, &event.event_type, &event.body],
968 |row| Ok((row.get::<_, i64>(0)?, row.get::<_, String>(1)?)),
969 )
970 .map_err(ThingdError::from)?;
971
972 event.sequence =
973 u64::try_from(sequence).map_err(|error| ThingdError::Storage(error.to_string()))?;
974 event.created_at = created_at;
975
976 if !event.idempotency_key.is_empty() {
978 self.event_idempotency_keys.insert(
979 (event.stream.clone(), event.idempotency_key.clone()),
980 event.sequence,
981 );
982 }
983
984 let text = extract_text_from_json(&event.body);
985 let seq_str = sequence.to_string();
986 transaction
987 .execute(
988 "INSERT INTO search_index (collection, id, kind, text) VALUES (?1, ?2, 'event', ?3)",
989 params![&event.stream, seq_str, text],
990 )
991 .map_err(ThingdError::from)?;
992
993 transaction.commit().map_err(ThingdError::from)?;
994
995 Ok(event)
996 }
997
998 fn append_events_batch(&mut self, events: Vec<MemoryEvent>) -> ThingdResult<Vec<MemoryEvent>> {
999 if events.is_empty() {
1000 return Ok(Vec::new());
1001 }
1002
1003 let transaction = self.connection.transaction().map_err(ThingdError::from)?;
1004
1005 let mut param_values: Vec<String> = Vec::new();
1006 let mut value_sql = String::new();
1007 let mut results: Vec<MemoryEvent> = Vec::with_capacity(events.len());
1008
1009 for (i, event) in events.into_iter().enumerate() {
1010 if i > 0 {
1011 value_sql.push_str(", ");
1012 }
1013 let si = i * 3 + 1;
1014 let ti = i * 3 + 2;
1015 let bi = i * 3 + 3;
1016 let _ = std::fmt::Write::write_fmt(
1017 &mut value_sql,
1018 format_args!("(?{si}, ?{ti}, ?{bi}, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'))"),
1019 );
1020 param_values.push(event.stream.clone());
1021 param_values.push(event.event_type.clone());
1022 param_values.push(event.body.clone());
1023 results.push(event);
1024 }
1025
1026 let sql = format!(
1027 "INSERT INTO events (stream, event_type, body, created_at) VALUES {value_sql} \
1028 RETURNING sequence, created_at"
1029 );
1030
1031 let param_slices: Vec<&dyn rusqlite::types::ToSql> = param_values
1032 .iter()
1033 .map(|s| s as &dyn rusqlite::types::ToSql)
1034 .collect();
1035
1036 let mut statement = transaction.prepare(&sql).map_err(ThingdError::from)?;
1037 let rows = statement
1038 .query_map(param_slices.as_slice(), |row| {
1039 Ok((row.get::<_, i64>(0)?, row.get::<_, String>(1)?))
1040 })
1041 .map_err(ThingdError::from)?;
1042
1043 for (result, row) in results.iter_mut().zip(rows) {
1044 let (sequence, created_at) = row.map_err(ThingdError::from)?;
1045 result.sequence =
1046 u64::try_from(sequence).map_err(|error| ThingdError::Storage(error.to_string()))?;
1047 result.created_at = created_at;
1048
1049 let text = extract_text_from_json(&result.body);
1050 let seq_str = sequence.to_string();
1051 transaction
1052 .execute(
1053 "INSERT INTO search_index (collection, id, kind, text) VALUES (?1, ?2, 'event', ?3)",
1054 params![&result.stream, seq_str, text],
1055 )
1056 .map_err(ThingdError::from)?;
1057 }
1058
1059 drop(statement);
1060 transaction.commit().map_err(ThingdError::from)?;
1061
1062 Ok(results)
1063 }
1064
1065 fn list_events(
1066 &self,
1067 stream: Option<&str>,
1068 options: ListEventsOptions,
1069 ) -> ThingdResult<Vec<MemoryEvent>> {
1070 let mut events = Vec::new();
1071
1072 let mut sql = String::from(
1073 "SELECT stream, event_type, body, sequence, created_at FROM events WHERE 1=1",
1074 );
1075 let mut param_values: Vec<Box<dyn rusqlite::types::ToSql>> = Vec::new();
1076
1077 if let Some(stream) = stream {
1078 let idx = param_values.len() + 1;
1079 write!(sql, " AND stream = ?{idx}").unwrap();
1080 param_values.push(Box::new(stream.to_string()));
1081 }
1082
1083 if let Some(from_sequence) = options.from_sequence {
1084 let idx = param_values.len() + 1;
1085 write!(sql, " AND sequence > ?{idx}").unwrap();
1086 param_values.push(Box::new(from_sequence.cast_signed()));
1087 }
1088
1089 sql.push_str(" ORDER BY sequence");
1090
1091 if let Some(limit) = options.limit {
1092 sql.push_str(" LIMIT ?");
1093 param_values.push(Box::new(limit.cast_signed()));
1094 }
1095
1096 let mut statement = self.connection.prepare(&sql).map_err(ThingdError::from)?;
1097
1098 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
1099 param_values.iter().map(AsRef::as_ref).collect();
1100
1101 let rows = statement
1102 .query_map(param_refs.as_slice(), row_to_event)
1103 .map_err(ThingdError::from)?;
1104
1105 for row in rows {
1106 events.push(row.map_err(ThingdError::from)?);
1107 }
1108
1109 Ok(events)
1110 }
1111
1112 fn count_events(&self) -> ThingdResult<u64> {
1113 let count: i64 = self
1114 .connection
1115 .query_row("SELECT COUNT(*) FROM events", [], |row| row.get(0))
1116 .map_err(ThingdError::from)?;
1117 Ok(u64::try_from(count).unwrap_or(0))
1118 }
1119
1120 fn list_streams(&self) -> ThingdResult<Vec<String>> {
1121 let mut statement = self
1122 .connection
1123 .prepare("SELECT DISTINCT stream FROM events ORDER BY stream")
1124 .map_err(ThingdError::from)?;
1125 let rows = statement
1126 .query_map([], |row| row.get::<_, String>(0))
1127 .map_err(ThingdError::from)?;
1128
1129 let mut streams = Vec::new();
1130 for row in rows {
1131 streams.push(row.map_err(ThingdError::from)?);
1132 }
1133 Ok(streams)
1134 }
1135
1136 fn delete_last_event(&mut self, stream: &str) -> ThingdResult<Option<MemoryEvent>> {
1137 if self.is_protected_stream(stream) {
1138 return Err(ThingdError::Protected(format!(
1139 "stream '{stream}' is protected and cannot be modified"
1140 )));
1141 }
1142 let transaction = self.connection.transaction().map_err(ThingdError::from)?;
1143
1144 let result = transaction
1145 .query_row(
1146 r"
1147 DELETE FROM events
1148 WHERE sequence = (
1149 SELECT MAX(sequence) FROM events WHERE stream = ?1
1150 )
1151 RETURNING stream, event_type, body, sequence, created_at
1152 ",
1153 params![stream],
1154 row_to_event,
1155 )
1156 .optional()
1157 .map_err(ThingdError::from)?;
1158
1159 if let Some(ref event) = result {
1160 let seq_str = event.sequence.to_string();
1162 transaction
1163 .execute(
1164 "DELETE FROM search_index WHERE collection = ?1 AND kind = 'event' AND id = ?2",
1165 params![stream, seq_str],
1166 )
1167 .map_err(ThingdError::from)?;
1168 }
1169
1170 transaction.commit().map_err(ThingdError::from)?;
1171 Ok(result)
1172 }
1173
1174 fn delete_stream(&mut self, stream: &str) -> ThingdResult<u64> {
1175 if self.is_protected_stream(stream) {
1176 return Err(ThingdError::Protected(format!(
1177 "stream '{stream}' is protected and cannot be modified"
1178 )));
1179 }
1180 let transaction = self.connection.transaction().map_err(ThingdError::from)?;
1181
1182 let count = transaction
1183 .execute("DELETE FROM events WHERE stream = ?1", params![stream])
1184 .map_err(ThingdError::from)?;
1185
1186 transaction
1187 .execute(
1188 "DELETE FROM search_index WHERE collection = ?1 AND kind = 'event'",
1189 params![stream],
1190 )
1191 .map_err(ThingdError::from)?;
1192
1193 transaction.commit().map_err(ThingdError::from)?;
1194 Ok(count as u64)
1195 }
1196}
1197
1198impl QueueStore for SqliteThingStore {
1199 fn push_job(&mut self, job: QueueJob) -> ThingdResult<QueueJob> {
1200 let transaction = self
1201 .connection
1202 .transaction_with_behavior(TransactionBehavior::Immediate)
1203 .map_err(ThingdError::from)?;
1204
1205 if let Some(existing) = transaction
1206 .query_row(
1207 &queue_job_select_sql("WHERE queue = ?1 AND id = ?2"),
1208 params![&job.queue, &job.id],
1209 row_to_queue_job,
1210 )
1211 .optional()
1212 .map_err(ThingdError::from)?
1213 {
1214 transaction.commit().map_err(ThingdError::from)?;
1215 return Ok(existing);
1216 }
1217
1218 let created_at: String = transaction
1220 .query_row(
1221 r"
1222 INSERT INTO queue_jobs (
1223 queue,
1224 id,
1225 body,
1226 attempts,
1227 max_attempts,
1228 status,
1229 available_at_ms,
1230 leased_at_ms,
1231 lease_expires_at_ms,
1232 completed_at_ms,
1233 dead_at_ms,
1234 created_at,
1235 updated_at
1236 )
1237 VALUES (
1238 ?1,
1239 ?2,
1240 ?3,
1241 ?4,
1242 ?5,
1243 ?6,
1244 ?7,
1245 ?8,
1246 ?9,
1247 ?10,
1248 ?11,
1249 strftime('%Y-%m-%dT%H:%M:%fZ', 'now'),
1250 strftime('%Y-%m-%dT%H:%M:%fZ', 'now')
1251 )
1252 RETURNING created_at
1253 ",
1254 params![
1255 &job.queue,
1256 &job.id,
1257 &job.body,
1258 u32_to_i64(job.attempts),
1259 u32_to_i64(job.max_attempts),
1260 status_to_str(job.status),
1261 job.available_at_ms,
1262 job.leased_at_ms,
1263 job.lease_expires_at_ms,
1264 job.completed_at_ms,
1265 job.dead_at_ms
1266 ],
1267 |row| row.get(0),
1268 )
1269 .map_err(ThingdError::from)?;
1270
1271 transaction.commit().map_err(ThingdError::from)?;
1272
1273 Ok(QueueJob { created_at, ..job })
1274 }
1275
1276 fn push_jobs_batch(&mut self, jobs: Vec<QueueJob>) -> ThingdResult<Vec<QueueJob>> {
1277 let transaction = self
1278 .connection
1279 .transaction_with_behavior(TransactionBehavior::Immediate)
1280 .map_err(ThingdError::from)?;
1281
1282 let mut results = Vec::with_capacity(jobs.len());
1283 for job in jobs {
1284 if let Some(existing) = transaction
1285 .query_row(
1286 &queue_job_select_sql("WHERE queue = ?1 AND id = ?2"),
1287 params![&job.queue, &job.id],
1288 row_to_queue_job,
1289 )
1290 .optional()
1291 .map_err(ThingdError::from)?
1292 {
1293 results.push(existing);
1294 continue;
1295 }
1296
1297 let created_at: String = transaction
1299 .query_row(
1300 r"
1301 INSERT INTO queue_jobs (
1302 queue,
1303 id,
1304 body,
1305 attempts,
1306 max_attempts,
1307 status,
1308 available_at_ms,
1309 leased_at_ms,
1310 lease_expires_at_ms,
1311 completed_at_ms,
1312 dead_at_ms,
1313 created_at,
1314 updated_at
1315 )
1316 VALUES (
1317 ?1,
1318 ?2,
1319 ?3,
1320 ?4,
1321 ?5,
1322 ?6,
1323 ?7,
1324 ?8,
1325 ?9,
1326 ?10,
1327 ?11,
1328 strftime('%Y-%m-%dT%H:%M:%fZ', 'now'),
1329 strftime('%Y-%m-%dT%H:%M:%fZ', 'now')
1330 )
1331 RETURNING created_at
1332 ",
1333 params![
1334 &job.queue,
1335 &job.id,
1336 &job.body,
1337 u32_to_i64(job.attempts),
1338 u32_to_i64(job.max_attempts),
1339 status_to_str(job.status),
1340 job.available_at_ms,
1341 job.leased_at_ms,
1342 job.lease_expires_at_ms,
1343 job.completed_at_ms,
1344 job.dead_at_ms
1345 ],
1346 |row| row.get(0),
1347 )
1348 .map_err(ThingdError::from)?;
1349
1350 results.push(QueueJob { created_at, ..job });
1351 }
1352
1353 transaction.commit().map_err(ThingdError::from)?;
1354
1355 Ok(results)
1356 }
1357
1358 fn claim_job_with_options(
1359 &mut self,
1360 queue: &str,
1361 options: QueueClaimOptions,
1362 ) -> ThingdResult<Option<QueueJob>> {
1363 let transaction = self
1364 .connection
1365 .transaction_with_behavior(TransactionBehavior::Immediate)
1366 .map_err(ThingdError::from)?;
1367
1368 release_expired_leases(&transaction, queue)?;
1369 let now = unix_timestamp_millis();
1370 let Some(mut job) = transaction
1371 .query_row(
1372 &queue_job_select_sql(
1373 "WHERE queue = ?1 AND status = 'ready' AND available_at_ms <= ?2 ORDER BY created_at LIMIT 1",
1374 ),
1375 params![queue, now],
1376 row_to_queue_job,
1377 )
1378 .optional()
1379 .map_err(ThingdError::from)?
1380 else {
1381 transaction.commit().map_err(ThingdError::from)?;
1382 return Ok(None);
1383 };
1384
1385 job.status = QueueJobStatus::Leased;
1386 job.attempts += 1;
1387 job.leased_at_ms = Some(now);
1388 job.lease_expires_at_ms = Some(now.saturating_add(u64_to_i64(options.lease_ms)));
1389
1390 transaction
1391 .execute(
1392 r"
1393 UPDATE queue_jobs
1394 SET attempts = ?3,
1395 status = ?4,
1396 leased_at_ms = ?5,
1397 lease_expires_at_ms = ?6,
1398 updated_at = strftime('%Y-%m-%dT%H:%M:%fZ', 'now')
1399 WHERE queue = ?1 AND id = ?2
1400 ",
1401 params![
1402 &job.queue,
1403 &job.id,
1404 u32_to_i64(job.attempts),
1405 status_to_str(job.status),
1406 job.leased_at_ms,
1407 job.lease_expires_at_ms
1408 ],
1409 )
1410 .map_err(ThingdError::from)?;
1411
1412 transaction.commit().map_err(ThingdError::from)?;
1413 Ok(Some(job))
1414 }
1415
1416 fn ack_job(&mut self, queue: &str, id: &str) -> ThingdResult<Option<QueueJob>> {
1417 let transaction = self
1418 .connection
1419 .transaction_with_behavior(TransactionBehavior::Immediate)
1420 .map_err(ThingdError::from)?;
1421
1422 let Some(mut job) = transaction
1423 .query_row(
1424 &queue_job_select_sql("WHERE queue = ?1 AND id = ?2"),
1425 params![queue, id],
1426 row_to_queue_job,
1427 )
1428 .optional()
1429 .map_err(ThingdError::from)?
1430 else {
1431 transaction.commit().map_err(ThingdError::from)?;
1432 return Ok(None);
1433 };
1434
1435 if job.status != QueueJobStatus::Leased {
1436 return Err(ThingdError::Conflict(format!(
1437 "job {id} must be leased before ack"
1438 )));
1439 }
1440
1441 job.status = QueueJobStatus::Completed;
1442 job.completed_at_ms = Some(unix_timestamp_millis());
1443 transaction
1444 .execute(
1445 r"
1446 UPDATE queue_jobs
1447 SET status = ?3,
1448 completed_at_ms = ?4,
1449 updated_at = strftime('%Y-%m-%dT%H:%M:%fZ', 'now')
1450 WHERE queue = ?1 AND id = ?2
1451 ",
1452 params![queue, id, status_to_str(job.status), job.completed_at_ms],
1453 )
1454 .map_err(ThingdError::from)?;
1455
1456 transaction.commit().map_err(ThingdError::from)?;
1457 Ok(Some(job))
1458 }
1459
1460 fn claim_and_ack(
1461 &mut self,
1462 queue: &str,
1463 options: QueueClaimOptions,
1464 ) -> ThingdResult<Option<QueueJob>> {
1465 let transaction = self
1466 .connection
1467 .transaction_with_behavior(TransactionBehavior::Immediate)
1468 .map_err(ThingdError::from)?;
1469
1470 release_expired_leases(&transaction, queue)?;
1471 let now = unix_timestamp_millis();
1472 let Some(mut job) = transaction
1473 .query_row(
1474 &queue_job_select_sql(
1475 "WHERE queue = ?1 AND status = 'ready' AND available_at_ms <= ?2 ORDER BY created_at LIMIT 1",
1476 ),
1477 params![queue, now],
1478 row_to_queue_job,
1479 )
1480 .optional()
1481 .map_err(ThingdError::from)?
1482 else {
1483 transaction.commit().map_err(ThingdError::from)?;
1484 return Ok(None);
1485 };
1486
1487 job.status = QueueJobStatus::Leased;
1489 job.attempts += 1;
1490 job.leased_at_ms = Some(now);
1491 job.lease_expires_at_ms = Some(now.saturating_add(u64_to_i64(options.lease_ms)));
1492
1493 transaction
1494 .execute(
1495 r"
1496 UPDATE queue_jobs
1497 SET attempts = ?3,
1498 status = ?4,
1499 leased_at_ms = ?5,
1500 lease_expires_at_ms = ?6,
1501 updated_at = strftime('%Y-%m-%dT%H:%M:%fZ', 'now')
1502 WHERE queue = ?1 AND id = ?2
1503 ",
1504 params![
1505 &job.queue,
1506 &job.id,
1507 u32_to_i64(job.attempts),
1508 status_to_str(job.status),
1509 job.leased_at_ms,
1510 job.lease_expires_at_ms
1511 ],
1512 )
1513 .map_err(ThingdError::from)?;
1514
1515 job.status = QueueJobStatus::Completed;
1517 job.completed_at_ms = Some(unix_timestamp_millis());
1518 transaction
1519 .execute(
1520 r"
1521 UPDATE queue_jobs
1522 SET status = ?3,
1523 completed_at_ms = ?4,
1524 updated_at = strftime('%Y-%m-%dT%H:%M:%fZ', 'now')
1525 WHERE queue = ?1 AND id = ?2
1526 ",
1527 params![
1528 &job.queue,
1529 &job.id,
1530 status_to_str(job.status),
1531 job.completed_at_ms
1532 ],
1533 )
1534 .map_err(ThingdError::from)?;
1535
1536 transaction.commit().map_err(ThingdError::from)?;
1537 Ok(Some(job))
1538 }
1539
1540 fn nack_job_with_options(
1541 &mut self,
1542 queue: &str,
1543 id: &str,
1544 options: QueueNackOptions,
1545 ) -> ThingdResult<Option<QueueJob>> {
1546 let transaction = self
1547 .connection
1548 .transaction_with_behavior(TransactionBehavior::Immediate)
1549 .map_err(ThingdError::from)?;
1550
1551 let Some(mut job) = transaction
1552 .query_row(
1553 &queue_job_select_sql("WHERE queue = ?1 AND id = ?2"),
1554 params![queue, id],
1555 row_to_queue_job,
1556 )
1557 .optional()
1558 .map_err(ThingdError::from)?
1559 else {
1560 transaction.commit().map_err(ThingdError::from)?;
1561 return Ok(None);
1562 };
1563
1564 if job.status != QueueJobStatus::Leased {
1565 return Err(ThingdError::Conflict(format!(
1566 "job {id} must be leased before nack"
1567 )));
1568 }
1569
1570 let now = unix_timestamp_millis();
1571 job.leased_at_ms = None;
1572 job.lease_expires_at_ms = None;
1573
1574 job.status = if job.attempts >= job.max_attempts {
1575 job.dead_at_ms = Some(now);
1576 QueueJobStatus::Dead
1577 } else {
1578 job.available_at_ms = now.saturating_add(u64_to_i64(options.delay_ms));
1579 QueueJobStatus::Ready
1580 };
1581
1582 if !options.error.is_empty() {
1583 job.last_error = options.error;
1584 }
1585
1586 transaction
1587 .execute(
1588 r"
1589 UPDATE queue_jobs
1590 SET attempts = ?3,
1591 status = ?4,
1592 available_at_ms = ?5,
1593 leased_at_ms = NULL,
1594 lease_expires_at_ms = NULL,
1595 dead_at_ms = ?6,
1596 last_error = ?7,
1597 updated_at = strftime('%Y-%m-%dT%H:%M:%fZ', 'now')
1598 WHERE queue = ?1 AND id = ?2
1599 ",
1600 params![
1601 queue,
1602 id,
1603 u32_to_i64(job.attempts),
1604 status_to_str(job.status),
1605 job.available_at_ms,
1606 job.dead_at_ms,
1607 job.last_error
1608 ],
1609 )
1610 .map_err(ThingdError::from)?;
1611
1612 transaction.commit().map_err(ThingdError::from)?;
1613 Ok(Some(job))
1614 }
1615
1616 fn list_jobs(&self, queue: &str) -> ThingdResult<Vec<QueueJob>> {
1617 let mut statement = self
1618 .connection
1619 .prepare(&queue_job_select_sql(
1620 "WHERE queue = ?1 ORDER BY created_at LIMIT 1000",
1621 ))
1622 .map_err(ThingdError::from)?;
1623 let rows = statement
1624 .query_map(params![queue], row_to_queue_job)
1625 .map_err(ThingdError::from)?;
1626
1627 let mut jobs = Vec::new();
1628 for row in rows {
1629 jobs.push(row.map_err(ThingdError::from)?);
1630 }
1631
1632 Ok(jobs)
1633 }
1634
1635 fn list_dead_jobs(&self, queue: &str) -> ThingdResult<Vec<QueueJob>> {
1636 let mut statement = self
1637 .connection
1638 .prepare(&queue_job_select_sql(
1639 "WHERE queue = ?1 AND status = 'dead' ORDER BY created_at LIMIT 1000",
1640 ))
1641 .map_err(ThingdError::from)?;
1642 let rows = statement
1643 .query_map(params![queue], row_to_queue_job)
1644 .map_err(ThingdError::from)?;
1645
1646 let mut jobs = Vec::new();
1647 for row in rows {
1648 jobs.push(row.map_err(ThingdError::from)?);
1649 }
1650
1651 Ok(jobs)
1652 }
1653
1654 fn list_queues(&self) -> ThingdResult<Vec<String>> {
1655 let mut statement = self
1656 .connection
1657 .prepare("SELECT DISTINCT queue FROM queue_jobs ORDER BY queue")
1658 .map_err(ThingdError::from)?;
1659 let rows = statement
1660 .query_map([], |row| row.get::<_, String>(0))
1661 .map_err(ThingdError::from)?;
1662
1663 let mut queues = Vec::new();
1664 for row in rows {
1665 queues.push(row.map_err(ThingdError::from)?);
1666 }
1667 Ok(queues)
1668 }
1669
1670 fn count_active_jobs(&self) -> ThingdResult<u64> {
1671 let count = self
1672 .connection
1673 .query_row(
1674 "SELECT COUNT(id) FROM queue_jobs WHERE status != 'dead'",
1675 [],
1676 |row| row.get::<_, i64>(0),
1677 )
1678 .map_err(ThingdError::from)?;
1679 Ok(u64::try_from(count).unwrap_or(0))
1680 }
1681
1682 fn count_dead_jobs(&self) -> ThingdResult<u64> {
1683 let count = self
1684 .connection
1685 .query_row(
1686 "SELECT COUNT(id) FROM queue_jobs WHERE status = 'dead'",
1687 [],
1688 |row| row.get::<_, i64>(0),
1689 )
1690 .map_err(ThingdError::from)?;
1691 Ok(u64::try_from(count).unwrap_or(0))
1692 }
1693}
1694
1695impl crate::store::Searcher for SqliteThingStore {
1696 #[allow(clippy::too_many_lines)]
1697 fn search(
1698 &self,
1699 query: &str,
1700 options: crate::SearchOptions,
1701 ) -> ThingdResult<Vec<crate::SearchHit>> {
1702 let sanitized = sanitize_fts_query(query);
1703 if sanitized.is_empty() {
1704 return Ok(Vec::new());
1705 }
1706 if query.chars().filter(|c| c.is_alphanumeric()).count() < 2 {
1708 return Ok(Vec::new());
1709 }
1710
1711 let mut sql = String::from(
1712 r"
1713 SELECT
1714 s.kind,
1715 s.collection,
1716 s.id,
1717 s.text,
1718 o.body AS object_body,
1719 o.version AS object_version,
1720 o.created_at AS object_created_at,
1721 o.updated_at AS object_updated_at,
1722 e.event_type AS event_type,
1723 e.body AS event_body,
1724 e.created_at AS event_created_at,
1725 bm25(search_index) AS bm25_score,
1726 (strftime('%s', 'now') - strftime('%s', coalesce(o.created_at, e.created_at))) AS age_seconds
1727 FROM search_index s
1728 LEFT JOIN objects o ON s.kind = 'object' AND s.collection = o.collection AND s.id = o.id
1729 LEFT JOIN events e ON s.kind = 'event' AND s.collection = e.stream AND s.id = CAST(e.sequence AS TEXT)
1730 WHERE search_index MATCH ?1
1731 AND (s.kind != 'object' OR o.collection IS NOT NULL)
1732 AND (s.kind != 'event' OR e.stream IS NOT NULL)
1733 ",
1734 );
1735
1736 let mut params: Vec<Box<dyn rusqlite::types::ToSql>> = vec![Box::new(sanitized)];
1737
1738 if let Some(ref collections) = options.collections
1740 && !collections.is_empty()
1741 {
1742 let placeholders: Vec<String> = (0..collections.len())
1743 .map(|i| format!("?{}", params.len() + i + 1))
1744 .collect();
1745 write!(sql, " AND s.collection IN ({})", placeholders.join(",")).unwrap();
1746 for coll in collections {
1747 params.push(Box::new(coll.clone()));
1748 }
1749 }
1750
1751 let effective_limit = options.limit.map_or(100, |l| l.min(1000));
1754
1755 sql.push_str(" ORDER BY bm25_score");
1756
1757 if options.filter.is_none() {
1759 write!(sql, " LIMIT ?{}", params.len() + 1).unwrap();
1760 params.push(Box::new(i64::try_from(effective_limit).unwrap_or(100)));
1761 } else {
1762 let fetch_limit = (effective_limit * 3).min(1000);
1764 write!(sql, " LIMIT ?{}", params.len() + 1).unwrap();
1765 params.push(Box::new(i64::try_from(fetch_limit).unwrap_or(1000)));
1766 }
1767
1768 let mut statement = self.connection.prepare(&sql).map_err(ThingdError::from)?;
1769
1770 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
1771 params.iter().map(AsRef::as_ref).collect();
1772
1773 let rows = statement
1774 .query_map(param_refs.as_slice(), |row| {
1775 let kind: String = row.get(0)?;
1776 let collection: String = row.get(1)?;
1777 let id: String = row.get(2)?;
1778 let text: String = row.get(3)?;
1779 let bm25_score: f64 = row.get(11)?;
1780 let age_seconds: Option<i64> = row.get(12)?;
1781
1782 let relevance_score = -bm25_score;
1783 let age =
1784 f64::from(i32::try_from(age_seconds.unwrap_or(0).max(0)).unwrap_or(i32::MAX));
1785 let recency_factor = 1.0 / (1.0 + age / 86400.0);
1786 let score = relevance_score * recency_factor;
1787
1788 let (body, version, created_at, updated_at, event_type) = if kind == "object" {
1789 let object_body: String = row.get(4)?;
1790 let object_version: i64 = row.get(5)?;
1791 let object_created_at: String = row.get(6)?;
1792 let object_updated_at: String = row.get(7)?;
1793 (
1794 object_body,
1795 Some(object_version.cast_unsigned()),
1796 object_created_at,
1797 Some(object_updated_at),
1798 None,
1799 )
1800 } else {
1801 let event_type_val: String = row.get(8)?;
1802 let event_body: String = row.get(9)?;
1803 let event_created_at: String = row.get(10)?;
1804 (
1805 event_body,
1806 None,
1807 event_created_at,
1808 None,
1809 Some(event_type_val),
1810 )
1811 };
1812
1813 Ok(crate::SearchHit {
1814 kind,
1815 collection,
1816 id,
1817 text,
1818 score,
1819 body,
1820 version,
1821 created_at,
1822 updated_at,
1823 event_type,
1824 })
1825 })
1826 .map_err(ThingdError::from)?;
1827
1828 let mut hits = Vec::new();
1829 for row in rows {
1830 let hit = row.map_err(ThingdError::from)?;
1831
1832 if let Some(ref filter) = options.filter
1834 && !matches_filter(&hit.body, filter)
1835 {
1836 continue;
1837 }
1838
1839 hits.push(hit);
1840 }
1841
1842 hits.sort_by(|a, b| {
1844 b.score
1845 .partial_cmp(&a.score)
1846 .unwrap_or(std::cmp::Ordering::Equal)
1847 });
1848
1849 if let Some(limit) = options.limit {
1851 hits.truncate(limit);
1852 }
1853
1854 Ok(hits)
1855 }
1856}
1857
1858impl crate::store::LinkStore for SqliteThingStore {
1859 fn create_link(&mut self, link: crate::Link) -> ThingdResult<crate::Link> {
1860 let id = uuid::Uuid::new_v4().to_string();
1861 self.connection
1862 .execute(
1863 r"
1864 INSERT INTO links (id, from_ref, type, to_ref, weight, metadata_json, created_at)
1865 VALUES (?1, ?2, ?3, ?4, ?5, ?6, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'))
1866 ",
1867 params![
1868 id,
1869 link.from_ref,
1870 link.link_type,
1871 link.to_ref,
1872 link.weight,
1873 link.metadata_json
1874 ],
1875 )
1876 .map_err(ThingdError::from)?;
1877
1878 let created_at: String = self
1879 .connection
1880 .query_row(
1881 "SELECT created_at FROM links WHERE id = ?1",
1882 params![id],
1883 |row| row.get(0),
1884 )
1885 .map_err(ThingdError::from)?;
1886
1887 Ok(crate::Link {
1888 id,
1889 from_ref: link.from_ref,
1890 link_type: link.link_type,
1891 to_ref: link.to_ref,
1892 weight: link.weight,
1893 metadata_json: link.metadata_json,
1894 created_at,
1895 })
1896 }
1897
1898 fn delete_link(&mut self, id: &str) -> ThingdResult<bool> {
1899 let changed = self
1900 .connection
1901 .execute("DELETE FROM links WHERE id = ?1", params![id])
1902 .map_err(ThingdError::from)?;
1903 Ok(changed > 0)
1904 }
1905
1906 fn get_link(&self, id: &str) -> ThingdResult<Option<crate::Link>> {
1907 self.connection
1908 .query_row(
1909 "SELECT id, from_ref, type, to_ref, weight, metadata_json, created_at FROM links WHERE id = ?1",
1910 params![id],
1911 row_to_link,
1912 )
1913 .optional()
1914 .map_err(ThingdError::from)
1915 }
1916
1917 fn get_neighbors(
1918 &self,
1919 reference: &str,
1920 direction: crate::LinkDirection,
1921 options: crate::LinkQueryOptions,
1922 ) -> ThingdResult<Vec<crate::Link>> {
1923 let (where_clause, param_value): (&str, String) = match direction {
1924 crate::LinkDirection::Outgoing => ("WHERE from_ref = ?1", reference.to_string()),
1925 crate::LinkDirection::Incoming => ("WHERE to_ref = ?1", reference.to_string()),
1926 crate::LinkDirection::Both => (
1927 "WHERE (from_ref = ?1 OR to_ref = ?1)",
1928 reference.to_string(),
1929 ),
1930 };
1931
1932 let (type_filter_sql, type_param) = options.link_type.as_ref().map_or_else(
1934 || (String::new(), None),
1935 |t| (" AND type = ?2".to_string(), Some(t.clone())),
1936 );
1937
1938 let limit_clause = options
1939 .limit
1940 .map_or_else(String::new, |_| " LIMIT ?".to_string());
1941
1942 let mut params: Vec<Box<dyn rusqlite::types::ToSql>> = Vec::new();
1943 params.push(Box::new(param_value));
1944 if let Some(ref t) = type_param {
1945 params.push(Box::new(t.clone()));
1946 }
1947 if let Some(l) = options.limit {
1948 params.push(Box::new(i64::try_from(l).unwrap_or(1000)));
1949 }
1950
1951 let sql = format!(
1952 "SELECT id, from_ref, type, to_ref, weight, metadata_json, created_at FROM links {where_clause}{type_filter_sql}{limit_clause}"
1953 );
1954
1955 let mut statement = self.connection.prepare(&sql).map_err(ThingdError::from)?;
1956
1957 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
1958 params.iter().map(AsRef::as_ref).collect();
1959
1960 let rows = statement
1961 .query_map(param_refs.as_slice(), row_to_link)
1962 .map_err(ThingdError::from)?;
1963
1964 let mut links = Vec::new();
1965 for row in rows {
1966 links.push(row.map_err(ThingdError::from)?);
1967 }
1968
1969 Ok(links)
1970 }
1971
1972 fn count_links(&self) -> ThingdResult<u64> {
1973 let count: i64 = self
1974 .connection
1975 .query_row("SELECT COUNT(*) FROM links", [], |row| row.get(0))
1976 .map_err(ThingdError::from)?;
1977 Ok(u64::try_from(count).unwrap_or(0))
1978 }
1979}
1980
1981fn row_to_link(row: &rusqlite::Row<'_>) -> rusqlite::Result<crate::Link> {
1982 Ok(crate::Link {
1983 id: row.get(0)?,
1984 from_ref: row.get(1)?,
1985 link_type: row.get(2)?,
1986 to_ref: row.get(3)?,
1987 weight: row.get(4)?,
1988 metadata_json: row.get(5)?,
1989 created_at: row.get(6)?,
1990 })
1991}
1992
1993fn row_to_object(row: &rusqlite::Row<'_>) -> rusqlite::Result<MemoryObject> {
1994 let version = row.get::<_, i64>(3)?;
1995
1996 Ok(MemoryObject {
1997 key: ObjectKey::new(row.get::<_, String>(0)?, row.get::<_, String>(1)?),
1998 body: row.get(2)?,
1999 version: u64::try_from(version).map_err(|error| {
2000 rusqlite::Error::FromSqlConversionFailure(
2001 3,
2002 rusqlite::types::Type::Integer,
2003 Box::new(error),
2004 )
2005 })?,
2006 created_at: row.get::<_, String>(4).unwrap_or_default(),
2007 updated_at: row.get::<_, String>(5).unwrap_or_default(),
2008 })
2009}
2010
2011fn row_to_event(row: &rusqlite::Row<'_>) -> rusqlite::Result<MemoryEvent> {
2012 let sequence = row.get::<_, i64>(3)?;
2013
2014 Ok(MemoryEvent {
2015 stream: row.get(0)?,
2016 event_type: row.get(1)?,
2017 body: row.get(2)?,
2018 sequence: u64::try_from(sequence).map_err(|error| {
2019 rusqlite::Error::FromSqlConversionFailure(
2020 3,
2021 rusqlite::types::Type::Integer,
2022 Box::new(error),
2023 )
2024 })?,
2025 created_at: row.get::<_, String>(4).unwrap_or_default(),
2026 idempotency_key: String::new(),
2027 })
2028}
2029
2030fn queue_job_select_sql(predicate: &str) -> String {
2031 format!(
2032 "SELECT queue, id, body, attempts, max_attempts, status, available_at_ms, leased_at_ms, lease_expires_at_ms, completed_at_ms, dead_at_ms, created_at, last_error FROM queue_jobs {predicate}"
2033 )
2034}
2035
2036fn row_to_queue_job(row: &rusqlite::Row<'_>) -> rusqlite::Result<QueueJob> {
2037 let attempts = row.get::<_, i64>(3)?;
2038 let max_attempts = row.get::<_, i64>(4)?;
2039 let status = row.get::<_, String>(5)?;
2040
2041 Ok(QueueJob {
2042 queue: row.get(0)?,
2043 id: row.get(1)?,
2044 body: row.get(2)?,
2045 attempts: u32::try_from(attempts).map_err(|error| {
2046 rusqlite::Error::FromSqlConversionFailure(
2047 3,
2048 rusqlite::types::Type::Integer,
2049 Box::new(error),
2050 )
2051 })?,
2052 max_attempts: u32::try_from(max_attempts).map_err(|error| {
2053 rusqlite::Error::FromSqlConversionFailure(
2054 4,
2055 rusqlite::types::Type::Integer,
2056 Box::new(error),
2057 )
2058 })?,
2059 status: match status.as_str() {
2060 "ready" => QueueJobStatus::Ready,
2061 "leased" => QueueJobStatus::Leased,
2062 "completed" => QueueJobStatus::Completed,
2063 "dead" => QueueJobStatus::Dead,
2064 _other => {
2065 return Err(rusqlite::Error::FromSqlConversionFailure(
2066 5,
2067 rusqlite::types::Type::Text,
2068 Box::new(std::fmt::Error),
2069 ));
2070 },
2071 },
2072 available_at_ms: row.get(6)?,
2073 leased_at_ms: row.get(7)?,
2074 lease_expires_at_ms: row.get(8)?,
2075 completed_at_ms: row.get(9)?,
2076 dead_at_ms: row.get(10)?,
2077 created_at: row.get::<_, String>(11).unwrap_or_default(),
2078 last_error: row.get::<_, String>(12).unwrap_or_default(),
2079 })
2080}
2081
2082fn release_expired_leases(connection: &rusqlite::Connection, queue: &str) -> ThingdResult<()> {
2083 connection
2084 .execute(
2085 r"
2086 UPDATE queue_jobs
2087 SET status = 'ready',
2088 leased_at_ms = NULL,
2089 lease_expires_at_ms = NULL,
2090 updated_at = strftime('%Y-%m-%dT%H:%M:%fZ', 'now')
2091 WHERE queue = ?1
2092 AND status = 'leased'
2093 AND lease_expires_at_ms IS NOT NULL
2094 AND lease_expires_at_ms <= ?2
2095 ",
2096 params![queue, unix_timestamp_millis()],
2097 )
2098 .map_err(ThingdError::from)?;
2099
2100 Ok(())
2101}
2102
2103const fn status_to_str(status: QueueJobStatus) -> &'static str {
2104 match status {
2105 QueueJobStatus::Ready => "ready",
2106 QueueJobStatus::Leased => "leased",
2107 QueueJobStatus::Completed => "completed",
2108 QueueJobStatus::Dead => "dead",
2109 }
2110}
2111
2112fn u32_to_i64(value: u32) -> i64 {
2113 i64::from(value)
2114}
2115
2116fn sanitize_fts_query(query: &str) -> String {
2117 let mut cleaned = String::new();
2118 let normalized: String = query
2119 .chars()
2120 .map(|c| {
2121 if c.is_alphanumeric() || c.is_whitespace() {
2122 c
2123 } else {
2124 ' '
2125 }
2126 })
2127 .collect();
2128
2129 for word in normalized.split_whitespace() {
2130 if !word.is_empty() {
2131 if !cleaned.is_empty() {
2132 cleaned.push(' ');
2133 }
2134 cleaned.push_str(word);
2135 cleaned.push('*');
2136 }
2137 }
2138 cleaned
2139}
2140
2141fn extract_text_from_json(json_str: &str) -> String {
2142 serde_json::from_str::<serde_json::Value>(json_str).map_or_else(
2143 |_| json_str.to_string(),
2144 |value| {
2145 let mut out = String::new();
2146 collect_strings(&value, &mut out);
2147 out.trim().to_string()
2148 },
2149 )
2150}
2151
2152fn collect_strings(value: &serde_json::Value, out: &mut String) {
2153 match value {
2154 serde_json::Value::String(s) => {
2155 if !out.is_empty() {
2156 out.push(' ');
2157 }
2158 out.push_str(s);
2159 },
2160 serde_json::Value::Array(arr) => {
2161 for val in arr {
2162 collect_strings(val, out);
2163 }
2164 },
2165 serde_json::Value::Object(obj) => {
2166 for (key, val) in obj {
2167 if !out.is_empty() {
2168 out.push(' ');
2169 }
2170 out.push_str(key);
2171 collect_strings(val, out);
2172 }
2173 },
2174 serde_json::Value::Number(num) => {
2175 if !out.is_empty() {
2176 out.push(' ');
2177 }
2178 out.push_str(&num.to_string());
2179 },
2180 serde_json::Value::Bool(b) => {
2181 if !out.is_empty() {
2182 out.push(' ');
2183 }
2184 out.push_str(&b.to_string());
2185 },
2186 serde_json::Value::Null => {},
2187 }
2188}
2189
2190fn matches_filter(body_str: &str, filter: &serde_json::Value) -> bool {
2191 let Ok(body) = serde_json::from_str::<serde_json::Value>(body_str) else {
2192 return false;
2193 };
2194
2195 let Some(filter_obj) = filter.as_object() else {
2196 return true;
2197 };
2198
2199 for (k, v) in filter_obj {
2200 if body.get(k) != Some(v) {
2201 return false;
2202 }
2203 }
2204 true
2205}
2206
2207#[cfg(test)]
2208mod tests {
2209 use rusqlite::Connection;
2210 use tempfile::NamedTempFile;
2211
2212 use super::*;
2213 use crate::store::Searcher;
2214 use crate::{ListObjectsOptions, SearchOptions};
2215
2216 #[test]
2217 fn records_schema_version_on_initialize() {
2218 let store = SqliteThingStore::open_in_memory().unwrap();
2219
2220 assert_eq!(store.schema_version().unwrap(), SQLITE_SCHEMA_VERSION);
2221 }
2222
2223 #[test]
2224 fn integrity_check_passes_on_fresh_store() {
2225 let store = SqliteThingStore::open_in_memory().unwrap();
2226 let ok: String = store
2227 .connection
2228 .query_row("PRAGMA quick_check", [], |row| row.get(0))
2229 .unwrap();
2230 assert_eq!(ok, "ok");
2231 }
2232
2233 #[test]
2234 fn wal_checkpoint_returns_zero_frames_on_in_memory() {
2235 let store = SqliteThingStore::open_in_memory().unwrap();
2236 let (_busy, _frames) = store.wal_checkpoint().unwrap();
2238 }
2239
2240 #[test]
2241 fn backup_to_creates_valid_database() {
2242 let mut store = SqliteThingStore::open_in_memory().unwrap();
2243 store
2244 .put_object(MemoryObject::new("test", "1", r#"{"v":1}"#))
2245 .unwrap();
2246
2247 let dir = tempfile::tempdir().unwrap();
2248 let backup_path = dir.path().join("backup.db");
2249 let path_str = backup_path.to_str().unwrap().to_string();
2250
2251 store.backup_to(&path_str).unwrap();
2252 assert!(backup_path.exists());
2253
2254 let backup = SqliteThingStore::open(&backup_path).unwrap();
2255 let obj = backup.get_object("test", "1").unwrap();
2256 assert!(obj.is_some());
2257 assert_eq!(obj.unwrap().key.id, "1");
2258 }
2259
2260 #[test]
2261 fn rejects_newer_schema_versions() {
2262 let file = NamedTempFile::new().unwrap();
2263 let connection = Connection::open(file.path()).unwrap();
2264 connection
2265 .execute_batch(
2266 r"
2267 CREATE TABLE thingd_schema_migrations (
2268 version INTEGER PRIMARY KEY,
2269 name TEXT NOT NULL,
2270 applied_at TEXT NOT NULL
2271 );
2272
2273 INSERT INTO thingd_schema_migrations (version, name, applied_at)
2274 VALUES (999, 'future', strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
2275 ",
2276 )
2277 .unwrap();
2278
2279 let Err(error) = SqliteThingStore::open(file.path()) else {
2280 panic!("expected newer schema version to be rejected");
2281 };
2282
2283 assert!(error.to_string().contains("newer than supported version"));
2284 }
2285
2286 #[test]
2287 fn stores_objects_across_reopen() {
2288 let file = NamedTempFile::new().unwrap();
2289
2290 {
2291 let mut store = SqliteThingStore::open(file.path()).unwrap();
2292 let object = store
2293 .put_object(MemoryObject::new(
2294 "decisions",
2295 "sqlite-backend",
2296 "{\"text\":\"Use SQLite\"}",
2297 ))
2298 .unwrap();
2299
2300 assert_eq!(object.version, 1);
2301 }
2302
2303 let store = SqliteThingStore::open(file.path()).unwrap();
2304 let object = store
2305 .get_object("decisions", "sqlite-backend")
2306 .unwrap()
2307 .unwrap();
2308
2309 assert_eq!(object.body, "{\"text\":\"Use SQLite\"}");
2310 assert_eq!(object.version, 1);
2311 }
2312
2313 #[test]
2314 fn increments_object_versions() {
2315 let mut store = SqliteThingStore::open_in_memory().unwrap();
2316
2317 let first = store
2318 .put_object(MemoryObject::new("decisions", "versioned", "{}"))
2319 .unwrap();
2320 let second = store
2321 .put_object(MemoryObject::new("decisions", "versioned", "{\"v\":2}"))
2322 .unwrap();
2323
2324 assert_eq!(first.version, 1);
2325 assert_eq!(second.version, 2);
2326 }
2327
2328 #[test]
2329 fn lists_objects_with_optional_collection_filter() {
2330 let mut store = SqliteThingStore::open_in_memory().unwrap();
2331
2332 store
2333 .put_object(MemoryObject::new("decisions", "sqlite-backend", "{}"))
2334 .unwrap();
2335 store
2336 .put_object(MemoryObject::new("notes", "agent-guide", "{}"))
2337 .unwrap();
2338
2339 let filtered = store
2340 .list_objects(
2341 Some(&["decisions".to_string()]),
2342 &ListObjectsOptions::default(),
2343 )
2344 .unwrap();
2345
2346 assert_eq!(
2347 store
2348 .list_objects(None, &ListObjectsOptions::default())
2349 .unwrap()
2350 .len(),
2351 2
2352 );
2353 assert_eq!(filtered.len(), 1);
2354 assert_eq!(filtered[0].key.collection, "decisions");
2355 }
2356
2357 #[test]
2358 fn stores_events_across_reopen() {
2359 let file = NamedTempFile::new().unwrap();
2360
2361 {
2362 let mut store = SqliteThingStore::open(file.path()).unwrap();
2363 let event = store
2364 .append_event(MemoryEvent::new(
2365 "project:thingd",
2366 "decision.made",
2367 "Use SQLite first",
2368 ))
2369 .unwrap();
2370
2371 assert_eq!(event.sequence, 1);
2372 }
2373
2374 let store = SqliteThingStore::open(file.path()).unwrap();
2375 let events = store
2376 .list_events(Some("project:thingd"), ListEventsOptions::default())
2377 .unwrap();
2378
2379 assert_eq!(events.len(), 1);
2380 assert_eq!(events[0].event_type, "decision.made");
2381 assert_eq!(events[0].sequence, 1);
2382 }
2383
2384 #[test]
2385 fn stores_queue_jobs_across_reopen() {
2386 let file = NamedTempFile::new().unwrap();
2387
2388 {
2389 let mut store = SqliteThingStore::open(file.path()).unwrap();
2390 let job = store
2391 .push_job(QueueJob::new("embed", "job-1", "{\"doc\":\"doc-1\"}", 3))
2392 .unwrap();
2393
2394 assert_eq!(job.status, QueueJobStatus::Ready);
2395 }
2396
2397 let store = SqliteThingStore::open(file.path()).unwrap();
2398 let jobs = store.list_jobs("embed").unwrap();
2399
2400 assert_eq!(jobs.len(), 1);
2401 assert_eq!(jobs[0].id, "job-1");
2402 assert_eq!(jobs[0].status, QueueJobStatus::Ready);
2403 }
2404
2405 #[test]
2406 fn returns_existing_queue_job_for_duplicate_push() {
2407 let mut store = SqliteThingStore::open_in_memory().unwrap();
2408
2409 let first = store
2410 .push_job(QueueJob::new("embed", "job-1", "{\"doc\":\"doc-1\"}", 3))
2411 .unwrap();
2412 let second = store
2413 .push_job(QueueJob::new("embed", "job-1", "{\"doc\":\"doc-2\"}", 3))
2414 .unwrap();
2415
2416 assert_eq!(first.body, "{\"doc\":\"doc-1\"}");
2417 assert_eq!(second.body, first.body);
2418 assert_eq!(store.list_jobs("embed").unwrap().len(), 1);
2419 }
2420
2421 #[test]
2422 fn claims_and_acks_queue_jobs() {
2423 let mut store = SqliteThingStore::open_in_memory().unwrap();
2424
2425 store
2426 .push_job(QueueJob::new("embed", "job-1", "{\"doc\":\"doc-1\"}", 3))
2427 .unwrap();
2428
2429 let claimed = store.claim_job("embed").unwrap().unwrap();
2430 let acked = store.ack_job("embed", "job-1").unwrap().unwrap();
2431
2432 assert_eq!(claimed.status, QueueJobStatus::Leased);
2433 assert_eq!(claimed.attempts, 1);
2434 assert!(claimed.leased_at_ms.is_some());
2435 assert!(claimed.lease_expires_at_ms.is_some());
2436 assert_eq!(acked.status, QueueJobStatus::Completed);
2437 assert!(acked.completed_at_ms.is_some());
2438 assert!(store.claim_job("embed").unwrap().is_none());
2439 }
2440
2441 #[test]
2442 fn nacks_queue_jobs_to_retry_then_dead_letter() {
2443 let mut store = SqliteThingStore::open_in_memory().unwrap();
2444
2445 store
2446 .push_job(QueueJob::new("embed", "job-1", "{\"doc\":\"doc-1\"}", 2))
2447 .unwrap();
2448
2449 store.claim_job("embed").unwrap().unwrap();
2450 let retried = store.nack_job("embed", "job-1").unwrap().unwrap();
2451 assert_eq!(retried.status, QueueJobStatus::Ready);
2452 assert_eq!(retried.attempts, 1);
2453
2454 store.claim_job("embed").unwrap().unwrap();
2455 let dead = store.nack_job("embed", "job-1").unwrap().unwrap();
2456 assert_eq!(dead.status, QueueJobStatus::Dead);
2457 assert_eq!(dead.attempts, 2);
2458 assert!(dead.dead_at_ms.is_some());
2459 assert_eq!(store.list_dead_jobs("embed").unwrap().len(), 1);
2460 }
2461
2462 #[test]
2463 fn does_not_claim_delayed_queue_jobs_before_available() {
2464 let mut store = SqliteThingStore::open_in_memory().unwrap();
2465
2466 store
2467 .push_job(QueueJob::new("embed", "job-1", "{\"doc\":\"doc-1\"}", 3).delay_by_ms(60_000))
2468 .unwrap();
2469
2470 assert!(store.claim_job("embed").unwrap().is_none());
2471 }
2472
2473 #[test]
2474 fn reclaims_queue_jobs_after_lease_expiration() {
2475 let mut store = SqliteThingStore::open_in_memory().unwrap();
2476
2477 store
2478 .push_job(QueueJob::new("embed", "job-1", "{\"doc\":\"doc-1\"}", 3))
2479 .unwrap();
2480
2481 let first = store
2482 .claim_job_with_options("embed", QueueClaimOptions::new(0))
2483 .unwrap()
2484 .unwrap();
2485 let second = store.claim_job("embed").unwrap().unwrap();
2486
2487 assert_eq!(first.status, QueueJobStatus::Leased);
2488 assert_eq!(second.status, QueueJobStatus::Leased);
2489 assert_eq!(second.attempts, 2);
2490 }
2491
2492 #[test]
2493 fn nacks_queue_jobs_with_retry_delay() {
2494 let mut store = SqliteThingStore::open_in_memory().unwrap();
2495
2496 store
2497 .push_job(QueueJob::new("embed", "job-1", "{\"doc\":\"doc-1\"}", 3))
2498 .unwrap();
2499
2500 store.claim_job("embed").unwrap().unwrap();
2501 let retried = store
2502 .nack_job_with_options("embed", "job-1", QueueNackOptions::new(60_000))
2503 .unwrap()
2504 .unwrap();
2505
2506 assert_eq!(retried.status, QueueJobStatus::Ready);
2507 assert!(store.claim_job("embed").unwrap().is_none());
2508 }
2509
2510 #[test]
2511 fn persists_completed_queue_jobs_across_reopen() {
2512 let file = NamedTempFile::new().unwrap();
2513
2514 {
2515 let mut store = SqliteThingStore::open(file.path()).unwrap();
2516 store
2517 .push_job(QueueJob::new("embed", "job-1", "{\"doc\":\"doc-1\"}", 3))
2518 .unwrap();
2519 store.claim_job("embed").unwrap().unwrap();
2520 store.ack_job("embed", "job-1").unwrap().unwrap();
2521 }
2522
2523 let store = SqliteThingStore::open(file.path()).unwrap();
2524 let jobs = store.list_jobs("embed").unwrap();
2525
2526 assert_eq!(jobs.len(), 1);
2527 assert_eq!(jobs[0].status, QueueJobStatus::Completed);
2528 assert_eq!(jobs[0].attempts, 1);
2529 }
2530
2531 #[test]
2532 fn test_fts5_search_indexing_and_stemming() {
2533 let mut store = SqliteThingStore::open_in_memory().unwrap();
2534
2535 store
2537 .put_object(MemoryObject::new(
2538 "decisions",
2539 "choice-1",
2540 "{\"text\":\"I choose this implementation plan because it has great benefits.\", \"status\":\"active\", \"priority\":1}",
2541 ))
2542 .unwrap();
2543
2544 store
2545 .put_object(MemoryObject::new(
2546 "decisions",
2547 "choice-2",
2548 "{\"text\":\"He chooses that plan.\", \"status\":\"draft\", \"priority\":2}",
2549 ))
2550 .unwrap();
2551
2552 let results = store
2554 .search("implementation", crate::SearchOptions::default())
2555 .unwrap();
2556 assert_eq!(results.len(), 1);
2557 assert_eq!(results[0].id, "choice-1");
2558
2559 let results_stem = store
2561 .search("choosing", crate::SearchOptions::default())
2562 .unwrap();
2563 assert_eq!(results_stem.len(), 2);
2564
2565 let options_col = crate::SearchOptions {
2567 collections: Some(vec!["unrelated_col".to_string()]),
2568 ..Default::default()
2569 };
2570 let results_col = store.search("choose", options_col).unwrap();
2571 assert_eq!(results_col.len(), 0);
2572
2573 let options_filter = crate::SearchOptions {
2575 filter: Some(serde_json::json!({"status": "active"})),
2576 ..Default::default()
2577 };
2578 let results_filter = store.search("choose", options_filter).unwrap();
2579 assert_eq!(results_filter.len(), 1);
2580 assert_eq!(results_filter[0].id, "choice-1");
2581
2582 store.delete_object("decisions", "choice-1").unwrap();
2584 let results_after_del = store
2585 .search("choose", crate::SearchOptions::default())
2586 .unwrap();
2587 assert_eq!(results_after_del.len(), 1);
2588 assert_eq!(results_after_del[0].id, "choice-2");
2589 }
2590
2591 #[test]
2592 fn search_consistent_after_batch_delete() {
2593 let mut store = SqliteThingStore::open_in_memory().unwrap();
2594
2595 store
2596 .put_object(MemoryObject::new(
2597 "col",
2598 "a",
2599 r#"{"label":"target","name":"alpha"}"#,
2600 ))
2601 .unwrap();
2602 store
2603 .put_object(MemoryObject::new(
2604 "col",
2605 "b",
2606 r#"{"label":"target","name":"bravo"}"#,
2607 ))
2608 .unwrap();
2609 store
2610 .put_object(MemoryObject::new(
2611 "col",
2612 "c",
2613 r#"{"label":"target","name":"charlie"}"#,
2614 ))
2615 .unwrap();
2616
2617 let results = store
2618 .search("target", crate::SearchOptions::default())
2619 .unwrap();
2620 assert_eq!(results.len(), 3);
2621
2622 store
2623 .delete_objects_batch(&[("col".into(), "a".into()), ("col".into(), "b".into())])
2624 .unwrap();
2625
2626 let results = store
2627 .search("target", crate::SearchOptions::default())
2628 .unwrap();
2629 assert_eq!(results.len(), 1);
2630 assert_eq!(results[0].id, "c");
2631 }
2632
2633 #[test]
2634 fn search_does_not_return_orphaned_fts_entries() {
2635 let store = SqliteThingStore::open_in_memory().unwrap();
2636
2637 store
2638 .connection
2639 .execute(
2640 "INSERT INTO search_index (collection, id, kind, text) VALUES (?1, ?2, 'object', ?3)",
2641 rusqlite::params!["orphan", "ghost", "this object does not exist"],
2642 )
2643 .unwrap();
2644
2645 let results = store
2646 .search("object", crate::SearchOptions::default())
2647 .unwrap();
2648 assert_eq!(results.len(), 0);
2649 }
2650
2651 #[test]
2652 fn get_and_search_consistent_after_delete() {
2653 let mut store = SqliteThingStore::open_in_memory().unwrap();
2654
2655 store
2656 .put_object(MemoryObject::new("col", "id-1", r#"{"name":"alice"}"#))
2657 .unwrap();
2658
2659 assert!(store.get_object("col", "id-1").unwrap().is_some());
2660
2661 let results = store
2662 .search("alice", crate::SearchOptions::default())
2663 .unwrap();
2664 assert_eq!(results.len(), 1);
2665
2666 store.delete_object("col", "id-1").unwrap();
2667
2668 assert!(store.get_object("col", "id-1").unwrap().is_none());
2669
2670 let results = store
2671 .search("alice", crate::SearchOptions::default())
2672 .unwrap();
2673 assert_eq!(results.len(), 0);
2674 }
2675
2676 #[test]
2677 fn counts_objects_correctly_after_deletions() {
2678 let mut store = SqliteThingStore::open_in_memory().unwrap();
2679
2680 assert_eq!(store.count_objects().unwrap(), 0);
2681
2682 store
2683 .put_object(MemoryObject::new("col1", "a", "{}"))
2684 .unwrap();
2685 store
2686 .put_object(MemoryObject::new("col1", "b", "{}"))
2687 .unwrap();
2688 store
2689 .put_object(MemoryObject::new("col2", "c", "{}"))
2690 .unwrap();
2691 assert_eq!(store.count_objects().unwrap(), 3);
2692
2693 store.delete_object("col1", "a").unwrap();
2694 assert_eq!(store.count_objects().unwrap(), 2);
2695
2696 store.delete_object("col1", "b").unwrap();
2697 assert_eq!(store.count_objects().unwrap(), 1);
2698
2699 store.delete_object("col2", "c").unwrap();
2700 assert_eq!(store.count_objects().unwrap(), 0);
2701 }
2702
2703 #[test]
2704 fn counts_events_correctly() {
2705 let mut store = SqliteThingStore::open_in_memory().unwrap();
2706
2707 assert_eq!(store.count_events().unwrap(), 0);
2708
2709 store
2710 .append_event(MemoryEvent::new("test", "a", ""))
2711 .unwrap();
2712 store
2713 .append_event(MemoryEvent::new("test", "b", ""))
2714 .unwrap();
2715 assert_eq!(store.count_events().unwrap(), 2);
2716 }
2717
2718 #[test]
2719 fn deletes_last_event_from_stream() {
2720 let mut store = SqliteThingStore::open_in_memory().unwrap();
2721
2722 store
2723 .append_event(MemoryEvent::new("match:1", "turn.recorded", "{}"))
2724 .unwrap();
2725 store
2726 .append_event(MemoryEvent::new("match:1", "turn.recorded", "{}"))
2727 .unwrap();
2728 store
2729 .append_event(MemoryEvent::new("match:2", "turn.recorded", "{}"))
2730 .unwrap();
2731
2732 let deleted = store.delete_last_event("match:1").unwrap().unwrap();
2733 assert_eq!(deleted.sequence, 2);
2734
2735 let remaining = store
2736 .list_events(Some("match:1"), ListEventsOptions::default())
2737 .unwrap();
2738 assert_eq!(remaining.len(), 1);
2739 assert_eq!(remaining[0].sequence, 1);
2740
2741 let match2 = store
2743 .list_events(Some("match:2"), ListEventsOptions::default())
2744 .unwrap();
2745 assert_eq!(match2.len(), 1);
2746 }
2747
2748 #[test]
2749 fn returns_none_when_delete_last_event_on_empty_stream() {
2750 let mut store = SqliteThingStore::open_in_memory().unwrap();
2751 assert!(store.delete_last_event("nonexistent").unwrap().is_none());
2752 }
2753
2754 #[test]
2755 fn deletes_stream_and_returns_count() {
2756 let mut store = SqliteThingStore::open_in_memory().unwrap();
2757
2758 store
2759 .append_event(MemoryEvent::new("match:1", "turn.recorded", "{}"))
2760 .unwrap();
2761 store
2762 .append_event(MemoryEvent::new("match:1", "turn.recorded", "{}"))
2763 .unwrap();
2764 store
2765 .append_event(MemoryEvent::new("match:2", "turn.recorded", "{}"))
2766 .unwrap();
2767
2768 let count = store.delete_stream("match:1").unwrap();
2769 assert_eq!(count, 2);
2770
2771 let remaining = store
2772 .list_events(Some("match:1"), ListEventsOptions::default())
2773 .unwrap();
2774 assert_eq!(remaining.len(), 0);
2775
2776 let match2 = store
2778 .list_events(Some("match:2"), ListEventsOptions::default())
2779 .unwrap();
2780 assert_eq!(match2.len(), 1);
2781 }
2782
2783 #[test]
2784 fn returns_zero_for_delete_stream_on_empty_stream() {
2785 let mut store = SqliteThingStore::open_in_memory().unwrap();
2786 assert_eq!(store.delete_stream("nonexistent").unwrap(), 0);
2787 }
2788
2789 #[test]
2790 fn counts_jobs_correctly() {
2791 let mut store = SqliteThingStore::open_in_memory().unwrap();
2792
2793 assert_eq!(store.count_active_jobs().unwrap(), 0);
2794 assert_eq!(store.count_dead_jobs().unwrap(), 0);
2795
2796 store
2797 .push_job(QueueJob::new("work", "j1", "p1", 3))
2798 .unwrap();
2799 store
2800 .push_job(QueueJob::new("work", "j2", "p2", 3))
2801 .unwrap();
2802 store
2803 .push_job(QueueJob::new("other", "j3", "p3", 1))
2804 .unwrap();
2805 assert_eq!(store.count_active_jobs().unwrap(), 3);
2806
2807 store.claim_job("other").unwrap();
2808 store.nack_job("other", "j3").unwrap();
2809 assert_eq!(store.count_dead_jobs().unwrap(), 1);
2810 assert_eq!(store.count_active_jobs().unwrap(), 2);
2811 }
2812
2813 #[test]
2814 fn lists_collections_streams_and_queues() {
2815 let mut store = SqliteThingStore::open_in_memory().unwrap();
2816
2817 assert!(store.list_collections().unwrap().is_empty());
2818 assert!(store.list_streams().unwrap().is_empty());
2819 assert!(store.list_queues().unwrap().is_empty());
2820
2821 store
2822 .put_object(MemoryObject::new("col-a", "x", "{}"))
2823 .unwrap();
2824 store
2825 .put_object(MemoryObject::new("col-b", "y", "{}"))
2826 .unwrap();
2827 store
2828 .put_object(MemoryObject::new("col-a", "z", "{}"))
2829 .unwrap();
2830 let collections = store.list_collections().unwrap();
2831 assert_eq!(collections, vec!["col-a", "col-b"]);
2832
2833 store
2834 .append_event(MemoryEvent::new("s1", "t", "e1"))
2835 .unwrap();
2836 store
2837 .append_event(MemoryEvent::new("s2", "t", "e2"))
2838 .unwrap();
2839 let streams = store.list_streams().unwrap();
2840 assert_eq!(streams, vec!["s1", "s2"]);
2841
2842 store
2843 .push_job(QueueJob::new("work", "j1", "p1", 3))
2844 .unwrap();
2845 store
2846 .push_job(QueueJob::new("jobs", "j2", "p2", 3))
2847 .unwrap();
2848 let queues = store.list_queues().unwrap();
2849 assert_eq!(queues, vec!["jobs", "work"]);
2850 }
2851
2852 #[test]
2853 fn search_respects_filter_and_limit() {
2854 let mut store = SqliteThingStore::open_in_memory().unwrap();
2855
2856 store
2857 .put_object(MemoryObject::new(
2858 "docs",
2859 "a",
2860 r#"{"text":"hello world","tag":"greeting"}"#,
2861 ))
2862 .unwrap();
2863 store
2864 .put_object(MemoryObject::new(
2865 "docs",
2866 "b",
2867 r#"{"text":"hello there","tag":"greeting"}"#,
2868 ))
2869 .unwrap();
2870 store
2871 .put_object(MemoryObject::new(
2872 "docs",
2873 "c",
2874 r#"{"text":"goodbye world","tag":"farewell"}"#,
2875 ))
2876 .unwrap();
2877
2878 let all = store.search("world", SearchOptions::default()).unwrap();
2879 assert_eq!(all.len(), 2);
2880
2881 let limited = store
2882 .search(
2883 "world",
2884 SearchOptions {
2885 limit: Some(1),
2886 ..Default::default()
2887 },
2888 )
2889 .unwrap();
2890 assert_eq!(limited.len(), 1);
2891
2892 let filtered = store
2893 .search(
2894 "hello",
2895 SearchOptions {
2896 collections: Some(vec!["docs".into()]),
2897 ..Default::default()
2898 },
2899 )
2900 .unwrap();
2901 assert_eq!(filtered.len(), 2);
2902 }
2903
2904 #[test]
2907 fn list_objects_filter_returns_matching_objects() {
2908 let mut store = SqliteThingStore::open_in_memory().unwrap();
2909
2910 store
2911 .put_object(MemoryObject::new("w", "a", r#"{"color":"red","size":1}"#))
2912 .unwrap();
2913 store
2914 .put_object(MemoryObject::new("w", "b", r#"{"color":"blue","size":2}"#))
2915 .unwrap();
2916 store
2917 .put_object(MemoryObject::new("w", "c", r#"{"color":"red","size":3}"#))
2918 .unwrap();
2919
2920 let opts = ListObjectsOptions {
2921 filter: vec![("color".into(), serde_json::json!("red"))],
2922 ..Default::default()
2923 };
2924 let results = store.list_objects(Some(&["w".to_string()]), &opts).unwrap();
2925 assert_eq!(results.len(), 2);
2926 assert!(results.iter().all(|o| o.body.contains("\"red\"")));
2927 }
2928
2929 #[test]
2930 fn list_objects_filter_no_match_returns_empty() {
2931 let mut store = SqliteThingStore::open_in_memory().unwrap();
2932
2933 store
2934 .put_object(MemoryObject::new("w", "a", r#"{"color":"red"}"#))
2935 .unwrap();
2936
2937 let opts = ListObjectsOptions {
2938 filter: vec![("color".into(), serde_json::json!("green"))],
2939 ..Default::default()
2940 };
2941 let results = store.list_objects(Some(&["w".to_string()]), &opts).unwrap();
2942 assert!(results.is_empty());
2943 }
2944
2945 #[test]
2946 fn list_objects_limit_truncates_results() {
2947 let mut store = SqliteThingStore::open_in_memory().unwrap();
2948
2949 for i in 0..5u32 {
2950 store
2951 .put_object(MemoryObject::new("col", format!("id-{i}"), "{}"))
2952 .unwrap();
2953 }
2954
2955 let opts = ListObjectsOptions {
2956 limit: Some(3),
2957 ..Default::default()
2958 };
2959 let results = store
2960 .list_objects(Some(&["col".to_string()]), &opts)
2961 .unwrap();
2962 assert_eq!(results.len(), 3);
2963 }
2964
2965 #[test]
2966 fn list_objects_offset_skips_results() {
2967 let mut store = SqliteThingStore::open_in_memory().unwrap();
2968
2969 for i in 0..5u32 {
2970 store
2971 .put_object(MemoryObject::new("col", format!("id-{i}"), "{}"))
2972 .unwrap();
2973 }
2974
2975 let opts = ListObjectsOptions {
2976 offset: Some(3),
2977 ..Default::default()
2978 };
2979 let results = store
2980 .list_objects(Some(&["col".to_string()]), &opts)
2981 .unwrap();
2982 assert_eq!(results.len(), 2);
2983 }
2984
2985 #[test]
2986 fn list_objects_filter_and_limit_combined() {
2987 let mut store = SqliteThingStore::open_in_memory().unwrap();
2988
2989 for i in 0..4u32 {
2990 store
2991 .put_object(MemoryObject::new(
2992 "col",
2993 format!("id-{i}"),
2994 r#"{"status":"active"}"#,
2995 ))
2996 .unwrap();
2997 }
2998 store
2999 .put_object(MemoryObject::new("col", "id-4", r#"{"status":"inactive"}"#))
3000 .unwrap();
3001
3002 let opts = ListObjectsOptions {
3003 filter: vec![("status".into(), serde_json::json!("active"))],
3004 limit: Some(2),
3005 ..Default::default()
3006 };
3007 let results = store
3008 .list_objects(Some(&["col".to_string()]), &opts)
3009 .unwrap();
3010 assert_eq!(results.len(), 2);
3011 assert!(results.iter().all(|o| o.body.contains("active")));
3012 }
3013
3014 #[test]
3015 fn list_objects_numeric_filter() {
3016 let mut store = SqliteThingStore::open_in_memory().unwrap();
3017
3018 store
3019 .put_object(MemoryObject::new("items", "a", r#"{"score":10,"tag":"x"}"#))
3020 .unwrap();
3021 store
3022 .put_object(MemoryObject::new("items", "b", r#"{"score":20,"tag":"x"}"#))
3023 .unwrap();
3024 store
3025 .put_object(MemoryObject::new("items", "c", r#"{"score":10,"tag":"y"}"#))
3026 .unwrap();
3027
3028 let opts = ListObjectsOptions {
3029 filter: vec![("score".into(), serde_json::json!(10))],
3030 ..Default::default()
3031 };
3032 let results = store
3033 .list_objects(Some(&["items".to_string()]), &opts)
3034 .unwrap();
3035 assert_eq!(results.len(), 2);
3036 }
3037
3038 #[test]
3041 fn append_event_returning_sets_sequence_and_timestamp() {
3042 let mut store = SqliteThingStore::open_in_memory().unwrap();
3043
3044 let first = store
3045 .append_event(MemoryEvent::new("s", "ev.first", r#"{"x":1}"#))
3046 .unwrap();
3047 let second = store
3048 .append_event(MemoryEvent::new("s", "ev.second", r#"{"x":2}"#))
3049 .unwrap();
3050
3051 assert_eq!(first.sequence, 1);
3052 assert_eq!(second.sequence, 2);
3053 assert!(
3054 !first.created_at.is_empty(),
3055 "created_at must be set by RETURNING"
3056 );
3057 assert!(
3058 !second.created_at.is_empty(),
3059 "created_at must be set by RETURNING"
3060 );
3061 }
3062
3063 #[test]
3064 fn append_event_sequence_monotonically_increases_across_streams() {
3065 let mut store = SqliteThingStore::open_in_memory().unwrap();
3066
3067 let a = store
3068 .append_event(MemoryEvent::new("stream-a", "t", "{}"))
3069 .unwrap();
3070 let b = store
3071 .append_event(MemoryEvent::new("stream-b", "t", "{}"))
3072 .unwrap();
3073 let c = store
3074 .append_event(MemoryEvent::new("stream-a", "t", "{}"))
3075 .unwrap();
3076
3077 assert_eq!(a.sequence, 1);
3078 assert_eq!(b.sequence, 2);
3079 assert_eq!(c.sequence, 3);
3080 }
3081
3082 #[test]
3085 fn cas_succeeds_on_matching_version() {
3086 let mut store = SqliteThingStore::open_in_memory().unwrap();
3087
3088 let stored = store
3089 .put_object(MemoryObject::new("col", "id", r#"{"v":1}"#))
3090 .unwrap();
3091 assert_eq!(stored.version, 1);
3092
3093 let opts = crate::PutObjectOptions {
3094 expected_version: Some(1),
3095 ..Default::default()
3096 };
3097 let updated = store
3098 .put_object_with_options(MemoryObject::new("col", "id", r#"{"v":2}"#), opts)
3099 .unwrap();
3100 assert_eq!(updated.version, 2);
3101 }
3102
3103 #[test]
3104 fn cas_fails_on_version_mismatch() {
3105 let mut store = SqliteThingStore::open_in_memory().unwrap();
3106
3107 store
3108 .put_object(MemoryObject::new("col", "id", r#"{"v":1}"#))
3109 .unwrap();
3110
3111 let opts = crate::PutObjectOptions {
3112 expected_version: Some(42),
3113 ..Default::default()
3114 };
3115 let err = store
3116 .put_object_with_options(MemoryObject::new("col", "id", r#"{"v":2}"#), opts)
3117 .unwrap_err();
3118 assert!(matches!(err, crate::ThingdError::Conflict(_)));
3119 }
3120
3121 #[test]
3122 fn cas_fails_on_nonexistent_object() {
3123 let mut store = SqliteThingStore::open_in_memory().unwrap();
3124
3125 let opts = crate::PutObjectOptions {
3126 expected_version: Some(1),
3127 ..Default::default()
3128 };
3129 let err = store
3130 .put_object_with_options(MemoryObject::new("col", "id", r#"{"v":1}"#), opts)
3131 .unwrap_err();
3132 assert!(matches!(err, crate::ThingdError::Conflict(_)));
3133 }
3134
3135 #[test]
3136 fn cas_none_skips_check() {
3137 let mut store = SqliteThingStore::open_in_memory().unwrap();
3138
3139 let stored = store
3140 .put_object_with_options(
3141 MemoryObject::new("col", "id", r#"{"v":1}"#),
3142 crate::PutObjectOptions::default(),
3143 )
3144 .unwrap();
3145 assert_eq!(stored.version, 1);
3146 }
3147
3148 #[test]
3151 fn put_objects_batch_atomicity() {
3152 let mut store = SqliteThingStore::open_in_memory().unwrap();
3153
3154 store
3156 .put_object(MemoryObject::new("col", "a", r#"{"name":"old-a"}"#))
3157 .unwrap();
3158 store
3159 .put_object(MemoryObject::new("col", "b", r#"{"name":"old-b"}"#))
3160 .unwrap();
3161
3162 let results = store
3164 .put_objects_batch(vec![
3165 MemoryObject::new("col", "a", r#"{"name":"new-a"}"#),
3166 MemoryObject::new("col", "c", r#"{"name":"new-c"}"#),
3167 ])
3168 .unwrap();
3169 assert_eq!(results.len(), 2);
3170 assert_eq!(results[0].version, 2); assert_eq!(results[1].version, 1); assert_eq!(store.count_objects().unwrap(), 3);
3175 }
3176
3177 #[test]
3178 fn delete_objects_batch_atomicity() {
3179 let mut store = SqliteThingStore::open_in_memory().unwrap();
3180
3181 for c in ["a", "b", "c"] {
3182 store
3183 .put_object(MemoryObject::new("col", c, r#"{"name":"x"}"#))
3184 .unwrap();
3185 }
3186
3187 assert_eq!(store.count_objects().unwrap(), 3);
3188
3189 let deleted = store
3190 .delete_objects_batch(&[("col".into(), "a".into()), ("col".into(), "b".into())])
3191 .unwrap();
3192 assert_eq!(deleted, 2);
3193 assert_eq!(store.count_objects().unwrap(), 1);
3194 assert!(store.get_object("col", "a").unwrap().is_none());
3195 assert!(store.get_object("col", "c").unwrap().is_some());
3196 }
3197
3198 #[test]
3199 fn batch_operations_empty_input() {
3200 let mut store = SqliteThingStore::open_in_memory().unwrap();
3201
3202 let results = store.put_objects_batch(vec![]).unwrap();
3203 assert!(results.is_empty());
3204
3205 let deleted = store.delete_objects_batch(&[]).unwrap();
3206 assert_eq!(deleted, 0);
3207 }
3208
3209 #[test]
3212 fn event_idempotency_returns_existing_event() {
3213 let mut store = SqliteThingStore::open_in_memory().unwrap();
3214
3215 let mut event = MemoryEvent::new("stream", "test", r#"{"key":"val"}"#);
3216 event.idempotency_key = "idem-1".to_string();
3217
3218 let first = store.append_event(event.clone()).unwrap();
3219 assert_eq!(first.sequence, 1);
3220
3221 let second = store.append_event(event).unwrap();
3222 assert_eq!(second.sequence, first.sequence);
3223 assert_eq!(second.body, first.body);
3224 }
3225
3226 #[test]
3227 fn event_idempotency_different_keys_are_distinct() {
3228 let mut store = SqliteThingStore::open_in_memory().unwrap();
3229
3230 let mut event_a = MemoryEvent::new("stream", "test", r#"{"key":"a"}"#);
3231 event_a.idempotency_key = "idem-a".to_string();
3232
3233 let mut event_b = MemoryEvent::new("stream", "test", r#"{"key":"b"}"#);
3234 event_b.idempotency_key = "idem-b".to_string();
3235
3236 let first = store.append_event(event_a).unwrap();
3237 let second = store.append_event(event_b).unwrap();
3238 assert_eq!(first.sequence, 1);
3239 assert_eq!(second.sequence, 2);
3240 }
3241}