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