use futures::stream::StreamExt;
use lunaris_core::{Episode, LunarisError, Scope, StoragePort, keyspace};
pub async fn recent_by_source(
storage: &dyn StoragePort,
scope: &Scope,
prefixes: &[String],
limit: usize,
) -> Result<Vec<Episode>, LunarisError> {
if limit == 0 {
return Ok(Vec::new());
}
let prefix = keyspace::episode_prefix(scope);
let mut stream =
storage.scan_range(scope, &prefix, None).await.map_err(LunarisError::Storage)?;
let mut matched: Vec<Episode> = Vec::new();
while let Some(item) = stream.next().await {
let (_key, value) = item.map_err(LunarisError::Storage)?;
let episode: Episode = match serde_json::from_slice(&value) {
Ok(ep) => ep,
Err(_) => continue,
};
if prefixes.iter().any(|p| episode.source.starts_with(p.as_str())) {
matched.push(episode);
}
}
matched.sort_by(|a, b| b.id.cmp(&a.id));
matched.truncate(limit);
Ok(matched)
}