use super::options::{
LibraryTransport, MirrorCounts, MirrorHandle, MirrorOptions, MirrorOutcome, MirrorProgress,
mirror_source_id,
};
use crate::{bridge::library::client::LibrarySession, settings::WorkerSettings};
use indicatrix_net::library::{AngleSettingWire, LibraryRequest, LibraryResponse, RangeFilterWire};
use indicatrix_vault::{
db::sqlite::Database,
model::{
angle::AngleSetting, detail::FacetDiagramDetail, entry::FacetDiagramEntry,
file::AttachedFile, mirror::MirrorState,
},
};
use slint::{ComponentHandle, Weak};
use std::{
sync::{
Arc, Mutex, PoisonError,
atomic::{AtomicBool, Ordering},
},
thread,
};
pub fn spawn_mirror_sync<T, P, D>(
ui_weak: Weak<T>,
db: Arc<Mutex<Database>>,
worker: WorkerSettings,
options: MirrorOptions,
on_progress: P,
on_done: D,
) -> MirrorHandle
where
T: ComponentHandle + 'static,
P: Fn(&T, MirrorProgress) + Send + 'static + Clone,
D: Fn(&T, MirrorOutcome) + Send + 'static,
{
let cancel = Arc::new(AtomicBool::new(false));
let cancel_worker = cancel.clone();
let ui_weak_done = ui_weak.clone();
thread::spawn(move || {
let source_id = mirror_source_id(&worker);
let session = LibrarySession::new(worker);
let progress_ui_weak = ui_weak;
let outcome = run_mirror_sync(
&db,
&session,
&source_id,
options,
&cancel_worker,
move |progress| {
let on_progress = on_progress.clone();
let _ =
progress_ui_weak.upgrade_in_event_loop(move |ui| on_progress(&ui, progress));
},
);
let _ = ui_weak_done.upgrade_in_event_loop(move |ui| on_done(&ui, outcome));
});
MirrorHandle { cancel }
}
fn enumerate_remote_catalogue(
transport: &impl LibraryTransport,
) -> Result<Vec<indicatrix_net::library::DesignSummary>, MirrorOutcome> {
let mut summaries = Vec::new();
let mut cursor = None;
loop {
let request = LibraryRequest::SearchPage {
query: String::new(),
shape_filter: "All".to_string(),
gear_filter: "All".to_string(),
range: RangeFilterWire::default(),
cursor,
};
match transport.request(&request) {
Ok(LibraryResponse::SearchResultsPage {
results,
next_cursor,
excluded_for_missing_curves: _,
}) => {
summaries.extend(results);
match next_cursor {
Some(c) => cursor = Some(c),
None => return Ok(summaries),
}
}
Ok(other) => {
return Err(MirrorOutcome::Failed(format!(
"remote worker replied to SearchPage with an unexpected message: {other:?}"
)));
}
Err(e) => {
return Err(MirrorOutcome::Failed(format!(
"could not list the remote library: {e}"
)));
}
}
}
}
#[must_use]
pub 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) {
return MirrorOutcome::Cancelled(counts);
}
let existing_state = {
let db = db.lock().unwrap_or_else(PoisonError::into_inner);
db.get_mirror_state(&summary.url).ok().flatten()
};
let unchanged = existing_state
.as_ref()
.is_some_and(|state| state.summary_version == summary.version);
if unchanged {
counts.skipped_unchanged += 1;
on_progress(MirrorProgress {
processed: processed + 1,
counts,
current_title: summary.title.clone(),
});
continue;
}
let is_new = existing_state.is_none();
match sync_one_design(db, transport, source_id, options, &summary, &mut counts) {
Ok(()) => {
if is_new {
counts.new_count += 1;
} else {
counts.updated_count += 1;
}
}
Err(()) => counts.failed += 1,
}
on_progress(MirrorProgress {
processed: processed + 1,
counts,
current_title: summary.title.clone(),
});
}
MirrorOutcome::Completed(counts)
}
fn sync_one_design(
db: &Arc<Mutex<Database>>,
transport: &impl LibraryTransport,
source_id: &str,
options: MirrorOptions,
summary: &indicatrix_net::library::DesignSummary,
counts: &mut MirrorCounts,
) -> Result<(), ()> {
let design = match transport.request(&LibraryRequest::FetchDesign {
entry_id: summary.entry_id,
}) {
Ok(LibraryResponse::Design(d)) => *d,
_ => return Err(()),
};
let mut attached_files = Vec::with_capacity(design.attachments.len());
for meta in &design.attachments {
if meta.size > options.max_attachment_bytes {
counts.attachments_skipped_too_large += 1;
continue;
}
match transport.request(&LibraryRequest::FetchAttachment {
attachment_id: meta.id,
}) {
Ok(LibraryResponse::Attachment { name, content }) => {
counts.attachment_bytes_fetched += content.len() as u64;
counts.attachments_fetched += 1;
attached_files.push(AttachedFile {
name,
url: meta.url.clone(),
content,
});
}
Ok(LibraryResponse::NotFound) => {
}
_ => return Err(()),
}
}
let entry = FacetDiagramEntry {
title: design.title.clone(),
url: design.url.clone(),
design_id: design.design_id.clone().unwrap_or_default(),
};
let detail = FacetDiagramDetail {
page_url: design.page_url.clone(),
diagram_image_name: design.diagram_image_name.clone(),
diagram_image_data: design.diagram_image_data.clone(),
angle_settings_table: design.angle_settings.iter().map(to_angle_setting).collect(),
attached_files,
competition_diagram: design.competition_diagram.clone(),
lw_ratio: design.lw_ratio.clone(),
refractive_index: design.refractive_index.clone(),
index_gear: design.index_gear.clone(),
volume: design.volume.clone(),
facets_count: design.facets_count.clone(),
shape: design.shape.clone(),
designer_info: design.designer_info.clone(),
..FacetDiagramDetail::default()
};
let local_entry_id = {
let db = db.lock().unwrap_or_else(PoisonError::into_inner);
let Ok(id) = db.save_diagram_entry(&entry, source_id) else {
return Err(());
};
if db.save_diagram_detail(&detail, id).is_err() {
return Err(());
}
let _ = db.upsert_mirror_state(&MirrorState {
url: design.url.clone(),
source_id: source_id.to_string(),
summary_version: summary.version,
design_version: design.version,
});
id
};
let _ = local_entry_id;
Ok(())
}
fn to_angle_setting(a: &AngleSettingWire) -> AngleSetting {
AngleSetting {
order_index: a.order_index,
facet: a.facet.clone(),
angle: a.angle.clone(),
index: a.index.clone(),
notes: a.notes.clone(),
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::bridge::library::client as library_client;
use indicatrix_net::library::{AttachedFileMeta, DesignRecord, DesignSummary};
use std::{cell::RefCell, collections::HashMap, sync::atomic::AtomicUsize};
struct FakeTransport {
search_result: Vec<DesignSummary>,
page_size: usize,
designs: HashMap<i64, DesignRecord>,
attachments: HashMap<i64, (String, Vec<u8>)>,
search_calls: AtomicUsize,
fetch_design_calls: AtomicUsize,
fetch_attachment_calls: AtomicUsize,
fetched_entry_ids: RefCell<Vec<i64>>,
fail_fetch_design_for: Option<i64>,
}
impl FakeTransport {
fn new(summaries: Vec<DesignSummary>, designs: Vec<DesignRecord>) -> Self {
Self {
search_result: summaries,
page_size: usize::MAX,
designs: designs.into_iter().map(|d| (d.entry_id, d)).collect(),
attachments: HashMap::new(),
search_calls: AtomicUsize::new(0),
fetch_design_calls: AtomicUsize::new(0),
fetch_attachment_calls: AtomicUsize::new(0),
fetched_entry_ids: RefCell::new(Vec::new()),
fail_fetch_design_for: None,
}
}
fn with_attachment(mut self, id: i64, name: &str, content: Vec<u8>) -> Self {
self.attachments.insert(id, (name.to_string(), content));
self
}
fn with_page_size(mut self, page_size: usize) -> Self {
self.page_size = page_size;
self
}
fn with_fetch_design_failure(mut self, entry_id: i64) -> Self {
self.fail_fetch_design_for = Some(entry_id);
self
}
}
impl LibraryTransport for FakeTransport {
fn request(
&self,
req: &LibraryRequest,
) -> Result<LibraryResponse, library_client::LibraryClientError> {
match req {
LibraryRequest::Search { .. } => {
unreachable!(
"mirror sync always pages via SearchPage now, never sends Search directly"
)
}
LibraryRequest::SearchPage { cursor, .. } => {
self.search_calls.fetch_add(1, Ordering::Relaxed);
let start = cursor.map_or(0, |after| {
self.search_result
.iter()
.position(|s| s.entry_id == after)
.map_or(self.search_result.len(), |i| i + 1)
});
let end = start
.saturating_add(self.page_size)
.min(self.search_result.len());
let page = self.search_result[start..end].to_vec();
let next_cursor = if page.len() == self.page_size {
page.last().map(|s| s.entry_id)
} else {
None
};
Ok(LibraryResponse::SearchResultsPage {
results: page,
next_cursor,
excluded_for_missing_curves: 0,
})
}
LibraryRequest::FetchDesign { entry_id } => {
self.fetch_design_calls.fetch_add(1, Ordering::Relaxed);
self.fetched_entry_ids.borrow_mut().push(*entry_id);
if self.fail_fetch_design_for == Some(*entry_id) {
return Err(library_client::LibraryClientError::Client(
indicatrix_net::client::ClientError::Net(
indicatrix_net::messages::NetError::Framing(
indicatrix_net::framing::FramingError::Io(std::io::Error::new(
std::io::ErrorKind::ConnectionReset,
"connection dropped and could not be re-established",
)),
),
),
));
}
Ok(self
.designs
.get(entry_id)
.map_or(LibraryResponse::NotFound, |d| {
LibraryResponse::Design(Box::new(d.clone()))
}))
}
LibraryRequest::FetchAttachment { attachment_id } => {
self.fetch_attachment_calls.fetch_add(1, Ordering::Relaxed);
match self.attachments.get(attachment_id) {
Some((name, content)) => Ok(LibraryResponse::Attachment {
name: name.clone(),
content: content.clone(),
}),
None => Ok(LibraryResponse::NotFound),
}
}
LibraryRequest::FilterOptions => unreachable!("mirror sync never sends this"),
LibraryRequest::FetchDesignSource { .. } => {
unreachable!(
"mirror sync never sends this -- only the editor's remote 'Load Selected' does"
)
}
}
}
}
fn temp_db() -> (Arc<Mutex<Database>>, std::path::PathBuf) {
static COUNTER: AtomicUsize = AtomicUsize::new(0);
let n = COUNTER.fetch_add(1, Ordering::Relaxed);
let path = std::env::temp_dir().join(format!(
"indicatrix-cut-library-mirror-test-{}-{n}.sqlite",
std::process::id()
));
let _ = std::fs::remove_file(&path);
let db = Database::new(Some(path.to_str().unwrap())).unwrap();
(Arc::new(Mutex::new(db)), path)
}
fn design_record(entry_id: i64, title: &str, url: &str, extra_field: &str) -> DesignRecord {
let mut record = DesignRecord {
entry_id,
title: title.to_string(),
url: url.to_string(),
design_id: Some(format!("D-{entry_id}")),
page_url: url.to_string(),
diagram_image_name: None,
diagram_image_data: None,
competition_diagram: None,
lw_ratio: Some(extra_field.to_string()),
refractive_index: Some("1.72".to_string()),
index_gear: Some("96".to_string()),
volume: Some("0.65".to_string()),
facets_count: Some("57".to_string()),
shape: Some("Round".to_string()),
designer_info: Some("Capps, Jerry".to_string()),
preview_material: None,
hw_ratio: None,
tw_ratio: None,
uw_ratio: None,
pw_ratio: None,
cw_ratio: None,
symmetry_order: None,
mirror_symmetry: None,
designer: None,
angle_settings: vec![AngleSettingWire {
order_index: 0,
facet: "P1".to_string(),
angle: "41.0".to_string(),
index: "96".to_string(),
notes: String::new(),
}],
attachments: Vec::new(),
version: [0u8; 32],
};
record.version = content_hash(&[title.as_bytes(), url.as_bytes(), extra_field.as_bytes()]);
record
}
fn design_summary(entry_id: i64, title: &str, url: &str, extra_field: &str) -> DesignSummary {
DesignSummary {
entry_id,
title: title.to_string(),
url: url.to_string(),
design_id: Some(format!("D-{entry_id}")),
shape: Some("Round".to_string()),
index_gear: Some("96".to_string()),
facets_count: Some("57".to_string()),
designer_info: Some("Capps, Jerry".to_string()),
lw_ratio: Some(extra_field.to_string()),
refractive_index: Some("1.72".to_string()),
volume: Some("0.65".to_string()),
competition_diagram: None,
ignored: false,
version: content_hash(&[title.as_bytes(), url.as_bytes(), extra_field.as_bytes()]),
}
}
fn content_hash(parts: &[&[u8]]) -> [u8; 32] {
use sha2::{Digest, Sha256};
let mut hasher = Sha256::new();
for p in parts {
hasher.update(p);
}
hasher.finalize().into()
}
#[test]
fn a_local_only_design_survives_a_mirror_sync() {
let (db, path) = temp_db();
let local_entry_id = {
let guard = db.lock().unwrap();
guard
.save_diagram_entry(
&FacetDiagramEntry {
title: "My Own Trichecker".to_string(),
url: "local://my_trichecker.asc".to_string(),
design_id: String::new(),
},
indicatrix_vault::local::LOCAL_SOURCE_ID,
)
.unwrap()
};
{
let guard = db.lock().unwrap();
guard
.save_diagram_detail(
&FacetDiagramDetail {
shape: Some("Trichecker".to_string()),
attached_files: vec![AttachedFile {
name: "my_trichecker.asc".to_string(),
url: String::new(),
content: b"a real user file".to_vec(),
}],
..FacetDiagramDetail::default()
},
local_entry_id,
)
.unwrap();
}
let remote_summary = design_summary(1, "Round Brilliant", "https://example.test/1", "v1");
let remote_design = design_record(1, "Round Brilliant", "https://example.test/1", "v1");
let transport = FakeTransport::new(vec![remote_summary], vec![remote_design]);
let outcome = run_mirror_sync(
&db,
&transport,
"remote-library:example.test:9443",
MirrorOptions::default(),
&AtomicBool::new(false),
|_| {},
);
assert!(matches!(outcome, MirrorOutcome::Completed(c) if c.new_count == 1));
let guard = db.lock().unwrap();
let local_full = guard.get_diagram_full(local_entry_id).unwrap().unwrap();
assert_eq!(local_full.title, "My Own Trichecker");
assert_eq!(local_full.url, "local://my_trichecker.asc");
assert_eq!(local_full.attached_files.len(), 1);
assert_eq!(local_full.attached_files[0].content, b"a real user file");
assert_eq!(guard.get_total_count().unwrap(), 2);
drop(guard);
std::fs::remove_file(&path).ok();
}
#[test]
fn a_new_design_is_saved_with_its_attachment() {
let (db, path) = temp_db();
let mut record = design_record(1, "Round Brilliant", "https://example.test/1", "v1");
record.attachments = vec![AttachedFileMeta {
id: 100,
name: "schedule.pdf".to_string(),
url: "https://example.test/schedule.pdf".to_string(),
size: 5,
}];
let summary = design_summary(1, "Round Brilliant", "https://example.test/1", "v1");
let transport = FakeTransport::new(vec![summary], vec![record]).with_attachment(
100,
"schedule.pdf",
vec![1, 2, 3, 4, 5],
);
let outcome = run_mirror_sync(
&db,
&transport,
"remote-library:w",
MirrorOptions::default(),
&AtomicBool::new(false),
|_| {},
);
let MirrorOutcome::Completed(counts) = outcome else {
panic!("expected Completed, got {outcome:?}");
};
assert_eq!(counts.new_count, 1);
assert_eq!(counts.attachments_fetched, 1);
assert_eq!(counts.attachment_bytes_fetched, 5);
let guard = db.lock().unwrap();
let full = guard
.search_diagrams(
"",
"All",
"All",
&indicatrix_vault::model::filter::RangeFilter::default(),
)
.unwrap();
assert_eq!(full.len(), 1);
let entry_id = full[0].id;
let record = guard.get_diagram_full(entry_id).unwrap().unwrap();
assert_eq!(record.attached_files.len(), 1);
assert_eq!(record.attached_files[0].content, vec![1, 2, 3, 4, 5]);
drop(guard);
std::fs::remove_file(&path).ok();
}
#[test]
fn a_second_sync_of_an_unchanged_catalogue_skips_every_design_and_makes_no_fetch_design_calls()
{
let (db, path) = temp_db();
let summary = design_summary(1, "Round Brilliant", "https://example.test/1", "v1");
let record = design_record(1, "Round Brilliant", "https://example.test/1", "v1");
let transport = FakeTransport::new(vec![summary], vec![record]);
let first = run_mirror_sync(
&db,
&transport,
"remote-library:w",
MirrorOptions::default(),
&AtomicBool::new(false),
|_| {},
);
assert!(matches!(first, MirrorOutcome::Completed(c) if c.new_count == 1));
assert_eq!(transport.fetch_design_calls.load(Ordering::Relaxed), 1);
let second = run_mirror_sync(
&db,
&transport,
"remote-library:w",
MirrorOptions::default(),
&AtomicBool::new(false),
|_| {},
);
let MirrorOutcome::Completed(counts) = second else {
panic!("expected Completed, got {second:?}");
};
assert_eq!(counts.skipped_unchanged, 1);
assert_eq!(counts.new_count, 0);
assert_eq!(counts.updated_count, 0);
assert_eq!(transport.search_calls.load(Ordering::Relaxed), 2);
assert_eq!(transport.fetch_design_calls.load(Ordering::Relaxed), 1);
assert_eq!(transport.fetch_attachment_calls.load(Ordering::Relaxed), 0);
std::fs::remove_file(&path).ok();
}
#[test]
fn a_changed_design_is_refetched_and_updated_in_place() {
let (db, path) = temp_db();
let summary_v1 = design_summary(1, "Round Brilliant", "https://example.test/1", "v1");
let record_v1 = design_record(1, "Round Brilliant", "https://example.test/1", "v1");
let transport_v1 = FakeTransport::new(vec![summary_v1], vec![record_v1]);
let _ = run_mirror_sync(
&db,
&transport_v1,
"remote-library:w",
MirrorOptions::default(),
&AtomicBool::new(false),
|_| {},
);
let summary_v2 =
design_summary(1, "Round Brilliant", "https://example.test/1", "v2-updated");
let record_v2 = design_record(1, "Round Brilliant", "https://example.test/1", "v2-updated");
let transport_v2 = FakeTransport::new(vec![summary_v2], vec![record_v2]);
let outcome = run_mirror_sync(
&db,
&transport_v2,
"remote-library:w",
MirrorOptions::default(),
&AtomicBool::new(false),
|_| {},
);
let MirrorOutcome::Completed(counts) = outcome else {
panic!("expected Completed, got {outcome:?}");
};
assert_eq!(counts.updated_count, 1);
assert_eq!(counts.new_count, 0);
let guard = db.lock().unwrap();
let items = guard
.search_diagrams(
"",
"All",
"All",
&indicatrix_vault::model::filter::RangeFilter::default(),
)
.unwrap();
assert_eq!(
items.len(),
1,
"must update the existing row, not add a second one"
);
assert_eq!(items[0].lw_ratio.as_deref(), Some("v2-updated"));
drop(guard);
std::fs::remove_file(&path).ok();
}
#[test]
fn cancelling_mid_sync_leaves_earlier_designs_committed_and_the_rest_untouched() {
let (db, path) = temp_db();
let summaries = vec![
design_summary(1, "First", "https://example.test/1", "v1"),
design_summary(2, "Second", "https://example.test/2", "v1"),
design_summary(3, "Third", "https://example.test/3", "v1"),
];
let designs = vec![
design_record(1, "First", "https://example.test/1", "v1"),
design_record(2, "Second", "https://example.test/2", "v1"),
design_record(3, "Third", "https://example.test/3", "v1"),
];
let transport = FakeTransport::new(summaries, designs);
let cancel = AtomicBool::new(false);
let mut processed_count = 0;
let outcome = run_mirror_sync(
&db,
&transport,
"remote-library:w",
MirrorOptions::default(),
&cancel,
|progress| {
processed_count = progress.processed;
if progress.processed == 1 {
cancel.store(true, Ordering::Relaxed);
}
},
);
assert_eq!(processed_count, 1);
let MirrorOutcome::Cancelled(counts) = outcome else {
panic!("expected Cancelled, got {outcome:?}");
};
assert_eq!(counts.new_count, 1);
let guard = db.lock().unwrap();
assert_eq!(guard.get_total_count().unwrap(), 1);
assert_eq!(transport.fetch_design_calls.load(Ordering::Relaxed), 1);
assert_eq!(*transport.fetched_entry_ids.borrow(), vec![1]);
drop(guard);
std::fs::remove_file(&path).ok();
}
#[test]
fn an_oversized_attachment_is_skipped_but_the_design_is_still_saved() {
let (db, path) = temp_db();
let mut record = design_record(1, "Round Brilliant", "https://example.test/1", "v1");
record.attachments = vec![
AttachedFileMeta {
id: 100,
name: "huge.pdf".to_string(),
url: String::new(),
size: 1000,
},
AttachedFileMeta {
id: 101,
name: "small.pdf".to_string(),
url: String::new(),
size: 5,
},
];
let summary = design_summary(1, "Round Brilliant", "https://example.test/1", "v1");
let transport = FakeTransport::new(vec![summary], vec![record])
.with_attachment(100, "huge.pdf", vec![0u8; 1000])
.with_attachment(101, "small.pdf", vec![9, 9, 9, 9, 9]);
let options = MirrorOptions {
max_attachment_bytes: 500,
};
let outcome = run_mirror_sync(
&db,
&transport,
"remote-library:w",
options,
&AtomicBool::new(false),
|_| {},
);
let MirrorOutcome::Completed(counts) = outcome else {
panic!("expected Completed, got {outcome:?}");
};
assert_eq!(counts.new_count, 1, "the design itself must still be saved");
assert_eq!(counts.attachments_fetched, 1);
assert_eq!(counts.attachments_skipped_too_large, 1);
assert_eq!(
transport.fetch_attachment_calls.load(Ordering::Relaxed),
1,
"the oversized attachment's bytes must never even be requested"
);
let guard = db.lock().unwrap();
let items = guard
.search_diagrams(
"",
"All",
"All",
&indicatrix_vault::model::filter::RangeFilter::default(),
)
.unwrap();
let full = guard.get_diagram_full(items[0].id).unwrap().unwrap();
assert_eq!(full.attached_files.len(), 1);
assert_eq!(full.attached_files[0].name, "small.pdf");
drop(guard);
std::fs::remove_file(&path).ok();
}
#[test]
fn a_failed_fetch_is_counted_and_never_marked_synced() {
let (db, path) = temp_db();
let summary = design_summary(1, "Ghost", "https://example.test/1", "v1");
let transport = FakeTransport::new(vec![summary], Vec::new());
let outcome = run_mirror_sync(
&db,
&transport,
"remote-library:w",
MirrorOptions::default(),
&AtomicBool::new(false),
|_| {},
);
let MirrorOutcome::Completed(counts) = outcome else {
panic!("expected Completed, got {outcome:?}");
};
assert_eq!(counts.failed, 1);
assert_eq!(counts.new_count, 0);
let guard = db.lock().unwrap();
assert_eq!(guard.get_total_count().unwrap(), 0);
assert_eq!(
guard.get_mirror_state("https://example.test/1").unwrap(),
None
);
drop(guard);
std::fs::remove_file(&path).ok();
}
#[test]
fn a_design_whose_connection_could_not_be_recovered_is_skipped_but_the_rest_of_the_sync_completes()
{
let (db, path) = temp_db();
let summaries = vec![
design_summary(1, "First", "https://example.test/1", "v1"),
design_summary(2, "Second", "https://example.test/2", "v1"),
design_summary(3, "Third", "https://example.test/3", "v1"),
];
let designs = vec![
design_record(1, "First", "https://example.test/1", "v1"),
design_record(2, "Second", "https://example.test/2", "v1"),
design_record(3, "Third", "https://example.test/3", "v1"),
];
let transport = FakeTransport::new(summaries, designs).with_fetch_design_failure(2);
let mut processed_count = 0;
let outcome = run_mirror_sync(
&db,
&transport,
"remote-library:w",
MirrorOptions::default(),
&AtomicBool::new(false),
|progress| processed_count = progress.processed,
);
let MirrorOutcome::Completed(counts) = outcome else {
panic!(
"an unrecoverable connection drop for ONE design must not abort the whole \
sync, got {outcome:?}"
);
};
assert_eq!(
processed_count, 3,
"every design was still examined, including after the drop"
);
assert_eq!(counts.total_found, 3);
assert_eq!(
counts.new_count, 2,
"designs 1 and 3 still land despite design 2's failure"
);
assert_eq!(counts.failed, 1);
let guard = db.lock().unwrap();
assert_eq!(guard.get_total_count().unwrap(), 2);
let items = guard
.search_diagrams(
"",
"All",
"All",
&indicatrix_vault::model::filter::RangeFilter::default(),
)
.unwrap();
let titles: std::collections::HashSet<&str> =
items.iter().map(|i| i.title.as_str()).collect();
assert!(titles.contains("First"));
assert!(titles.contains("Third"));
assert!(
!titles.contains("Second"),
"the design whose connection could not be recovered must never be half-saved, \
titles were: {titles:?}"
);
assert_eq!(
guard.get_mirror_state("https://example.test/2").unwrap(),
None,
"never marked synced, so the next sync retries it -- same as any other failed fetch"
);
drop(guard);
std::fs::remove_file(&path).ok();
}
#[test]
fn mirror_source_id_is_distinct_per_worker_address() {
let a = WorkerSettings {
address: "10.0.0.5:9443".to_string(),
..WorkerSettings::default()
};
let b = WorkerSettings {
address: "10.0.0.6:9443".to_string(),
..WorkerSettings::default()
};
assert_ne!(mirror_source_id(&a), mirror_source_id(&b));
}
#[test]
fn a_multi_page_mirror_reaches_designs_beyond_the_first_page() {
let (db, path) = temp_db();
let entries: Vec<(i64, &str)> = vec![
(1, "First"),
(2, "Second"),
(3, "Third"),
(4, "Fourth"),
(5, "Fifth"),
];
let summaries: Vec<DesignSummary> = entries
.iter()
.map(|(id, title)| {
design_summary(*id, title, &format!("https://example.test/{id}"), "v1")
})
.collect();
let designs: Vec<DesignRecord> = entries
.iter()
.map(|(id, title)| {
design_record(*id, title, &format!("https://example.test/{id}"), "v1")
})
.collect();
let transport = FakeTransport::new(summaries, designs).with_page_size(2);
let outcome = run_mirror_sync(
&db,
&transport,
"remote-library:w",
MirrorOptions::default(),
&AtomicBool::new(false),
|_| {},
);
let MirrorOutcome::Completed(counts) = outcome else {
panic!("expected Completed, got {outcome:?}");
};
assert_eq!(
counts.total_found, 5,
"every design across every page must be counted, not just the first page's"
);
assert_eq!(counts.new_count, 5);
assert_eq!(transport.search_calls.load(Ordering::Relaxed), 3);
assert_eq!(transport.fetch_design_calls.load(Ordering::Relaxed), 5);
let guard = db.lock().unwrap();
assert_eq!(guard.get_total_count().unwrap(), 5);
let items = guard
.search_diagrams(
"",
"All",
"All",
&indicatrix_vault::model::filter::RangeFilter::default(),
)
.unwrap();
let titles: std::collections::HashSet<&str> =
items.iter().map(|i| i.title.as_str()).collect();
assert!(titles.contains("Fourth"), "titles were: {titles:?}");
assert!(titles.contains("Fifth"), "titles were: {titles:?}");
drop(guard);
std::fs::remove_file(&path).ok();
}
}