aboutsummaryrefslogtreecommitdiffstats
diff options
context:
space:
mode:
Diffstat
-rw-r--r--Cargo.lock1+1 −0
-rw-r--r--Cargo.toml1+1 −0
-rw-r--r--src/backends/go.rs526+171 −355
-rw-r--r--src/backends/mod.rs487+441 −46
-rw-r--r--src/backends/zig.rs653+240 −413
-rw-r--r--src/controller_backend.rs26+13 −13
-rw-r--r--src/controller_web.rs42+17 −25
-rw-r--r--src/main.rs66+23 −43
-rw-r--r--src/proxy.rs23+20 −3
-rw-r--r--src/storage.rs430+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)]