1use std::num::NonZeroUsize;
16use std::path::{Path, PathBuf};
17use std::sync::{Arc, Condvar, Mutex, RwLock};
18use std::time::{Duration, Instant};
19use std::{fs::File, io::BufReader};
20
21use exif::{In, Reader as ExifReader, Tag};
22use lru::LruCache;
23
24use image::codecs::jpeg::JpegEncoder;
25use image::imageops::FilterType;
26use image::{DynamicImage, GenericImageView, ImageReader};
27
28use crate::services::image_utils::apply_exif_orientation;
29
30const THUMBNAIL_TIMEOUT: Duration = Duration::from_secs(10);
32
33#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
35pub enum ThumbnailSize {
36 Small, Medium, Large, }
40
41impl ThumbnailSize {
42 pub fn pixels(&self) -> u32 {
43 match self {
47 ThumbnailSize::Small => 260,
48 ThumbnailSize::Medium => 430,
49 ThumbnailSize::Large => 860,
50 }
51 }
52
53 pub fn dir_name(&self) -> &'static str {
54 match self {
55 ThumbnailSize::Small => "small",
56 ThumbnailSize::Medium => "medium",
57 ThumbnailSize::Large => "large",
58 }
59 }
60}
61
62#[derive(Debug, Clone)]
69pub struct ThumbnailEntry {
70 pub path: PathBuf,
71 pub file_size: u64,
72}
73
74const CACHE_ENTRY_CAPACITY: usize = 10_000;
80
81struct ConcurrencyLimiter {
87 max: usize,
88 max_background: usize,
89 state: Mutex<LimiterState>,
90 cv: Condvar,
91}
92
93#[derive(Debug, Default)]
94struct LimiterState {
95 foreground_active: usize,
96 background_active: usize,
97 foreground_waiting: usize,
98}
99
100#[derive(Debug, Clone, Copy)]
101enum GenerationPriority {
102 Foreground,
103 Background,
104}
105
106impl ConcurrencyLimiter {
107 fn new(max: usize, max_background: usize) -> Self {
108 Self {
109 max,
110 max_background,
111 state: Mutex::new(LimiterState::default()),
112 cv: Condvar::new(),
113 }
114 }
115
116 fn acquire(self: &Arc<Self>, priority: GenerationPriority) -> ConcurrencyPermit {
120 let mut state = self.state.lock().unwrap_or_else(|e| e.into_inner());
121 if matches!(priority, GenerationPriority::Foreground) {
122 state.foreground_waiting += 1;
123 }
124
125 loop {
126 let active = state.foreground_active + state.background_active;
127 let can_acquire = match priority {
128 GenerationPriority::Foreground => active < self.max,
129 GenerationPriority::Background => {
130 active < self.max
131 && state.background_active < self.max_background
132 && state.foreground_waiting == 0
133 }
134 };
135
136 if can_acquire {
137 match priority {
138 GenerationPriority::Foreground => {
139 state.foreground_waiting = state.foreground_waiting.saturating_sub(1);
140 state.foreground_active += 1;
141 }
142 GenerationPriority::Background => {
143 state.background_active += 1;
144 }
145 }
146 break;
147 }
148
149 state = self.cv.wait(state).unwrap_or_else(|e| e.into_inner());
150 }
151
152 ConcurrencyPermit {
153 limiter: self.clone(),
154 priority,
155 }
156 }
157
158 fn release(&self, priority: GenerationPriority) {
160 {
161 let mut state = self.state.lock().unwrap_or_else(|e| e.into_inner());
162 match priority {
163 GenerationPriority::Foreground => {
164 state.foreground_active = state.foreground_active.saturating_sub(1);
165 }
166 GenerationPriority::Background => {
167 state.background_active = state.background_active.saturating_sub(1);
168 }
169 }
170 }
171 self.cv.notify_all();
172 }
173}
174
175struct ConcurrencyPermit {
176 limiter: Arc<ConcurrencyLimiter>,
177 priority: GenerationPriority,
178}
179
180impl Drop for ConcurrencyPermit {
181 fn drop(&mut self) {
182 self.limiter.release(self.priority);
183 }
184}
185
186struct GenerationDeduper {
187 active: Mutex<std::collections::HashSet<String>>,
188 cv: Condvar,
189}
190
191impl GenerationDeduper {
192 fn new() -> Self {
193 Self {
194 active: Mutex::new(std::collections::HashSet::new()),
195 cv: Condvar::new(),
196 }
197 }
198
199 fn enter(self: &Arc<Self>, key: String) -> GenerationGuard {
200 let mut active = self.active.lock().unwrap_or_else(|e| e.into_inner());
201 while active.contains(&key) {
202 active = self.cv.wait(active).unwrap_or_else(|e| e.into_inner());
203 }
204 active.insert(key.clone());
205 GenerationGuard {
206 deduper: self.clone(),
207 key,
208 }
209 }
210
211 fn leave(&self, key: &str) {
212 {
213 let mut active = self.active.lock().unwrap_or_else(|e| e.into_inner());
214 active.remove(key);
215 }
216 self.cv.notify_all();
217 }
218}
219
220struct GenerationGuard {
221 deduper: Arc<GenerationDeduper>,
222 key: String,
223}
224
225impl Drop for GenerationGuard {
226 fn drop(&mut self) {
227 self.deduper.leave(&self.key);
228 }
229}
230
231pub struct ThumbnailService {
233 _drive_root: PathBuf,
235
236 cache_dir: PathBuf,
238
239 cache: Arc<RwLock<LruCache<(String, ThumbnailSize), ThumbnailEntry>>>,
245
246 max_cache_bytes: u64,
248
249 current_cache_bytes: Arc<RwLock<u64>>,
251
252 generation_limiter: Arc<ConcurrencyLimiter>,
254
255 generating: Arc<GenerationDeduper>,
258}
259
260impl ThumbnailService {
261 pub fn new<P: AsRef<Path>>(drive_root: P, max_cache_gb: f64) -> std::io::Result<Self> {
263 let drive_root = drive_root.as_ref().to_path_buf();
264 let cache_dir = drive_root.join(".photovault").join("thumbnails");
265
266 for size in [
268 ThumbnailSize::Small,
269 ThumbnailSize::Medium,
270 ThumbnailSize::Large,
271 ] {
272 std::fs::create_dir_all(cache_dir.join(size.dir_name()))?;
273 }
274
275 let max_cache_bytes = (max_cache_gb * 1024.0 * 1024.0 * 1024.0) as u64;
276
277 let capacity = NonZeroUsize::new(CACHE_ENTRY_CAPACITY)
278 .expect("CACHE_ENTRY_CAPACITY is a non-zero compile-time constant");
279
280 Ok(Self {
281 _drive_root: drive_root,
282 cache_dir,
283 cache: Arc::new(RwLock::new(LruCache::new(capacity))),
284 max_cache_bytes,
285 current_cache_bytes: Arc::new(RwLock::new(0)),
286 generation_limiter: Arc::new(ConcurrencyLimiter::new(4, 1)),
290 generating: Arc::new(GenerationDeduper::new()),
291 })
292 }
293
294 pub fn prewarm_small(
303 &self,
304 items: &[(PathBuf, String, i32)],
305 cancel: &std::sync::atomic::AtomicBool,
306 ) -> usize {
307 use rayon::prelude::*;
308 use std::sync::atomic::{AtomicUsize, Ordering};
309 let size = ThumbnailSize::Small;
310 let generated = AtomicUsize::new(0);
311 items
312 .par_iter()
313 .for_each(|(photo_path, file_hash, orientation)| {
314 if cancel.load(Ordering::Relaxed) {
315 return;
316 }
317 let cached = self.thumbnail_path(file_hash, size);
321 if cached.exists() {
322 return;
323 }
324 match self.generate_thumbnail_background(photo_path, file_hash, *orientation, size)
325 {
326 Ok(_) => {
327 generated.fetch_add(1, Ordering::Relaxed);
328 }
329 Err(e) => tracing::trace!("prewarm skipped {}: {}", photo_path.display(), e),
330 }
331 });
332 generated.load(Ordering::Relaxed)
333 }
334
335 pub fn generate_thumbnail(
339 &self,
340 photo_path: &Path,
341 file_hash: &str,
342 orientation: i32,
343 size: ThumbnailSize,
344 ) -> Result<PathBuf, String> {
345 self.generate_thumbnail_with_priority(
346 photo_path,
347 file_hash,
348 orientation,
349 size,
350 GenerationPriority::Foreground,
351 )
352 }
353
354 pub fn generate_thumbnail_background(
355 &self,
356 photo_path: &Path,
357 file_hash: &str,
358 orientation: i32,
359 size: ThumbnailSize,
360 ) -> Result<PathBuf, String> {
361 self.generate_thumbnail_with_priority(
362 photo_path,
363 file_hash,
364 orientation,
365 size,
366 GenerationPriority::Background,
367 )
368 }
369
370 fn generate_thumbnail_with_priority(
371 &self,
372 photo_path: &Path,
373 file_hash: &str,
374 orientation: i32,
375 size: ThumbnailSize,
376 priority: GenerationPriority,
377 ) -> Result<PathBuf, String> {
378 let thumb_path = self.thumbnail_path(file_hash, size);
379 if self.try_existing_thumbnail(file_hash, size, &thumb_path) {
380 return Ok(thumb_path);
381 }
382
383 let key = format!("{}:{file_hash}", size.dir_name());
384 let _generation = self.generating.enter(key);
385 if self.try_existing_thumbnail(file_hash, size, &thumb_path) {
386 return Ok(thumb_path);
387 }
388
389 let _permit = self.generation_limiter.acquire(priority);
390 let start = Instant::now();
391 self.generate_thumbnail_inner(photo_path, file_hash, orientation, size, start)
392 }
393
394 fn generate_thumbnail_inner(
396 &self,
397 photo_path: &Path,
398 file_hash: &str,
399 orientation: i32,
400 size: ThumbnailSize,
401 start: Instant,
402 ) -> Result<PathBuf, String> {
403 let thumb_path = self.thumbnail_path(file_hash, size);
405 if self.try_existing_thumbnail(file_hash, size, &thumb_path) {
406 return Ok(thumb_path);
407 }
408
409 if thumb_path.exists() {
411 let _ = std::fs::remove_file(&thumb_path);
412 }
413
414 if let Some(parent) = thumb_path.parent() {
416 std::fs::create_dir_all(parent)
417 .map_err(|e| format!("Failed to create thumbnail subdirectory: {}", e))?;
418 }
419
420 if start.elapsed() > THUMBNAIL_TIMEOUT {
422 return Err("Thumbnail generation timed out before decode".to_string());
423 }
424
425 let img = Self::decode_image_fast(photo_path, size)?;
427 let img = apply_exif_orientation(img, orientation);
428
429 if start.elapsed() > THUMBNAIL_TIMEOUT {
431 return Err("Thumbnail generation timed out after decode".to_string());
432 }
433
434 let thumb = self.create_thumbnail(&img, size);
436
437 let quality = match size {
441 ThumbnailSize::Small => 78,
442 ThumbnailSize::Medium => 83,
443 ThumbnailSize::Large => 88,
444 };
445 let mut out = std::fs::File::create(&thumb_path)
446 .map_err(|e| format!("Failed to create thumbnail file: {}", e))?;
447 let mut encoder = JpegEncoder::new_with_quality(&mut out, quality);
448 encoder
449 .encode_image(&thumb)
450 .map_err(|e| format!("Failed to encode thumbnail: {}", e))?;
451
452 self.add_to_cache(file_hash, size, &thumb_path);
453 self.track_cache_size(&thumb_path);
454 self.evict_if_needed();
455
456 Ok(thumb_path)
457 }
458
459 fn try_existing_thumbnail(
460 &self,
461 file_hash: &str,
462 size: ThumbnailSize,
463 thumb_path: &Path,
464 ) -> bool {
465 if !thumb_path.exists() {
466 return false;
467 }
468 let file_size = std::fs::metadata(thumb_path).map(|m| m.len()).unwrap_or(0);
469 if file_size == 0 {
470 return false;
471 }
472 if file_size <= 5000 && !Self::thumbnail_dimensions_are_adequate(thumb_path, size) {
477 return false;
478 }
479 self.add_to_cache(file_hash, size, thumb_path);
480 true
481 }
482
483 fn thumbnail_dimensions_are_adequate(thumb_path: &Path, size: ThumbnailSize) -> bool {
484 ImageReader::open(thumb_path)
485 .ok()
486 .and_then(|r| r.with_guessed_format().ok())
487 .and_then(|r| r.into_dimensions().ok())
488 .map(|(width, height)| width.max(height) >= size.pixels())
489 .unwrap_or(false)
490 }
491
492 fn decode_image_fast(
498 photo_path: &Path,
499 target_size: ThumbnailSize,
500 ) -> Result<DynamicImage, String> {
501 let ext = photo_path
507 .extension()
508 .and_then(|e| e.to_str())
509 .map(|e| e.to_ascii_lowercase());
510 if matches!(ext.as_deref(), Some("heic") | Some("heif")) {
511 return crate::services::image_io::open_image(photo_path);
512 }
513
514 if matches!(ext.as_deref(), Some("jpg") | Some("jpeg")) {
515 if let Some(img) = Self::decode_embedded_jpeg_thumbnail(photo_path, target_size) {
516 return Ok(img);
517 }
518 }
519
520 let reader =
521 ImageReader::open(photo_path).map_err(|e| format!("Failed to open image: {}", e))?;
522
523 let reader = reader
525 .with_guessed_format()
526 .map_err(|e| format!("Failed to guess image format: {}", e))?;
527
528 let mut limits = image::Limits::default();
530 limits.max_alloc = Some(200 * 1024 * 1024);
531
532 let mut reader = reader;
533 reader.limits(limits);
534
535 let img = reader
536 .decode()
537 .map_err(|e| format!("Failed to decode image: {}", e))?;
538
539 Ok(img)
540 }
541
542 fn decode_embedded_jpeg_thumbnail(
543 photo_path: &Path,
544 target_size: ThumbnailSize,
545 ) -> Option<DynamicImage> {
546 let file = File::open(photo_path).ok()?;
547 let mut reader = BufReader::new(file);
548 let exif = ExifReader::new().read_from_container(&mut reader).ok()?;
549 let offset = exif
550 .get_field(Tag::JPEGInterchangeFormat, In::THUMBNAIL)
551 .and_then(|f| f.value.get_uint(0))? as usize;
552 let len = exif
553 .get_field(Tag::JPEGInterchangeFormatLength, In::THUMBNAIL)
554 .and_then(|f| f.value.get_uint(0))? as usize;
555 let end = offset.checked_add(len)?;
556 let bytes = exif.buf().get(offset..end)?;
557 let img = image::load_from_memory(bytes).ok()?;
558 let (width, height) = img.dimensions();
559 if width.max(height) < target_size.pixels() {
560 return None;
561 }
562 Some(img)
563 }
564
565 fn track_cache_size(&self, thumb_path: &Path) {
567 let file_size = std::fs::metadata(thumb_path).map(|m| m.len()).unwrap_or(0);
568 if let Ok(mut current) = self.current_cache_bytes.write() {
569 *current += file_size;
570 }
571 }
572
573 fn create_thumbnail(&self, img: &DynamicImage, size: ThumbnailSize) -> DynamicImage {
581 let max_dim = size.pixels();
582 let (width, height) = img.dimensions();
583
584 let (new_width, new_height) = if width > height {
586 let ratio = max_dim as f64 / width as f64;
587 (max_dim, (height as f64 * ratio) as u32)
588 } else {
589 let ratio = max_dim as f64 / height as f64;
590 ((width as f64 * ratio) as u32, max_dim)
591 };
592
593 let filter = match size {
594 ThumbnailSize::Small => FilterType::Triangle,
595 ThumbnailSize::Medium => FilterType::CatmullRom,
596 ThumbnailSize::Large => FilterType::Lanczos3,
597 };
598 img.resize(new_width, new_height, filter)
599 }
600
601 fn thumbnail_path(&self, file_hash: &str, size: ThumbnailSize) -> PathBuf {
604 let subdir = &file_hash[..2.min(file_hash.len())];
606
607 self.cache_dir
608 .join(size.dir_name())
609 .join("v2")
610 .join(subdir)
611 .join(format!("{}.jpg", file_hash))
612 }
613
614 fn add_to_cache(&self, file_hash: &str, size: ThumbnailSize, path: &Path) {
616 let file_size = std::fs::metadata(path).map(|m| m.len()).unwrap_or(0);
617
618 let entry = ThumbnailEntry {
619 path: path.to_path_buf(),
620 file_size,
621 };
622
623 if let Ok(mut cache) = self.cache.write() {
624 cache.put((file_hash.to_string(), size), entry);
625 }
626 }
627
628 fn evict_if_needed(&self) {
633 let current = match self.current_cache_bytes.read() {
634 Ok(v) => *v,
635 Err(_) => return,
636 };
637
638 if current <= self.max_cache_bytes {
639 return;
640 }
641
642 let target = self.max_cache_bytes * 80 / 100;
643 let mut freed = 0u64;
644
645 loop {
649 let popped = {
650 let mut cache = match self.cache.write() {
651 Ok(c) => c,
652 Err(_) => break,
653 };
654 cache.pop_lru()
655 };
656
657 let Some((_, entry)) = popped else { break };
658
659 if std::fs::remove_file(&entry.path).is_ok() {
660 freed += entry.file_size;
661 }
662
663 if current.saturating_sub(freed) <= target {
664 break;
665 }
666 }
667
668 if let Ok(mut current) = self.current_cache_bytes.write() {
669 *current = current.saturating_sub(freed);
670 }
671
672 tracing::info!("Evicted {} bytes from thumbnail cache", freed);
673 }
674
675 pub fn load_existing_thumbnails(&self) -> std::io::Result<()> {
677 let mut total_size = 0u64;
678
679 for size in [
680 ThumbnailSize::Small,
681 ThumbnailSize::Medium,
682 ThumbnailSize::Large,
683 ] {
684 let size_dir = self.cache_dir.join(size.dir_name()).join("v2");
685
686 if !size_dir.exists() {
687 continue;
688 }
689
690 for subdir_entry in std::fs::read_dir(&size_dir)? {
691 let subdir_entry = subdir_entry?;
692 if !subdir_entry.file_type()?.is_dir() {
693 continue;
694 }
695
696 for file_entry in std::fs::read_dir(subdir_entry.path())? {
697 let file_entry = file_entry?;
698 let path = file_entry.path();
699
700 if path.extension().map(|e| e == "jpg").unwrap_or(false) {
701 if let Some(stem) = path.file_stem() {
702 let hash = stem.to_string_lossy().to_string();
703 let file_size = file_entry.metadata()?.len();
704
705 total_size += file_size;
706
707 let entry = ThumbnailEntry {
708 path: path.clone(),
709 file_size,
710 };
711
712 if let Ok(mut cache) = self.cache.write() {
713 cache.put((hash, size), entry);
714 }
715 }
716 }
717 }
718 }
719 }
720
721 if let Ok(mut current) = self.current_cache_bytes.write() {
722 *current = total_size;
723 }
724
725 let count = self.cache.read().map(|c| c.len()).unwrap_or(0);
726 tracing::info!(
727 "Loaded {} existing thumbnails ({} MB)",
728 count,
729 total_size / 1024 / 1024
730 );
731
732 Ok(())
733 }
734}
735
736#[cfg(test)]
737mod tests {
738 use super::*;
739 use tempfile::tempdir;
740
741 #[test]
742 fn test_thumbnail_path_generation() {
743 let temp = tempdir().unwrap();
744 let service = ThumbnailService::new(temp.path(), 1.0).unwrap();
745
746 let path = service.thumbnail_path("abcdef123456", ThumbnailSize::Medium);
747
748 assert!(path.to_string_lossy().contains("medium"));
749 assert!(path.to_string_lossy().contains("ab")); assert!(path.to_string_lossy().contains("abcdef123456.jpg"));
751 }
752
753 #[test]
754 fn test_thumbnail_sizes() {
755 assert_eq!(ThumbnailSize::Small.pixels(), 260);
756 assert_eq!(ThumbnailSize::Medium.pixels(), 430);
757 assert_eq!(ThumbnailSize::Large.pixels(), 860);
758 }
759
760 #[test]
761 fn test_cache_directories_created() {
762 let temp = tempdir().unwrap();
763 let _service = ThumbnailService::new(temp.path(), 1.0).unwrap();
764
765 assert!(temp.path().join(".photovault/thumbnails/small").exists());
766 assert!(temp.path().join(".photovault/thumbnails/medium").exists());
767 assert!(temp.path().join(".photovault/thumbnails/large").exists());
768 }
769
770 #[test]
771 fn small_but_dimensionally_valid_thumbnail_is_reused() {
772 let temp = tempdir().unwrap();
773 let service = ThumbnailService::new(temp.path(), 1.0).unwrap();
774 let path = service.thumbnail_path("abcdef123456", ThumbnailSize::Small);
775 std::fs::create_dir_all(path.parent().unwrap()).unwrap();
776 DynamicImage::new_rgb8(ThumbnailSize::Small.pixels(), ThumbnailSize::Small.pixels())
777 .save(&path)
778 .unwrap();
779
780 assert!(std::fs::metadata(&path).unwrap().len() <= 5000);
781 assert!(service.try_existing_thumbnail("abcdef123456", ThumbnailSize::Small, &path));
782 }
783
784 #[test]
785 fn test_generation_deduper_serializes_same_key() {
786 let deduper = Arc::new(GenerationDeduper::new());
787 let first = deduper.enter("medium:abc".into());
788 let (tx, rx) = std::sync::mpsc::channel();
789 let deduper_for_thread = deduper.clone();
790
791 let handle = std::thread::spawn(move || {
792 let _second = deduper_for_thread.enter("medium:abc".into());
793 tx.send(()).unwrap();
794 });
795
796 assert!(rx.recv_timeout(Duration::from_millis(50)).is_err());
797 drop(first);
798 assert!(rx.recv_timeout(Duration::from_secs(1)).is_ok());
799 handle.join().unwrap();
800 }
801
802 #[test]
803 fn test_concurrency_limiter_blocks_until_release() {
804 use std::sync::Arc;
805 use std::thread;
806 use std::time::Duration;
807
808 let limiter = Arc::new(ConcurrencyLimiter::new(1, 1));
809 let permit = limiter.acquire(GenerationPriority::Foreground);
810
811 let l = limiter.clone();
815 let started = Arc::new(std::sync::atomic::AtomicBool::new(false));
816 let s = started.clone();
817 let handle = thread::spawn(move || {
818 l.acquire(GenerationPriority::Foreground);
819 s.store(true, std::sync::atomic::Ordering::Release);
820 });
821
822 thread::sleep(Duration::from_millis(20));
823 assert!(!started.load(std::sync::atomic::Ordering::Acquire));
824
825 drop(permit);
826 handle.join().unwrap();
827 assert!(started.load(std::sync::atomic::Ordering::Acquire));
828 }
829
830 #[test]
831 fn test_background_generation_is_capped() {
832 use std::sync::mpsc;
833 use std::sync::Arc;
834 use std::thread;
835 use std::time::Duration;
836
837 let limiter = Arc::new(ConcurrencyLimiter::new(4, 1));
838 let permit = limiter.acquire(GenerationPriority::Background);
839
840 let (tx, rx) = mpsc::channel();
841 let l = limiter.clone();
842 let handle = thread::spawn(move || {
843 let _permit = l.acquire(GenerationPriority::Background);
844 tx.send(()).unwrap();
845 });
846
847 thread::sleep(Duration::from_millis(20));
848 assert!(rx.try_recv().is_err());
849
850 drop(permit);
851 rx.recv_timeout(Duration::from_secs(1)).unwrap();
852 handle.join().unwrap();
853 }
854
855 #[test]
856 fn test_waiting_foreground_generation_preempts_background() {
857 use std::sync::mpsc;
858 use std::sync::Arc;
859 use std::thread;
860 use std::time::Duration;
861
862 let limiter = Arc::new(ConcurrencyLimiter::new(1, 1));
863 let permit = limiter.acquire(GenerationPriority::Background);
864
865 let (tx, rx) = mpsc::channel();
866 let l = limiter.clone();
867 let tx_foreground = tx.clone();
868 let foreground = thread::spawn(move || {
869 let _permit = l.acquire(GenerationPriority::Foreground);
870 tx_foreground.send("foreground").unwrap();
871 });
872
873 thread::sleep(Duration::from_millis(20));
874 {
875 let state = limiter.state.lock().unwrap();
876 assert_eq!(state.foreground_waiting, 1);
877 }
878
879 let l = limiter.clone();
880 let background = thread::spawn(move || {
881 let _permit = l.acquire(GenerationPriority::Background);
882 tx.send("background").unwrap();
883 });
884
885 thread::sleep(Duration::from_millis(20));
886 assert!(rx.try_recv().is_err());
887
888 drop(permit);
889 assert_eq!(
890 rx.recv_timeout(Duration::from_secs(1)).unwrap(),
891 "foreground"
892 );
893 assert_eq!(
894 rx.recv_timeout(Duration::from_secs(1)).unwrap(),
895 "background"
896 );
897
898 foreground.join().unwrap();
899 background.join().unwrap();
900 }
901}