//! TRACES: FR-MRG-1 | FR-MRG-2 | FR-MRG-3 | FR-MRG-5 | FR-MRG-7 //! The merge job: from a set of files to a linear DNG beside them. //! //! The orchestration of a panorama, with no interface types in it: it runs //! on a thread of its own and reports through a channel, on the pattern //! [`crate::export`]'s batch established (FR-MRG-7). The interface drains //! the channel from a timer; a headless example drains it from a loop. //! //! # Stages //! //! 1. **Load.** Each file is read and decoded to sensor data, and its edit //! graph built the way a session would — orientation, lens profile — //! because those are the two things the camera-space tap uses //! (FR-MRG-2). //! 2. **Proxies.** Each frame is demosaiced and rendered through the tap //! at proxy size, upright; the keypoint detector reads that. //! 3. **Alignment.** `dr_pano::align`. A frame it could not place stops //! the job with the frame named (FR-MRG-5). //! 4. **Gain.** One scalar per frame from the proxies' overlaps, so the //! stop of exposure drift a hand-held sweep collects does not band. //! 5. **Merge.** `dr_gpu::MergePass`, chunk by chunk, into a DNG written //! strip by strip beside the first source (FR-MRG-3). //! //! The demosaiced frames are the memory: 160 MB each at f16 for a 20 MP //! body, so at most two are resident and the rest are demosaiced again //! when a band needs them. That is the trade FR-MRG-11 asks for — the //! sensor data stays on the CPU at a fifth of the size — and it is what //! makes a twelve-frame set fit. use std::collections::VecDeque; use std::path::{Path, PathBuf}; use std::sync::mpsc::{Receiver, Sender}; use std::sync::Arc; use std::time::Instant; use dr_decode::RawImage; use dr_gpu::{AdjustPass, DemosaicedImage, Demosaicer, GpuContext, MergeFrame, MergeOutput, MergePass}; use dr_pano::bundle::Cameras; use dr_pano::projection::{self, Projection}; use dr_pano::{Alignment, Gray}; use dr_pipeline::EditGraph; pub use crate::export::Cancel; /// One frame as the job receives it: the file's bytes, already fetched, /// and a name for messages and for the composite's own name. #[derive(Debug, Clone)] pub struct MergeInput { /// The file's name, `_MG_8320.CR2`. pub name: String, pub bytes: Arc>, } impl MergeInput { pub fn read(path: &Path) -> Result { Ok(MergeInput { name: path .file_name() .map(|n| n.to_string_lossy().into_owned()) .unwrap_or_else(|| "frame".into()), bytes: Arc::new(std::fs::read(path).map_err(|e| format!("{}: {e}", path.display()))?), }) } } /// Where the composite goes (FR-MRG-3). #[derive(Debug, Clone)] pub enum MergeDestination { /// A folder on this device: written there directly. Local(PathBuf), /// Beside its sources in the library, through the export outbox: the /// file is staged with a destination record and the drain uploads it — /// the one path that works for a folder library and a server alike, and /// the one FR-MRG-3 names for a folder the job cannot write itself. Outbox { outbox: PathBuf, /// The sources' folder, relative to the library root. remote_dir: String, }, } /// What the job is told. #[derive(Debug, Clone)] pub struct MergeRequest { /// The frames, in capture order. pub frames: Vec, pub destination: MergeDestination, /// `None` for the projection the field of view suggests. pub projection: Option, /// Pixels over which a frame's weight ramps up from its edge. pub feather_px: f32, /// GPU work unit; also the DNG strip height. pub chunk: (u32, u32), } impl MergeRequest { pub fn new(frames: Vec, destination: MergeDestination) -> Self { MergeRequest { frames, destination, projection: None, feather_px: 200.0, chunk: (2048, 512), } } } /// What the photographer decides once the alignment is shown (FR-MRG-1: /// never automatic). #[derive(Debug, Clone, Copy, PartialEq)] pub enum Decision { Merge { projection: Option }, Abandon, } /// What the job says while it runs. #[derive(Debug, Clone)] pub enum MergeEvent { /// A stage, with progress within it. Progress { stage: &'static str, done: usize, total: usize, }, /// The alignment, before any pixel is written: what the interface shows /// for the photographer to confirm, and what names a failure. Aligned(AlignmentReport), /// The composite is written: on the device at `path`, or staged in the /// outbox for the drain to upload, in which case `staged` is true and /// the library learns of it when the folder is next scanned. Done { path: PathBuf, staged: bool, width: u32, height: u32, }, Failed(String), Cancelled, } /// The alignment, described for a panel. #[derive(Debug, Clone)] pub struct AlignmentReport { /// Equivalent focal length in millimetres on full frame, from the fit. pub focal_mm: f64, pub rms_px: f64, pub projection: Projection, pub width: u32, pub height: u32, /// Per frame, in input order: yaw and pitch in degrees if aligned, or /// why not. pub frames: Vec>, pub links: usize, /// The aligned set drawn on the suggested surface at proxy resolution, /// `(width, height, rgba)` — what the photographer confirms. pub preview: Option<(u32, u32, Vec)>, } /// Run the whole job on the calling thread, reporting on `events`. /// /// After `Aligned` the job waits on `decision` — the photographer's /// confirmation, or the headless caller's immediate `Merge` — and returns /// when the file is written, the job failed, or `cancel` was seen. pub fn run( ctx: GpuContext, request: MergeRequest, events: Sender, decision: Receiver, cancel: Cancel, ) { let send = |e: MergeEvent| { let _ = events.send(e); }; match run_inner(&ctx, &request, &events, &decision, &cancel) { Ok(Some(done)) => send(done), Ok(None) => send(MergeEvent::Cancelled), Err(e) => send(MergeEvent::Failed(e)), } } /// The frames as loaded: sensor data on the CPU and the graph the tap uses. struct Loaded { raw: RawImage, meta: dr_decode::Metadata, graph: Arc, /// The upright size, which is what the cameras are in. size: (u32, u32), } fn run_inner( ctx: &GpuContext, request: &MergeRequest, events: &Sender, decision: &Receiver, cancel: &Cancel, ) -> Result, String> { if request.frames.len() < 2 { return Err("a panorama needs at least two frames".into()); } let progress = |stage: &'static str, done: usize, total: usize| { let _ = events.send(MergeEvent::Progress { stage, done, total }); }; // 1. Load. let t = Instant::now(); let mut frames: Vec = Vec::with_capacity(request.frames.len()); for (i, input) in request.frames.iter().enumerate() { if cancel.is_cancelled() { return Ok(None); } progress("Reading", i, request.frames.len()); let bytes = &input.bytes; let meta = dr_decode::metadata(bytes).unwrap_or_default(); let raw = dr_decode::decode(bytes).map_err(|e| format!("{}: {e}", input.name))?; let orientation = meta.orientation.unwrap_or_default(); let mut graph = EditGraph::default_chain(); graph.set_orientation(orientation); graph.set_lens_profile(crate::develop::DevelopSession::profile_for(&meta)); let (w, h) = (raw.crop.width, raw.crop.height); let size = if orientation.quarter_turns % 2 == 1 { (h, w) } else { (w, h) }; frames.push(Loaded { raw, meta, graph: Arc::new(graph), size, }); } log::info!("merge: {} frames read in {:?}", frames.len(), t.elapsed()); // 2. Proxies and keypoints. let t = Instant::now(); let demosaicer = Demosaicer::new(ctx).map_err(|e| e.to_string())?; let mut adjust = AdjustPass::new(ctx); let mut detector = dr_pano::xfeat::XFeat::embedded().map_err(|e| e.to_string())?; let mut features = Vec::with_capacity(frames.len()); let mut proxies: Vec = Vec::with_capacity(frames.len()); let mut colour: Vec> = Vec::with_capacity(frames.len()); for (i, f) in frames.iter().enumerate() { if cancel.is_cancelled() { return Ok(None); } progress("Finding features", i, frames.len()); let image = demosaicer.run(&f.raw).map_err(|e| e.to_string())?; let (proxy, rgb) = camera_proxy(&mut adjust, &image, &f.graph, f.size)?; features.push(detector.detect(&proxy).map_err(|e| e.to_string())?); proxies.push(proxy); colour.push(rgb); } log::info!("merge: features in {:?}", t.elapsed()); // 3. Alignment. let t = Instant::now(); progress("Aligning", 0, 1); let alignment = dr_pano::align(&features, &dr_pano::AlignOptions::default()) .map_err(|e| e.to_string())?; log::info!( "merge: aligned in {:?}, focal {:.1} px, rms {:.2} px, {} links", t.elapsed(), alignment.focal, alignment.rms_px, alignment.links.len() ); // The geometry at full resolution: the proxy's long edge against the // frame's. let proxy_long = proxies[0].width.max(proxies[0].height) as f64; let full_long = frames[0].size.0.max(frames[0].size.1) as f64; let focal_full = alignment.focal * full_long / proxy_long; let cameras = Cameras { rotations: alignment .rotations .iter() .map(|r| r.unwrap_or(dr_pano::linalg::Mat3::IDENTITY)) .collect(), focal: focal_full, }; let frame_size = (frames[0].size.0 as f64, frames[0].size.1 as f64); let (hfov, vfov) = field_of_view(&cameras, frame_size); let projection = request .projection .unwrap_or_else(|| Projection::suggest(hfov, vfov)); let bounds = projection::bounds(projection, focal_full, &cameras, frame_size) .ok_or("the frames project nowhere")?; // 4. Gains, before the report so the preview shows them. let gains = if alignment.is_complete() { gains(&proxies, &alignment) } else { vec![1.0; frames.len()] }; let preview = if alignment.is_complete() { progress("Drawing the preview", 0, 1); Some(preview(&colour, &proxies, &alignment, &gains, projection, &frames[0].raw)) } else { None }; let report = AlignmentReport { focal_mm: focal_full * 36.0 / full_long, rms_px: alignment.rms_px, projection, width: bounds.width().ceil() as u32, height: bounds.height().ceil() as u32, frames: describe(&alignment), links: alignment.links.len(), preview, }; let _ = events.send(MergeEvent::Aligned(report)); if !alignment.is_complete() { let names: Vec = alignment .unaligned .iter() .map(|(k, why)| format!("{}: {why}", request.frames[*k].name)) .collect(); return Err(format!("not every frame could be placed — {}", names.join("; "))); } // Never automatic (FR-MRG-1): nothing is written until the alignment // has been seen and confirmed. Polled, so a cancel while waiting is // seen within a moment. let projection = loop { if cancel.is_cancelled() { return Ok(None); } match decision.recv_timeout(std::time::Duration::from_millis(100)) { Ok(Decision::Merge { projection: p }) => break p.unwrap_or(projection), Ok(Decision::Abandon) => return Ok(None), Err(std::sync::mpsc::RecvTimeoutError::Timeout) => continue, Err(std::sync::mpsc::RecvTimeoutError::Disconnected) => return Ok(None), } }; let bounds = projection::bounds(projection, focal_full, &cameras, frame_size) .ok_or("the frames project nowhere")?; // 5. Merge, into a DNG beside the first frame. let name = format!("{}-pano.dng", stem(Path::new(&request.frames[0].name))); let (out_path, staged) = match &request.destination { MergeDestination::Local(dir) => (unused_name(dir, &name), false), MergeDestination::Outbox { outbox, remote_dir } => { std::fs::create_dir_all(outbox).map_err(|e| format!("{}: {e}", outbox.display()))?; let path = unused_name(outbox, &name); // The record first here, unlike an export: the payload is // written over minutes and a record naming a half-written file // is worse than a payload with no record, so the record is // removed again if the merge fails. The drain skips payloads // without records. let record = crate::export::destination_record(&path); let final_name = path .file_name() .map(|n| n.to_string_lossy().into_owned()) .unwrap_or(name.clone()); std::fs::write(&record, format!("{remote_dir}\n{final_name}\n")) .map_err(|e| format!("{}: {e}", record.display()))?; (path, true) } }; let cleanup = |path: &Path| { let _ = std::fs::remove_file(path); if staged { let _ = std::fs::remove_file(crate::export::destination_record(path)); } }; let (out_w, out_h) = (report_size(&bounds).0, report_size(&bounds).1); let first = &frames[0]; let black = first.raw.black_level[0]; let white_level = u32::from(first.raw.white_level.saturating_sub(black)).max(1); let profile = dng_profile(first, white_level); let (_, carried) = crate::export::header_for_file(&first.meta); let output = MergeOutput { projection, scale: focal_full, bounds, feather: request.feather_px, chunk: request.chunk, sample_scale: white_level as f32, }; let merge_frames: Vec = frames .iter() .zip(&gains) .map(|(f, &gain)| MergeFrame { graph: f.graph.clone(), gain, }) .collect(); // The writer pulls strips; the merge pushes bands. A channel between // them, and the writer on its own thread, so neither waits on the // other's pace more than one band. let (band_tx, band_rx) = std::sync::mpsc::sync_channel::>(1); let rows_per_strip = request.chunk.1.max(1); let file = std::fs::File::create(&out_path).map_err(|e| format!("{}: {e}", out_path.display()))?; let writer = std::thread::spawn(move || -> Result<(), String> { let mut file = std::io::BufWriter::new(file); dr_export::write_linear_dng( &mut file, out_w, out_h, rows_per_strip, &profile, Some(&carried), |_, buf| { let band = band_rx .recv() .map_err(|_| dr_export::ExportError::Encode("the merge stopped early".into()))?; buf.extend_from_slice(&band); Ok(()) }, ) .map_err(|e| e.to_string()) }); let t = Instant::now(); let mut pass = MergePass::new(ctx).map_err(|e| e.to_string())?; let mut resident: VecDeque<(usize, Arc)> = VecDeque::new(); let total_bands = out_h.div_ceil(rows_per_strip) as usize; let mut bands_done = 0usize; progress("Merging", 0, total_bands); let merged = pass.merge( &mut adjust, &merge_frames, &cameras, (frames[0].size.0, frames[0].size.1), &output, |k| { if let Some((_, img)) = resident.iter().find(|(i, _)| *i == k) { return Ok(img.clone()); } let img = Arc::new(demosaicer.run(&frames[k].raw)?); resident.push_back((k, img.clone())); while resident.len() > 2 { resident.pop_front(); } Ok(img) }, |band| { bands_done += 1; progress("Merging", bands_done, total_bands); band_tx .send(band.rgb.to_vec()) .map_err(|_| dr_gpu::GpuError::Readback("the writer stopped".into())) }, || cancel.is_cancelled(), ); drop(band_tx); let written = writer.join().unwrap_or_else(|_| Err("the writer panicked".into())); match merged { Ok(()) => {} Err(e) if cancel.is_cancelled() => { cleanup(&out_path); log::info!("merge cancelled: {e}"); return Ok(None); } Err(e) => { cleanup(&out_path); return Err(e.to_string()); } } if let Err(e) = written { cleanup(&out_path); return Err(e); } log::info!("merge: {}×{} written to {} in {:?}", out_w, out_h, out_path.display(), t.elapsed()); Ok(Some(MergeEvent::Done { path: out_path, staged, width: out_w, height: out_h, })) } /// The proxy the detector reads: the frame through the camera-space tap at /// proxy size, upright, as gamma-encoded grey. /// /// Through the tap rather than the embedded preview so that the keypoints /// are in the *undistorted* frame the tiles will be rendered in — a lens /// profile moves the corners by tens of pixels at full resolution, and an /// alignment measured on the distorted preview would be wrong by that much /// at the seams. fn camera_proxy( adjust: &mut AdjustPass, image: &DemosaicedImage, graph: &EditGraph, upright: (u32, u32), ) -> Result<(Gray, Vec), String> { let long = dr_pano::xfeat::INPUT_LONG_EDGE as f64; let scale = (long / f64::from(upright.0.max(upright.1))).min(1.0); let w = ((f64::from(upright.0) * scale).round() as u32).max(1); let h = ((f64::from(upright.1) * scale).round() as u32).max(1); let shader = graph.compose_camera_linear(dr_pipeline::CropRect::default()); adjust .render_camera_linear(image, &shader, w, h) .map_err(|e| e.to_string())?; let (rgba, rw, rh) = adjust.read_camera_linear().map_err(|e| e.to_string())?; // Camera RGB is unbalanced — green-heavy on every Bayer body — and // linear. Luma weights would over-count green further; an equal mean is // as good a grey as any for corners and edges, and the gamma is what // gives the detector the contrast range it was trained on. let data: Vec = rgba .chunks_exact(4) .map(|p| ((p[0] + p[1] + p[2]) / 3.0).clamp(0.0, 1.0).powf(1.0 / 2.2)) .collect(); let rgb: Vec = rgba.chunks_exact(4).flat_map(|p| [p[0], p[1], p[2]]).collect(); Ok(( Gray { width: rw as usize, height: rh as usize, data, }, rgb, )) } /// The aligned set on its surface, in colour, for the page. /// /// A quick look, not the pipeline: the first frame's white balance and /// matrix, a gamma, and the frames averaged where they overlap with their /// gains applied. Ghosting here is the alignment's error and banding is /// the gains', which is exactly what the photographer is being asked to /// look at. Fitted to 1600 px across. fn preview( colour: &[Vec], proxies: &[Gray], alignment: &Alignment, gains: &[f32], projection: Projection, first: &RawImage, ) -> (u32, u32, Vec) { let cameras = alignment.cameras(); let (fw, fh) = (proxies[0].width as f64, proxies[0].height as f64); let scale = alignment.focal; let Some(bounds) = projection::bounds(projection, scale, &cameras, (fw, fh)) else { return (0, 0, Vec::new()); }; let out_w = 1600usize.min(bounds.width().ceil() as usize).max(1); let px = bounds.width() / out_w as f64; let out_h = ((bounds.height() / px).ceil() as usize).max(1); let wb = first.wb_coeffs; let m = first .color_matrix .unwrap_or([1.0, 0.0, 0.0, 0.0, 1.0, 0.0, 0.0, 0.0, 1.0]); let mut out = vec![0u8; out_w * out_h * 4]; for oy in 0..out_h { for ox in 0..out_w { let u = bounds.min_u + (ox as f64 + 0.5) * px; let v = bounds.min_v + (oy as f64 + 0.5) * px; let d = projection.to_direction(scale, u, v); let mut sum = [0.0f32; 3]; let mut n = 0u32; for (k, g) in proxies.iter().enumerate() { let Some((x, y)) = cameras.project(k, d) else { continue }; let (x, y) = (x + fw / 2.0, y + fh / 2.0); if x < 0.0 || y < 0.0 || x >= fw - 1.0 || y >= fh - 1.0 { continue; } let i = (y as usize * g.width + x as usize) * 3; for c in 0..3 { sum[c] += colour[k][i + c] * gains[k]; } n += 1; } let o = (oy * out_w + ox) * 4; if n == 0 { out[o + 3] = 255; continue; } let cam = [sum[0] / n as f32 * wb[0], sum[1] / n as f32 * wb[1], sum[2] / n as f32 * wb[2]]; for c in 0..3 { let lin = m[c * 3] * cam[0] + m[c * 3 + 1] * cam[1] + m[c * 3 + 2] * cam[2]; out[o + c] = (lin.clamp(0.0, 1.0).powf(1.0 / 2.2) * 255.0) as u8; } out[o + 3] = 255; } } (out_w as u32, out_h as u32, out) } /// The horizontal and vertical field of view the cameras span, in radians: /// the angle between the extreme frame centres plus one frame's own field. fn field_of_view(cameras: &Cameras, frame: (f64, f64)) -> (f64, f64) { let f = cameras.focal; let own_h = 2.0 * (frame.0 / (2.0 * f)).atan(); let own_v = 2.0 * (frame.1 / (2.0 * f)).atan(); let (mut min_yaw, mut max_yaw, mut min_pitch, mut max_pitch) = (0.0f64, 0.0f64, 0.0f64, 0.0f64); for r in &cameras.rotations { let d = *r * dr_pano::linalg::Vec3::new(0.0, 0.0, 1.0); let yaw = d.x().atan2(d.z()); let pitch = d.y().asin(); min_yaw = min_yaw.min(yaw); max_yaw = max_yaw.max(yaw); min_pitch = min_pitch.min(pitch); max_pitch = max_pitch.max(pitch); } (max_yaw - min_yaw + own_h, max_pitch - min_pitch + own_v) } fn report_size(b: &projection::Bounds) -> (u32, u32) { (b.width().ceil().max(1.0) as u32, b.height().ceil().max(1.0) as u32) } fn describe(a: &Alignment) -> Vec> { a.rotations .iter() .enumerate() .map(|(k, r)| match r { Some(r) => { let yaw = r.0[0][2].atan2(r.0[2][2]).to_degrees(); let pitch = (-r.0[1][2]).asin().to_degrees(); Ok((yaw, pitch)) } None => Err(a .unaligned .iter() .find(|(i, _)| *i == k) .map(|(_, why)| why.to_string()) .unwrap_or_else(|| "not aligned".into())), }) .collect() } /// One gain per frame, from the proxies' overlaps. /// /// For each link, the mean grey of each frame over the region both see is /// compared; the log-gains that best reconcile every link are solved in /// least squares with the reference frame held at 1. Grey here is the /// proxy's gamma-encoded value raised back to linear, so the gain is a /// linear multiplier as the shader applies it. fn gains(proxies: &[Gray], a: &Alignment) -> Vec { let n = proxies.len(); let cameras = a.cameras(); let (w, h) = (proxies[0].width as f64, proxies[0].height as f64); let linear = |g: &Gray, x: usize, y: usize| g.data[y * g.width + x].powf(2.2) as f64; // Ratios per link. let mut ratios: Vec<(usize, usize, f64)> = Vec::new(); for l in &a.links { let (mut sum_i, mut sum_j, mut count) = (0.0, 0.0, 0usize); let step = 8; for y in (0..proxies[l.i].height).step_by(step) { for x in (0..proxies[l.i].width).step_by(step) { let p = (x as f64 - w / 2.0, y as f64 - h / 2.0); let d = cameras.bearing(l.i, p); let Some((qx, qy)) = cameras.project(l.j, d) else { continue }; let (qx, qy) = (qx + w / 2.0, qy + h / 2.0); if qx < 0.0 || qy < 0.0 || qx >= w - 1.0 || qy >= h - 1.0 { continue; } sum_i += linear(&proxies[l.i], x, y); sum_j += linear(&proxies[l.j], qx as usize, qy as usize); count += 1; } } if count >= 50 && sum_i > 0.0 && sum_j > 0.0 { ratios.push((l.i, l.j, (sum_i / sum_j).ln())); } } if ratios.is_empty() { return vec![1.0; n]; } // Least squares on log gains: g_j - g_i = ln(mean_i / mean_j), with // the reference frame's gain fixed at 0 by a strong prior. let root = a .rotations .iter() .position(|r| *r == Some(dr_pano::linalg::Mat3::IDENTITY)) .unwrap_or(0); let mut ata = dr_pano::linalg::DMat::zeros(n); let mut atb = vec![0.0f64; n]; for &(i, j, r) in &ratios { // Row: +1 at j, −1 at i, rhs r. ata[(j, j)] += 1.0; ata[(i, i)] += 1.0; ata[(i, j)] -= 1.0; ata[(j, i)] -= 1.0; atb[j] += r; atb[i] -= r; } ata[(root, root)] += 1e6; // A tiny ridge keeps an unlinked frame (there are none if the alignment // is complete) from making the system singular. for k in 0..n { ata[(k, k)] += 1e-9; } match ata.solve_spd(&atb) { Some(g) => g.iter().map(|v| v.exp() as f32).collect(), None => vec![1.0; n], } } fn dng_profile(first: &Loaded, white_level: u32) -> dr_export::DngProfile { let wb = first.raw.wb_coeffs; let neutral_from_wb = [1.0 / wb[0].max(1e-3), 1.0 / wb[1].max(1e-3), 1.0 / wb[2].max(1e-3)]; let (calibrations, as_shot_neutral) = match &first.raw.profile { Some(p) => (p.dng_calibrations(), p.neutral().unwrap_or(neutral_from_wb)), None => (Vec::new(), neutral_from_wb), }; let unique_model = match (first.raw.make.trim(), first.raw.model.trim()) { ("", "") => "DarkRoom panorama".to_string(), (make, model) if model.starts_with(make) => model.to_string(), (make, model) => format!("{make} {model}"), }; dr_export::DngProfile { unique_model, calibrations, as_shot_neutral, white_level, } } fn stem(p: &Path) -> String { p.file_stem() .map(|s| s.to_string_lossy().into_owned()) .unwrap_or_else(|| p.display().to_string()) } /// `name` in `dir`, numbered if that name is taken: a merge never /// overwrites (FR-MRG-3). fn unused_name(dir: &Path, name: &str) -> PathBuf { let mut candidate = dir.join(name); let base = stem(Path::new(name)); let mut n = 2; while candidate.exists() { candidate = dir.join(format!("{base}-{n}.dng")); n += 1; } candidate } /// Drain everything a job has said so far. pub fn drain(rx: &Receiver) -> Vec { let mut out = Vec::new(); while let Ok(e) = rx.try_recv() { out.push(e); } out }