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