This repository has no description
1package repo
2
3import (
4 "context"
5 "database/sql"
6 "errors"
7 "fmt"
8 "log/slog"
9 "net/http"
10 "net/url"
11 "slices"
12 "strings"
13 "time"
14
15 "tangled.org/core/appview/cloudflare"
16 "tangled.org/core/appview/codesearch"
17
18 "tangled.org/core/api/tangled"
19 "tangled.org/core/appview/config"
20 "tangled.org/core/appview/db"
21 "tangled.org/core/appview/knotacl"
22 "tangled.org/core/appview/knotcompat"
23 "tangled.org/core/appview/models"
24 "tangled.org/core/appview/notify"
25 "tangled.org/core/appview/oauth"
26 "tangled.org/core/appview/pages"
27 "tangled.org/core/appview/pagination"
28 "tangled.org/core/appview/reporesolver"
29 "tangled.org/core/appview/sites"
30 "tangled.org/core/consts"
31 "tangled.org/core/idresolver"
32 "tangled.org/core/ogre"
33 "tangled.org/core/orm"
34 "tangled.org/core/rbac"
35 "tangled.org/core/tid"
36 "tangled.org/core/xrpc/serviceauth"
37 xrpcclient "tangled.org/core/xrpc/xrpcclient"
38
39 comatproto "github.com/bluesky-social/indigo/api/atproto"
40 "github.com/bluesky-social/indigo/atproto/atclient"
41 "github.com/bluesky-social/indigo/atproto/syntax"
42 lexutil "github.com/bluesky-social/indigo/lex/util"
43 indigoxrpc "github.com/bluesky-social/indigo/xrpc"
44 "github.com/go-chi/chi/v5"
45)
46
47type Repo struct {
48 repoResolver *reporesolver.RepoResolver
49 idResolver *idresolver.Resolver
50 config *config.Config
51 oauth *oauth.OAuth
52 pages *pages.Pages
53 db *db.DB
54 enforcer *rbac.Enforcer
55 acl *knotacl.Service
56 notifier notify.Notifier
57 logger *slog.Logger
58 serviceAuth *serviceauth.ServiceAuth
59 cfClient *cloudflare.Client
60 ogreClient *ogre.Client
61 codesearch *codesearch.CodeSearch
62
63 knotMirrorXRPC *indigoxrpc.Client
64 archiveClient *http.Client
65}
66
67func New(
68 oauth *oauth.OAuth,
69 repoResolver *reporesolver.RepoResolver,
70 pages *pages.Pages,
71 idResolver *idresolver.Resolver,
72 db *db.DB,
73 config *config.Config,
74 notifier notify.Notifier,
75 enforcer *rbac.Enforcer,
76 acl *knotacl.Service,
77 logger *slog.Logger,
78 cfClient *cloudflare.Client,
79 codesearch *codesearch.CodeSearch,
80) *Repo {
81 return &Repo{
82 oauth: oauth,
83 repoResolver: repoResolver,
84 pages: pages,
85 idResolver: idResolver,
86 config: config,
87 db: db,
88 notifier: notifier,
89 enforcer: enforcer,
90 acl: acl,
91 logger: logger,
92 cfClient: cfClient,
93 ogreClient: ogre.NewClient(config.Ogre.Host),
94 codesearch: codesearch,
95
96 knotMirrorXRPC: newKnotMirrorXRPCClient(config.KnotMirror.Url),
97 archiveClient: newArchiveClient(config.KnotMirror.ArchiveHeaderTimeout),
98 }
99}
100
101// modify the spindle configured for this repo
102func (rp *Repo) EditSpindle(w http.ResponseWriter, r *http.Request) {
103 user := rp.oauth.GetMultiAccountUser(r)
104 l := rp.logger.With("handler", "EditSpindle")
105 l = l.With("did", user.Did)
106
107 errorId := "operation-error"
108 fail := func(msg string, err error) {
109 l.Error(msg, "err", err)
110 rp.pages.Notice(w, errorId, msg)
111 }
112
113 f, err := rp.repoResolver.Resolve(r)
114 if err != nil {
115 fail("Failed to resolve repo. Try again later", err)
116 return
117 }
118
119 newSpindle := r.FormValue("spindle")
120 removingSpindle := newSpindle == "[[none]]" // see pages/templates/repo/settings/pipelines.html for more info on why we use this value
121 client, err := rp.oauth.AuthorizedClient(r)
122 if err != nil {
123 fail("Failed to authorize. Try again later.", err)
124 return
125 }
126
127 if !removingSpindle {
128 // ensure that this is a valid spindle for this user
129 validSpindles, err := rp.enforcer.GetSpindlesForUser(user.Did)
130 if err != nil {
131 fail("Failed to find spindles. Try again later.", err)
132 return
133 }
134
135 if !slices.Contains(validSpindles, newSpindle) {
136 fail("Failed to configure spindle.", fmt.Errorf("%s is not a valid spindle: %q", newSpindle, validSpindles))
137 return
138 }
139 }
140
141 newRepo := *f
142 newRepo.Spindle = newSpindle
143 record := newRepo.AsRecord()
144
145 spindlePtr := &newSpindle
146 if removingSpindle {
147 spindlePtr = nil
148 newRepo.Spindle = ""
149 }
150
151 // optimistic update
152 err = db.UpdateSpindle(rp.db, newRepo.RepoDid, spindlePtr)
153 if err != nil {
154 fail("Failed to update spindle. Try again later.", err)
155 return
156 }
157
158 ex, err := comatproto.RepoGetRecord(r.Context(), client, "", tangled.RepoNSID, newRepo.Did, newRepo.Rkey)
159 if err != nil {
160 fail("Failed to update spindle, no record found on PDS.", err)
161 return
162 }
163 _, err = comatproto.RepoPutRecord(r.Context(), client, &comatproto.RepoPutRecord_Input{
164 Collection: tangled.RepoNSID,
165 Repo: newRepo.Did,
166 Rkey: newRepo.Rkey,
167 SwapRecord: ex.Cid,
168 Record: &lexutil.LexiconTypeDecoder{
169 Val: &record,
170 },
171 })
172
173 if err != nil {
174 fail("Failed to update spindle, unable to save to PDS.", err)
175 return
176 }
177
178 rp.pages.HxRefresh(w)
179}
180
181func (rp *Repo) AddLabelDef(w http.ResponseWriter, r *http.Request) {
182 user := rp.oauth.GetMultiAccountUser(r)
183 l := rp.logger.With("handler", "AddLabel")
184 l = l.With("did", user.Did)
185
186 f, err := rp.repoResolver.Resolve(r)
187 if err != nil {
188 l.Error("failed to get repo and knot", "err", err)
189 return
190 }
191
192 errorId := "add-label-error"
193 fail := func(msg string, err error) {
194 l.Error(msg, "err", err)
195 rp.pages.Notice(w, errorId, msg)
196 }
197
198 // get form values for label definition
199 name := r.FormValue("name")
200 concreteType := r.FormValue("valueType")
201 valueFormat := r.FormValue("valueFormat")
202 enumValues := r.FormValue("enumValues")
203 scope := r.Form["scope"]
204 color := r.FormValue("color")
205 multiple := r.FormValue("multiple") == "true"
206
207 var variants []string
208 for part := range strings.SplitSeq(enumValues, ",") {
209 if part = strings.TrimSpace(part); part != "" {
210 variants = append(variants, part)
211 }
212 }
213
214 if concreteType == "" {
215 concreteType = "null"
216 }
217
218 format := models.ValueTypeFormatAny
219 if valueFormat == "did" {
220 format = models.ValueTypeFormatDid
221 }
222
223 valueType := models.ValueType{
224 Type: models.ConcreteType(concreteType),
225 Format: format,
226 Enum: variants,
227 }
228
229 label := models.LabelDefinition{
230 Did: user.Did,
231 Rkey: tid.TID(),
232 Name: name,
233 ValueType: valueType,
234 Scope: scope,
235 Color: &color,
236 Multiple: multiple,
237 Created: time.Now(),
238 }
239 if err := label.Validate(); err != nil {
240 fail(err.Error(), err)
241 return
242 }
243
244 // announce this relation into the firehose, store into owners' pds
245 client, err := rp.oauth.AuthorizedClient(r)
246 if err != nil {
247 fail(err.Error(), err)
248 return
249 }
250
251 // emit a labelRecord
252 labelRecord := label.AsRecord()
253 resp, err := comatproto.RepoPutRecord(r.Context(), client, &comatproto.RepoPutRecord_Input{
254 Collection: tangled.LabelDefinitionNSID,
255 Repo: label.Did,
256 Rkey: label.Rkey,
257 Record: &lexutil.LexiconTypeDecoder{
258 Val: &labelRecord,
259 },
260 })
261 // invalid record
262 if err != nil {
263 fail("Failed to write record to PDS.", err)
264 return
265 }
266
267 aturi := resp.Uri
268 l = l.With("at-uri", aturi)
269 l.Info("wrote label record to PDS")
270
271 // update the repo to subscribe to this label
272 newRepo := *f
273 newRepo.Labels = append(newRepo.Labels, aturi)
274 repoRecord := newRepo.AsRecord()
275
276 ex, err := comatproto.RepoGetRecord(r.Context(), client, "", tangled.RepoNSID, newRepo.Did, newRepo.Rkey)
277 if err != nil {
278 fail("Failed to update labels, no record found on PDS.", err)
279 return
280 }
281 _, err = comatproto.RepoPutRecord(r.Context(), client, &comatproto.RepoPutRecord_Input{
282 Collection: tangled.RepoNSID,
283 Repo: newRepo.Did,
284 Rkey: newRepo.Rkey,
285 SwapRecord: ex.Cid,
286 Record: &lexutil.LexiconTypeDecoder{
287 Val: &repoRecord,
288 },
289 })
290 if err != nil {
291 fail("Failed to update labels for repo.", err)
292 return
293 }
294
295 tx, err := rp.db.BeginTx(r.Context(), nil)
296 if err != nil {
297 fail("Failed to add label.", err)
298 return
299 }
300
301 rollback := func() {
302 err1 := tx.Rollback()
303 err2 := rollbackRecord(context.Background(), aturi, client)
304
305 // ignore txn complete errors, this is okay
306 if errors.Is(err1, sql.ErrTxDone) {
307 err1 = nil
308 }
309
310 if errs := errors.Join(err1, err2); errs != nil {
311 l.Error("failed to rollback changes", "errs", errs)
312 return
313 }
314 }
315 defer rollback()
316
317 _, err = db.AddLabelDefinition(tx, &label)
318 if err != nil {
319 fail("Failed to add label.", err)
320 return
321 }
322
323 if err = db.SubscribeLabel(tx, &models.RepoLabel{
324 RepoDid: syntax.DID(f.RepoDid),
325 LabelAt: label.AtUri(),
326 }); err != nil {
327 fail("Failed to subscribe to label.", err)
328 return
329 }
330
331 err = tx.Commit()
332 if err != nil {
333 fail("Failed to add label.", err)
334 return
335 }
336
337 // clear aturi when everything is successful
338 aturi = ""
339
340 rp.pages.HxRefresh(w)
341}
342
343func (rp *Repo) DeleteLabelDef(w http.ResponseWriter, r *http.Request) {
344 user := rp.oauth.GetMultiAccountUser(r)
345 l := rp.logger.With("handler", "DeleteLabel")
346 l = l.With("did", user.Did)
347
348 f, err := rp.repoResolver.Resolve(r)
349 if err != nil {
350 l.Error("failed to get repo and knot", "err", err)
351 return
352 }
353
354 errorId := "label-operation"
355 fail := func(msg string, err error) {
356 l.Error(msg, "err", err)
357 rp.pages.Notice(w, errorId, msg)
358 }
359
360 // get form values
361 labelId := r.FormValue("label-id")
362
363 label, err := db.GetLabelDefinition(rp.db, orm.FilterEq("id", labelId))
364 if err != nil {
365 fail("Failed to find label definition.", err)
366 return
367 }
368
369 client, err := rp.oauth.AuthorizedClient(r)
370 if err != nil {
371 fail(err.Error(), err)
372 return
373 }
374
375 // delete label record from PDS
376 _, err = comatproto.RepoDeleteRecord(r.Context(), client, &comatproto.RepoDeleteRecord_Input{
377 Collection: tangled.LabelDefinitionNSID,
378 Repo: label.Did,
379 Rkey: label.Rkey,
380 })
381 if err != nil {
382 fail("Failed to delete label record from PDS.", err)
383 return
384 }
385
386 // update repo record to remove the label reference
387 newRepo := *f
388 var updated []string
389 removedAt := label.AtUri().String()
390 for _, l := range newRepo.Labels {
391 if l != removedAt {
392 updated = append(updated, l)
393 }
394 }
395 newRepo.Labels = updated
396 repoRecord := newRepo.AsRecord()
397
398 ex, err := comatproto.RepoGetRecord(r.Context(), client, "", tangled.RepoNSID, newRepo.Did, newRepo.Rkey)
399 if err != nil {
400 fail("Failed to update labels, no record found on PDS.", err)
401 return
402 }
403 _, err = comatproto.RepoPutRecord(r.Context(), client, &comatproto.RepoPutRecord_Input{
404 Collection: tangled.RepoNSID,
405 Repo: newRepo.Did,
406 Rkey: newRepo.Rkey,
407 SwapRecord: ex.Cid,
408 Record: &lexutil.LexiconTypeDecoder{
409 Val: &repoRecord,
410 },
411 })
412 if err != nil {
413 fail("Failed to update repo record.", err)
414 return
415 }
416
417 // transaction for DB changes
418 tx, err := rp.db.BeginTx(r.Context(), nil)
419 if err != nil {
420 fail("Failed to delete label.", err)
421 return
422 }
423 defer tx.Rollback()
424
425 err = db.UnsubscribeLabel(
426 tx,
427 orm.FilterEq("repo_did", f.RepoDid),
428 orm.FilterEq("label_at", removedAt),
429 )
430 if err != nil {
431 fail("Failed to unsubscribe label.", err)
432 return
433 }
434
435 err = db.DeleteLabelDefinition(tx, orm.FilterEq("id", label.Id))
436 if err != nil {
437 fail("Failed to delete label definition.", err)
438 return
439 }
440
441 err = tx.Commit()
442 if err != nil {
443 fail("Failed to delete label.", err)
444 return
445 }
446
447 // everything succeeded
448 rp.pages.HxRefresh(w)
449}
450
451func (rp *Repo) SubscribeLabel(w http.ResponseWriter, r *http.Request) {
452 user := rp.oauth.GetMultiAccountUser(r)
453 l := rp.logger.With("handler", "SubscribeLabel")
454 l = l.With("did", user.Did)
455
456 f, err := rp.repoResolver.Resolve(r)
457 if err != nil {
458 l.Error("failed to get repo and knot", "err", err)
459 return
460 }
461
462 if err := r.ParseForm(); err != nil {
463 l.Error("invalid form", "err", err)
464 return
465 }
466
467 errorId := "default-label-operation"
468 fail := func(msg string, err error) {
469 l.Error(msg, "err", err)
470 rp.pages.Notice(w, errorId, msg)
471 }
472
473 labelAts := r.Form["label"]
474 _, err = db.GetLabelDefinitions(rp.db, orm.FilterIn("at_uri", labelAts))
475 if err != nil {
476 fail("Failed to subscribe to label.", err)
477 return
478 }
479
480 newRepo := *f
481 newRepo.Labels = append(newRepo.Labels, labelAts...)
482
483 // dedup
484 slices.Sort(newRepo.Labels)
485 newRepo.Labels = slices.Compact(newRepo.Labels)
486
487 repoRecord := newRepo.AsRecord()
488
489 client, err := rp.oauth.AuthorizedClient(r)
490 if err != nil {
491 fail(err.Error(), err)
492 return
493 }
494
495 ex, err := comatproto.RepoGetRecord(r.Context(), client, "", tangled.RepoNSID, f.Did, f.Rkey)
496 if err != nil {
497 fail("Failed to update labels, no record found on PDS.", err)
498 return
499 }
500 _, err = comatproto.RepoPutRecord(r.Context(), client, &comatproto.RepoPutRecord_Input{
501 Collection: tangled.RepoNSID,
502 Repo: newRepo.Did,
503 Rkey: newRepo.Rkey,
504 SwapRecord: ex.Cid,
505 Record: &lexutil.LexiconTypeDecoder{
506 Val: &repoRecord,
507 },
508 })
509
510 tx, err := rp.db.Begin()
511 if err != nil {
512 fail("Failed to subscribe to label.", err)
513 return
514 }
515 defer tx.Rollback()
516
517 for _, l := range labelAts {
518 err = db.SubscribeLabel(tx, &models.RepoLabel{
519 RepoDid: syntax.DID(f.RepoDid),
520 LabelAt: syntax.ATURI(l),
521 })
522 if err != nil {
523 fail("Failed to subscribe to label.", err)
524 return
525 }
526 }
527
528 if err := tx.Commit(); err != nil {
529 fail("Failed to subscribe to label.", err)
530 return
531 }
532
533 // everything succeeded
534 rp.pages.HxRefresh(w)
535}
536
537func (rp *Repo) UnsubscribeLabel(w http.ResponseWriter, r *http.Request) {
538 user := rp.oauth.GetMultiAccountUser(r)
539 l := rp.logger.With("handler", "UnsubscribeLabel")
540 l = l.With("did", user.Did)
541
542 f, err := rp.repoResolver.Resolve(r)
543 if err != nil {
544 l.Error("failed to get repo and knot", "err", err)
545 return
546 }
547
548 if err := r.ParseForm(); err != nil {
549 l.Error("invalid form", "err", err)
550 return
551 }
552
553 errorId := "default-label-operation"
554 fail := func(msg string, err error) {
555 l.Error(msg, "err", err)
556 rp.pages.Notice(w, errorId, msg)
557 }
558
559 labelAts := r.Form["label"]
560 _, err = db.GetLabelDefinitions(rp.db, orm.FilterIn("at_uri", labelAts))
561 if err != nil {
562 fail("Failed to unsubscribe to label.", err)
563 return
564 }
565
566 // update repo record to remove the label reference
567 newRepo := *f
568 var updated []string
569 for _, l := range newRepo.Labels {
570 if !slices.Contains(labelAts, l) {
571 updated = append(updated, l)
572 }
573 }
574 newRepo.Labels = updated
575 repoRecord := newRepo.AsRecord()
576
577 client, err := rp.oauth.AuthorizedClient(r)
578 if err != nil {
579 fail(err.Error(), err)
580 return
581 }
582
583 ex, err := comatproto.RepoGetRecord(r.Context(), client, "", tangled.RepoNSID, f.Did, f.Rkey)
584 if err != nil {
585 fail("Failed to update labels, no record found on PDS.", err)
586 return
587 }
588 _, err = comatproto.RepoPutRecord(r.Context(), client, &comatproto.RepoPutRecord_Input{
589 Collection: tangled.RepoNSID,
590 Repo: newRepo.Did,
591 Rkey: newRepo.Rkey,
592 SwapRecord: ex.Cid,
593 Record: &lexutil.LexiconTypeDecoder{
594 Val: &repoRecord,
595 },
596 })
597
598 err = db.UnsubscribeLabel(
599 rp.db,
600 orm.FilterEq("repo_did", f.RepoDid),
601 orm.FilterIn("label_at", labelAts),
602 )
603 if err != nil {
604 fail("Failed to unsubscribe label.", err)
605 return
606 }
607
608 // everything succeeded
609 rp.pages.HxRefresh(w)
610}
611
612func (rp *Repo) LabelPanel(w http.ResponseWriter, r *http.Request) {
613 l := rp.logger.With("handler", "LabelPanel")
614
615 f, err := rp.repoResolver.Resolve(r)
616 if err != nil {
617 l.Error("failed to get repo and knot", "err", err)
618 return
619 }
620
621 subjectStr := r.FormValue("subject")
622 subject, err := syntax.ParseATURI(subjectStr)
623 if err != nil {
624 l.Error("failed to get repo and knot", "err", err)
625 return
626 }
627
628 labelDefs, err := db.GetLabelDefinitions(
629 rp.db,
630 orm.FilterIn("at_uri", f.Labels),
631 orm.FilterContains("scope", subject.Collection().String()),
632 )
633 if err != nil {
634 l.Error("failed to fetch label defs", "err", err)
635 return
636 }
637
638 defs := make(map[string]*models.LabelDefinition)
639 for _, l := range labelDefs {
640 defs[l.AtUri().String()] = &l
641 }
642
643 states, err := db.GetLabels(rp.db, orm.FilterEq("subject", subject))
644 if err != nil {
645 l.Error("failed to build label state", "err", err)
646 return
647 }
648 state := states[subject]
649
650 user := rp.oauth.GetMultiAccountUser(r)
651 rp.pages.LabelPanel(w, pages.LabelPanelParams{
652 BaseParams: pages.BaseParamsFromContext(r.Context()),
653 RepoInfo: rp.repoResolver.GetRepoInfo(r, user),
654 Defs: defs,
655 Subject: subject.String(),
656 State: state,
657 })
658}
659
660func (rp *Repo) EditLabelPanel(w http.ResponseWriter, r *http.Request) {
661 l := rp.logger.With("handler", "EditLabelPanel")
662
663 f, err := rp.repoResolver.Resolve(r)
664 if err != nil {
665 l.Error("failed to get repo and knot", "err", err)
666 return
667 }
668
669 subjectStr := r.FormValue("subject")
670 subject, err := syntax.ParseATURI(subjectStr)
671 if err != nil {
672 l.Error("failed to get repo and knot", "err", err)
673 return
674 }
675
676 labelDefs, err := db.GetLabelDefinitions(
677 rp.db,
678 orm.FilterIn("at_uri", f.Labels),
679 orm.FilterContains("scope", subject.Collection().String()),
680 )
681 if err != nil {
682 l.Error("failed to fetch labels", "err", err)
683 return
684 }
685
686 defs := make(map[string]*models.LabelDefinition)
687 for _, l := range labelDefs {
688 defs[l.AtUri().String()] = &l
689 }
690
691 states, err := db.GetLabels(rp.db, orm.FilterEq("subject", subject))
692 if err != nil {
693 l.Error("failed to build label state", "err", err)
694 return
695 }
696 state := states[subject]
697
698 user := rp.oauth.GetMultiAccountUser(r)
699 rp.pages.EditLabelPanel(w, pages.EditLabelPanelParams{
700 BaseParams: pages.BaseParamsFromContext(r.Context()),
701 RepoInfo: rp.repoResolver.GetRepoInfo(r, user),
702 Defs: defs,
703 Subject: subject.String(),
704 State: state,
705 })
706}
707
708func (rp *Repo) AddCollaborator(w http.ResponseWriter, r *http.Request) {
709 user := rp.oauth.GetMultiAccountUser(r)
710 l := rp.logger.With("handler", "AddCollaborator")
711 l = l.With("did", user.Did)
712
713 f, err := rp.repoResolver.Resolve(r)
714 if err != nil {
715 l.Error("failed to get repo and knot", "err", err)
716 return
717 }
718
719 errorId := "add-collaborator-error"
720 fail := func(msg string, err error) {
721 l.Error(msg, "err", err)
722 rp.pages.Notice(w, errorId, msg)
723 }
724
725 collaborator := r.FormValue("collaborator")
726 if collaborator == "" {
727 fail("Invalid form.", nil)
728 return
729 }
730
731 // remove a single leading `@`, to make @handle work with ResolveIdent
732 collaborator = strings.TrimPrefix(collaborator, "@")
733
734 collaboratorIdent, err := rp.idResolver.ResolveIdent(r.Context(), collaborator)
735 if err != nil {
736 fail(fmt.Sprintf("'%s' is not a valid DID/handle.", collaborator), err)
737 return
738 }
739
740 if collaboratorIdent.DID.String() == user.Did {
741 fail("You seem to be adding yourself as a collaborator.", nil)
742 return
743 }
744 l = l.With("collaborator", collaboratorIdent.Handle)
745 l = l.With("knot", f.Knot)
746
747 capStatus := knotcompat.KnotCapability(r.Context(), f.Knot, rp.config.Core.Dev, consts.CapKnotACL)
748 if capStatus == knotcompat.CapUnknown {
749 fail("Could not reach the knot to add the collaborator. Try again later.", nil)
750 return
751 }
752 if capStatus == knotcompat.CapPresent {
753 if f.RepoDid == "" {
754 fail("This repository is missing its DID and cannot manage collaborators.", nil)
755 return
756 }
757
758 client, err := rp.oauth.ServiceClient(
759 r,
760 oauth.WithService(f.Knot),
761 oauth.WithLxm(tangled.RepoAddCollaboratorNSID),
762 oauth.WithDev(rp.config.Core.Dev),
763 )
764 if err != nil {
765 fail("Failed to connect to knot server.", err)
766 return
767 }
768
769 err = tangled.RepoAddCollaborator(r.Context(), client, &tangled.RepoAddCollaborator_Input{
770 Repo: f.RepoDid,
771 Subject: collaboratorIdent.DID.String(),
772 })
773 if xrpcerr := xrpcclient.HandleXrpcErr(err); xrpcerr != nil {
774 l.Error("failed to call XRPC repo.addCollaborator", "xrpcerr", xrpcerr, "err", err)
775 rp.pages.Notice(w, errorId, xrpcerr.Error())
776 return
777 }
778
779 rp.acl.InvalidateCollaborators(f.Knot, f.RepoDid)
780
781 rp.pages.HxRefresh(w)
782 return
783 }
784
785 existing, err := db.GetCollaborators(rp.db,
786 orm.FilterEq("repo_did", f.RepoDid),
787 orm.FilterEq("subject_did", collaboratorIdent.DID.String()),
788 )
789 if err != nil {
790 fail("Failed to check existing collaborators.", err)
791 return
792 }
793 if len(existing) > 0 {
794 fail(fmt.Sprintf("%s is already a collaborator.", collaboratorIdent.Handle), nil)
795 return
796 }
797
798 // announce this relation into the firehose, store into owners' pds
799 client, err := rp.oauth.AuthorizedClient(r)
800 if err != nil {
801 fail("Failed to write to PDS.", err)
802 return
803 }
804
805 // emit a record
806 currentUser := rp.oauth.GetMultiAccountUser(r)
807 rkey := tid.TID()
808 createdAt := time.Now()
809 resp, err := comatproto.RepoPutRecord(r.Context(), client, &comatproto.RepoPutRecord_Input{
810 Collection: tangled.RepoCollaboratorNSID,
811 Repo: currentUser.Did,
812 Rkey: rkey,
813 Record: knotcompat.Collaborator(repoCollaboratorRecord(f, collaboratorIdent.DID.String(), createdAt)),
814 })
815 // invalid record
816 if err != nil {
817 fail("Failed to write record to PDS.", err)
818 return
819 }
820
821 aturi := resp.Uri
822 l = l.With("at-uri", aturi)
823 l.Info("wrote record to PDS")
824
825 tx, err := rp.db.BeginTx(r.Context(), nil)
826 if err != nil {
827 fail("Failed to add collaborator.", err)
828 return
829 }
830
831 rollback := func() {
832 err1 := tx.Rollback()
833 err2 := rp.enforcer.E.LoadPolicy()
834 err3 := rollbackRecord(context.Background(), aturi, client)
835
836 // ignore txn complete errors, this is okay
837 if errors.Is(err1, sql.ErrTxDone) {
838 err1 = nil
839 }
840
841 if errs := errors.Join(err1, err2, err3); errs != nil {
842 l.Error("failed to rollback changes", "errs", errs)
843 return
844 }
845 }
846 defer rollback()
847
848 err = rp.enforcer.AddCollaborator(collaboratorIdent.DID.String(), f.Knot, f.RepoIdentifier())
849 if err != nil {
850 fail("Failed to add collaborator permissions.", err)
851 return
852 }
853
854 err = db.AddCollaborator(tx, models.Collaborator{
855 Did: syntax.DID(currentUser.Did),
856 Rkey: sql.NullString{String: rkey, Valid: true},
857 SubjectDid: collaboratorIdent.DID,
858 RepoDid: syntax.DID(f.RepoDid),
859 Created: createdAt,
860 })
861 if err != nil {
862 fail("Failed to add collaborator.", err)
863 return
864 }
865
866 err = tx.Commit()
867 if err != nil {
868 fail("Failed to add collaborator.", err)
869 return
870 }
871
872 err = rp.enforcer.E.SavePolicy()
873 if err != nil {
874 fail("Failed to update collaborator permissions.", err)
875 return
876 }
877
878 // clear aturi to when everything is successful
879 aturi = ""
880
881 rp.pages.HxRefresh(w)
882}
883
884func (rp *Repo) RemoveCollaborator(w http.ResponseWriter, r *http.Request) {
885 user := rp.oauth.GetMultiAccountUser(r)
886 l := rp.logger.With("handler", "RemoveCollaborator")
887 l = l.With("did", user.Did)
888
889 f, err := rp.repoResolver.Resolve(r)
890 if err != nil {
891 l.Error("failed to get repo and knot", "err", err)
892 return
893 }
894
895 errorId := "collaborator-error"
896 fail := func(msg string, err error) {
897 l.Error(msg, "err", err)
898 rp.pages.Notice(w, errorId, msg)
899 }
900
901 collaborator := r.FormValue("collaborator")
902 if collaborator == "" {
903 fail("Invalid form.", nil)
904 return
905 }
906 collaborator = strings.TrimPrefix(collaborator, "@")
907
908 collaboratorIdent, err := rp.idResolver.ResolveIdent(r.Context(), collaborator)
909 if err != nil {
910 fail(fmt.Sprintf("'%s' is not a valid DID/handle.", collaborator), err)
911 return
912 }
913 l = l.With("collaborator", collaboratorIdent.Handle, "knot", f.Knot)
914
915 if collaboratorIdent.DID.String() == f.Did {
916 fail("Cannot remove the repository owner.", nil)
917 return
918 }
919
920 capStatus := knotcompat.KnotCapability(r.Context(), f.Knot, rp.config.Core.Dev, consts.CapKnotACL)
921 if capStatus == knotcompat.CapUnknown {
922 fail("Could not reach the knot to remove the collaborator. Try again later.", nil)
923 return
924 }
925 if capStatus == knotcompat.CapPresent {
926 if f.RepoDid == "" {
927 fail("This repository is missing its DID and cannot manage collaborators.", nil)
928 return
929 }
930
931 client, err := rp.oauth.ServiceClient(
932 r,
933 oauth.WithService(f.Knot),
934 oauth.WithLxm(tangled.RepoRemoveCollaboratorNSID),
935 oauth.WithDev(rp.config.Core.Dev),
936 )
937 if err != nil {
938 fail("Failed to connect to knot server.", err)
939 return
940 }
941
942 err = tangled.RepoRemoveCollaborator(r.Context(), client, &tangled.RepoRemoveCollaborator_Input{
943 Repo: f.RepoDid,
944 Subject: collaboratorIdent.DID.String(),
945 })
946 if xrpcerr := xrpcclient.HandleXrpcErr(err); xrpcerr != nil {
947 l.Error("failed to call XRPC repo.removeCollaborator", "xrpcerr", xrpcerr, "err", err)
948 rp.pages.Notice(w, errorId, xrpcerr.Error())
949 return
950 }
951
952 rp.acl.InvalidateCollaborators(f.Knot, f.RepoDid)
953
954 rp.pages.HxRefresh(w)
955 return
956 }
957
958 existing, err := db.GetCollaborators(rp.db,
959 orm.FilterEq("repo_did", f.RepoDid),
960 orm.FilterEq("subject_did", collaboratorIdent.DID.String()),
961 )
962 if err != nil {
963 fail("Failed to look up collaborator.", err)
964 return
965 }
966 if len(existing) == 0 {
967 fail(fmt.Sprintf("%s is not a collaborator.", collaboratorIdent.Handle), nil)
968 return
969 }
970 row := existing[0]
971
972 client, err := rp.oauth.AuthorizedClient(r)
973 if err != nil {
974 fail("Failed to write to PDS.", err)
975 return
976 }
977
978 tx, err := rp.db.BeginTx(r.Context(), nil)
979 if err != nil {
980 fail("Failed to remove collaborator.", err)
981 return
982 }
983 committed := false
984 defer func() {
985 if !committed {
986 tx.Rollback()
987 if err := rp.enforcer.E.LoadPolicy(); err != nil {
988 l.Error("failed to reload policy after rollback", "err", err)
989 }
990 }
991 }()
992
993 if err := rp.enforcer.RemoveCollaborator(collaboratorIdent.DID.String(), f.Knot, f.RepoIdentifier()); err != nil {
994 fail("Failed to remove collaborator permissions.", err)
995 return
996 }
997
998 if err := db.DeleteCollaborator(tx,
999 orm.FilterEq("repo_did", f.RepoDid),
1000 orm.FilterEq("subject_did", collaboratorIdent.DID.String()),
1001 ); err != nil {
1002 fail("Failed to remove collaborator.", err)
1003 return
1004 }
1005
1006 if row.Rkey.Valid && row.Rkey.String != "" {
1007 if _, err := comatproto.RepoDeleteRecord(r.Context(), client, &comatproto.RepoDeleteRecord_Input{
1008 Collection: tangled.RepoCollaboratorNSID,
1009 Repo: row.Did.String(),
1010 Rkey: row.Rkey.String,
1011 }); err != nil {
1012 fail("Failed to delete collaborator record from PDS.", err)
1013 return
1014 }
1015 }
1016
1017 if err := tx.Commit(); err != nil {
1018 fail("Failed to remove collaborator.", err)
1019 return
1020 }
1021 committed = true
1022
1023 if err := rp.enforcer.E.SavePolicy(); err != nil {
1024 fail("Failed to update collaborator permissions.", err)
1025 return
1026 }
1027
1028 rp.pages.HxRefresh(w)
1029}
1030
1031func (rp *Repo) RenameRepo(w http.ResponseWriter, r *http.Request) {
1032 l := rp.logger.With("handler", "RenameRepo")
1033 noticeId := "rename-repo-error"
1034
1035 user := rp.oauth.GetMultiAccountUser(r)
1036 f, err := rp.repoResolver.Resolve(r)
1037 if err != nil {
1038 l.Error("failed to get repo and knot", "err", err)
1039 rp.pages.Notice(w, noticeId, "Failed to load repository.")
1040 return
1041 }
1042 l = l.With("did", user.Did, "rkey", f.Rkey, "oldName", f.Name)
1043
1044 if f.RepoDid == "" {
1045 rp.pages.Notice(w, noticeId, "This repository's knot has not completed the DID migration; rename is unavailable.")
1046 return
1047 }
1048
1049 if !knotcompat.KnotSupports114(r.Context(), f.Knot, rp.config.Core.Dev) {
1050 rp.pages.Notice(w, noticeId, "This repository's knot is below v1.14 and does not yet support renames. Ask the knot operator to upgrade.")
1051 return
1052 }
1053
1054 newName, err := validateRenameInput(f.Name, f.Rkey, r.FormValue("name"))
1055 if err != nil {
1056 rp.pages.Notice(w, noticeId, err.Error())
1057 return
1058 }
1059 newRkey := strings.ToLower(newName)
1060 l = l.With("newName", newName, "newRkey", newRkey)
1061
1062 atpClient, err := rp.oauth.AuthorizedClient(r)
1063 if err != nil {
1064 l.Error("failed to get authorized client", "err", err)
1065 rp.pages.Notice(w, noticeId, "Failed to authorize. Try again later.")
1066 return
1067 }
1068
1069 newRepo := *f
1070 newRepo.Name = newName
1071 newRepo.Rkey = newRkey
1072 newRepo.Created = time.Now()
1073 record := newRepo.AsRecord()
1074
1075 if newRkey == f.Rkey {
1076 ex, err := comatproto.RepoGetRecord(r.Context(), atpClient, "", tangled.RepoNSID, f.Did, f.Rkey)
1077 if err != nil {
1078 l.Error("failed to fetch existing record", "err", err)
1079 rp.pages.Notice(w, noticeId, "Failed to read repository record from PDS.")
1080 return
1081 }
1082
1083 _, err = comatproto.RepoPutRecord(r.Context(), atpClient, &comatproto.RepoPutRecord_Input{
1084 Collection: tangled.RepoNSID,
1085 Repo: f.Did,
1086 Rkey: f.Rkey,
1087 SwapRecord: ex.Cid,
1088 Record: &lexutil.LexiconTypeDecoder{
1089 Val: &record,
1090 },
1091 })
1092 if err != nil {
1093 l.Error("failed to update display name on PDS", "err", err)
1094 rp.pages.Notice(w, noticeId, "Failed to save display name to PDS.")
1095 return
1096 }
1097 l.Info("updated display name on PDS")
1098
1099 if err := db.UpdateRepoDisplayName(rp.db, f.Did, f.Rkey, newName); err != nil {
1100 l.Error("optimistic display name update failed", "err", err)
1101 }
1102 } else {
1103 ex, getErr := comatproto.RepoGetRecord(r.Context(), atpClient, "", tangled.RepoNSID, f.Did, newRkey)
1104 switch {
1105 case getErr != nil:
1106 _, err = comatproto.RepoCreateRecord(r.Context(), atpClient, &comatproto.RepoCreateRecord_Input{
1107 Collection: tangled.RepoNSID,
1108 Repo: f.Did,
1109 Rkey: &newRkey,
1110 Record: &lexutil.LexiconTypeDecoder{Val: &record},
1111 })
1112 if err != nil {
1113 l.Error("failed to write rename to PDS", "err", err)
1114 rp.pages.Notice(w, noticeId, "Failed to save renamed repository to PDS.")
1115 return
1116 }
1117 l.Info("wrote rename-create to PDS; old record retained as alias")
1118
1119 default:
1120 existing, ok := ex.Value.Val.(*tangled.Repo)
1121 if !ok || existing.RepoDid == nil || *existing.RepoDid != f.RepoDid {
1122 rp.pages.Notice(w, noticeId, fmt.Sprintf("You already have a repository named %q.", newRkey))
1123 return
1124 }
1125 _, err = comatproto.RepoPutRecord(r.Context(), atpClient, &comatproto.RepoPutRecord_Input{
1126 Collection: tangled.RepoNSID,
1127 Repo: f.Did,
1128 Rkey: newRkey,
1129 SwapRecord: ex.Cid,
1130 Record: &lexutil.LexiconTypeDecoder{Val: &record},
1131 })
1132 if err != nil {
1133 l.Error("failed to rewrite rename-back record on PDS", "err", err)
1134 rp.pages.Notice(w, noticeId, "Failed to save renamed repository to PDS.")
1135 return
1136 }
1137 l.Info("rewrote rename-back record on PDS over prior alias")
1138 }
1139
1140 tx, err := rp.db.Begin()
1141 if err != nil {
1142 l.Error("failed to begin rename tx", "err", err)
1143 rp.pages.HxLocation(w, fmt.Sprintf("/%s", f.RepoDid))
1144 return
1145 }
1146 defer tx.Rollback()
1147
1148 if err := db.RenameRepo(tx, f.Did, f.Rkey, newRkey, newName); err != nil {
1149 l.Error("optimistic rename failed", "err", err)
1150 rp.pages.HxLocation(w, fmt.Sprintf("/%s", f.RepoDid))
1151 return
1152 }
1153 if err := db.RecordRepoRename(tx, f.Did, f.Rkey, f.RepoDid); err != nil {
1154 l.Error("failed to record rename history", "err", err)
1155 }
1156 if err := db.DeleteRepoRename(tx, f.Did, newRkey); err != nil {
1157 l.Error("failed to clear stale rename hint", "err", err)
1158 }
1159 if err := tx.Commit(); err != nil {
1160 l.Error("failed to commit rename tx", "err", err)
1161 rp.pages.HxLocation(w, fmt.Sprintf("/%s", f.RepoDid))
1162 return
1163 }
1164 }
1165
1166 oldRepo := *f
1167 rp.notifier.RenameRepo(r.Context(), syntax.DID(user.Did), &oldRepo, &newRepo)
1168
1169 if newRkey != f.Rkey {
1170 rp.migrateSiteOnRename(r.Context(), f, newName, newRkey)
1171 }
1172
1173 rp.pages.HxLocation(w, fmt.Sprintf("/%s", f.RepoDid))
1174}
1175
1176func validateRenameInput(currentName, currentRkey, raw string) (string, error) {
1177 newName := strings.TrimSpace(raw)
1178 if newName == "" {
1179 return "", errors.New("Repository name cannot be empty.")
1180 }
1181 if err := models.ValidateRepoName(newName); err != nil {
1182 return "", err
1183 }
1184 newName = models.StripGitExt(newName)
1185 if newName == currentName {
1186 if _, tidErr := syntax.ParseTID(currentRkey); tidErr == nil {
1187 return newName, nil
1188 }
1189 return "", errors.New("New name matches the current name.")
1190 }
1191 return newName, nil
1192}
1193
1194func (rp *Repo) migrateSiteOnRename(ctx context.Context, oldRepo *models.Repo, newName, newRkey string) {
1195 l := rp.logger.With("handler", "migrateSiteOnRename", "repo_did", oldRepo.RepoDid)
1196
1197 siteConfig, err := db.GetRepoSiteConfig(rp.db, oldRepo.RepoDid)
1198 if err != nil || siteConfig == nil {
1199 return
1200 }
1201
1202 if !rp.cfClient.Enabled() {
1203 return
1204 }
1205
1206 ownerClaim, _ := db.GetActiveDomainClaimForDid(rp.db, oldRepo.Did)
1207
1208 go func() {
1209 bgCtx := context.Background()
1210 oldRkey := oldRepo.Rkey
1211 oldName := oldRepo.Name
1212
1213 if err := sites.Delete(bgCtx, rp.cfClient, oldRepo.Did, oldRkey); err != nil {
1214 l.Error("sites: failed to delete old R2 prefix", "oldRkey", oldRkey, "err", err)
1215 }
1216
1217 newRepo := *oldRepo
1218 newRepo.Name = newName
1219 newRepo.Rkey = newRkey
1220 if deployErr := sites.Deploy(bgCtx, rp.cfClient, rp.config, &newRepo, siteConfig.Branch, siteConfig.Dir); deployErr != nil {
1221 l.Error("sites: redeploy after rename failed", "err", deployErr)
1222 }
1223
1224 if ownerClaim != nil {
1225 // drop the old name's entry when the name actually changed.
1226 if oldName != newName {
1227 if err := sites.DeleteDomainMapping(bgCtx, rp.cfClient, ownerClaim.Domain, oldName); err != nil {
1228 l.Error("sites: failed to remove old KV mapping", "oldName", oldName, "err", err)
1229 }
1230 }
1231 if err := sites.PutDomainMapping(bgCtx, rp.cfClient, ownerClaim.Domain, oldRepo.Did, newName, newRkey, siteConfig.IsIndex); err != nil {
1232 l.Error("sites: failed to write new KV mapping", "newName", newName, "newRkey", newRkey, "err", err)
1233 }
1234 }
1235
1236 l.Info("sites: migrated on rename", "oldName", oldName, "oldRkey", oldRkey, "newName", newName, "newRkey", newRkey)
1237 }()
1238}
1239
1240func (rp *Repo) DeleteRepo(w http.ResponseWriter, r *http.Request) {
1241 user := rp.oauth.GetMultiAccountUser(r)
1242 l := rp.logger.With("handler", "DeleteRepo")
1243
1244 noticeId := "operation-error"
1245 f, err := rp.repoResolver.Resolve(r)
1246 if err != nil {
1247 l.Error("failed to get repo and knot", "err", err)
1248 return
1249 }
1250
1251 // remove record from pds
1252 atpClient, err := rp.oauth.AuthorizedClient(r)
1253 if err != nil {
1254 l.Error("failed to get authorized client", "err", err)
1255 return
1256 }
1257 _, err = comatproto.RepoDeleteRecord(r.Context(), atpClient, &comatproto.RepoDeleteRecord_Input{
1258 Collection: tangled.RepoNSID,
1259 Repo: user.Did,
1260 Rkey: f.Rkey,
1261 })
1262 if err != nil {
1263 l.Error("failed to delete record", "err", err)
1264 rp.pages.Notice(w, noticeId, "Failed to delete repository from PDS.")
1265 return
1266 }
1267 l.Info("removed repo record", "aturi", f.RepoAt().String())
1268
1269 client, err := rp.oauth.ServiceClient(
1270 r,
1271 oauth.WithService(f.Knot),
1272 oauth.WithLxm(tangled.RepoDeleteNSID),
1273 oauth.WithDev(rp.config.Core.Dev),
1274 )
1275 if err != nil {
1276 l.Error("failed to connect to knot server", "err", err)
1277 return
1278 }
1279
1280 err = tangled.RepoDelete(
1281 r.Context(),
1282 client,
1283 &tangled.RepoDelete_Input{
1284 Did: f.Did,
1285 Name: f.Name,
1286 Rkey: f.Rkey,
1287 },
1288 )
1289 if xrpcerr := xrpcclient.HandleXrpcErr(err); xrpcerr != nil {
1290 l.Error("failed to call XRPC repo.delete", "xrpcerr", xrpcerr, "err", err)
1291 rp.pages.Notice(w, noticeId, xrpcerr.Error())
1292 return
1293 }
1294 l.Info("deleted repo from knot")
1295
1296 tx, err := rp.db.BeginTx(r.Context(), nil)
1297 if err != nil {
1298 l.Error("failed to start tx")
1299 w.Write(fmt.Append(nil, "failed to add collaborator: ", err))
1300 return
1301 }
1302 defer func() {
1303 tx.Rollback()
1304 err = rp.enforcer.E.LoadPolicy()
1305 if err != nil {
1306 l.Error("failed to rollback policies")
1307 }
1308 }()
1309
1310 // remove collaborator RBAC
1311 repoCollaborators, err := rp.enforcer.E.GetImplicitUsersForResourceByDomain(f.RepoIdentifier(), f.Knot)
1312 if err != nil {
1313 rp.pages.Notice(w, noticeId, "Failed to remove collaborators")
1314 return
1315 }
1316 for _, c := range repoCollaborators {
1317 did := c[0]
1318 rp.enforcer.RemoveCollaborator(did, f.Knot, f.RepoIdentifier())
1319 }
1320 l.Info("removed collaborators")
1321
1322 // remove repo RBAC
1323 err = rp.enforcer.RemoveRepo(f.Did, f.Knot, f.RepoIdentifier())
1324 if err != nil {
1325 rp.pages.Notice(w, noticeId, "Failed to update RBAC rules")
1326 return
1327 }
1328
1329 // remove repo from db
1330 err = db.RemoveRepo(tx, f.Did, f.Rkey)
1331 if err != nil {
1332 rp.pages.Notice(w, noticeId, "Failed to update appview")
1333 return
1334 }
1335 l.Info("removed repo from db")
1336
1337 err = tx.Commit()
1338 if err != nil {
1339 l.Error("failed to commit changes", "err", err)
1340 http.Error(w, err.Error(), http.StatusInternalServerError)
1341 return
1342 }
1343
1344 err = rp.enforcer.E.SavePolicy()
1345 if err != nil {
1346 l.Error("failed to update ACLs", "err", err)
1347 http.Error(w, err.Error(), http.StatusInternalServerError)
1348 return
1349 }
1350
1351 rp.notifier.DeleteRepo(r.Context(), f)
1352 rp.pages.HxRedirect(w, fmt.Sprintf("/%s", f.Did))
1353}
1354
1355func (rp *Repo) SyncRepoFork(w http.ResponseWriter, r *http.Request) {
1356 l := rp.logger.With("handler", "SyncRepoFork")
1357
1358 ref := chi.URLParam(r, "ref")
1359 ref, _ = url.PathUnescape(ref)
1360
1361 user := rp.oauth.GetMultiAccountUser(r)
1362 f, err := rp.repoResolver.Resolve(r)
1363 if err != nil {
1364 l.Error("failed to resolve source repo", "err", err)
1365 return
1366 }
1367
1368 switch r.Method {
1369 case http.MethodPost:
1370 client, err := rp.oauth.ServiceClient(
1371 r,
1372 oauth.WithService(f.Knot),
1373 oauth.WithLxm(tangled.RepoForkSyncNSID),
1374 oauth.WithDev(rp.config.Core.Dev),
1375 )
1376 if err != nil {
1377 rp.pages.Notice(w, "repo", "Failed to connect to knot server.")
1378 return
1379 }
1380
1381 if f.Source == "" {
1382 rp.pages.Notice(w, "repo", "This repository is not a fork.")
1383 return
1384 }
1385
1386 err = tangled.RepoForkSync(
1387 r.Context(),
1388 client,
1389 &tangled.RepoForkSync_Input{
1390 Did: user.Did,
1391 Name: f.Name,
1392 Repo: f.RepoDidPtr(),
1393 Source: f.Source,
1394 Branch: ref,
1395 },
1396 )
1397 if xrpcerr := xrpcclient.HandleXrpcErr(err); xrpcerr != nil {
1398 l.Error("failed to call XRPC repo.forkSync", "xrpcerr", xrpcerr, "err", err)
1399 rp.pages.Notice(w, "repo", err.Error())
1400 return
1401 }
1402
1403 rp.pages.HxRefresh(w)
1404 return
1405 }
1406}
1407
1408func (rp *Repo) ForkRepo(w http.ResponseWriter, r *http.Request) {
1409 l := rp.logger.With("handler", "ForkRepo")
1410
1411 user := rp.oauth.GetMultiAccountUser(r)
1412 f, err := rp.repoResolver.Resolve(r)
1413 if err != nil {
1414 l.Error("failed to resolve source repo", "err", err)
1415 return
1416 }
1417
1418 switch r.Method {
1419 case http.MethodGet:
1420 user := rp.oauth.GetMultiAccountUser(r)
1421 knots := rp.acl.KnotsForUser(r.Context(), user.Did)
1422
1423 spindles, err := rp.enforcer.GetSpindlesForUser(user.Did)
1424 if err != nil {
1425 l.Error("failed to fetch spindles", "err", err)
1426 }
1427
1428 rp.pages.ForkRepo(w, pages.ForkRepoParams{
1429 BaseParams: pages.BaseParamsFromContext(r.Context()),
1430 Knots: knots,
1431 Spindles: spindles,
1432 RepoInfo: rp.repoResolver.GetRepoInfo(r, user),
1433 })
1434
1435 case http.MethodPost:
1436 l := rp.logger.With("handler", "ForkRepo")
1437
1438 targetKnot := r.FormValue("knot")
1439 if targetKnot == "" {
1440 rp.pages.Notice(w, "repo", "Invalid form submission—missing knot domain.")
1441 return
1442 }
1443 l = l.With("targetKnot", targetKnot)
1444
1445 if !rp.acl.IsRepoCreateAllowed(r.Context(), targetKnot, user.Did) {
1446 rp.pages.Notice(w, "repo", "You do not have permission to create a repo in this knot.")
1447 return
1448 }
1449
1450 // optional spindle selection; validate the user is a member if provided
1451 spindle := r.FormValue("spindle")
1452 if spindle != "" {
1453 validSpindles, err := rp.enforcer.GetSpindlesForUser(user.Did)
1454 if err != nil {
1455 l.Error("failed to fetch spindles", "err", err)
1456 rp.pages.Notice(w, "repo", "Failed to configure spindle. Try again later.")
1457 return
1458 }
1459 if !slices.Contains(validSpindles, spindle) {
1460 rp.pages.Notice(w, "repo", "Invalid spindle selection.")
1461 return
1462 }
1463 }
1464
1465 // choose a name for a fork
1466 forkName := strings.ToLower(r.FormValue("repo_name"))
1467 if forkName == "" {
1468 rp.pages.Notice(w, "repo", "Repository name cannot be empty.")
1469 return
1470 }
1471
1472 // this check is *only* to see if the forked repo name already exists
1473 // in the user's account.
1474 existingRepo, err := db.GetRepo(
1475 rp.db,
1476 orm.FilterEq("did", user.Did),
1477 orm.FilterEq("name", forkName),
1478 )
1479 if err != nil {
1480 if !errors.Is(err, sql.ErrNoRows) {
1481 l.Error("error fetching existing repo from db", "err", err)
1482 rp.pages.Notice(w, "repo", "Failed to fork this repository. Try again later.")
1483 return
1484 }
1485 } else if existingRepo != nil {
1486 // repo with this name already exists
1487 rp.pages.Notice(w, "repo", "A repository with this name already exists.")
1488 return
1489 }
1490 l = l.With("forkName", forkName)
1491
1492 uri := "https"
1493 if rp.config.Core.Dev {
1494 uri = "http"
1495 }
1496
1497 forkSourceUrl := fmt.Sprintf("%s://%s/%s", uri, f.Knot, f.RepoIdentifier())
1498 l = l.With("cloneUrl", forkSourceUrl)
1499
1500 rkey := strings.ToLower(forkName)
1501
1502 // TODO: this could coordinate better with the knot to receive a clone status
1503 client, err := rp.oauth.ServiceClient(
1504 r,
1505 oauth.WithService(targetKnot),
1506 oauth.WithLxm(tangled.RepoCreateNSID),
1507 oauth.WithDev(rp.config.Core.Dev),
1508 oauth.WithTimeout(time.Second*20),
1509 )
1510 if err != nil {
1511 l.Error("could not create service client", "err", err)
1512 rp.pages.Notice(w, "repo", "Failed to connect to knot server.")
1513 return
1514 }
1515
1516 forkInput := &tangled.RepoCreate_Input{
1517 Rkey: rkey,
1518 Name: rkey,
1519 Source: &forkSourceUrl,
1520 }
1521 createResp, err := tangled.RepoCreate(
1522 r.Context(),
1523 client,
1524 forkInput,
1525 )
1526 if xrpcerr := xrpcclient.HandleXrpcErr(err); xrpcerr != nil {
1527 l.Error("failed to call XRPC repo.create", "xrpcerr", xrpcerr, "err", err)
1528 rp.pages.Notice(w, "repo", xrpcerr.Error())
1529 return
1530 }
1531
1532 var repoDid string
1533 if createResp != nil && createResp.RepoDid != nil {
1534 repoDid = *createResp.RepoDid
1535 }
1536 if repoDid == "" {
1537 l.Error("knot returned empty repo DID for fork")
1538 rp.pages.Notice(w, "repo", "Knot failed to mint a repo DID. The knot may need to be upgraded.")
1539 return
1540 }
1541
1542 forkSource := f.RepoAt().String()
1543 if f.RepoDid != "" {
1544 forkSource = f.RepoDid
1545 }
1546
1547 forkDescription := r.Form.Get("description")
1548
1549 repo := &models.Repo{
1550 Did: user.Did,
1551 Name: rkey,
1552 Knot: targetKnot,
1553 Rkey: rkey,
1554 Source: forkSource,
1555 Description: forkDescription,
1556 Spindle: spindle,
1557 Created: time.Now(),
1558 Labels: rp.config.Label.DefaultLabelDefs,
1559 RepoDid: repoDid,
1560 }
1561 record := repo.AsRecord()
1562
1563 cleanupKnot := func() {
1564 go func() {
1565 delays := []time.Duration{0, 2 * time.Second, 5 * time.Second}
1566 for attempt, delay := range delays {
1567 time.Sleep(delay)
1568 deleteClient, dErr := rp.oauth.ServiceClient(
1569 r,
1570 oauth.WithService(targetKnot),
1571 oauth.WithLxm(tangled.RepoDeleteNSID),
1572 oauth.WithDev(rp.config.Core.Dev),
1573 )
1574 if dErr != nil {
1575 l.Error("failed to create delete client for knot cleanup", "attempt", attempt+1, "err", dErr)
1576 continue
1577 }
1578 ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
1579 if dErr := tangled.RepoDelete(ctx, deleteClient, &tangled.RepoDelete_Input{
1580 Did: user.Did,
1581 Name: forkName,
1582 Rkey: rkey,
1583 }); dErr != nil {
1584 cancel()
1585 l.Error("failed to clean up fork on knot after rollback", "attempt", attempt+1, "err", dErr)
1586 continue
1587 }
1588 cancel()
1589 l.Info("successfully cleaned up fork on knot after rollback", "attempt", attempt+1)
1590 return
1591 }
1592 l.Error("exhausted retries for knot cleanup, fork may be orphaned",
1593 "did", user.Did, "fork", forkName, "knot", targetKnot)
1594 }()
1595 }
1596
1597 atpClient, err := rp.oauth.AuthorizedClient(r)
1598 if err != nil {
1599 l.Error("failed to create xrpcclient", "err", err)
1600 cleanupKnot()
1601 rp.pages.Notice(w, "repo", "Failed to fork repository.")
1602 return
1603 }
1604
1605 atresp, err := comatproto.RepoPutRecord(r.Context(), atpClient, &comatproto.RepoPutRecord_Input{
1606 Collection: tangled.RepoNSID,
1607 Repo: user.Did,
1608 Rkey: rkey,
1609 Record: &lexutil.LexiconTypeDecoder{
1610 Val: &record,
1611 },
1612 })
1613 if err != nil {
1614 l.Error("failed to write to PDS", "err", err)
1615 cleanupKnot()
1616 rp.pages.Notice(w, "repo", "Failed to announce repository creation.")
1617 return
1618 }
1619
1620 aturi := atresp.Uri
1621 l = l.With("aturi", aturi)
1622 l.Info("wrote to PDS")
1623
1624 tx, err := rp.db.BeginTx(r.Context(), nil)
1625 if err != nil {
1626 l.Info("txn failed", "err", err)
1627 rp.pages.Notice(w, "repo", "Failed to save repository information.")
1628 return
1629 }
1630
1631 rollback := func() {
1632 err1 := tx.Rollback()
1633 err2 := rp.enforcer.E.LoadPolicy()
1634 err3 := rollbackRecord(context.Background(), aturi, atpClient)
1635
1636 if errors.Is(err1, sql.ErrTxDone) {
1637 err1 = nil
1638 }
1639
1640 if errs := errors.Join(err1, err2, err3); errs != nil {
1641 l.Error("failed to rollback changes", "errs", errs)
1642 }
1643
1644 if aturi != "" {
1645 cleanupKnot()
1646 }
1647 }
1648 defer rollback()
1649
1650 err = db.AddRepo(tx, repo)
1651 if err != nil {
1652 l.Error("failed to AddRepo", "err", err)
1653 rp.pages.Notice(w, "repo", "Failed to save repository information.")
1654 return
1655 }
1656
1657 rbacPath := repo.RepoIdentifier()
1658 err = rp.enforcer.AddRepo(user.Did, targetKnot, rbacPath)
1659 if err != nil {
1660 l.Error("failed to add ACLs", "err", err)
1661 rp.pages.Notice(w, "repo", "Failed to set up repository permissions.")
1662 return
1663 }
1664
1665 err = tx.Commit()
1666 if err != nil {
1667 l.Error("failed to commit changes", "err", err)
1668 http.Error(w, err.Error(), http.StatusInternalServerError)
1669 return
1670 }
1671
1672 err = rp.enforcer.E.SavePolicy()
1673 if err != nil {
1674 l.Error("failed to update ACLs", "err", err)
1675 http.Error(w, err.Error(), http.StatusInternalServerError)
1676 return
1677 }
1678
1679 aturi = ""
1680
1681 rp.notifier.NewRepo(r.Context(), repo)
1682 if repoDid != "" {
1683 rp.pages.HxLocation(w, fmt.Sprintf("/%s", repoDid))
1684 } else {
1685 rp.pages.HxLocation(w, fmt.Sprintf("/%s/%s", user.Did, forkName))
1686 }
1687 }
1688}
1689
1690func (rp *Repo) Stars(w http.ResponseWriter, r *http.Request) {
1691 l := rp.logger.With("handler", "Stars")
1692
1693 user := rp.oauth.GetMultiAccountUser(r)
1694 f, err := rp.repoResolver.Resolve(r)
1695 if err != nil {
1696 l.Error("failed to resolve source repo", "err", err)
1697 return
1698 }
1699
1700 page := pagination.FromContext(r.Context())
1701 if page.Limit > 30 || page.Limit <= 0 {
1702 page.Limit = 30
1703 }
1704
1705 starrers, err := db.GetStars(rp.db, string(f.RepoDid), page)
1706 if err != nil {
1707 l.Error("failed to fetch starrers", "err", err, "repoDid", f.RepoDid)
1708 return
1709 }
1710
1711 totalCount, err := db.GetStarCount(rp.db, models.StarSubjectRepo, string(f.RepoDid))
1712 if err != nil {
1713 l.Error("failed to fetch star count", "err", err, "repoDid", f.RepoDid)
1714 return
1715 }
1716
1717 rp.pages.RepoStars(w, pages.RepoStarsParams{
1718 BaseParams: pages.BaseParamsFromContext(r.Context()),
1719 RepoInfo: rp.repoResolver.GetRepoInfo(r, user),
1720 Starrers: starrers,
1721 Page: page,
1722 TotalCount: totalCount,
1723 })
1724}
1725
1726func (rp *Repo) Forks(w http.ResponseWriter, r *http.Request) {
1727 l := rp.logger.With("handler", "Forks")
1728
1729 user := rp.oauth.GetMultiAccountUser(r)
1730 f, err := rp.repoResolver.Resolve(r)
1731 if err != nil {
1732 l.Error("failed to resolve source repo", "err", err)
1733 return
1734 }
1735
1736 var forks []models.Repo
1737 totalCount := 0
1738 page := pagination.FromContext(r.Context())
1739 if f.RepoDid != "" {
1740 forks, err = db.GetReposPaginated(rp.db, page, orm.FilterEq("source", f.RepoDid))
1741 if err != nil {
1742 l.Error("failed to fetch forks", "err", err, "repoAt", f.RepoAt())
1743 return
1744 }
1745
1746 totalCount, err = db.GetForkCount(rp.db, f.RepoDid)
1747 if err != nil {
1748 l.Error("failed to fetch fork count", "err", err, "repoAt", f.RepoAt())
1749 return
1750 }
1751 }
1752
1753 err = rp.pages.RepoForks(w, pages.RepoForksParams{
1754 BaseParams: pages.BaseParamsFromContext(r.Context()),
1755 RepoInfo: rp.repoResolver.GetRepoInfo(r, user),
1756 Forks: forks,
1757 Page: page,
1758 TotalCount: totalCount,
1759 })
1760 if err != nil {
1761 l.Error("failed to render page", "err", err)
1762 }
1763}
1764
1765// this is used to rollback changes made to the PDS
1766//
1767// it is a no-op if the provided ATURI is empty
1768func rollbackRecord(ctx context.Context, aturi string, client *atclient.APIClient) error {
1769 if aturi == "" {
1770 return nil
1771 }
1772
1773 parsed := syntax.ATURI(aturi)
1774
1775 collection := parsed.Collection().String()
1776 repo := parsed.Authority().String()
1777 rkey := parsed.RecordKey().String()
1778
1779 _, err := comatproto.RepoDeleteRecord(ctx, client, &comatproto.RepoDeleteRecord_Input{
1780 Collection: collection,
1781 Repo: repo,
1782 Rkey: rkey,
1783 })
1784 return err
1785}
1786
1787func repoCollaboratorRecord(f *models.Repo, subject string, createdAt time.Time) *tangled.RepoCollaborator {
1788 return &tangled.RepoCollaborator{
1789 Subject: subject,
1790 CreatedAt: createdAt.Format(time.RFC3339),
1791 Repo: f.RepoDid,
1792 }
1793}