This repository has no description
1mod blocklist;
2mod body;
3mod branches;
4mod cob;
5mod collaborators;
6mod error;
7mod events;
8mod forks;
9mod lfs;
10mod lists;
11mod locks;
12mod members;
13mod merge;
14mod patchtext;
15mod query;
16mod reads;
17mod receive;
18mod repos;
19mod reservations;
20mod service;
21mod sniff;
22mod wire;
23
24#[cfg(test)]
25mod tests;
26
27pub use error::XrpcError;
28pub use knot_pack::MaxWireBytes;
29pub use knot_postreceive::LanguagesPushBudget;
30pub use knot_resource::{
31 Burst, GlobalInflight, LimitConfig, PerPeerInflight, PreAuthLimiter, RateLimit, RefillMicros,
32};
33pub use lfs::LfsWeb;
34pub use locks::CobLocks;
35pub use merge::Committer;
36pub use receive::advertiser as receive_advertiser;
37pub use reservations::{GlobalQuota, PerActorQuota, ReservationTtl, Reservations};
38
39use std::collections::BTreeSet;
40use std::net::IpAddr;
41use std::path::PathBuf;
42use std::sync::Arc;
43use std::time::{Duration, Instant};
44
45use axum::Json;
46use axum::Router;
47use axum::body::Bytes;
48use axum::extract::{DefaultBodyLimit, FromRequestParts, MatchedPath, Request, State};
49use axum::middleware::{Next, from_fn_with_state};
50use axum::response::{IntoResponse, Response};
51use axum::routing::{get, post};
52use http::request::Parts;
53use http::{HeaderMap, HeaderValue, StatusCode, header::AUTHORIZATION};
54use serde::de::DeserializeOwned;
55use serde_json::json;
56
57use knot_atproto::{Atproto, AtprotoError, ServiceJwt};
58use knot_events::{EventLog, SubscriberGate};
59use knot_git::Layout;
60use knot_index::{Index, Resolved};
61use knot_maintenance::MaintenanceHandle;
62use knot_resource::Slots;
63use knot_runtime::{Clock, Entropy, HttpTransport};
64use knot_secrets::SealedStore;
65use knot_types::{
66 AccountDid, AdmissionPolicy, AppviewEndpoint, CiLogsAddr, KnotHostname, KnotId, KnotServiceUrl,
67 Nsid, OwnerDid, OwnerRef, RepoDid, RepoRkey, UnixSeconds,
68};
69
70use base64::Engine;
71use knot_pack::SocketPeer;
72use knot_resource::{AdmitGuard, Refusal};
73
74pub(crate) const PUSH_NSID: &str = "sh.tangled.repo.push";
75
76#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
77pub enum ReadBudget {
78 Within(Duration),
79 Unbounded,
80}
81
82impl ReadBudget {
83 pub fn deadline(self) -> Option<Instant> {
84 match self {
85 ReadBudget::Within(budget) => Some(Instant::now() + budget),
86 ReadBudget::Unbounded => None,
87 }
88 }
89}
90
91// `XrpcState` keeps a bunch of these side by side,
92// some usize & some u64.
93// Within each group every one of them typechecked in every other one's slot.
94knot_types::scalar_newtype! {
95 pub struct BodyLimit(usize);
96 pub struct PatchLimit(usize);
97 pub struct PatchDecompressedLimit(u64);
98 pub struct ResponseLimit(usize);
99 pub struct ArchiveLimit(u64);
100 pub struct ForkPackLimit(u64);
101 pub struct TreeReadBudget(ReadBudget);
102 pub struct BlobReadBudget(ReadBudget);
103 pub struct LanguagesReadBudget(ReadBudget);
104}
105
106#[derive(Debug, Clone, Copy, PartialEq, Eq)]
107pub struct ByteLimits {
108 pub body: BodyLimit,
109 pub patch: PatchLimit,
110 pub patch_decompressed: PatchDecompressedLimit,
111 pub response: ResponseLimit,
112 pub archive: ArchiveLimit,
113 pub fork_pack: ForkPackLimit,
114 pub pack: MaxWireBytes,
115}
116
117impl Default for ByteLimits {
118 fn default() -> Self {
119 Self {
120 body: BodyLimit::new(64 * 1024),
121 patch: PatchLimit::new(16 * 1024 * 1024),
122 patch_decompressed: PatchDecompressedLimit::new(128 * 1024 * 1024),
123 response: ResponseLimit::new(5 * 1024 * 1024),
124 archive: ArchiveLimit::new(1024 * 1024 * 1024),
125 fork_pack: ForkPackLimit::new(1024 * 1024 * 1024),
126 pack: MaxWireBytes::new(8 * 1024 * 1024 * 1024),
127 }
128 }
129}
130
131#[derive(Debug, Clone, Copy, PartialEq, Eq)]
132pub struct Budgets {
133 pub tree_last_commit: TreeReadBudget,
134 pub blob_last_commit: BlobReadBudget,
135 pub languages: LanguagesReadBudget,
136 pub languages_push: LanguagesPushBudget,
137}
138
139impl Default for Budgets {
140 fn default() -> Self {
141 Self {
142 tree_last_commit: TreeReadBudget::new(ReadBudget::Within(Duration::from_millis(300))),
143 blob_last_commit: BlobReadBudget::new(ReadBudget::Within(Duration::from_millis(2_000))),
144 languages: LanguagesReadBudget::new(ReadBudget::Within(Duration::from_millis(1_000))),
145 languages_push: LanguagesPushBudget::new(Duration::from_millis(2_000)),
146 }
147 }
148}
149
150pub struct XrpcState<H, C> {
151 pub layout: Layout,
152 pub index: Arc<Index>,
153 pub atproto: Arc<Atproto<H, C>>,
154 pub secrets: Arc<SealedStore>,
155 pub entropy: Arc<dyn Entropy>,
156 pub admins: BTreeSet<AccountDid>,
157 pub admission: AdmissionPolicy,
158 pub knot_did: KnotId,
159 pub knot_hostname: KnotHostname,
160 pub ci_logs: Option<CiLogsAddr>,
161 pub meta_path: PathBuf,
162 pub knot_service_url: KnotServiceUrl,
163 pub limiter: Arc<PreAuthLimiter>,
164 pub cob_locks: Arc<CobLocks>,
165 pub reservations: Arc<Reservations>,
166 pub trusted_proxy_header: Option<http::HeaderName>,
167 pub committer: Committer,
168 pub byte_limits: ByteLimits,
169 pub budgets: Budgets,
170 pub git_http: Arc<dyn HttpTransport>,
171 pub pack_limits: knot_pack::PackLimits,
172 pub service_owner: AccountDid,
173 pub events: Arc<EventLog<C>>,
174 pub subscriber_gate: Arc<SubscriberGate>,
175 pub maintenance: MaintenanceHandle,
176 pub appview: AppviewEndpoint,
177 pub slots: Slots,
178 pub lfs: Option<LfsWeb>,
179 pub catalog: Arc<knot_messages::Catalog>,
180}
181
182impl<H: HttpTransport, C: Clock> XrpcState<H, C> {
183 pub fn now(&self) -> UnixSeconds {
184 UnixSeconds::new((self.atproto.now().get() / 1_000_000) as i64)
185 }
186
187 pub(crate) fn knot_authority(&self) -> &str {
188 self.knot_service_url.authority()
189 }
190
191 pub(crate) async fn authenticate(
192 &self,
193 headers: &HeaderMap,
194 method: &Method,
195 ) -> Result<AccountDid, XrpcError> {
196 let token = bearer(headers)?;
197 self.atproto
198 .verify_service_jwt(&token, method.nsid())
199 .await
200 .map_err(map_verify_error)
201 }
202
203 pub(crate) async fn authenticate_push(
204 &self,
205 headers: &HeaderMap,
206 ) -> Result<AccountDid, XrpcError> {
207 let token = push_credential(headers)?;
208 let method = Nsid::new_owned(PUSH_NSID).expect("push nsid is always a valid nsid");
209 self.atproto
210 .verify_service_jwt_guarded(
211 &token,
212 &method,
213 knot_atproto::ReplayGuard::ReusableUntilExpiry,
214 )
215 .await
216 .map_err(map_verify_error)
217 }
218}
219
220fn map_verify_error(error: AtprotoError) -> XrpcError {
221 if error.is_transient() {
222 XrpcError::upstream_unavailable(error.to_string())
223 } else {
224 XrpcError::auth_required(error.to_string())
225 }
226}
227
228pub(crate) struct Method(Nsid);
229
230impl Method {
231 fn nsid(&self) -> &Nsid {
232 &self.0
233 }
234
235 #[cfg(test)]
236 pub(crate) fn from_nsid(nsid: &str) -> Self {
237 Self(Nsid::new_owned(nsid).expect("test route nsid parses"))
238 }
239}
240
241impl<S: Send + Sync> FromRequestParts<S> for Method {
242 type Rejection = XrpcError;
243
244 async fn from_request_parts(parts: &mut Parts, state: &S) -> Result<Self, Self::Rejection> {
245 let matched = MatchedPath::from_request_parts(parts, state)
246 .await
247 .map_err(|_| XrpcError::internal("xrpc handler reached without a matched route"))?;
248 let nsid = matched
249 .as_str()
250 .strip_prefix("/xrpc/")
251 .ok_or_else(|| XrpcError::internal("xrpc route paths are prefixed with /xrpc/"))?;
252 Nsid::new_owned(nsid)
253 .map(Self)
254 .map_err(|_| XrpcError::internal("route nsid is always a valid nsid"))
255 }
256}
257
258pub fn router<H: HttpTransport, C: Clock>(state: Arc<XrpcState<H, C>>) -> Router {
259 let merge_routes = Router::new()
260 .route(merge::MERGE_ROUTE, post(merge::merge::<H, C>))
261 .route(merge::MERGE_CHECK_ROUTE, post(merge::merge_check::<H, C>))
262 .layer(DefaultBodyLimit::max(state.byte_limits.patch.get()));
263 Router::new()
264 .merge(merge_routes)
265 .route(members::ADD_ROUTE, post(members::add_member::<H, C>))
266 .route(members::REMOVE_ROUTE, post(members::remove_member::<H, C>))
267 .route(blocklist::BAN_ROUTE, post(blocklist::ban::<H, C>))
268 .route(blocklist::UNBAN_ROUTE, post(blocklist::unban::<H, C>))
269 .route(
270 collaborators::ADD_ROUTE,
271 post(collaborators::add_collaborator::<H, C>),
272 )
273 .route(
274 collaborators::REMOVE_ROUTE,
275 post(collaborators::remove_collaborator::<H, C>),
276 )
277 .route(repos::CREATE_ROUTE, post(repos::create_repo::<H, C>))
278 .route(repos::DELETE_ROUTE, post(repos::delete_repo::<H, C>))
279 .route(repos::RENAME_ROUTE, post(repos::rename_repo::<H, C>))
280 .route(repos::RESERVE_ROUTE, post(repos::reserve_key::<H, C>))
281 .route(
282 branches::SET_DEFAULT_ROUTE,
283 post(branches::set_default_branch::<H, C>),
284 )
285 .route(
286 branches::DELETE_ROUTE,
287 post(branches::delete_branch::<H, C>),
288 )
289 .route(forks::STATUS_ROUTE, post(forks::fork_status::<H, C>))
290 .route(forks::SYNC_ROUTE, post(forks::fork_sync::<H, C>))
291 .route(forks::HIDDEN_REF_ROUTE, post(forks::hidden_ref::<H, C>))
292 .route(reads::TREE_ROUTE, get(reads::repo_tree::<H, C>))
293 .route(reads::LOG_ROUTE, get(reads::repo_log::<H, C>))
294 .route(reads::BRANCHES_ROUTE, get(reads::repo_branches::<H, C>))
295 .route(reads::BRANCH_ROUTE, get(reads::repo_branch::<H, C>))
296 .route(reads::TAGS_ROUTE, get(reads::repo_tags::<H, C>))
297 .route(reads::TAG_ROUTE, get(reads::repo_tag::<H, C>))
298 .route(reads::BLOB_ROUTE, get(reads::repo_blob::<H, C>))
299 .route(reads::DIFF_ROUTE, get(reads::repo_diff::<H, C>))
300 .route(reads::COMPARE_ROUTE, get(reads::repo_compare::<H, C>))
301 .route(reads::ARCHIVE_ROUTE, get(reads::repo_archive::<H, C>))
302 .route(reads::LANGUAGES_ROUTE, get(reads::repo_languages::<H, C>))
303 .route(
304 reads::GET_DEFAULT_BRANCH_ROUTE,
305 get(reads::repo_get_default_branch::<H, C>),
306 )
307 .route(
308 reads::DESCRIBE_REPO_ROUTE,
309 get(reads::repo_describe_repo::<H, C>),
310 )
311 .route(reads::LIST_REFS_ROUTE, get(reads::git_list_refs::<H, C>))
312 .route(reads::LIST_REPOS_ROUTE, get(reads::sync_list_repos::<H, C>))
313 .route(lists::LIST_MEMBERS_ROUTE, get(lists::list_members::<H, C>))
314 .route(
315 lists::LIST_COLLABORATORS_ROUTE,
316 get(lists::list_collaborators::<H, C>),
317 )
318 .route(service::VERSION_ROUTE, get(service::version))
319 .route(service::OWNER_ROUTE, get(service::owner::<H, C>))
320 .layer(DefaultBodyLimit::max(state.byte_limits.body.get()))
321 .layer(from_fn_with_state(
322 Arc::clone(&state),
323 enforce_pre_auth_limit::<H, C>,
324 ))
325 .merge(lfs::routes::<H, C>())
326 .merge(receive::routes::<H, C>())
327 .route(service::HEALTH_ROUTE, get(service::health::<H, C>))
328 .route(events::EVENTS_ROUTE, get(events::events::<H, C>))
329 .with_state(state)
330}
331
332async fn enforce_pre_auth_limit<H: HttpTransport, C: Clock>(
333 State(state): State<Arc<XrpcState<H, C>>>,
334 socket: SocketPeer,
335 request: Request,
336 next: Next,
337) -> Response {
338 let peer = effective_peer(&state, socket, request.headers());
339 match admit_pre_auth(&state, peer) {
340 Ok(guard) => {
341 let response = next.run(request).await;
342 drop(guard);
343 response
344 }
345 Err(error) => error.into_response(),
346 }
347}
348
349pub(crate) fn effective_peer<H: HttpTransport, C: Clock>(
350 state: &XrpcState<H, C>,
351 socket: SocketPeer,
352 headers: &HeaderMap,
353) -> Option<IpAddr> {
354 state
355 .trusted_proxy_header
356 .as_ref()
357 .and_then(|header| knot_types::forwarded_peer(headers, header))
358 .or(socket.ip())
359}
360
361pub(crate) fn admit_pre_auth<H: HttpTransport, C: Clock>(
362 state: &XrpcState<H, C>,
363 peer: Option<IpAddr>,
364) -> Result<AdmitGuard, XrpcError> {
365 state
366 .limiter
367 .admit(peer, state.atproto.now())
368 .map_err(|refusal| match refusal {
369 Refusal::RateLimited => {
370 XrpcError::rate_limited("too many pre-authentication requests, retry shortly")
371 }
372 Refusal::Saturated => {
373 XrpcError::overloaded("knot is shedding pre-authentication load, retry shortly")
374 }
375 })
376}
377
378pub(crate) const BASIC_CHALLENGE: HeaderValue = HeaderValue::from_static("Basic realm=\"knot\"");
379
380fn strip_bearer(value: &str) -> Option<&str> {
381 let (scheme, rest) = value.split_once(' ')?;
382 scheme.eq_ignore_ascii_case("Bearer").then_some(rest)
383}
384
385fn bearer(headers: &HeaderMap) -> Result<ServiceJwt, XrpcError> {
386 headers
387 .get(AUTHORIZATION)
388 .and_then(|value| value.to_str().ok())
389 .and_then(strip_bearer)
390 .map(str::trim)
391 .and_then(|token| ServiceJwt::new(token).ok())
392 .ok_or_else(|| XrpcError::auth_required("missing or malformed Bearer authorization header"))
393}
394
395fn strip_basic(value: &str) -> Option<String> {
396 let (scheme, rest) = value.split_once(' ')?;
397 if !scheme.eq_ignore_ascii_case("Basic") {
398 return None;
399 }
400 let decoded = base64::engine::general_purpose::STANDARD
401 .decode(rest.trim())
402 .ok()?;
403 let text = String::from_utf8(decoded).ok()?;
404 let (_user, password) = text.split_once(':')?;
405 (!password.is_empty()).then(|| password.to_string())
406}
407
408fn push_credential(headers: &HeaderMap) -> Result<ServiceJwt, XrpcError> {
409 let value = headers
410 .get(AUTHORIZATION)
411 .and_then(|value| value.to_str().ok())
412 .ok_or_else(|| XrpcError::auth_required("missing authorization header"))?;
413 strip_bearer(value)
414 .map(str::trim)
415 .map(str::to_string)
416 .or_else(|| strip_basic(value))
417 .and_then(|token| ServiceJwt::new(token).ok())
418 .ok_or_else(|| {
419 XrpcError::auth_required("authorization isn't a bearer token or basic credential")
420 })
421}
422
423pub(crate) fn decode<T: DeserializeOwned>(body: &Bytes) -> Result<T, XrpcError> {
424 serde_json::from_slice(body)
425 .map_err(|error| XrpcError::invalid_request(format!("invalid request body: {error}")))
426}
427
428pub(crate) fn ok_empty() -> Response {
429 (StatusCode::OK, Json(json!({}))).into_response()
430}
431
432pub(crate) fn current_owner<H: HttpTransport, C: Clock>(
433 state: &XrpcState<H, C>,
434 repo: &RepoDid,
435) -> Option<OwnerDid> {
436 match state.index.owner_of(repo) {
437 Resolved::Ready(owner) => owner,
438 Resolved::Warming => None,
439 }
440}
441
442pub(crate) async fn fold_collaborators<H: HttpTransport, C: Clock>(
443 state: &XrpcState<H, C>,
444 repo: &RepoDid,
445) {
446 let index = Arc::clone(&state.index);
447 let target = repo.clone();
448 let _ = run_blocking(move || Ok(index.ensure_collaborators(&target))).await;
449}
450
451pub(crate) async fn authorize_push<H: HttpTransport, C: Clock>(
452 state: &XrpcState<H, C>,
453 actor: &AccountDid,
454 repo: &RepoDid,
455 denied: &str,
456) -> Result<(), XrpcError> {
457 fold_collaborators(state, repo).await;
458 let acl = knot_acl::KnotAcl::new(&state.admins, state.admission, &state.index);
459 if knot_acl::can_push(&acl, actor, repo).is_allowed() {
460 Ok(())
461 } else {
462 Err(XrpcError::forbidden(denied))
463 }
464}
465
466pub(crate) async fn authenticate_and_authorize_push<H: HttpTransport, C: Clock>(
467 state: &XrpcState<H, C>,
468 socket: SocketPeer,
469 headers: &HeaderMap,
470 repo: &RepoDid,
471 denied: &str,
472) -> Result<AccountDid, XrpcError> {
473 let peer = effective_peer(state, socket, headers);
474 let guard = admit_pre_auth(state, peer)?;
475 let actor = state.authenticate_push(headers).await?;
476 guard.refund();
477 authorize_push(state, &actor, repo, denied).await?;
478 Ok(actor)
479}
480
481pub(crate) async fn run_blocking<T, F>(task: F) -> Result<T, XrpcError>
482where
483 F: FnOnce() -> Result<T, XrpcError> + Send + 'static,
484 T: Send + 'static,
485{
486 match tokio::task::spawn_blocking(task).await {
487 Ok(result) => result,
488 Err(_) => Err(XrpcError::internal("blocking task failed to complete")),
489 }
490}
491
492#[derive(serde::Deserialize)]
493#[serde(transparent)]
494pub(crate) struct OwnerSegment(String);
495
496#[derive(serde::Deserialize)]
497#[serde(transparent)]
498pub(crate) struct RepoNameSegment(String);
499
500impl OwnerSegment {
501 pub(crate) fn as_str(&self) -> &str {
502 &self.0
503 }
504}
505
506impl RepoNameSegment {
507 pub(crate) fn as_str(&self) -> &str {
508 &self.0
509 }
510}
511
512#[derive(serde::Deserialize)]
513#[serde(transparent)]
514pub(crate) struct RepoDidSegment(String);
515
516impl RepoDidSegment {
517 pub(crate) fn as_str(&self) -> &str {
518 &self.0
519 }
520}
521
522#[derive(serde::Deserialize)]
523pub(crate) struct RepoPathParams {
524 pub(crate) did: OwnerSegment,
525 pub(crate) name: RepoNameSegment,
526}
527
528pub(crate) fn resolve_repo_did<H: HttpTransport, C: Clock>(
529 state: &XrpcState<H, C>,
530 segment: &RepoDidSegment,
531) -> Result<RepoDid, XrpcError> {
532 let raw = segment.as_str();
533 let trimmed = raw.strip_suffix(".git").unwrap_or(raw);
534 let did = RepoDid::new(trimmed).map_err(|_| XrpcError::not_found("repository not found"))?;
535 match state.index.owner_of(&did) {
536 Resolved::Ready(Some(_)) => Ok(did),
537 Resolved::Ready(None) => Err(XrpcError::not_found("repository not found")),
538 Resolved::Warming => Err(XrpcError::warming(
539 "registry projection is still warming, retry shortly",
540 )),
541 }
542}
543
544pub(crate) async fn resolve_repo_named<H: HttpTransport, C: Clock>(
545 state: &XrpcState<H, C>,
546 owner: &OwnerSegment,
547 name: &RepoNameSegment,
548) -> Result<RepoDid, XrpcError> {
549 let owner = resolve_owner_segment(state, owner).await?;
550 RepoRkey::clone_path_candidates(name.as_str())
551 .find_map(|rkey| match state.index.resolve_repo(&owner, &rkey) {
552 Resolved::Ready(Some(did)) => Some(Ok(did)),
553 Resolved::Ready(None) => None,
554 Resolved::Warming => Some(Err(XrpcError::warming(
555 "registry projection is still warming, retry shortly",
556 ))),
557 })
558 .unwrap_or_else(|| Err(XrpcError::not_found("repository not found")))
559}
560
561async fn resolve_owner_segment<H: HttpTransport, C: Clock>(
562 state: &XrpcState<H, C>,
563 owner: &OwnerSegment,
564) -> Result<OwnerDid, XrpcError> {
565 let not_found = || XrpcError::not_found("repository not found");
566 match OwnerRef::parse(owner.as_str()).ok_or_else(not_found)? {
567 OwnerRef::Did(did) => Ok(did),
568 OwnerRef::Handle(handle) => state
569 .atproto
570 .resolve_handle_to_did(&handle)
571 .await
572 .map(OwnerDid::from)
573 .map_err(|_| not_found()),
574 }
575}
576
577#[cfg(test)]
578mod credential_tests {
579 use super::push_credential;
580 use base64::Engine;
581 use http::{HeaderMap, HeaderValue, header::AUTHORIZATION};
582
583 fn with(value: &str) -> HeaderMap {
584 let mut headers = HeaderMap::new();
585 headers.insert(AUTHORIZATION, HeaderValue::from_str(value).unwrap());
586 headers
587 }
588
589 fn basic(user_pass: &str) -> String {
590 format!(
591 "Basic {}",
592 base64::engine::general_purpose::STANDARD.encode(user_pass)
593 )
594 }
595
596 #[test]
597 fn a_bearer_token_is_taken_verbatim() {
598 assert_eq!(
599 push_credential(&with("Bearer jwt.abc.def"))
600 .unwrap()
601 .as_str(),
602 "jwt.abc.def"
603 );
604 assert_eq!(
605 push_credential(&with("bearer jwt.abc.def"))
606 .unwrap()
607 .as_str(),
608 "jwt.abc.def"
609 );
610 }
611
612 #[test]
613 fn a_basic_credential_yields_the_password_after_the_first_colon() {
614 assert_eq!(
615 push_credential(&with(&basic("x-tangled-token:jwt.abc.def")))
616 .unwrap()
617 .as_str(),
618 "jwt.abc.def",
619 "RFC 7617 puts the token in the password half, so the username stays colon-free"
620 );
621 }
622
623 #[test]
624 fn malformed_or_empty_credentials_are_rejected() {
625 assert!(push_credential(&HeaderMap::new()).is_err());
626 assert!(push_credential(&with("Bearer ")).is_err());
627 assert!(push_credential(&with(&basic("x-tangled-token:"))).is_err());
628 assert!(push_credential(&with(&basic("no-colon"))).is_err());
629 assert!(push_credential(&with("Basic !!!not-base64")).is_err());
630 assert!(push_credential(&with("Digest whatever")).is_err());
631 }
632}