use super::{
catalogue::enumerate_remote_catalogue,
design::{SyncedDesign, sync_one_design},
};
use crate::bridge::library::mirror::options::{
LibraryTransport, MirrorCounts, MirrorOptions, MirrorOutcome, MirrorProgress,
};
use indicatrix_net::library::DesignSummary;
use indicatrix_vault::{db::sqlite::Database, model::mirror::MirrorState};
use std::sync::{
Arc, Mutex, PoisonError,
atomic::{AtomicBool, Ordering},
};
use tracing::warn;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(super) enum Decision {
SkipDeleted,
SkipUnchanged,
Sync,
}
pub(super) fn remote_design_version(summary: &DesignSummary) -> Option<[u8; 32]> {
(summary.design_version != [0u8; 32]).then_some(summary.design_version)
}
pub(super) fn decide(
state: Option<&MirrorState>,
summary: &DesignSummary,
remote_design: Option<[u8; 32]>,
) -> Decision {
let Some(state) = state else {
return Decision::Sync;
};
if state.deleted_locally {
return Decision::SkipDeleted;
}
let summary_moved = state.summary_version != summary.version;
let design_moved = remote_design.is_none_or(|version| version != state.design_version);
if summary_moved || design_moved {
Decision::Sync
} else {
Decision::SkipUnchanged
}
}
fn load_mirror_state(
db: &Arc<Mutex<Database>>,
summary: &DesignSummary,
) -> Result<Option<MirrorState>, ()> {
db.lock()
.unwrap_or_else(PoisonError::into_inner)
.get_mirror_state(&summary.url)
.map_err(|e| {
warn!(
"Mirror sync: cannot read the mirror state of {} ({e}); leaving it for the \
next sync",
summary.url
);
})
}
fn progress_after(
processed: usize,
counts: MirrorCounts,
summary: &DesignSummary,
) -> MirrorProgress {
MirrorProgress {
processed: processed + 1,
counts,
current_title: summary.title.clone(),
}
}
#[must_use]
pub(super) fn run_mirror_sync(
db: &Arc<Mutex<Database>>,
transport: &impl LibraryTransport,
source_id: &str,
options: MirrorOptions,
cancel: &AtomicBool,
mut on_progress: impl FnMut(MirrorProgress),
) -> MirrorOutcome {
let summaries = match enumerate_remote_catalogue(transport) {
Ok(list) => list,
Err(outcome) => return outcome,
};
let mut counts = MirrorCounts {
total_found: summaries.len(),
..MirrorCounts::default()
};
for (processed, summary) in summaries.into_iter().enumerate() {
if cancel.load(Ordering::Relaxed) {
counts.orphaned_mirror_states = count_orphaned_mirror_states(db);
return MirrorOutcome::Cancelled(counts);
}
let Ok(existing_state) = load_mirror_state(db, &summary) else {
counts.failed += 1;
on_progress(progress_after(processed, counts, &summary));
continue;
};
match decide(
existing_state.as_ref(),
&summary,
remote_design_version(&summary),
) {
Decision::SkipUnchanged => {
counts.skipped_unchanged += 1;
on_progress(progress_after(processed, counts, &summary));
continue;
}
Decision::SkipDeleted => {
counts.skipped_deleted += 1;
on_progress(progress_after(processed, counts, &summary));
continue;
}
Decision::Sync => {}
}
let is_new = existing_state.is_none();
match sync_one_design(
db,
transport,
source_id,
options,
&summary,
is_new,
&mut counts,
) {
Ok(SyncedDesign::Saved) => {
if is_new {
counts.new_count += 1;
} else {
counts.updated_count += 1;
}
}
Ok(SyncedDesign::SkippedLocalConflict) => counts.local_conflicts_skipped += 1,
Err(()) => counts.failed += 1,
}
on_progress(progress_after(processed, counts, &summary));
}
counts.orphaned_mirror_states = count_orphaned_mirror_states(db);
MirrorOutcome::Completed(counts)
}
fn count_orphaned_mirror_states(db: &Arc<Mutex<Database>>) -> u64 {
db.lock()
.unwrap_or_else(PoisonError::into_inner)
.count_mirror_states_without_entry()
.unwrap_or(0)
}