use std::time::Instant;
use tokio::sync::mpsc;
use crate::sources;
pub fn spawn_news_content_worker(
mut news_content_req_rx: mpsc::UnboundedReceiver<String>,
news_content_res_tx: mpsc::UnboundedSender<(String, String)>,
) {
tokio::spawn(async move {
while let Some(mut url) = news_content_req_rx.recv().await {
let mut skipped = 0usize;
while let Ok(newer_url) = news_content_req_rx.try_recv() {
skipped += 1;
url = newer_url;
}
if skipped > 0 {
tracing::debug!(
skipped,
url = %url,
"news_content_worker: drained stale requests, processing most recent"
);
}
let url_clone = url.clone();
let started = Instant::now();
tracing::info!(url = %url_clone, "news_content_worker: fetch start");
match sources::fetch_news_content(&url).await {
Ok(content) => {
tracing::debug!(
url = %url_clone,
elapsed_ms = started.elapsed().as_millis(),
len = content.len(),
"news_content_worker: fetch success"
);
let _ = news_content_res_tx.send((url_clone, content));
}
Err(e) => {
tracing::warn!(
error = %e,
url = %url_clone,
elapsed_ms = started.elapsed().as_millis(),
"news_content_worker: fetch failed"
);
let _ = news_content_res_tx
.send((url_clone, format!("Failed to load content: {e}")));
}
}
}
});
}
#[cfg(test)]
mod tests {
#[test]
fn test_news_content_worker_error_format() {
let error = "Network error";
let error_msg = format!("Failed to load content: {error}");
assert!(error_msg.contains("Failed to load content"));
assert!(error_msg.contains(error));
}
}