Verify the catalog snapshot before it is sent, and after it lands
Two checks around the upload, both cheap next to what they prevent. Before: the snapshot is quick_checked before it leaves. It is the copy every other device merges from, and a damaged one costs each of them a download, a failed merge and a refusal to push. After: the staged upload's size on the server is compared to the bytes sent before it is rotated into place. A chunked upload is assembled server-side, and an assembly that goes wrong is a file of plausible size no device can open — caught here, on the device that caused it, for one listing; otherwise on every other device, after the fact. A mismatch, or a size the server will not confirm, discards the upload and leaves the current copy and its generations untouched.
This commit is contained in:
@@ -716,11 +716,40 @@ async fn sync_catalog(
|
||||
// this whole scheme exists to survive. Under its own name a failure leaves
|
||||
// the current copy untouched and costs one stray file, retried next pass.
|
||||
let staging = RemotePath::new(format!("{}/{UPLOAD_NAME}", base.as_str()));
|
||||
let sent = bytes.len() as u64;
|
||||
if let Err(e) = backend.put(&staging, bytes, None).await {
|
||||
log::warn!("uploading catalog: {e}");
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
// TRACES: NFR-R2
|
||||
// Confirm the server holds what was sent before it becomes the copy every
|
||||
// other device reads. A chunked upload is assembled server-side, and an
|
||||
// assembly that went wrong is a file of plausible size that no device can
|
||||
// open — the one failure the generations exist to survive, and cheaper
|
||||
// to catch here, on the device that caused it, than on every other one
|
||||
// after. One listing; the size is what the server can vouch for without
|
||||
// reading the file back.
|
||||
match remote_size(backend, base, UPLOAD_NAME).await {
|
||||
Some(held) if held == sent => {}
|
||||
Some(held) => {
|
||||
log::warn!(
|
||||
"not replacing the catalog: sent {sent} bytes but the server holds {held}; \
|
||||
the upload is discarded and retried next pass"
|
||||
);
|
||||
let _ = backend.delete(&RemoteId::Path(staging), None).await;
|
||||
return Ok(());
|
||||
}
|
||||
None => {
|
||||
log::warn!(
|
||||
"not replacing the catalog: the server would not confirm the upload's size; \
|
||||
it is discarded and retried next pass"
|
||||
);
|
||||
let _ = backend.delete(&RemoteId::Path(staging), None).await;
|
||||
return Ok(());
|
||||
}
|
||||
}
|
||||
|
||||
// Rotate, then move the upload into place. Every step here is a rename on
|
||||
// the server, and every destination is empty by the time it is written to
|
||||
// — a `move_to` will not overwrite, by design — so a failure at any point
|
||||
@@ -1223,6 +1252,9 @@ mod derived_guard_tests {
|
||||
bodies: std::collections::HashMap<String, Vec<u8>>,
|
||||
/// Every `move_to`, as (from, to) names, in order.
|
||||
moves: Arc<std::sync::Mutex<Vec<(String, String)>>>,
|
||||
/// What `list` claims the staged upload's size is, when a test wants
|
||||
/// the server to have assembled it wrongly. `None` reports the truth.
|
||||
staged_size: Option<u64>,
|
||||
}
|
||||
|
||||
impl Fussy {
|
||||
@@ -1247,6 +1279,7 @@ mod derived_guard_tests {
|
||||
advertise: None,
|
||||
bodies: Default::default(),
|
||||
moves: Default::default(),
|
||||
staged_size: None,
|
||||
},
|
||||
puts,
|
||||
last_put,
|
||||
@@ -1267,20 +1300,33 @@ mod derived_guard_tests {
|
||||
dir: &RemotePath,
|
||||
_since: Option<&dr_sync::Validator>,
|
||||
) -> Result<Vec<dr_sync::RemoteEntry>, RemoteError> {
|
||||
let Some(size) = self.advertise else {
|
||||
return Ok(Vec::new());
|
||||
let entry = |name: &str, size: u64| {
|
||||
let path = RemotePath::new(format!("{}/{name}", dir.as_str()));
|
||||
dr_sync::RemoteEntry {
|
||||
id: RemoteId::Path(path.clone()),
|
||||
path,
|
||||
kind: dr_sync::EntryKind::File,
|
||||
validator: dr_sync::Validator::new("v"),
|
||||
size,
|
||||
modified: None,
|
||||
has_preview: false,
|
||||
materialised: true,
|
||||
}
|
||||
};
|
||||
let path = RemotePath::new(format!("{}/catalog.sqlite", dir.as_str()));
|
||||
Ok(vec![dr_sync::RemoteEntry {
|
||||
id: RemoteId::Path(path.clone()),
|
||||
path,
|
||||
kind: dr_sync::EntryKind::File,
|
||||
validator: dr_sync::Validator::new("v"),
|
||||
size,
|
||||
modified: None,
|
||||
has_preview: false,
|
||||
materialised: true,
|
||||
}])
|
||||
let mut out = Vec::new();
|
||||
if let Some(size) = self.advertise {
|
||||
out.push(entry(CATALOG_NAME, size));
|
||||
}
|
||||
// The staged upload lists at the size of what was last put, as a
|
||||
// server that assembled it correctly would report — unless a test
|
||||
// says the assembly went wrong.
|
||||
if self.puts.load(Ordering::SeqCst) > 0 {
|
||||
let held = self
|
||||
.staged_size
|
||||
.unwrap_or(self.last_put.lock().unwrap().len() as u64);
|
||||
out.push(entry(UPLOAD_NAME, held));
|
||||
}
|
||||
Ok(out)
|
||||
}
|
||||
async fn dir_validator(
|
||||
&self,
|
||||
@@ -1566,6 +1612,30 @@ mod derived_guard_tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn an_upload_the_server_holds_at_the_wrong_size_is_not_rotated_in() {
|
||||
// The assembly went wrong on the server. Rotating that into place would
|
||||
// hand every other device the damaged file the generations exist to
|
||||
// survive; discarding it costs one retry.
|
||||
let (catalog_path, scratch) = fixture("wrong-size");
|
||||
let (mut backend, puts) = Fussy::reading(Some(RemoteError::NotFound("nope".into())));
|
||||
backend.staged_size = Some(7);
|
||||
let moves = backend.moves.clone();
|
||||
let mut report = SyncReport::default();
|
||||
sync_catalog(
|
||||
&backend,
|
||||
&RemotePath::new(".darkroom-derived"),
|
||||
&catalog_path,
|
||||
&scratch,
|
||||
&mut report,
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(puts.load(Ordering::SeqCst), 1, "it was uploaded");
|
||||
assert!(!report.catalog_uploaded, "but not accepted");
|
||||
assert!(moves.lock().unwrap().is_empty(), "and nothing was rotated");
|
||||
}
|
||||
|
||||
async fn run_with_moves(
|
||||
fail_with: Option<RemoteError>,
|
||||
name: &str,
|
||||
|
||||
Reference in New Issue
Block a user