use std::collections::HashSet;
use std::hash::BuildHasher;
use std::time::{Duration, Instant};
use crate::state::types::{NewsFeedItem, NewsFeedSource};
use tracing::{info, warn};
use super::Result;
use super::cache::{
CACHE_TTL_SECONDS, CacheEntry, NEWS_CACHE, load_from_disk_cache, save_to_disk_cache,
};
use super::rate_limit::{
extract_retry_after_from_error, increase_archlinux_backoff, rate_limit, rate_limit_archlinux,
reset_archlinux_backoff, retry_with_backoff, set_network_error,
};
pub(super) async fn append_arch_news(
limit: usize,
cutoff_date: Option<&str>,
) -> Result<Vec<NewsFeedItem>> {
const SOURCE: &str = "arch_news";
if cutoff_date.is_none()
&& let Ok(cache) = NEWS_CACHE.lock()
&& let Some(entry) = cache.get(SOURCE)
&& entry.timestamp.elapsed().as_secs() < CACHE_TTL_SECONDS
{
info!("using in-memory cached arch news");
return Ok(entry.data.clone());
}
if cutoff_date.is_none()
&& let Some(disk_data) = load_from_disk_cache(SOURCE)
{
if let Ok(mut cache) = NEWS_CACHE.lock() {
cache.insert(
SOURCE.to_string(),
CacheEntry {
data: disk_data.clone(),
timestamp: Instant::now(),
},
);
}
return Ok(disk_data);
}
let fetch_result: std::result::Result<
Vec<crate::state::types::NewsItem>,
Box<dyn std::error::Error + Send + Sync>,
> = {
let _permit = rate_limit_archlinux().await;
let result = crate::sources::fetch_arch_news(limit, cutoff_date).await;
if let Err(ref e) = result {
let error_str = e.to_string();
let retry_after_seconds = extract_retry_after_from_error(&error_str);
if error_str.contains("429") || error_str.contains("503") {
if let Some(retry_after) = retry_after_seconds {
warn!(
retry_after_seconds = retry_after,
"HTTP {} detected, noting Retry-After for future requests",
if error_str.contains("429") {
"429"
} else {
"503"
}
);
increase_archlinux_backoff(Some(retry_after));
} else {
warn!(
"HTTP {} detected, increasing backoff for future requests",
if error_str.contains("429") {
"429"
} else {
"503"
}
);
increase_archlinux_backoff(None);
}
} else {
increase_archlinux_backoff(None);
}
}
result
};
match fetch_result {
Ok(news) => {
reset_archlinux_backoff();
let items: Vec<NewsFeedItem> = news
.into_iter()
.map(|n| NewsFeedItem {
id: n.url.clone(),
date: n.date,
title: n.title,
summary: None,
url: Some(n.url),
source: NewsFeedSource::ArchNews,
severity: None,
packages: Vec::new(),
})
.collect();
if cutoff_date.is_none() {
if let Ok(mut cache) = NEWS_CACHE.lock() {
cache.insert(
SOURCE.to_string(),
CacheEntry {
data: items.clone(),
timestamp: Instant::now(),
},
);
}
save_to_disk_cache(SOURCE, &items);
}
Ok(items)
}
Err(e) => {
warn!(error = %e, "arch news fetch failed");
set_network_error();
increase_archlinux_backoff(None);
if let Ok(cache) = NEWS_CACHE.lock()
&& let Some(entry) = cache.get(SOURCE)
{
info!(
cached_items = entry.data.len(),
age_secs = entry.timestamp.elapsed().as_secs(),
"using stale in-memory cached arch news due to fetch failure"
);
return Ok(entry.data.clone());
}
if let Some(disk_data) = load_from_disk_cache(SOURCE) {
info!(
cached_items = disk_data.len(),
"using disk cached arch news due to fetch failure"
);
return Ok(disk_data);
}
Err(e)
}
}
}
pub(super) async fn append_advisories<S>(
limit: usize,
installed_filter: Option<&HashSet<String, S>>,
installed_only: bool,
cutoff_date: Option<&str>,
) -> Result<Vec<NewsFeedItem>>
where
S: BuildHasher + Send + Sync + 'static,
{
const SOURCE: &str = "advisories";
if cutoff_date.is_none()
&& !installed_only
&& let Ok(cache) = NEWS_CACHE.lock()
&& let Some(entry) = cache.get(SOURCE)
&& entry.timestamp.elapsed().as_secs() < CACHE_TTL_SECONDS
{
info!("using in-memory cached advisories");
return Ok(entry.data.clone());
}
if cutoff_date.is_none()
&& !installed_only
&& let Some(disk_data) = load_from_disk_cache(SOURCE)
{
if let Ok(mut cache) = NEWS_CACHE.lock() {
cache.insert(
SOURCE.to_string(),
CacheEntry {
data: disk_data.clone(),
timestamp: Instant::now(),
},
);
}
return Ok(disk_data);
}
rate_limit().await;
let fetch_result = retry_with_backoff(
|| async {
rate_limit().await;
crate::sources::fetch_security_advisories(limit, cutoff_date).await
},
2, )
.await;
match fetch_result {
Ok(advisories) => {
let mut filtered = Vec::new();
for adv in advisories {
if installed_only
&& let Some(set) = installed_filter
&& !adv.packages.iter().any(|p| set.contains(p))
{
continue;
}
filtered.push(adv);
}
if cutoff_date.is_none() && !installed_only {
if let Ok(mut cache) = NEWS_CACHE.lock() {
cache.insert(
SOURCE.to_string(),
CacheEntry {
data: filtered.clone(),
timestamp: Instant::now(),
},
);
}
save_to_disk_cache(SOURCE, &filtered);
}
Ok(filtered)
}
Err(e) => {
warn!(error = %e, "security advisories fetch failed");
set_network_error();
if let Ok(cache) = NEWS_CACHE.lock()
&& let Some(entry) = cache.get(SOURCE)
{
info!(
cached_items = entry.data.len(),
age_secs = entry.timestamp.elapsed().as_secs(),
"using stale in-memory cached advisories due to fetch failure"
);
return Ok(entry.data.clone());
}
if let Some(disk_data) = load_from_disk_cache(SOURCE) {
info!(
cached_items = disk_data.len(),
"using disk cached advisories due to fetch failure"
);
return Ok(disk_data);
}
Err(e)
}
}
}
pub(super) async fn fetch_slow_sources<HS>(
include_arch_news: bool,
include_advisories: bool,
limit: usize,
installed_filter: Option<&HashSet<String, HS>>,
installed_only: bool,
cutoff_date: Option<&str>,
) -> (
std::result::Result<Vec<NewsFeedItem>, Box<dyn std::error::Error + Send + Sync>>,
std::result::Result<Vec<NewsFeedItem>, Box<dyn std::error::Error + Send + Sync>>,
)
where
HS: BuildHasher + Send + Sync + 'static,
{
let arch_result: std::result::Result<
Vec<NewsFeedItem>,
Box<dyn std::error::Error + Send + Sync>,
> = if include_arch_news {
info!("fetching arch news...");
tokio::time::timeout(
Duration::from_secs(30),
append_arch_news(limit, cutoff_date),
)
.await
.map_or_else(
|_| {
warn!("arch news fetch timed out after 30s, continuing without arch news");
Err("Arch news fetch timeout".into())
},
|result| {
info!(
"arch news fetch completed: items={}",
result.as_ref().map(Vec::len).unwrap_or(0)
);
result
},
)
} else {
Ok(Vec::new())
};
let advisories_result: std::result::Result<
Vec<NewsFeedItem>,
Box<dyn std::error::Error + Send + Sync>,
> = if include_advisories {
info!("fetching advisories...");
tokio::time::timeout(
Duration::from_secs(30),
append_advisories(limit, installed_filter, installed_only, cutoff_date),
)
.await
.map_or_else(
|_| {
warn!("advisories fetch timed out after 30s, continuing without advisories");
Err("Advisories fetch timeout".into())
},
|result| {
info!(
"advisories fetch completed: items={}",
result.as_ref().map(Vec::len).unwrap_or(0)
);
result
},
)
} else {
Ok(Vec::new())
};
(arch_result, advisories_result)
}