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 head = input.default_branch.map(|branch| branch.head_ref());
171
172 let owner = OwnerDid::new(actor.as_str()).expect("account DID is always a valid owner DID");
173 match state.index.resolve_repo(&owner, &input.rkey) {
174 Resolved::Ready(Some(existing)) => match state.index.rkey_of(&existing) {
175 Resolved::Ready(Some(canonical)) if canonical == input.rkey => {
176 return Err(XrpcError::conflict(
177 "repository with that record key already exists for this owner",
178 ));
179 }
180 Resolved::Warming => {
181 return Err(XrpcError::warming("registry projection is still warming"));
182 }
183 _ => {}
184 },
185 Resolved::Warming => {
186 return Err(XrpcError::warming("registry projection is still warming"));
187 }
188 Resolved::Ready(None) => {}
189 }
190
191 let fork = match &input.source {
192 Some(origin) => Some(crate::forks::ForkSource::resolve(&state, origin).await?),
193 None => None,
194 };
195 let object_format = fork.as_ref().map(|fork| fork.refs.object_format);
196
197 let (repo_did, provisioning) =
198 provision_repo_did(&state, &actor, &input.rkey, input.repo_did).await?;
199 let now = state.now();
200 let knot_signer = state.secrets.signer(&state.knot_did)?;
201 let key = ActorId::from_secp256k1(knot_signer.public_key().as_bytes());
202
203 let registration = Registration {
204 owner,
205 rkey: input.rkey,
206 name: input.name,
207 repo: repo_did.clone(),
208 created_at: now,
209 };
210
211 let lfs_store = state.lfs.as_ref().map(|web| Arc::clone(&web.handle.store));
212 let layout = state.layout.clone();
213 let placed = repo_did.clone();
214 let lfs = lfs_store.clone();
215 let provisioning = run_blocking(move || {
216 let repo = match object_format {
217 Some(format) => layout.create_with_format(&placed, format),
218 None => layout.create(&placed),
219 }
220 .map_err(|error| match error {
221 GitError::AlreadyExists(_) => XrpcError::conflict("repository already exists on disk"),
222 GitError::ReservedDid(_) => {
223 XrpcError::invalid_request("repoDid mustn't be knot's own identity")
224 }
225 other => XrpcError::internal(other.to_string()),
226 })?;
227 let outcome = stage_repo(&repo, head.as_ref());
228 if outcome.is_err() {
229 rollback_local(&layout, lfs.as_deref(), &placed);
230 }
231 outcome.map(|()| provisioning)
232 })
233 .await?;
234
235 let lfs_missing = match &fork {
236 Some(fork) => match populate(&state, &repo_did, fork).await {
237 Ok(missing) => missing,
238 Err(error) => {
239 let layout = state.layout.clone();
240 let placed = repo_did.clone();
241 let lfs = lfs_store.clone();
242 let _ = run_blocking(move || {
243 rollback_local(&layout, lfs.as_deref(), &placed);
244 Ok(())
245 })
246 .await;
247 return Err(error);
248 }
249 },
250 None => Vec::new(),
251 };
252
253 let submission = match &provisioning {
254 Provisioned::Minted(prepared) => state
255 .atproto
256 .submit_plc_operation(prepared)
257 .await
258 .map(|_| ()),
259 Provisioned::Reserved => Ok(()),
260 };
261 if let Err(error) = submission {
262 let layout = state.layout.clone();
263 let placed = repo_did.clone();
264 let lfs = lfs_store.clone();
265 let _ = run_blocking(move || {
266 rollback_local(&layout, lfs.as_deref(), &placed);
267 Ok(())
268 })
269 .await;
270 return Err(error.into());
271 }
272
273 let reserved = matches!(provisioning, Provisioned::Reserved);
274 let layout = state.layout.clone();
275 let cob_locks = Arc::clone(&state.cob_locks);
276 let meta_path = state.meta_path.clone();
277 let placed = repo_did.clone();
278 let lfs = lfs_store.clone();
279 let home = CobHome::from(&state.knot_did);
280 run_blocking(move || {
281 let outcome = {
282 let _guard = cob_locks.meta();
283 Repo::open(&meta_path)
284 .map_err(XrpcError::from)
285 .and_then(|meta| {
286 register(
287 &CobStore::new(&meta),
288 &home,
289 registration,
290 &knot_signer,
291 now,
292 )
293 })
294 };
295 if outcome.is_err() {
296 rollback_local(&layout, lfs.as_deref(), &placed);
297 if !reserved {
298 tracing::error!(
299 repo = %placed,
300 "registration failed after did:plc submitted to PLC directory"
301 );
302 }
303 }
304 outcome
305 })
306 .await?;
307
308 if reserved {
309 state.reservations.release(&repo_did);
310 }
311
312 let index = Arc::clone(&state.index);
313 let refreshed = repo_did.clone();
314 run_blocking(move || {
315 index.refresh_registry().map_err(|error| {
316 XrpcError::internal(format!(
317 "repo {refreshed} was created and registered but registry projection refresh failed: {error}"
318 ))
319 })
320 })
321 .await?;
322
323 Ok((
324 http::StatusCode::OK,
325 Json(CreateOutput {
326 repo_did,
327 key,
328 lfs_missing,
329 }),
330 )
331 .into_response())
332}
333
334async fn populate<H: HttpTransport, C: Clock>(
335 state: &Arc<XrpcState<H, C>>,
336 repo_did: &RepoDid,
337 fork: &crate::forks::ForkSource,
338) -> Result<Vec<knot_lfs::LfsOid>, XrpcError> {
339 let tips = fork.refs.tips();
340 let pack = crate::forks::upstream_pack(
341 state,
342 &fork.upstream,
343 knot_pack::WantOids::new(tips.clone()),
344 knot_pack::HaveOids::default(),
345 )
346 .await?;
347 let layout = state.layout.clone();
348 let placed = repo_did.clone();
349 let origin = fork.origin.clone();
350 let refs = fork.refs.clone();
351 run_blocking(move || {
352 let repo = layout.open(&placed)?;
353 crate::forks::populate_fork(&repo, &refs, &pack, &origin)
354 })
355 .await?;
356 crate::lfs::mirror_fork_objects(
357 Arc::clone(state),
358 fork.upstream.clone(),
359 repo_did.clone(),
360 knot_pack::WantOids::new(tips),
361 knot_pack::HaveOids::default(),
362 )
363 .await
364}
365
366fn rollback_local(layout: &Layout, lfs: Option<&knot_lfs::DiskStore>, repo_did: &RepoDid) {
367 if let Some(store) = lfs
368 && let Err(error) = store.remove_repo(repo_did)
369 {
370 tracing::error!(repo = %repo_did, %error, "couldn't roll back lfs prefix");
371 }
372 if let Err(error) = layout.remove(repo_did) {
373 tracing::error!(repo = %repo_did, %error, "couldn't roll back on-disk repo :3");
374 }
375}
376
377enum Provisioned {
378 Minted(PreparedRepoDid),
379 Reserved,
380}
381
382async fn provision_repo_did<H: HttpTransport, C: Clock>(
383 state: &XrpcState<H, C>,
384 actor: &knot_types::AccountDid,
385 rkey: &RepoRkey,
386 provided: Option<RepoDid>,
387) -> Result<(RepoDid, Provisioned), XrpcError> {
388 match provided {
389 Some(did) if did.as_str().starts_with("did:web:") => {
390 match state.index.owner_of(&did) {
391 Resolved::Ready(Some(_)) => {
392 return Err(XrpcError::conflict(
393 "that repo DID is already hosted on this knot",
394 ));
395 }
396 Resolved::Warming => {
397 return Err(XrpcError::warming("registry projection is still warming"));
398 }
399 Resolved::Ready(None) => {}
400 }
401 if !state.reservations.holder_is(&did, actor, state.now()) {
402 return Err(XrpcError::invalid_request(
403 "reserve this did:web for your own account via sh.tangled.repo.reserveKey before creating it",
404 ));
405 }
406 let knot_public = state.secrets.public_key(&state.knot_did)?;
407 state
408 .atproto
409 .verify_did_web_publishes_key(&did, &knot_public)
410 .await
411 .map_err(XrpcError::from)?;
412 Ok((did, Provisioned::Reserved))
413 }
414 Some(_) => Err(XrpcError::invalid_request(
415 "repoDid must be did:web hosted on your own domain. Omit it to mint a did:plc.",
416 )),
417 None => {
418 let knot_signer = state.secrets.signer(&state.knot_did)?;
419 let owner =
420 OwnerDid::new(actor.as_str()).expect("account DID is always a valid owner DID");
421 let nonce = knot_atproto::MintNonce::mint(&*state.entropy, &owner, rkey);
422 let prepared =
423 knot_atproto::prepare_repo_did(&knot_signer, &state.knot_service_url, &nonce)
424 .map_err(XrpcError::from)?;
425 Ok((prepared.did.clone(), Provisioned::Minted(prepared)))
426 }
427 }
428}
429
430fn stage_repo(repo: &Repo, head: Option<&RefName>) -> Result<(), XrpcError> {
431 if let Some(refname) = head {
432 repo.set_head(refname)?;
433 }
434 Ok(())
435}
436
437fn register(
438 store: &CobStore,
439 home: &CobHome,
440 registration: Registration,
441 signer: &dyn Signer,
442 now: UnixSeconds,
443) -> Result<(), XrpcError> {
444 match store
445 .list::<RepoRegistryCob>()
446 .map_err(XrpcError::from)?
447 .as_slice()
448 {
449 [] => store
450 .create(home, &RegistryChange::Register(registration), signer, now)
451 .map(|_| ())
452 .map_err(XrpcError::from),
453 [object] => register_repo(store, home, *object, registration, signer, now)
454 .map(|_| ())
455 .map_err(XrpcError::from),
456 many => Err(XrpcError::internal(format!(
457 "{} repo registry objects share meta-repo",
458 many.len()
459 ))),
460 }
461}
462
463pub(crate) async fn delete_repo<H: HttpTransport, C: Clock>(
464 State(state): State<Arc<XrpcState<H, C>>>,
465 headers: HeaderMap,
466 method: crate::Method,
467 body: Bytes,
468) -> Result<Response, XrpcError> {
469 let actor = state.authenticate(&headers, &method).await?;
470 let DeleteInput {
471 did, rkey, force, ..
472 } = decode(&body)?;
473
474 let repo_did = match state.index.resolve_repo(&did, &rkey) {
475 Resolved::Ready(Some(repo_did)) => repo_did,
476 Resolved::Ready(None) => return Ok(ok_empty()),
477 Resolved::Warming => {
478 return Err(XrpcError::warming("registry projection is still warming"));
479 }
480 };
481
482 let acl = KnotAcl::new(&state.admins, state.admission, &state.index);
483 if !can_delete_repo(&acl, &actor, &repo_did).is_allowed() {
484 return Err(XrpcError::forbidden(
485 "only repository owner or a knot admin may delete it",
486 ));
487 }
488
489 if force {
490 if !can_admin_knot(&acl, &actor).is_allowed() {
491 return Err(XrpcError::forbidden(
492 "only knot admin may force a delete past the PDS record check",
493 ));
494 }
495 } else {
496 let owner = AccountDid::from(did.clone());
497 match state.atproto.repo_record_present(&owner, &rkey).await {
498 Ok(RecordPresence::Present) => {
499 return Err(XrpcError::conflict(
500 "sh.tangled.repo record still exists on the owner's PDS. Remove it there first or force the delete.",
501 ));
502 }
503 Ok(RecordPresence::Absent) => {}
504 Err(error) => {
505 tracing::warn!(
506 repo = %repo_did,
507 %error,
508 "proceeding w/ best-effort delete despite unconfirmed owner PDS record :p"
509 );
510 }
511 }
512 }
513
514 let now = state.now();
515 let knot_signer = state.secrets.signer(&state.knot_did)?;
516 let layout = state.layout.clone();
517 let index = Arc::clone(&state.index);
518 let cob_locks = Arc::clone(&state.cob_locks);
519 let meta_path = state.meta_path.clone();
520 let target = RepoRef { owner: did, rkey };
521 let deleted = repo_did.clone();
522 let home = CobHome::from(&state.knot_did);
523 let lfs_store = state.lfs.as_ref().map(|lfs| Arc::clone(&lfs.handle.store));
524 run_blocking(move || {
525 let _repo_guard = cob_locks.repo(&deleted);
526 let _meta_guard = cob_locks.meta();
527 let meta = Repo::open(&meta_path)?;
528 let store = CobStore::new(&meta);
529 deregister(&store, &home, target, &deleted, &knot_signer, now)?;
530 let removal = layout.remove(&deleted);
531 index.refresh_registry().map_err(|error| {
532 XrpcError::internal(format!(
533 "repo {deleted} was deregistered but registry projection refresh failed: {error}"
534 ))
535 })?;
536 if let Some(store) = &lfs_store
537 && let Err(error) = store.remove_repo(&deleted)
538 {
539 tracing::warn!(
540 repo = %deleted,
541 %error,
542 "lfs prefix removal on delete failed, orphan sweep will reclaim"
543 );
544 }
545 removal.map_err(|error| {
546 XrpcError::internal(format!(
547 "repo {deleted} was deregistered but its on-disk directory couldn't be removed: {error}"
548 ))
549 })
550 })
551 .await?;
552
553 Ok(ok_empty())
554}
555
556fn deregister(
557 store: &CobStore,
558 home: &CobHome,
559 target: RepoRef,
560 expected: &RepoDid,
561 signer: &dyn Signer,
562 now: UnixSeconds,
563) -> Result<(), XrpcError> {
564 match store
565 .list::<RepoRegistryCob>()
566 .map_err(XrpcError::from)?
567 .as_slice()
568 {
569 [] => Ok(()),
570 [object] => deregister_repo(store, home, *object, target, expected.clone(), signer, now)
571 .map(|_| ())
572 .map_err(XrpcError::from),
573 many => Err(XrpcError::internal(format!(
574 "{} repo registry objects share meta-repo",
575 many.len()
576 ))),
577 }
578}
579
580pub(crate) async fn rename_repo<H: HttpTransport, C: Clock>(
581 State(state): State<Arc<XrpcState<H, C>>>,
582 headers: HeaderMap,
583 method: crate::Method,
584 body: Bytes,
585) -> Result<Response, XrpcError> {
586 let actor = state.authenticate(&headers, &method).await?;
587 let RenameInput { repo, rkey, name } = decode(&body)?;
588
589 let owner = match state.index.owner_of(&repo) {
590 Resolved::Ready(Some(owner)) => owner,
591 Resolved::Ready(None) => {
592 return Err(XrpcError::not_found("no such repository on this knot"));
593 }
594 Resolved::Warming => {
595 return Err(XrpcError::warming("registry projection is still warming"));
596 }
597 };
598
599 crate::authorize_push(
600 &state,
601 &actor,
602 &repo,
603 "only repository owner or a collaborator may rename it",
604 )
605 .await?;
606
607 let now = state.now();
608 let knot_signer = state.secrets.signer(&state.knot_did)?;
609 let index = Arc::clone(&state.index);
610 let cob_locks = Arc::clone(&state.cob_locks);
611 let meta_path = state.meta_path.clone();
612 let renamed = repo.clone();
613 let home = CobHome::from(&state.knot_did);
614 run_blocking(move || {
615 {
616 let _guard = cob_locks.meta();
617 let meta = Repo::open(&meta_path)?;
618 let store = CobStore::new(&meta);
619 match store
620 .list::<RepoRegistryCob>()
621 .map_err(XrpcError::from)?
622 .as_slice()
623 {
624 [] => {
625 return Err(XrpcError::not_found(
626 "no repositories are registered on this knot",
627 ));
628 }
629 [object] => {
630 knot_cobs::rename_repo(
631 &store,
632 &home,
633 *object,
634 Rename {
635 owner,
636 rkey,
637 name,
638 repo: renamed.clone(),
639 },
640 &knot_signer,
641 now,
642 )
643 .map(|_| ())
644 .map_err(XrpcError::from)?;
645 }
646 many => {
647 return Err(XrpcError::internal(format!(
648 "{} repo registry objects share meta-repo",
649 many.len()
650 )));
651 }
652 }
653 }
654 index.refresh_registry().map_err(|error| {
655 XrpcError::internal(format!(
656 "repo {renamed} was renamed but registry projection refresh failed: {error}"
657 ))
658 })
659 })
660 .await?;
661
662 Ok(ok_empty())
663}