diff --git a/bobbin/crates/xrpc/src/lib.rs b/bobbin/crates/xrpc/src/lib.rs index b6f786c94..3af834add 100644 --- a/bobbin/crates/xrpc/src/lib.rs +++ b/bobbin/crates/xrpc/src/lib.rs @@ -8,14 +8,15 @@ use bobbin_resolver::{ use axum::{ Router, - body::Body, + body::{Body, Bytes}, extract::{FromRequestParts, Query, RawQuery, State, rejection::QueryRejection}, http::{ HeaderMap, HeaderName, StatusCode, header::{ - ACCEPT_RANGES, CACHE_CONTROL, CONTENT_DISPOSITION, CONTENT_ENCODING, CONTENT_LENGTH, - CONTENT_RANGE, CONTENT_SECURITY_POLICY, CONTENT_TYPE, ETAG, IF_MODIFIED_SINCE, - IF_NONE_MATCH, IF_RANGE, LAST_MODIFIED, RANGE, X_CONTENT_TYPE_OPTIONS, + ACCEPT_RANGES, AUTHORIZATION, CACHE_CONTROL, CONTENT_DISPOSITION, CONTENT_ENCODING, + CONTENT_LENGTH, CONTENT_RANGE, CONTENT_SECURITY_POLICY, CONTENT_TYPE, ETAG, + IF_MODIFIED_SINCE, IF_NONE_MATCH, IF_RANGE, LAST_MODIFIED, RANGE, + X_CONTENT_TYPE_OPTIONS, }, request::Parts, }, @@ -659,11 +660,21 @@ const KNOT_PROXIED_NSIDS: &[&str] = &[ ]; const MIRROR_V2_PROXIED_NSIDS: &[&str] = &[ - "sh.tangled.git.temp2.listCommits", "sh.tangled.git.temp2.getBlame", "sh.tangled.git.temp2.getDiff", "sh.tangled.git.temp2.getInterdiff", + "sh.tangled.git.temp2.listCommits", "sh.tangled.git.temp2.mergeCheck", + "org.tangled.temp.git.getBlob", + "org.tangled.temp.git.getBranch", + "org.tangled.temp.git.getEntry", + "org.tangled.temp.git.getMergeBase", + "org.tangled.temp.git.getTree", +]; + +const MIRROR_V2_WRITE_NSIDS: &[&str] = &[ + "sh.tangled.git.mergeCommit", + "org.tangled.temp.git.deleteBranch", ]; const PASSTHROUGH_HEADERS: &[&HeaderName] = &[ @@ -691,7 +702,24 @@ type ProxyParams = Vec<(String, String)>; fn knot_proxied_routes() -> Router { let with_repo = register_proxied(Router::new(), REPO_PROXIED_NSIDS, proxy_repo_handler); let with_knot = register_proxied(with_repo, KNOT_PROXIED_NSIDS, proxy_knot_handler); - register_proxied(with_knot, MIRROR_V2_PROXIED_NSIDS, proxy_mirror_v2_handler) + let with_mirror = register_proxied(with_knot, MIRROR_V2_PROXIED_NSIDS, proxy_mirror_v2_handler); + register_mirror_v2_writes(with_mirror) +} + +fn register_mirror_v2_writes(router: Router) -> Router { + MIRROR_V2_WRITE_NSIDS + .iter() + .fold(router, |router, &nsid_lit| { + let nsid = nsid_static(nsid_lit); + router.route( + &format!("/xrpc/{nsid_lit}"), + axum::routing::post( + move |State(state): State, headers: HeaderMap, body: Bytes| { + proxy_mirror_v2_write_handler(state, headers, body, nsid.clone()) + }, + ), + ) + }) } fn register_proxied( @@ -4207,12 +4235,6 @@ async fn proxy_knot_handler( dispatch_knot(&state, &nsid, &host, &forward, allowed).await } -/// Hand the request to the v2 mirror unchanged. Deliberately shorter than its siblings: the origin -/// is operator-configured rather than resolved per-repo, and the mirror answers the same nsid we -/// were called on, so there is nothing to extract, rewrite, or validate. -/// -/// Modelled on [`dispatch_knot`], not [`dispatch_mirror`]: there is no knot to fall back to, so a -/// 4xx streams straight through and the mirror's own lexicon errors reach the client intact. async fn proxy_mirror_v2_handler( state: AppState, headers: HeaderMap, @@ -4235,3 +4257,50 @@ async fn proxy_mirror_v2_handler( .map(upstream_to_axum) .map_err(map_proxy_error) } + +async fn proxy_mirror_v2_write_handler( + state: AppState, + headers: HeaderMap, + body: Bytes, + nsid: Nsid, +) -> Result { + let mirror_v2 = state + .mirror_v2 + .as_ref() + .ok_or_else(|| XrpcError::UpstreamUnavailable("mirror_v2.url is unset".into()))?; + let token = headers + .get(AUTHORIZATION) + .ok_or_else(|| XrpcError::AuthRequired("the mirror needs a service-auth token".into()))?; + + let url = format!("{}xrpc/{nsid}", mirror_v2.host().url().as_str(),); + let request = reqwest::Client::new() + .post(url) + .header(AUTHORIZATION, token) + .header( + CONTENT_TYPE, + headers + .get(CONTENT_TYPE) + .cloned() + .unwrap_or(axum::http::HeaderValue::from_static("application/json")), + ) + .body(body); + + let upstream = request + .send() + .await + .map_err(|e| XrpcError::UpstreamUnavailable(format!("mirror: {e}")))?; + + let status = upstream.status(); + let upstream_headers = upstream.headers().clone(); + let mut response = Response::builder() + .status(status) + .body(Body::from_stream(upstream.bytes_stream())) + .expect("response body construction must succeed"); + let response_headers = response.headers_mut(); + PASSTHROUGH_HEADERS.iter().for_each(|name| { + if let Some(value) = upstream_headers.get(*name) { + response_headers.insert((*name).clone(), value.clone()); + } + }); + Ok(response) +} diff --git a/bobbin/crates/xrpc/tests/mirror_proxy.rs b/bobbin/crates/xrpc/tests/mirror_proxy.rs index a9d74d485..6c706dc6f 100644 --- a/bobbin/crates/xrpc/tests/mirror_proxy.rs +++ b/bobbin/crates/xrpc/tests/mirror_proxy.rs @@ -117,7 +117,8 @@ impl Harness { Arc::new(RepoIdResolver::detached(RuntimeHasher::default())), Arc::new(bobbin_xrpc::default_directory()), ) - .with_mirror(mirror_proxy); + .with_mirror(mirror_proxy.clone()) + .with_mirror_v2(mirror_proxy); let mut record = json!({ "$type": "sh.tangled.repo", @@ -358,3 +359,59 @@ async fn the_knot_keyed_endpoints_never_read_the_mirror() { assert_eq!(resp.status(), StatusCode::OK); assert!(paths(&h.mirror).await.is_empty()); } + +const DELETE_BRANCH: &str = "org.tangled.temp.git.deleteBranch"; + +async fn post_delete_branch(h: &Harness, token: Option<&str>) -> (StatusCode, String) { + let builder = Request::builder() + .method("POST") + .uri(format!("/xrpc/{DELETE_BRANCH}")) + .header("content-type", "application/json") + .extension(ConnectInfo(SOCKET)); + let builder = match token { + Some(value) => builder.header("authorization", value), + None => builder, + }; + let body = json!({ "repo": REPO_DID, "branch": "scratch" }).to_string(); + let resp = router(h.state.clone()) + .oneshot(builder.body(Body::from(body)).unwrap()) + .await + .expect("router infallible"); + let status = resp.status(); + let body = to_bytes(resp.into_body(), 64 * 1024).await.unwrap(); + (status, String::from_utf8(body.to_vec()).unwrap()) +} + +#[tokio::test] +async fn a_mirror_procedure_forwards_the_body_and_the_bearer() { + let h = Harness::with_mirror().await; + Mock::given(method("POST")) + .and(path(format!("/xrpc/{DELETE_BRANCH}"))) + .respond_with(ok_from_mirror()) + .mount(&h.mirror) + .await; + + let (status, body) = post_delete_branch(&h, Some("Bearer push-token")).await; + assert_eq!(status, StatusCode::OK); + assert_eq!(body, FROM_MIRROR); + + let sent = h.mirror.received_requests().await.unwrap(); + let sent = sent.last().expect("the mirror was asked"); + assert_eq!(sent.method, wiremock::http::Method::POST); + assert_eq!( + sent.headers.get("authorization").unwrap(), + "Bearer push-token" + ); + assert_eq!( + serde_json::from_slice::(&sent.body).unwrap(), + json!({ "repo": REPO_DID, "branch": "scratch" }) + ); +} + +#[tokio::test] +async fn a_mirror_procedure_without_a_bearer_never_reaches_the_mirror() { + let h = Harness::with_mirror().await; + let (status, _) = post_delete_branch(&h, None).await; + assert_eq!(status, StatusCode::UNAUTHORIZED); + assert!(paths(&h.mirror).await.is_empty()); +} diff --git a/knotmirror/xrpc/xrpc.go b/knotmirror/xrpc/xrpc.go index 6868be449..4d8732058 100644 --- a/knotmirror/xrpc/xrpc.go +++ b/knotmirror/xrpc/xrpc.go @@ -119,6 +119,11 @@ func (x *Xrpc) Router() http.Handler { r.Get("/sh.tangled.git.temp2.getInterdiff", x.proxyV2) r.Get("/sh.tangled.git.temp2.listCommits", x.proxyV2) r.Get("/sh.tangled.git.temp2.mergeCheck", x.proxyV2) + r.Get("/org.tangled.temp.git.getBlob", x.proxyV2) + r.Get("/org.tangled.temp.git.getBranch", x.proxyV2) + r.Get("/org.tangled.temp.git.getEntry", x.proxyV2) + r.Get("/org.tangled.temp.git.getMergeBase", x.proxyV2) + r.Get("/org.tangled.temp.git.getTree", x.proxyV2) r.Post("/sh.tangled.git.mergeCommit", x.proxyV2) r.Post("/org.tangled.temp.git.mergeCommit", x.proxyV2) r.Post("/org.tangled.temp.git.deleteBranch", x.proxyV2)