| 1 | //! `dogfood-worker` — the job runner (spec §5). |
| 2 | //! |
| 3 | //! Jobs are claimed with `SELECT … FOR UPDATE SKIP LOCKED`. A `NOTIFY` on |
| 4 | //! insert wakes the worker so it does not poll on the hot path; a 5-second poll |
| 5 | //! remains as a fallback for missed notifications. |
| 6 | |
| 7 | use std::sync::Arc; |
| 8 | use std::time::Duration; |
| 9 | |
| 10 | use anyhow::{Context, Result}; |
| 11 | use df_store::RepoStore; |
| 12 | use sqlx::postgres::PgListener; |
| 13 | use sqlx::PgPool; |
| 14 | use uuid::Uuid; |
| 15 | |
| 16 | mod anchors; |
| 17 | mod index_push; |
| 18 | mod patch_id; |
| 19 | |
| 20 | /// Identifies this worker in `jobs.locked_by`, so a stuck job can be traced |
| 21 | /// back to the process that took it. |
| 22 | fn worker_id() -> String { |
| 23 | format!( |
| 24 | "{}-{}", |
| 25 | std::env::var("HOSTNAME").unwrap_or_else(|_| "worker".into()), |
| 26 | std::process::id() |
| 27 | ) |
| 28 | } |
| 29 | |
| 30 | #[tokio::main] |
| 31 | async fn main() -> Result<()> { |
| 32 | let _ = dotenvy::dotenv(); |
| 33 | init_tracing(); |
| 34 | |
| 35 | let database_url = std::env::var("DATABASE_URL").context("DATABASE_URL must be set")?; |
| 36 | let repo_root = std::env::var("REPO_ROOT").unwrap_or_else(|_| "/srv/repos".into()); |
| 37 | let max_conn = std::env::var("DATABASE_MAX_CONNECTIONS") |
| 38 | .ok() |
| 39 | .and_then(|v| v.parse().ok()) |
| 40 | .unwrap_or(5); |
| 41 | |
| 42 | let db = df_db::connect(&database_url, max_conn) |
| 43 | .await |
| 44 | .context("connecting to the database")?; |
| 45 | |
| 46 | // The worker does NOT run migrations — `web` owns that (spec §10), and two |
| 47 | // processes racing to migrate is exactly what the advisory lock exists to |
| 48 | // prevent. If the schema is behind, the queries below fail loudly. |
| 49 | let store: Arc<dyn RepoStore> = Arc::new(df_store::GitStore::new(&repo_root)); |
| 50 | |
| 51 | let me = worker_id(); |
| 52 | tracing::info!(worker = %me, root = %repo_root, "dogfood-worker started"); |
| 53 | |
| 54 | // LISTEN for wake-ups. If the listener cannot be established the worker |
| 55 | // still functions on the poll fallback, so this is a warning, not fatal. |
| 56 | let mut listener = match PgListener::connect(&database_url).await { |
| 57 | Ok(mut l) => match l.listen("dogfood_jobs").await { |
| 58 | Ok(()) => Some(l), |
| 59 | Err(e) => { |
| 60 | tracing::warn!("LISTEN failed, falling back to polling: {e}"); |
| 61 | None |
| 62 | } |
| 63 | }, |
| 64 | Err(e) => { |
| 65 | tracing::warn!("could not open a listener, falling back to polling: {e}"); |
| 66 | None |
| 67 | } |
| 68 | }; |
| 69 | |
| 70 | let mut shutdown = std::pin::pin!(shutdown_signal()); |
| 71 | |
| 72 | loop { |
| 73 | // Drain everything currently runnable before sleeping. |
| 74 | loop { |
| 75 | match claim_and_run(&db, store.as_ref(), &me).await { |
| 76 | Ok(true) => continue, |
| 77 | Ok(false) => break, |
| 78 | Err(e) => { |
| 79 | tracing::error!("job loop error: {e:#}"); |
| 80 | break; |
| 81 | } |
| 82 | } |
| 83 | } |
| 84 | |
| 85 | // Wait for a notification, the poll interval, or shutdown. |
| 86 | let wake = async { |
| 87 | match listener.as_mut() { |
| 88 | Some(l) => { |
| 89 | let _ = l.recv().await; |
| 90 | } |
| 91 | None => std::future::pending::<()>().await, |
| 92 | } |
| 93 | }; |
| 94 | |
| 95 | tokio::select! { |
| 96 | _ = wake => {} |
| 97 | _ = tokio::time::sleep(Duration::from_secs(5)) => {} |
| 98 | _ = &mut shutdown => { |
| 99 | tracing::info!("shutting down"); |
| 100 | return Ok(()); |
| 101 | } |
| 102 | } |
| 103 | } |
| 104 | } |
| 105 | |
| 106 | /// Claim one job and run it. Returns whether a job was found. |
| 107 | async fn claim_and_run(db: &PgPool, store: &dyn RepoStore, me: &str) -> Result<bool> { |
| 108 | let mut tx = db.begin().await?; |
| 109 | |
| 110 | // SKIP LOCKED lets several workers share the queue without blocking. |
| 111 | let job: Option<(Uuid, String, serde_json::Value, i32, i32)> = sqlx::query_as( |
| 112 | "SELECT id, kind, payload, attempts, max_attempts |
| 113 | FROM jobs |
| 114 | WHERE locked_at IS NULL AND run_at <= now() |
| 115 | ORDER BY run_at |
| 116 | FOR UPDATE SKIP LOCKED |
| 117 | LIMIT 1", |
| 118 | ) |
| 119 | .fetch_optional(&mut *tx) |
| 120 | .await?; |
| 121 | |
| 122 | let Some((id, kind, payload, attempts, max_attempts)) = job else { |
| 123 | tx.rollback().await?; |
| 124 | return Ok(false); |
| 125 | }; |
| 126 | |
| 127 | sqlx::query("UPDATE jobs SET locked_at = now(), locked_by = $2, attempts = attempts + 1 WHERE id = $1") |
| 128 | .bind(id) |
| 129 | .bind(me) |
| 130 | .execute(&mut *tx) |
| 131 | .await?; |
| 132 | tx.commit().await?; |
| 133 | |
| 134 | tracing::info!(job = %id, %kind, attempt = attempts + 1, "running job"); |
| 135 | let started = std::time::Instant::now(); |
| 136 | |
| 137 | let result = match kind.as_str() { |
| 138 | "index_push" => run_index_push(db, store, &payload).await, |
| 139 | other => Err(anyhow::anyhow!("unknown job kind: {other}")), |
| 140 | }; |
| 141 | |
| 142 | match result { |
| 143 | Ok(()) => { |
| 144 | sqlx::query("DELETE FROM jobs WHERE id = $1") |
| 145 | .bind(id) |
| 146 | .execute(db) |
| 147 | .await?; |
| 148 | tracing::info!(job = %id, ms = started.elapsed().as_millis(), "job complete"); |
| 149 | } |
| 150 | Err(e) => { |
| 151 | let attempts_now = attempts + 1; |
| 152 | let give_up = attempts_now >= max_attempts; |
| 153 | tracing::error!(job = %id, attempt = attempts_now, "job failed: {e:#}"); |
| 154 | |
| 155 | if give_up { |
| 156 | // Leave the row, unlocked and exhausted, so it is visible for |
| 157 | // inspection rather than silently vanishing. |
| 158 | sqlx::query( |
| 159 | "UPDATE jobs SET locked_at = NULL, locked_by = NULL, last_error = $2 |
| 160 | WHERE id = $1", |
| 161 | ) |
| 162 | .bind(id) |
| 163 | .bind(format!("{e:#}")) |
| 164 | .execute(db) |
| 165 | .await?; |
| 166 | tracing::error!(job = %id, "job exhausted its attempts"); |
| 167 | } else { |
| 168 | // Exponential backoff, capped. |
| 169 | let delay = Duration::from_secs(2u64.saturating_pow(attempts_now as u32).min(300)); |
| 170 | sqlx::query( |
| 171 | "UPDATE jobs SET locked_at = NULL, locked_by = NULL, last_error = $2, |
| 172 | run_at = now() + $3 |
| 173 | WHERE id = $1", |
| 174 | ) |
| 175 | .bind(id) |
| 176 | .bind(format!("{e:#}")) |
| 177 | .bind(delay) |
| 178 | .execute(db) |
| 179 | .await?; |
| 180 | } |
| 181 | } |
| 182 | } |
| 183 | |
| 184 | Ok(true) |
| 185 | } |
| 186 | |
| 187 | async fn run_index_push( |
| 188 | db: &PgPool, |
| 189 | store: &dyn RepoStore, |
| 190 | payload: &serde_json::Value, |
| 191 | ) -> Result<()> { |
| 192 | let repo_id: Uuid = payload |
| 193 | .get("repo_id") |
| 194 | .and_then(|v| v.as_str()) |
| 195 | .and_then(|s| s.parse().ok()) |
| 196 | .context("index_push payload is missing a valid repo_id")?; |
| 197 | |
| 198 | let pushed_by: Option<Uuid> = payload |
| 199 | .get("pushed_by") |
| 200 | .and_then(|v| v.as_str()) |
| 201 | .and_then(|s| s.parse().ok()); |
| 202 | |
| 203 | let outcome = index_push::run(db, store, repo_id, pushed_by).await?; |
| 204 | tracing::info!( |
| 205 | repo = %repo_id, |
| 206 | changes = outcome.changes_seen, |
| 207 | revisions = outcome.revisions_added, |
| 208 | "indexed push" |
| 209 | ); |
| 210 | Ok(()) |
| 211 | } |
| 212 | |
| 213 | fn init_tracing() { |
| 214 | use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt, EnvFilter}; |
| 215 | |
| 216 | let filter = EnvFilter::try_from_default_env() |
| 217 | .unwrap_or_else(|_| EnvFilter::new("info,df_worker=debug,df_index=debug")); |
| 218 | let json = !std::io::IsTerminal::is_terminal(&std::io::stdout()); |
| 219 | |
| 220 | let registry = tracing_subscriber::registry().with(filter); |
| 221 | if json { |
| 222 | registry.with(tracing_subscriber::fmt::layer().json()).init(); |
| 223 | } else { |
| 224 | registry.with(tracing_subscriber::fmt::layer()).init(); |
| 225 | } |
| 226 | } |
| 227 | |
| 228 | async fn shutdown_signal() { |
| 229 | let ctrl_c = async { |
| 230 | tokio::signal::ctrl_c().await.expect("Ctrl+C handler"); |
| 231 | }; |
| 232 | #[cfg(unix)] |
| 233 | let terminate = async { |
| 234 | tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate()) |
| 235 | .expect("SIGTERM handler") |
| 236 | .recv() |
| 237 | .await; |
| 238 | }; |
| 239 | #[cfg(not(unix))] |
| 240 | let terminate = std::future::pending::<()>(); |
| 241 | |
| 242 | tokio::select! { |
| 243 | _ = ctrl_c => {} |
| 244 | _ = terminate => {} |
| 245 | } |
| 246 | } |
246 lines · Rust