MetonymyQT icon

Tokio Rust Concurrency Attempt 3

MetonymyQT | PRO | 04/06/25 01:39:45 PM UTC | 0 ⭐ | 8021 👁️ | Never ⏰ | [rust, Concurrency]
Rust |

4.95 KB

|

Software

|

0 👍

/

0 👎

use anyhow::{Result, anyhow};
use log::{debug, error, info};
use signal_hook::consts::{SIGINT, SIGTERM};
use signal_hook::iterator::Signals;
use std::ops::Deref;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::thread;
use std::time::Duration;
use tokio::fs::File;
use tokio::io::AsyncWriteExt;
use tokio::sync::mpsc;
use tokio::sync::mpsc::Sender;
use tokio::sync::oneshot;
use tokio::task::JoinHandle;
use tokio::time::{sleep, timeout};
 
fn setup_graceful_shutdown(cancelled: oneshot::Sender<bool>) -> JoinHandle<()> {
    tokio::spawn(async move {
        let signals = Signals::new([SIGINT, SIGTERM]);
        match signals {
            Ok(mut signal_info) => {
                loop {
                    if signal_info.pending().next().is_some() {
                        info!("user request cancellation");
                        let _ = cancelled.send(true);
                        return;
                    }
                    sleep(Duration::from_millis(250)).await;
                }
            }
            Err(error) => {
                error!("Failed to setup signal handler: {error}")
            }
        }
    })
}
 
async fn download_image(tx: Sender<Result<bytes::Bytes, anyhow::Error>>, url: &str) {
    info!("Downloading image for {}", url);
    let resp = reqwest::get(url).await;
    if let Ok(resp) = resp {
        let image_bytes = resp.bytes().await;
        match image_bytes {
            Ok(data) => {
                tx.send(Ok(data)).await.expect("failed to send message");
            }
            Err(err) => {
                let send_result = tx.send(Err(anyhow!("failed to get image data: {}", err)))
                    .await;
 
                debug!("download image ok {:?} for url {}", send_result, url)
            }
        }
    } else {
        let send_result = tx.send(Err(anyhow!("failed to get response")))
            .await;
 
        debug!("download image error {:?}  for url {}", send_result, url)
    }
}
 
#[tokio::main]
async fn main() -> Result<()> {
    env_logger::init();
 
    let task_timeout = Duration::from_millis(60000);
 
    let (tx, mut rx) = mpsc::channel(32);
    let (cancel_tx, mut cancel_rx) = oneshot::channel::<bool>();
    let (done_tx, mut done_rx) = oneshot::channel::<bool>();
    let signal_wait_handle = setup_graceful_shutdown(cancel_tx);
 
    let urls: Vec<&str> = vec![
        "https://images.unsplash.com/photo-1739117956532-8a37c76093dd?fm=jpg",
        "https://images.unsplash.com/photo-1735930371721-d1243e8122ab?fm=jpg",
        "https://images.unsplash.com/photo-1732204662871-f853329ce638?fm=jpg",
        "https://images.unsplash.com/photo-1712677925320-5f93d7f853f5?fm=jpg",
        "https://images.unsplash.com/photo-1712677925900-12574cb3cf15?fm=jpg",
        "https://images.unsplash.com/photo-1710436000845-bb707af976a6?fm=jpg",
        "https://images.unsplash.com/photo-1700386277812-5ff74eb03650?fm=jpg",
        "https://images.unsplash.com/photo-1700386277415-c9c5270ee8b8?fm=jpg",
        "https://images.unsplash.com/photo-1691947563165-28011f40d786?fm=jpg",
        "https://images.unsplash.com/photo-1691782359338-97fcbcce8ce8?fm=jpg",
    ];
 
    let mut handles: Vec<JoinHandle<()>> = Vec::new();
    for url in urls {
        let tx = tx.clone();
        let handle = tokio::spawn(async move {
            let timeout_status = timeout(task_timeout, download_image(tx, url)).await;
            debug!("timeout status = {:?} for url {}", timeout_status, url)
        });
        handles.push(handle);
    }
 
    let wait_task = tokio::spawn(async move {
        for handle in handles {
            let _ = handle.await;
        }
        drop(tx);
    });
 
    let write_task = tokio::spawn(async move {
        let mut counter = 1;
        while let Some(message) = rx.recv().await {
            if let Ok(image_bytes) = message {
                let result = timeout(
                    task_timeout,
                    tokio::spawn(async move {
                        let file = File::create(format!("{}.jpg", counter)).await;
                        if let Ok(mut file) = file {
                            let _ = file.write_all(&image_bytes).await;
                            info!("Wrote image {}", counter);
                        } else {
                            error!("Failed to write image {}", counter)
                        }
                    }),
                )
                .await;
 
                debug!("write task timeout result {:?}", task_timeout);
                counter += 1;
            }
        }
        let _ = done_tx.send(true);
    });
 
    tokio::select! {
        _ = cancel_rx => {
            info!("computation cancelled");
            wait_task.abort();
            write_task.abort();
        }
        _ = done_rx => {
            info!("operation completed");
        }
    }
    signal_wait_handle.abort();
 
    Ok(())
}
 

Comments

  •  icon
    01/01/70 12:00:00 AM UTC
    Plain Text |

    0 B

    |

    👍

    /

    👎

    
        
  •  icon
    01/01/70 12:00:00 AM UTC
    Plain Text |

    0 B

    |

    👍

    /

    👎