Feeds were already persisted, but nothing ever re-fetched them after the initial subscribe — scoring.rs had a TODO where feed polling was supposed to go. Add a feed_polling job that re-fetches every subscribed feed, relies on articles.url's unique constraint to skip ones already seen, and stamps last_fetched_at. Runs once on startup (so reopening the app catches up immediately) and every 15 minutes after, ahead of the scoring pass. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01QmwX8eafbMstJvt8XPSqft
74 lines
2.2 KiB
Rust
74 lines
2.2 KiB
Rust
mod affinity_decay;
|
|
mod feed_polling;
|
|
mod scoring;
|
|
|
|
use anyhow::Result;
|
|
use feedsignal_db::Db;
|
|
use feedsignal_llm::Llm;
|
|
use std::sync::Arc;
|
|
use tokio_cron_scheduler::{Job, JobScheduler};
|
|
|
|
/// Wires up the recurring jobs described in the design discussion: polling
|
|
/// subscribed feeds for new articles, scoring them, and decaying topic
|
|
/// affinities once a day so stale signals fade. Call once at server
|
|
/// startup, and also runs once immediately so a freshly (re)opened app
|
|
/// doesn't wait 15 minutes for its first check.
|
|
pub async fn start_scheduler(db: Arc<Db>, llm: Arc<Llm>) -> Result<JobScheduler> {
|
|
let scheduler = JobScheduler::new().await?;
|
|
|
|
{
|
|
let db = db.clone();
|
|
tokio::spawn(async move {
|
|
if let Err(err) = feed_polling::run(&db).await {
|
|
tracing::error!(?err, "initial feed poll failed");
|
|
}
|
|
});
|
|
}
|
|
|
|
{
|
|
let db = db.clone();
|
|
scheduler
|
|
.add(Job::new_async("0 */15 * * * *", move |_uuid, _lock| {
|
|
let db = db.clone();
|
|
Box::pin(async move {
|
|
if let Err(err) = feed_polling::run(&db).await {
|
|
tracing::error!(?err, "feed polling run failed");
|
|
}
|
|
})
|
|
})?)
|
|
.await?;
|
|
}
|
|
|
|
{
|
|
let db = db.clone();
|
|
let llm = llm.clone();
|
|
scheduler
|
|
.add(Job::new_async("0 */15 * * * *", move |_uuid, _lock| {
|
|
let db = db.clone();
|
|
let llm = llm.clone();
|
|
Box::pin(async move {
|
|
if let Err(err) = scoring::run(&db, &llm).await {
|
|
tracing::error!(?err, "scoring pipeline run failed");
|
|
}
|
|
})
|
|
})?)
|
|
.await?;
|
|
}
|
|
|
|
{
|
|
let db = db.clone();
|
|
scheduler
|
|
.add(Job::new_async("0 0 4 * * *", move |_uuid, _lock| {
|
|
let db = db.clone();
|
|
Box::pin(async move {
|
|
if let Err(err) = affinity_decay::run(&db).await {
|
|
tracing::error!(?err, "affinity decay run failed");
|
|
}
|
|
})
|
|
})?)
|
|
.await?;
|
|
}
|
|
|
|
scheduler.start().await?;
|
|
Ok(scheduler)
|
|
}
|