From 5dcf89baba6493d6445e0682c10377ba5bf0f6f3 Mon Sep 17 00:00:00 2001 From: Seongmin Lee Date: Thu, 3 Sep 2026 02:55:04 +0900 Subject: [PATCH] gitmirror: split into sub-crates Signed-off-by: Seongmin Lee --- Cargo.lock | 69 +- Cargo.toml | 9 +- crates/gix-receive-pack/src/lib.rs | 37 -- crates/tangled-axum/Cargo.toml | 16 + crates/tangled-axum/src/atproto.rs | 181 +++++ crates/tangled-axum/src/lib.rs | 1 + gitmirror/crates/gitmirror-git/Cargo.toml | 18 + gitmirror/crates/gitmirror-git/src/layout.rs | 61 ++ gitmirror/crates/gitmirror-git/src/lib.rs | 3 + .../gitmirror-git/src/repo_ext.rs} | 7 +- gitmirror/crates/gitmirror-git/src/scratch.rs | 30 + gitmirror/crates/gitmirror-ingest/Cargo.toml | 11 + gitmirror/crates/gitmirror-ingest/src/lib.rs | 0 .../{ => crates/gitmirror-xrpc}/Cargo.toml | 28 +- .../gitmirror-xrpc/src/did_ext.rs} | 0 .../{ => crates/gitmirror-xrpc}/src/diff.rs | 0 gitmirror/crates/gitmirror-xrpc/src/error.rs | 91 +++ .../gitmirror-xrpc}/src/git_transport.rs | 0 gitmirror/crates/gitmirror-xrpc/src/lib.rs | 91 +++ .../{ => crates/gitmirror-xrpc}/src/merge.rs | 2 +- .../gitmirror-xrpc/src/routes/git.rs} | 558 +++++++++++++++- .../crates/gitmirror-xrpc/src/routes/mod.rs | 1 + gitmirror/crates/gitmirror/Cargo.toml | 34 + gitmirror/crates/gitmirror/src/config.rs | 111 ++++ gitmirror/crates/gitmirror/src/main.rs | 223 +++++++ gitmirror/src/auth.rs | 215 ------ gitmirror/src/git.rs | 140 ---- gitmirror/src/layout.rs | 58 -- gitmirror/src/main.rs | 95 --- gitmirror/src/xrpc.rs | 623 ------------------ nix/Cargo.nix | 288 ++++++-- 31 files changed, 1717 insertions(+), 1284 deletions(-) create mode 100644 crates/tangled-axum/Cargo.toml create mode 100644 crates/tangled-axum/src/atproto.rs create mode 100644 crates/tangled-axum/src/lib.rs create mode 100644 gitmirror/crates/gitmirror-git/Cargo.toml create mode 100644 gitmirror/crates/gitmirror-git/src/layout.rs create mode 100644 gitmirror/crates/gitmirror-git/src/lib.rs rename gitmirror/{src/git_repo_ext.rs => crates/gitmirror-git/src/repo_ext.rs} (95%) create mode 100644 gitmirror/crates/gitmirror-git/src/scratch.rs create mode 100644 gitmirror/crates/gitmirror-ingest/Cargo.toml create mode 100644 gitmirror/crates/gitmirror-ingest/src/lib.rs rename gitmirror/{ => crates/gitmirror-xrpc}/Cargo.toml (68%) rename gitmirror/{src/repo_did.rs => crates/gitmirror-xrpc/src/did_ext.rs} (100%) rename gitmirror/{ => crates/gitmirror-xrpc}/src/diff.rs (100%) create mode 100644 gitmirror/crates/gitmirror-xrpc/src/error.rs rename gitmirror/{ => crates/gitmirror-xrpc}/src/git_transport.rs (100%) create mode 100644 gitmirror/crates/gitmirror-xrpc/src/lib.rs rename gitmirror/{ => crates/gitmirror-xrpc}/src/merge.rs (99%) rename gitmirror/{src/xrpc_git.rs => crates/gitmirror-xrpc/src/routes/git.rs} (59%) create mode 100644 gitmirror/crates/gitmirror-xrpc/src/routes/mod.rs create mode 100644 gitmirror/crates/gitmirror/Cargo.toml create mode 100644 gitmirror/crates/gitmirror/src/config.rs create mode 100644 gitmirror/crates/gitmirror/src/main.rs delete mode 100644 gitmirror/src/auth.rs delete mode 100644 gitmirror/src/git.rs delete mode 100644 gitmirror/src/layout.rs delete mode 100644 gitmirror/src/main.rs delete mode 100644 gitmirror/src/xrpc.rs diff --git a/Cargo.lock b/Cargo.lock index 98f9880b1..40ecff27f 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2741,39 +2741,71 @@ dependencies = [ [[package]] name = "gitmirror" -version = "0.1.0" +version = "0.0.1" +dependencies = [ + "anyhow", + "axum", + "bobbin-runtime", + "clap", + "confique", + "futures", + "gitmirror-xrpc", + "jacquard-axum", + "jacquard-common", + "jacquard-identity", + "reqwest 0.13.1", + "rustls", + "socket2", + "thiserror 2.0.18", + "tokio", + "tokio-util", + "tracing", + "tracing-subscriber", +] + +[[package]] +name = "gitmirror-git" +version = "0.0.1" +dependencies = [ + "gix", + "gix-rebase", + "itertools 0.14.0", + "jacquard-common", + "tempfile", + "thiserror 2.0.18", +] + +[[package]] +name = "gitmirror-ingest" +version = "0.0.1" + +[[package]] +name = "gitmirror-xrpc" +version = "0.0.1" dependencies = [ "anyhow", "axum", - "base64", "bobbin-knot-proxy", "bobbin-runtime", - "bobbin-types", - "bytes", "chrono", - "clap", "futures-lite", + "gitmirror-git", "gix", "gix-pack", - "gix-rebase", "gix-receive-pack", "gix-transport", "jacquard-axum", "jacquard-common", "jacquard-identity", - "k256", + "lexicons", "line-numbers", - "multibase", "reqwest 0.13.1", "rustc-hash", - "serde", "serde_json", + "tangled-axum", "tempfile", - "thiserror 2.0.18", "tokio", - "tokio-stream", "tracing", - "tracing-subscriber", ] [[package]] @@ -8593,6 +8625,18 @@ version = "0.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7b2093cf4c8eb1e67749a6762251bc9cd836b6fc171623bd0a9d324d37af2417" +[[package]] +name = "tangled-axum" +version = "0.0.1" +dependencies = [ + "axum", + "bobbin-runtime", + "jacquard-axum", + "jacquard-common", + "jacquard-identity", + "tokio", +] + [[package]] name = "tantivy" version = "0.26.1" @@ -8977,7 +9021,6 @@ dependencies = [ "futures-core", "pin-project-lite", "tokio", - "tokio-util", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml index d5b928654..76c0bc733 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,7 +1,7 @@ [workspace] resolver = "3" -members = ["bobbin/crates/*", "crates/*", "shuttle", "knot2/crates/*", "gitmirror"] -default-members = ["bobbin/crates/*", "crates/*", "knot2/crates/*", "gitmirror"] +members = ["bobbin/crates/*", "crates/*", "shuttle", "knot2/crates/*", "gitmirror/crates/*"] +default-members = ["bobbin/crates/*", "crates/*", "knot2/crates/*", "gitmirror/crates/*"] exclude = ["sites"] [workspace.package] @@ -31,6 +31,7 @@ print_stderr = "warn" [workspace.dependencies] knot-capability = { path = "crates/knot-capability" } lexicons = { path = "crates/lexicons" } +tangled-axum = { path = "crates/tangled-axum" } trusted-proxies = { path = "crates/trusted-proxies" } gix-rebase = { path = "crates/gix-rebase" } gix-receive-pack = { path = "crates/gix-receive-pack" } @@ -81,6 +82,10 @@ knot-maintenance = { path = "knot2/crates/knot-maintenance" } knot-sim = { path = "knot2/crates/knot-sim" } knot-edge = { path = "knot2/crates/knot-edge" } +gitmirror-git = { path = "gitmirror/crates/gitmirror-git" } +gitmirror-ingest = { path = "gitmirror/crates/gitmirror-ingest" } +gitmirror-xrpc = { path = "gitmirror/crates/gitmirror-xrpc" } + jacquard-api = "0.12.1" jacquard-axum = "0.12.1" jacquard-common = "0.12.1" diff --git a/crates/gix-receive-pack/src/lib.rs b/crates/gix-receive-pack/src/lib.rs index 0d9c437ad..d0866f805 100644 --- a/crates/gix-receive-pack/src/lib.rs +++ b/crates/gix-receive-pack/src/lib.rs @@ -1,40 +1,3 @@ -//! The client side of `git-receive-pack`: turn a set of reference updates plus a packfile into the -//! request a git server expects, and read back what it did. -//! -//! ### Layering -//! -//! ```text -//! caller (gitmirror) implements gix_transport::client::async_io::Transport over its own HTTP -//! client; builds the packfile from the advertisement this crate exposes -//! | -//! gix-transport handshake(): info/refs GET, content-type check, `# service=` validation, -//! redirect re-basing, Capabilities parsing -//! | -//! THIS CRATE advertisement -> ReceivePackAdvertisement, capability selection, -//! command-block encoding, report-status parsing -//! ``` -//! -//! Object walking and pack generation are deliberately *not* here. They are the only CPU-bound -//! stage, they are already interruptible (`gix_pack::data::output::count::objects` takes a -//! `should_interrupt`), and they belong wherever the caller can bound them — a blocking pool, its -//! own cancellation. Everything in this crate is either pure byte work or an `await` on the -//! transport, so it is safe to run inline on an async task. -//! -//! ### Staging -//! -//! [`Session::connect`] performs the handshake; the caller then reads [`Session::advertised`] to -//! decide its updates and build its pack; [`Session::upload_pack`] sends them. The split is a data -//! dependency, not a style: `expected_old` comes from the advertised refs, and the pack's exclusion -//! set comes from the advertised tips plus `.have` lines, so neither exists before the handshake. -//! Staging it this way also puts the caller's CPU work *between* two awaits, where `spawn_blocking` -//! and cancellation are trivial and stay out of this crate. -//! -//! ### Protocol version -//! -//! v0 only. Protocol v2 defined `ls-refs`, `fetch` and `object-info` and no push command, so -//! `git-receive-pack` answers v0 regardless of what the client asks for. Construct the transport -//! with `Protocol::V1`, which is the only value that sends no `Git-Protocol` header at all. - pub mod advertisement; pub mod command; pub mod report; diff --git a/crates/tangled-axum/Cargo.toml b/crates/tangled-axum/Cargo.toml new file mode 100644 index 000000000..1acd971c5 --- /dev/null +++ b/crates/tangled-axum/Cargo.toml @@ -0,0 +1,16 @@ +[package] +name = "tangled-axum" +version.workspace = true +edition.workspace = true +license.workspace = true +rust-version.workspace = true + +[dependencies] +axum = { workspace = true } +bobbin-runtime = { workspace = true } +jacquard-axum = { workspace = true } +jacquard-common = { workspace = true } +jacquard-identity = { workspace = true } + +[dev-dependencies] +tokio = { workspace = true, features = ["macros", "rt"] } diff --git a/crates/tangled-axum/src/atproto.rs b/crates/tangled-axum/src/atproto.rs new file mode 100644 index 000000000..055466d8f --- /dev/null +++ b/crates/tangled-axum/src/atproto.rs @@ -0,0 +1,181 @@ +use std::future::Future; + +use axum::extract::FromRequestParts; +use axum::http::header; +use axum::http::request::Parts; +use bobbin_runtime::UnixMicros; +use jacquard_axum::service_auth::{ + ReplayKey, ReplayStore, ServiceAuth, ServiceAuthConfig, ServiceAuthError, +}; +use jacquard_common::IntoStatic as _; +use jacquard_common::SmolStr; +use jacquard_common::service_auth; +use jacquard_common::types::crypto::KeyCodec; +use jacquard_common::types::string::{Did, DidService}; +use jacquard_identity::resolver::IdentityResolver; + +pub trait AtprotoService: Send + Sync { + type Resolver: IdentityResolver + Send + Sync; + + fn resolver(&self) -> &Self::Resolver; + + fn verify_claims( + &self, + claims: &service_auth::ServiceAuthClaims, + ) -> impl Future> + Send; +} + +/// Proof that claims passed a [`ClaimsPolicy`]. +#[derive(Debug)] +pub struct Verified(()); + +pub struct ClaimsPolicy<'a> { + audience: Option>, + allowed_services: &'a [SmolStr], + require_lxm: bool, + replay: Option<&'a dyn ReplayStore>, +} + +impl<'a> ClaimsPolicy<'a> { + pub fn skip_audience(mut self) -> Self { + self.audience = None; + self + } + + pub fn require_lxm(mut self, require: bool) -> Self { + self.require_lxm = require; + self + } + + pub async fn check( + &self, + claims: &service_auth::ServiceAuthClaims, + now: UnixMicros, + ) -> Result { + if let Some(expected) = self.audience.as_ref() { + if claims.aud.audience().as_str() != expected.as_str() { + return Err(service_auth::ServiceAuthError::AudienceMismatch { + expected: expected.clone().into_static(), + actual: DidService::new_owned(claims.aud.as_str()).unwrap(), + } + .into()); + } + + if !self.allowed_services.is_empty() + && let Some(service) = claims.aud.service() + && !self + .allowed_services + .iter() + .any(|allowed| allowed.as_str() == service) + { + return Err(service_auth::ServiceAuthError::ServiceIdMismatch { + allowed: self.allowed_services.to_vec(), + actual: Some(SmolStr::new(service)), + } + .into()); + } + } + + let now = i64::try_from(now.raw() / 1_000_000).unwrap_or(i64::MAX); + if claims.exp <= now { + return Err(service_auth::ServiceAuthError::Expired { + exp: claims.exp, + now, + } + .into()); + } + + if self.require_lxm && claims.lxm.is_none() { + return Err(ServiceAuthError::MethodBindingRequired); + } + + if let Some(store) = self.replay { + let jti = claims.jti.as_ref().ok_or(ServiceAuthError::MissingJti)?; + let key = ReplayKey::new( + claims.iss.clone().into_static(), + claims.aud.clone().into_static(), + jti.clone(), + ); + store.check_and_insert(key, claims.exp).await?; + } + + Ok(Verified(())) + } +} + +impl<'a, R: IdentityResolver> From<&'a ServiceAuthConfig> for ClaimsPolicy<'a> { + fn from(config: &'a ServiceAuthConfig) -> Self { + Self { + audience: Some(ServiceAuth::service_did(config)), + allowed_services: ServiceAuth::allowed_services(config), + require_lxm: ServiceAuth::require_lxm(config), + replay: ServiceAuth::replay_protection_enabled(config) + .then(|| ServiceAuth::replay_store(config)), + } + } +} + +pub async fn resolve_signing_key( + resolver: &R, + iss: &Did, +) -> Result +where + R: IdentityResolver + Sync, +{ + let doc = resolver.resolve_did_doc_owned(iss).await.map_err(|e| { + ServiceAuthError::DidResolutionFailed { + did: iss.clone().into_static(), + source: Box::new(e), + } + })?; + + let signing_key = doc + .atproto_public_key() + .map_err(|e| ServiceAuthError::InvalidKey(e.to_string()))? + .ok_or_else(|| ServiceAuthError::NoSigningKey(iss.clone().into_static())) + .map(|key| match key.codec { + KeyCodec::P256 => Ok(service_auth::PublicKey::from_p256_bytes(&key.bytes)), + KeyCodec::Secp256k1 => Ok(service_auth::PublicKey::from_k256_bytes(&key.bytes)), + codec => { + Err(ServiceAuthError::InvalidKey(format!( + "{codec:?} does not sign service auth tokens" + ))) + } + })???; + + Ok(signing_key) +} + +/// [`jacquard_axum::service_auth::ExtractServiceAuth`] with the claims policy delegated +/// to the state, plus the raw token. +pub struct ExtractAnyServiceAuth(pub service_auth::ServiceAuthClaims, pub String); + +impl FromRequestParts for ExtractAnyServiceAuth { + type Rejection = ServiceAuthError; + + async fn from_request_parts(parts: &mut Parts, state: &S) -> Result { + let token = bearer_token_from_parts(parts)?.ok_or(ServiceAuthError::MissingAuthHeader)?; + let parsed = service_auth::parse_jwt(token)?; + + let signing_key = resolve_signing_key(state.resolver(), &parsed.claims().iss).await?; + service_auth::verify_signature(&parsed, &signing_key)?; + + let Verified(()) = state.verify_claims(parsed.claims()).await?; + + Ok(Self(parsed.into_claims(), token.to_owned())) + } +} + +fn bearer_token_from_parts(parts: &Parts) -> Result, ServiceAuthError> { + let Some(auth_header) = parts.headers.get(header::AUTHORIZATION) else { + return Ok(None); + }; + + let auth_str = auth_header + .to_str() + .map_err(|_| ServiceAuthError::InvalidAuthHeader)?; + let token = auth_str + .strip_prefix("Bearer ") + .ok_or(ServiceAuthError::InvalidAuthHeader)?; + Ok(Some(token)) +} diff --git a/crates/tangled-axum/src/lib.rs b/crates/tangled-axum/src/lib.rs new file mode 100644 index 000000000..224322fa0 --- /dev/null +++ b/crates/tangled-axum/src/lib.rs @@ -0,0 +1 @@ +pub mod atproto; diff --git a/gitmirror/crates/gitmirror-git/Cargo.toml b/gitmirror/crates/gitmirror-git/Cargo.toml new file mode 100644 index 000000000..d03db00e5 --- /dev/null +++ b/gitmirror/crates/gitmirror-git/Cargo.toml @@ -0,0 +1,18 @@ +[package] +name = "gitmirror-git" +version.workspace = true +edition.workspace = true +license.workspace = true +rust-version.workspace = true + +[dependencies] +gix-rebase = { workspace = true } + +gix = { workspace = true } +tempfile = { workspace = true } +thiserror = { workspace = true } +itertools = { workspace = true } +jacquard-common = { workspace = true } + +[lints] +workspace = true diff --git a/gitmirror/crates/gitmirror-git/src/layout.rs b/gitmirror/crates/gitmirror-git/src/layout.rs new file mode 100644 index 000000000..7c5245a9c --- /dev/null +++ b/gitmirror/crates/gitmirror-git/src/layout.rs @@ -0,0 +1,61 @@ +use std::path::PathBuf; + +use itertools::Itertools as _; +use jacquard_common::types::did::Did; + +use crate::scratch::TempRepository; + +#[derive(Debug, Clone)] +pub struct Layout { + scan_path: PathBuf, +} + +#[derive(Debug, thiserror::Error)] +pub enum Error { + #[error("repo not found")] + RepoNotFound, + #[error("open_scratch requires at least one repo")] + NotEnoughRepos, + #[error("scratch repo: {0}")] + Io(#[from] std::io::Error), + #[error("init scratch repo: {0}")] + Init(#[from] gix::init::Error), + #[error("open scratch repo: {0}")] + Open(#[from] gix::open::Error), +} + +impl Layout { + pub fn new(scan_path: PathBuf) -> Self { + Self { scan_path } + } + + pub fn repo_path(&self, repo: &Did) -> PathBuf { + // TODO: normalize DID like knot_git::Layout does. + self.scan_path.join(repo.as_str()) + } + + pub fn open_scratch(&self, dids: &[&Did]) -> Result { + if dids.is_empty() { + return Err(Error::NotEnoughRepos); + } + let alternates = dids + .iter() + .copied() + .unique() + .map(|did| { + std::fs::canonicalize(self.repo_path(did).join("objects")) + .map(|dir| format!("{}\n", dir.display())) + .map_err(|_| Error::RepoNotFound) + }) + .collect::>()?; + + let scratch = tempfile::tempdir()?; + gix::init_bare(scratch.path())?; + + let info_dir = scratch.path().join("objects").join("info"); + std::fs::create_dir_all(&info_dir)?; + std::fs::write(info_dir.join("alternates"), alternates)?; + + Ok(TempRepository::open(scratch)?) + } +} diff --git a/gitmirror/crates/gitmirror-git/src/lib.rs b/gitmirror/crates/gitmirror-git/src/lib.rs new file mode 100644 index 000000000..0ef9a130d --- /dev/null +++ b/gitmirror/crates/gitmirror-git/src/lib.rs @@ -0,0 +1,3 @@ +pub mod layout; +pub mod repo_ext; +pub mod scratch; diff --git a/gitmirror/src/git_repo_ext.rs b/gitmirror/crates/gitmirror-git/src/repo_ext.rs similarity index 95% rename from gitmirror/src/git_repo_ext.rs rename to gitmirror/crates/gitmirror-git/src/repo_ext.rs index 7d0cc46bd..2deaf41e8 100644 --- a/gitmirror/src/git_repo_ext.rs +++ b/gitmirror/crates/gitmirror-git/src/repo_ext.rs @@ -1,10 +1,9 @@ -//! Merge operations layered onto [`gix::Repository`], so they read as repository methods rather -//! than as free functions taking one. +//! Extended [`gix::Repository`]. use std::collections::{HashMap, HashSet}; #[derive(Debug, thiserror::Error)] -pub(crate) enum MergeError { +pub enum MergeError { /// Not a failure so much as an answer: these paths couldn't be reconciled, and nothing /// reachable was written. #[error("the merge conflicted in {}", .0.join(", "))] @@ -21,7 +20,7 @@ pub(crate) enum MergeError { Replay(#[from] gix_rebase::replay::Error), } -pub(crate) trait RepositoryExt { +pub trait RepositoryExt { /// Replay `source` and the ancestors of it that `onto` doesn't already have on top of `onto`. fn rebase( &self, diff --git a/gitmirror/crates/gitmirror-git/src/scratch.rs b/gitmirror/crates/gitmirror-git/src/scratch.rs new file mode 100644 index 000000000..56c0aff8b --- /dev/null +++ b/gitmirror/crates/gitmirror-git/src/scratch.rs @@ -0,0 +1,30 @@ +use tempfile::TempDir; + +/// A [`gix::Repository`] backed by a temporary directory that is removed on drop. Derefs to +/// `gix::Repository`, so it is used just like one; keeping it alive keeps the scratch dir alive. +pub struct TempRepository { + repo: gix::Repository, + // Declared AFTER `repo` so `repo` drops first: any file handles into the scratch dir close + // before the dir itself is removed (Rust drops struct fields in declaration order). + _dir: TempDir, +} + +impl TempRepository { + pub fn open(dir: TempDir) -> Result { + let repo = gix::open(dir.path())?; + Ok(Self { repo, _dir: dir }) + } + + #[allow(dead_code)] + pub fn with_object_memory(mut self) -> Self { + self.repo.objects.enable_object_memory(); + self + } +} + +impl std::ops::Deref for TempRepository { + type Target = gix::Repository; + fn deref(&self) -> &gix::Repository { + &self.repo + } +} diff --git a/gitmirror/crates/gitmirror-ingest/Cargo.toml b/gitmirror/crates/gitmirror-ingest/Cargo.toml new file mode 100644 index 000000000..6c1caf617 --- /dev/null +++ b/gitmirror/crates/gitmirror-ingest/Cargo.toml @@ -0,0 +1,11 @@ +[package] +name = "gitmirror-ingest" +version.workspace = true +edition.workspace = true +license.workspace = true +rust-version.workspace = true + +[dependencies] + +[lints] +workspace = true diff --git a/gitmirror/crates/gitmirror-ingest/src/lib.rs b/gitmirror/crates/gitmirror-ingest/src/lib.rs new file mode 100644 index 000000000..e69de29bb diff --git a/gitmirror/Cargo.toml b/gitmirror/crates/gitmirror-xrpc/Cargo.toml similarity index 68% rename from gitmirror/Cargo.toml rename to gitmirror/crates/gitmirror-xrpc/Cargo.toml index 71d8f93f8..a9a6aef57 100644 --- a/gitmirror/Cargo.toml +++ b/gitmirror/crates/gitmirror-xrpc/Cargo.toml @@ -1,44 +1,38 @@ [package] -name = "gitmirror" -version = "0.1.0" +name = "gitmirror-xrpc" +version.workspace = true edition.workspace = true license.workspace = true rust-version.workspace = true [dependencies] -anyhow = { workspace = true } -axum = { workspace = true } bobbin-knot-proxy = { workspace = true } bobbin-runtime = { workspace = true } -bobbin-types = { workspace = true } +gitmirror-git = { workspace = true } +gix-receive-pack = { workspace = true } +lexicons = { workspace = true } +tangled-axum = { workspace = true } + +anyhow = { workspace = true } +axum = { workspace = true } chrono = { workspace = true } -clap = { workspace = true } jacquard-axum = { workspace = true } jacquard-common = { workspace = true } jacquard-identity = { workspace = true } -serde = { workspace = true } serde_json = { workspace = true } tokio = { workspace = true, features = ["net", "rt-multi-thread", "sync"] } -tokio-stream = { workspace = true, features = ["sync"] } tracing = "0.1" -tracing-subscriber = { version = "0.3", features = ["env-filter"] } line-numbers = "0.4.0" rustc-hash = "2.1.2" gix = { version = "0.84", features = ["parallel", "blob-diff", "merge", "sha1", "sha256", "revision", "tree-editor"] } gix-pack = { workspace = true } -gix-receive-pack = { workspace = true } gix-transport = { workspace = true, features = ["async-client"] } -gix-rebase = { workspace = true } futures-lite = { workspace = true } reqwest = { workspace = true } -tempfile = "3" -thiserror = { workspace = true } [dev-dependencies] -base64 = { workspace = true } -bytes = { workspace = true } -k256 = { workspace = true } -multibase = "0.9" +tempfile = "3" +tokio = { workspace = true, features = ["macros", "io-util"] } [lints] workspace = true diff --git a/gitmirror/src/repo_did.rs b/gitmirror/crates/gitmirror-xrpc/src/did_ext.rs similarity index 100% rename from gitmirror/src/repo_did.rs rename to gitmirror/crates/gitmirror-xrpc/src/did_ext.rs diff --git a/gitmirror/src/diff.rs b/gitmirror/crates/gitmirror-xrpc/src/diff.rs similarity index 100% rename from gitmirror/src/diff.rs rename to gitmirror/crates/gitmirror-xrpc/src/diff.rs diff --git a/gitmirror/crates/gitmirror-xrpc/src/error.rs b/gitmirror/crates/gitmirror-xrpc/src/error.rs new file mode 100644 index 000000000..011d6dd24 --- /dev/null +++ b/gitmirror/crates/gitmirror-xrpc/src/error.rs @@ -0,0 +1,91 @@ +use axum::{Json, response::{IntoResponse, Response}}; +use reqwest::StatusCode; +use serde_json::json; +use tracing::error; + +#[derive(Debug)] +pub(crate) enum XrpcError { + InvalidRequest(String), + RepoNotFound { detail: String }, + RefNotFound { rev: String, detail: String }, + RevisionNotFound { rev: String }, + CompareError(String), + MergeConflict(Vec), + PushRejected(String), + Unauthorized(String), + Internal(String), +} + +impl IntoResponse for XrpcError { + fn into_response(self) -> Response { + let (status, error, message) = match self { + Self::InvalidRequest(m) => (StatusCode::BAD_REQUEST, "InvalidRequest", m), + Self::RepoNotFound { detail } => { + error!(error = %detail, "repo not found"); + ( + StatusCode::NOT_FOUND, + "RepoNotFound", + "repository not found".to_owned(), + ) + } + Self::RefNotFound { rev, detail } => { + error!(error = %detail, rev = %rev, "revision not found"); + ( + StatusCode::NOT_FOUND, + "RefNotFound", + format!("revision not found: {rev}"), + ) + } + Self::RevisionNotFound { rev } => ( + StatusCode::NOT_FOUND, + "RevisionNotFound", + format!("commit not found: {rev}"), + ), + Self::CompareError(m) => { + error!(error = %m, "compare failed"); + ( + StatusCode::INTERNAL_SERVER_ERROR, + "CompareError", + "failed to compare revisions".to_owned(), + ) + } + Self::MergeConflict(paths) => ( + StatusCode::CONFLICT, + "MergeConflict", + format!("merge produced conflicts in {}", paths.join(", ")), + ), + Self::PushRejected(detail) => { + error!(error = %detail, "push rejected"); + ( + StatusCode::CONFLICT, + "PushRejected", + "knot rejected the push".to_owned(), + ) + } + Self::Unauthorized(detail) => { + error!(error = %detail, "authentication failed"); + (StatusCode::UNAUTHORIZED, "AuthenticationRequired", detail) + } + Self::Internal(m) => { + error!(error = %m, "xrpc request failed"); + ( + StatusCode::INTERNAL_SERVER_ERROR, + "InternalServerError", + "internal error".to_owned(), + ) + } + }; + (status, Json(json!({ "error": error, "message": message }))).into_response() + } +} + +impl From for XrpcError { + fn from(e: gitmirror_git::layout::Error) -> Self { + match e { + gitmirror_git::layout::Error::RepoNotFound => Self::RepoNotFound { + detail: e.to_string(), + }, + _ => Self::Internal(e.to_string()), + } + } +} diff --git a/gitmirror/src/git_transport.rs b/gitmirror/crates/gitmirror-xrpc/src/git_transport.rs similarity index 100% rename from gitmirror/src/git_transport.rs rename to gitmirror/crates/gitmirror-xrpc/src/git_transport.rs diff --git a/gitmirror/crates/gitmirror-xrpc/src/lib.rs b/gitmirror/crates/gitmirror-xrpc/src/lib.rs new file mode 100644 index 000000000..58485aa2a --- /dev/null +++ b/gitmirror/crates/gitmirror-xrpc/src/lib.rs @@ -0,0 +1,91 @@ +mod did_ext; +mod diff; +mod error; +mod git_transport; +mod merge; +mod routes; + +use std::{path::PathBuf, sync::Arc}; + +use axum::{ + routing::{get, post}, + Router, +}; +use bobbin_runtime::{Clock, ReqwestHttp}; +use tangled_axum::atproto::{AtprotoService, ClaimsPolicy, Verified}; +use gitmirror_git::layout::Layout; +use jacquard_axum::service_auth::{self, ServiceAuthError}; +use jacquard_common::service_auth::ServiceAuthClaims; +use jacquard_identity::JacquardResolver; + +pub(crate) type Directory = JacquardResolver; + +#[derive(Clone)] +pub struct AppState { + pub(crate) layout: Arc, + pub(crate) http: reqwest::Client, + pub(crate) service_auth: service_auth::ServiceAuthConfig>, + pub(crate) knot_policy: KnotPolicy, + pub(crate) clock: Arc, +} + +impl AppState { + #[allow(clippy::too_many_arguments)] + pub fn new( + repo_base: Arc, + http: reqwest::Client, + service_auth: service_auth::ServiceAuthConfig>, + knot_policy: KnotPolicy, + clock: Arc, + ) -> Self { + let layout = Arc::new(Layout::new((*repo_base).clone())); + Self { + layout, + http, + service_auth, + knot_policy, + clock, + } + } +} + +impl AtprotoService for AppState { + type Resolver = Directory; + + fn resolver(&self) -> &Self::Resolver { + self.service_auth.resolver() + } + + async fn verify_claims(&self, claims: &ServiceAuthClaims) -> Result { + ClaimsPolicy::from(&self.service_auth) + .skip_audience() + .check(claims, self.clock.now_unix_micros()) + .await + } +} + +pub fn router(state: AppState) -> Router { + use routes::*; + Router::new() + .route("/xrpc/sh.tangled.git.temp2.getDiff", get(git::get_diff)) + .route( + "/xrpc/sh.tangled.git.temp2.getInterdiff", + get(git::get_interdiff), + ) + .route( + "/xrpc/sh.tangled.git.temp2.listCommits", + get(git::list_commits), + ) + .route( + "/xrpc/sh.tangled.git.temp2.mergeCheck", + get(git::merge_check), + ) + .route("/xrpc/sh.tangled.git.mergeCommit", post(git::merge_commit)) + .with_state(state) +} + +#[derive(Debug, Clone, Copy)] +pub struct KnotPolicy { + pub allow_private: bool, + pub require_https: bool, +} diff --git a/gitmirror/src/merge.rs b/gitmirror/crates/gitmirror-xrpc/src/merge.rs similarity index 99% rename from gitmirror/src/merge.rs rename to gitmirror/crates/gitmirror-xrpc/src/merge.rs index 5135e488c..24656d65f 100644 --- a/gitmirror/src/merge.rs +++ b/gitmirror/crates/gitmirror-xrpc/src/merge.rs @@ -2,7 +2,7 @@ use gix::bstr::ByteSlice as _; use gix::diff::tree_with_rewrites::Change; use gix::merge::tree::{Conflict, TreatAsUnresolved}; -use bobbin_types::sh_tangled::git::temp2::merge_check::{ +use lexicons::sh_tangled::git::temp2::merge_check::{ Conflict as MergeConflict, MergeCheckOutput, }; diff --git a/gitmirror/src/xrpc_git.rs b/gitmirror/crates/gitmirror-xrpc/src/routes/git.rs similarity index 59% rename from gitmirror/src/xrpc_git.rs rename to gitmirror/crates/gitmirror-xrpc/src/routes/git.rs index 400100188..509ca84e2 100644 --- a/gitmirror/src/xrpc_git.rs +++ b/gitmirror/crates/gitmirror-xrpc/src/routes/git.rs @@ -1,15 +1,26 @@ use axum::extract::State; use bobbin_knot_proxy::KnotHost; -use bobbin_types::sh_tangled::git::merge_commit; -use jacquard_axum::service_auth::ServiceAuth; +use gitmirror_git::repo_ext::{MergeError, RepositoryExt as _}; +use gix::ObjectId; +use gix::bstr::ByteSlice as _; +use gix::revision::plumbing::Spec as RevSpec; +use gix::revision::walk::Sorting; use jacquard_axum::{ExtractXrpc, XrpcResponse}; +use jacquard_common::ToSmolStr; +use jacquard_common::types::string::Datetime; +use lexicons::sh_tangled; +use lexicons::sh_tangled::git::{ + merge_commit, temp2::{get_diff, get_interdiff, list_commits, merge_check} +}; use jacquard_identity::resolver::IdentityResolver as _; +use tangled_axum::atproto::{AtprotoService as _, ExtractAnyServiceAuth}; use tracing::info; -use crate::auth::ExtractAnyServiceAuth; -use crate::git_repo_ext::{MergeError, RepositoryExt as _}; -use crate::repo_did::DidDocumentExt as _; -use crate::xrpc::{find_commit, XrpcError, XrpcState}; +use crate::did_ext::DidDocumentExt as _; +use crate::{AppState, KnotPolicy, diff, error::XrpcError}; + +const DEFAULT_LIMIT: u32 = 50; +const MAX_LIMIT: u32 = 100; const AGENT: &str = concat!("gitmirror/", env!("CARGO_PKG_VERSION")); @@ -36,8 +47,360 @@ impl MergeStyle { } } +struct GixSignature<'a>(gix::actor::SignatureRef<'a>); + +impl TryFrom> for sh_tangled::git::Signature { + type Error = anyhow::Error; + + fn try_from(GixSignature(sig): GixSignature) -> Result { + let time = sig.time()?; + let offset = chrono::FixedOffset::east_opt(time.offset).ok_or_else(|| { + anyhow::anyhow!("commit timezone offset out of range: {}", time.offset) + })?; + let when = chrono::DateTime::from_timestamp(time.seconds, 0) + .ok_or_else(|| anyhow::anyhow!("commit timestamp out of range: {}", time.seconds))? + .with_timezone(&offset); + Ok(Self { + name: sig.name.to_smolstr(), + email: sig.email.to_smolstr(), + when: Datetime::new(when), + extra_data: Default::default(), + }) + } +} + +struct GixCommit<'a>(gix::Commit<'a>); + +impl TryFrom> for list_commits::Commit { + type Error = anyhow::Error; + + fn try_from(GixCommit(commit): GixCommit<'_>) -> Result { + let decoded = commit.decode()?; + Ok(Self { + oid: commit.id.to_smolstr(), + parents: decoded + .parents + .iter() + .map(|parent| parent.to_smolstr()) + .collect(), + tree: decoded.tree.to_smolstr(), + author: GixSignature(decoded.author()?).try_into()?, + committer: GixSignature(decoded.committer()?).try_into()?, + extra_headers: decoded + .extra_headers + .iter() + .map(|(key, value)| list_commits::Header { + key: key.to_smolstr(), + value: value.to_smolstr(), + extra_data: Default::default(), + }) + .collect(), + message: decoded.message.to_smolstr(), + extra_data: Default::default(), + }) + } +} + +impl From for sh_tangled::git::DiffSrc { + fn from(f: crate::diff::FileContent) -> Self { + Self { + path: f.path.into(), + oid: f.oid.into(), + size: f.size as i64, + is_binary: f.is_binary, + is_submodule: f.is_submodule, + content: f + .content + .map(|b| String::from_utf8_lossy(&b).into_owned().into()), + extra_data: Default::default(), + } + } +} + +impl From for sh_tangled::git::DiffHunk { + fn from(h: crate::diff::Hunk) -> Self { + // The novel sets are hash sets; sort them so the output is stable across runs. + let sorted = |set: rustc_hash::FxHashSet| -> Vec { + let mut v: Vec = set.into_iter().map(|n| i64::from(n.0)).collect(); + v.sort_unstable(); + v + }; + Self { + novel_lhs: sorted(h.novel_lhs), + novel_rhs: sorted(h.novel_rhs), + lines: h + .lines + .into_iter() + .map(|(lhs, rhs)| sh_tangled::git::LinePair { + lhs: lhs.map(|n| i64::from(n.0)), + rhs: rhs.map(|n| i64::from(n.0)), + extra_data: Default::default(), + }) + .collect(), + extra_data: Default::default(), + } + } +} + +impl From for sh_tangled::git::FileDiff { + fn from(d: crate::diff::Diff) -> Self { + Self { + lhs_src: d.lhs_src.into(), + rhs_src: d.rhs_src.into(), + hunks: d.hunks.into_iter().map(Into::into).collect(), + has_byte_changes: d + .has_byte_changes + .map(|(lhs, rhs)| sh_tangled::git::ByteChanges { + lhs: lhs as i64, + rhs: rhs as i64, + extra_data: Default::default(), + }), + has_syntactic_changes: d.has_syntactic_changes, + extra_data: Default::default(), + } + } +} + +/// Resolve a full hex oid to a commit that actually exists in `repo`. +pub(crate) fn find_commit(repo: &gix::Repository, sha: &str) -> Result { + let oid = gix::ObjectId::from_hex(sha.as_bytes()) + .map_err(|e| XrpcError::InvalidRequest(format!("bad commit sha {sha:?}: {e}")))?; + repo.find_commit(oid) + .map_err(|_| XrpcError::RevisionNotFound { + rev: sha.to_owned(), + })?; + Ok(oid) +} + +fn get_diff_inner( + repo: &gix::Repository, + base: gix::ObjectId, + head: gix::ObjectId, +) -> Result, XrpcError> { + let compare = || -> anyhow::Result> { + let merge_base = repo.merge_base(base, head)?.detach(); + let old = repo.find_tree(repo.find_commit(merge_base)?.tree_id()?)?; + let new = repo.find_tree(repo.find_commit(head)?.tree_id()?)?; + diff::diff(repo, &old, &new, false)? + .map(|d| d.map(Into::into)) + .collect() + }; + compare().map_err(|e| XrpcError::CompareError(e.to_string())) +} + +#[axum::debug_handler] +pub(crate) async fn get_diff( + State(state): State, + ExtractXrpc(args): ExtractXrpc, +) -> Result, XrpcError> { + let scratch = state + .layout + .open_scratch(&[&args.head_repo, &args.base_repo])?; + let base = find_commit(&scratch, args.base_commit.as_ref())?; + let head = find_commit(&scratch, args.head_commit.as_ref())?; + + tokio::task::spawn_blocking(move || get_diff_inner(&scratch, base, head)) + .await + .map_err(|e| XrpcError::Internal(e.to_string()))? + .map(|diffs| { + XrpcResponse(get_diff::GetDiffOutput { + diffs, + extra_data: Default::default(), + }) + }) +} + +fn get_interdiff_inner( + repo: &gix::Repository, + (from_base_id, from_head_id): (ObjectId, ObjectId), + (to_base_id, to_head_id): (ObjectId, ObjectId), +) -> Result, XrpcError> { + let to_head = repo + .find_commit(to_head_id) + .map_err(|e| XrpcError::Internal(e.to_string()))?; + let to_head_tree = to_head + .tree() + .map_err(|e| XrpcError::Internal(e.to_string()))?; + let rebased_tree = diff::prepare_interdiff(repo, (from_base_id, from_head_id), to_base_id) + .map_err(|e| XrpcError::Internal(e.to_string()))?; + let compare = || -> anyhow::Result> { + diff::diff(repo, &rebased_tree, &to_head_tree, true)? + .map(|d| d.map(Into::into)) + .collect() + }; + compare().map_err(|e| XrpcError::CompareError(e.to_string())) +} + +#[axum::debug_handler] +pub(crate) async fn get_interdiff( + State(state): State, + ExtractXrpc(args): ExtractXrpc, +) -> Result, XrpcError> { + let from_base_id = ObjectId::from_hex(args.base_commit1.as_bytes()) + .map_err(|e| XrpcError::InvalidRequest(e.to_string()))?; + let from_head_id = ObjectId::from_hex(args.base_commit2.as_bytes()) + .map_err(|e| XrpcError::InvalidRequest(e.to_string()))?; + let to_base_id = ObjectId::from_hex(args.head_commit1.as_bytes()) + .map_err(|e| XrpcError::InvalidRequest(e.to_string()))?; + let to_head_id = ObjectId::from_hex(args.head_commit2.as_bytes()) + .map_err(|e| XrpcError::InvalidRequest(e.to_string()))?; + + let scratch = state + .layout + .open_scratch(&[&args.head_repo, &args.base_repo])?; + + tokio::task::spawn_blocking(move || { + get_interdiff_inner( + &scratch, + (from_base_id, from_head_id), + (to_base_id, to_head_id), + ) + }) + .await + .map_err(|e| XrpcError::Internal(e.to_string()))? + .map(|diffs| { + XrpcResponse(get_interdiff::GetInterdiffOutput { + diffs, + extra_data: Default::default(), + }) + }) +} + +fn commit_log_walk<'repo>( + repo: &'repo gix::Repository, + ranges: &[Vec], + all_refs: bool, +) -> anyhow::Result> { + let mut tips = Vec::new(); + let mut hidden = Vec::new(); + + if all_refs { + for r in repo.references()?.all()? { + let mut r = r.map_err(|e| anyhow::anyhow!(e))?; + if let Ok(commit) = r.peel_to_commit() { + tips.push(commit.id); + } + } + } else { + for range in ranges { + let revspec = repo.rev_parse(range.as_bstr())?; + let spec = revspec.detach(); + match spec { + RevSpec::Include(id) => tips.push(id), + RevSpec::Range { from, to } => { + tips.push(to); + hidden.push(from); + } + _ => { + anyhow::bail!("The spec isn't currently supported: {spec:?}") + } + } + } + } + + Ok(repo + .rev_walk(tips) + .sorting(Sorting::ByCommitTime(Default::default())) + .with_hidden(hidden) + .all()?) +} + +fn list_commits_inner( + repo: &gix::Repository, + args: list_commits::ListCommits, +) -> Result, XrpcError> { + let ranges: Vec> = args + .ranges + .clone() + .unwrap_or_default() + .iter() + .map(|revspec| revspec.as_bytes().to_vec()) + .collect(); + let walk = commit_log_walk(repo, &ranges, args.all_refs.unwrap_or(false)).map_err(|e| { + XrpcError::RefNotFound { + rev: args.ranges.unwrap_or_default().join(", "), + detail: e.to_string(), + } + })?; + + let limit = args + .limit + .map(|limit| limit as u32) + .unwrap_or(DEFAULT_LIMIT); + if limit == 0 || limit > MAX_LIMIT { + return Err(XrpcError::InvalidRequest(format!( + "limit must be between 1 and {MAX_LIMIT}" + ))); + } + + let mut commits: Vec = Vec::with_capacity(limit as usize); + let mut skipped = 0usize; + + for info in walk { + let info = info.map_err(|e| XrpcError::Internal(e.to_string()))?; + + if (skipped as i64) < args.skip.unwrap_or(0) { + skipped += 1; + continue; + } + if (commits.len() as i64) == args.limit.unwrap_or(50) { + break; + } + + let commit = info + .object() + .map_err(|e| XrpcError::Internal(e.to_string()))?; + commits.push( + list_commits::Commit::try_from(GixCommit(commit)) + .map_err(|e| XrpcError::Internal(e.to_string()))?, + ); + } + + Ok(commits) +} + +#[axum::debug_handler] +pub(crate) async fn list_commits( + State(state): State, + ExtractXrpc(args): ExtractXrpc, +) -> Result, XrpcError> { + let path = state.layout.repo_path(&args.repo); + let repo = gix::open(path).map_err(|e| XrpcError::RepoNotFound { + detail: e.to_string(), + })?; + + tokio::task::spawn_blocking(move || list_commits_inner(&repo, args)) + .await + .map_err(|e| XrpcError::Internal(e.to_string()))? + .map(|commits| { + XrpcResponse(list_commits::ListCommitsOutput { + commits, + extra_data: Default::default(), + }) + }) +} + +#[axum::debug_handler] +pub(crate) async fn merge_check( + State(state): State, + ExtractXrpc(params): ExtractXrpc, +) -> Result, XrpcError> { + let scratch = state + .layout + .open_scratch(&[¶ms.target_repo, ¶ms.source_repo])?; + let target = find_commit(&scratch, params.target_commit.as_ref())?; + let source = find_commit(&scratch, params.source_commit.as_ref())?; + + tokio::task::spawn_blocking(move || crate::merge::merge_check(&scratch, target, source)) + .await + .map_err(|e| XrpcError::Internal(e.to_string()))? + .map(XrpcResponse) + .map_err(|e| XrpcError::Internal(e.to_string())) +} + +#[axum::debug_handler] pub(crate) async fn merge_commit( - State(state): State, + State(state): State, ExtractAnyServiceAuth(auth, token): ExtractAnyServiceAuth, ExtractXrpc(input): ExtractXrpc, ) -> Result, XrpcError> { @@ -68,9 +431,7 @@ pub(crate) async fn merge_commit( // 2. load target & source repos as scratch repo. let scratch = state .layout - .open_scratch(&[&input.target.repo, &input.source.repo]) - .map_err(|e| XrpcError::Internal(e.to_string()))? - .ok_or_else(|| XrpcError::RepoNotFound { detail: "repo not found".to_string() })?; + .open_scratch(&[&input.target.repo, &input.source.repo])?; let source_commit_id = find_commit(&scratch, input.source.commit.as_str())?; let target_commit_id = find_commit(&scratch, input.target.commit.as_str())?; @@ -133,7 +494,8 @@ pub(crate) async fn merge_commit( })?; // 4. push the merged tip to the target repository - let mut transport = crate::git_transport::HttpTransport::new(http, remote.clone(), token); + let mut transport = + crate::git_transport::HttpTransport::new(http, remote.clone(), token); let mut session = runtime .block_on(gix_receive_pack::Session::handshake( &mut transport, @@ -181,12 +543,6 @@ pub(crate) async fn merge_commit( Ok(XrpcResponse(())) } -#[derive(Debug, Clone, Copy)] -pub(crate) struct KnotPolicy { - pub allow_private: bool, - pub require_https: bool, -} - fn safe_knot_host(endpoint: &str, policy: KnotPolicy) -> Result { let host = KnotHost::parse(endpoint).map_err(|_| "is not a usable base URL")?; if policy.require_https && host.url().scheme() != "https" { @@ -224,14 +580,14 @@ fn knot_service_did(knot: &KnotHost) -> String { /// end and stages nothing (`knot-pack::receive::stage_pack_bytes` treats empty as no pack), so the /// body simply stops after the command block. fn build_pack( - scratch: &crate::git::TempRepository, + scratch: &gitmirror_git::scratch::TempRepository, advertised: &gix_receive_pack::ReceivePackAdvertisement, new_commit_id: gix::ObjectId, hidden: Vec, ) -> anyhow::Result>> { - use gix::parallel::{reduce::Finalize as _, InOrderIter}; + use gix::parallel::{InOrderIter, reduce::Finalize as _}; use gix::progress::Discard; - use gix_pack::data::{output, Version}; + use gix_pack::data::{Version, output}; const CHUNK_SIZE: usize = 50; @@ -328,14 +684,169 @@ fn build_pack( } #[cfg(test)] -mod tests { +mod diff_tests { + use super::*; + use gitmirror_git::layout::Layout; + use jacquard_common::types::did::Did; + use std::path::Path; + use std::process::Command; + + const BASE_DID: &str = "did:plc:upstream"; + const HEAD_DID: &str = "did:plc:fork"; + + fn git(dir: &Path, args: &[&str]) { + let status = Command::new("git") + .args(args) + .current_dir(dir) + .env("GIT_AUTHOR_NAME", "t") + .env("GIT_AUTHOR_EMAIL", "t@t") + .env("GIT_COMMITTER_NAME", "t") + .env("GIT_COMMITTER_EMAIL", "t@t") + .status() + .expect("run git"); + assert!(status.success(), "git {args:?} failed"); + } + + fn rev_parse(dir: &Path, rev: &str) -> String { + let out = Command::new("git") + .args(["rev-parse", rev]) + .current_dir(dir) + .output() + .expect("rev-parse"); + String::from_utf8(out.stdout).unwrap().trim().to_owned() + } + + /// The real fork shape: the base tip lives only in the upstream mirror and the head tip only + /// in the fork's, with a shared ancestor. Returns `(repo_base, base_commit, head_commit)`. + fn fork_fixture() -> (tempfile::TempDir, String, String) { + let root = tempfile::tempdir().unwrap(); + let repo_base = root.path().join("repos"); + std::fs::create_dir_all(&repo_base).unwrap(); + + let upstream = root.path().join("upstream"); + std::fs::create_dir_all(&upstream).unwrap(); + git(&upstream, &["init", "-q", "-b", "main"]); + std::fs::write(upstream.join("a.txt"), "line1\nline2\nline3\n").unwrap(); + std::fs::write(upstream.join("b.txt"), "keep\n").unwrap(); + git(&upstream, &["add", "."]); + git(&upstream, &["commit", "-q", "-m", "shared ancestor"]); + + // Fork before upstream moves on, so neither tip is reachable from the other. + let fork = root.path().join("fork"); + git( + root.path(), + &[ + "clone", + "-q", + upstream.to_str().unwrap(), + fork.to_str().unwrap(), + ], + ); + + std::fs::write(upstream.join("upstream.txt"), "theirs\n").unwrap(); + git(&upstream, &["add", "."]); + git(&upstream, &["commit", "-q", "-m", "upstream only"]); + let base_commit = rev_parse(&upstream, "HEAD"); + + std::fs::write(fork.join("a.txt"), "line1\nCHANGED\nline3\n").unwrap(); + git(&fork, &["mv", "b.txt", "c.txt"]); + git(&fork, &["add", "-A"]); + git(&fork, &["commit", "-q", "-m", "fork only"]); + let head_commit = rev_parse(&fork, "HEAD"); + + for (did, src) in [(BASE_DID, &upstream), (HEAD_DID, &fork)] { + git( + root.path(), + &[ + "clone", + "-q", + "--bare", + src.to_str().unwrap(), + repo_base.join(did).to_str().unwrap(), + ], + ); + } + + (root, base_commit, head_commit) + } + + fn diff_fixture( + root: &Path, + base_commit: &str, + head_commit: &str, + ) -> Result, XrpcError> { + let repo_base = root.join("repos"); + let scratch = Layout::new(repo_base) + .open_scratch(&[&Did::raw(HEAD_DID.into()), &Did::raw(BASE_DID.into())])?; + let base = find_commit(&scratch, base_commit)?; + let head = find_commit(&scratch, head_commit)?; + get_diff_inner(&scratch, base, head) + } + + #[test] + fn a_cross_repo_diff_reports_only_what_the_head_side_changed() { + let (root, base, head) = fork_fixture(); + let diffs = diff_fixture(root.path(), &base, &head).unwrap(); + + let mut paths: Vec = diffs + .iter() + .flat_map(|d| [d.lhs_src.path.to_string(), d.rhs_src.path.to_string()]) + .collect(); + paths.sort(); + paths.dedup(); + // `upstream.txt` is on the base side only: three-dot semantics exclude it. + assert_eq!(paths, ["a.txt", "b.txt", "c.txt"]); + + let a = diffs + .iter() + .find(|d| d.rhs_src.path == "a.txt") + .expect("a.txt is in the diff"); + // Only line 2 (0-based: 1) changed. + assert_eq!(a.hunks.len(), 1); + assert_eq!(a.hunks[0].novel_lhs, [1]); + assert_eq!(a.hunks[0].novel_rhs, [1]); + assert_eq!(a.hunks[0].lines, { + vec![sh_tangled::git::LinePair { + lhs: Some(1), + rhs: Some(1), + extra_data: Default::default(), + }] + }); + // "line2\n" -> "CHANGED\n" is two bytes longer. + assert_eq!( + a.has_byte_changes, + Some(sh_tangled::git::ByteChanges { + lhs: 18, + rhs: 20, + extra_data: Default::default(), + }) + ); + } + + #[test] + fn a_malformed_sha_is_a_client_error_and_an_absent_one_is_a_miss() { + let (root, base, head) = fork_fixture(); + + assert!(matches!( + diff_fixture(root.path(), &base, "not-a-sha"), + Err(XrpcError::InvalidRequest(_)) + )); + assert!(matches!( + diff_fixture(root.path(), &base, &"0".repeat(head.len())), + Err(XrpcError::RevisionNotFound { .. }) + )); + } +} + +#[cfg(test)] +mod merge_commit_tests { use super::*; use jacquard_common::types::did::Did; use std::io::Write as _; use std::path::{Path, PathBuf}; use std::process::{Command, Stdio}; - use crate::layout::Layout; + use gitmirror_git::layout::Layout; const TARGET_DID: &str = "did:plc:upstream"; const SOURCE_DID: &str = "did:plc:fork"; @@ -511,12 +1022,11 @@ mod tests { } } - fn scratch(fixture: &Fixture) -> crate::git::TempRepository { + fn scratch(fixture: &Fixture) -> gitmirror_git::scratch::TempRepository { fixture .layout .open_scratch(&[&Did::raw(TARGET_DID.into()), &Did::raw(SOURCE_DID.into())]) .unwrap() - .unwrap() } #[test] diff --git a/gitmirror/crates/gitmirror-xrpc/src/routes/mod.rs b/gitmirror/crates/gitmirror-xrpc/src/routes/mod.rs new file mode 100644 index 000000000..c2bf1c3ee --- /dev/null +++ b/gitmirror/crates/gitmirror-xrpc/src/routes/mod.rs @@ -0,0 +1 @@ +pub mod git; diff --git a/gitmirror/crates/gitmirror/Cargo.toml b/gitmirror/crates/gitmirror/Cargo.toml new file mode 100644 index 000000000..e92a4ad42 --- /dev/null +++ b/gitmirror/crates/gitmirror/Cargo.toml @@ -0,0 +1,34 @@ +[package] +name = "gitmirror" +version.workspace = true +edition.workspace = true +license.workspace = true +rust-version.workspace = true + +[[bin]] +name = "gitmirror" +path = "src/main.rs" + +[dependencies] +bobbin-runtime = { workspace = true } +gitmirror-xrpc = { workspace = true } + +anyhow = { workspace = true } +axum = { workspace = true } +clap = { workspace = true } +confique = { workspace = true } +futures = { workspace = true } +jacquard-axum = { workspace = true } +jacquard-common = { workspace = true } +jacquard-identity = { workspace = true } +reqwest = { workspace = true } +rustls = { workspace = true } +socket2 = "0.6" +thiserror = { workspace = true } +tokio = { workspace = true } +tokio-util = { workspace = true } +tracing = { workspace = true } +tracing-subscriber = { workspace = true } + +[lints] +workspace = true diff --git a/gitmirror/crates/gitmirror/src/config.rs b/gitmirror/crates/gitmirror/src/config.rs new file mode 100644 index 000000000..a353500b0 --- /dev/null +++ b/gitmirror/crates/gitmirror/src/config.rs @@ -0,0 +1,111 @@ +use std::{net::SocketAddr, path::PathBuf, str::FromStr as _}; + +use anyhow::Context as _; +use confique::Config; +use jacquard_common::types::did::Did; + +const SYSTEM_CONFIG_PATH: &str = "/etc/gitmirror/config.toml"; + +#[derive(Debug, Config)] +pub struct MirrorConfig { + #[config(nested)] + pub server: ServerConfig, + #[config(nested)] + pub repo: RepoConfig, + + #[config(nested)] + pub service_auth: ServiceAuthConfig, + #[config(nested)] + pub identity: IdentityConfig, + + #[config(nested)] + pub knot: KnotConfig, + #[config(nested)] + pub log: LogConfig, +} + +#[derive(Debug, Config)] +pub struct ServerConfig { + /// Addresses the XRPC server listens on. When using as an env var, comma-separated. + #[config( + env = "GITMIRROR_BIND", + parse_env = parse_binds, + default = ["127.0.0.1:9001", "[::1]:9001"] + )] + pub binds: Vec, +} + +// TODO: share codes with bobbin + +#[derive(Debug, thiserror::Error)] +pub enum BindParseError { + #[error("GITMIRROR_BIND must list at least one address")] + Empty, + #[error("invalid bind entry `{0}`: {1}")] + Invalid(String, std::net::AddrParseError), +} + +fn parse_binds(raw: &str) -> Result, BindParseError> { + let addrs: Vec = raw + .split(',') + .map(str::trim) + .filter(|s| !s.is_empty()) + .map(|s| SocketAddr::from_str(s).map_err(|e| BindParseError::Invalid(s.to_owned(), e))) + .collect::>()?; + if addrs.is_empty() { + return Err(BindParseError::Empty); + } + Ok(addrs) +} + +#[derive(Debug, Config)] +pub struct RepoConfig { + #[config(env = "GITMIRROR_SCAN_PATH")] + pub scan_path: PathBuf, +} + +#[derive(Debug, Config)] +pub struct ServiceAuthConfig { + /// This service's own DID, e.g. `did:web:mirror.tangled.network`. + #[config(env = "GITMIRROR_SERVICE_DID")] + pub did: Did, +} + +#[derive(Debug, Config)] +pub struct IdentityConfig { + #[config(env = "GITMIRROR_PLC_URL", default = "https://plc.directory")] + pub plc_url: String, +} + +#[derive(Debug, Config)] +pub struct KnotConfig { + /// Allow knot endpoints on private or loopback addresses. + #[config(env = "GITMIRROR_KNOT_ALLOW_PRIVATE", default = false)] + pub allow_private: bool, + + /// Require https on knot endpoints. + #[config(env = "GITMIRROR_KNOT_REQUIRE_HTTPS", default = true)] + pub require_https: bool, +} + +#[derive(Debug, Config)] +pub struct LogConfig { + /// `tracing-subscriber` env-filter directive, e.g. `gitmirror_xrpc=debug,info`. + #[config(env = "GITMIRROR_LOG", default = "info")] + pub filter: String, +} + +pub fn load(path: Option<&PathBuf>) -> anyhow::Result { + let mut builder = MirrorConfig::builder().env(); + if let Some(p) = path { + builder = builder.file(p); + } + builder + .file(SYSTEM_CONFIG_PATH) + .load() + .context("load configuration") +} + +pub fn template() -> String { + confique::toml::template::(confique::toml::FormatOptions::default()) +} diff --git a/gitmirror/crates/gitmirror/src/main.rs b/gitmirror/crates/gitmirror/src/main.rs new file mode 100644 index 000000000..b0c4f17ad --- /dev/null +++ b/gitmirror/crates/gitmirror/src/main.rs @@ -0,0 +1,223 @@ +use std::{net::SocketAddr, path::PathBuf, process::ExitCode, sync::Arc, time::Duration}; + +use anyhow::{Context as _, anyhow}; +use bobbin_runtime::{ReqwestHttp, SystemClock}; +use clap::{Parser, Subcommand}; +use jacquard_axum::service_auth::ServiceAuthConfig; +use jacquard_common::deps::fluent_uri::Uri; +use jacquard_identity::{JacquardResolver, resolver::PlcSource}; +use tokio::signal::unix::{SignalKind, signal}; +use tokio_util::sync::CancellationToken; +use tracing::level_filters::LevelFilter; +use tracing_subscriber::EnvFilter; + +mod config; + +use config::MirrorConfig; +use gitmirror_xrpc::{AppState, KnotPolicy, router}; + +#[derive(Parser)] +#[command(name = "gitmirror", about = "Git mirror service")] +struct Cli { + /// Path to a TOML config file. Environment variables override file values + /// for any `GITMIRROR_*` setting; `/etc/gitmirror/config.toml` is consulted as a + /// final fallback so distro packaging can drop a default in place. + #[arg(short, long, value_name = "FILE", env = "GITMIRROR_CONFIG")] + config: Option, + + #[command(subcommand)] + command: Command, +} + +#[derive(Subcommand)] +enum Command { + /// Print a fully-commented TOML template to stdout. Use this to seed + /// `config.toml` for a fresh deploy. + ConfigTemplate, + /// Load and validate the configuration without starting the server. + Validate, + /// Run the XRPC server. + Serve, +} + +#[tokio::main] +async fn main() -> ExitCode { + let _ = rustls::crypto::aws_lc_rs::default_provider().install_default(); + + let cli = Cli::parse(); + + if let Command::ConfigTemplate = cli.command { + print!("{}", config::template()); + return ExitCode::SUCCESS; + } + + let cfg = match config::load(cli.config.as_ref()) { + Ok(c) => c, + Err(e) => { + eprintln!("failed to load configuration: {e:#}"); + return ExitCode::FAILURE; + } + }; + + if let Err(e) = init_tracing(&cfg) { + eprintln!("failed to install tracing subscriber: {e}"); + return ExitCode::FAILURE; + } + + match cli.command { + Command::Validate => { + println!("configuration is valid"); + ExitCode::SUCCESS + } + Command::Serve => match serve(cfg).await { + Ok(()) => ExitCode::SUCCESS, + Err(e) => { + tracing::error!(error = ?e, "fatal"); + ExitCode::FAILURE + } + }, + // Handled before the config is loaded. + Command::ConfigTemplate => ExitCode::SUCCESS, + } +} + +fn init_tracing(cfg: &MirrorConfig) -> Result<(), String> { + let combined = format!("{},{}", LevelFilter::INFO, cfg.log.filter); + let filter = EnvFilter::try_new(&combined) + .map_err(|e| format!("invalid log filter `{}`: {e}", cfg.log.filter))?; + tracing_subscriber::fmt() + .with_env_filter(filter) + .try_init() + .map_err(|e| e.to_string()) +} + +async fn serve(cfg: MirrorConfig) -> anyhow::Result<()> { + let http = reqwest::Client::builder() + .connect_timeout(Duration::from_secs(10)) + .timeout(Duration::from_mins(5)) + .build()?; + let directory = Arc::new( + JacquardResolver::new(ReqwestHttp::new(http.clone()), Default::default()) + .with_plc_source(plc_source(&cfg.identity.plc_url)?), + ); + let service_auth = + ServiceAuthConfig::new(cfg.service_auth.did.clone(), directory).disable_replay_protection(); + let state = AppState::new( + Arc::new(cfg.repo.scan_path), + http, + service_auth, + KnotPolicy { + allow_private: cfg.knot.allow_private, + require_https: cfg.knot.require_https, + }, + Arc::new(SystemClock::new()), + ); + let app = router(state); + + // TODO: add debug routes + + let binds = cfg.server.binds.clone(); + let bind_display = binds + .iter() + .map(|bind| bind.to_string()) + .collect::>() + .join(","); + tracing::info!(binds = %bind_display, "gitmirror listening"); + + serve_all(binds, app).await +} + +async fn serve_all(binds: Vec, app: axum::Router) -> anyhow::Result<()> { + let listeners = binds + .into_iter() + .map(|addr| { + let listener = bind_listener(addr).with_context(|| format!("bind {addr}"))?; + tracing::info!(%addr, "gitmirror listener bound"); + anyhow::Ok(listener) + }) + .collect::>>()?; + + let cancel = CancellationToken::new(); + let trigger = cancel.clone(); + let signal_task = tokio::spawn(async move { + wait_for_shutdown().await; + tracing::info!("shutdown signal received, draining server"); + trigger.cancel(); + }); + + let result = futures::future::try_join_all(listeners.into_iter().map(|listener| { + let app = app.clone(); + let cancel = cancel.clone(); + async move { + axum::serve(listener, app) + .with_graceful_shutdown(async move { cancel.cancelled().await }) + .await + } + })) + .await; + signal_task.abort(); + result.map(|_| ()).context("axum server failed") +} + +fn bind_listener(addr: SocketAddr) -> std::io::Result { + let domain = match addr { + SocketAddr::V4(_) => socket2::Domain::IPV4, + SocketAddr::V6(_) => socket2::Domain::IPV6, + }; + let socket = socket2::Socket::new(domain, socket2::Type::STREAM, Some(socket2::Protocol::TCP))?; + // Otherwise a dual-stack `[::]` bind swallows the port and the sibling `0.0.0.0` bind fails. + if matches!(addr, SocketAddr::V6(_)) { + socket.set_only_v6(true)?; + } + socket.set_reuse_address(true)?; + socket.set_nonblocking(true)?; + socket.bind(&addr.into())?; + socket.listen(1024)?; + tokio::net::TcpListener::from_std(socket.into()) +} + +async fn wait_for_shutdown() { + let ctrl_c = tokio::signal::ctrl_c(); + let mut sigterm = match signal(SignalKind::terminate()) { + Ok(s) => s, + Err(e) => { + tracing::warn!( + ?e, + "could not install SIGTERM handler, shutdown will only honor ctrl-c" + ); + ctrl_c.await.ok(); + return; + } + }; + tokio::select! { + _ = ctrl_c => {} + _ = sigterm.recv() => {} + } +} + +/// `jacquard-identity`'s `DidStep::PlcHttp` joins the base and the DID by concatenation, so a base +/// missing its trailing slash would ask for `https://plc.exampledid:plc:...`. +fn plc_source(url: &str) -> anyhow::Result { + let base = format!("{}/", url.trim_end_matches('/')); + Ok(PlcSource::PlcDirectory { + base: Uri::parse(base.as_str()) + .map_err(|e| anyhow!("{url:?} is not a usable PLC directory: {e}"))? + .to_owned(), + }) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn a_plc_base_carries_the_separator_the_resolver_leaves_out() { + let base = |url| match plc_source(url).unwrap() { + PlcSource::PlcDirectory { base } => base.to_string(), + other => panic!("{other:?}"), + }; + assert_eq!(base("https://plc.example"), "https://plc.example/"); + assert_eq!(base("https://plc.example/"), "https://plc.example/"); + assert!(plc_source("not a url").is_err()); + } +} diff --git a/gitmirror/src/auth.rs b/gitmirror/src/auth.rs deleted file mode 100644 index 9c3e6fb9d..000000000 --- a/gitmirror/src/auth.rs +++ /dev/null @@ -1,215 +0,0 @@ -use axum::extract::FromRequestParts; -use axum::http::header; -use axum::http::request::Parts; -use bobbin_runtime::ReqwestHttp; -use jacquard_axum::service_auth::{ - ReplayKey, ReplayStore, ServiceAuth, ServiceAuthConfig, ServiceAuthError, -}; -use jacquard_common::IntoStatic as _; -use jacquard_common::deps::fluent_uri::Uri; -use jacquard_common::deps::smol_str::SmolStr; -use jacquard_common::service_auth; -use jacquard_common::types::crypto::KeyCodec; -use jacquard_common::types::string::Did; -use jacquard_identity::JacquardResolver; -use jacquard_identity::resolver::{IdentityResolver as _, PlcSource}; - -use crate::xrpc::XrpcState; - -/// jacquard's own `HttpClient` is implemented for reqwest 0.12; the workspace is on 0.13, and -/// `bobbin_runtime::ReqwestHttp` is the adapter that already bridges the two. -pub(crate) type Directory = JacquardResolver; - -/// Service auth against `audience`, which is the knot the token will be forwarded to. `plc` is the -/// directory `did:plc:` documents are read from — the token's issuer, and the repo whose knot is -/// pushed to. -pub(crate) fn config( - http: reqwest::Client, - audience: Did, - plc: &str, -) -> anyhow::Result> { - Ok(ServiceAuthConfig::new( - audience, - JacquardResolver::new(ReqwestHttp::new(http), Default::default()) - .with_plc_source(plc_source(plc)?), - ) - .disable_replay_protection()) -} - -/// `jacquard-identity`'s `DidStep::PlcHttp` joins the base and the DID by concatenation, so a base -/// missing its trailing slash would ask for `https://plc.exampledid:plc:...`. -fn plc_source(url: &str) -> anyhow::Result { - let base = format!("{}/", url.trim_end_matches('/')); - Ok(PlcSource::PlcDirectory { - base: Uri::parse(base.as_str()) - .map_err(|e| anyhow::anyhow!("{url:?} is not a usable PLC directory: {e}"))? - .to_owned(), - }) -} - -impl ServiceAuth for XrpcState { - type Resolver = Directory; - - fn service_did(&self) -> Did<&str> { - self.service_auth.service_did() - } - - fn resolver(&self) -> &Self::Resolver { - self.service_auth.resolver() - } - - fn require_lxm(&self) -> bool { - // Qualified: `ServiceAuthConfig`'s inherent `require_lxm` is the builder setter. - ServiceAuth::require_lxm(&*self.service_auth) - } - - fn allowed_services(&self) -> &[SmolStr] { - self.service_auth.allowed_services() - } - - fn replay_protection_enabled(&self) -> bool { - self.service_auth.replay_protection_enabled() - } - - fn replay_store(&self) -> &dyn ReplayStore { - self.service_auth.replay_store() - } -} - -fn bearer_token_from_parts(parts: &Parts) -> Result, ServiceAuthError> { - let Some(auth_header) = parts.headers.get(header::AUTHORIZATION) else { - return Ok(None); - }; - - let auth_str = auth_header - .to_str() - .map_err(|_| ServiceAuthError::InvalidAuthHeader)?; - let token = auth_str - .strip_prefix("Bearer ") - .ok_or(ServiceAuthError::InvalidAuthHeader)?; - Ok(Some(token)) -} - -/// [`jacquard_axum::service_auth::ExtractServiceAuth`] without the audience check, plus the raw -/// token -pub struct ExtractAnyServiceAuth(pub service_auth::ServiceAuthClaims, pub String); - -impl FromRequestParts for ExtractAnyServiceAuth -where - S: ServiceAuth + Send + Sync, - S::Resolver: Send + Sync, -{ - type Rejection = ServiceAuthError; - - async fn from_request_parts(parts: &mut Parts, state: &S) -> Result { - let token = bearer_token_from_parts(&parts)?.ok_or(ServiceAuthError::MissingAuthHeader)?; - let claims = verify_service_auth(state, token).await?; - - Ok(Self(claims, token.to_owned())) - } -} - -async fn verify_service_auth( - state: &S, - token: &str, -) -> Result -where - S: ServiceAuth + Send + Sync, - S::Resolver: Send + Sync, -{ - let parsed = service_auth::parse_jwt(token)?; - let claims = parsed.claims(); - - let doc = state - .resolver() - .resolve_did_doc_owned(&claims.iss) - .await - .map_err(|e| ServiceAuthError::DidResolutionFailed { - did: claims.iss.clone().into_static(), - source: Box::new(e), - })?; - - let signing_key = doc - .atproto_public_key() - .map_err(|e| ServiceAuthError::InvalidKey(e.to_string()))? - .ok_or_else(|| ServiceAuthError::NoSigningKey(claims.iss.clone().into_static())) - .map(|key| match key.codec { - KeyCodec::P256 => Ok(service_auth::PublicKey::from_p256_bytes(&key.bytes)), - KeyCodec::Secp256k1 => Ok(service_auth::PublicKey::from_k256_bytes(&key.bytes)), - codec => { - Err(ServiceAuthError::InvalidKey(format!( - "{codec:?} does not sign service auth tokens" - ))) - } - })???; - - service_auth::verify_signature(&parsed, &signing_key)?; - let now = chrono::Utc::now().timestamp(); - if claims.exp <= now { - return Err(service_auth::ServiceAuthError::Expired { exp: claims.exp, now }.into()); - } - - if state.require_lxm() && claims.lxm.is_none() { - return Err(ServiceAuthError::MethodBindingRequired); - } - - if state.replay_protection_enabled() { - let jti = claims.jti.as_ref().ok_or(ServiceAuthError::MissingJti)?; - let key = ReplayKey::new( - claims.iss.clone().into_static(), - claims.aud.clone().into_static(), - jti.clone(), - ); - state - .replay_store() - .check_and_insert(key, claims.exp) - .await?; - } - - Ok(parsed.into_claims()) -} - -// // TODO(boltless): make more flexible AtprotoService to replace -// // `jacquard_axum::service_auth::ServiceAuth`. -// // It will only verify the parsed signature, `claim.exp` and replayed state. -// pub trait AtprotoService { -// /// The identity resolver type -// type Resolver: IdentityResolver; -// -// /// Get the service DID (expected audience) -// fn service_did(&self) -> Did<&str>; -// -// /// Get a reference to the identity resolver -// fn resolver(&self) -> &Self::Resolver; -// -// /// Replay store used by replay protection. -// fn replay_store(&self) -> Option<&dyn ReplayStore>; -// } -// -// #[derive(Debug, Clone)] -// pub struct VerifiedServiceAuth<'a> { -// did: Did, -// aud: DidService, -// lxm: Option, -// jti: Option>, -// raw: CowStr<'a>, -// } -// -// pub struct ExtractServiceAuth(pub VerifiedServiceAuth<'static>); -// pub struct ExtractOptionalServiceAuth(pub Option>); - -#[cfg(test)] -mod tests { - use super::*; - - #[test] - fn a_plc_base_carries_the_separator_the_resolver_leaves_out() { - let base = |url| match plc_source(url).unwrap() { - PlcSource::PlcDirectory { base } => base.to_string(), - other => panic!("{other:?}"), - }; - assert_eq!(base("https://plc.example"), "https://plc.example/"); - assert_eq!(base("https://plc.example/"), "https://plc.example/"); - assert!(plc_source("not a url").is_err()); - } -} diff --git a/gitmirror/src/git.rs b/gitmirror/src/git.rs deleted file mode 100644 index a66529947..000000000 --- a/gitmirror/src/git.rs +++ /dev/null @@ -1,140 +0,0 @@ -use gix::bstr::ByteSlice as _; -use gix::revision::plumbing::Spec as RevSpec; -use gix::revision::walk::Sorting; -use jacquard_common::types::did::validate_did; -use tempfile::TempDir; - -use crate::xrpc::XrpcError; - -/// A [`gix::Repository`] backed by a temporary directory that is removed on drop. Derefs to -/// `gix::Repository`, so it is used just like one; keeping it alive keeps the scratch dir alive. -pub(crate) struct TempRepository { - repo: gix::Repository, - // Declared AFTER `repo` so `repo` drops first: any file handles into the scratch dir close - // before the dir itself is removed (Rust drops struct fields in declaration order). - _dir: TempDir, -} - -impl TempRepository { - pub fn open(dir: TempDir) -> Result { - let repo = gix::open(dir.path())?; - Ok(Self { repo, _dir: dir }) - } - - #[allow(dead_code)] - pub fn with_object_memory(mut self) -> Self { - self.repo.objects.enable_object_memory(); - self - } -} - -impl std::ops::Deref for TempRepository { - type Target = gix::Repository; - fn deref(&self) -> &gix::Repository { - &self.repo - } -} - -/// Build a throwaway bare repo whose `objects/info/alternates` points read-only at each of the -/// given git repositories, so a single `gix::Repository` can see objects from all of them -/// without ever mutating them. -pub(crate) fn open_scratch( - repo_base: &std::path::Path, - dids: &[&str], -) -> Result { - let mut seen: Vec<&str> = Vec::new(); - let mut object_dirs = Vec::new(); - for &did in dids { - if seen.contains(&did) { - continue; - } - // A DID has no `/`, so this is also what keeps `repo_base.join(did)` inside `repo_base`. - validate_did(did) - .map_err(|_| XrpcError::InvalidRequest("repo must be a DID".to_owned()))?; - seen.push(did); - let object_dir = - std::fs::canonicalize(repo_base.join(did).join("objects")).map_err(|e| { - XrpcError::RepoNotFound { - detail: format!("repo not found: {did}: {e}"), - } - })?; - object_dirs.push(object_dir); - } - if object_dirs.is_empty() { - return Err(XrpcError::Internal( - "open_scratch requires at least one repo".to_owned(), - )); - } - - let internal = |e: std::io::Error| XrpcError::Internal(e.to_string()); - let scratch = tempfile::tempdir().map_err(internal)?; - gix::init_bare(scratch.path()).map_err(|e| XrpcError::Internal(e.to_string()))?; - - let info_dir = scratch.path().join("objects").join("info"); - std::fs::create_dir_all(&info_dir).map_err(internal)?; - let alternates = object_dirs - .iter() - .map(|p| p.display().to_string()) - .collect::>() - .join("\n"); - std::fs::write(info_dir.join("alternates"), format!("{alternates}\n")).map_err(internal)?; - - TempRepository::open(scratch).map_err(|e| XrpcError::Internal(e.to_string())) -} - -pub(crate) fn commit_log_walk<'repo>( - repo: &'repo gix::Repository, - ranges: &[Vec], - all_refs: bool, -) -> anyhow::Result> { - let mut tips = Vec::new(); - let mut hidden = Vec::new(); - - if all_refs { - for r in repo.references()?.all()? { - let mut r = r.map_err(|e| anyhow::anyhow!(e))?; - if let Ok(commit) = r.peel_to_commit() { - tips.push(commit.id); - } - } - } else { - for range in ranges { - let revspec = repo.rev_parse(range.as_bstr())?; - let spec = revspec.detach(); - match spec { - RevSpec::Include(id) => tips.push(id), - RevSpec::Range { from, to } => { - tips.push(to); - hidden.push(from); - } - _ => { - anyhow::bail!("The spec isn't currently supported: {spec:?}") - } - } - } - } - - Ok(repo - .rev_walk(tips) - .sorting(Sorting::ByCommitTime(Default::default())) - .with_hidden(hidden) - .all()?) -} - -#[cfg(test)] -mod tests { - use super::*; - - #[test] - fn temp_repository_removes_scratch_dir_on_drop() { - let dir = tempfile::tempdir().unwrap(); - gix::init_bare(dir.path()).unwrap(); - let temp = TempRepository::open(dir).unwrap(); - let path = temp.git_dir().to_path_buf(); - assert!(path.exists()); - // Derefs to gix::Repository. - assert!(temp.object_hash() == gix::hash::Kind::Sha1); - drop(temp); - assert!(!path.exists(), "scratch dir should be gone after drop"); - } -} diff --git a/gitmirror/src/layout.rs b/gitmirror/src/layout.rs deleted file mode 100644 index d0a097a1d..000000000 --- a/gitmirror/src/layout.rs +++ /dev/null @@ -1,58 +0,0 @@ -use std::path::PathBuf; - -use jacquard_common::types::did::Did; -use tracing::error; - -use crate::git::TempRepository; - -#[derive(Debug, Clone)] -pub(crate) struct Layout { - scan_path: PathBuf, -} - -impl Layout { - pub fn new(scan_path: PathBuf) -> Self { - Self { scan_path } - } - - pub fn repo_path(&self, repo: &Did) -> PathBuf { - self.scan_path.join(repo.to_lowercase().as_str()) - } - - pub fn open_scratch(&self, dids: &[&Did]) -> anyhow::Result> { - let mut seen: Vec<&Did> = Vec::new(); - let mut object_dirs = Vec::new(); - for &did in dids { - if seen.contains(&did) { - continue; - } - seen.push(did); - let object_dir = match std::fs::canonicalize(self.repo_path(did).join("objects")) { - Ok(path) => path, - Err(err) => { - error!(error = %err, "repo not found"); - return Ok(None); - }, - }; - object_dirs.push(object_dir); - } - if object_dirs.is_empty() { - return Err(anyhow::anyhow!( - "open_scratch requires at least one repo".to_owned(), - )); - } - let scratch = tempfile::tempdir()?; - gix::init_bare(scratch.path())?; - - let info_dir = scratch.path().join("objects").join("info"); - std::fs::create_dir_all(&info_dir)?; - let alternates = object_dirs - .iter() - .map(|p| p.display().to_string()) - .collect::>() - .join("\n"); - std::fs::write(info_dir.join("alternates"), format!("{alternates}\n"))?; - - Ok(Some(TempRepository::open(scratch)?)) - } -} diff --git a/gitmirror/src/main.rs b/gitmirror/src/main.rs deleted file mode 100644 index 1c2e5cd61..000000000 --- a/gitmirror/src/main.rs +++ /dev/null @@ -1,95 +0,0 @@ -use std::net::SocketAddr; -use std::path::PathBuf; - -use clap::{Parser, Subcommand}; - -mod auth; -mod diff; -mod git; -mod git_repo_ext; -mod git_transport; -mod layout; -mod merge; -mod repo_did; -mod xrpc; -mod xrpc_git; - -#[derive(Parser)] -#[command(name = "gitmirror", about = "Git mirror service")] -struct Cli { - #[command(subcommand)] - cmd: Command, -} - -#[derive(Subcommand)] -enum Command { - /// Run the XRPC server. - Serve(ServeArgs), -} - -#[derive(clap::Args)] -struct ServeArgs { - /// Address to bind the HTTP XRPC server to. - #[arg(long, env = "GITMIRROR_ADDR", default_value = "127.0.0.1:9001")] - xrpc_addr: SocketAddr, - - /// Base directory holding bare mirror repos, one per DID (/). - #[arg(long, env = "GITMIRROR_REPO_BASE", default_value = "repos")] - repo_base: PathBuf, - - /// This service's own DID, e.g. `did:web:mirror.tangled.network`. - #[arg(long, env = "GITMIRROR_SERVICE_DID")] - service_did: String, - - /// Allow knot endpoints on private or loopback addresses. Off in production; on for local - /// testing against a knot on localhost. - #[arg( - long, - env = "GITMIRROR_KNOT_ALLOW_PRIVATE", - default_value_t = false, - default_missing_value = "true" - )] - knot_allow_private: bool, - - /// PLC directory to read `did:plc:` documents from. - #[arg(long, env = "GITMIRROR_PLC_URL", default_value = "https://plc.directory")] - plc_url: String, - - /// Require https on knot endpoints. Disable only when pushing to a local knot. - #[arg( - long, - env = "GITMIRROR_KNOT_REQUIRE_HTTPS", - default_value_t = true, - default_missing_value = "true" - )] - knot_require_https: bool, -} - -#[tokio::main] -async fn main() -> anyhow::Result<()> { - tracing_subscriber::fmt() - .with_env_filter( - tracing_subscriber::EnvFilter::try_from_default_env().unwrap_or_else(|_| "info".into()), - ) - .init(); - - match Cli::parse().cmd { - Command::Serve(args) => { - // Parsed here so a malformed DID is a startup failure rather than a surprise on the - // first request that reaches for it. - let service_did = jacquard_common::types::did::Did::new_owned(&args.service_did) - .map_err(|e| anyhow::anyhow!("--service-did {:?}: {e}", args.service_did))?; - xrpc::serve( - args.xrpc_addr, - args.repo_base, - service_did, - xrpc_git::KnotPolicy { - allow_private: args.knot_allow_private, - require_https: args.knot_require_https, - }, - args.plc_url, - ) - .await - } - } -} diff --git a/gitmirror/src/xrpc.rs b/gitmirror/src/xrpc.rs deleted file mode 100644 index 910a79ca6..000000000 --- a/gitmirror/src/xrpc.rs +++ /dev/null @@ -1,623 +0,0 @@ -use std::net::SocketAddr; -use std::path::PathBuf; -use std::sync::Arc; -use std::time::Duration; - -use axum::extract::State; -use axum::http::StatusCode; -use axum::response::{IntoResponse, Response}; -use axum::routing::{get, post}; -use axum::{Json, Router}; -use bobbin_types::sh_tangled; -use bobbin_types::sh_tangled::git::temp2::{get_diff, get_interdiff, list_commits, merge_check}; -use gix::ObjectId; -use jacquard_axum::{ExtractXrpc, XrpcResponse}; -use jacquard_common::ToSmolStr; -use jacquard_common::types::did::Did; -use jacquard_common::types::string::Datetime; -use serde_json::json; -use tracing::{error, info}; - -use crate::diff; -use crate::git::{commit_log_walk, open_scratch}; -use crate::layout::Layout; - -const DEFAULT_LIMIT: u32 = 50; -const MAX_LIMIT: u32 = 100; - -#[derive(Clone)] -pub(crate) struct XrpcState { - // TODO: deprecate in favor of `.layout` - repo_base: Arc, - pub(crate) layout: Arc, - pub(crate) http: reqwest::Client, - /// Service auth for the endpoints that take a credential; see [`crate::auth`]. - pub(crate) service_auth: Arc>, - pub(crate) knot_policy: crate::xrpc_git::KnotPolicy, -} - -#[derive(Debug)] -pub(crate) enum XrpcError { - InvalidRequest(String), - RepoNotFound { detail: String }, - RefNotFound { rev: String, detail: String }, - RevisionNotFound { rev: String }, - CompareError(String), - MergeConflict(Vec), - PushRejected(String), - Unauthorized(String), - Internal(String), -} - -impl IntoResponse for XrpcError { - fn into_response(self) -> Response { - let (status, error, message) = match self { - Self::InvalidRequest(m) => (StatusCode::BAD_REQUEST, "InvalidRequest", m), - Self::RepoNotFound { detail } => { - error!(error = %detail, "repo not found"); - ( - StatusCode::NOT_FOUND, - "RepoNotFound", - "repository not found".to_owned(), - ) - } - Self::RefNotFound { rev, detail } => { - error!(error = %detail, rev = %rev, "revision not found"); - ( - StatusCode::NOT_FOUND, - "RefNotFound", - format!("revision not found: {rev}"), - ) - } - Self::RevisionNotFound { rev } => ( - StatusCode::NOT_FOUND, - "RevisionNotFound", - format!("commit not found: {rev}"), - ), - Self::CompareError(m) => { - error!(error = %m, "compare failed"); - ( - StatusCode::INTERNAL_SERVER_ERROR, - "CompareError", - "failed to compare revisions".to_owned(), - ) - } - Self::MergeConflict(paths) => ( - StatusCode::CONFLICT, - "MergeConflict", - format!("merge produced conflicts in {}", paths.join(", ")), - ), - Self::PushRejected(detail) => { - error!(error = %detail, "push rejected"); - ( - StatusCode::CONFLICT, - "PushRejected", - "knot rejected the push".to_owned(), - ) - } - Self::Unauthorized(detail) => { - error!(error = %detail, "authentication failed"); - (StatusCode::UNAUTHORIZED, "AuthenticationRequired", detail) - } - Self::Internal(m) => { - error!(error = %m, "xrpc request failed"); - ( - StatusCode::INTERNAL_SERVER_ERROR, - "InternalServerError", - "internal error".to_owned(), - ) - } - }; - (status, Json(json!({ "error": error, "message": message }))).into_response() - } -} - -struct GixSignature<'a>(gix::actor::SignatureRef<'a>); - -impl TryFrom> for sh_tangled::git::Signature { - type Error = anyhow::Error; - - fn try_from(GixSignature(sig): GixSignature) -> Result { - let time = sig.time()?; - let offset = chrono::FixedOffset::east_opt(time.offset).ok_or_else(|| { - anyhow::anyhow!("commit timezone offset out of range: {}", time.offset) - })?; - let when = chrono::DateTime::from_timestamp(time.seconds, 0) - .ok_or_else(|| anyhow::anyhow!("commit timestamp out of range: {}", time.seconds))? - .with_timezone(&offset); - Ok(Self { - name: sig.name.to_smolstr(), - email: sig.email.to_smolstr(), - when: Datetime::new(when), - extra_data: Default::default(), - }) - } -} - -struct GixCommit<'a>(gix::Commit<'a>); - -impl TryFrom> for list_commits::Commit { - type Error = anyhow::Error; - - fn try_from(GixCommit(commit): GixCommit<'_>) -> Result { - let decoded = commit.decode()?; - Ok(Self { - oid: commit.id.to_smolstr(), - parents: decoded - .parents - .iter() - .map(|parent| parent.to_smolstr()) - .collect(), - tree: decoded.tree.to_smolstr(), - author: GixSignature(decoded.author()?).try_into()?, - committer: GixSignature(decoded.committer()?).try_into()?, - extra_headers: decoded - .extra_headers - .iter() - .map(|(key, value)| list_commits::Header { - key: key.to_smolstr(), - value: value.to_smolstr(), - extra_data: Default::default(), - }) - .collect(), - message: decoded.message.to_smolstr(), - extra_data: Default::default(), - }) - } -} - -impl From for sh_tangled::git::DiffSrc { - fn from(f: crate::diff::FileContent) -> Self { - Self { - path: f.path.into(), - oid: f.oid.into(), - size: f.size as i64, - is_binary: f.is_binary, - is_submodule: f.is_submodule, - content: f - .content - .map(|b| String::from_utf8_lossy(&b).into_owned().into()), - extra_data: Default::default(), - } - } -} - -impl From for sh_tangled::git::DiffHunk { - fn from(h: crate::diff::Hunk) -> Self { - // The novel sets are hash sets; sort them so the output is stable across runs. - let sorted = |set: rustc_hash::FxHashSet| -> Vec { - let mut v: Vec = set.into_iter().map(|n| i64::from(n.0)).collect(); - v.sort_unstable(); - v - }; - Self { - novel_lhs: sorted(h.novel_lhs), - novel_rhs: sorted(h.novel_rhs), - lines: h - .lines - .into_iter() - .map(|(lhs, rhs)| sh_tangled::git::LinePair { - lhs: lhs.map(|n| i64::from(n.0)), - rhs: rhs.map(|n| i64::from(n.0)), - extra_data: Default::default(), - }) - .collect(), - extra_data: Default::default(), - } - } -} - -impl From for sh_tangled::git::FileDiff { - fn from(d: crate::diff::Diff) -> Self { - Self { - lhs_src: d.lhs_src.into(), - rhs_src: d.rhs_src.into(), - hunks: d.hunks.into_iter().map(Into::into).collect(), - has_byte_changes: d - .has_byte_changes - .map(|(lhs, rhs)| sh_tangled::git::ByteChanges { - lhs: lhs as i64, - rhs: rhs as i64, - extra_data: Default::default(), - }), - has_syntactic_changes: d.has_syntactic_changes, - extra_data: Default::default(), - } - } -} - -/// Resolve a full hex oid to a commit that actually exists in `repo`. -pub(crate) fn find_commit(repo: &gix::Repository, sha: &str) -> Result { - let oid = gix::ObjectId::from_hex(sha.as_bytes()) - .map_err(|e| XrpcError::InvalidRequest(format!("bad commit sha {sha:?}: {e}")))?; - repo.find_commit(oid) - .map_err(|_| XrpcError::RevisionNotFound { - rev: sha.to_owned(), - })?; - Ok(oid) -} - -fn get_diff_inner( - repo: &gix::Repository, - base: gix::ObjectId, - head: gix::ObjectId, -) -> Result, XrpcError> { - let compare = || -> anyhow::Result> { - let merge_base = repo.merge_base(base, head)?.detach(); - let old = repo.find_tree(repo.find_commit(merge_base)?.tree_id()?)?; - let new = repo.find_tree(repo.find_commit(head)?.tree_id()?)?; - diff::diff(repo, &old, &new, false)? - .map(|d| d.map(Into::into)) - .collect() - }; - compare().map_err(|e| XrpcError::CompareError(e.to_string())) -} - -async fn get_diff( - State(state): State, - ExtractXrpc(args): ExtractXrpc, -) -> Result, XrpcError> { - let scratch = open_scratch( - &state.repo_base, - &[args.head_repo.as_str(), args.base_repo.as_str()], - )?; - let base = find_commit(&scratch, args.base_commit.as_ref())?; - let head = find_commit(&scratch, args.head_commit.as_ref())?; - - tokio::task::spawn_blocking(move || get_diff_inner(&scratch, base, head)) - .await - .map_err(|e| XrpcError::Internal(e.to_string()))? - .map(|diffs| { - XrpcResponse(get_diff::GetDiffOutput { - diffs, - extra_data: Default::default(), - }) - }) -} - -async fn merge_check( - State(state): State, - ExtractXrpc(params): ExtractXrpc, -) -> Result, XrpcError> { - let scratch = open_scratch( - &state.repo_base, - &[params.target_repo.as_str(), params.source_repo.as_str()], - )?; - let target = find_commit(&scratch, params.target_commit.as_ref())?; - let source = find_commit(&scratch, params.source_commit.as_ref())?; - - tokio::task::spawn_blocking(move || crate::merge::merge_check(&scratch, target, source)) - .await - .map_err(|e| XrpcError::Internal(e.to_string()))? - .map(XrpcResponse) - .map_err(|e| XrpcError::Internal(e.to_string())) -} - -fn get_interdiff_inner( - repo: &gix::Repository, - (from_base_id, from_head_id): (ObjectId, ObjectId), - (to_base_id, to_head_id): (ObjectId, ObjectId), -) -> Result, XrpcError> { - let to_head = repo - .find_commit(to_head_id) - .map_err(|e| XrpcError::Internal(e.to_string()))?; - let to_head_tree = to_head - .tree() - .map_err(|e| XrpcError::Internal(e.to_string()))?; - let rebased_tree = diff::prepare_interdiff(&repo, (from_base_id, from_head_id), to_base_id) - .map_err(|e| XrpcError::Internal(e.to_string()))?; - let compare = || -> anyhow::Result> { - diff::diff(&repo, &rebased_tree, &to_head_tree, true)? - .map(|d| d.map(Into::into)) - .collect() - }; - compare().map_err(|e| XrpcError::CompareError(e.to_string())) -} - -async fn get_interdiff( - State(state): State, - ExtractXrpc(args): ExtractXrpc, -) -> Result, XrpcError> { - let from_base_id = ObjectId::from_hex(args.base_commit1.as_bytes()) - .map_err(|e| XrpcError::InvalidRequest(e.to_string()))?; - let from_head_id = ObjectId::from_hex(args.base_commit2.as_bytes()) - .map_err(|e| XrpcError::InvalidRequest(e.to_string()))?; - let to_base_id = ObjectId::from_hex(args.head_commit1.as_bytes()) - .map_err(|e| XrpcError::InvalidRequest(e.to_string()))?; - let to_head_id = ObjectId::from_hex(args.head_commit2.as_bytes()) - .map_err(|e| XrpcError::InvalidRequest(e.to_string()))?; - - let scratch = open_scratch( - &state.repo_base, - &[args.head_repo.as_str(), args.base_repo.as_str()], - )?; - - tokio::task::spawn_blocking(move || { - get_interdiff_inner( - &scratch, - (from_base_id, from_head_id), - (to_base_id, to_head_id), - ) - }) - .await - .map_err(|e| XrpcError::Internal(e.to_string()))? - .map(|diffs| { - XrpcResponse(get_interdiff::GetInterdiffOutput { - diffs, - extra_data: Default::default(), - }) - }) -} - -fn list_commits_inner( - repo: &gix::Repository, - args: list_commits::ListCommits, -) -> Result, XrpcError> { - let ranges: Vec> = args - .ranges - .clone() - .unwrap_or_default() - .iter() - .map(|revspec| revspec.as_bytes().to_vec()) - .collect(); - let walk = commit_log_walk(repo, &ranges, args.all_refs.unwrap_or(false)).map_err(|e| { - XrpcError::RefNotFound { - rev: args.ranges.unwrap_or_default().join(", "), - detail: e.to_string(), - } - })?; - - let limit = args - .limit - .map(|limit| limit as u32) - .unwrap_or(DEFAULT_LIMIT); - if limit == 0 || limit > MAX_LIMIT { - return Err(XrpcError::InvalidRequest(format!( - "limit must be between 1 and {MAX_LIMIT}" - ))); - } - - let mut commits: Vec = Vec::with_capacity(limit as usize); - let mut skipped = 0usize; - - for info in walk { - let info = info.map_err(|e| XrpcError::Internal(e.to_string()))?; - - if (skipped as i64) < args.skip.unwrap_or(0) { - skipped += 1; - continue; - } - if (commits.len() as i64) == args.limit.unwrap_or(50) { - break; - } - - let commit = info - .object() - .map_err(|e| XrpcError::Internal(e.to_string()))?; - commits.push( - list_commits::Commit::try_from(GixCommit(commit)) - .map_err(|e| XrpcError::Internal(e.to_string()))?, - ); - } - - Ok(commits) -} - -async fn list_commits( - State(state): State, - ExtractXrpc(args): ExtractXrpc, -) -> Result, XrpcError> { - let path = state.repo_base.join(args.repo.as_str()); - let repo = gix::open(path).map_err(|e| XrpcError::RepoNotFound { - detail: e.to_string(), - })?; - - tokio::task::spawn_blocking(move || list_commits_inner(&repo, args)) - .await - .map_err(|e| XrpcError::Internal(e.to_string()))? - .map(|commits| { - XrpcResponse(list_commits::ListCommitsOutput { - commits, - extra_data: Default::default(), - }) - }) -} - -pub async fn serve( - addr: SocketAddr, - repo_base: PathBuf, - service_did: Did, - knot_policy: crate::xrpc_git::KnotPolicy, - plc: String, -) -> anyhow::Result<()> { - let http = reqwest::Client::builder() - .connect_timeout(Duration::from_secs(10)) - .timeout(Duration::from_secs(300)) - .build()?; - let app = Router::new() - .route("/xrpc/sh.tangled.git.temp2.getDiff", get(get_diff)) - .route( - "/xrpc/sh.tangled.git.temp2.getInterdiff", - get(get_interdiff), - ) - .route("/xrpc/sh.tangled.git.temp2.listCommits", get(list_commits)) - .route("/xrpc/sh.tangled.git.temp2.mergeCheck", get(merge_check)) - .route( - "/xrpc/sh.tangled.git.mergeCommit", - post(crate::xrpc_git::merge_commit), - ) - .with_state(XrpcState { - repo_base: Arc::new(repo_base.clone()), - layout: Arc::new(Layout::new(repo_base)), - http: http.clone(), - service_auth: Arc::new(crate::auth::config(http, service_did, &plc)?), - knot_policy, - }); - - let listener = tokio::net::TcpListener::bind(addr).await?; - info!(addr = %addr, "gitmirror XRPC server listening"); - axum::serve(listener, app).await?; - Ok(()) -} - -#[cfg(test)] -mod tests { - use super::*; - use std::path::Path; - use std::process::Command; - - const BASE_DID: &str = "did:plc:upstream"; - const HEAD_DID: &str = "did:plc:fork"; - - fn git(dir: &Path, args: &[&str]) { - let status = Command::new("git") - .args(args) - .current_dir(dir) - .env("GIT_AUTHOR_NAME", "t") - .env("GIT_AUTHOR_EMAIL", "t@t") - .env("GIT_COMMITTER_NAME", "t") - .env("GIT_COMMITTER_EMAIL", "t@t") - .status() - .expect("run git"); - assert!(status.success(), "git {args:?} failed"); - } - - fn rev_parse(dir: &Path, rev: &str) -> String { - let out = Command::new("git") - .args(["rev-parse", rev]) - .current_dir(dir) - .output() - .expect("rev-parse"); - String::from_utf8(out.stdout).unwrap().trim().to_owned() - } - - /// The real fork shape: the base tip lives only in the upstream mirror and the head tip only - /// in the fork's, with a shared ancestor. Returns `(repo_base, base_commit, head_commit)`. - fn fork_fixture() -> (tempfile::TempDir, String, String) { - let root = tempfile::tempdir().unwrap(); - let repo_base = root.path().join("repos"); - std::fs::create_dir_all(&repo_base).unwrap(); - - let upstream = root.path().join("upstream"); - std::fs::create_dir_all(&upstream).unwrap(); - git(&upstream, &["init", "-q", "-b", "main"]); - std::fs::write(upstream.join("a.txt"), "line1\nline2\nline3\n").unwrap(); - std::fs::write(upstream.join("b.txt"), "keep\n").unwrap(); - git(&upstream, &["add", "."]); - git(&upstream, &["commit", "-q", "-m", "shared ancestor"]); - - // Fork before upstream moves on, so neither tip is reachable from the other. - let fork = root.path().join("fork"); - git( - root.path(), - &[ - "clone", - "-q", - upstream.to_str().unwrap(), - fork.to_str().unwrap(), - ], - ); - - std::fs::write(upstream.join("upstream.txt"), "theirs\n").unwrap(); - git(&upstream, &["add", "."]); - git(&upstream, &["commit", "-q", "-m", "upstream only"]); - let base_commit = rev_parse(&upstream, "HEAD"); - - std::fs::write(fork.join("a.txt"), "line1\nCHANGED\nline3\n").unwrap(); - git(&fork, &["mv", "b.txt", "c.txt"]); - git(&fork, &["add", "-A"]); - git(&fork, &["commit", "-q", "-m", "fork only"]); - let head_commit = rev_parse(&fork, "HEAD"); - - for (did, src) in [(BASE_DID, &upstream), (HEAD_DID, &fork)] { - git( - root.path(), - &[ - "clone", - "-q", - "--bare", - src.to_str().unwrap(), - repo_base.join(did).to_str().unwrap(), - ], - ); - } - - (root, base_commit, head_commit) - } - - fn diff_fixture( - root: &Path, - base_commit: &str, - head_commit: &str, - ) -> Result, XrpcError> { - let repo_base = root.join("repos"); - let scratch = open_scratch(&repo_base, &[HEAD_DID, BASE_DID])?; - let base = find_commit(&scratch, base_commit)?; - let head = find_commit(&scratch, head_commit)?; - get_diff_inner(&scratch, base, head) - } - - #[test] - fn a_cross_repo_diff_reports_only_what_the_head_side_changed() { - let (root, base, head) = fork_fixture(); - let diffs = diff_fixture(root.path(), &base, &head).unwrap(); - - let mut paths: Vec = diffs - .iter() - .flat_map(|d| [d.lhs_src.path.to_string(), d.rhs_src.path.to_string()]) - .collect(); - paths.sort(); - paths.dedup(); - // `upstream.txt` is on the base side only: three-dot semantics exclude it. - assert_eq!(paths, ["a.txt", "b.txt", "c.txt"]); - - let a = diffs - .iter() - .find(|d| d.rhs_src.path == "a.txt") - .expect("a.txt is in the diff"); - // Only line 2 (0-based: 1) changed. - assert_eq!(a.hunks.len(), 1); - assert_eq!(a.hunks[0].novel_lhs, [1]); - assert_eq!(a.hunks[0].novel_rhs, [1]); - assert_eq!(a.hunks[0].lines, { - vec![sh_tangled::git::LinePair { - lhs: Some(1), - rhs: Some(1), - extra_data: Default::default(), - }] - }); - // "line2\n" -> "CHANGED\n" is two bytes longer. - assert_eq!( - a.has_byte_changes, - Some(sh_tangled::git::ByteChanges { - lhs: 18, - rhs: 20, - extra_data: Default::default(), - }) - ); - } - - #[test] - fn a_malformed_sha_is_a_client_error_and_an_absent_one_is_a_miss() { - let (root, base, head) = fork_fixture(); - - assert!(matches!( - diff_fixture(root.path(), &base, "not-a-sha"), - Err(XrpcError::InvalidRequest(_)) - )); - assert!(matches!( - diff_fixture(root.path(), &base, &"0".repeat(head.len())), - Err(XrpcError::RevisionNotFound { .. }) - )); - } - - #[test] - fn a_repo_that_is_not_a_did_never_reaches_the_filesystem() { - let (root, _, _) = fork_fixture(); - let repo_base = root.path().join("repos"); - assert!(matches!( - open_scratch(&repo_base, &["../../etc"]).map(|_| ()), - Err(XrpcError::InvalidRequest(_)) - )); - } -} diff --git a/nix/Cargo.nix b/nix/Cargo.nix index d3a5c93fd..97f10bbc6 100644 --- a/nix/Cargo.nix +++ b/nix/Cargo.nix @@ -193,6 +193,36 @@ rec { # File a bug if you depend on any for non-debug work! debug = internal.debugCrate { inherit packageId; }; }; + "gitmirror-git" = rec { + packageId = "gitmirror-git"; + build = internal.buildRustCrateWithFeatures { + packageId = "gitmirror-git"; + }; + + # Debug support which might change between releases. + # File a bug if you depend on any for non-debug work! + debug = internal.debugCrate { inherit packageId; }; + }; + "gitmirror-ingest" = rec { + packageId = "gitmirror-ingest"; + build = internal.buildRustCrateWithFeatures { + packageId = "gitmirror-ingest"; + }; + + # Debug support which might change between releases. + # File a bug if you depend on any for non-debug work! + debug = internal.debugCrate { inherit packageId; }; + }; + "gitmirror-xrpc" = rec { + packageId = "gitmirror-xrpc"; + build = internal.buildRustCrateWithFeatures { + packageId = "gitmirror-xrpc"; + }; + + # Debug support which might change between releases. + # File a bug if you depend on any for non-debug work! + debug = internal.debugCrate { inherit packageId; }; + }; "gix-rebase" = rec { packageId = "gix-rebase"; build = internal.buildRustCrateWithFeatures { @@ -573,6 +603,16 @@ rec { # File a bug if you depend on any for non-debug work! debug = internal.debugCrate { inherit packageId; }; }; + "tangled-axum" = rec { + packageId = "tangled-axum"; + build = internal.buildRustCrateWithFeatures { + packageId = "tangled-axum"; + }; + + # Debug support which might change between releases. + # File a bug if you depend on any for non-debug work! + debug = internal.debugCrate { inherit packageId; }; + }; "trusted-proxies" = rec { packageId = "trusted-proxies"; build = internal.buildRustCrateWithFeatures { @@ -9533,7 +9573,7 @@ rec { }; "gitmirror" = rec { crateName = "gitmirror"; - version = "0.1.0"; + version = "0.0.1"; edition = "2024"; crateBin = [ { @@ -9542,7 +9582,143 @@ rec { requiredFeatures = [ ]; } ]; - src = lib.cleanSourceWith { filter = sourceFilter; src = ../gitmirror; }; + src = lib.cleanSourceWith { filter = sourceFilter; src = ../gitmirror/crates/gitmirror; }; + dependencies = [ + { + name = "anyhow"; + packageId = "anyhow"; + } + { + name = "axum"; + packageId = "axum"; + features = [ "macros" ]; + } + { + name = "bobbin-runtime"; + packageId = "bobbin-runtime"; + } + { + name = "clap"; + packageId = "clap"; + features = [ "derive" "env" ]; + } + { + name = "confique"; + packageId = "confique"; + usesDefaultFeatures = false; + features = [ "toml" ]; + } + { + name = "futures"; + packageId = "futures"; + } + { + name = "gitmirror-xrpc"; + packageId = "gitmirror-xrpc"; + } + { + name = "jacquard-axum"; + packageId = "jacquard-axum"; + } + { + name = "jacquard-common"; + packageId = "jacquard-common"; + } + { + name = "jacquard-identity"; + packageId = "jacquard-identity"; + features = [ "cache" ]; + } + { + name = "reqwest"; + packageId = "reqwest 0.13.1"; + usesDefaultFeatures = false; + features = [ "rustls" "webpki-roots" "http2" "json" "gzip" "stream" ]; + } + { + name = "rustls"; + packageId = "rustls"; + features = [ "aws_lc_rs" "prefer-post-quantum" ]; + } + { + name = "socket2"; + packageId = "socket2"; + } + { + name = "thiserror"; + packageId = "thiserror 2.0.18"; + } + { + name = "tokio"; + packageId = "tokio"; + features = [ "macros" "rt-multi-thread" "time" "signal" "io-util" "net" "sync" ]; + } + { + name = "tokio-util"; + packageId = "tokio-util"; + features = [ "rt" ]; + } + { + name = "tracing"; + packageId = "tracing"; + } + { + name = "tracing-subscriber"; + packageId = "tracing-subscriber"; + features = [ "env-filter" "fmt" "json" ]; + } + ]; + + }; + "gitmirror-git" = rec { + crateName = "gitmirror-git"; + version = "0.0.1"; + edition = "2024"; + src = lib.cleanSourceWith { filter = sourceFilter; src = ../gitmirror/crates/gitmirror-git; }; + libName = "gitmirror_git"; + dependencies = [ + { + name = "gix"; + packageId = "gix"; + features = [ "parallel" "revision" "blob-diff" "worktree-archive" "tree-editor" "merge" "sha1" "sha256" ]; + } + { + name = "gix-rebase"; + packageId = "gix-rebase"; + } + { + name = "itertools"; + packageId = "itertools 0.14.0"; + } + { + name = "jacquard-common"; + packageId = "jacquard-common"; + } + { + name = "tempfile"; + packageId = "tempfile"; + } + { + name = "thiserror"; + packageId = "thiserror 2.0.18"; + } + ]; + + }; + "gitmirror-ingest" = rec { + crateName = "gitmirror-ingest"; + version = "0.0.1"; + edition = "2024"; + src = lib.cleanSourceWith { filter = sourceFilter; src = ../gitmirror/crates/gitmirror-ingest; }; + libName = "gitmirror_ingest"; + + }; + "gitmirror-xrpc" = rec { + crateName = "gitmirror-xrpc"; + version = "0.0.1"; + edition = "2024"; + src = lib.cleanSourceWith { filter = sourceFilter; src = ../gitmirror/crates/gitmirror-xrpc; }; + libName = "gitmirror_xrpc"; dependencies = [ { name = "anyhow"; @@ -9561,24 +9737,19 @@ rec { name = "bobbin-runtime"; packageId = "bobbin-runtime"; } - { - name = "bobbin-types"; - packageId = "bobbin-types"; - } { name = "chrono"; packageId = "chrono"; features = [ "serde" ]; } - { - name = "clap"; - packageId = "clap"; - features = [ "derive" "env" ]; - } { name = "futures-lite"; packageId = "futures-lite"; } + { + name = "gitmirror-git"; + packageId = "gitmirror-git"; + } { name = "gix"; packageId = "gix"; @@ -9590,10 +9761,6 @@ rec { usesDefaultFeatures = false; features = [ "generate" "streaming-input" "sha1" "sha256" ]; } - { - name = "gix-rebase"; - packageId = "gix-rebase"; - } { name = "gix-receive-pack"; packageId = "gix-receive-pack"; @@ -9617,6 +9784,10 @@ rec { packageId = "jacquard-identity"; features = [ "cache" ]; } + { + name = "lexicons"; + packageId = "lexicons"; + } { name = "line-numbers"; packageId = "line-numbers"; @@ -9631,61 +9802,34 @@ rec { name = "rustc-hash"; packageId = "rustc-hash"; } - { - name = "serde"; - packageId = "serde"; - features = [ "derive" ]; - } { name = "serde_json"; packageId = "serde_json"; features = [ "raw_value" ]; } { - name = "tempfile"; - packageId = "tempfile"; - } - { - name = "thiserror"; - packageId = "thiserror 2.0.18"; + name = "tangled-axum"; + packageId = "tangled-axum"; } { name = "tokio"; packageId = "tokio"; features = [ "macros" "rt-multi-thread" "time" "signal" "io-util" "net" "sync" "net" "rt-multi-thread" "sync" ]; } - { - name = "tokio-stream"; - packageId = "tokio-stream"; - features = [ "sync" ]; - } { name = "tracing"; packageId = "tracing"; } - { - name = "tracing-subscriber"; - packageId = "tracing-subscriber"; - features = [ "env-filter" ]; - } ]; devDependencies = [ { - name = "base64"; - packageId = "base64"; - } - { - name = "bytes"; - packageId = "bytes"; - } - { - name = "k256"; - packageId = "k256"; - features = [ "ecdsa" ]; + name = "tempfile"; + packageId = "tempfile"; } { - name = "multibase"; - packageId = "multibase"; + name = "tokio"; + packageId = "tokio"; + features = [ "macros" "rt-multi-thread" "time" "signal" "io-util" "net" "sync" "macros" "io-util" ]; } ]; @@ -30074,6 +30218,45 @@ rec { "Oliver Giersch" ]; + }; + "tangled-axum" = rec { + crateName = "tangled-axum"; + version = "0.0.1"; + edition = "2024"; + src = lib.cleanSourceWith { filter = sourceFilter; src = ../crates/tangled-axum; }; + libName = "tangled_axum"; + dependencies = [ + { + name = "axum"; + packageId = "axum"; + features = [ "macros" ]; + } + { + name = "bobbin-runtime"; + packageId = "bobbin-runtime"; + } + { + name = "jacquard-axum"; + packageId = "jacquard-axum"; + } + { + name = "jacquard-common"; + packageId = "jacquard-common"; + } + { + name = "jacquard-identity"; + packageId = "jacquard-identity"; + features = [ "cache" ]; + } + ]; + devDependencies = [ + { + name = "tokio"; + packageId = "tokio"; + features = [ "macros" "rt-multi-thread" "time" "signal" "io-util" "net" "sync" "macros" "rt" ]; + } + ]; + }; "tantivy" = rec { crateName = "tantivy"; @@ -31307,11 +31490,6 @@ rec { packageId = "tokio"; features = [ "sync" ]; } - { - name = "tokio-util"; - packageId = "tokio-util"; - optional = true; - } ]; devDependencies = [ { @@ -31331,7 +31509,7 @@ rec { "time" = [ "tokio/time" ]; "tokio-util" = [ "dep:tokio-util" ]; }; - resolvedDefaultFeatures = [ "default" "sync" "time" "tokio-util" ]; + resolvedDefaultFeatures = [ "default" "time" ]; }; "tokio-tungstenite 0.24.0" = rec { crateName = "tokio-tungstenite"; -- 2.51.2