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