use std::collections::{HashMap, HashSet}; use std::fs; use std::path::{Path, PathBuf}; use std::sync::Arc; use std::time::Instant; use anyhow::{Context, Result}; use axum::body::Bytes; use axum::extract::{Path as AxumPath, State}; use axum::http::{HeaderMap, HeaderValue, StatusCode, header}; use axum::response::{IntoResponse, Response}; use tokio::sync::{Mutex, watch}; use crate::catalog::{self, Repository}; use crate::generate::{self, Config as GenerateConfig, STATE_FILE}; use crate::render::catalog_index; use super::{AppState, Inner, WebResult, internal}; #[derive(Default)] pub(super) struct Snapshot { pub(super) catalog_html: Bytes, pub(super) repos: HashMap, } #[derive(Clone)] pub(super) struct RepoEntry { pub(super) repository: Repository, build_gate: Arc>, } pub(super) async fn scan_periodically(state: AppState) { let mut tick = tokio::time::interval(state.config.check_interval); tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); tick.tick().await; // the first tick is immediate; run() already scanned loop { tick.tick().await; if let Err(error) = state.rescan().await { eprintln!("sorceryd: rescan: {error}"); } } } impl Inner { /// Discover repositories and publish a fresh immutable catalog snapshot. pub(super) async fn rescan(&self) -> Result<()> { let config = self.config.clone(); let (repositories, freshness) = tokio::task::spawn_blocking(move || -> Result<_> { let repositories = catalog::discover(&config.repositories)?; let freshness = repositories .iter() .map(|repository| is_current(&generate_config(&config, repository))) .collect::>(); Ok((repositories, freshness)) }) .await??; let old = self.snapshot.load_full(); let mut repos = HashMap::with_capacity(repositories.len()); let mut stale = Vec::new(); for (repository, current) in repositories.iter().zip(freshness) { let current = match current { Ok(current) => current, Err(error) => { eprintln!( "sorceryd: skipping {}/{}: {error}", repository.user, repository.name ); continue; } }; let key = format!("{}/{}", repository.user, repository.name); let build_gate = old .repos .get(&key) .map(|entry| entry.build_gate.clone()) .unwrap_or_else(|| Arc::new(Mutex::new(()))); let entry = RepoEntry { repository: repository.clone(), build_gate, }; if !current { stale.push(entry.clone()); } repos.insert(key, entry); } self.snapshot.store(Arc::new(Snapshot { catalog_html: catalog_index(&self.config.instance_name, &repositories).into(), repos, })); self.background.send_replace(stale); Ok(()) } pub(super) async fn build(&self, entry: &RepoEntry) -> Result { let _gate = entry.build_gate.lock().await; let _cache = self.cache_lock.read().await; let config = generate_config(&self.config, &entry.repository); let highlights = self.highlights.clone(); tokio::task::spawn_blocking(move || ensure_current(&config, highlights)).await? } } pub(super) async fn background_worker(state: AppState, mut stale: watch::Receiver>) { while stale.changed().await.is_ok() { let entries = stale.borrow_and_update().clone(); for entry in entries { let started = Instant::now(); let name = format!("{}/{}", entry.repository.user, entry.repository.name); match state.build(&entry).await { Ok(true) => eprintln!( "sorceryd: rebuilt {name} in {:.2}s", started.elapsed().as_secs_f64(), ), Ok(false) => {} Err(error) => eprintln!("sorceryd: background rebuild of {name}: {error}"), } } match garbage_collect(&state).await { Ok(None | Some(0)) => {} Ok(Some(count)) => { eprintln!("sorceryd: removed {count} stale cache directories") } Err(error) => eprintln!("sorceryd: cache garbage collection: {error}"), } } } async fn garbage_collect(state: &Inner) -> Result> { let Ok(_cache) = state.cache_lock.try_write() else { return Ok(None); }; let cache = state.config.cache.clone(); let live = state .snapshot .load() .repos .values() .map(|entry| state.config.repo_cache(&entry.repository)) .collect::>(); tokio::task::spawn_blocking(move || garbage_collect_cache(&cache, &live)) .await? .map(Some) } fn garbage_collect_cache(cache: &Path, live: &HashSet) -> Result { let mut removed = 0; for user in fs::read_dir(cache)? { let user = user?; if !user.file_type()?.is_dir() { continue; } for repo in fs::read_dir(user.path())? { let repo = repo?; if !repo.file_type()?.is_dir() { continue; } let path = repo.path(); if path.join(STATE_FILE).is_file() && !live.contains(&path) { fs::remove_dir_all(path)?; removed += 1; } } } Ok(removed) } fn generate_config(config: &super::Config, repository: &Repository) -> GenerateConfig { let name = format!("{}/{}", repository.user, repository.name); GenerateConfig { repo: repository.path.clone(), out: config.repo_cache(repository), instance_name: config.instance_name.clone(), name: Some(name.clone()), clone_url: config .clone_url_base .as_ref() .map(|base| format!("{}/{name}", base.trim_end_matches('/'))), } } fn is_current(config: &GenerateConfig) -> Result { let expected = generate::expected_state(config)?; Ok(fs::read_to_string(config.out.join(STATE_FILE)) .ok() .as_deref() == Some(&expected)) } /// Regenerate the cache directory if its recorded state differs from the /// repository's, building into a staging sibling and swapping atomically. fn ensure_current( config: &GenerateConfig, highlights: Arc, ) -> Result { if is_current(config)? { return Ok(false); } let parent = config .out .parent() .context("cache repository has no parent")?; let name = config .out .file_name() .context("cache repository has no name")?; fs::create_dir_all(parent)?; let staging = parent.join(format!( ".{}.new-{}", name.to_string_lossy(), std::process::id() )); if staging.exists() { fs::remove_dir_all(&staging)?; } let mut staging_config = config.clone(); staging_config.out = staging.clone(); generate::full(&staging_config, Some(&config.out), highlights)?; // keep info/refs + objects/info/packs current for static git-dir access let status = std::process::Command::new("git") .arg("-C") .arg(&config.repo) .arg("update-server-info") .status()?; if !status.success() { eprintln!( "sorceryd: update-server-info failed for {}", config.repo.display() ); } generate::swap_dir(&staging, &config.out)?; Ok(true) } pub(super) async fn refresh( State(state): State, AxumPath((user, repo)): AxumPath<(String, String)>, headers: HeaderMap, ) -> WebResult { let Some(expected) = state.refresh_token.as_deref() else { return Err((StatusCode::NOT_FOUND, "not found".into())); }; let authorized = headers .get(header::AUTHORIZATION) .map(HeaderValue::as_bytes) .and_then(|value| { let (scheme, token) = value.split_at_checked(7)?; scheme.eq_ignore_ascii_case(b"bearer ").then_some(token) }) .is_some_and(|provided| constant_time_eq(expected, provided)); if !authorized { return Err((StatusCode::UNAUTHORIZED, "unauthorized".into())); } state.rescan().await.map_err(internal)?; let entry = state.resolve(&user, &repo)?; state.build(&entry).await.map_err(internal)?; Ok((StatusCode::OK, "refreshed\n").into_response()) } fn constant_time_eq(expected: &[u8], provided: &[u8]) -> bool { let mut difference = expected.len() ^ provided.len(); for (i, byte) in expected.iter().enumerate() { difference |= usize::from(*byte ^ provided.get(i).copied().unwrap_or(0)); } difference == 0 } #[cfg(test)] mod tests { use std::collections::HashSet; use std::fs; use anyhow::Result; use crate::generate::STATE_FILE; use crate::testutil::TempDir; use super::garbage_collect_cache; #[test] fn garbage_collects_only_inactive_site_caches() -> Result<()> { let root = TempDir::new("cache-gc"); let live = root.join("alice/live"); let stale = root.join("alice/stale"); let staging = root.join("alice/.live.new-1"); let asset = root.join("fonts/inter"); for path in [&live, &stale, &staging] { fs::create_dir_all(path)?; fs::write(path.join(STATE_FILE), "state")?; } fs::create_dir_all(&asset)?; assert_eq!( garbage_collect_cache(&root, &HashSet::from([live.clone()]))?, 2 ); assert!(live.is_dir()); assert!(!stale.exists()); assert!(!staging.exists()); assert!(asset.is_dir()); Ok(()) } }