diff --git a/crates/codex-db/src/repositories/read_progress.rs b/crates/codex-db/src/repositories/read_progress.rs index 5816c73e..9ae1a1ba 100644 --- a/crates/codex-db/src/repositories/read_progress.rs +++ b/crates/codex-db/src/repositories/read_progress.rs @@ -8,7 +8,7 @@ use crate::entities::reading_sessions::SessionKind; use crate::entities::{read_progress, read_progress::Entity as ReadProgress}; use crate::repositories::ReadCompletionRepository; use crate::repositories::reading_sessions::{ - AppendOutcome, DeviceContext, NewSession, ReadingSessionRepository, fold, + AppendOutcome, DeviceContext, NewSession, ProgressProjection, ReadingSessionRepository, fold, }; use anyhow::{Result, anyhow}; use chrono::Utc; @@ -16,6 +16,35 @@ use sea_orm::*; use std::collections::HashMap; use uuid::Uuid; +/// How many times [`ReadProgressRepository::record_session`] retries a +/// transaction that lost the SQLite write lock. With the backoff doubling from +/// 50ms, the last attempt starts about 1.5s after the first. +const SQLITE_BUSY_RETRIES: u32 = 5; + +/// Whether an error is SQLite refusing a write because another connection holds +/// the lock (`SQLITE_BUSY`, including `SQLITE_BUSY_SNAPSHOT` under WAL, or +/// `SQLITE_LOCKED`). Matched on the primary result code, since the extended +/// codes vary with the journal mode. +fn is_sqlite_busy(err: &anyhow::Error) -> bool { + err.chain().any(|cause| { + let Some(DbErr::Exec(runtime) | DbErr::Query(runtime) | DbErr::Conn(runtime)) = + cause.downcast_ref::() + else { + return false; + }; + let RuntimeErr::SqlxError(sea_orm::sqlx::Error::Database(db_err)) = runtime else { + return false; + }; + db_err + .try_downcast_ref::() + .is_some() + && db_err + .code() + .and_then(|code| code.parse::().ok()) + .is_some_and(|code| matches!(code & 0xff, 5 | 6)) + }) +} + pub struct ReadProgressRepository; impl ReadProgressRepository { @@ -242,9 +271,38 @@ impl ReadProgressRepository { /// The entry point for the sessions API. Everything runs in one transaction /// so the log and its projections cannot disagree, and a replayed id is /// reported rather than applied twice. + /// + /// On SQLite a transaction that loses the write lock to a concurrent one is + /// retried from the start. A client sending the same queue twice at once + /// hits exactly this: both requests read before either writes, and the one + /// that loses cannot proceed on what it read. Running it again re-reads, + /// finds the sibling's committed row, and reports the replay as a + /// duplicate. PostgreSQL waits on the row instead, so it never gets here. pub async fn record_session( db: &DatabaseConnection, session: NewSession, + ) -> Result<(AppendOutcome, Option)> { + let mut retries = 0; + loop { + match Self::record_session_once(db, session.clone()).await { + Err(err) if retries < SQLITE_BUSY_RETRIES && is_sqlite_busy(&err) => { + retries += 1; + let backoff = std::time::Duration::from_millis(50 << (retries - 1)); + tracing::warn!( + session_id = %session.id, + retry = retries, + "reading session lost the SQLite write lock, retrying: {err}" + ); + tokio::time::sleep(backoff).await; + } + result => return result, + } + } + } + + async fn record_session_once( + db: &DatabaseConnection, + session: NewSession, ) -> Result<(AppendOutcome, Option)> { let txn = db.begin().await?; let now = Utc::now(); @@ -272,6 +330,9 @@ impl ReadProgressRepository { /// be one. That happens after a reset with no reading since: marking a book /// unread has always removed the row outright rather than zeroing it, and /// callers check for its absence. + /// + /// When nothing in the current pass moved the reader (only time spent was + /// reported), the row is left exactly as it is and returned as found. async fn refold( db: &C, user_id: Uuid, @@ -282,9 +343,13 @@ impl ReadProgressRepository { let sessions = ReadingSessionRepository::load_current_pass(db, user_id, book_id).await?; let folded = fold(&sessions); - let Some(projected) = folded.progress else { - Self::delete_row_in(db, user_id, book_id).await?; - return Ok(None); + let projected = match folded.progress { + ProgressProjection::Row(projected) => projected, + ProgressProjection::Absent => { + Self::delete_row_in(db, user_id, book_id).await?; + return Ok(None); + } + ProgressProjection::Unchanged => return Self::get_in(db, user_id, book_id).await, }; let existing = Self::get_in(db, user_id, book_id).await?; @@ -703,6 +768,51 @@ mod tests { use codex_models::ScanningStrategy; use codex_utils::password; + /// Only lock contention is retried. Anything else, a unique violation + /// above all, must surface on the first attempt rather than be re-run. + #[test] + fn only_sqlite_lock_contention_counts_as_busy() { + assert!(!is_sqlite_busy(&anyhow!("database is locked"))); + assert!(!is_sqlite_busy(&anyhow::Error::from( + DbErr::RecordNotFound("gone".to_string()) + ))); + assert!(!is_sqlite_busy(&anyhow::Error::from( + DbErr::RecordNotInserted + ))); + } + + /// The positive case, from a real lock: one connection holds the write + /// lock, and another tries to write without waiting for it. + #[tokio::test] + async fn a_held_sqlite_write_lock_counts_as_busy() { + let dir = tempfile::tempdir().unwrap(); + let url = format!("sqlite://{}?mode=rwc", dir.path().join("busy.db").display()); + let holder = sea_orm::Database::connect(&url).await.unwrap(); + holder + .execute_unprepared("CREATE TABLE t (x INTEGER)") + .await + .unwrap(); + let held = holder.begin().await.unwrap(); + held.execute_unprepared("INSERT INTO t VALUES (1)") + .await + .unwrap(); + + let mut options = sea_orm::ConnectOptions::new(&url); + options.max_connections(1); + let other = sea_orm::Database::connect(options).await.unwrap(); + other + .execute_unprepared("PRAGMA busy_timeout = 0") + .await + .unwrap(); + let err = other + .execute_unprepared("INSERT INTO t VALUES (2)") + .await + .expect_err("the write lock is held"); + + assert!(is_sqlite_busy(&anyhow::Error::from(err))); + held.rollback().await.unwrap(); + } + async fn create_test_user(db: &DatabaseConnection) -> users::Model { let password_hash = password::hash_password("password").unwrap(); let user = users::Model { diff --git a/crates/codex-db/src/repositories/reading_sessions/fold.rs b/crates/codex-db/src/repositories/reading_sessions/fold.rs index 0dfdf9f3..a2962da4 100644 --- a/crates/codex-db/src/repositories/reading_sessions/fold.rs +++ b/crates/codex-db/src/repositories/reading_sessions/fold.rs @@ -55,15 +55,41 @@ pub struct FoldedCompletion { pub completed_at: DateTime, } -/// The result of folding one `(user, book)` slice. +/// What the fold says about the `read_progress` row. #[derive(Clone, Debug, PartialEq, Default)] -pub struct Fold { - /// `None` means the `read_progress` row must **not exist**. +pub enum ProgressProjection { + /// The row should hold exactly this. + Row(FoldedProgress), + /// The row must **not exist**. /// /// Marking a book unread deletes that row today, and callers assert on its /// absence rather than on a zeroed row, so a pass consisting only of a /// reset has to project to nothing at all. - pub progress: Option, + #[default] + Absent, + /// Nothing in the pass moved the reader, so whatever the row holds stands. + /// + /// Distinct from [`Self::Absent`] because the row is not always the fold + /// of its log: importing a reading-progress export writes `read_progress` + /// directly, and a session that only reports time spent must not delete a + /// row like that. + Unchanged, +} + +impl ProgressProjection { + /// The projected row, when the fold produced one. + pub fn row(&self) -> Option<&FoldedProgress> { + match self { + Self::Row(progress) => Some(progress), + Self::Absent | Self::Unchanged => None, + } + } +} + +/// The result of folding one `(user, book)` slice. +#[derive(Clone, Debug, PartialEq, Default)] +pub struct Fold { + pub progress: ProgressProjection, /// Set when the current pass has finished. Earlier passes are not /// re-reported: their completions were banked while they were current. pub completion: Option, @@ -92,16 +118,27 @@ pub fn fold(sessions: &[Session]) -> Fold { let mut current: Vec<&Session> = sessions.iter().filter(|s| s.pass == current_pass).collect(); current.sort_by_key(|s| sort_key(s)); - // A pass whose only events are resets is a book marked unread and not read - // since. That must project to no row, not to a zeroed one. - let mut reading = current - .iter() - .filter(|s| s.session_kind() != SessionKind::Reset) - .peekable(); - if reading.peek().is_none() { - return Fold::default(); + // Only sessions that moved the reader shape the row. One that reports time + // spent and nothing else (a reader left open, an unattended auto-advance) + // still counts towards reading statistics, which read the log directly, + // but it says nothing about where the reader is. + let reading: Vec<&&Session> = current.iter().filter(|s| moves_the_reader(s)).collect(); + if reading.is_empty() { + // A pass that holds a reset is a book marked unread and not read since, + // which must project to no row, not to a zeroed one. Without a reset, + // nothing here has any say over the row. + let reset = current + .iter() + .any(|s| s.session_kind() == SessionKind::Reset); + return Fold { + progress: if reset { + ProgressProjection::Absent + } else { + ProgressProjection::Unchanged + }, + completion: None, + }; } - let reading: Vec<&&Session> = reading.collect(); // The pass began when its first reading happened. This delimits the pass // for the completion guard, and it has to survive a back-tap (which is just @@ -111,12 +148,15 @@ pub fn fold(sessions: &[Session]) -> Fold { .iter() .map(|s| s.client_started_at) .min() - .expect("non-empty by the peek above"); - let updated_at = current + .expect("non-empty, checked above"); + // Arrival time rather than client time: this orders the shelves by when the + // server last learned something new. Resets do not count, since a pass that + // has reading in it began after its reset. + let updated_at = reading .iter() .map(|s| s.server_recorded_at) .max() - .expect("non-empty because reading is non-empty"); + .expect("non-empty, checked above"); let mut current_page = 0; let mut progress_percentage = None; @@ -154,7 +194,7 @@ pub fn fold(sessions: &[Session]) -> Fold { completed = false; completed_at = None; } - SessionKind::Reset => unreachable!("resets are filtered out above"), + SessionKind::Reset => unreachable!("resets do not move the reader"), } } @@ -164,7 +204,7 @@ pub fn fold(sessions: &[Session]) -> Fold { }); Fold { - progress: Some(FoldedProgress { + progress: ProgressProjection::Row(FoldedProgress { pass: current_pass, current_page, progress_percentage, @@ -178,6 +218,21 @@ pub fn fold(sessions: &[Session]) -> Fold { } } +/// Whether a session changes what the progress row should say: it carries a +/// position, or it is a completion. Resets are handled as pass boundaries, and +/// everything else reports time only. +fn moves_the_reader(session: &Session) -> bool { + match session.session_kind() { + SessionKind::Reset => false, + SessionKind::Completed => true, + SessionKind::Progress => { + session.to_page.is_some() + || session.to_percentage.is_some() + || session.r2_progression.is_some() + } + } +} + #[cfg(test)] mod tests { use super::*; @@ -238,6 +293,12 @@ mod tests { Self::new(SessionKind::Reset, ended) } + /// A session that reports time spent and nothing else: no page, no + /// percentage, no progression. A reader left open unattended. + fn time_only(ended: DateTime) -> Self { + Self::new(SessionKind::Progress, ended) + } + fn pass(mut self, pass: i32) -> Self { self.pass = pass; self @@ -298,7 +359,11 @@ mod tests { } fn progress_of(sessions: Vec) -> FoldedProgress { - fold_of(sessions).progress.expect("expected a progress row") + fold_of(sessions) + .progress + .row() + .cloned() + .expect("expected a progress row") } // ------------------------------------------------------------------ @@ -353,8 +418,8 @@ mod tests { let mut reversed: Vec = sessions().into_iter().map(S::build).collect(); reversed.reverse(); - assert_eq!(forward.progress.unwrap().current_page, 40); - assert_eq!(fold(&reversed).progress.unwrap().current_page, 40); + assert_eq!(forward.progress.row().unwrap().current_page, 40); + assert_eq!(fold(&reversed).progress.row().unwrap().current_page, 40); } /// Simultaneous client times fall back to arrival order rather than being @@ -446,7 +511,7 @@ mod tests { S::reset(at(20)).pass(2), ]); - assert_eq!(folded.progress, None); + assert_eq!(folded.progress, ProgressProjection::Absent); assert_eq!(folded.completion, None); } @@ -458,7 +523,7 @@ mod tests { S::reset(at(20)).pass(3), ]); - assert_eq!(folded.progress, None); + assert_eq!(folded.progress, ProgressProjection::Absent); } /// Reading after a reset starts fresh: the new pass's position, and a @@ -471,7 +536,11 @@ mod tests { S::progress(5, at(30)).pass(2), ]); - let progress = folded.progress.expect("expected a progress row"); + let progress = folded + .progress + .row() + .cloned() + .expect("expected a progress row"); assert_eq!(progress.pass, 2); assert_eq!(progress.current_page, 5); assert!(!progress.completed); @@ -511,7 +580,7 @@ mod tests { S::progress(3, at(20)).pass(2), ]); - let progress = folded.progress.unwrap(); + let progress = folded.progress.row().cloned().unwrap(); assert_eq!(progress.current_page, 3); assert_eq!(folded.completion, None); } @@ -546,13 +615,122 @@ mod tests { } #[test] - fn updated_at_is_the_latest_arrival() { + fn updated_at_is_the_latest_arrival_that_moved_something() { let progress = progress_of(vec![ S::progress(10, at(10)).recorded(at(15)), S::progress(20, at(20)).recorded(at(90)), + S::time_only(at(30)).recorded(at(120)), ]); - assert_eq!(progress.updated_at, at(90)); + assert_eq!( + progress.updated_at, + at(90), + "a session that moved nothing must not bump the row up the shelves" + ); + } + + // ------------------------------------------------------------------ + // Sessions that report time only + // ------------------------------------------------------------------ + + /// A reader left open on a book reports time spent with no position. That + /// is reading time, not reading progress. + #[test] + fn a_time_only_session_leaves_the_position_alone() { + let progress = progress_of(vec![ + S::progress(2, at(10)) + .percentage(0.1) + .r2(r#"{"locator":"a"}"#), + S::time_only(at(20)), + ]); + + assert_eq!(progress.current_page, 2); + assert_eq!(progress.progress_percentage, Some(0.1)); + assert_eq!( + progress.r2_progression.as_deref(), + Some(r#"{"locator":"a"}"#) + ); + } + + /// Opening a finished book and leaving it is not reading it again. + #[test] + fn a_time_only_session_does_not_reopen_a_finished_book() { + let folded = fold_of(vec![S::completed(50, at(10)), S::time_only(at(20))]); + + let progress = folded + .progress + .row() + .cloned() + .expect("expected a progress row"); + assert!(progress.completed); + assert_eq!(progress.completed_at, Some(at(10))); + assert_eq!( + folded.completion, + Some(FoldedCompletion { + started_at: at(10), + completed_at: at(10), + }) + ); + } + + /// Nothing moved, so nothing is projected, and the row (if any) stands. + /// Not `Absent`: that would delete a row this pass has no say over. + #[test] + fn a_pass_of_only_time_only_sessions_leaves_the_row_alone() { + let folded = fold_of(vec![S::time_only(at(10)), S::time_only(at(20))]); + + assert_eq!(folded.progress, ProgressProjection::Unchanged); + assert_eq!(folded.completion, None); + } + + /// After a reset the row must not exist, and opening the book without + /// reading does not bring it back. + #[test] + fn a_time_only_session_after_a_reset_still_projects_no_row() { + let folded = fold_of(vec![ + S::completed(50, at(10)).pass(1), + S::reset(at(20)).pass(2), + S::time_only(at(30)).pass(2), + ]); + + assert_eq!(folded.progress, ProgressProjection::Absent); + } + + #[test] + fn a_positioned_session_after_a_time_only_one_wins() { + let progress = progress_of(vec![ + S::progress(2, at(10)), + S::time_only(at(20)), + S::progress(7, at(30)).recorded(at(40)), + ]); + + assert_eq!(progress.current_page, 7); + assert_eq!(progress.updated_at, at(40)); + } + + /// A stray open before the real start must not backdate the pass. + #[test] + fn started_at_ignores_time_only_sessions() { + let progress = progress_of(vec![ + S::time_only(at(5)).started(at(5)), + S::progress(10, at(30)).started(at(25)), + ]); + + assert_eq!(progress.started_at, at(25)); + } + + /// A completion is a change of state even without a position. + #[test] + fn a_completion_without_a_position_still_moves_the_row() { + let folded = fold_of(vec![S::new(SessionKind::Completed, at(10))]); + + let progress = folded + .progress + .row() + .cloned() + .expect("expected a progress row"); + assert!(progress.completed); + assert_eq!(progress.completed_at, Some(at(10))); } // ------------------------------------------------------------------ @@ -638,7 +816,7 @@ mod tests { S::progress(20, at(30)).device("ipad"), ]); - let progress = folded.progress.unwrap(); + let progress = folded.progress.row().cloned().unwrap(); assert!(!progress.completed); assert_eq!(progress.current_page, 20); assert_eq!(folded.completion, None); @@ -656,7 +834,11 @@ mod tests { session.kind = "some-future-kind".to_string(); let folded = fold(&[session]); - let progress = folded.progress.expect("expected a progress row"); + let progress = folded + .progress + .row() + .cloned() + .expect("expected a progress row"); assert_eq!(progress.current_page, 10); assert!(!progress.completed); } diff --git a/crates/codex-db/src/repositories/reading_sessions/mod.rs b/crates/codex-db/src/repositories/reading_sessions/mod.rs index 8e21e865..969a39a5 100644 --- a/crates/codex-db/src/repositories/reading_sessions/mod.rs +++ b/crates/codex-db/src/repositories/reading_sessions/mod.rs @@ -7,7 +7,7 @@ pub mod fold; -pub use fold::{Fold, FoldedCompletion, FoldedProgress, fold}; +pub use fold::{Fold, FoldedCompletion, FoldedProgress, ProgressProjection, fold}; use crate::entities::{ reading_sessions, @@ -15,6 +15,7 @@ use crate::entities::{ }; use anyhow::Result; use chrono::{DateTime, Duration, Utc}; +use sea_orm::sea_query::OnConflict; use sea_orm::*; use uuid::Uuid; @@ -470,7 +471,32 @@ impl ReadingSessionRepository { client_ended_at: Set(session.client_ended_at), server_recorded_at: Set(now), }; - let inserted = model.insert(db).await?; + // The check above reads before it writes, so it cannot see a concurrent + // request's uncommitted row with the same id: a client that sends its + // queue twice at once gets past it on both requests. The insert has to + // tolerate the collision itself. PostgreSQL waits for the sibling's + // transaction and then skips the row; SQLite cannot get that far (the + // second writer loses the lock instead), which `record_session` retries. + // + // `exec` rather than `exec_with_returning`: only `exec` reports an + // insert that did nothing as a conflict. Every column is set here, so + // the stored row is known without reading it back. + let inserted = model.clone().try_into_model()?; + match ReadingSessions::insert(model) + .on_conflict( + OnConflict::column(reading_sessions::Column::Id) + .do_nothing() + .to_owned(), + ) + .do_nothing() + .exec(db) + .await? + { + TryInsertResult::Inserted(_) => {} + TryInsertResult::Conflicted | TryInsertResult::Empty => { + return Ok(AppendOutcome::Duplicate); + } + } // A measured session describes a whole sitting, including everything // the position-only writes made during it already said. Fold them in diff --git a/tests/api/reading_sessions.rs b/tests/api/reading_sessions.rs index 2e71765b..e01e9ff7 100644 --- a/tests/api/reading_sessions.rs +++ b/tests/api/reading_sessions.rs @@ -671,3 +671,232 @@ async fn a_batch_spanning_books_returns_progress_for_each() { ); assert_eq!(page_for(second.id), 25); } + +// --------------------------------------------------------------------------- +// Sessions that report time spent and nothing else +// --------------------------------------------------------------------------- + +async fn post_sessions( + state: &std::sync::Arc, + token: &str, + sessions: Vec, +) -> RecordReadingSessionsResponse { + let app = create_test_router(state.clone()).await; + let request = post_json_request_with_auth( + "/api/v1/reading-sessions", + &json!({ "sessions": sessions }), + token, + ); + let (status, response): (StatusCode, Option) = + make_json_request(app, request).await; + assert_eq!(status, StatusCode::OK); + response.expect("expected a JSON body") +} + +/// A reader left open: time spent, no position. The shape an unattended iPad +/// queues and uploads the next morning. +fn time_only(id: Uuid, book_id: Uuid, from: i64, to: i64) -> serde_json::Value { + json!({ + "id": id, + "bookId": book_id, + "deviceId": "ipad", + "kind": "progress", + "activeDurationMs": 953, + "pagesRead": 1, + "clientStartedAt": at(from), + "clientEndedAt": at(to), + }) +} + +/// The reported failure: an unattended open must not move a book to the top of +/// Keep Reading. The row is left exactly as it was, and the time still lands +/// in the log for the statistics. +#[tokio::test] +async fn a_time_only_session_leaves_progress_exactly_as_it_was() { + let (db, _temp_dir) = setup_test_db().await; + let book = create_book(&db, "/test/book.cbz").await; + let state = create_test_auth_state(db.clone()).await; + let (user_id, token) = create_admin_and_token(&db, &state).await; + + post_sessions( + &state, + &token, + vec![session(Uuid::new_v4(), book.id, "phone", 2, 0, 10)], + ) + .await; + let before = ReadProgressRepository::get_by_user_and_book(&db, user_id, book.id) + .await + .unwrap() + .unwrap(); + + let id = Uuid::new_v4(); + let response = post_sessions(&state, &token, vec![time_only(id, book.id, 60, 61)]).await; + assert_eq!(response.accepted, vec![id]); + + let after = ReadProgressRepository::get_by_user_and_book(&db, user_id, book.id) + .await + .unwrap() + .unwrap(); + assert_eq!( + after, before, + "a session that moved nothing must not touch the row" + ); + + let sessions = + codex::db::repositories::ReadingSessionRepository::load_for_book(&db, user_id, book.id) + .await + .unwrap(); + let stored = sessions + .iter() + .find(|s| s.id == id) + .expect("the session is still recorded"); + assert_eq!(stored.active_duration_ms, Some(953)); + assert_eq!(stored.pages_read, Some(1)); +} + +/// Opening a finished book and leaving it there does not un-finish it. +#[tokio::test] +async fn a_time_only_session_does_not_reopen_a_finished_book() { + let (db, _temp_dir) = setup_test_db().await; + let book = create_book(&db, "/test/book.cbz").await; + let state = create_test_auth_state(db.clone()).await; + let (user_id, token) = create_admin_and_token(&db, &state).await; + + let mut finished = session(Uuid::new_v4(), book.id, "phone", 100, 0, 10); + finished["kind"] = json!("completed"); + post_sessions(&state, &token, vec![finished]).await; + let before = ReadProgressRepository::get_by_user_and_book(&db, user_id, book.id) + .await + .unwrap() + .unwrap(); + assert!(before.completed); + + post_sessions( + &state, + &token, + vec![time_only(Uuid::new_v4(), book.id, 60, 61)], + ) + .await; + + let after = ReadProgressRepository::get_by_user_and_book(&db, user_id, book.id) + .await + .unwrap() + .unwrap(); + assert_eq!(after, before); +} + +/// Opening an unread book without reading it does not put it in Keep Reading. +#[tokio::test] +async fn a_time_only_session_on_an_unread_book_creates_no_progress() { + let (db, _temp_dir) = setup_test_db().await; + let book = create_book(&db, "/test/book.cbz").await; + let state = create_test_auth_state(db.clone()).await; + let (user_id, token) = create_admin_and_token(&db, &state).await; + + let id = Uuid::new_v4(); + let response = post_sessions(&state, &token, vec![time_only(id, book.id, 0, 1)]).await; + assert_eq!(response.accepted, vec![id]); + + let progress = ReadProgressRepository::get_by_user_and_book(&db, user_id, book.id) + .await + .unwrap(); + assert!( + progress.is_none(), + "no row, so the book stays off the in-progress shelf" + ); +} + +/// A row can exist with no positioned session behind it: importing a +/// reading-progress export writes `read_progress` directly. A time-only session +/// must not delete or rewrite it. +#[tokio::test] +async fn a_time_only_session_keeps_a_row_the_log_does_not_explain() { + use codex::db::entities::read_progress; + use sea_orm::{ActiveModelTrait, Set}; + + let (db, _temp_dir) = setup_test_db().await; + let book = create_book(&db, "/test/book.cbz").await; + let state = create_test_auth_state(db.clone()).await; + let (user_id, token) = create_admin_and_token(&db, &state).await; + + let imported = read_progress::ActiveModel { + id: Set(Uuid::new_v4()), + user_id: Set(user_id), + book_id: Set(book.id), + current_page: Set(42), + progress_percentage: Set(None), + completed: Set(false), + started_at: Set(at(0)), + updated_at: Set(at(10)), + completed_at: Set(None), + r2_progression: Set(None), + } + .insert(&db) + .await + .unwrap(); + + post_sessions( + &state, + &token, + vec![time_only(Uuid::new_v4(), book.id, 60, 61)], + ) + .await; + + let after = ReadProgressRepository::get_by_user_and_book(&db, user_id, book.id) + .await + .unwrap(); + assert_eq!(after, Some(imported)); +} + +/// The boundary of the rule. A measured session with no position of its own +/// absorbs the position writes made during its sitting, and so does move the +/// reader: those writes said where they were. +#[tokio::test] +async fn a_measured_session_that_absorbs_position_writes_moves_the_row() { + let (db, _temp_dir) = setup_test_db().await; + let book = create_book(&db, "/test/book.cbz").await; + let state = create_test_auth_state(db.clone()).await; + let (user_id, token) = create_admin_and_token(&db, &state).await; + + let started = Utc::now() - Duration::minutes(20); + let app = create_test_router(state.clone()).await; + let mut request = put_json_request_with_auth( + &format!("/api/v1/books/{}/progress", book.id), + &json!({ "currentPage": 9 }), + &token, + ); + request + .headers_mut() + .insert("x-codex-device-id", "browser-abc".parse().unwrap()); + let (status, _body) = make_request(app, request).await; + assert_eq!(status, StatusCode::OK); + + post_sessions( + &state, + &token, + vec![json!({ + "id": Uuid::new_v4(), + "bookId": book.id, + "deviceId": "browser-abc", + "kind": "progress", + "activeDurationMs": 600_000, + "clientStartedAt": started, + "clientEndedAt": Utc::now(), + })], + ) + .await; + + let sessions = + codex::db::repositories::ReadingSessionRepository::load_for_book(&db, user_id, book.id) + .await + .unwrap(); + assert_eq!(sessions.len(), 1, "the position write was absorbed"); + assert_eq!(sessions[0].to_page, Some(9)); + + let progress = ReadProgressRepository::get_by_user_and_book(&db, user_id, book.id) + .await + .unwrap() + .unwrap(); + assert_eq!(progress.current_page, 9); + assert_eq!(progress.updated_at, sessions[0].server_recorded_at); +} diff --git a/tests/db/reading_sessions.rs b/tests/db/reading_sessions.rs index d11c616a..629a75bb 100644 --- a/tests/db/reading_sessions.rs +++ b/tests/db/reading_sessions.rs @@ -17,11 +17,11 @@ use codex::db::ScanningStrategy; use codex::db::entities::reading_sessions; use codex::db::entities::reading_sessions::SessionKind; use codex::db::repositories::{ - BookRepository, DeviceContext, LibraryRepository, NewSession, ReadProgressRepository, - ReadingSessionRepository, SeriesRepository, UserRepository, + AppendOutcome, BookRepository, DeviceContext, LibraryRepository, NewSession, + ReadProgressRepository, ReadingSessionRepository, SeriesRepository, UserRepository, }; use common::*; -use sea_orm::DatabaseConnection; +use sea_orm::{DatabaseConnection, TransactionTrait}; use uuid::Uuid; const DEVICE: &str = "browser-1"; @@ -265,6 +265,66 @@ async fn compat_client_unaffected_sqlite() { exercise_compat_client_unaffected(&db).await; } +/// A client that sends the same queue twice at once: both requests carry the +/// same session id, and the second arrives while the first is still inside its +/// transaction. +/// +/// The replay check reads before it writes, so it cannot see a sibling's +/// uncommitted row. Holding the first transaction open makes the race happen +/// every run instead of by timing. On PostgreSQL the second insert waits on the +/// first's row and then collides with it; on SQLite it loses the write lock. +/// Either way the replay must come back as a duplicate, not an error. +async fn exercise_concurrent_replay_is_a_duplicate(db: &DatabaseConnection) { + let user = persist_user(db).await; + let book = persist_book(db).await; + let id = Uuid::new_v4(); + let now = Utc::now(); + let session = move || { + NewSession::from_client( + id, + user, + book, + DEVICE, + None, + SessionKind::Progress, + Some(60_000), + Some(1), + now - Duration::minutes(2), + now - Duration::minutes(1), + ) + .with_page(5) + }; + + let first = db.begin().await.unwrap(); + let outcome = ReadingSessionRepository::append(&first, session(), Utc::now()) + .await + .unwrap(); + assert_eq!(outcome, AppendOutcome::Inserted); + + let replay_db = db.clone(); + let replay = + tokio::spawn( + async move { ReadProgressRepository::record_session(&replay_db, session()).await }, + ); + + // Long enough for the replay to have read and reached its insert. + tokio::time::sleep(std::time::Duration::from_millis(300)).await; + first.commit().await.unwrap(); + + let (outcome, _) = replay + .await + .unwrap() + .expect("a concurrent replay must not fail"); + assert_eq!(outcome, AppendOutcome::Duplicate); + assert_eq!(sessions(db, user, book).await.len(), 1); +} + +#[tokio::test] +async fn concurrent_replay_is_a_duplicate_sqlite() { + let (db, _t) = setup_test_db().await; + exercise_concurrent_replay_is_a_duplicate(&db).await; +} + /// All of the above against PostgreSQL, sequenced in one test on purpose: /// `setup_test_db_postgres` truncates a database shared by the whole run, so /// two PostgreSQL tests running at once delete each other's fixtures. @@ -281,4 +341,5 @@ async fn reading_sessions_postgres() { exercise_other_device_is_untouched(&db).await; exercise_stale_position_save_is_not_absorbed(&db).await; exercise_compat_client_unaffected(&db).await; + exercise_concurrent_replay_is_a_duplicate(&db).await; }