diff options
Diffstat
| -rw-r--r-- | Cargo.lock | 1 | +1 −0 |
| -rw-r--r-- | Cargo.toml | 1 | +1 −0 |
| -rw-r--r-- | src/backends/go.rs | 526 | +171 −355 |
| -rw-r--r-- | src/backends/mod.rs | 487 | +441 −46 |
| -rw-r--r-- | src/backends/zig.rs | 653 | +240 −413 |
| -rw-r--r-- | src/controller_backend.rs | 26 | +13 −13 |
| -rw-r--r-- | src/controller_web.rs | 42 | +17 −25 |
| -rw-r--r-- | src/main.rs | 66 | +23 −43 |
| -rw-r--r-- | src/proxy.rs | 23 | +20 −3 |
| -rw-r--r-- | src/storage.rs | 430 | +404 −26 |
10 files changed, 1331 insertions, 924 deletions
diff --git a/Cargo.lock b/Cargo.lock index ee311d6..348a0d4 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -5844,6 +5844,7 @@ dependencies = [ "serde", "serde_json", "sqlx", + "strum", "tempfile", "thiserror 2.0.18", "tokio", diff --git a/Cargo.toml b/Cargo.toml index da0ac9b..d23a462 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -34,6 +34,7 @@ rustls = { version = "0.23", default-features = false, features = [ semver = "1.0" serde = { version = "1.0", features = ["derive"] } serde_json = "1.0" +strum = { version = "0.27", features = ["derive"] } sqlx = { version = "0.8", features = [ "sqlite", "runtime-tokio", diff --git a/src/backends/go.rs b/src/backends/go.rs index f11a057..fa758c6 100644 --- a/src/backends/go.rs +++ b/src/backends/go.rs @@ -1,15 +1,12 @@ // SPDX-FileCopyrightText: 2026 Nikolay Govorov <me@govorov.online> // SPDX-License-Identifier: AGPL-3.0-or-later -use std::sync::Arc; use std::time::Duration; -use mime_guess::mime; use serde::{Deserialize, Serialize}; -use thiserror::Error; use url::Url; -use super::{Archive, Backend, BackendDelegate, FileKind, IndexError, ResolveError, ResolvedFile}; +use super::*; use crate::utils::deserialize_duration; #[derive(Debug, Clone, Deserialize)] @@ -20,6 +17,7 @@ pub struct GoConfig { #[serde(deserialize_with = "deserialize_duration")] pub refresh_interval: Duration, } + impl Default for GoConfig { fn default() -> Self { Self { @@ -30,6 +28,116 @@ impl Default for GoConfig { } } +impl BackendConfig for GoConfig { + fn enabled(&self) -> bool { + self.enabled + } + fn refresh_interval(&self) -> Duration { + self.refresh_interval + } +} + +/// Typed metadata for Go files +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct GoFileMeta { + pub kind: FileKind, +} + +/// Typed metadata for Go releases +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct GoReleaseMeta { + pub stable: bool, +} + +/// Upstream API file format (only used for parsing JSON from go.dev) +#[derive(Deserialize)] +struct UpstreamFile { + filename: String, + os: Option<String>, + arch: Option<String>, + sha256: String, + size: u64, + kind: FileKind, +} + +/// Upstream API release format (only used for parsing JSON from go.dev) +#[derive(Deserialize)] +struct UpstreamRelease { + version: String, + stable: bool, + files: Vec<UpstreamFile>, +} + +/// Go backend specification +pub struct GoSpec; + +#[async_trait::async_trait] +impl BackendSpec for GoSpec { + const ID: &'static str = "go"; + + type Config = GoConfig; + type ReleaseMeta = GoReleaseMeta; + type FileMeta = GoFileMeta; + type Filename<'a> = GoFilename<'a>; + + async fn fetch_index( + config: &Self::Config, + network: &dyn BackendNetwork, + ) -> Result<Vec<(RawRelease, Vec<RawReleaseFile>)>, BackendError> { + let mut url = config.upstream.clone(); + url.query_pairs_mut() + .append_pair("mode", "json") + .append_pair("include", "all"); + + let bytes = network.http_get(&url).await?; + let upstream: Vec<UpstreamRelease> = + serde_json::from_slice(&bytes).map_err(|e| BackendError::Upstream(e.to_string()))?; + + let mut releases = Vec::with_capacity(upstream.len()); + for u in upstream { + let sort_key = match GoVersion::parse(&u.version) { + Ok(v) => v.sort_key(), + Err(_) => { + tracing::error!(version = u.version, "invalid go version, skipping"); + continue; + } + }; + + let release: Release<GoReleaseMeta> = Release { + backend: Self::ID.to_string(), + version: u.version.clone(), + sort_key, + meta: GoReleaseMeta { stable: u.stable }, + }; + + let files: Vec<RawReleaseFile> = u + .files + .iter() + .map(|f| { + let file: ReleaseFile<GoFileMeta> = ReleaseFile { + backend: Self::ID.to_string(), + version: u.version.clone(), + filename: f.filename.clone(), + checksum: f.sha256.clone(), + size: f.size as i64, + os: f.os.as_ref().and_then(|s| s.parse().ok()), + arch: f.arch.as_ref().and_then(|s| s.parse().ok()), + meta: GoFileMeta { kind: f.kind }, + }; + file.to_raw() + }) + .collect(); + + releases.push((release.to_raw(), files)); + } + + Ok(releases) + } +} + +/// Type alias for Go backend +pub type GoBackend = Backend<GoSpec>; + #[derive(Debug, Clone, Copy, PartialEq, Eq)] enum ReleaseType { Stable, @@ -37,12 +145,8 @@ enum ReleaseType { Beta(u64), } -#[derive(Debug, Clone, PartialEq, Eq, Error)] -#[error("invalid tarball filename")] -struct ParseError; - #[derive(Debug, Clone, PartialEq, Eq)] -struct GoVersion { +pub struct GoVersion { major: u64, minor: u64, patch: Option<u64>, @@ -51,23 +155,23 @@ struct GoVersion { impl GoVersion { /// Parses version string like "go1", "go1.22.3", "go1.23rc1", or "go1.9.2rc2". - fn parse(s: &str) -> Result<Self, ParseError> { - let s = s.strip_prefix("go").ok_or(ParseError)?; + fn parse(s: &str) -> Result<Self, BackendError> { + let s = s.strip_prefix("go").ok_or(BackendError::NotFound)?; let parts: Vec<&str> = s.split('.').collect(); let (version, consumed) = Self::from_parts(&parts)?; if consumed != parts.len() { - return Err(ParseError); + return Err(BackendError::NotFound); } Ok(version) } /// Parses version from dot-separated parts, returns (version, parts_consumed). - fn from_parts(parts: &[&str]) -> Result<(Self, usize), ParseError> { + fn from_parts(parts: &[&str]) -> Result<(Self, usize), BackendError> { let major = parts .first() - .ok_or(ParseError)? + .ok_or(BackendError::NotFound)? .parse() - .map_err(|_| ParseError)?; + .map_err(|_| BackendError::NotFound)?; let Some(minor_str) = parts.get(1) else { return Ok(( @@ -129,18 +233,26 @@ impl GoVersion { /// "25" -> (25, None) /// "2rc2" -> (2, Some(ReleaseCandidate(2))) /// "1beta1" -> (1, Some(Beta(1))) - fn parse_version_part(s: &str) -> Result<(u64, Option<ReleaseType>), ParseError> { + fn parse_version_part(s: &str) -> Result<(u64, Option<ReleaseType>), BackendError> { if let Some(idx) = s.find("rc") { - let num = s[..idx].parse::<u64>().map_err(|_| ParseError)?; - let rc_num = s[idx + 2..].parse::<u64>().map_err(|_| ParseError)?; + let num = s[..idx] + .parse::<u64>() + .map_err(|_| BackendError::NotFound)?; + let rc_num = s[idx + 2..] + .parse::<u64>() + .map_err(|_| BackendError::NotFound)?; return Ok((num, Some(ReleaseType::ReleaseCandidate(rc_num)))); } if let Some(idx) = s.find("beta") { - let num = s[..idx].parse::<u64>().map_err(|_| ParseError)?; - let beta_num = s[idx + 4..].parse::<u64>().map_err(|_| ParseError)?; + let num = s[..idx] + .parse::<u64>() + .map_err(|_| BackendError::NotFound)?; + let beta_num = s[idx + 4..] + .parse::<u64>() + .map_err(|_| BackendError::NotFound)?; return Ok((num, Some(ReleaseType::Beta(beta_num)))); } - let num = s.parse::<u64>().map_err(|_| ParseError)?; + let num = s.parse::<u64>().map_err(|_| BackendError::NotFound)?; Ok((num, None)) } @@ -148,14 +260,14 @@ impl GoVersion { /// "25" -> (25, Stable) /// "26rc2" -> (26, ReleaseCandidate(2)) /// "26beta1" -> (26, Beta(1)) - fn parse_minor_with_release(s: &str) -> Result<(u64, ReleaseType), ParseError> { + fn parse_minor_with_release(s: &str) -> Result<(u64, ReleaseType), BackendError> { let (num, release) = GoVersion::parse_version_part(s)?; Ok((num, release.unwrap_or(ReleaseType::Stable))) } } #[cfg(test)] -mod version_tests { +mod tests_go_version { use super::*; #[test] @@ -188,7 +300,7 @@ mod version_tests { /// Describes a single file stored at `dl.google.com/go/`. #[derive(Debug, Clone, PartialEq, Eq)] -struct GoFile<'a> { +pub struct GoFilename<'a> { filename: &'a str, version: GoVersion, os: Option<&'a str>, @@ -198,14 +310,14 @@ struct GoFile<'a> { sha256: bool, } -impl<'a> GoFile<'a> { - pub fn parse(filename: &'a str) -> Result<Self, ParseError> { +impl<'a> BackendFilename<'a, GoConfig> for GoFilename<'a> { + fn parse(filename: &'a str) -> Result<Self, BackendError> { let mut buffer = filename; let mut sha256 = false; let archive; // Strip "go" prefix - buffer = buffer.strip_prefix("go").ok_or(ParseError)?; + buffer = buffer.strip_prefix("go").ok_or(BackendError::NotFound)?; // Check for .sha256 suffix if let Some(it) = buffer.strip_suffix(".sha256") { @@ -227,24 +339,24 @@ impl<'a> GoFile<'a> { buffer = it; archive = Archive::Pkg; } else { - return Err(ParseError); + return Err(BackendError::NotFound); } if buffer.is_empty() { - return Err(ParseError); + return Err(BackendError::NotFound); } // Split by dots: "1.25.6.linux-amd64" -> ["1", "25", "6", "linux-amd64"] let parts: Vec<&str> = buffer.split('.').collect(); if parts.len() < 3 { - return Err(ParseError); + return Err(BackendError::NotFound); } // Parse version, get how many parts were consumed let (version, consumed) = GoVersion::from_parts(&parts)?; // Remainder must exist and be either "src" or "os-arch" - let remainder = parts.get(consumed).ok_or(ParseError)?; + let remainder = parts.get(consumed).ok_or(BackendError::NotFound)?; let (os, arch, kind) = if *remainder == "src" { (None, None, FileKind::Source) } else if let Some((os, arch)) = remainder.split_once('-') { @@ -254,10 +366,10 @@ impl<'a> GoFile<'a> { }; (Some(os), Some(arch), kind) } else { - return Err(ParseError); + return Err(BackendError::NotFound); }; - Ok(GoFile { + Ok(GoFilename { filename, version, os, @@ -268,25 +380,28 @@ impl<'a> GoFile<'a> { }) } - /// Builds the upstream URL for this tarball. - pub fn upstream_url(&self, upstream: &Url, source: &str) -> Result<Url, ()> { - let mut url = upstream.clone(); + fn upstream_url(&self, config: &GoConfig, source: &str) -> Result<Url, BackendError> { + let mut url = config.upstream.clone(); url.path_segments_mut() - .map_err(|_| ())? + .map_err(|()| BackendError::Internal("cannot build upstream URL".into()))? .pop_if_empty() .push(self.filename); url.query_pairs_mut().append_pair("source", source); Ok(url) } + + fn is_sha256(&self) -> bool { + self.sha256 + } } #[cfg(test)] -mod file_tests { +mod tests_go_filename { use super::*; #[test] fn test_parse_stable_binary() { - let t = GoFile::parse("go1.25.6.linux-amd64.tar.gz").unwrap(); + let t = GoFilename::parse("go1.25.6.linux-amd64.tar.gz").unwrap(); assert_eq!(t.version.major, 1); assert_eq!(t.version.minor, 25); assert_eq!(t.version.patch, Some(6)); @@ -300,7 +415,7 @@ mod file_tests { #[test] fn test_parse_first_minor_release() { - let t = GoFile::parse("go1.25.linux-amd64.tar.gz").unwrap(); + let t = GoFilename::parse("go1.25.linux-amd64.tar.gz").unwrap(); assert_eq!(t.version.major, 1); assert_eq!(t.version.minor, 25); assert_eq!(t.version.patch, None); @@ -309,7 +424,7 @@ mod file_tests { #[test] fn test_parse_rc() { - let t = GoFile::parse("go1.26rc2.linux-amd64.tar.gz").unwrap(); + let t = GoFilename::parse("go1.26rc2.linux-amd64.tar.gz").unwrap(); assert_eq!(t.version.major, 1); assert_eq!(t.version.minor, 26); assert_eq!(t.version.patch, None); @@ -321,7 +436,7 @@ mod file_tests { #[test] fn test_parse_beta() { - let t = GoFile::parse("go1.26beta1.darwin-arm64.tar.gz").unwrap(); + let t = GoFilename::parse("go1.26beta1.darwin-arm64.tar.gz").unwrap(); assert_eq!(t.version.major, 1); assert_eq!(t.version.minor, 26); assert!(matches!(t.version.release_type, ReleaseType::Beta(1))); @@ -329,7 +444,7 @@ mod file_tests { #[test] fn test_parse_source() { - let t = GoFile::parse("go1.25.6.src.tar.gz").unwrap(); + let t = GoFilename::parse("go1.25.6.src.tar.gz").unwrap(); assert!(matches!(t.kind, FileKind::Source)); assert_eq!(t.os, None); assert_eq!(t.arch, None); @@ -338,7 +453,7 @@ mod file_tests { #[test] fn test_parse_windows_zip() { - let t = GoFile::parse("go1.25.6.windows-amd64.zip").unwrap(); + let t = GoFilename::parse("go1.25.6.windows-amd64.zip").unwrap(); assert_eq!(t.archive, Archive::Zip); assert_eq!(t.os, Some("windows")); assert_eq!(t.arch, Some("amd64")); @@ -347,21 +462,21 @@ mod file_tests { #[test] fn test_parse_msi() { - let t = GoFile::parse("go1.25.6.windows-amd64.msi").unwrap(); + let t = GoFilename::parse("go1.25.6.windows-amd64.msi").unwrap(); assert_eq!(t.archive, Archive::Msi); assert!(matches!(t.kind, FileKind::Installer)); } #[test] fn test_parse_pkg() { - let t = GoFile::parse("go1.25.6.darwin-arm64.pkg").unwrap(); + let t = GoFilename::parse("go1.25.6.darwin-arm64.pkg").unwrap(); assert_eq!(t.archive, Archive::Pkg); assert!(matches!(t.kind, FileKind::Installer)); } #[test] fn test_parse_sha256() { - let t = GoFile::parse("go1.25.6.linux-amd64.tar.gz.sha256").unwrap(); + let t = GoFilename::parse("go1.25.6.linux-amd64.tar.gz.sha256").unwrap(); assert!(t.sha256); assert_eq!(t.archive, Archive::TarGz); assert_eq!(t.version.patch, Some(6)); @@ -369,9 +484,12 @@ mod file_tests { #[test] fn test_upstream_url() { - let t = GoFile::parse("go1.25.6.linux-amd64.tar.gz").unwrap(); - let upstream = Url::parse("https://dl.google.com/go/").unwrap(); - let url = t.upstream_url(&upstream, "zorian:test").unwrap(); + let t = GoFilename::parse("go1.25.6.linux-amd64.tar.gz").unwrap(); + let config = GoConfig { + upstream: Url::parse("https://dl.google.com/go/").unwrap(), + ..Default::default() + }; + let url = t.upstream_url(&config, "zorian:test").unwrap(); assert_eq!( url.as_str(), "https://dl.google.com/go/go1.25.6.linux-amd64.tar.gz?source=zorian%3Atest" @@ -380,7 +498,7 @@ mod file_tests { #[test] fn test_parse_patch_rc() { - let t = GoFile::parse("go1.9.2rc2.linux-amd64.tar.gz").unwrap(); + let t = GoFilename::parse("go1.9.2rc2.linux-amd64.tar.gz").unwrap(); assert_eq!(t.version.major, 1); assert_eq!(t.version.minor, 9); assert_eq!(t.version.patch, Some(2)); @@ -392,313 +510,11 @@ mod file_tests { #[test] fn test_invalid_prefix() { - assert!(GoFile::parse("rust1.25.6.linux-amd64.tar.gz").is_err()); + assert!(GoFilename::parse("rust1.25.6.linux-amd64.tar.gz").is_err()); } #[test] fn test_invalid_extension() { - assert!(GoFile::parse("go1.25.6.linux-amd64.tar.bz2").is_err()); - } -} - -#[derive(Clone, Debug, Serialize, Deserialize)] -pub struct GoTarball { - pub filename: String, - pub os: Option<String>, - pub arch: Option<String>, - pub version: Option<String>, - pub sha256: String, - pub size: u64, - pub kind: FileKind, -} - -#[derive(Clone, Debug, Serialize, Deserialize)] -pub struct GoRelease { - pub version: String, - pub stable: bool, - pub files: Vec<GoTarball>, -} - -pub struct GoBackend { - config: GoConfig, - source: String, - delegate: Arc<dyn BackendDelegate>, -} -impl GoBackend { - pub fn new(config: GoConfig, source: String, delegate: Arc<dyn BackendDelegate>) -> Self { - Self { - config, - source, - delegate, - } - } -} -#[async_trait::async_trait] -impl Backend for GoBackend { - const ID: &'static str = "go"; - type Release = self::GoRelease; - - fn enabled(&self) -> bool { - self.config.enabled - } - - fn refresh_interval(&self) -> std::time::Duration { - self.config.refresh_interval - } - - async fn resolve_file(&self, filename: &str) -> Result<ResolvedFile, ResolveError> { - // For .sha256 files, return hash directly from the index - if let Some(base) = filename.strip_suffix(".sha256") { - let result: Result<Option<String>, _> = - sqlx::query_scalar("SELECT sha256 FROM go_files WHERE filename = ?1") - .bind(base) - .fetch_optional(self.delegate.db()) - .await; - - match result { - Ok(Some(hash)) => { - return Ok(ResolvedFile::Content { - data: hash.into(), - mime: mime::TEXT_PLAIN, - }); - } - Ok(None) => return Err(ResolveError::NotFound), - Err(e) => { - tracing::error!(filename, "failed to query sha256: {e}"); - return Err(ResolveError::Internal); - } - } - } - - // Check that file exists in index before proxying to upstream - let exists: Result<Option<i32>, _> = - sqlx::query_scalar("SELECT 1 FROM go_files WHERE filename = ?1") - .bind(filename) - .fetch_optional(self.delegate.db()) - .await; - - match exists { - Ok(None) => return Err(ResolveError::NotFound), - Err(e) => { - tracing::error!(filename, "failed to check file existence: {e}"); - return Err(ResolveError::Internal); - } - Ok(Some(_)) => {} - } - - let file = GoFile::parse(filename).map_err(|_| ResolveError::NotFound)?; - let url = file - .upstream_url(&self.config.upstream, &self.source) - .map_err(|_| ResolveError::Internal)?; - Ok(ResolvedFile::Upstream { - url, - mime: mime::APPLICATION_OCTET_STREAM, - }) - } - - async fn migrate(&self) -> Result<(), IndexError> { - sqlx::query( - "CREATE TABLE IF NOT EXISTS go_versions ( - id INTEGER PRIMARY KEY, - version TEXT NOT NULL UNIQUE, - stable INTEGER NOT NULL - ) STRICT", - ) - .execute(self.delegate.db()) - .await - .map_err(|e| IndexError::Database(e.to_string()))?; - - sqlx::query( - "CREATE TABLE IF NOT EXISTS go_files ( - version TEXT NOT NULL, - filename TEXT NOT NULL, - os TEXT, - arch TEXT, - sha256 TEXT NOT NULL, - size INTEGER NOT NULL, - kind TEXT NOT NULL, - PRIMARY KEY (version, filename), - FOREIGN KEY (version) REFERENCES go_versions(version) - ) STRICT", - ) - .execute(self.delegate.db()) - .await - .map_err(|e| IndexError::Database(e.to_string()))?; - - Ok(()) - } - - async fn fetch_index(&self) -> Result<(), IndexError> { - let mut url = self.config.upstream.clone(); - url.query_pairs_mut() - .append_pair("mode", "json") - .append_pair("include", "all"); - let bytes = self.delegate.http_get(&url).await?; - - let versions: Vec<GoRelease> = - serde_json::from_slice(&bytes).map_err(|e| IndexError::Parse(e.to_string()))?; - - for version in versions { - if let Err(e) = self.insert_version(&version).await { - tracing::error!(version = version.version, "failed to index version: {e}"); - } - } - - Ok(()) - } - - async fn get_versions(&self) -> Result<Vec<Self::Release>, IndexError> { - use futures::{StreamExt, TryStreamExt}; - - sqlx::query_as( - " - SELECT - v.version, v.stable, - COALESCE( - json_group_array(json_object( - 'filename', f.filename, 'os', f.os, 'arch', f.arch, - 'version', f.version, 'sha256', f.sha256, 'size', f.size, 'kind', f.kind - )) FILTER (WHERE f.filename IS NOT NULL), - '[]' - ) - FROM go_versions v - LEFT JOIN go_files f ON v.version = f.version - GROUP BY v.version - ORDER BY v.id ASC - ", - ) - .fetch(self.delegate.db()) - .map(|row| { - let (version, stable, files_json): (String, bool, String) = - row.map_err(|e| IndexError::Database(e.to_string()))?; - let files = - serde_json::from_str(&files_json).map_err(|e| IndexError::Parse(e.to_string()))?; - Ok(GoRelease { - version, - stable, - files, - }) - }) - .try_collect() - .await - } -} -impl GoBackend { - async fn insert_version(&self, version: &GoRelease) -> Result<(), IndexError> { - let id = GoVersion::parse(&version.version) - .map(|v| v.sort_key()) - .map_err(|_| IndexError::Parse(format!("invalid go version: {}", version.version)))?; - - let mut tx = self - .delegate - .db() - .begin() - .await - .map_err(|e| IndexError::Database(e.to_string()))?; - - sqlx::query( - "INSERT INTO go_versions (id, version, stable) - VALUES (?1, ?2, ?3) - ON CONFLICT(version) DO UPDATE SET - id = excluded.id, - stable = excluded.stable - WHERE id IS NOT excluded.id OR stable IS NOT excluded.stable", - ) - .bind(id) - .bind(&version.version) - .bind(version.stable) - .execute(&mut *tx) - .await - .map_err(|e| IndexError::Database(e.to_string()))?; - - for file in &version.files { - Self::insert_file(&mut tx, &version.version, file).await?; - } - - tx.commit() - .await - .map_err(|e| IndexError::Database(e.to_string()))?; - Ok(()) - } - - async fn insert_file( - tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>, - version: &str, - file: &GoTarball, - ) -> Result<(), IndexError> { - let kind = match file.kind { - FileKind::Source => "source", - FileKind::Archive => "archive", - FileKind::Installer => "installer", - FileKind::Bootstrap => unreachable!("go does not have bootstrap files"), - }; - - let os = file.os.as_deref().filter(|s| !s.is_empty()); - - let exists: Option<i32> = - sqlx::query_scalar("SELECT 1 FROM go_files WHERE version = ?1 AND filename = ?2") - .bind(version) - .bind(&file.filename) - .fetch_optional(&mut **tx) - .await - .map_err(|e| IndexError::Database(e.to_string()))?; - - let changed: Option<(i32,)> = sqlx::query_as( - "INSERT INTO go_files (version, filename, os, arch, sha256, size, kind) - VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7) - ON CONFLICT(version, filename) DO UPDATE SET - os = excluded.os, arch = excluded.arch, sha256 = excluded.sha256, - size = excluded.size, kind = excluded.kind - WHERE os IS NOT excluded.os OR arch IS NOT excluded.arch - OR sha256 IS NOT excluded.sha256 OR size IS NOT excluded.size - OR kind IS NOT excluded.kind - RETURNING 1", - ) - .bind(version) - .bind(&file.filename) - .bind(os) - .bind(&file.arch) - .bind(&file.sha256) - .bind(file.size as i64) - .bind(kind) - .fetch_optional(&mut **tx) - .await - .map_err(|e| IndexError::Database(e.to_string()))?; - - if exists.is_some() && changed.is_some() { - tracing::warn!(version, filename = file.filename, "go index file changed"); - } - - Ok(()) - } -} - -#[cfg(test)] -mod api_tests { - use super::*; - - #[test] - fn test_deserialize_go_release() { - let json = r#"{ - "version": "go1.22.0", - "stable": true, - "files": [ - { - "filename": "go1.22.0.windows-amd64.msi", - "os": "windows", - "arch": "amd64", - "version": "go1.22.0", - "sha256": "11a47de052db9971359e8c2f3a1667f8d56fa4c6bbec0687cf4cf2403a07628a", - "size": 63172608, - "kind": "installer" - } - ] - }"#; - - let release: GoRelease = serde_json::from_str(json).unwrap(); - assert_eq!(release.version, "go1.22.0"); - assert_eq!(release.files.len(), 1); - assert_eq!(release.files[0].filename, "go1.22.0.windows-amd64.msi"); - assert_eq!(release.files[0].size, 63172608); + assert!(GoFilename::parse("go1.25.6.linux-amd64.tar.bz2").is_err()); } } diff --git a/src/backends/mod.rs b/src/backends/mod.rs index 1bda2d9..d6344f7 100644 --- a/src/backends/mod.rs +++ b/src/backends/mod.rs @@ -4,30 +4,121 @@ pub mod go; pub mod zig; +use std::collections::HashMap; +use std::sync::Arc; use std::time::Duration; use async_trait::async_trait; use bytes::Bytes; use mime_guess::mime; -use serde::{Deserialize, Serialize}; -use sqlx::{Pool, Sqlite}; +use serde::{Deserialize, Serialize, de::DeserializeOwned}; use thiserror::Error; use url::Url; pub use go::{GoBackend, GoConfig}; pub use zig::{ZigBackend, ZigConfig}; -/// Error during index operations. -#[derive(Debug, Error)] -pub enum IndexError { - #[error("fetch error: {0}")] - Fetch(String), +/// Backend operation error. +#[derive(Debug, Clone, Error)] +pub enum BackendError { + /// Resource not found + #[error("not found")] + NotFound, + + /// Upstream network error + #[error("network: {0}")] + Network(String), + + /// Internal storage error + #[error("storage: {0}")] + Storage(String), + + /// Malformed data from upstream + #[error("upstream: {0}")] + Upstream(String), + + /// Internal logic error + #[error("internal: {0}")] + Internal(String), +} - #[error("database error: {0}")] - Database(String), +/// Operating system +#[derive( + Debug, + Clone, + Copy, + PartialEq, + Eq, + Serialize, + Deserialize, + strum::Display, + strum::EnumString, + strum::AsRefStr, +)] +#[serde(rename_all = "lowercase")] +#[strum(serialize_all = "lowercase")] +pub enum Os { + Linux, + Windows, + Darwin, + #[serde(alias = "macos")] + Macos, + Freebsd, + Netbsd, + Openbsd, + Illumos, + Plan9, + Aix, + Solaris, + Dragonfly, + Android, + Ios, + Js, + Wasip1, + Wasi, +} - #[error("parse error: {0}")] - Parse(String), +/// CPU architecture +#[derive( + Debug, + Clone, + Copy, + PartialEq, + Eq, + Serialize, + Deserialize, + strum::Display, + strum::EnumString, + strum::AsRefStr, +)] +#[serde(rename_all = "lowercase")] +#[strum(serialize_all = "lowercase")] +pub enum Arch { + Amd64, + #[serde(alias = "x86_64")] + X86_64, + Arm64, + #[serde(alias = "aarch64")] + Aarch64, + #[serde(rename = "386")] + #[strum(serialize = "386")] + I386, + Arm, + Armv6l, + Armv7a, + Loong64, + Mips, + Mips64, + Mips64le, + Mipsle, + Ppc64, + Ppc64le, + Riscv64, + S390x, + Wasm32, + Powerpc, + Powerpc64, + Powerpc64le, } #[derive(Debug, Clone, Copy, PartialEq, Eq)] @@ -48,13 +139,13 @@ pub enum FileKind { Installer, } -/// Error during file resolution. -#[derive(Debug, Clone, Copy)] -pub enum ResolveError { - /// File not found (404) - NotFound, - /// Internal error (500) - Internal, +/// Version type for ordering. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum VersionType { + Stable, + Rc(u64), + Beta(u64), + Dev(u64), } /// Result of resolving a file request. @@ -68,51 +159,355 @@ pub enum ResolvedFile { }, } -/// Delegate provides I/O primitives to backends. +/// Release entry with typed metadata +#[derive(Debug, Clone)] +pub struct Release<M> { + pub backend: String, + pub version: String, + pub sort_key: i64, + pub meta: M, +} + +/// Type alias for storage (untyped meta) +pub type RawRelease = Release<Option<serde_json::Value>>; + +impl<M: Serialize> Release<M> { + pub fn to_raw(&self) -> RawRelease { + Release { + backend: self.backend.clone(), + version: self.version.clone(), + sort_key: self.sort_key, + meta: Some(serde_json::to_value(&self.meta).expect("meta serialization")), + } + } +} + +impl RawRelease { + pub fn try_into_typed<M: DeserializeOwned>(self) -> Result<Release<M>, BackendError> { + let meta = self + .meta + .ok_or_else(|| BackendError::Internal("missing release meta".into()))?; + let meta = serde_json::from_value(meta) + .map_err(|e| BackendError::Internal(format!("release meta: {e}")))?; + Ok(Release { + backend: self.backend, + version: self.version, + sort_key: self.sort_key, + meta, + }) + } +} + +/// Release file entry with typed metadata +#[derive(Debug, Clone)] +pub struct ReleaseFile<M> { + pub backend: String, + pub version: String, + pub filename: String, + pub checksum: String, + pub size: i64, + pub os: Option<Os>, + pub arch: Option<Arch>, + pub meta: M, +} + +/// Type alias for storage (untyped meta) +pub type RawReleaseFile = ReleaseFile<Option<serde_json::Value>>; + +impl<M: Serialize> ReleaseFile<M> { + pub fn to_raw(&self) -> RawReleaseFile { + ReleaseFile { + backend: self.backend.clone(), + version: self.version.clone(), + filename: self.filename.clone(), + checksum: self.checksum.clone(), + size: self.size, + os: self.os, + arch: self.arch, + meta: Some(serde_json::to_value(&self.meta).expect("meta serialization")), + } + } +} + +impl RawReleaseFile { + pub fn try_into_typed<M: DeserializeOwned>(self) -> Result<ReleaseFile<M>, BackendError> { + let meta = self + .meta + .ok_or_else(|| BackendError::Internal("missing file meta".into()))?; + let meta = serde_json::from_value(meta) + .map_err(|e| BackendError::Internal(format!("file meta: {e}")))?; + Ok(ReleaseFile { + backend: self.backend, + version: self.version, + filename: self.filename, + checksum: self.checksum, + size: self.size, + os: self.os, + arch: self.arch, + meta, + }) + } +} + +/// Release view for display (version + metadata + files) +#[derive(Debug, Clone)] +pub struct ReleaseView<M, F> { + pub version: String, + pub meta: M, + pub files: Vec<F>, +} + +/// Trait for backend configuration with common fields. +pub trait BackendConfig { + fn enabled(&self) -> bool; + fn refresh_interval(&self) -> Duration; +} + +/// Trait for parsing filenames and handling checksums/signatures. +pub trait BackendFilename<'a, C>: Sized { + /// Backend-specific signature suffix (e.g., ".minisig"). + /// SHA256 is handled automatically by the base. + const SIGNATURE_SUFFIX: Option<&'static str> = None; + + fn parse(filename: &'a str) -> Result<Self, BackendError>; + + fn upstream_url(&self, config: &C, source: &str) -> Result<Url, BackendError>; + + /// Is this a .sha256 checksum file? + fn is_sha256(&self) -> bool; + + /// Is this a signature file (minisig, gpg)? + fn is_signature(&self) -> bool { + false + } + + /// Can bypass index check? (for dev builds) + fn can_bypass_index(&self) -> bool { + false + } + + /// Get signature content from file meta (not sha256) + fn signature_content(&self, _file: &RawReleaseFile) -> Option<Bytes> { + None + } +} + +/// Storage interface for backend version/file index #[async_trait] -pub trait BackendDelegate: Send + Sync { - /// Get SQLite pool - fn db(&self) -> &Pool<Sqlite>; +pub trait BackendStorage: Send + Sync { + /// Query versions by backend, ordered by sort_key. + async fn query_releases(&self, backend: &str) -> Result<Vec<RawRelease>, BackendError>; + + /// Insert or update releases with their files in a single transaction. + async fn insert_releases( + &self, + releases: &[(RawRelease, Vec<RawReleaseFile>)], + ) -> Result<(), BackendError>; + + /// Query files by backend with optional filters. + async fn query_files( + &self, + backend: &str, + version: Option<&str>, + filename: Option<&str>, + meta_null_field: Option<&str>, + ) -> Result<Vec<RawReleaseFile>, BackendError>; + + /// Update a single file. + async fn update_file(&self, file: &RawReleaseFile) -> Result<bool, BackendError>; +} +/// Delegate provides I/O primitives to backends. +#[async_trait] +pub trait BackendNetwork: Send + Sync { /// HTTP GET request - async fn http_get(&self, url: &Url) -> Result<Bytes, IndexError>; + async fn http_get(&self, url: &Url) -> Result<Bytes, BackendError>; } -/// Trait for backend-specific logic (parsing, URL building, version indexing). +/// Trait for backend-specific types and methods. #[async_trait] -pub trait Backend: Send + Sync + 'static { +pub trait BackendSpec: Send + Sync + 'static { /// Fixed unique identifier for storage const ID: &'static str; - /// Backend-specific release representation - type Release; + /// Meta field where signature is stored (for query_files with meta_null_field) + /// None = signatures not supported by this backend + const SIGNATURE_META_FIELD: Option<&'static str> = None; + + type Config: BackendConfig + Send + Sync; + type Filename<'a>: BackendFilename<'a, Self::Config> + Send; + + type FileMeta: DeserializeOwned; + type ReleaseMeta: DeserializeOwned; + + /// Fetch release index from upstream and convert to storage format + async fn fetch_index( + config: &Self::Config, + network: &dyn BackendNetwork, + ) -> Result<Vec<(RawRelease, Vec<RawReleaseFile>)>, BackendError>; + + /// Fetch signature for a file and return updated file + async fn fetch_signature( + _file: &RawReleaseFile, + _config: &Self::Config, + _source: &str, + _network: &dyn BackendNetwork, + ) -> Result<RawReleaseFile, BackendError> { + unimplemented!("fetch_signature unimplemented"); // default: not supported + } +} + +/// Generic backend implementation +pub struct Backend<S: BackendSpec> { + pub config: S::Config, + pub source: String, + pub storage: Arc<dyn BackendStorage>, + pub network: Arc<dyn BackendNetwork>, +} - /// Whether the backend is enabled - fn enabled(&self) -> bool; +impl<S: BackendSpec> Backend<S> { + pub fn new( + config: S::Config, + source: String, + storage: Arc<dyn BackendStorage>, + network: Arc<dyn BackendNetwork>, + ) -> Self { + Self { + config, + source, + storage, + network, + } + } - /// Index refresh interval - fn refresh_interval(&self) -> Duration; + pub fn enabled(&self) -> bool { + self.config.enabled() + } - /// Create tables for this backend (called at startup) - async fn migrate(&self) -> Result<(), IndexError>; + pub fn refresh_interval(&self) -> Duration { + self.config.refresh_interval() + } /// Fetch index from upstream and store in DB - async fn fetch_index(&self) -> Result<(), IndexError>; + pub async fn refresh(&self) -> Result<(), BackendError> { + let releases = S::fetch_index(&self.config, &*self.network).await?; + + self.storage.insert_releases(&releases).await?; + + // Fetch signatures if supported + if let Some(field) = S::SIGNATURE_META_FIELD { + let files = self + .storage + .query_files(S::ID, None, None, Some(field)) + .await?; + + for file in files { + match S::fetch_signature(&file, &self.config, &self.source, &*self.network).await { + Ok(updated) => { + if let Err(e) = self.storage.update_file(&updated).await { + tracing::debug!( + filename = file.filename, + "failed to store signature: {e}" + ); + } + } + Err(e) => { + tracing::debug!(filename = file.filename, "failed to fetch signature: {e}"); + } + } + } + } - /// Resolves filename to upstream URL or direct content. - async fn resolve_file(&self, filename: &str) -> Result<ResolvedFile, ResolveError>; + Ok(()) + } - /// Load versions from DB - async fn get_versions(&self) -> Result<Vec<Self::Release>, IndexError>; -} + /// Load releases from DB + pub async fn get_releases( + &self, + ) -> Result<Vec<ReleaseView<S::ReleaseMeta, ReleaseFile<S::FileMeta>>>, BackendError> { + let raw_releases = self.storage.query_releases(S::ID).await?; + let raw_files = self.storage.query_files(S::ID, None, None, None).await?; + + let mut files_by_version: HashMap<String, Vec<ReleaseFile<S::FileMeta>>> = HashMap::new(); + for raw_file in raw_files { + if let Ok(file) = raw_file.try_into_typed::<S::FileMeta>() { + files_by_version + .entry(file.version.clone()) + .or_default() + .push(file); + } + } -/// Version type for ordering. -#[derive(Debug, Clone, Copy, PartialEq, Eq)] -pub enum VersionType { - Stable, - Rc(u64), - Beta(u64), - Dev(u64), + let releases = raw_releases + .into_iter() + .filter_map(|raw| { + let typed: Release<S::ReleaseMeta> = raw.try_into_typed().ok()?; + let files = files_by_version.remove(&typed.version).unwrap_or_default(); + Some(ReleaseView { + version: typed.version, + meta: typed.meta, + files, + }) + }) + .collect(); + + Ok(releases) + } + + /// Resolves filename to upstream URL or direct content + pub async fn resolve_file(&self, filename: &str) -> Result<ResolvedFile, BackendError> { + let parsed = S::Filename::parse(filename)?; + + let is_checksum_or_sig = parsed.is_sha256() || parsed.is_signature(); + let mime = if is_checksum_or_sig { + mime::TEXT_PLAIN + } else { + mime::APPLICATION_OCTET_STREAM + }; + + let url = parsed.upstream_url(&self.config, &self.source)?; + + // Strip suffix to get base filename + let base = filename + .strip_suffix(".sha256") + .or_else(|| { + <S::Filename<'_> as BackendFilename<'_, S::Config>>::SIGNATURE_SUFFIX + .and_then(|s| filename.strip_suffix(s)) + }) + .unwrap_or(filename); + + let files = self + .storage + .query_files(S::ID, None, Some(base), None) + .await?; + + let Some(file) = files.into_iter().next() else { + return if parsed.can_bypass_index() { + Ok(ResolvedFile::Upstream { url, mime }) + } else { + Err(BackendError::NotFound) + }; + }; + + // SHA256 — universal, always in file.checksum + if parsed.is_sha256() { + return Ok(ResolvedFile::Content { + mime, + data: file.checksum.clone().into(), + }); + } + + // Signature — backend supports it, may or may not be indexed + if parsed.is_signature() { + if let Some(data) = parsed.signature_content(&file) { + return Ok(ResolvedFile::Content { mime, data }); + } + // Not indexed yet, proxy to upstream + return Ok(ResolvedFile::Upstream { url, mime }); + } + + Ok(ResolvedFile::Upstream { url, mime }) + } } /// Computes numeric sort key for correct version ordering. diff --git a/src/backends/zig.rs b/src/backends/zig.rs index 9f95c1a..a4b0bef 100644 --- a/src/backends/zig.rs +++ b/src/backends/zig.rs @@ -2,19 +2,15 @@ // SPDX-License-Identifier: AGPL-3.0-or-later use std::collections::HashMap; -use std::sync::Arc; use std::time::Duration; -use mime_guess::mime; use semver::Version as SemVersion; use serde::{Deserialize, Serialize}; -use thiserror::Error; use url::Url; -use super::{ - Archive, Backend, BackendDelegate, FileKind, IndexError, ResolveError, ResolvedFile, - VersionType, stable_version, -}; +use bytes::Bytes; + +use super::*; use crate::utils::{deserialize_duration, deserialize_size}; #[derive(Debug, Clone, Deserialize)] @@ -25,6 +21,7 @@ pub struct ZigConfig { #[serde(deserialize_with = "deserialize_duration")] pub refresh_interval: Duration, } + impl Default for ZigConfig { fn default() -> Self { Self { @@ -35,19 +32,203 @@ impl Default for ZigConfig { } } +impl BackendConfig for ZigConfig { + fn enabled(&self) -> bool { + self.enabled + } + fn refresh_interval(&self) -> Duration { + self.refresh_interval + } +} + +/// Typed metadata for Zig files +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct ZigFileMeta { + pub target: String, + + #[serde(default)] + pub minisig: Option<String>, +} + +/// Typed metadata for Zig releases (stored in DB) +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct ZigReleaseMeta { + pub date: Option<String>, + pub docs: Option<String>, + #[serde(alias = "stdDocs")] + pub std_docs: Option<String>, + pub notes: Option<String>, +} + +/// Upstream API tarball format (only used for parsing JSON from ziglang.org) +#[derive(Deserialize)] +struct UpstreamFile { + #[serde(alias = "tarball")] + filename: String, + shasum: String, + #[serde(deserialize_with = "deserialize_size")] + size: u64, +} + +/// Upstream API release format (only used for parsing JSON from ziglang.org) +#[derive(Deserialize)] +#[allow(dead_code)] +struct UpstreamRelease { + #[serde(flatten)] + meta: ZigReleaseMeta, + /// Present only in the "master" entry + version: Option<String>, + src: Option<UpstreamFile>, + bootstrap: Option<UpstreamFile>, + #[serde(flatten)] + targets: HashMap<String, UpstreamFile>, +} + +/// Zig file with metadata +pub type ZigFile = ReleaseFile<ZigFileMeta>; + +impl ZigFile { + fn from_upstream( + version: &str, + target: &str, + tarball: &UpstreamFile, + ) -> Option<RawReleaseFile> { + let url = Url::parse(&tarball.filename).ok()?; + let filename = url + .path_segments() + .and_then(|mut s| s.next_back()) + .filter(|s| !s.is_empty())?; + + let parsed = ZigFilename::parse(filename).ok()?; + + let file: ReleaseFile<ZigFileMeta> = ReleaseFile { + backend: ZigSpec::ID.to_string(), + version: version.to_string(), + filename: filename.to_string(), + checksum: tarball.shasum.clone(), + size: tarball.size as i64, + os: parsed.os.and_then(|s| s.parse().ok()), + arch: parsed.arch.and_then(|s| s.parse().ok()), + meta: ZigFileMeta { + minisig: None, + target: target.to_string(), + }, + }; + Some(file.to_raw()) + } +} + +/// Zig backend specification +pub struct ZigSpec; + +#[async_trait::async_trait] +impl BackendSpec for ZigSpec { + const ID: &'static str = "zig"; + const SIGNATURE_META_FIELD: Option<&'static str> = Some("minisig"); + + type Config = ZigConfig; + type ReleaseMeta = ZigReleaseMeta; + type FileMeta = ZigFileMeta; + type Filename<'a> = ZigFilename<'a>; + + async fn fetch_index( + config: &Self::Config, + network: &dyn BackendNetwork, + ) -> Result<Vec<(RawRelease, Vec<RawReleaseFile>)>, BackendError> { + let mut url = config.upstream.clone(); + url.path_segments_mut() + .map_err(|_| BackendError::Internal("cannot-be-a-base URL".into()))? + .pop_if_empty() + .extend(["download", "index.json"]); + + let bytes = network.http_get(&url).await?; + let map: HashMap<String, UpstreamRelease> = + serde_json::from_slice(&bytes).map_err(|e| BackendError::Upstream(e.to_string()))?; + + let mut releases = Vec::with_capacity(map.len()); + for (version_str, upstream) in map { + let sort_key = match ZigVersion::parse(&version_str) { + Ok(v) => v.sort_key(), + Err(e) => { + tracing::error!(version = version_str, "invalid zig version, skipping: {e}"); + continue; + } + }; + + let release: Release<ZigReleaseMeta> = Release { + backend: Self::ID.to_string(), + version: version_str.to_string(), + sort_key, + meta: upstream.meta.clone(), + }; + + let mut files = Vec::new(); + + if let Some(ref tarball) = upstream.src + && let Some(f) = ZigFile::from_upstream(&version_str, "src", tarball) + { + files.push(f); + } + if let Some(ref tarball) = upstream.bootstrap + && let Some(f) = ZigFile::from_upstream(&version_str, "bootstrap", tarball) + { + files.push(f); + } + for (target, tarball) in &upstream.targets { + if let Some(f) = ZigFile::from_upstream(&version_str, target, tarball) { + files.push(f); + } + } + + releases.push((release.to_raw(), files)); + } + + Ok(releases) + } + + async fn fetch_signature( + file: &RawReleaseFile, + config: &Self::Config, + source: &str, + network: &dyn BackendNetwork, + ) -> Result<RawReleaseFile, BackendError> { + let minisig_filename = format!("{}.minisig", file.filename); + let parsed = ZigFilename::parse(&minisig_filename) + .map_err(|_| BackendError::Internal(format!("invalid filename: {}", file.filename)))?; + + let url = parsed + .upstream_url(config, source) + .map_err(|_| BackendError::Internal("cannot build URL".into()))?; + + let bytes = network.http_get(&url).await?; + let minisig = + String::from_utf8(bytes.to_vec()).map_err(|e| BackendError::Upstream(e.to_string()))?; + + let mut typed: ReleaseFile<ZigFileMeta> = file.clone().try_into_typed()?; + typed.meta.minisig = Some(minisig); + + tracing::debug!(filename = file.filename, "cached minisig"); + Ok(typed.to_raw()) + } +} + +/// Type alias for Zig backend +pub type ZigBackend = Backend<ZigSpec>; + /// Wrapper for sort key computation. -enum ZigVersion { +pub enum ZigVersion { Master, Semver(SemVersion), } + impl ZigVersion { - fn parse(s: &str) -> Result<Self, IndexError> { + fn parse(s: &str) -> Result<Self, BackendError> { if s == "master" { return Ok(Self::Master); } SemVersion::parse(s) .map(Self::Semver) - .map_err(|_| IndexError::Parse(format!("invalid zig version: {s}"))) + .map_err(|_| BackendError::Upstream(format!("invalid zig version: {s}"))) } fn sort_key(&self) -> i64 { @@ -69,17 +250,13 @@ impl ZigVersion { } } -#[derive(Debug, Clone, PartialEq, Eq, Error)] -#[error("invalid tarball filename")] -struct ParseError; - /// Describes a single file stored at `ziglang.org/download/`. /// /// The tarball naming has changed several times. When parsing, /// we standardize the files, but for the reverse operation /// (getting a string from a tarball), we preserve the original path. #[derive(Debug, Clone, PartialEq, Eq)] -struct ZigFile<'a> { +pub struct ZigFilename<'a> { filename: &'a str, os: Option<&'a str>, arch: Option<&'a str>, @@ -90,8 +267,11 @@ struct ZigFile<'a> { development: bool, } -impl<'a> ZigFile<'a> { - pub fn parse(filename: &'a str) -> Result<Self, ParseError> { +impl<'a> BackendFilename<'a, ZigConfig> for ZigFilename<'a> { + const SIGNATURE_SUFFIX: Option<&'static str> = Some(".minisig"); + + /// Builds the upstream URL for this tarball. + fn parse(filename: &'a str) -> Result<Self, BackendError> { let mut buffer = filename; let mut minisig = false; let archive; @@ -99,7 +279,7 @@ impl<'a> ZigFile<'a> { // (?:|-bootstrap|-[a-zA-Z0-9_]+-[a-zA-Z0-9_]+)-( // \d+\.\d+\.\d+(?:-dev\.\d+\+[0-9a-f]+)? // )\.(?:tar\.xz|zip)(?:\.minisig)? - buffer = buffer.strip_prefix("zig-").ok_or(ParseError)?; + buffer = buffer.strip_prefix("zig-").ok_or(BackendError::NotFound)?; // (?:|bootstrap|[a-zA-Z0-9_]+-[a-zA-Z0-9_]+)-( // \d+\.\d+\.\d+(?:-dev\.\d+\+[0-9a-f]+)? @@ -119,25 +299,25 @@ impl<'a> ZigFile<'a> { buffer = it; archive = Archive::TarXz; } else { - return Err(ParseError); + return Err(BackendError::NotFound); } if buffer.is_empty() { - return Err(ParseError); + return Err(BackendError::NotFound); } let mut it = buffer.rsplit('-'); - let last = it.next().ok_or(ParseError)?; + let last = it.next().ok_or(BackendError::NotFound)?; let development = last.starts_with("dev"); let version = if !development { - SemVersion::parse(last).map_err(|_| ParseError)? + SemVersion::parse(last).map_err(|_| BackendError::NotFound)? } else { - let semver = it.next().ok_or(ParseError)?; + let semver = it.next().ok_or(BackendError::NotFound)?; let devver = last; let version_str = format!("{}-{}", semver, devver); - SemVersion::parse(&version_str).map_err(|_| ParseError)? + SemVersion::parse(&version_str).map_err(|_| BackendError::NotFound)? }; let (os, arch, kind) = if let Some(payload) = it.next() { @@ -151,9 +331,9 @@ impl<'a> ZigFile<'a> { let (os, arch) = if version <= SemVersion::new(0, 2, 0) && payload == "win64" { ("windows", "x86_64") } else if version <= SemVersion::new(0, 14, 0) { - (it.next().ok_or(ParseError)?, payload) + (it.next().ok_or(BackendError::NotFound)?, payload) } else { - (payload, it.next().ok_or(ParseError)?) + (payload, it.next().ok_or(BackendError::NotFound)?) }; (Some(os), Some(arch), FileKind::Archive) } @@ -162,10 +342,10 @@ impl<'a> ZigFile<'a> { }; if it.next().is_some() { - return Err(ParseError); + return Err(BackendError::NotFound); } - Ok(ZigFile { + Ok(ZigFilename { filename, os, arch, @@ -177,11 +357,13 @@ impl<'a> ZigFile<'a> { }) } - /// Builds the upstream URL for this tarball. - pub fn upstream_url(&self, upstream: &Url, source: &str) -> Result<Url, ()> { - let mut url = upstream.clone(); + fn upstream_url(&self, config: &ZigConfig, source: &str) -> Result<Url, BackendError> { + let mut url = config.upstream.clone(); + { - let mut segments = url.path_segments_mut().map_err(|_| ())?; + let mut segments = url + .path_segments_mut() + .map_err(|()| BackendError::Internal("cannot build upstream URL".into()))?; segments.pop_if_empty(); if self.development { segments.push("builds"); @@ -190,402 +372,47 @@ impl<'a> ZigFile<'a> { } segments.push(self.filename); } + url.query_pairs_mut().append_pair("source", source); Ok(url) } -} - -#[derive(Clone, Debug, Serialize, Deserialize)] -pub struct ZigTarball { - /// e.g. "zig-x86_64-linux-0.15.2.tar.xz" - /// Note: upstream API returns full URL in "tarball" field, we extract filename when storing - #[serde(alias = "tarball")] - pub filename: String, - - /// e.g. "02aa270f183da276e5b5920b1dac44a63f1a49e55050ebde3aecc9eb82f93239" - pub shasum: String, - - /// e.g. 53733924 - #[serde(deserialize_with = "deserialize_size")] - pub size: u64, -} -#[derive(Clone, Debug, Serialize, Deserialize)] -pub struct ZigRelease { - /// e.g. "0.15.2" (older releases don't have this field) - #[serde(default)] - pub version: String, - - /// e.g. "2025-10-11" - pub date: Option<String>, - - /// e.g. "https://ziglang.org/documentation/0.15.2/" - pub docs: Option<String>, - - /// e.g. "https://ziglang.org/documentation/0.15.2/std/" - #[serde(rename = "stdDocs")] - pub std_docs: Option<String>, - - /// e.g. "https://ziglang.org/download/0.15.2/release-notes.html" - pub notes: Option<String>, - - /// Source tarball - pub src: Option<ZigTarball>, - - /// Bootstrap tarball - pub bootstrap: Option<ZigTarball>, - - /// Platform-specific files (e.g., "x86_64-linux", "aarch64-macos") - #[serde(flatten)] - pub targets: HashMap<String, ZigTarball>, -} - -pub struct ZigBackend { - config: ZigConfig, - source: String, - delegate: Arc<dyn BackendDelegate>, -} -impl ZigBackend { - pub fn new(config: ZigConfig, source: String, delegate: Arc<dyn BackendDelegate>) -> Self { - Self { - config, - source, - delegate, - } - } -} -#[async_trait::async_trait] -impl Backend for ZigBackend { - const ID: &'static str = "zig"; - type Release = self::ZigRelease; - - fn enabled(&self) -> bool { - self.config.enabled - } - - fn refresh_interval(&self) -> std::time::Duration { - self.config.refresh_interval - } - - async fn resolve_file(&self, filename: &str) -> Result<ResolvedFile, ResolveError> { - let file = ZigFile::parse(filename).map_err(|_| ResolveError::NotFound)?; - let mime = if file.minisig { - mime::TEXT_PLAIN - } else { - mime::APPLICATION_OCTET_STREAM - }; - let url = file - .upstream_url(&self.config.upstream, &self.source) - .map_err(|_| ResolveError::Internal)?; - - // For stable builds, check index - let base_filename = filename.strip_suffix(".minisig").unwrap_or(filename); - let row: Option<Option<String>> = - sqlx::query_scalar("SELECT minisig FROM zig_files WHERE filename = ?1") - .bind(base_filename) - .fetch_optional(self.delegate.db()) - .await - .map_err(|e| { - tracing::error!(filename, "failed to query file: {e}"); - ResolveError::Internal - })?; - - // File not in index - if row.is_none() { - // Dev builds go directly to upstream without index check - return if file.development { - Ok(ResolvedFile::Upstream { url, mime }) - } else { - Err(ResolveError::NotFound) - }; - } - - // Return cached minisig if available - if file.minisig - && let Some(Some(data)) = row - { - return Ok(ResolvedFile::Content { - data: data.into(), - mime: mime::TEXT_PLAIN, - }); - } - - Ok(ResolvedFile::Upstream { url, mime }) - } - - async fn migrate(&self) -> Result<(), IndexError> { - sqlx::query( - "CREATE TABLE IF NOT EXISTS zig_versions ( - id INTEGER PRIMARY KEY, - version TEXT NOT NULL UNIQUE, - date TEXT, - docs TEXT, - std_docs TEXT, - notes TEXT - ) STRICT", - ) - .execute(self.delegate.db()) - .await - .map_err(|e| IndexError::Database(e.to_string()))?; - - sqlx::query( - "CREATE TABLE IF NOT EXISTS zig_files ( - version TEXT NOT NULL, - target TEXT NOT NULL, - filename TEXT NOT NULL, - shasum TEXT NOT NULL, - size INTEGER NOT NULL, - minisig TEXT, - PRIMARY KEY (version, target), - FOREIGN KEY (version) REFERENCES zig_versions(version) - ) STRICT", - ) - .execute(self.delegate.db()) - .await - .map_err(|e| IndexError::Database(e.to_string()))?; - - // Migration: add minisig column to existing tables - match sqlx::query("ALTER TABLE zig_files ADD COLUMN minisig TEXT") - .execute(self.delegate.db()) - .await - { - Ok(_) => {} - Err(sqlx::Error::Database(e)) if e.message().contains("duplicate column") => {} - Err(e) => return Err(IndexError::Database(e.to_string())), - } - - Ok(()) - } - - async fn fetch_index(&self) -> Result<(), IndexError> { - let mut url = self.config.upstream.clone(); - url.path_segments_mut() - .map_err(|_| IndexError::Parse("cannot-be-a-base URL".into()))? - .pop_if_empty() - .extend(["download", "index.json"]); - let bytes = self.delegate.http_get(&url).await?; - - let index: HashMap<String, ZigRelease> = - serde_json::from_slice(&bytes).map_err(|e| IndexError::Parse(e.to_string()))?; - - for (version_str, version) in index { - if let Err(e) = self.insert_version(&version_str, &version).await { - tracing::error!(version = version_str, "failed to index version: {e}"); - } - } - - self.fetch_minisigs().await?; - - Ok(()) + fn is_sha256(&self) -> bool { + false // Zig uses minisig, not sha256 } - async fn get_versions(&self) -> Result<Vec<Self::Release>, IndexError> { - use futures::{StreamExt, TryStreamExt}; - - #[derive(Deserialize)] - struct FileRow { - target: String, - filename: String, - shasum: String, - size: u64, - } - - sqlx::query_as(" - SELECT - v.version, v.date, v.docs, v.std_docs, v.notes, - COALESCE( - json_group_array(json_object('target', f.target, 'filename', f.filename, 'shasum', f.shasum, 'size', f.size)) FILTER (WHERE f.version IS NOT NULL), - '[]' - ) - FROM zig_versions v - LEFT JOIN zig_files f ON v.version = f.version - GROUP BY v.version - ORDER BY v.id ASC - ") - .fetch(self.delegate.db()) - .map(|row| { - let (version, date, docs, std_docs, notes, files_json): - (String, Option<String>, Option<String>, Option<String>, Option<String>, String) = - row.map_err(|e| IndexError::Database(e.to_string()))?; - - let file_rows: Vec<FileRow> = serde_json::from_str(&files_json) - .map_err(|e| IndexError::Parse(e.to_string()))?; - - let mut src = None; - let mut bootstrap = None; - let mut targets = HashMap::new(); - - for f in file_rows { - let file = ZigTarball { filename: f.filename, shasum: f.shasum, size: f.size }; - match f.target.as_str() { - "src" => src = Some(file), - "bootstrap" => bootstrap = Some(file), - _ => { targets.insert(f.target, file); } - } - } - - Ok(ZigRelease { version, date, docs, std_docs, notes, src, bootstrap, targets }) - }) - .try_collect() - .await - } -} -impl ZigBackend { - async fn insert_version( - &self, - version_str: &str, - version: &ZigRelease, - ) -> Result<(), IndexError> { - let id = ZigVersion::parse(version_str)?.sort_key(); - - let mut tx = self - .delegate - .db() - .begin() - .await - .map_err(|e| IndexError::Database(e.to_string()))?; - - sqlx::query( - "INSERT INTO zig_versions (id, version, date, docs, std_docs, notes) - VALUES (?1, ?2, ?3, ?4, ?5, ?6) - ON CONFLICT(version) DO UPDATE SET - id = excluded.id, - date = excluded.date, docs = excluded.docs, - std_docs = excluded.std_docs, notes = excluded.notes - WHERE id IS NOT excluded.id OR date IS NOT excluded.date - OR docs IS NOT excluded.docs OR std_docs IS NOT excluded.std_docs - OR notes IS NOT excluded.notes", - ) - .bind(id) - .bind(version_str) - .bind(&version.date) - .bind(&version.docs) - .bind(&version.std_docs) - .bind(&version.notes) - .execute(&mut *tx) - .await - .map_err(|e| IndexError::Database(e.to_string()))?; - - if let Some(ref file) = version.src { - Self::insert_file(&mut tx, version_str, "src", file).await?; - } - if let Some(ref file) = version.bootstrap { - Self::insert_file(&mut tx, version_str, "bootstrap", file).await?; - } - for (target, file) in &version.targets { - Self::insert_file(&mut tx, version_str, target, file).await?; - } - - tx.commit() - .await - .map_err(|e| IndexError::Database(e.to_string()))?; - Ok(()) + fn is_signature(&self) -> bool { + self.minisig } - async fn insert_file( - tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>, - version: &str, - target: &str, - file: &ZigTarball, - ) -> Result<(), IndexError> { - let url = Url::parse(&file.filename) - .map_err(|e| IndexError::Parse(format!("invalid tarball URL: {e}")))?; - let filename = url - .path_segments() - .and_then(|mut s| s.next_back()) - .filter(|s| !s.is_empty()) - .ok_or_else(|| IndexError::Parse(format!("no filename in URL: {}", file.filename)))?; - - let exists: Option<i32> = - sqlx::query_scalar("SELECT 1 FROM zig_files WHERE version = ?1 AND target = ?2") - .bind(version) - .bind(target) - .fetch_optional(&mut **tx) - .await - .map_err(|e| IndexError::Database(e.to_string()))?; - - let changed: Option<(i32,)> = sqlx::query_as( - "INSERT INTO zig_files (version, target, filename, shasum, size) - VALUES (?1, ?2, ?3, ?4, ?5) - ON CONFLICT(version, target) DO UPDATE SET - filename = excluded.filename, shasum = excluded.shasum, size = excluded.size - WHERE filename IS NOT excluded.filename - OR shasum IS NOT excluded.shasum - OR size IS NOT excluded.size - RETURNING 1", - ) - .bind(version) - .bind(target) - .bind(filename) - .bind(&file.shasum) - .bind(file.size as i64) - .fetch_optional(&mut **tx) - .await - .map_err(|e| IndexError::Database(e.to_string()))?; - - if exists.is_some() && changed.is_some() { - tracing::warn!(version, target, "zig index file changed"); - } - - Ok(()) + fn can_bypass_index(&self) -> bool { + self.development } - async fn fetch_minisigs(&self) -> Result<(), IndexError> { - let files: Vec<(String,)> = - sqlx::query_as("SELECT filename FROM zig_files WHERE minisig IS NULL") - .fetch_all(self.delegate.db()) - .await - .map_err(|e| IndexError::Database(e.to_string()))?; - - for (filename,) in files { - if let Err(e) = self.fetch_minisig(&filename).await { - tracing::debug!(filename, "failed to fetch minisig: {e}"); - } + fn signature_content(&self, file: &RawReleaseFile) -> Option<Bytes> { + if !self.minisig { + return None; } - - Ok(()) - } - - async fn fetch_minisig(&self, filename: &str) -> Result<(), IndexError> { - let minisig_filename = format!("{}.minisig", filename); - let file = ZigFile::parse(&minisig_filename) - .map_err(|_| IndexError::Parse(format!("invalid filename: {}", filename)))?; - - let url = file - .upstream_url(&self.config.upstream, &self.source) - .map_err(|_| IndexError::Parse("cannot build URL".into()))?; - - let bytes = self.delegate.http_get(&url).await?; - let minisig = - String::from_utf8(bytes.to_vec()).map_err(|e| IndexError::Parse(e.to_string()))?; - - sqlx::query("UPDATE zig_files SET minisig = ?1 WHERE filename = ?2") - .bind(&minisig) - .bind(filename) - .execute(self.delegate.db()) - .await - .map_err(|e| IndexError::Database(e.to_string()))?; - - tracing::debug!(filename, "cached minisig"); - Ok(()) + let typed: ReleaseFile<ZigFileMeta> = file.clone().try_into_typed().ok()?; + typed.meta.minisig.map(|s| s.into_bytes().into()) } } #[cfg(test)] -mod tests { +mod tests_zig_filename { use super::*; #[test] fn parse_old_combined_platform() { // <= 0.2.0: zig-PLATFORM-VERSION (only Windows used this format) - let file = ZigFile::parse("zig-win64-0.1.1.zip").unwrap(); + let file = ZigFilename::parse("zig-win64-0.1.1.zip").unwrap(); assert_eq!(file.os, Some("windows")); assert_eq!(file.arch, Some("x86_64")); assert_eq!(file.version, SemVersion::new(0, 1, 1)); assert_eq!(file.archive, Archive::Zip); assert_eq!(file.kind, FileKind::Archive); - let file = ZigFile::parse("zig-win64-0.2.0.zip").unwrap(); + let file = ZigFilename::parse("zig-win64-0.2.0.zip").unwrap(); assert_eq!(file.os, Some("windows")); assert_eq!(file.arch, Some("x86_64")); assert_eq!(file.version, SemVersion::new(0, 2, 0)); @@ -594,13 +421,13 @@ mod tests { #[test] fn parse_middle_os_arch_format() { // 0.2.0 to 0.14.0: zig-OS-ARCH-VERSION - let file = ZigFile::parse("zig-linux-x86_64-0.13.0.tar.xz").unwrap(); + let file = ZigFilename::parse("zig-linux-x86_64-0.13.0.tar.xz").unwrap(); assert_eq!(file.os, Some("linux")); assert_eq!(file.arch, Some("x86_64")); assert_eq!(file.version, SemVersion::new(0, 13, 0)); assert_eq!(file.archive, Archive::TarXz); - let file = ZigFile::parse("zig-windows-x86_64-0.10.0.zip").unwrap(); + let file = ZigFilename::parse("zig-windows-x86_64-0.10.0.zip").unwrap(); assert_eq!(file.os, Some("windows")); assert_eq!(file.arch, Some("x86_64")); } @@ -608,19 +435,19 @@ mod tests { #[test] fn parse_new_arch_os_format() { // > 0.14.0: zig-ARCH-OS-VERSION - let file = ZigFile::parse("zig-x86_64-linux-0.15.0.tar.xz").unwrap(); + let file = ZigFilename::parse("zig-x86_64-linux-0.15.0.tar.xz").unwrap(); assert_eq!(file.os, Some("linux")); assert_eq!(file.arch, Some("x86_64")); assert_eq!(file.version, SemVersion::new(0, 15, 0)); - let file = ZigFile::parse("zig-aarch64-macos-0.15.0.tar.xz").unwrap(); + let file = ZigFilename::parse("zig-aarch64-macos-0.15.0.tar.xz").unwrap(); assert_eq!(file.os, Some("macos")); assert_eq!(file.arch, Some("aarch64")); } #[test] fn parse_dev_version() { - let file = ZigFile::parse("zig-x86_64-linux-0.14.0-dev.123+abc123.tar.xz").unwrap(); + let file = ZigFilename::parse("zig-x86_64-linux-0.14.0-dev.123+abc123.tar.xz").unwrap(); assert!(file.development); assert_eq!(file.version.major, 0); assert_eq!(file.version.minor, 14); @@ -629,7 +456,7 @@ mod tests { #[test] fn parse_source_tarball() { - let file = ZigFile::parse("zig-0.13.0.tar.xz").unwrap(); + let file = ZigFilename::parse("zig-0.13.0.tar.xz").unwrap(); assert_eq!(file.os, None); assert_eq!(file.arch, None); assert_eq!(file.kind, FileKind::Source); @@ -637,18 +464,18 @@ mod tests { #[test] fn parse_bootstrap() { - let file = ZigFile::parse("zig-bootstrap-0.13.0.tar.xz").unwrap(); + let file = ZigFilename::parse("zig-bootstrap-0.13.0.tar.xz").unwrap(); assert_eq!(file.kind, FileKind::Bootstrap); } #[test] fn parse_minisig() { - let file = ZigFile::parse("zig-win64-0.1.1.zip.minisig").unwrap(); + let file = ZigFilename::parse("zig-win64-0.1.1.zip.minisig").unwrap(); assert!(file.minisig); assert_eq!(file.os, Some("windows")); assert_eq!(file.arch, Some("x86_64")); - let file = ZigFile::parse("zig-x86_64-linux-0.15.0.tar.xz.minisig").unwrap(); + let file = ZigFilename::parse("zig-x86_64-linux-0.15.0.tar.xz.minisig").unwrap(); assert!(file.minisig); assert_eq!(file.os, Some("linux")); } @@ -656,12 +483,12 @@ mod tests { #[test] fn parse_boundary_version() { // 0.14.0 should use OS-ARCH format - let file = ZigFile::parse("zig-linux-x86_64-0.14.0.tar.xz").unwrap(); + let file = ZigFilename::parse("zig-linux-x86_64-0.14.0.tar.xz").unwrap(); assert_eq!(file.os, Some("linux")); assert_eq!(file.arch, Some("x86_64")); // 0.14.1 should use ARCH-OS format - let file = ZigFile::parse("zig-x86_64-linux-0.14.1.tar.xz").unwrap(); + let file = ZigFilename::parse("zig-x86_64-linux-0.14.1.tar.xz").unwrap(); assert_eq!(file.os, Some("linux")); assert_eq!(file.arch, Some("x86_64")); } diff --git a/src/controller_backend.rs b/src/controller_backend.rs index 23ebbbd..98d1f88 100644 --- a/src/controller_backend.rs +++ b/src/controller_backend.rs @@ -7,20 +7,20 @@ use axum::{Router, body, extract, http, response, routing}; use mime_guess::mime; use tracing::error; -use crate::backends::{Backend, ResolveError, ResolvedFile}; +use crate::backends::{Backend, BackendError, BackendSpec, ResolvedFile}; use crate::proxy; use crate::storage; /// Generic controller for backend HTTP handling. -pub struct BackendController<B: Backend> { - backend: Arc<B>, +pub struct BackendController<S: BackendSpec> { + backend: Arc<Backend<S>>, storage: Arc<storage::StorageService>, upstream: Arc<proxy::ProxyService>, } -impl<B: Backend> BackendController<B> { +impl<S: BackendSpec> BackendController<S> { pub fn new( - backend: Arc<B>, + backend: Arc<Backend<S>>, storage: Arc<storage::StorageService>, upstream: Arc<proxy::ProxyService>, ) -> Self { @@ -46,17 +46,17 @@ impl<B: Backend> BackendController<B> { return Ok(Self::build_response(http::StatusCode::OK, data, mime)); } Ok(ResolvedFile::Upstream { url, mime }) => (url, mime), - Err(ResolveError::NotFound) => { - error!(backend = B::ID, filename, "file not found"); + Err(BackendError::NotFound) => { + error!(backend = S::ID, filename, "file not found"); return Err(http::StatusCode::NOT_FOUND); } - Err(ResolveError::Internal) => { - error!(backend = B::ID, filename, "internal error resolving file"); + Err(e) => { + error!(backend = S::ID, filename, "error resolving file: {e}"); return Err(http::StatusCode::INTERNAL_SERVER_ERROR); } }; - match controller.storage.get(B::ID, &filename).await { + match controller.storage.get(S::ID, &filename).await { Ok(Some(entry)) => { return Ok(Self::build_response( http::StatusCode::OK, @@ -67,7 +67,7 @@ impl<B: Backend> BackendController<B> { Ok(None) => {} Err(err) => { error!( - backend = B::ID, + backend = S::ID, filename, "failed to get file from storage: {err}" ); return Err(http::StatusCode::INTERNAL_SERVER_ERROR); @@ -79,11 +79,11 @@ impl<B: Backend> BackendController<B> { .fetch(proxy::DownloadRequest { url }) .await?; - match controller.storage.put(B::ID, &filename, &entry.bytes).await { + match controller.storage.put(S::ID, &filename, &entry.bytes).await { Ok(()) => {} Err(err) => { error!( - backend = B::ID, + backend = S::ID, filename, "failed to put file to storage: {err}" ); return Err(http::StatusCode::INTERNAL_SERVER_ERROR); diff --git a/src/controller_web.rs b/src/controller_web.rs index 1c33296..885a13b 100644 --- a/src/controller_web.rs +++ b/src/controller_web.rs @@ -16,7 +16,7 @@ use sqlx::types::chrono; use tower::{Layer, Service}; use tracing::error; -use crate::backends::{Backend, GoBackend, ZigBackend}; +use crate::backends::{GoBackend, ZigBackend}; #[derive(Embed)] #[folder = "src/assets/"] @@ -46,16 +46,16 @@ impl WebController { )) .with_state(self.clone()); - let assets = axum::Router::new().route("/assets/{*path}", routing::get(Self::assets)); + let assets = axum::Router::new() + .route("/assets/{*path}", routing::get(Self::assets)) + .layer(LastModifiedLayer {}); - axum::Router::new() - .merge(pages) - .merge(assets) - .layer(LastModifiedLayer {}) - .layer(tower_http::set_header::SetResponseHeaderLayer::overriding( + axum::Router::new().merge(pages).merge(assets).layer( + tower_http::set_header::SetResponseHeaderLayer::overriding( header::CONTENT_SECURITY_POLICY, HeaderValue::from_static(CSP), - )) + ), + ) } async fn assets(extract::Path(path): extract::Path<String>) -> Response<Body> { @@ -87,7 +87,7 @@ impl WebController { async fn index(extract::State(ctrl): extract::State<Arc<Self>>) -> Markup { let zig_versions = if let Some(ref backend) = ctrl.zig { - match backend.get_versions().await { + match backend.get_releases().await { Ok(v) => v, Err(e) => { error!("failed to get zig versions: {e}"); @@ -99,7 +99,7 @@ impl WebController { }; let go_versions = if let Some(ref backend) = ctrl.go { - match backend.get_versions().await { + match backend.get_releases().await { Ok(v) => v, Err(e) => { error!("failed to get go versions: {e}"); @@ -165,31 +165,23 @@ impl WebController { @for v in zig_versions.iter().rev() { tr { td { (v.version) } - td { (v.date.as_deref().unwrap_or("-")) } + td { (v.meta.date.as_deref().unwrap_or("-")) } td { - @if let Some(ref url) = v.docs { + @if let Some(ref url) = v.meta.docs { a href=(url) { "docs" } } " " - @if let Some(ref url) = v.std_docs { + @if let Some(ref url) = v.meta.std_docs { a href=(url) { "std" } } " " - @if let Some(ref url) = v.notes { + @if let Some(ref url) = v.meta.notes { a href=(url) { "notes" } } } td { - @if let Some(ref src) = v.src { - a href=(format!("/zig/{}", src.filename)) { "src" } - " " - } - @if let Some(ref bootstrap) = v.bootstrap { - a href=(format!("/zig/{}", bootstrap.filename)) { "bootstrap" } - " " - } - @for (target, tarball) in v.targets.iter() { - a href=(format!("/zig/{}", tarball.filename)) { (target) } + @for file in &v.files { + a href=(format!("/zig/{}", file.filename)) { (&file.meta.target) } " " } } @@ -240,7 +232,7 @@ impl WebController { @for v in go_versions.iter().rev() { tr { td { (v.version) } - td { @if v.stable { "✓" } @else { "" } } + td { @if v.meta.stable { "✓" } @else { "" } } td { @for file in &v.files { a href=(format!("/go/{}", file.filename)) { (file.filename) } diff --git a/src/main.rs b/src/main.rs index d708124..1dd5b55 100644 --- a/src/main.rs +++ b/src/main.rs @@ -31,48 +31,23 @@ use tokio::signal; use tracing::{error, info, trace}; use tracing_subscriber::registry::LookupSpan; -use crate::backends::{Backend, BackendDelegate, GoBackend, IndexError, ZigBackend}; +use crate::backends::{Backend, BackendSpec, GoBackend, ZigBackend}; use crate::controller_backend::BackendController; use crate::controller_web::WebController; -/// Implementation of BackendDelegate for the application. -struct AppDelegate { - proxy: Arc<proxy::ProxyService>, - storage: Arc<storage::StorageService>, -} - -#[async_trait::async_trait] -impl BackendDelegate for AppDelegate { - async fn http_get(&self, url: &url::Url) -> Result<bytes::Bytes, IndexError> { - self.proxy - .fetch(proxy::DownloadRequest { url: url.clone() }) - .await - .map(|f| f.bytes) - .map_err(|e| IndexError::Fetch(e.to_string())) - } - - fn db(&self) -> &sqlx::Pool<sqlx::Sqlite> { - self.storage.db() - } -} - -async fn init_backend<B: Backend>( - backend: B, +async fn init_backend<S: BackendSpec>( + backend: Backend<S>, index_tasks: &mut tokio::task::JoinSet<()>, index_cancel: tokio_util::sync::CancellationToken, -) -> Option<Arc<B>> { +) -> Option<Arc<Backend<S>>> { if !backend.enabled() { return None; } let backend = Arc::new(backend); - if let Err(e) = backend.migrate().await { - error!(backend = B::ID, "migration failed: {e}"); - std::process::exit(1); - } let interval = backend.refresh_interval(); if !interval.is_zero() { index_tasks.spawn(run_index_refresh( - B::ID, + S::ID, backend.clone(), interval, index_cancel, @@ -81,9 +56,9 @@ async fn init_backend<B: Backend>( Some(backend) } -async fn run_index_refresh<B: backends::Backend>( +async fn run_index_refresh<S: BackendSpec>( name: &'static str, - backend: Arc<B>, + backend: Arc<Backend<S>>, interval: std::time::Duration, cancel: tokio_util::sync::CancellationToken, ) { @@ -98,7 +73,7 @@ async fn run_index_refresh<B: backends::Backend>( } _ = ticker.tick() => { info!(backend = name, "refreshing index"); - match backend.fetch_index().await { + match backend.refresh().await { Ok(()) => info!(backend = name, "index refreshed"), Err(e) => error!(backend = name, "index refresh failed: {e}"), } @@ -180,16 +155,11 @@ async fn main() { telemetry::TelemetryService::init(config.telemetry(), config.appname(), VERSION); let storage = Arc::new(storage::StorageService::new(config.clone()).await.unwrap()); - let upstream = Arc::new(proxy::ProxyService::new()); + let network = Arc::new(proxy::ProxyService::new()); let source = format!("zorian:{}", config.appname()); let backends = config.backends(); - let delegate: Arc<dyn BackendDelegate> = Arc::new(AppDelegate { - proxy: upstream.clone(), - storage: storage.clone(), - }); - const REQUEST_ID_HEADER: http::HeaderName = http::HeaderName::from_static("x-request-id"); let trace_layer = tower_http::trace::TraceLayer::new_for_http() @@ -288,14 +258,24 @@ async fn main() { let index_cancel = tokio_util::sync::CancellationToken::new(); let zig_backend = init_backend( - ZigBackend::new(backends.zig.clone(), source.clone(), delegate.clone()), + ZigBackend::new( + backends.zig.clone(), + source.clone(), + storage.clone(), + network.clone(), + ), &mut index_tasks, index_cancel.clone(), ) .await; let go_backend = init_backend( - GoBackend::new(backends.go.clone(), source.clone(), delegate.clone()), + GoBackend::new( + backends.go.clone(), + source.clone(), + storage.clone(), + network.clone(), + ), &mut index_tasks, index_cancel.clone(), ) @@ -308,7 +288,7 @@ async fn main() { let ctrl = Arc::new(BackendController::new( backend.clone(), storage.clone(), - upstream.clone(), + network.clone(), )); app = app.nest("/zig", ctrl.router()); } @@ -317,7 +297,7 @@ async fn main() { let ctrl = Arc::new(BackendController::new( backend.clone(), storage.clone(), - upstream.clone(), + network.clone(), )); app = app.nest("/go", ctrl.router()); } diff --git a/src/proxy.rs b/src/proxy.rs index e627323..0794682 100644 --- a/src/proxy.rs +++ b/src/proxy.rs @@ -7,8 +7,12 @@ use hyper::{Request, http}; use hyper_tls::HttpsConnector; use hyper_util::client::legacy::{Client, connect::HttpConnector}; use hyper_util::rt::TokioExecutor; +use tower::ServiceExt; +use tower_http::follow_redirect::FollowRedirect; use url::Url; +use super::backends::{BackendError, BackendNetwork}; + #[derive(Clone)] pub struct DownloadRequest { pub url: Url, @@ -20,14 +24,16 @@ pub struct File { } pub struct ProxyService { - client: Client<HttpsConnector<HttpConnector>, Empty<Bytes>>, + client: FollowRedirect<Client<HttpsConnector<HttpConnector>, Empty<Bytes>>>, } impl ProxyService { pub fn new() -> Self { let https = HttpsConnector::new(); let client = Client::builder(TokioExecutor::new()).build(https); - Self { client } + Self { + client: FollowRedirect::new(client), + } } pub async fn fetch(&self, request: DownloadRequest) -> Result<File, http::StatusCode> { @@ -40,7 +46,8 @@ impl ProxyService { let response = self .client - .request(request) + .clone() + .oneshot(request) .await .map_err(|_| http::StatusCode::GATEWAY_TIMEOUT)?; @@ -59,3 +66,13 @@ impl ProxyService { Ok(File { bytes }) } } + +#[async_trait::async_trait] +impl BackendNetwork for ProxyService { + async fn http_get(&self, url: &url::Url) -> Result<bytes::Bytes, BackendError> { + self.fetch(DownloadRequest { url: url.clone() }) + .await + .map(|f| f.bytes) + .map_err(|e| BackendError::Network(e.to_string())) + } +} diff --git a/src/storage.rs b/src/storage.rs index ae5b056..a557f30 100644 --- a/src/storage.rs +++ b/src/storage.rs @@ -35,27 +35,25 @@ use sqlx::{FromRow, Pool, Sqlite, query, query_as, query_scalar, sqlite}; use thiserror::Error; use tokio::fs; use tokio::io::AsyncWriteExt; -use tracing::{debug, instrument, warn}; +use tracing::{debug, error, instrument, warn}; use uuid::Uuid; +use super::backends::{BackendError, BackendStorage, RawRelease, RawReleaseFile, Release}; use super::config::ConfigService; const SQLITE_POOL_SIZE: u32 = 16; const INLINE_THRESHOLD: usize = 256 * 1024; // 256 KB -// You cannot change this ID, this will desynchronize the records in the database and the files on the disk. -const UUID_ROOT_NAMESPACE: Uuid = Uuid::from_bytes([ - 0x8b, 0x06, 0x3c, 0x4c, 0x6b, 0x5c, 0x4a, 0x8b, 0x92, 0x8f, 0x75, 0x8b, 0x0e, 0x63, 0xc3, 0x5d, -]); - #[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] pub struct Id(pub Uuid); + impl std::ops::Deref for Id { type Target = Uuid; fn deref(&self) -> &Self::Target { &self.0 } } + impl sqlx::Type<Sqlite> for Id { fn type_info() -> sqlite::SqliteTypeInfo { <Vec<u8> as sqlx::Type<Sqlite>>::type_info() @@ -65,6 +63,7 @@ impl sqlx::Type<Sqlite> for Id { <Vec<u8> as sqlx::Type<Sqlite>>::compatible(ty) } } + impl<'r> sqlx::decode::Decode<'r, Sqlite> for Id { fn decode( value: sqlite::SqliteValueRef<'r>, @@ -74,6 +73,7 @@ impl<'r> sqlx::decode::Decode<'r, Sqlite> for Id { Ok(Id(uuid)) } } + impl<'q> Encode<'q, Sqlite> for Id { fn encode_by_ref( &self, @@ -86,12 +86,14 @@ impl<'q> Encode<'q, Sqlite> for Id { #[derive(Debug, Clone)] pub struct Blob(pub Bytes); + impl std::ops::Deref for Blob { type Target = Bytes; fn deref(&self) -> &Self::Target { &self.0 } } + impl sqlx::Type<Sqlite> for Blob { fn type_info() -> sqlite::SqliteTypeInfo { <Vec<u8> as sqlx::Type<Sqlite>>::type_info() // BLOB @@ -101,6 +103,7 @@ impl sqlx::Type<Sqlite> for Blob { <Vec<u8> as sqlx::Type<Sqlite>>::compatible(ty) } } + impl<'r> sqlx::decode::Decode<'r, Sqlite> for Blob { fn decode( value: sqlite::SqliteValueRef<'r>, @@ -128,9 +131,10 @@ pub enum StorageError { IoError(#[from] std::io::Error), } +/// Cached file entry (datafiles table) #[allow(unused)] #[derive(Debug, Clone, FromRow)] -pub struct File { +pub struct Object { pub id: Id, pub scope: String, pub created_at: chrono::DateTime<chrono::Utc>, @@ -141,9 +145,15 @@ pub struct File { pub inlined: bool, } -impl File { +impl Object { + // You cannot change this ID, this will desynchronize the records in the database and the files on the disk. + const UUID_ROOT_NAMESPACE: Uuid = Uuid::from_bytes([ + 0x8b, 0x06, 0x3c, 0x4c, 0x6b, 0x5c, 0x4a, 0x8b, 0x92, 0x8f, 0x75, 0x8b, 0x0e, 0x63, 0xc3, + 0x5d, + ]); + fn uuid(scope: &str, file: &str) -> Id { - let scope_ns = Uuid::new_v5(&UUID_ROOT_NAMESPACE, scope.as_bytes()); + let scope_ns = Uuid::new_v5(&Object::UUID_ROOT_NAMESPACE, scope.as_bytes()); Id(Uuid::new_v5(&scope_ns, file.as_bytes())) } @@ -164,6 +174,64 @@ impl File { } } +/// Builder for file query SQL +pub struct FileQuery<'a> { + backend: &'a str, + version: Option<&'a str>, + filename: Option<&'a str>, + meta_null: Option<&'a str>, +} + +impl<'a> FileQuery<'a> { + pub fn new(backend: &'a str) -> Self { + Self { + backend, + version: None, + filename: None, + meta_null: None, + } + } + + pub fn version(mut self, version: &'a str) -> Self { + self.version = Some(version); + self + } + + pub fn filename(mut self, filename: &'a str) -> Self { + self.filename = Some(filename); + self + } + + pub fn where_meta_null(mut self, field: &'a str) -> Self { + self.meta_null = Some(field); + self + } + + pub fn build_sql(&self) -> String { + let mut sql = String::from( + "SELECT backend, version, filename, checksum, size, os, arch, meta FROM files WHERE backend = ?1", + ); + let mut param_idx = 1; + + if self.version.is_some() { + param_idx += 1; + sql.push_str(&format!(" AND version = ?{}", param_idx)); + } + if self.filename.is_some() { + param_idx += 1; + sql.push_str(&format!(" AND filename = ?{}", param_idx)); + } + if let Some(field) = self.meta_null { + sql.push_str(&format!( + " AND (meta IS NULL OR json_extract(meta, '$.{}') IS NULL)", + field + )); + } + + sql + } +} + struct FileSystem { objects: PathBuf, database: PathBuf, @@ -186,7 +254,7 @@ impl FileSystem { } fn object(&self, scope: &str, file: &str) -> PathBuf { - let id = File::uuid(scope, file); + let id = Object::uuid(scope, file); let hash_hex = hex::encode(id.0.as_bytes()); self.objects_root() .join(&hash_hex[0..2]) @@ -259,11 +327,6 @@ impl StorageService { Ok(storage) } - /// Get SQLite pool for use by backends - pub fn db(&self) -> &Pool<Sqlite> { - &self.sqlite - } - /// Synchronously traverses the tree and removes temporary files. /// Must run before the application starts. async fn doctor(&self) -> Result<(), StorageError> { @@ -325,9 +388,9 @@ impl StorageService { } async fn migrations(&self) -> Result<(), StorageError> { + // Object storage table query( - " - CREATE TABLE IF NOT EXISTS datafiles( + "CREATE TABLE IF NOT EXISTS datafiles( id BLOB PRIMARY KEY CHECK (length(id) = 16), scope TEXT NOT NULL, created_at TEXT DEFAULT (datetime('now')), @@ -335,20 +398,68 @@ impl StorageService { file_size INTEGER NOT NULL, file_bytes BLOB, inlined INTEGER NOT NULL, - UNIQUE (scope, file_name) - ) STRICT; - ", + ) STRICT", ) .execute(&self.sqlite) .await?; + // Drop old backend-specific tables (data will be re-fetched from upstream) + query("DROP TABLE IF EXISTS go_files") + .execute(&self.sqlite) + .await?; + query("DROP TABLE IF EXISTS go_versions") + .execute(&self.sqlite) + .await?; + query("DROP TABLE IF EXISTS zig_files") + .execute(&self.sqlite) + .await?; + query("DROP TABLE IF EXISTS zig_versions") + .execute(&self.sqlite) + .await?; + + // Unified versions table + query( + "CREATE TABLE IF NOT EXISTS versions ( + backend TEXT NOT NULL, + version TEXT NOT NULL, + sort_key INTEGER NOT NULL, + meta BLOB CHECK (meta IS NULL OR (json_valid(meta) AND json_type(meta) = 'object')), + PRIMARY KEY (backend, version) + ) STRICT", + ) + .execute(&self.sqlite) + .await?; + + // Unified files table + query( + "CREATE TABLE IF NOT EXISTS files ( + backend TEXT NOT NULL, + version TEXT NOT NULL, + filename TEXT NOT NULL, + checksum TEXT NOT NULL, + size INTEGER NOT NULL, + os TEXT, + arch TEXT, + meta BLOB CHECK (meta IS NULL OR (json_valid(meta) AND json_type(meta) = 'object')), + PRIMARY KEY (backend, version, filename), + FOREIGN KEY (backend, version) REFERENCES versions(backend, version) + ) STRICT", + ) + .execute(&self.sqlite) + .await?; + + // Index for fast filename lookups + query("CREATE INDEX IF NOT EXISTS idx_files_backend_filename ON files(backend, filename)") + .execute(&self.sqlite) + .await?; + Ok(()) } #[instrument(skip(self))] - pub async fn get(&self, scope: &str, filename: &str) -> Result<Option<File>, StorageError> { - let file: Option<File> = + pub async fn get(&self, scope: &str, filename: &str) -> Result<Option<Object>, StorageError> { + let file: Option<Object> = query_as("SELECT * FROM datafiles WHERE scope = ?1 AND file_name = ?2") .bind(scope) .bind(filename) @@ -370,7 +481,7 @@ impl StorageService { Err(e) => return Err(e.into()), }; - let hash = File::hash(scope, filename, &bytes); + let hash = Object::hash(scope, filename, &bytes); if hash != file.file_bytes.to_vec() { return Err(StorageError::IntegrityError); } @@ -418,12 +529,12 @@ impl StorageService { let payload: Vec<u8> = if inlined { bytes.to_vec() } else { - File::hash(scope, filename, bytes) + Object::hash(scope, filename, bytes) }; let result = async { let mut tx = self.sqlite.begin().await?; - let id = File::uuid(scope, filename); + let id = Object::uuid(scope, filename); let result = query( " @@ -458,7 +569,7 @@ impl StorageService { Ok(()) } Err(sqlx::Error::Database(ref db_err)) if db_err.is_unique_violation() => { - let existing: Option<File> = + let existing: Option<Object> = query_as("SELECT * FROM datafiles WHERE scope = ?1 AND file_name = ?2") .bind(scope) .bind(filename) @@ -499,6 +610,273 @@ impl StorageService { result } + + /// Insert or update releases with their files in a single transaction. + pub async fn insert_releases( + &self, + releases: &[(RawRelease, Vec<RawReleaseFile>)], + ) -> Result<(), StorageError> { + for (release, files) in releases { + let mut tx = match self.sqlite.begin().await { + Ok(tx) => tx, + Err(e) => { + error!( + backend = release.backend, + version = release.version, + "failed to start transaction: {e}" + ); + return Err(e.into()); + } + }; + + let meta_json = release + .meta + .as_ref() + .map(|m| serde_json::to_vec(m).unwrap()); + + let result = async { + // Insert/update version + query( + "INSERT INTO versions (backend, version, sort_key, meta) + VALUES (?1, ?2, ?3, ?4) + ON CONFLICT(backend, version) DO UPDATE SET + sort_key = excluded.sort_key, + meta = excluded.meta + WHERE sort_key IS NOT excluded.sort_key OR meta IS NOT excluded.meta", + ) + .bind(&release.backend) + .bind(&release.version) + .bind(release.sort_key) + .bind(meta_json.as_deref()) + .execute(&mut *tx) + .await?; + + // Insert/update files + for file in files { + let changed = + Self::insert_file(&mut tx, &release.backend, &release.version, file) + .await?; + if changed { + warn!( + backend = release.backend, + version = release.version, + filename = file.filename, + "index file changed" + ); + } + } + + Ok::<(), StorageError>(()) + } + .await; + + match result { + Ok(()) => { + if let Err(e) = tx.commit().await { + error!( + backend = release.backend, + version = release.version, + "failed to commit release: {e}" + ); + } + } + Err(e) => { + error!( + backend = release.backend, + version = release.version, + "failed to store release, skipping: {e}" + ); + let _ = tx.rollback().await; + } + } + } + + Ok(()) + } + + /// Insert or update a single file within a transaction. + /// Returns `true` if an existing file was modified. + async fn insert_file( + tx: &mut sqlx::Transaction<'_, Sqlite>, + backend: &str, + version: &str, + file: &RawReleaseFile, + ) -> Result<bool, StorageError> { + let meta_json = file.meta.as_ref().map(|m| serde_json::to_vec(m).unwrap()); + let os_str = file.os.as_ref().map(AsRef::as_ref); + let arch_str = file.arch.as_ref().map(AsRef::as_ref); + + let existed: Option<(i32,)> = + query_as("SELECT 1 FROM files WHERE backend = ?1 AND version = ?2 AND filename = ?3") + .bind(backend) + .bind(version) + .bind(&file.filename) + .fetch_optional(&mut **tx) + .await?; + + let changed: Option<(i32,)> = query_as( + "INSERT INTO files (backend, version, filename, checksum, size, os, arch, meta) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8) + ON CONFLICT(backend, version, filename) DO UPDATE SET + checksum = excluded.checksum, + size = excluded.size, + os = excluded.os, + arch = excluded.arch, + meta = excluded.meta + WHERE checksum IS NOT excluded.checksum + OR size IS NOT excluded.size + OR os IS NOT excluded.os + OR arch IS NOT excluded.arch + OR meta IS NOT excluded.meta + RETURNING 1", + ) + .bind(backend) + .bind(version) + .bind(&file.filename) + .bind(&file.checksum) + .bind(file.size) + .bind(os_str) + .bind(arch_str) + .bind(meta_json.as_deref()) + .fetch_optional(&mut **tx) + .await?; + + Ok(existed.is_some() && changed.is_some()) + } + + /// Execute a file query. + pub async fn query_files_filtered( + &self, + q: FileQuery<'_>, + ) -> Result<Vec<RawReleaseFile>, StorageError> { + let sql = q.build_sql(); + + let mut query = sqlx::query_as::< + _, + ( + String, + String, + String, + String, + i64, + Option<String>, + Option<String>, + Option<Vec<u8>>, + ), + >(&sql) + .bind(q.backend); + + if let Some(v) = q.version { + query = query.bind(v); + } + if let Some(f) = q.filename { + query = query.bind(f); + } + + let rows = query.fetch_all(&self.sqlite).await?; + + rows.into_iter() + .map( + |(backend, version, filename, checksum, size, os_str, arch_str, meta_bytes)| { + let meta = meta_bytes + .map(|b| serde_json::from_slice(&b)) + .transpose() + .map_err(|e| StorageError::DbError(sqlx::Error::Decode(Box::new(e))))?; + + Ok(RawReleaseFile { + backend, + version, + filename, + checksum, + size, + os: os_str.and_then(|s| s.parse().ok()), + arch: arch_str.and_then(|s| s.parse().ok()), + meta, + }) + }, + ) + .collect() + } + + /// Query all versions for a backend, ordered by sort_key. + pub async fn query_releases(&self, backend: &str) -> Result<Vec<RawRelease>, StorageError> { + let rows = query_as::<_, (String, String, i64, Option<Vec<u8>>)>( + "SELECT backend, version, sort_key, meta FROM versions WHERE backend = ?1 ORDER BY sort_key", + ) + .bind(backend) + .fetch_all(&self.sqlite) + .await?; + + rows.into_iter() + .map(|(backend, version, sort_key, meta_bytes)| { + let meta = meta_bytes + .map(|b| serde_json::from_slice(&b)) + .transpose() + .map_err(|e| StorageError::DbError(sqlx::Error::Decode(Box::new(e))))?; + + Ok(Release { + backend, + version, + sort_key, + meta, + }) + }) + .collect() + } + + /// Update a single file entry. + pub async fn update_file(&self, file: &RawReleaseFile) -> Result<bool, StorageError> { + let mut tx = self.sqlite.begin().await?; + let changed = Self::insert_file(&mut tx, &file.backend, &file.version, file).await?; + tx.commit().await?; + Ok(changed) + } +} + +#[async_trait::async_trait] +impl BackendStorage for StorageService { + async fn insert_releases( + &self, + releases: &[(RawRelease, Vec<RawReleaseFile>)], + ) -> Result<(), BackendError> { + StorageService::insert_releases(self, releases) + .await + .map_err(|e| BackendError::Storage(e.to_string())) + } + + async fn query_releases(&self, backend: &str) -> Result<Vec<RawRelease>, BackendError> { + StorageService::query_releases(self, backend) + .await + .map_err(|e| BackendError::Storage(e.to_string())) + } + + async fn query_files( + &self, + backend: &str, + version: Option<&str>, + filename: Option<&str>, + meta_null_field: Option<&str>, + ) -> Result<Vec<RawReleaseFile>, BackendError> { + let mut q = FileQuery::new(backend); + if let Some(v) = version { + q = q.version(v); + } + if let Some(f) = filename { + q = q.filename(f); + } + if let Some(field) = meta_null_field { + q = q.where_meta_null(field); + } + self.query_files_filtered(q) + .await + .map_err(|e| BackendError::Storage(e.to_string())) + } + + async fn update_file(&self, file: &RawReleaseFile) -> Result<bool, BackendError> { + StorageService::update_file(self, file) + .await + .map_err(|e| BackendError::Storage(e.to_string())) + } } #[cfg(test)] |
