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