Skip to main content

autopulse_database/
conn.rs

1use crate::models::{NewScanEvent, ScanEvent};
2use anyhow::Context;
3use autopulse_utils::sify;
4use diesel::connection::SimpleConnection;
5use diesel::r2d2::{ConnectionManager, Pool, PooledConnection};
6use diesel::SelectableHelper;
7use diesel::{Connection, RunQueryDsl};
8use diesel_migrations::{embed_migrations, EmbeddedMigrations, MigrationHarness};
9use serde::Deserialize;
10use std::fs::OpenOptions;
11use std::path::PathBuf;
12use std::time::{SystemTime, UNIX_EPOCH};
13use tracing::{info, warn};
14
15#[doc(hidden)]
16#[cfg(feature = "postgres")]
17const POSTGRES_MIGRATIONS: EmbeddedMigrations = embed_migrations!("migrations/postgres");
18
19#[doc(hidden)]
20#[cfg(feature = "sqlite")]
21const SQLITE_MIGRATIONS: EmbeddedMigrations = embed_migrations!("migrations/sqlite");
22
23#[derive(Deserialize, Debug)]
24#[serde(rename_all = "lowercase")]
25#[derive(Default)]
26pub enum DatabaseType {
27    #[cfg(feature = "sqlite")]
28    #[cfg_attr(feature = "sqlite", default)]
29    Sqlite,
30    #[cfg(feature = "postgres")]
31    #[cfg_attr(not(feature = "sqlite"), default)]
32    Postgres,
33}
34
35impl DatabaseType {
36    pub fn default_url(&self) -> String {
37        match self {
38            #[cfg(feature = "sqlite")]
39            Self::Sqlite => "sqlite://data/autopulse.db".to_string(),
40            #[cfg(feature = "postgres")]
41            Self::Postgres => "postgres://autopulse:autopulse@localhost:5432/autopulse".to_string(),
42        }
43    }
44}
45
46/// Represents a connection to either a `PostgreSQL` or `SQLite` database.
47#[derive(diesel::MultiConnection)]
48pub enum AnyConnection {
49    /// A connection to a `PostgreSQL` database.
50    ///
51    /// This is used when the `database_url` is a `PostgreSQL` URL.
52    ///
53    /// # Example
54    ///
55    /// ```md
56    /// postgres://user:password@localhost:5432/database
57    /// ```
58    #[cfg(feature = "postgres")]
59    Postgresql(diesel::PgConnection),
60    // Mysql(diesel::MysqlConnection),
61    /// A connection to a `SQLite` database.
62    ///
63    /// This is used when the `database_url` is a `SQLite` URL.
64    ///
65    /// Note: The directory where the database is stored will also be populated with a WAL file and a journal file.
66    ///
67    /// # Example
68    ///
69    /// ```bash
70    /// # Relative path
71    /// sqlite://database.db
72    /// sqlite://data/database.db
73    ///
74    /// # Absolute path
75    /// sqlite:///data/database.db
76    ///
77    /// # In-memory database
78    /// sqlite://:memory: # In-memory database
79    /// ```
80    #[cfg(feature = "sqlite")]
81    Sqlite(diesel::SqliteConnection),
82}
83
84#[doc(hidden)]
85#[derive(Debug, Default)]
86pub struct AcquireHook {
87    pub setup: bool,
88}
89
90impl diesel::r2d2::CustomizeConnection<AnyConnection, diesel::r2d2::Error> for AcquireHook {
91    fn on_acquire(&self, conn: &mut AnyConnection) -> Result<(), diesel::r2d2::Error> {
92        (|| {
93            match conn {
94                #[cfg(feature = "sqlite")]
95                AnyConnection::Sqlite(ref mut conn) => {
96                    conn.batch_execute("PRAGMA busy_timeout = 5000")?;
97                    conn.batch_execute("PRAGMA synchronous = NORMAL;")?;
98                    conn.batch_execute("PRAGMA wal_autocheckpoint = 1000;")?;
99                    conn.batch_execute("PRAGMA foreign_keys = ON;")?;
100
101                    if self.setup {
102                        conn.batch_execute("PRAGMA journal_mode = WAL;")?;
103                    }
104                }
105                #[cfg(feature = "postgres")]
106                AnyConnection::Postgresql(_) => {}
107            }
108            Ok(())
109        })()
110        .map_err(diesel::r2d2::Error::QueryError)
111    }
112}
113
114impl AnyConnection {
115    pub fn pre_init(database_url: &str) -> anyhow::Result<()> {
116        if database_url.starts_with("sqlite://") && !database_url.contains(":memory:") {
117            let path = database_url
118                .strip_prefix("sqlite://")
119                .expect("already checked prefix");
120
121            let path = PathBuf::from(path);
122
123            let Some(parent) = path.parent().filter(|p| !p.as_os_str().is_empty()) else {
124                return Ok(());
125            };
126
127            // Create directory if it doesn't exist
128            if !parent.exists() {
129                std::fs::create_dir_all(parent).with_context(|| {
130                    format!("failed to create database directory: {}", parent.display())
131                })?;
132            }
133
134            let timestamp = SystemTime::now()
135                .duration_since(UNIX_EPOCH)
136                .map(|duration| duration.as_nanos())
137                .unwrap_or_default();
138            let probe = parent.join(format!(
139                ".autopulse-db-write-test-{}-{timestamp}",
140                std::process::id()
141            ));
142
143            let file = OpenOptions::new()
144                .write(true)
145                .create_new(true)
146                .open(&probe)
147                .with_context(|| {
148                    format!("database directory is not writable: {}", parent.display())
149                })?;
150            drop(file);
151
152            std::fs::remove_file(&probe).with_context(|| {
153                format!(
154                    "failed to remove database directory write test file: {}",
155                    probe.display()
156                )
157            })?;
158        }
159
160        Ok(())
161    }
162
163    pub fn migrate(&mut self) -> anyhow::Result<()> {
164        let migrations_applied = match self {
165            #[cfg(feature = "postgres")]
166            Self::Postgresql(conn) => conn.run_pending_migrations(POSTGRES_MIGRATIONS),
167            #[cfg(feature = "sqlite")]
168            Self::Sqlite(conn) => conn.run_pending_migrations(SQLITE_MIGRATIONS),
169        }
170        // Preserve `e.source()` chain via anyhow::Error::from_boxed; the
171        // previous `anyhow!("...{e}")` flattened it to Display only.
172        .map_err(|e| anyhow::Error::from_boxed(e).context("failed to run migrations"))?;
173
174        if !migrations_applied.is_empty() {
175            info!(
176                "Applied {} migration{}",
177                migrations_applied.len(),
178                sify(&migrations_applied)
179            );
180        }
181
182        Ok(())
183    }
184
185    pub fn close(&mut self) -> anyhow::Result<()> {
186        match self {
187            #[cfg(feature = "postgres")]
188            Self::Postgresql(_) => {}
189            #[cfg(feature = "sqlite")]
190            Self::Sqlite(conn) => {
191                // Should cleanup spare wal/shm files
192                conn.batch_execute("PRAGMA wal_checkpoint(TRUNCATE);")
193                    .context("failed to checkpoint WAL")?;
194            }
195        }
196
197        Ok(())
198    }
199
200    pub fn update_found(
201        &mut self,
202        previous: &ScanEvent,
203        status: &str,
204        at: chrono::NaiveDateTime,
205    ) -> anyhow::Result<Option<ScanEvent>> {
206        use crate::schema::scan_events::dsl::*;
207        use diesel::prelude::*;
208        use diesel::sql_types::Bool;
209
210        let query = diesel::update(
211            scan_events
212                .find(&previous.id)
213                .filter(process_status.eq("pending"))
214                .filter(found_status.eq(&previous.found_status))
215                .filter(
216                    file_hash.eq(&previous.file_hash).or(file_hash
217                        .is_null()
218                        .and(previous.file_hash.is_none().into_sql::<Bool>())),
219                ),
220        )
221        .set((
222            found_status.eq(status),
223            found_at.eq(Some(at)),
224            updated_at.eq(at),
225        ));
226        match self {
227            #[cfg(feature = "postgres")]
228            Self::Postgresql(conn) => query.returning(ScanEvent::as_returning()).get_result(conn),
229            #[cfg(feature = "sqlite")]
230            Self::Sqlite(conn) => query.returning(ScanEvent::as_returning()).get_result(conn),
231        }
232        .optional()
233        .map_err(Into::into)
234    }
235
236    pub fn update_process(
237        &mut self,
238        previous: &ScanEvent,
239        updated: &ScanEvent,
240    ) -> anyhow::Result<Option<ScanEvent>> {
241        use crate::schema::scan_events::dsl::*;
242        use diesel::prelude::*;
243        use diesel::sql_types::Bool;
244
245        // Dedupe and file checks own other columns and must survive target awaits.
246        let query = diesel::update(
247            scan_events
248                .find(&previous.id)
249                .filter(process_status.eq(&previous.process_status))
250                .filter(process_status.eq_any(["pending", "retry"]))
251                .filter(failed_times.eq(previous.failed_times))
252                .filter(targets_hit.eq(&previous.targets_hit))
253                // A manual acceleration can race with selection as a retry
254                // becomes due. It must not invalidate work already dispatched.
255                .filter(
256                    next_retry_at.le(previous.next_retry_at).or(next_retry_at
257                        .is_null()
258                        .and(previous.next_retry_at.is_none().into_sql::<Bool>())),
259                ),
260        )
261        .set((
262            process_status.eq(&updated.process_status),
263            failed_times.eq(updated.failed_times),
264            next_retry_at.eq(updated.next_retry_at),
265            targets_hit.eq(&updated.targets_hit),
266            processed_at.eq(updated.processed_at),
267            updated_at.eq(updated.updated_at),
268        ));
269        match self {
270            #[cfg(feature = "postgres")]
271            Self::Postgresql(conn) => query.returning(ScanEvent::as_returning()).get_result(conn),
272            #[cfg(feature = "sqlite")]
273            Self::Sqlite(conn) => query.returning(ScanEvent::as_returning()).get_result(conn),
274        }
275        .optional()
276        .map_err(Into::into)
277    }
278
279    pub fn insert_and_return(&mut self, ev: &NewScanEvent) -> anyhow::Result<ScanEvent> {
280        match self {
281            #[cfg(feature = "postgres")]
282            Self::Postgresql(conn) => diesel::insert_into(crate::schema::scan_events::table)
283                .values(ev)
284                .returning(ScanEvent::as_returning())
285                .get_result::<ScanEvent>(conn)
286                .map_err(Into::into),
287            #[cfg(feature = "sqlite")]
288            Self::Sqlite(conn) => diesel::insert_into(crate::schema::scan_events::table)
289                .values(ev)
290                .returning(ScanEvent::as_returning())
291                .get_result::<ScanEvent>(conn)
292                .map_err(Into::into),
293        }
294    }
295
296    /// Inserts a queued event, or updates the existing pending/retry row for the path.
297    pub fn upsert_pending(
298        &mut self,
299        ev: &NewScanEvent,
300        now: chrono::NaiveDateTime,
301    ) -> anyhow::Result<ScanEvent> {
302        match self {
303            #[cfg(feature = "postgres")]
304            Self::Postgresql(conn) => upsert_pending_pg(conn, ev, now),
305            #[cfg(feature = "sqlite")]
306            Self::Sqlite(conn) => upsert_pending_sqlite(conn, ev, now),
307        }
308    }
309}
310
311#[cfg(feature = "postgres")]
312fn upsert_pending_pg(
313    conn: &mut diesel::PgConnection,
314    ev: &NewScanEvent,
315    now: chrono::NaiveDateTime,
316) -> anyhow::Result<ScanEvent> {
317    use crate::models::ProcessStatus;
318    use crate::schema::scan_events::dsl::{
319        can_process, file_hash, file_path, process_status, updated_at,
320    };
321    use diesel::dsl::case_when;
322    use diesel::upsert::{excluded, DecoratableTarget};
323    use diesel::ExpressionMethods;
324
325    // Keep this predicate aligned with the partial index; Postgres checks that
326    // match at runtime, and the smoke test covers it.
327    let pending: String = ProcessStatus::Pending.into();
328    let retry: String = ProcessStatus::Retry.into();
329
330    diesel::insert_into(crate::schema::scan_events::table)
331        .values(ev)
332        .on_conflict(file_path)
333        .filter_target(process_status.eq_any([pending, retry]))
334        .do_update()
335        .set((
336            updated_at.eq(now),
337            can_process.eq(
338                case_when(can_process.lt(excluded(can_process)), excluded(can_process))
339                    .otherwise(can_process),
340            ),
341            file_hash.eq(case_when(file_hash.is_null(), excluded(file_hash)).otherwise(file_hash)),
342        ))
343        .returning(ScanEvent::as_returning())
344        .get_result::<ScanEvent>(conn)
345        .map_err(Into::into)
346}
347
348#[cfg(feature = "sqlite")]
349fn upsert_pending_sqlite(
350    conn: &mut diesel::SqliteConnection,
351    ev: &NewScanEvent,
352    now: chrono::NaiveDateTime,
353) -> anyhow::Result<ScanEvent> {
354    use crate::models::ProcessStatus;
355    use crate::schema::scan_events::dsl::{
356        can_process, file_hash, file_path, process_status, scan_events, updated_at,
357    };
358    use diesel::{ExpressionMethods, QueryDsl};
359    use diesel::{OptionalExtension, SelectableHelper};
360
361    // Acquire the write lock before reading, including across independent pools.
362    conn.immediate_transaction(|conn| {
363        let pending: String = ProcessStatus::Pending.into();
364        let retry: String = ProcessStatus::Retry.into();
365
366        let existing: Option<ScanEvent> = scan_events
367            .filter(file_path.eq(&ev.file_path))
368            .filter(process_status.eq_any([pending, retry]))
369            .select(ScanEvent::as_select())
370            .first::<ScanEvent>(conn)
371            .optional()?;
372
373        if let Some(existing) = existing {
374            let later_can_process = std::cmp::max(existing.can_process, ev.can_process);
375            let file_hash_value = existing.file_hash.clone().or_else(|| ev.file_hash.clone());
376            diesel::update(&existing)
377                .set((
378                    updated_at.eq(now),
379                    can_process.eq(later_can_process),
380                    file_hash.eq(file_hash_value),
381                ))
382                .get_result::<ScanEvent>(conn)
383                .map_err(Into::into)
384        } else {
385            diesel::insert_into(crate::schema::scan_events::table)
386                .values(ev)
387                .returning(ScanEvent::as_returning())
388                .get_result::<ScanEvent>(conn)
389                .map_err(Into::into)
390        }
391    })
392}
393
394#[doc(hidden)]
395pub type DbPool = Pool<ConnectionManager<AnyConnection>>;
396
397#[doc(hidden)]
398pub fn get_conn(
399    pool: &Pool<ConnectionManager<AnyConnection>>,
400) -> anyhow::Result<PooledConnection<ConnectionManager<AnyConnection>>> {
401    pool.get().context("failed to get connection from pool")
402}
403
404pub fn close_pool(pool: &Pool<ConnectionManager<AnyConnection>>) {
405    match pool.get() {
406        Ok(mut conn) => {
407            if let Err(e) = conn.close() {
408                warn!("failed to close database connection cleanly: {e}");
409            }
410        }
411        Err(e) => {
412            warn!("failed to get connection for pool shutdown: {e}");
413        }
414    }
415}
416
417#[doc(hidden)]
418pub fn get_pool(database_url: &String) -> anyhow::Result<Pool<ConnectionManager<AnyConnection>>> {
419    // First pool fires `AcquireHook { setup: true }` once (VACUUM/WAL), then dropped.
420    let manager = ConnectionManager::<AnyConnection>::new(database_url);
421
422    let setup_pool = Pool::builder()
423        .max_size(1)
424        .connection_customizer(Box::new(AcquireHook { setup: true }))
425        .build(manager)
426        .context("failed to create setup pool")?;
427
428    drop(setup_pool);
429
430    let manager = ConnectionManager::<AnyConnection>::new(database_url);
431
432    let builder = Pool::builder().connection_customizer(Box::new(AcquireHook::default()));
433
434    #[cfg(feature = "sqlite")]
435    let builder = if database_url.starts_with("sqlite://") {
436        builder.max_size(1)
437    } else {
438        builder
439    };
440
441    builder.build(manager).context("failed to create pool")
442}
443
444#[cfg(test)]
445mod tests {
446    use super::*;
447    use std::fs;
448    use tempfile::tempdir;
449
450    #[test]
451    #[cfg(feature = "sqlite")]
452    fn independent_sqlite_pools_coalesce_concurrent_arrivals() {
453        use diesel::QueryDsl;
454        use std::sync::{Arc, Barrier};
455        let dir = tempdir().unwrap();
456        let url = format!("sqlite://{}", dir.path().join("concurrent.db").display());
457        let first = get_pool(&url).unwrap();
458        get_conn(&first).unwrap().migrate().unwrap();
459        let pools = (0..8).map(|_| get_pool(&url).unwrap()).collect::<Vec<_>>();
460        let barrier = Arc::new(Barrier::new(pools.len()));
461        let start = chrono::Utc::now().naive_utc();
462        let handles = pools
463            .into_iter()
464            .enumerate()
465            .map(|(index, pool)| {
466                let barrier = barrier.clone();
467                std::thread::spawn(move || {
468                    barrier.wait();
469                    get_conn(&pool)
470                        .unwrap()
471                        .upsert_pending(
472                            &NewScanEvent {
473                                can_process: start + chrono::Duration::seconds(index as i64),
474                                file_hash: (index == 3).then(|| "hash".into()),
475                                ..Default::default()
476                            },
477                            start,
478                        )
479                        .unwrap()
480                        .id
481                })
482            })
483            .collect::<Vec<_>>();
484        let mut ids = handles
485            .into_iter()
486            .map(|handle| handle.join().unwrap())
487            .collect::<Vec<_>>();
488        ids.sort();
489        ids.dedup();
490        assert_eq!(ids.len(), 1);
491        let saved = crate::schema::scan_events::table
492            .find(&ids[0])
493            .first::<ScanEvent>(&mut get_conn(&first).unwrap())
494            .unwrap();
495        assert_eq!(saved.can_process, start + chrono::Duration::seconds(7));
496        assert_eq!(saved.file_hash.as_deref(), Some("hash"));
497    }
498
499    #[test]
500    #[cfg(feature = "sqlite")]
501    fn found_update_preserves_dedupe_timer() {
502        let mut conn = AnyConnection::establish(":memory:").unwrap();
503        conn.migrate().unwrap();
504        let event = NewScanEvent::default();
505        let snapshot = conn.insert_and_return(&event).unwrap();
506        let later = NewScanEvent {
507            can_process: event.can_process + chrono::Duration::seconds(60),
508            ..event
509        };
510        conn.upsert_pending(&later, chrono::Utc::now().naive_utc())
511            .unwrap();
512        let saved = conn
513            .update_found(&snapshot, "found", chrono::Utc::now().naive_utc())
514            .unwrap()
515            .unwrap();
516        assert_eq!(saved.can_process, later.can_process);
517        assert_eq!(saved.found_status, "found");
518    }
519
520    #[test]
521    #[cfg(feature = "sqlite")]
522    fn terminal_update_clears_retry_timestamp() {
523        use crate::schema::scan_events::dsl::*;
524        use diesel::{ExpressionMethods, QueryDsl};
525
526        let mut conn = AnyConnection::establish(":memory:").unwrap();
527        conn.migrate().unwrap();
528        let event = conn.insert_and_return(&NewScanEvent::default()).unwrap();
529        diesel::update(scan_events.find(&event.id))
530            .set((
531                process_status.eq("retry"),
532                next_retry_at.eq(Some(event.created_at)),
533            ))
534            .execute(&mut conn)
535            .unwrap();
536        let mut snapshot = scan_events
537            .find(&event.id)
538            .first::<ScanEvent>(&mut conn)
539            .unwrap();
540        let previous = snapshot.clone();
541        snapshot.process_status = "failed".into();
542        snapshot.next_retry_at = None;
543        let saved = conn.update_process(&previous, &snapshot).unwrap().unwrap();
544        assert_eq!(saved.process_status, "failed");
545        assert_eq!(saved.next_retry_at, None);
546    }
547
548    #[test]
549    fn test_pre_init_memory_db_skipped() {
550        let result = AnyConnection::pre_init("sqlite://:memory:");
551        assert!(result.is_ok());
552    }
553
554    #[test]
555    fn test_pre_init_creates_directory() {
556        let tmp = tempdir().unwrap();
557        let db_path = tmp.path().join("subdir").join("test.db");
558        let url = format!("sqlite://{}", db_path.display());
559
560        let result = AnyConnection::pre_init(&url);
561        assert!(result.is_ok());
562        assert!(db_path.parent().unwrap().exists());
563    }
564
565    #[test]
566    fn test_pre_init_no_parent_directory() {
567        let result = AnyConnection::pre_init("sqlite://test.db");
568        assert!(result.is_ok());
569    }
570
571    #[test]
572    fn test_pre_init_writable_directory_succeeds() {
573        let tmp = tempdir().unwrap();
574        let subdir = tmp.path().join("writable");
575        fs::create_dir(&subdir).unwrap();
576
577        let db_path = subdir.join("test.db");
578        let url = format!("sqlite://{}", db_path.display());
579
580        let result = AnyConnection::pre_init(&url);
581        assert!(result.is_ok());
582    }
583
584    #[cfg(unix)]
585    #[test]
586    fn test_pre_init_existing_unwritable_directory_fails_with_context() {
587        use std::os::unix::fs::PermissionsExt;
588
589        let tmp = tempdir().unwrap();
590        let subdir = tmp.path().join("readonly");
591        fs::create_dir(&subdir).unwrap();
592        fs::set_permissions(&subdir, fs::Permissions::from_mode(0o555)).unwrap();
593
594        let db_path = subdir.join("test.db");
595        let url = format!("sqlite://{}", db_path.display());
596
597        let result = AnyConnection::pre_init(&url);
598
599        fs::set_permissions(&subdir, fs::Permissions::from_mode(0o755)).unwrap();
600        let err = result.expect_err("unwritable database directory should fail pre-init");
601        let err = err.to_string();
602        assert!(err.contains("database directory is not writable"));
603        assert!(err.contains(&subdir.display().to_string()));
604    }
605
606    #[test]
607    fn test_pre_init_postgres_skipped() {
608        let result = AnyConnection::pre_init("postgres://localhost/test");
609        assert!(result.is_ok());
610    }
611
612    #[test]
613    #[cfg(feature = "sqlite")]
614    fn test_close_pool_cleans_up_wal_files() {
615        let tmp = tempdir().unwrap();
616        let db_path = tmp.path().join("test.db");
617        let url = format!("sqlite://{}", db_path.display());
618
619        AnyConnection::pre_init(&url).unwrap();
620        let pool = get_pool(&url).unwrap();
621
622        // Get a connection to trigger WAL mode and create the db
623        {
624            let mut conn = get_conn(&pool).unwrap();
625            conn.migrate().unwrap();
626        }
627
628        // WAL files may exist at this point
629        close_pool(&pool);
630        drop(pool);
631
632        // Verify no WAL files remain
633        let wal_path = tmp.path().join("test.db-wal");
634        let shm_path = tmp.path().join("test.db-shm");
635        assert!(!wal_path.exists(), "WAL file should be cleaned up");
636        assert!(!shm_path.exists(), "SHM file should be cleaned up");
637    }
638
639    #[test]
640    #[cfg(feature = "sqlite")]
641    fn dedupe_migration_merges_max_can_process_into_survivor() {
642        use crate::models::ProcessStatus;
643        use crate::schema::scan_events::dsl::{file_path, process_status, scan_events};
644        use chrono::{NaiveDate, NaiveDateTime, NaiveTime};
645        use diesel::{ExpressionMethods, QueryDsl, RunQueryDsl};
646
647        let tmp = tempdir().unwrap();
648        let db_path = tmp.path().join("test.db");
649        let url = format!("sqlite://{}", db_path.display());
650
651        AnyConnection::pre_init(&url).unwrap();
652        let pool = get_pool(&url).unwrap();
653        let mut conn = get_conn(&pool).unwrap();
654
655        conn.batch_execute(
656            r#"
657            CREATE TABLE scan_events (
658                id TEXT PRIMARY KEY NOT NULL,
659                event_source TEXT NOT NULL,
660                event_timestamp TIMESTAMP DEFAULT CURRENT_TIMESTAMP NOT NULL,
661                file_path TEXT NOT NULL,
662                file_hash TEXT,
663                process_status TEXT NOT NULL DEFAULT 'pending',
664                found_status TEXT NOT NULL DEFAULT 'not_found',
665                failed_times INTEGER DEFAULT 0 NOT NULL,
666                next_retry_at TIMESTAMP,
667                targets_hit TEXT DEFAULT '' NOT NULL,
668                found_at TIMESTAMP,
669                processed_at TIMESTAMP,
670                created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP NOT NULL,
671                updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP NOT NULL,
672                can_process TIMESTAMP NOT NULL DEFAULT "2024-10-14T12:00:00.000"
673            );
674
675            CREATE TABLE __diesel_schema_migrations (
676                version VARCHAR(50) PRIMARY KEY NOT NULL,
677                run_on TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP
678            );
679
680            INSERT INTO __diesel_schema_migrations (version) VALUES
681                ('20240829125750'),
682                ('20240905143749'),
683                ('20240906161345'),
684                ('20241012130403'),
685                ('20241205114327'),
686                ('20241205115656'),
687                ('202512300005460000'),
688                ('20260519000001');
689
690            INSERT INTO scan_events (
691                id, event_source, file_path, file_hash, process_status,
692                updated_at, created_at, event_timestamp, can_process
693            ) VALUES
694                (
695                    'older-long-wait', 'sonarr', '/media/migrate.mkv', 'sha256:migrate', 'pending',
696                    '2026-01-01 00:00:00', '2026-01-01 00:00:00',
697                    '2026-01-01 00:00:00', '2026-01-01 03:00:00'
698                ),
699                (
700                    'newer-short-wait', 'notify', '/media/migrate.mkv', NULL, 'retry',
701                    '2026-01-01 01:00:00', '2026-01-01 01:00:00',
702                    '2026-01-01 01:00:00', '2026-01-01 02:00:00'
703                );
704            "#,
705        )
706        .unwrap();
707
708        conn.migrate().unwrap();
709
710        let pending: String = ProcessStatus::Pending.into();
711        let retry: String = ProcessStatus::Retry.into();
712        let rows = scan_events
713            .filter(file_path.eq("/media/migrate.mkv"))
714            .filter(process_status.eq_any([pending, retry]))
715            .load::<ScanEvent>(&mut conn)
716            .unwrap();
717
718        assert_eq!(rows.len(), 1, "migration should leave one non-terminal row");
719        assert_eq!(rows[0].id, "newer-short-wait", "newest row should survive");
720        assert_eq!(
721            rows[0].can_process,
722            NaiveDateTime::new(
723                NaiveDate::from_ymd_opt(2026, 1, 1).unwrap(),
724                NaiveTime::from_hms_opt(3, 0, 0).unwrap(),
725            ),
726            "survivor should inherit the duplicate group's longest wait"
727        );
728        assert_eq!(
729            rows[0].file_hash,
730            Some("sha256:migrate".to_string()),
731            "survivor should inherit a duplicate's hash when it has none"
732        );
733    }
734}