This repository has no description
1use std::sync::Arc;
2
3use axum::Json;
4use axum::body::Bytes;
5use axum::extract::State;
6use axum::response::{IntoResponse, Response};
7use http::HeaderMap;
8use serde::{Deserialize, Serialize};
9
10use knot_acl::{KnotAcl, can_admin_knot, can_create_repo, can_delete_repo};
11use knot_atproto::{PreparedRepoDid, RecordPresence};
12use knot_cob::{CobHome, CobStore};
13use knot_cobs::{
14 Registration, RegistryChange, Rename, RepoRef, RepoRegistryCob, deregister_repo, register_repo,
15};
16use knot_git::{GitError, Layout, Repo};
17use knot_index::Resolved;
18use knot_runtime::{Clock, HttpTransport, Signer};
19use knot_types::{
20 AccountDid, ActorId, BranchName, OwnerDid, RefName, RepoDid, RepoName, RepoRkey, UnixSeconds,
21};
22
23use crate::body::SourceUrl;
24use crate::error::XrpcError;
25use crate::reservations::ReserveDecision;
26use crate::{XrpcState, decode, ok_empty, run_blocking};
27
28pub(crate) const CREATE_ROUTE: &str = "/xrpc/sh.tangled.repo.create";
29pub(crate) const DELETE_ROUTE: &str = "/xrpc/sh.tangled.repo.delete";
30pub(crate) const RENAME_ROUTE: &str = "/xrpc/sh.tangled.repo.rename";
31pub(crate) const RESERVE_ROUTE: &str = "/xrpc/sh.tangled.repo.reserveKey";
32
33#[derive(Deserialize)]
34struct CreateInput {
35 rkey: RepoRkey,
36 name: RepoName,
37 #[serde(rename = "defaultBranch")]
38 default_branch: Option<BranchName>,
39 #[serde(default, deserialize_with = "crate::body::optional_source_url")]
40 source: Option<SourceUrl>,
41 #[serde(rename = "repoDid")]
42 repo_did: Option<RepoDid>,
43}
44
45#[derive(Serialize)]
46struct CreateOutput {
47 #[serde(rename = "repoDid")]
48 repo_did: RepoDid,
49 key: ActorId,
50 #[serde(rename = "lfsMissing", skip_serializing_if = "Vec::is_empty")]
51 lfs_missing: Vec<knot_lfs::LfsOid>,
52}
53
54#[derive(Deserialize)]
55struct DeleteInput {
56 did: OwnerDid,
57 rkey: RepoRkey,
58 #[serde(rename = "name")]
59 _name: RepoName,
60 #[serde(default)]
61 force: bool,
62}
63
64#[derive(Deserialize)]
65struct RenameInput {
66 repo: RepoDid,
67 rkey: RepoRkey,
68 name: RepoName,
69}
70
71#[derive(Deserialize)]
72struct ReserveInput {
73 #[serde(rename = "repoDid")]
74 repo_did: RepoDid,
75}
76
77#[derive(Serialize)]
78struct ReserveOutput {
79 #[serde(rename = "repoDid")]
80 repo_did: RepoDid,
81 key: ActorId,
82}
83
84pub(crate) async fn reserve_key<H: HttpTransport, C: Clock>(
85 State(state): State<Arc<XrpcState<H, C>>>,
86 headers: HeaderMap,
87 method: crate::Method,
88 body: Bytes,
89) -> Result<Response, XrpcError> {
90 let actor = state.authenticate(&headers, &method).await?;
91 let acl = KnotAcl::new(&state.admins, state.admission, &state.index);
92 if !can_create_repo(&acl, &actor).is_allowed() {
93 return Err(XrpcError::forbidden(
94 "only knot admin or member may reserve a repository key",
95 ));
96 }
97
98 let ReserveInput { repo_did } = decode(&body)?;
99 if !repo_did.as_str().starts_with("did:web:") {
100 return Err(XrpcError::invalid_request(
101 "only did:web repo identity needs a reserved key. Omit repoDid on create to mint a did:plc.",
102 ));
103 }
104 if repo_did.as_str() == state.knot_did.as_str() {
105 return Err(XrpcError::invalid_request(
106 "repoDid mustn't be knot's own identity",
107 ));
108 }
109 match state.index.owner_of(&repo_did) {
110 Resolved::Ready(Some(_)) => {
111 return Err(XrpcError::conflict(
112 "that repo DID is already hosted on this knot",
113 ));
114 }
115 Resolved::Warming => {
116 return Err(XrpcError::warming("registry projection is still warming"));
117 }
118 Resolved::Ready(None) => {}
119 }
120
121 let now = state.now();
122 state.reservations.prune(now);
123
124 match state.reservations.try_reserve(&repo_did, &actor, now) {
125 ReserveDecision::HeldByOther => {
126 return Err(XrpcError::conflict(
127 "that repo DID is reserved by another account",
128 ));
129 }
130 ReserveDecision::PerActorFull => {
131 return Err(XrpcError::rate_limited(
132 "you are holding maximum number of reserved repository keys awaiting creation",
133 ));
134 }
135 ReserveDecision::GlobalFull => {
136 return Err(XrpcError::rate_limited(
137 "knot is holding maximum number of reserved repository keys awaiting creation",
138 ));
139 }
140 ReserveDecision::Fresh | ReserveDecision::Renewed => {}
141 }
142
143 let public = match state.secrets.public_key(&state.knot_did) {
144 Ok(public) => public,
145 Err(error) => {
146 state.reservations.release(&repo_did);
147 return Err(error.into());
148 }
149 };
150 let key = ActorId::from_secp256k1(public.as_bytes());
151
152 Ok((http::StatusCode::OK, Json(ReserveOutput { repo_did, key })).into_response())
153}
154
155pub(crate) async fn create_repo<H: HttpTransport, C: Clock>(
156 State(state): State<Arc<XrpcState<H, C>>>,
157 headers: HeaderMap,
158 method: crate::Method,
159 body: Bytes,
160) -> Result<Response, XrpcError> {
161 let actor = state.authenticate(&headers, &method).await?;
162 let acl = KnotAcl::new(&state.admins, state.admission, &state.index);
163 if !can_create_repo(&acl, &actor).is_allowed() {
164 return Err(XrpcError::forbidden(
165 "only knot admin or member may create repositories",
166 ));
167 }
168
169 let input: CreateInput = decode(&body)?;
170 let source = input
171 .source
172 .as_ref()
173 .map(|source| {
174 crate::forks::resolve_upstream(&state, source).map(|upstream| (source, upstream))
175 })
176 .transpose()?;
177 let head = input.default_branch.map(|branch| branch.head_ref());
178
179 let owner = OwnerDid::new(actor.as_str()).expect("account DID is always a valid owner DID");
180 match state.index.resolve_repo(&owner, &input.rkey) {
181 Resolved::Ready(Some(existing)) => match state.index.rkey_of(&existing) {
182 Resolved::Ready(Some(canonical)) if canonical == input.rkey => {
183 return Err(XrpcError::conflict(
184 "repository with that record key already exists for this owner",
185 ));
186 }
187 Resolved::Warming => {
188 return Err(XrpcError::warming("registry projection is still warming"));
189 }
190 _ => {}
191 },
192 Resolved::Warming => {
193 return Err(XrpcError::warming("registry projection is still warming"));
194 }
195 Resolved::Ready(None) => {}
196 }
197
198 let (repo_did, provisioning) =
199 provision_repo_did(&state, &actor, &input.rkey, input.repo_did).await?;
200 let now = state.now();
201 let knot_signer = state.secrets.signer(&state.knot_did)?;
202 let key = ActorId::from_secp256k1(knot_signer.public_key().as_bytes());
203
204 let registration = Registration {
205 owner,
206 rkey: input.rkey,
207 name: input.name,
208 repo: repo_did.clone(),
209 created_at: now,
210 };
211
212 let lfs_store = state.lfs.as_ref().map(|web| Arc::clone(&web.handle.store));
213 let layout = state.layout.clone();
214 let placed = repo_did.clone();
215 let lfs = lfs_store.clone();
216 let provisioning = run_blocking(move || {
217 let repo = layout.create(&placed).map_err(|error| match error {
218 GitError::AlreadyExists(_) => XrpcError::conflict("repository already exists on disk"),
219 GitError::ReservedDid(_) => {
220 XrpcError::invalid_request("repoDid mustn't be knot's own identity")
221 }
222 other => XrpcError::internal(other.to_string()),
223 })?;
224 let outcome = stage_repo(&repo, head.as_ref());
225 if outcome.is_err() {
226 rollback_local(&layout, lfs.as_deref(), &placed);
227 }
228 outcome.map(|()| provisioning)
229 })
230 .await?;
231
232 let lfs_missing = match &source {
233 Some((source_url, upstream)) => {
234 match populate(&state, &repo_did, source_url, upstream).await {
235 Ok(missing) => missing,
236 Err(error) => {
237 let layout = state.layout.clone();
238 let placed = repo_did.clone();
239 let lfs = lfs_store.clone();
240 let _ = run_blocking(move || {
241 rollback_local(&layout, lfs.as_deref(), &placed);
242 Ok(())
243 })
244 .await;
245 return Err(error);
246 }
247 }
248 }
249 None => Vec::new(),
250 };
251
252 let submission = match &provisioning {
253 Provisioned::Minted(prepared) => state
254 .atproto
255 .submit_plc_operation(prepared)
256 .await
257 .map(|_| ()),
258 Provisioned::Reserved => Ok(()),
259 };
260 if let Err(error) = submission {
261 let layout = state.layout.clone();
262 let placed = repo_did.clone();
263 let lfs = lfs_store.clone();
264 let _ = run_blocking(move || {
265 rollback_local(&layout, lfs.as_deref(), &placed);
266 Ok(())
267 })
268 .await;
269 return Err(error.into());
270 }
271
272 let reserved = matches!(provisioning, Provisioned::Reserved);
273 let layout = state.layout.clone();
274 let cob_locks = Arc::clone(&state.cob_locks);
275 let meta_path = state.meta_path.clone();
276 let placed = repo_did.clone();
277 let lfs = lfs_store.clone();
278 let home = CobHome::from(&state.knot_did);
279 run_blocking(move || {
280 let outcome = {
281 let _guard = cob_locks.meta();
282 Repo::open(&meta_path)
283 .map_err(XrpcError::from)
284 .and_then(|meta| {
285 register(
286 &CobStore::new(&meta),
287 &home,
288 registration,
289 &knot_signer,
290 now,
291 )
292 })
293 };
294 if outcome.is_err() {
295 rollback_local(&layout, lfs.as_deref(), &placed);
296 if !reserved {
297 tracing::error!(
298 repo = %placed,
299 "registration failed after did:plc submitted to PLC directory"
300 );
301 }
302 }
303 outcome
304 })
305 .await?;
306
307 if reserved {
308 state.reservations.release(&repo_did);
309 }
310
311 let index = Arc::clone(&state.index);
312 let refreshed = repo_did.clone();
313 run_blocking(move || {
314 index.refresh_registry().map_err(|error| {
315 XrpcError::internal(format!(
316 "repo {refreshed} was created and registered but registry projection refresh failed: {error}"
317 ))
318 })
319 })
320 .await?;
321
322 Ok((
323 http::StatusCode::OK,
324 Json(CreateOutput {
325 repo_did,
326 key,
327 lfs_missing,
328 }),
329 )
330 .into_response())
331}
332
333async fn populate<H: HttpTransport, C: Clock>(
334 state: &Arc<XrpcState<H, C>>,
335 repo_did: &RepoDid,
336 source: &SourceUrl,
337 upstream: &crate::forks::Upstream,
338) -> Result<Vec<knot_lfs::LfsOid>, XrpcError> {
339 let prefixes = vec![
340 "HEAD".to_string(),
341 "refs/heads/".to_string(),
342 "refs/tags/".to_string(),
343 ];
344 let refs = crate::forks::upstream_refs(state, upstream, prefixes).await?;
345 let tips = refs.tips();
346 let pack = crate::forks::upstream_pack(
347 state,
348 upstream,
349 knot_pack::WantOids::new(refs.tips()),
350 knot_pack::HaveOids::default(),
351 )
352 .await?;
353 let layout = state.layout.clone();
354 let placed = repo_did.clone();
355 let origin = source.clone();
356 run_blocking(move || {
357 let repo = layout.open(&placed)?;
358 crate::forks::populate_fork(&repo, &refs, &pack, &origin)
359 })
360 .await?;
361 crate::lfs::mirror_fork_objects(
362 Arc::clone(state),
363 upstream.clone(),
364 repo_did.clone(),
365 knot_pack::WantOids::new(tips),
366 knot_pack::HaveOids::default(),
367 )
368 .await
369}
370
371fn rollback_local(layout: &Layout, lfs: Option<&knot_lfs::DiskStore>, repo_did: &RepoDid) {
372 if let Some(store) = lfs
373 && let Err(error) = store.remove_repo(repo_did)
374 {
375 tracing::error!(repo = %repo_did, %error, "couldn't roll back lfs prefix");
376 }
377 if let Err(error) = layout.remove(repo_did) {
378 tracing::error!(repo = %repo_did, %error, "couldn't roll back on-disk repo :3");
379 }
380}
381
382enum Provisioned {
383 Minted(PreparedRepoDid),
384 Reserved,
385}
386
387async fn provision_repo_did<H: HttpTransport, C: Clock>(
388 state: &XrpcState<H, C>,
389 actor: &knot_types::AccountDid,
390 rkey: &RepoRkey,
391 provided: Option<RepoDid>,
392) -> Result<(RepoDid, Provisioned), XrpcError> {
393 match provided {
394 Some(did) if did.as_str().starts_with("did:web:") => {
395 match state.index.owner_of(&did) {
396 Resolved::Ready(Some(_)) => {
397 return Err(XrpcError::conflict(
398 "that repo DID is already hosted on this knot",
399 ));
400 }
401 Resolved::Warming => {
402 return Err(XrpcError::warming("registry projection is still warming"));
403 }
404 Resolved::Ready(None) => {}
405 }
406 if !state.reservations.holder_is(&did, actor, state.now()) {
407 return Err(XrpcError::invalid_request(
408 "reserve this did:web for your own account via sh.tangled.repo.reserveKey before creating it",
409 ));
410 }
411 let knot_public = state.secrets.public_key(&state.knot_did)?;
412 state
413 .atproto
414 .verify_did_web_publishes_key(&did, &knot_public)
415 .await
416 .map_err(XrpcError::from)?;
417 Ok((did, Provisioned::Reserved))
418 }
419 Some(_) => Err(XrpcError::invalid_request(
420 "repoDid must be did:web hosted on your own domain. Omit it to mint a did:plc.",
421 )),
422 None => {
423 let knot_signer = state.secrets.signer(&state.knot_did)?;
424 let owner =
425 OwnerDid::new(actor.as_str()).expect("account DID is always a valid owner DID");
426 let nonce = knot_atproto::MintNonce::mint(&*state.entropy, &owner, rkey);
427 let prepared =
428 knot_atproto::prepare_repo_did(&knot_signer, &state.knot_service_url, &nonce)
429 .map_err(XrpcError::from)?;
430 Ok((prepared.did.clone(), Provisioned::Minted(prepared)))
431 }
432 }
433}
434
435fn stage_repo(repo: &Repo, head: Option<&RefName>) -> Result<(), XrpcError> {
436 if let Some(refname) = head {
437 repo.set_head(refname)?;
438 }
439 Ok(())
440}
441
442fn register(
443 store: &CobStore,
444 home: &CobHome,
445 registration: Registration,
446 signer: &dyn Signer,
447 now: UnixSeconds,
448) -> Result<(), XrpcError> {
449 match store
450 .list::<RepoRegistryCob>()
451 .map_err(XrpcError::from)?
452 .as_slice()
453 {
454 [] => store
455 .create(home, &RegistryChange::Register(registration), signer, now)
456 .map(|_| ())
457 .map_err(XrpcError::from),
458 [object] => register_repo(store, home, *object, registration, signer, now)
459 .map(|_| ())
460 .map_err(XrpcError::from),
461 many => Err(XrpcError::internal(format!(
462 "{} repo registry objects share meta-repo",
463 many.len()
464 ))),
465 }
466}
467
468pub(crate) async fn delete_repo<H: HttpTransport, C: Clock>(
469 State(state): State<Arc<XrpcState<H, C>>>,
470 headers: HeaderMap,
471 method: crate::Method,
472 body: Bytes,
473) -> Result<Response, XrpcError> {
474 let actor = state.authenticate(&headers, &method).await?;
475 let DeleteInput {
476 did, rkey, force, ..
477 } = decode(&body)?;
478
479 let repo_did = match state.index.resolve_repo(&did, &rkey) {
480 Resolved::Ready(Some(repo_did)) => repo_did,
481 Resolved::Ready(None) => return Ok(ok_empty()),
482 Resolved::Warming => {
483 return Err(XrpcError::warming("registry projection is still warming"));
484 }
485 };
486
487 let acl = KnotAcl::new(&state.admins, state.admission, &state.index);
488 if !can_delete_repo(&acl, &actor, &repo_did).is_allowed() {
489 return Err(XrpcError::forbidden(
490 "only repository owner or a knot admin may delete it",
491 ));
492 }
493
494 if force {
495 if !can_admin_knot(&acl, &actor).is_allowed() {
496 return Err(XrpcError::forbidden(
497 "only knot admin may force a delete past the PDS record check",
498 ));
499 }
500 } else {
501 let owner = AccountDid::from(did.clone());
502 match state.atproto.repo_record_present(&owner, &rkey).await {
503 Ok(RecordPresence::Present) => {
504 return Err(XrpcError::conflict(
505 "sh.tangled.repo record still exists on the owner's PDS. Remove it there first or force the delete.",
506 ));
507 }
508 Ok(RecordPresence::Absent) => {}
509 Err(error) => {
510 tracing::warn!(
511 repo = %repo_did,
512 %error,
513 "proceeding w/ best-effort delete despite unconfirmed owner PDS record :p"
514 );
515 }
516 }
517 }
518
519 let now = state.now();
520 let knot_signer = state.secrets.signer(&state.knot_did)?;
521 let layout = state.layout.clone();
522 let index = Arc::clone(&state.index);
523 let cob_locks = Arc::clone(&state.cob_locks);
524 let meta_path = state.meta_path.clone();
525 let target = RepoRef { owner: did, rkey };
526 let deleted = repo_did.clone();
527 let home = CobHome::from(&state.knot_did);
528 let lfs_store = state.lfs.as_ref().map(|lfs| Arc::clone(&lfs.handle.store));
529 run_blocking(move || {
530 let _repo_guard = cob_locks.repo(&deleted);
531 let _meta_guard = cob_locks.meta();
532 let meta = Repo::open(&meta_path)?;
533 let store = CobStore::new(&meta);
534 deregister(&store, &home, target, &deleted, &knot_signer, now)?;
535 let removal = layout.remove(&deleted);
536 index.refresh_registry().map_err(|error| {
537 XrpcError::internal(format!(
538 "repo {deleted} was deregistered but registry projection refresh failed: {error}"
539 ))
540 })?;
541 if let Some(store) = &lfs_store
542 && let Err(error) = store.remove_repo(&deleted)
543 {
544 tracing::warn!(
545 repo = %deleted,
546 %error,
547 "lfs prefix removal on delete failed, orphan sweep will reclaim"
548 );
549 }
550 removal.map_err(|error| {
551 XrpcError::internal(format!(
552 "repo {deleted} was deregistered but its on-disk directory couldn't be removed: {error}"
553 ))
554 })
555 })
556 .await?;
557
558 Ok(ok_empty())
559}
560
561fn deregister(
562 store: &CobStore,
563 home: &CobHome,
564 target: RepoRef,
565 expected: &RepoDid,
566 signer: &dyn Signer,
567 now: UnixSeconds,
568) -> Result<(), XrpcError> {
569 match store
570 .list::<RepoRegistryCob>()
571 .map_err(XrpcError::from)?
572 .as_slice()
573 {
574 [] => Ok(()),
575 [object] => deregister_repo(store, home, *object, target, expected.clone(), signer, now)
576 .map(|_| ())
577 .map_err(XrpcError::from),
578 many => Err(XrpcError::internal(format!(
579 "{} repo registry objects share meta-repo",
580 many.len()
581 ))),
582 }
583}
584
585pub(crate) async fn rename_repo<H: HttpTransport, C: Clock>(
586 State(state): State<Arc<XrpcState<H, C>>>,
587 headers: HeaderMap,
588 method: crate::Method,
589 body: Bytes,
590) -> Result<Response, XrpcError> {
591 let actor = state.authenticate(&headers, &method).await?;
592 let RenameInput { repo, rkey, name } = decode(&body)?;
593
594 let owner = match state.index.owner_of(&repo) {
595 Resolved::Ready(Some(owner)) => owner,
596 Resolved::Ready(None) => {
597 return Err(XrpcError::not_found("no such repository on this knot"));
598 }
599 Resolved::Warming => {
600 return Err(XrpcError::warming("registry projection is still warming"));
601 }
602 };
603
604 crate::authorize_push(
605 &state,
606 &actor,
607 &repo,
608 "only repository owner or a collaborator may rename it",
609 )
610 .await?;
611
612 let now = state.now();
613 let knot_signer = state.secrets.signer(&state.knot_did)?;
614 let index = Arc::clone(&state.index);
615 let cob_locks = Arc::clone(&state.cob_locks);
616 let meta_path = state.meta_path.clone();
617 let renamed = repo.clone();
618 let home = CobHome::from(&state.knot_did);
619 run_blocking(move || {
620 {
621 let _guard = cob_locks.meta();
622 let meta = Repo::open(&meta_path)?;
623 let store = CobStore::new(&meta);
624 match store
625 .list::<RepoRegistryCob>()
626 .map_err(XrpcError::from)?
627 .as_slice()
628 {
629 [] => {
630 return Err(XrpcError::not_found(
631 "no repositories are registered on this knot",
632 ));
633 }
634 [object] => {
635 knot_cobs::rename_repo(
636 &store,
637 &home,
638 *object,
639 Rename {
640 owner,
641 rkey,
642 name,
643 repo: renamed.clone(),
644 },
645 &knot_signer,
646 now,
647 )
648 .map(|_| ())
649 .map_err(XrpcError::from)?;
650 }
651 many => {
652 return Err(XrpcError::internal(format!(
653 "{} repo registry objects share meta-repo",
654 many.len()
655 )));
656 }
657 }
658 }
659 index.refresh_registry().map_err(|error| {
660 XrpcError::internal(format!(
661 "repo {renamed} was renamed but registry projection refresh failed: {error}"
662 ))
663 })
664 })
665 .await?;
666
667 Ok(ok_empty())
668}