Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
28 changes: 28 additions & 0 deletions crates/codex-api/src/routes/v1/dto/reading_progress_transfer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ pub use codex_services::reading_transfer::model::*;

use serde::Deserialize;
use utoipa::ToSchema;
use uuid::Uuid;

fn default_true() -> bool {
true
Expand All @@ -24,12 +25,39 @@ pub struct ExportReadingProgressQuery {
/// notice until the numbers are gone.
#[serde(default = "default_true")]
pub include_sessions: bool,

/// Comma-separated library ids to export, e.g.
/// `?libraryIds=<uuid>,<uuid>`. Omitted exports every library the reader
/// has state for.
///
/// Comma-separated rather than a repeated key because axum's `Query`
/// extractor deserializes with `serde_urlencoded`, which collapses a
/// repeated key instead of collecting it into a `Vec`.
pub library_ids: Option<String>,
}

impl ExportReadingProgressQuery {
/// Parse `library_ids`, rejecting anything that is not a uuid rather than
/// silently exporting more than the caller asked for.
pub fn parsed_library_ids(&self) -> Result<Option<Vec<Uuid>>, uuid::Error> {
let Some(raw) = self.library_ids.as_deref() else {
return Ok(None);
};
let ids = raw
.split(',')
.map(str::trim)
.filter(|part| !part.is_empty())
.map(Uuid::parse_str)
.collect::<Result<Vec<_>, _>>()?;
Ok(Some(ids))
}
}

impl Default for ExportReadingProgressQuery {
fn default() -> Self {
Self {
include_sessions: true,
library_ids: None,
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,8 @@ use codex_services::reading_transfer::import::{ImportError, ImportOptions};
get,
path = "/api/v1/reading-progress/export",
params(
("includeSessions" = Option<bool>, Query, description = "Include the reading-session log (default: true). Sessions are the only source of every reading statistic, so this is opt-out rather than opt-in.")
("includeSessions" = Option<bool>, Query, description = "Include the reading-session log (default: true). Sessions are the only source of every reading statistic, so this is opt-out rather than opt-in."),
("libraryIds" = Option<String>, Query, description = "Comma-separated library ids to export. Omitted exports every library the reader has state for.")
),
responses(
(status = 200, description = "The export document", body = codex_services::reading_transfer::model::ReadingProgressExportDocument),
Expand All @@ -60,10 +61,18 @@ pub async fn export_reading_progress(
) -> Result<Response, ApiError> {
auth.require_permission(&Permission::ProgressRead)?;

let document =
codex_services::export_reading_progress(&state.db, auth.user_id, query.include_sessions)
.await
.map_err(|e| ApiError::Internal(format!("Failed to export reading progress: {e}")))?;
let library_ids = query
.parsed_library_ids()
.map_err(|e| ApiError::BadRequest(format!("Invalid libraryIds: {e}")))?;

let document = codex_services::export_reading_progress(
&state.db,
auth.user_id,
query.include_sessions,
library_ids.as_deref(),
)
.await
.map_err(|e| ApiError::Internal(format!("Failed to export reading progress: {e}")))?;

let filename = format!(
"codex-reading-progress-{}.json",
Expand Down Expand Up @@ -130,6 +139,7 @@ pub async fn import_reading_progress(
conflict_policy: request.conflict_policy,
reattach_sessions: request.reattach_sessions,
accept_stem_matches: request.accept_stem_matches,
library_ids: request.library_ids,
};

let response =
Expand Down
36 changes: 30 additions & 6 deletions crates/codex-services/src/reading_transfer/export.rs
Original file line number Diff line number Diff line change
Expand Up @@ -95,10 +95,15 @@ async fn sessions_for_user(
/// `include_sessions = false` omits the `sessions` key entirely on every book
/// rather than emitting empty arrays, so a client can tell "not exported"
/// apart from "exported, and there were none".
/// `library_ids` of `None` exports everything the reader has state for.
/// `Some` narrows to those libraries, which is what a split wants: carrying
/// an entire reading history when only one library is being reorganised makes
/// the file larger and the import's matching job harder for no gain.
pub async fn export_reading_progress(
db: &DatabaseConnection,
user_id: Uuid,
include_sessions: bool,
library_ids: Option<&[Uuid]>,
) -> Result<ReadingProgressExportDocument> {
let progress_rows = ReadProgressRepository::get_by_user(db, user_id).await?;
let completion_rows = completions_for_user(db, user_id).await?;
Expand Down Expand Up @@ -133,8 +138,15 @@ pub async fn export_reading_progress(
let series_ids: Vec<Uuid> = series_id_set.into_iter().collect();

let mut series_rows = SeriesRepository::get_by_ids(db, &series_ids).await?;
if let Some(wanted) = library_ids {
series_rows.retain(|s| wanted.contains(&s.library_id));
}
series_rows.sort_by(|a, b| a.path.cmp(&b.path).then_with(|| a.id.cmp(&b.id)));

// Re-derived after the filter so nothing downstream loads or emits a
// series the caller excluded.
let series_ids: Vec<Uuid> = series_rows.iter().map(|s| s.id).collect();

let library_ids: Vec<Uuid> = series_rows
.iter()
.map(|s| s.library_id)
Expand Down Expand Up @@ -381,7 +393,9 @@ mod tests {
.await
.unwrap();

let doc = export_reading_progress(conn, user, true).await.unwrap();
let doc = export_reading_progress(conn, user, true, None)
.await
.unwrap();

assert_eq!(doc.format, READING_PROGRESS_FORMAT);
assert_eq!(doc.series.len(), 1);
Expand Down Expand Up @@ -449,7 +463,9 @@ mod tests {
missing.deleted = Set(true);
missing.update(conn).await.unwrap();

let doc = export_reading_progress(conn, user, true).await.unwrap();
let doc = export_reading_progress(conn, user, true, None)
.await
.unwrap();

assert_eq!(
doc.series.len(),
Expand Down Expand Up @@ -520,15 +536,19 @@ mod tests {
.await
.unwrap();

let doc_without = export_reading_progress(conn, user, false).await.unwrap();
let doc_without = export_reading_progress(conn, user, false, None)
.await
.unwrap();
assert!(!doc_without.includes_sessions);
let book_doc = &doc_without.series[0].books[0];
assert!(book_doc.sessions.is_none());
// Completions and progress are unaffected by include_sessions.
assert_eq!(book_doc.completions.len(), 1);
assert!(book_doc.progress.is_some());

let doc_with = export_reading_progress(conn, user, true).await.unwrap();
let doc_with = export_reading_progress(conn, user, true, None)
.await
.unwrap();
assert!(doc_with.includes_sessions);
let book_doc = &doc_with.series[0].books[0];
assert!(book_doc.sessions.is_some());
Expand Down Expand Up @@ -684,7 +704,9 @@ mod tests {
.unwrap();
}

let doc = export_reading_progress(conn, user, true).await.unwrap();
let doc = export_reading_progress(conn, user, true, None)
.await
.unwrap();
assert_eq!(doc.series.len(), SERIES_COUNT);

let json = serde_json::to_vec(&doc).unwrap();
Expand Down Expand Up @@ -746,7 +768,9 @@ mod tests {
.await
.unwrap();

let doc = export_reading_progress(conn, user_a, true).await.unwrap();
let doc = export_reading_progress(conn, user_a, true, None)
.await
.unwrap();
assert_eq!(doc.series.len(), 1);
assert_eq!(doc.series[0].books.len(), 1);
assert_eq!(
Expand Down
72 changes: 67 additions & 5 deletions crates/codex-services/src/reading_transfer/import.rs
Original file line number Diff line number Diff line change
Expand Up @@ -44,10 +44,16 @@ use super::model::{
pub struct ImportOptions {
pub dry_run: bool,
pub hash_mode: HashMode,
pub source_preference: Vec<String>,
/// `None` means every source the exported series carries. See
/// [`ImportReadingProgressRequest::source_preference`].
pub source_preference: Option<Vec<String>>,
pub conflict_policy: ConflictPolicy,
pub reattach_sessions: bool,
pub accept_stem_matches: bool,
/// Which libraries a series may match into. `None` searches every library
/// the reader can see; `Some` narrows the search, which is how an import
/// can run while the old copy of a series is still present.
pub library_ids: Option<Vec<Uuid>>,
}

/// A rejection worth a 400, versus every other failure which is a 500.
Expand Down Expand Up @@ -304,6 +310,11 @@ enum RowDecision {
/// the file, or belongs to someone else (a UUID collision that should
/// never happen, handled by leaving it alone).
Skip,
/// On a *different* live book, in a library that still has the file. The
/// write is the same no-op as [`Self::Skip`], but it means something very
/// different: the progress moved and the reading history did not. Counted
/// apart so the report can say so.
SkipStranded,
}

/// An existing history row, as much as a decision needs.
Expand Down Expand Up @@ -440,12 +451,15 @@ fn decide_row(
match existing {
None => RowDecision::Insert,
Some(row) if row.user_id == user_id => {
let on_another_live_book = matches!(row.book_id, Some(book) if book != target_book && !soft_deleted.contains(&book));
let off_a_live_book = match row.book_id {
None => true,
Some(book) => book != target_book && soft_deleted.contains(&book),
};
if off_a_live_book && reattach {
RowDecision::Reattach
} else if on_another_live_book {
RowDecision::SkipStranded
} else {
RowDecision::Skip
}
Expand Down Expand Up @@ -591,6 +605,10 @@ fn book_report_from_plan(plan: &BookPlan<'_>) -> ImportBookReport {
RowDecision::Insert => completions.inserted += 1,
RowDecision::Reattach => completions.reattached += 1,
RowDecision::Skip => completions.skipped += 1,
RowDecision::SkipStranded => {
completions.skipped += 1;
completions.stranded += 1;
}
}
}
let mut sessions = WriteCounts::default();
Expand All @@ -599,6 +617,10 @@ fn book_report_from_plan(plan: &BookPlan<'_>) -> ImportBookReport {
RowDecision::Insert => sessions.inserted += 1,
RowDecision::Reattach => sessions.reattached += 1,
RowDecision::Skip => sessions.skipped += 1,
RowDecision::SkipStranded => {
sessions.skipped += 1;
sessions.stranded += 1;
}
}
}

Expand Down Expand Up @@ -651,6 +673,7 @@ fn tally_book_report(summary: &mut ImportSummary, report: &ImportBookReport, cou
summary.completions_reattached += report.completions.reattached;
summary.sessions_inserted += report.sessions.inserted;
summary.sessions_reattached += report.sessions.reattached;
summary.rows_stranded += report.completions.stranded + report.sessions.stranded;
}

// ---------------------------------------------------------------------------
Expand Down Expand Up @@ -710,7 +733,7 @@ async fn apply_completion(
decision: &RowDecision,
) -> Result<()> {
match decision {
RowDecision::Skip => Ok(()),
RowDecision::Skip | RowDecision::SkipStranded => Ok(()),
RowDecision::Insert => {
read_completions::ActiveModel {
id: Set(doc.id),
Expand Down Expand Up @@ -743,7 +766,7 @@ async fn apply_session(
decision: &RowDecision,
) -> Result<()> {
match decision {
RowDecision::Skip => Ok(()),
RowDecision::Skip | RowDecision::SkipStranded => Ok(()),
RowDecision::Insert => {
reading_sessions::ActiveModel {
id: Set(doc.id),
Expand Down Expand Up @@ -1042,7 +1065,8 @@ pub async fn import_reading_progress(
db,
&content_filter,
series_doc,
&options.source_preference,
options.source_preference.as_deref(),
options.library_ids.as_deref(),
)
.await
{
Expand Down Expand Up @@ -1160,6 +1184,20 @@ pub async fn import_reading_progress(
}
}

// A stranded row means the progress moved and the reading history did
// not, which the per-book counts show but the summary otherwise reads as
// success. Say it plainly: the reader can still fix it by removing or
// rescanning the other library and importing again.
if summary.rows_stranded > 0 {
notices.push(format!(
"{} session/completion rows were left behind: their books are still \
live in another library. Progress moved but reading history did not. \
Delete or rescan that library so its books are no longer on disk, \
then import again to bring the history across.",
summary.rows_stranded
));
}

Ok(ImportReadingProgressResponse {
dry_run: options.dry_run,
sessions_in_file: document.includes_sessions,
Expand Down Expand Up @@ -1481,8 +1519,13 @@ mod tests {
assert_eq!(decision, RowDecision::Reattach);
}

/// Left alone, but not for the harmless reason a plain `Skip` means.
/// The row is on a live book in another library: reattaching would strip
/// history from a library the reader may still be using, and skipping it
/// silently would hide that the progress moved without it. Hence its own
/// variant, which the report turns into a notice.
#[test]
fn row_on_a_live_book_is_left_alone() {
fn row_on_a_live_book_elsewhere_is_left_alone_but_flagged() {
let user = Uuid::new_v4();
let live = Uuid::new_v4();
let mut planned = PlannedState::default();
Expand All @@ -1495,6 +1538,25 @@ mod tests {
true,
&mut planned,
);
assert_eq!(decision, RowDecision::SkipStranded);
}

/// The genuinely harmless skip: the row is already on the book being
/// imported onto, which is what makes re-importing the same file a no-op.
#[test]
fn row_already_on_the_target_book_is_a_plain_skip() {
let user = Uuid::new_v4();
let target = Uuid::new_v4();
let mut planned = PlannedState::default();
let decision = decide_row(
Uuid::new_v4(),
Some(&row(user, Some(target))),
target,
user,
&HashSet::new(),
true,
&mut planned,
);
assert_eq!(decision, RowDecision::Skip);
}

Expand Down
Loading
Loading