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