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