This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / knot2 / crates / knot-xrpc / src / receive.rs
9.2 kB 291 lines
1use std::future::Future; 2use std::io::Read; 3use std::pin::Pin; 4use std::sync::Arc; 5 6use axum::Router; 7use axum::body::Body; 8use axum::extract::{DefaultBodyLimit, Path, State}; 9use axum::http::{HeaderMap, HeaderValue, StatusCode, header}; 10use axum::response::{IntoResponse, Response}; 11use axum::routing::post; 12use futures::TryStreamExt; 13use knot_messages::{ErrorKey, HttpMessages}; 14use knot_pack::{PackError, PackReceiver, ReceivedPack, SocketPeer}; 15use knot_receive::{Push, land}; 16use knot_runtime::{Clock, HttpTransport}; 17use knot_types::{AccountDid, ActorId, ObjectFormat, RepoDid}; 18use tokio_util::io::{StreamReader, SyncIoBridge}; 19 20use crate::{XrpcError, XrpcState, authenticate_and_authorize_push, run_blocking}; 21 22const READ_CHUNK: usize = 64 * 1024; 23const ADVERTISEMENT: &str = "application/x-git-receive-pack-advertisement"; 24const RESULT: &str = "application/x-git-receive-pack-result"; 25 26pub(crate) fn routes<H: HttpTransport, C: Clock>() -> Router<Arc<XrpcState<H, C>>> { 27 Router::new() 28 .route( 29 "/{did}/{name}/git-receive-pack", 30 post(receive_named::<H, C>), 31 ) 32 .route("/{did}/git-receive-pack", post(receive_did::<H, C>)) 33 .layer(DefaultBodyLimit::disable()) 34} 35 36pub fn advertiser<H: HttpTransport, C: Clock>( 37 state: Arc<XrpcState<H, C>>, 38) -> Arc<dyn knot_pack::ReceiveAdvertiser> { 39 Arc::new(ReceiveGate { state }) 40} 41 42struct ReceiveGate<H, C> { 43 state: Arc<XrpcState<H, C>>, 44} 45 46impl<H: HttpTransport, C: Clock> knot_pack::ReceiveAdvertiser for ReceiveGate<H, C> { 47 fn advertise( 48 &self, 49 repo: RepoDid, 50 peer: SocketPeer, 51 headers: HeaderMap, 52 ) -> Pin<Box<dyn Future<Output = Response> + Send + '_>> { 53 Box::pin(async move { serve_advertisement(&self.state, repo, peer, &headers).await }) 54 } 55} 56 57async fn serve_advertisement<H: HttpTransport, C: Clock>( 58 state: &Arc<XrpcState<H, C>>, 59 repo: RepoDid, 60 peer: SocketPeer, 61 headers: &HeaderMap, 62) -> Response { 63 if let Err(response) = authorized_pusher(state, peer, headers, &repo).await { 64 return response; 65 } 66 let layout = state.layout.clone(); 67 let target = repo.clone(); 68 let missing = state.catalog.http.repo_not_found.text(); 69 let advert = run_blocking(move || { 70 let repo = layout 71 .open(&target) 72 .map_err(|_| XrpcError::not_found(missing))?; 73 knot_pack::advertise_receive(&repo).map_err(map_pack) 74 }) 75 .await; 76 match advert { 77 Ok(body) => git_response(ADVERTISEMENT, body), 78 Err(error) => error.into_response(), 79 } 80} 81 82async fn receive_named<H: HttpTransport, C: Clock>( 83 State(state): State<Arc<XrpcState<H, C>>>, 84 Path(crate::RepoPathParams { did, name }): Path<crate::RepoPathParams>, 85 peer: SocketPeer, 86 headers: HeaderMap, 87 body: Body, 88) -> Response { 89 let repo = match crate::resolve_repo_named(&state, &did, &name).await { 90 Ok(repo) => repo, 91 Err(error) => return error.into_response(), 92 }; 93 serve_receive(&state, repo, peer, &headers, body).await 94} 95 96async fn receive_did<H: HttpTransport, C: Clock>( 97 State(state): State<Arc<XrpcState<H, C>>>, 98 Path(did): Path<crate::RepoDidSegment>, 99 peer: SocketPeer, 100 headers: HeaderMap, 101 body: Body, 102) -> Response { 103 let repo = match crate::resolve_repo_did(&state, &did) { 104 Ok(repo) => repo, 105 Err(error) => return error.into_response(), 106 }; 107 serve_receive(&state, repo, peer, &headers, body).await 108} 109 110async fn serve_receive<H: HttpTransport, C: Clock>( 111 state: &Arc<XrpcState<H, C>>, 112 repo_did: RepoDid, 113 peer: SocketPeer, 114 headers: &HeaderMap, 115 body: Body, 116) -> Response { 117 let pusher = match authorized_pusher(state, peer, headers, &repo_did).await { 118 Ok(pusher) => pusher, 119 Err(response) => return response, 120 }; 121 122 let knot_actor = match state.secrets.public_key(&state.knot_did) { 123 Ok(public) => ActorId::from_secp256k1(public.as_bytes()), 124 Err(error) => return XrpcError::from(error).into_response(), 125 }; 126 127 let format = { 128 let layout = state.layout.clone(); 129 let target = repo_did.clone(); 130 let missing = state.catalog.http.repo_not_found.text(); 131 match run_blocking(move || { 132 layout 133 .open(&target) 134 .map(|repo| repo.object_format()) 135 .map_err(|_| XrpcError::not_found(missing)) 136 }) 137 .await 138 { 139 Ok(format) => format, 140 Err(error) => return error.into_response(), 141 } 142 }; 143 144 let _receive_permit = state.slots.receive.acquire().await; 145 146 let received = { 147 let scratch = state.layout.scratch_dir().to_path_buf(); 148 let limits = state.pack_limits; 149 let limit = state.byte_limits.pack; 150 let catalog = Arc::clone(&state.catalog); 151 let reader = StreamReader::new(body.into_data_stream().map_err(std::io::Error::other)); 152 run_blocking(move || { 153 drain_pack( 154 SyncIoBridge::new(reader), 155 &scratch, 156 limit, 157 limits, 158 format, 159 &catalog.http, 160 ) 161 }) 162 .await 163 }; 164 let received = match received { 165 Ok(received) => received, 166 Err(error) => return error.into_response(), 167 }; 168 if received.is_empty() { 169 return git_response(RESULT, Vec::new()); 170 } 171 172 let landed = land(Push { 173 layout: &state.layout, 174 repo_did: &repo_did, 175 received, 176 limits: state.pack_limits, 177 knot_actor, 178 committer: pusher, 179 events: Arc::clone(&state.events), 180 index: &state.index, 181 atproto: &state.atproto, 182 resolve_slots: &state.slots.resolve, 183 appview: &state.appview, 184 maintenance: &state.maintenance, 185 hostname: &state.knot_hostname, 186 languages_push_budget: state.budgets.languages_push, 187 catalog: Arc::clone(&state.catalog), 188 ci_logs: state.ci_logs.clone(), 189 }) 190 .await; 191 match landed { 192 Ok(framed) => git_response(RESULT, framed), 193 Err(error) => { 194 tracing::warn!(repo = repo_did.as_str(), %error, "http receive-pack failed"); 195 error.into_response() 196 } 197 } 198} 199 200fn drain_pack<R: Read>( 201 mut reader: R, 202 scratch: &std::path::Path, 203 limit: knot_pack::MaxWireBytes, 204 limits: knot_pack::PackLimits, 205 format: ObjectFormat, 206 messages: &HttpMessages, 207) -> Result<ReceivedPack, XrpcError> { 208 let mut receiver = PackReceiver::new(scratch, limit, limits, format.kind()) 209 .map_err(|error| XrpcError::internal(format!("receive staging failed: {error}")))?; 210 let mut buffer = [0u8; READ_CHUNK]; 211 loop { 212 let read = reader 213 .read(&mut buffer) 214 .map_err(|error| XrpcError::invalid_request(format!("receive read error: {error}")))?; 215 if read == 0 { 216 break; 217 } 218 if receiver 219 .write(&buffer[..read]) 220 .map_err(|error| map_receive_read(error, messages))? 221 { 222 break; 223 } 224 } 225 receiver 226 .finish() 227 .map_err(|error| map_receive_read(error, messages)) 228} 229 230async fn authorized_pusher<H: HttpTransport, C: Clock>( 231 state: &Arc<XrpcState<H, C>>, 232 peer: SocketPeer, 233 headers: &HeaderMap, 234 repo: &RepoDid, 235) -> Result<AccountDid, Response> { 236 let denied = state.catalog.http.push_denied.text(); 237 authenticate_and_authorize_push(state, peer, headers, repo, &denied) 238 .await 239 .map_err(challenge) 240} 241 242fn challenge(error: XrpcError) -> Response { 243 if error.status() == StatusCode::UNAUTHORIZED { 244 ( 245 StatusCode::UNAUTHORIZED, 246 [(header::WWW_AUTHENTICATE, crate::BASIC_CHALLENGE)], 247 error.to_string(), 248 ) 249 .into_response() 250 } else { 251 error.into_response() 252 } 253} 254 255fn git_response(content_type: &'static str, body: Vec<u8>) -> Response { 256 ( 257 StatusCode::OK, 258 [ 259 (header::CONTENT_TYPE, HeaderValue::from_static(content_type)), 260 ( 261 header::CACHE_CONTROL, 262 HeaderValue::from_static("no-cache, max-age=0, must-revalidate"), 263 ), 264 ], 265 body, 266 ) 267 .into_response() 268} 269 270fn map_pack(error: PackError) -> XrpcError { 271 XrpcError::named(error.http_status(), "PackError", error.to_string()) 272} 273 274fn map_receive_read(error: knot_pack::ReceiveReadError, messages: &HttpMessages) -> XrpcError { 275 match error { 276 knot_pack::ReceiveReadError::TooLarge => { 277 XrpcError::request_too_large(messages.push_too_large.text()) 278 } 279 knot_pack::ReceiveReadError::Truncated => { 280 XrpcError::invalid_request(messages.receive_ended_early.text()) 281 } 282 knot_pack::ReceiveReadError::Pack(error) => XrpcError::invalid_request( 283 messages 284 .malformed_pack 285 .line(|ErrorKey::Error| error.to_string()), 286 ), 287 knot_pack::ReceiveReadError::Io(error) => { 288 XrpcError::internal(format!("receive io error: {error}")) 289 } 290 } 291}