Let one outbox drain run at a time
A finished merge drained the outbox and the sync pass that followed drained it again beside it: each read the 800 MB composite into memory and sent it, and the later one found its record cleared underneath it and logged the file as missing. Drains now take a lock and the one that waited finds the queue empty.
This commit is contained in:
@@ -400,6 +400,9 @@ pub enum UploadMessage {
|
|||||||
},
|
},
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Held by whichever drain is sending the outbox. See [`spawn_upload`].
|
||||||
|
static OUTBOX_DRAIN: std::sync::Mutex<()> = std::sync::Mutex::new(());
|
||||||
|
|
||||||
/// Drain the outbox to the server.
|
/// Drain the outbox to the server.
|
||||||
///
|
///
|
||||||
/// Its own thread with its own runtime, like every other network path here —
|
/// Its own thread with its own runtime, like every other network path here —
|
||||||
@@ -416,6 +419,16 @@ pub fn spawn_upload(
|
|||||||
let (tx, rx) = std::sync::mpsc::channel();
|
let (tx, rx) = std::sync::mpsc::channel();
|
||||||
|
|
||||||
executors::spawn(Executor::Network, "upload", move || {
|
executors::spawn(Executor::Network, "upload", move || {
|
||||||
|
// One drain at a time. Three things start one — a sync pass, an
|
||||||
|
// export, a finished merge — and they used to run side by side over
|
||||||
|
// the same queue: each read the same 800 MB composite into memory and
|
||||||
|
// sent it, and the one that came second found its record cleared
|
||||||
|
// under it and reported the file missing. Waiting here costs nothing:
|
||||||
|
// the drain that holds the lock is already sending everything queued,
|
||||||
|
// and the one that waited finds the queue empty and finishes.
|
||||||
|
let _draining = OUTBOX_DRAIN
|
||||||
|
.lock()
|
||||||
|
.unwrap_or_else(std::sync::PoisonError::into_inner);
|
||||||
let rt = match crate::net_runtime::build() {
|
let rt = match crate::net_runtime::build() {
|
||||||
Ok(rt) => rt,
|
Ok(rt) => rt,
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
|
|||||||
Reference in New Issue
Block a user