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