char/sorcery

static-files based git repo viewer

git clone https://git.t4t.associates/char/sorcery

Charlotte Somsimplify + split www daemon3665650

main
9.9 KiB306 linesraw
1use std::collections::{HashMap, HashSet};
2use std::fs;
3use std::path::{Path, PathBuf};
4use std::sync::Arc;
5use std::time::Instant;
6
7use anyhow::{Context, Result};
8use axum::body::Bytes;
9use axum::extract::{Path as AxumPath, State};
10use axum::http::{HeaderMap, HeaderValue, StatusCode, header};
11use axum::response::{IntoResponse, Response};
12use tokio::sync::{Mutex, watch};
13
14use crate::catalog::{self, Repository};
15use crate::generate::{self, Config as GenerateConfig, STATE_FILE};
16use crate::render::catalog_index;
17
18use super::{AppState, Inner, WebResult, internal};
19
20#[derive(Default)]
21pub(super) struct Snapshot {
22    pub(super) catalog_html: Bytes,
23    pub(super) repos: HashMap<String, RepoEntry>,
24}
25
26#[derive(Clone)]
27pub(super) struct RepoEntry {
28    pub(super) repository: Repository,
29    build_gate: Arc<Mutex<()>>,
30}
31
32pub(super) async fn scan_periodically(state: AppState) {
33    let mut tick = tokio::time::interval(state.config.check_interval);
34    tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
35    tick.tick().await; // the first tick is immediate; run() already scanned
36    loop {
37        tick.tick().await;
38        if let Err(error) = state.rescan().await {
39            eprintln!("sorceryd: rescan: {error}");
40        }
41    }
42}
43
44impl Inner {
45    /// Discover repositories and publish a fresh immutable catalog snapshot.
46    pub(super) async fn rescan(&self) -> Result<()> {
47        let config = self.config.clone();
48        let (repositories, freshness) = tokio::task::spawn_blocking(move || -> Result<_> {
49            let repositories = catalog::discover(&config.repositories)?;
50            let freshness = repositories
51                .iter()
52                .map(|repository| is_current(&generate_config(&config, repository)))
53                .collect::<Vec<_>>();
54            Ok((repositories, freshness))
55        })
56        .await??;
57
58        let old = self.snapshot.load_full();
59        let mut repos = HashMap::with_capacity(repositories.len());
60        let mut stale = Vec::new();
61        for (repository, current) in repositories.iter().zip(freshness) {
62            let current = match current {
63                Ok(current) => current,
64                Err(error) => {
65                    eprintln!(
66                        "sorceryd: skipping {}/{}: {error}",
67                        repository.user, repository.name
68                    );
69                    continue;
70                }
71            };
72            let key = format!("{}/{}", repository.user, repository.name);
73            let build_gate = old
74                .repos
75                .get(&key)
76                .map(|entry| entry.build_gate.clone())
77                .unwrap_or_else(|| Arc::new(Mutex::new(())));
78            let entry = RepoEntry {
79                repository: repository.clone(),
80                build_gate,
81            };
82            if !current {
83                stale.push(entry.clone());
84            }
85            repos.insert(key, entry);
86        }
87
88        self.snapshot.store(Arc::new(Snapshot {
89            catalog_html: catalog_index(&self.config.instance_name, &repositories).into(),
90            repos,
91        }));
92        self.background.send_replace(stale);
93        Ok(())
94    }
95
96    pub(super) async fn build(&self, entry: &RepoEntry) -> Result<bool> {
97        let _gate = entry.build_gate.lock().await;
98        let _cache = self.cache_lock.read().await;
99        let config = generate_config(&self.config, &entry.repository);
100        let highlights = self.highlights.clone();
101        tokio::task::spawn_blocking(move || ensure_current(&config, highlights)).await?
102    }
103}
104
105pub(super) async fn background_worker(state: AppState, mut stale: watch::Receiver<Vec<RepoEntry>>) {
106    while stale.changed().await.is_ok() {
107        let entries = stale.borrow_and_update().clone();
108        for entry in entries {
109            let started = Instant::now();
110            let name = format!("{}/{}", entry.repository.user, entry.repository.name);
111            match state.build(&entry).await {
112                Ok(true) => eprintln!(
113                    "sorceryd: rebuilt {name} in {:.2}s",
114                    started.elapsed().as_secs_f64(),
115                ),
116                Ok(false) => {}
117                Err(error) => eprintln!("sorceryd: background rebuild of {name}: {error}"),
118            }
119        }
120        match garbage_collect(&state).await {
121            Ok(None | Some(0)) => {}
122            Ok(Some(count)) => {
123                eprintln!("sorceryd: removed {count} stale cache directories")
124            }
125            Err(error) => eprintln!("sorceryd: cache garbage collection: {error}"),
126        }
127    }
128}
129
130async fn garbage_collect(state: &Inner) -> Result<Option<usize>> {
131    let Ok(_cache) = state.cache_lock.try_write() else {
132        return Ok(None);
133    };
134    let cache = state.config.cache.clone();
135    let live = state
136        .snapshot
137        .load()
138        .repos
139        .values()
140        .map(|entry| state.config.repo_cache(&entry.repository))
141        .collect::<HashSet<_>>();
142    tokio::task::spawn_blocking(move || garbage_collect_cache(&cache, &live))
143        .await?
144        .map(Some)
145}
146
147fn garbage_collect_cache(cache: &Path, live: &HashSet<PathBuf>) -> Result<usize> {
148    let mut removed = 0;
149    for user in fs::read_dir(cache)? {
150        let user = user?;
151        if !user.file_type()?.is_dir() {
152            continue;
153        }
154        for repo in fs::read_dir(user.path())? {
155            let repo = repo?;
156            if !repo.file_type()?.is_dir() {
157                continue;
158            }
159            let path = repo.path();
160            if path.join(STATE_FILE).is_file() && !live.contains(&path) {
161                fs::remove_dir_all(path)?;
162                removed += 1;
163            }
164        }
165    }
166    Ok(removed)
167}
168
169fn generate_config(config: &super::Config, repository: &Repository) -> GenerateConfig {
170    let name = format!("{}/{}", repository.user, repository.name);
171    GenerateConfig {
172        repo: repository.path.clone(),
173        out: config.repo_cache(repository),
174        instance_name: config.instance_name.clone(),
175        name: Some(name.clone()),
176        clone_url: config
177            .clone_url_base
178            .as_ref()
179            .map(|base| format!("{}/{name}", base.trim_end_matches('/'))),
180    }
181}
182
183fn is_current(config: &GenerateConfig) -> Result<bool> {
184    let expected = generate::expected_state(config)?;
185    Ok(fs::read_to_string(config.out.join(STATE_FILE))
186        .ok()
187        .as_deref()
188        == Some(&expected))
189}
190
191/// Regenerate the cache directory if its recorded state differs from the
192/// repository's, building into a staging sibling and swapping atomically.
193fn ensure_current(
194    config: &GenerateConfig,
195    highlights: Arc<crate::highlight::Cache>,
196) -> Result<bool> {
197    if is_current(config)? {
198        return Ok(false);
199    }
200
201    let parent = config
202        .out
203        .parent()
204        .context("cache repository has no parent")?;
205    let name = config
206        .out
207        .file_name()
208        .context("cache repository has no name")?;
209    fs::create_dir_all(parent)?;
210    let staging = parent.join(format!(
211        ".{}.new-{}",
212        name.to_string_lossy(),
213        std::process::id()
214    ));
215    if staging.exists() {
216        fs::remove_dir_all(&staging)?;
217    }
218    let mut staging_config = config.clone();
219    staging_config.out = staging.clone();
220    generate::full(&staging_config, Some(&config.out), highlights)?;
221    // keep info/refs + objects/info/packs current for static git-dir access
222    let status = std::process::Command::new("git")
223        .arg("-C")
224        .arg(&config.repo)
225        .arg("update-server-info")
226        .status()?;
227    if !status.success() {
228        eprintln!(
229            "sorceryd: update-server-info failed for {}",
230            config.repo.display()
231        );
232    }
233    generate::swap_dir(&staging, &config.out)?;
234    Ok(true)
235}
236
237pub(super) async fn refresh(
238    State(state): State<AppState>,
239    AxumPath((user, repo)): AxumPath<(String, String)>,
240    headers: HeaderMap,
241) -> WebResult<Response> {
242    let Some(expected) = state.refresh_token.as_deref() else {
243        return Err((StatusCode::NOT_FOUND, "not found".into()));
244    };
245    let authorized = headers
246        .get(header::AUTHORIZATION)
247        .map(HeaderValue::as_bytes)
248        .and_then(|value| {
249            let (scheme, token) = value.split_at_checked(7)?;
250            scheme.eq_ignore_ascii_case(b"bearer ").then_some(token)
251        })
252        .is_some_and(|provided| constant_time_eq(expected, provided));
253    if !authorized {
254        return Err((StatusCode::UNAUTHORIZED, "unauthorized".into()));
255    }
256
257    state.rescan().await.map_err(internal)?;
258    let entry = state.resolve(&user, &repo)?;
259    state.build(&entry).await.map_err(internal)?;
260    Ok((StatusCode::OK, "refreshed\n").into_response())
261}
262
263fn constant_time_eq(expected: &[u8], provided: &[u8]) -> bool {
264    let mut difference = expected.len() ^ provided.len();
265    for (i, byte) in expected.iter().enumerate() {
266        difference |= usize::from(*byte ^ provided.get(i).copied().unwrap_or(0));
267    }
268    difference == 0
269}
270
271#[cfg(test)]
272mod tests {
273    use std::collections::HashSet;
274    use std::fs;
275
276    use anyhow::Result;
277
278    use crate::generate::STATE_FILE;
279    use crate::testutil::TempDir;
280
281    use super::garbage_collect_cache;
282
283    #[test]
284    fn garbage_collects_only_inactive_site_caches() -> Result<()> {
285        let root = TempDir::new("cache-gc");
286        let live = root.join("alice/live");
287        let stale = root.join("alice/stale");
288        let staging = root.join("alice/.live.new-1");
289        let asset = root.join("fonts/inter");
290        for path in [&live, &stale, &staging] {
291            fs::create_dir_all(path)?;
292            fs::write(path.join(STATE_FILE), "state")?;
293        }
294        fs::create_dir_all(&asset)?;
295
296        assert_eq!(
297            garbage_collect_cache(&root, &HashSet::from([live.clone()]))?,
298            2
299        );
300        assert!(live.is_dir());
301        assert!(!stale.exists());
302        assert!(!staging.exists());
303        assert!(asset.is_dir());
304        Ok(())
305    }
306}