Skip to main content

thingd/
sqlite.rs

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