This repository has no description
1package repo
2
3import (
4 "context"
5 "database/sql"
6 "errors"
7 "fmt"
8 "log/slog"
9 "net/http"
10 "net/url"
11 "slices"
12 "strings"
13 "time"
14
15 "tangled.org/core/appview/cloudflare"
16
17 "tangled.org/core/api/tangled"
18 "tangled.org/core/appview/config"
19 "tangled.org/core/appview/db"
20 "tangled.org/core/appview/models"
21 "tangled.org/core/appview/notify"
22 "tangled.org/core/appview/oauth"
23 "tangled.org/core/appview/pages"
24 "tangled.org/core/appview/reporesolver"
25 "tangled.org/core/appview/validator"
26 xrpcclient "tangled.org/core/appview/xrpcclient"
27 "tangled.org/core/eventconsumer"
28 "tangled.org/core/idresolver"
29 "tangled.org/core/ogre"
30 "tangled.org/core/orm"
31 "tangled.org/core/rbac"
32 "tangled.org/core/tid"
33 "tangled.org/core/xrpc/serviceauth"
34
35 comatproto "github.com/bluesky-social/indigo/api/atproto"
36 "github.com/bluesky-social/indigo/atproto/atclient"
37 "github.com/bluesky-social/indigo/atproto/syntax"
38 lexutil "github.com/bluesky-social/indigo/lex/util"
39
40 "github.com/go-chi/chi/v5"
41)
42
43type Repo struct {
44 repoResolver *reporesolver.RepoResolver
45 idResolver *idresolver.Resolver
46 config *config.Config
47 oauth *oauth.OAuth
48 pages *pages.Pages
49 spindlestream *eventconsumer.Consumer
50 db *db.DB
51 enforcer *rbac.Enforcer
52 notifier notify.Notifier
53 logger *slog.Logger
54 serviceAuth *serviceauth.ServiceAuth
55 validator *validator.Validator
56 cfClient *cloudflare.Client
57 ogreClient *ogre.Client
58}
59
60func New(
61 oauth *oauth.OAuth,
62 repoResolver *reporesolver.RepoResolver,
63 pages *pages.Pages,
64 spindlestream *eventconsumer.Consumer,
65 idResolver *idresolver.Resolver,
66 db *db.DB,
67 config *config.Config,
68 notifier notify.Notifier,
69 enforcer *rbac.Enforcer,
70 logger *slog.Logger,
71 validator *validator.Validator,
72 cfClient *cloudflare.Client,
73) *Repo {
74 return &Repo{
75 oauth: oauth,
76 repoResolver: repoResolver,
77 pages: pages,
78 idResolver: idResolver,
79 config: config,
80 spindlestream: spindlestream,
81 db: db,
82 notifier: notifier,
83 enforcer: enforcer,
84 logger: logger,
85 validator: validator,
86 cfClient: cfClient,
87 ogreClient: ogre.NewClient(config.Ogre.Host),
88 }
89}
90
91// modify the spindle configured for this repo
92func (rp *Repo) EditSpindle(w http.ResponseWriter, r *http.Request) {
93 user := rp.oauth.GetMultiAccountUser(r)
94 l := rp.logger.With("handler", "EditSpindle")
95 l = l.With("did", user.Did)
96
97 errorId := "operation-error"
98 fail := func(msg string, err error) {
99 l.Error(msg, "err", err)
100 rp.pages.Notice(w, errorId, msg)
101 }
102
103 f, err := rp.repoResolver.Resolve(r)
104 if err != nil {
105 fail("Failed to resolve repo. Try again later", err)
106 return
107 }
108
109 newSpindle := r.FormValue("spindle")
110 removingSpindle := newSpindle == "[[none]]" // see pages/templates/repo/settings/pipelines.html for more info on why we use this value
111 client, err := rp.oauth.AuthorizedClient(r)
112 if err != nil {
113 fail("Failed to authorize. Try again later.", err)
114 return
115 }
116
117 if !removingSpindle {
118 // ensure that this is a valid spindle for this user
119 validSpindles, err := rp.enforcer.GetSpindlesForUser(user.Did)
120 if err != nil {
121 fail("Failed to find spindles. Try again later.", err)
122 return
123 }
124
125 if !slices.Contains(validSpindles, newSpindle) {
126 fail("Failed to configure spindle.", fmt.Errorf("%s is not a valid spindle: %q", newSpindle, validSpindles))
127 return
128 }
129 }
130
131 newRepo := *f
132 newRepo.Spindle = newSpindle
133 record := newRepo.AsRecord()
134
135 spindlePtr := &newSpindle
136 if removingSpindle {
137 spindlePtr = nil
138 newRepo.Spindle = ""
139 }
140
141 // optimistic update
142 err = db.UpdateSpindle(rp.db, newRepo.RepoAt().String(), spindlePtr)
143 if err != nil {
144 fail("Failed to update spindle. Try again later.", err)
145 return
146 }
147
148 ex, err := comatproto.RepoGetRecord(r.Context(), client, "", tangled.RepoNSID, newRepo.Did, newRepo.Rkey)
149 if err != nil {
150 fail("Failed to update spindle, no record found on PDS.", err)
151 return
152 }
153 _, err = comatproto.RepoPutRecord(r.Context(), client, &comatproto.RepoPutRecord_Input{
154 Collection: tangled.RepoNSID,
155 Repo: newRepo.Did,
156 Rkey: newRepo.Rkey,
157 SwapRecord: ex.Cid,
158 Record: &lexutil.LexiconTypeDecoder{
159 Val: &record,
160 },
161 })
162
163 if err != nil {
164 fail("Failed to update spindle, unable to save to PDS.", err)
165 return
166 }
167
168 if !removingSpindle {
169 // add this spindle to spindle stream
170 rp.spindlestream.AddSource(
171 context.Background(),
172 eventconsumer.NewSpindleSource(newSpindle),
173 )
174 }
175
176 rp.pages.HxRefresh(w)
177}
178
179func (rp *Repo) AddLabelDef(w http.ResponseWriter, r *http.Request) {
180 user := rp.oauth.GetMultiAccountUser(r)
181 l := rp.logger.With("handler", "AddLabel")
182 l = l.With("did", user.Did)
183
184 f, err := rp.repoResolver.Resolve(r)
185 if err != nil {
186 l.Error("failed to get repo and knot", "err", err)
187 return
188 }
189
190 errorId := "add-label-error"
191 fail := func(msg string, err error) {
192 l.Error(msg, "err", err)
193 rp.pages.Notice(w, errorId, msg)
194 }
195
196 // get form values for label definition
197 name := r.FormValue("name")
198 concreteType := r.FormValue("valueType")
199 valueFormat := r.FormValue("valueFormat")
200 enumValues := r.FormValue("enumValues")
201 scope := r.Form["scope"]
202 color := r.FormValue("color")
203 multiple := r.FormValue("multiple") == "true"
204
205 var variants []string
206 for part := range strings.SplitSeq(enumValues, ",") {
207 if part = strings.TrimSpace(part); part != "" {
208 variants = append(variants, part)
209 }
210 }
211
212 if concreteType == "" {
213 concreteType = "null"
214 }
215
216 format := models.ValueTypeFormatAny
217 if valueFormat == "did" {
218 format = models.ValueTypeFormatDid
219 }
220
221 valueType := models.ValueType{
222 Type: models.ConcreteType(concreteType),
223 Format: format,
224 Enum: variants,
225 }
226
227 label := models.LabelDefinition{
228 Did: user.Did,
229 Rkey: tid.TID(),
230 Name: name,
231 ValueType: valueType,
232 Scope: scope,
233 Color: &color,
234 Multiple: multiple,
235 Created: time.Now(),
236 }
237 if err := rp.validator.ValidateLabelDefinition(&label); err != nil {
238 fail(err.Error(), err)
239 return
240 }
241
242 // announce this relation into the firehose, store into owners' pds
243 client, err := rp.oauth.AuthorizedClient(r)
244 if err != nil {
245 fail(err.Error(), err)
246 return
247 }
248
249 // emit a labelRecord
250 labelRecord := label.AsRecord()
251 resp, err := comatproto.RepoPutRecord(r.Context(), client, &comatproto.RepoPutRecord_Input{
252 Collection: tangled.LabelDefinitionNSID,
253 Repo: label.Did,
254 Rkey: label.Rkey,
255 Record: &lexutil.LexiconTypeDecoder{
256 Val: &labelRecord,
257 },
258 })
259 // invalid record
260 if err != nil {
261 fail("Failed to write record to PDS.", err)
262 return
263 }
264
265 aturi := resp.Uri
266 l = l.With("at-uri", aturi)
267 l.Info("wrote label record to PDS")
268
269 // update the repo to subscribe to this label
270 newRepo := *f
271 newRepo.Labels = append(newRepo.Labels, aturi)
272 repoRecord := newRepo.AsRecord()
273
274 ex, err := comatproto.RepoGetRecord(r.Context(), client, "", tangled.RepoNSID, newRepo.Did, newRepo.Rkey)
275 if err != nil {
276 fail("Failed to update labels, no record found on PDS.", err)
277 return
278 }
279 _, err = comatproto.RepoPutRecord(r.Context(), client, &comatproto.RepoPutRecord_Input{
280 Collection: tangled.RepoNSID,
281 Repo: newRepo.Did,
282 Rkey: newRepo.Rkey,
283 SwapRecord: ex.Cid,
284 Record: &lexutil.LexiconTypeDecoder{
285 Val: &repoRecord,
286 },
287 })
288 if err != nil {
289 fail("Failed to update labels for repo.", err)
290 return
291 }
292
293 tx, err := rp.db.BeginTx(r.Context(), nil)
294 if err != nil {
295 fail("Failed to add label.", err)
296 return
297 }
298
299 rollback := func() {
300 err1 := tx.Rollback()
301 err2 := rollbackRecord(context.Background(), aturi, client)
302
303 // ignore txn complete errors, this is okay
304 if errors.Is(err1, sql.ErrTxDone) {
305 err1 = nil
306 }
307
308 if errs := errors.Join(err1, err2); errs != nil {
309 l.Error("failed to rollback changes", "errs", errs)
310 return
311 }
312 }
313 defer rollback()
314
315 _, err = db.AddLabelDefinition(tx, &label)
316 if err != nil {
317 fail("Failed to add label.", err)
318 return
319 }
320
321 if err = db.SubscribeLabel(tx, &models.RepoLabel{
322 RepoAt: f.RepoAt(),
323 LabelAt: label.AtUri(),
324 }); err != nil {
325 fail("Failed to subscribe to label.", err)
326 return
327 }
328
329 err = tx.Commit()
330 if err != nil {
331 fail("Failed to add label.", err)
332 return
333 }
334
335 // clear aturi when everything is successful
336 aturi = ""
337
338 rp.pages.HxRefresh(w)
339}
340
341func (rp *Repo) DeleteLabelDef(w http.ResponseWriter, r *http.Request) {
342 user := rp.oauth.GetMultiAccountUser(r)
343 l := rp.logger.With("handler", "DeleteLabel")
344 l = l.With("did", user.Did)
345
346 f, err := rp.repoResolver.Resolve(r)
347 if err != nil {
348 l.Error("failed to get repo and knot", "err", err)
349 return
350 }
351
352 errorId := "label-operation"
353 fail := func(msg string, err error) {
354 l.Error(msg, "err", err)
355 rp.pages.Notice(w, errorId, msg)
356 }
357
358 // get form values
359 labelId := r.FormValue("label-id")
360
361 label, err := db.GetLabelDefinition(rp.db, orm.FilterEq("id", labelId))
362 if err != nil {
363 fail("Failed to find label definition.", err)
364 return
365 }
366
367 client, err := rp.oauth.AuthorizedClient(r)
368 if err != nil {
369 fail(err.Error(), err)
370 return
371 }
372
373 // delete label record from PDS
374 _, err = comatproto.RepoDeleteRecord(r.Context(), client, &comatproto.RepoDeleteRecord_Input{
375 Collection: tangled.LabelDefinitionNSID,
376 Repo: label.Did,
377 Rkey: label.Rkey,
378 })
379 if err != nil {
380 fail("Failed to delete label record from PDS.", err)
381 return
382 }
383
384 // update repo record to remove the label reference
385 newRepo := *f
386 var updated []string
387 removedAt := label.AtUri().String()
388 for _, l := range newRepo.Labels {
389 if l != removedAt {
390 updated = append(updated, l)
391 }
392 }
393 newRepo.Labels = updated
394 repoRecord := newRepo.AsRecord()
395
396 ex, err := comatproto.RepoGetRecord(r.Context(), client, "", tangled.RepoNSID, newRepo.Did, newRepo.Rkey)
397 if err != nil {
398 fail("Failed to update labels, no record found on PDS.", err)
399 return
400 }
401 _, err = comatproto.RepoPutRecord(r.Context(), client, &comatproto.RepoPutRecord_Input{
402 Collection: tangled.RepoNSID,
403 Repo: newRepo.Did,
404 Rkey: newRepo.Rkey,
405 SwapRecord: ex.Cid,
406 Record: &lexutil.LexiconTypeDecoder{
407 Val: &repoRecord,
408 },
409 })
410 if err != nil {
411 fail("Failed to update repo record.", err)
412 return
413 }
414
415 // transaction for DB changes
416 tx, err := rp.db.BeginTx(r.Context(), nil)
417 if err != nil {
418 fail("Failed to delete label.", err)
419 return
420 }
421 defer tx.Rollback()
422
423 err = db.UnsubscribeLabel(
424 tx,
425 orm.FilterEq("repo_at", f.RepoAt()),
426 orm.FilterEq("label_at", removedAt),
427 )
428 if err != nil {
429 fail("Failed to unsubscribe label.", err)
430 return
431 }
432
433 err = db.DeleteLabelDefinition(tx, orm.FilterEq("id", label.Id))
434 if err != nil {
435 fail("Failed to delete label definition.", err)
436 return
437 }
438
439 err = tx.Commit()
440 if err != nil {
441 fail("Failed to delete label.", err)
442 return
443 }
444
445 // everything succeeded
446 rp.pages.HxRefresh(w)
447}
448
449func (rp *Repo) SubscribeLabel(w http.ResponseWriter, r *http.Request) {
450 user := rp.oauth.GetMultiAccountUser(r)
451 l := rp.logger.With("handler", "SubscribeLabel")
452 l = l.With("did", user.Did)
453
454 f, err := rp.repoResolver.Resolve(r)
455 if err != nil {
456 l.Error("failed to get repo and knot", "err", err)
457 return
458 }
459
460 if err := r.ParseForm(); err != nil {
461 l.Error("invalid form", "err", err)
462 return
463 }
464
465 errorId := "default-label-operation"
466 fail := func(msg string, err error) {
467 l.Error(msg, "err", err)
468 rp.pages.Notice(w, errorId, msg)
469 }
470
471 labelAts := r.Form["label"]
472 _, err = db.GetLabelDefinitions(rp.db, orm.FilterIn("at_uri", labelAts))
473 if err != nil {
474 fail("Failed to subscribe to label.", err)
475 return
476 }
477
478 newRepo := *f
479 newRepo.Labels = append(newRepo.Labels, labelAts...)
480
481 // dedup
482 slices.Sort(newRepo.Labels)
483 newRepo.Labels = slices.Compact(newRepo.Labels)
484
485 repoRecord := newRepo.AsRecord()
486
487 client, err := rp.oauth.AuthorizedClient(r)
488 if err != nil {
489 fail(err.Error(), err)
490 return
491 }
492
493 ex, err := comatproto.RepoGetRecord(r.Context(), client, "", tangled.RepoNSID, f.Did, f.Rkey)
494 if err != nil {
495 fail("Failed to update labels, no record found on PDS.", err)
496 return
497 }
498 _, err = comatproto.RepoPutRecord(r.Context(), client, &comatproto.RepoPutRecord_Input{
499 Collection: tangled.RepoNSID,
500 Repo: newRepo.Did,
501 Rkey: newRepo.Rkey,
502 SwapRecord: ex.Cid,
503 Record: &lexutil.LexiconTypeDecoder{
504 Val: &repoRecord,
505 },
506 })
507
508 tx, err := rp.db.Begin()
509 if err != nil {
510 fail("Failed to subscribe to label.", err)
511 return
512 }
513 defer tx.Rollback()
514
515 for _, l := range labelAts {
516 err = db.SubscribeLabel(tx, &models.RepoLabel{
517 RepoAt: f.RepoAt(),
518 LabelAt: syntax.ATURI(l),
519 })
520 if err != nil {
521 fail("Failed to subscribe to label.", err)
522 return
523 }
524 }
525
526 if err := tx.Commit(); err != nil {
527 fail("Failed to subscribe to label.", err)
528 return
529 }
530
531 // everything succeeded
532 rp.pages.HxRefresh(w)
533}
534
535func (rp *Repo) UnsubscribeLabel(w http.ResponseWriter, r *http.Request) {
536 user := rp.oauth.GetMultiAccountUser(r)
537 l := rp.logger.With("handler", "UnsubscribeLabel")
538 l = l.With("did", user.Did)
539
540 f, err := rp.repoResolver.Resolve(r)
541 if err != nil {
542 l.Error("failed to get repo and knot", "err", err)
543 return
544 }
545
546 if err := r.ParseForm(); err != nil {
547 l.Error("invalid form", "err", err)
548 return
549 }
550
551 errorId := "default-label-operation"
552 fail := func(msg string, err error) {
553 l.Error(msg, "err", err)
554 rp.pages.Notice(w, errorId, msg)
555 }
556
557 labelAts := r.Form["label"]
558 _, err = db.GetLabelDefinitions(rp.db, orm.FilterIn("at_uri", labelAts))
559 if err != nil {
560 fail("Failed to unsubscribe to label.", err)
561 return
562 }
563
564 // update repo record to remove the label reference
565 newRepo := *f
566 var updated []string
567 for _, l := range newRepo.Labels {
568 if !slices.Contains(labelAts, l) {
569 updated = append(updated, l)
570 }
571 }
572 newRepo.Labels = updated
573 repoRecord := newRepo.AsRecord()
574
575 client, err := rp.oauth.AuthorizedClient(r)
576 if err != nil {
577 fail(err.Error(), err)
578 return
579 }
580
581 ex, err := comatproto.RepoGetRecord(r.Context(), client, "", tangled.RepoNSID, f.Did, f.Rkey)
582 if err != nil {
583 fail("Failed to update labels, no record found on PDS.", err)
584 return
585 }
586 _, err = comatproto.RepoPutRecord(r.Context(), client, &comatproto.RepoPutRecord_Input{
587 Collection: tangled.RepoNSID,
588 Repo: newRepo.Did,
589 Rkey: newRepo.Rkey,
590 SwapRecord: ex.Cid,
591 Record: &lexutil.LexiconTypeDecoder{
592 Val: &repoRecord,
593 },
594 })
595
596 err = db.UnsubscribeLabel(
597 rp.db,
598 orm.FilterEq("repo_at", f.RepoAt()),
599 orm.FilterIn("label_at", labelAts),
600 )
601 if err != nil {
602 fail("Failed to unsubscribe label.", err)
603 return
604 }
605
606 // everything succeeded
607 rp.pages.HxRefresh(w)
608}
609
610func (rp *Repo) LabelPanel(w http.ResponseWriter, r *http.Request) {
611 l := rp.logger.With("handler", "LabelPanel")
612
613 f, err := rp.repoResolver.Resolve(r)
614 if err != nil {
615 l.Error("failed to get repo and knot", "err", err)
616 return
617 }
618
619 subjectStr := r.FormValue("subject")
620 subject, err := syntax.ParseATURI(subjectStr)
621 if err != nil {
622 l.Error("failed to get repo and knot", "err", err)
623 return
624 }
625
626 labelDefs, err := db.GetLabelDefinitions(
627 rp.db,
628 orm.FilterIn("at_uri", f.Labels),
629 orm.FilterContains("scope", subject.Collection().String()),
630 )
631 if err != nil {
632 l.Error("failed to fetch label defs", "err", err)
633 return
634 }
635
636 defs := make(map[string]*models.LabelDefinition)
637 for _, l := range labelDefs {
638 defs[l.AtUri().String()] = &l
639 }
640
641 states, err := db.GetLabels(rp.db, orm.FilterEq("subject", subject))
642 if err != nil {
643 l.Error("failed to build label state", "err", err)
644 return
645 }
646 state := states[subject]
647
648 user := rp.oauth.GetMultiAccountUser(r)
649 rp.pages.LabelPanel(w, pages.LabelPanelParams{
650 LoggedInUser: user,
651 RepoInfo: rp.repoResolver.GetRepoInfo(r, user),
652 Defs: defs,
653 Subject: subject.String(),
654 State: state,
655 })
656}
657
658func (rp *Repo) EditLabelPanel(w http.ResponseWriter, r *http.Request) {
659 l := rp.logger.With("handler", "EditLabelPanel")
660
661 f, err := rp.repoResolver.Resolve(r)
662 if err != nil {
663 l.Error("failed to get repo and knot", "err", err)
664 return
665 }
666
667 subjectStr := r.FormValue("subject")
668 subject, err := syntax.ParseATURI(subjectStr)
669 if err != nil {
670 l.Error("failed to get repo and knot", "err", err)
671 return
672 }
673
674 labelDefs, err := db.GetLabelDefinitions(
675 rp.db,
676 orm.FilterIn("at_uri", f.Labels),
677 orm.FilterContains("scope", subject.Collection().String()),
678 )
679 if err != nil {
680 l.Error("failed to fetch labels", "err", err)
681 return
682 }
683
684 defs := make(map[string]*models.LabelDefinition)
685 for _, l := range labelDefs {
686 defs[l.AtUri().String()] = &l
687 }
688
689 states, err := db.GetLabels(rp.db, orm.FilterEq("subject", subject))
690 if err != nil {
691 l.Error("failed to build label state", "err", err)
692 return
693 }
694 state := states[subject]
695
696 user := rp.oauth.GetMultiAccountUser(r)
697 rp.pages.EditLabelPanel(w, pages.EditLabelPanelParams{
698 LoggedInUser: user,
699 RepoInfo: rp.repoResolver.GetRepoInfo(r, user),
700 Defs: defs,
701 Subject: subject.String(),
702 State: state,
703 })
704}
705
706func (rp *Repo) AddCollaborator(w http.ResponseWriter, r *http.Request) {
707 user := rp.oauth.GetMultiAccountUser(r)
708 l := rp.logger.With("handler", "AddCollaborator")
709 l = l.With("did", user.Did)
710
711 f, err := rp.repoResolver.Resolve(r)
712 if err != nil {
713 l.Error("failed to get repo and knot", "err", err)
714 return
715 }
716
717 errorId := "add-collaborator-error"
718 fail := func(msg string, err error) {
719 l.Error(msg, "err", err)
720 rp.pages.Notice(w, errorId, msg)
721 }
722
723 collaborator := r.FormValue("collaborator")
724 if collaborator == "" {
725 fail("Invalid form.", nil)
726 return
727 }
728
729 // remove a single leading `@`, to make @handle work with ResolveIdent
730 collaborator = strings.TrimPrefix(collaborator, "@")
731
732 collaboratorIdent, err := rp.idResolver.ResolveIdent(r.Context(), collaborator)
733 if err != nil {
734 fail(fmt.Sprintf("'%s' is not a valid DID/handle.", collaborator), err)
735 return
736 }
737
738 if collaboratorIdent.DID.String() == user.Did {
739 fail("You seem to be adding yourself as a collaborator.", nil)
740 return
741 }
742 l = l.With("collaborator", collaboratorIdent.Handle)
743 l = l.With("knot", f.Knot)
744
745 // announce this relation into the firehose, store into owners' pds
746 client, err := rp.oauth.AuthorizedClient(r)
747 if err != nil {
748 fail("Failed to write to PDS.", err)
749 return
750 }
751
752 // emit a record
753 currentUser := rp.oauth.GetMultiAccountUser(r)
754 rkey := tid.TID()
755 createdAt := time.Now()
756 resp, err := comatproto.RepoPutRecord(r.Context(), client, &comatproto.RepoPutRecord_Input{
757 Collection: tangled.RepoCollaboratorNSID,
758 Repo: currentUser.Did,
759 Rkey: rkey,
760 Record: &lexutil.LexiconTypeDecoder{
761 Val: repoCollaboratorRecord(f, collaboratorIdent.DID.String(), createdAt),
762 },
763 })
764 // invalid record
765 if err != nil {
766 fail("Failed to write record to PDS.", err)
767 return
768 }
769
770 aturi := resp.Uri
771 l = l.With("at-uri", aturi)
772 l.Info("wrote record to PDS")
773
774 tx, err := rp.db.BeginTx(r.Context(), nil)
775 if err != nil {
776 fail("Failed to add collaborator.", err)
777 return
778 }
779
780 rollback := func() {
781 err1 := tx.Rollback()
782 err2 := rp.enforcer.E.LoadPolicy()
783 err3 := rollbackRecord(context.Background(), aturi, client)
784
785 // ignore txn complete errors, this is okay
786 if errors.Is(err1, sql.ErrTxDone) {
787 err1 = nil
788 }
789
790 if errs := errors.Join(err1, err2, err3); errs != nil {
791 l.Error("failed to rollback changes", "errs", errs)
792 return
793 }
794 }
795 defer rollback()
796
797 err = rp.enforcer.AddCollaborator(collaboratorIdent.DID.String(), f.Knot, f.RepoIdentifier())
798 if err != nil {
799 fail("Failed to add collaborator permissions.", err)
800 return
801 }
802
803 err = db.AddCollaborator(tx, models.Collaborator{
804 Did: syntax.DID(currentUser.Did),
805 Rkey: rkey,
806 SubjectDid: collaboratorIdent.DID,
807 RepoAt: f.RepoAt(),
808 Created: createdAt,
809 })
810 if err != nil {
811 fail("Failed to add collaborator.", err)
812 return
813 }
814
815 err = tx.Commit()
816 if err != nil {
817 fail("Failed to add collaborator.", err)
818 return
819 }
820
821 err = rp.enforcer.E.SavePolicy()
822 if err != nil {
823 fail("Failed to update collaborator permissions.", err)
824 return
825 }
826
827 // clear aturi to when everything is successful
828 aturi = ""
829
830 rp.pages.HxRefresh(w)
831}
832
833func (rp *Repo) DeleteRepo(w http.ResponseWriter, r *http.Request) {
834 user := rp.oauth.GetMultiAccountUser(r)
835 l := rp.logger.With("handler", "DeleteRepo")
836
837 noticeId := "operation-error"
838 f, err := rp.repoResolver.Resolve(r)
839 if err != nil {
840 l.Error("failed to get repo and knot", "err", err)
841 return
842 }
843
844 // remove record from pds
845 atpClient, err := rp.oauth.AuthorizedClient(r)
846 if err != nil {
847 l.Error("failed to get authorized client", "err", err)
848 return
849 }
850 _, err = comatproto.RepoDeleteRecord(r.Context(), atpClient, &comatproto.RepoDeleteRecord_Input{
851 Collection: tangled.RepoNSID,
852 Repo: user.Did,
853 Rkey: f.Rkey,
854 })
855 if err != nil {
856 l.Error("failed to delete record", "err", err)
857 rp.pages.Notice(w, noticeId, "Failed to delete repository from PDS.")
858 return
859 }
860 l.Info("removed repo record", "aturi", f.RepoAt().String())
861
862 client, err := rp.oauth.ServiceClient(
863 r,
864 oauth.WithService(f.Knot),
865 oauth.WithLxm(tangled.RepoDeleteNSID),
866 oauth.WithDev(rp.config.Core.Dev),
867 )
868 if err != nil {
869 l.Error("failed to connect to knot server", "err", err)
870 return
871 }
872
873 err = tangled.RepoDelete(
874 r.Context(),
875 client,
876 &tangled.RepoDelete_Input{
877 Did: f.Did,
878 Name: f.Name,
879 Rkey: f.Rkey,
880 },
881 )
882 if xrpcerr := xrpcclient.HandleXrpcErr(err); xrpcerr != nil {
883 l.Error("failed to call XRPC repo.delete", "xrpcerr", xrpcerr, "err", err)
884 rp.pages.Notice(w, noticeId, xrpcerr.Error())
885 return
886 }
887 l.Info("deleted repo from knot")
888
889 tx, err := rp.db.BeginTx(r.Context(), nil)
890 if err != nil {
891 l.Error("failed to start tx")
892 w.Write(fmt.Append(nil, "failed to add collaborator: ", err))
893 return
894 }
895 defer func() {
896 tx.Rollback()
897 err = rp.enforcer.E.LoadPolicy()
898 if err != nil {
899 l.Error("failed to rollback policies")
900 }
901 }()
902
903 // remove collaborator RBAC
904 repoCollaborators, err := rp.enforcer.E.GetImplicitUsersForResourceByDomain(f.RepoIdentifier(), f.Knot)
905 if err != nil {
906 rp.pages.Notice(w, noticeId, "Failed to remove collaborators")
907 return
908 }
909 for _, c := range repoCollaborators {
910 did := c[0]
911 rp.enforcer.RemoveCollaborator(did, f.Knot, f.RepoIdentifier())
912 }
913 l.Info("removed collaborators")
914
915 // remove repo RBAC
916 err = rp.enforcer.RemoveRepo(f.Did, f.Knot, f.RepoIdentifier())
917 if err != nil {
918 rp.pages.Notice(w, noticeId, "Failed to update RBAC rules")
919 return
920 }
921
922 // remove repo from db
923 err = db.RemoveRepo(tx, f.Did, f.Name)
924 if err != nil {
925 rp.pages.Notice(w, noticeId, "Failed to update appview")
926 return
927 }
928 l.Info("removed repo from db")
929
930 err = tx.Commit()
931 if err != nil {
932 l.Error("failed to commit changes", "err", err)
933 http.Error(w, err.Error(), http.StatusInternalServerError)
934 return
935 }
936
937 err = rp.enforcer.E.SavePolicy()
938 if err != nil {
939 l.Error("failed to update ACLs", "err", err)
940 http.Error(w, err.Error(), http.StatusInternalServerError)
941 return
942 }
943
944 rp.notifier.DeleteRepo(r.Context(), f)
945 rp.pages.HxRedirect(w, fmt.Sprintf("/%s", f.Did))
946}
947
948func (rp *Repo) SyncRepoFork(w http.ResponseWriter, r *http.Request) {
949 l := rp.logger.With("handler", "SyncRepoFork")
950
951 ref := chi.URLParam(r, "ref")
952 ref, _ = url.PathUnescape(ref)
953
954 user := rp.oauth.GetMultiAccountUser(r)
955 f, err := rp.repoResolver.Resolve(r)
956 if err != nil {
957 l.Error("failed to resolve source repo", "err", err)
958 return
959 }
960
961 switch r.Method {
962 case http.MethodPost:
963 client, err := rp.oauth.ServiceClient(
964 r,
965 oauth.WithService(f.Knot),
966 oauth.WithLxm(tangled.RepoForkSyncNSID),
967 oauth.WithDev(rp.config.Core.Dev),
968 )
969 if err != nil {
970 rp.pages.Notice(w, "repo", "Failed to connect to knot server.")
971 return
972 }
973
974 if f.Source == "" {
975 rp.pages.Notice(w, "repo", "This repository is not a fork.")
976 return
977 }
978
979 err = tangled.RepoForkSync(
980 r.Context(),
981 client,
982 &tangled.RepoForkSync_Input{
983 Did: user.Did,
984 Name: f.Name,
985 Source: f.Source,
986 Branch: ref,
987 },
988 )
989 if xrpcerr := xrpcclient.HandleXrpcErr(err); xrpcerr != nil {
990 l.Error("failed to call XRPC repo.forkSync", "xrpcerr", xrpcerr, "err", err)
991 rp.pages.Notice(w, "repo", err.Error())
992 return
993 }
994
995 rp.pages.HxRefresh(w)
996 return
997 }
998}
999
1000func (rp *Repo) ForkRepo(w http.ResponseWriter, r *http.Request) {
1001 l := rp.logger.With("handler", "ForkRepo")
1002
1003 user := rp.oauth.GetMultiAccountUser(r)
1004 f, err := rp.repoResolver.Resolve(r)
1005 if err != nil {
1006 l.Error("failed to resolve source repo", "err", err)
1007 return
1008 }
1009
1010 switch r.Method {
1011 case http.MethodGet:
1012 user := rp.oauth.GetMultiAccountUser(r)
1013 knots, err := rp.enforcer.GetKnotsForUser(user.Did)
1014 if err != nil {
1015 rp.pages.Notice(w, "repo", "Invalid user account.")
1016 return
1017 }
1018
1019 rp.pages.ForkRepo(w, pages.ForkRepoParams{
1020 LoggedInUser: user,
1021 Knots: knots,
1022 RepoInfo: rp.repoResolver.GetRepoInfo(r, user),
1023 })
1024
1025 case http.MethodPost:
1026 l := rp.logger.With("handler", "ForkRepo")
1027
1028 targetKnot := r.FormValue("knot")
1029 if targetKnot == "" {
1030 rp.pages.Notice(w, "repo", "Invalid form submission—missing knot domain.")
1031 return
1032 }
1033 l = l.With("targetKnot", targetKnot)
1034
1035 ok, err := rp.enforcer.E.Enforce(user.Did, targetKnot, targetKnot, "repo:create")
1036 if err != nil || !ok {
1037 rp.pages.Notice(w, "repo", "You do not have permission to create a repo in this knot.")
1038 return
1039 }
1040
1041 // choose a name for a fork
1042 forkName := r.FormValue("repo_name")
1043 if forkName == "" {
1044 rp.pages.Notice(w, "repo", "Repository name cannot be empty.")
1045 return
1046 }
1047
1048 // this check is *only* to see if the forked repo name already exists
1049 // in the user's account.
1050 existingRepo, err := db.GetRepo(
1051 rp.db,
1052 orm.FilterEq("did", user.Did),
1053 orm.FilterEq("name", forkName),
1054 )
1055 if err != nil {
1056 if !errors.Is(err, sql.ErrNoRows) {
1057 l.Error("error fetching existing repo from db", "err", err)
1058 rp.pages.Notice(w, "repo", "Failed to fork this repository. Try again later.")
1059 return
1060 }
1061 } else if existingRepo != nil {
1062 // repo with this name already exists
1063 rp.pages.Notice(w, "repo", "A repository with this name already exists.")
1064 return
1065 }
1066 l = l.With("forkName", forkName)
1067
1068 uri := "https"
1069 if rp.config.Core.Dev {
1070 uri = "http"
1071 }
1072
1073 forkSourceUrl := fmt.Sprintf("%s://%s/%s", uri, f.Knot, f.RepoIdentifier())
1074 l = l.With("cloneUrl", forkSourceUrl)
1075
1076 rkey := tid.TID()
1077
1078 // TODO: this could coordinate better with the knot to receive a clone status
1079 client, err := rp.oauth.ServiceClient(
1080 r,
1081 oauth.WithService(targetKnot),
1082 oauth.WithLxm(tangled.RepoCreateNSID),
1083 oauth.WithDev(rp.config.Core.Dev),
1084 oauth.WithTimeout(time.Second*20),
1085 )
1086 if err != nil {
1087 l.Error("could not create service client", "err", err)
1088 rp.pages.Notice(w, "repo", "Failed to connect to knot server.")
1089 return
1090 }
1091
1092 forkInput := &tangled.RepoCreate_Input{
1093 Rkey: rkey,
1094 Name: forkName,
1095 Source: &forkSourceUrl,
1096 }
1097 createResp, err := tangled.RepoCreate(
1098 r.Context(),
1099 client,
1100 forkInput,
1101 )
1102 if xrpcerr := xrpcclient.HandleXrpcErr(err); xrpcerr != nil {
1103 l.Error("failed to call XRPC repo.create", "xrpcerr", xrpcerr, "err", err)
1104 rp.pages.Notice(w, "repo", xrpcerr.Error())
1105 return
1106 }
1107
1108 var repoDid string
1109 if createResp != nil && createResp.RepoDid != nil {
1110 repoDid = *createResp.RepoDid
1111 }
1112 if repoDid == "" {
1113 l.Error("knot returned empty repo DID for fork")
1114 rp.pages.Notice(w, "repo", "Knot failed to mint a repo DID. The knot may need to be upgraded.")
1115 return
1116 }
1117
1118 forkSource := f.RepoAt().String()
1119 if f.RepoDid != "" {
1120 forkSource = f.RepoDid
1121 }
1122
1123 repo := &models.Repo{
1124 Did: user.Did,
1125 Name: forkName,
1126 Knot: targetKnot,
1127 Rkey: rkey,
1128 Source: forkSource,
1129 Description: f.Description,
1130 Created: time.Now(),
1131 Labels: rp.config.Label.DefaultLabelDefs,
1132 RepoDid: repoDid,
1133 }
1134 record := repo.AsRecord()
1135
1136 cleanupKnot := func() {
1137 go func() {
1138 delays := []time.Duration{0, 2 * time.Second, 5 * time.Second}
1139 for attempt, delay := range delays {
1140 time.Sleep(delay)
1141 deleteClient, dErr := rp.oauth.ServiceClient(
1142 r,
1143 oauth.WithService(targetKnot),
1144 oauth.WithLxm(tangled.RepoDeleteNSID),
1145 oauth.WithDev(rp.config.Core.Dev),
1146 )
1147 if dErr != nil {
1148 l.Error("failed to create delete client for knot cleanup", "attempt", attempt+1, "err", dErr)
1149 continue
1150 }
1151 ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
1152 if dErr := tangled.RepoDelete(ctx, deleteClient, &tangled.RepoDelete_Input{
1153 Did: user.Did,
1154 Name: forkName,
1155 Rkey: rkey,
1156 }); dErr != nil {
1157 cancel()
1158 l.Error("failed to clean up fork on knot after rollback", "attempt", attempt+1, "err", dErr)
1159 continue
1160 }
1161 cancel()
1162 l.Info("successfully cleaned up fork on knot after rollback", "attempt", attempt+1)
1163 return
1164 }
1165 l.Error("exhausted retries for knot cleanup, fork may be orphaned",
1166 "did", user.Did, "fork", forkName, "knot", targetKnot)
1167 }()
1168 }
1169
1170 atpClient, err := rp.oauth.AuthorizedClient(r)
1171 if err != nil {
1172 l.Error("failed to create xrpcclient", "err", err)
1173 cleanupKnot()
1174 rp.pages.Notice(w, "repo", "Failed to fork repository.")
1175 return
1176 }
1177
1178 atresp, err := comatproto.RepoPutRecord(r.Context(), atpClient, &comatproto.RepoPutRecord_Input{
1179 Collection: tangled.RepoNSID,
1180 Repo: user.Did,
1181 Rkey: rkey,
1182 Record: &lexutil.LexiconTypeDecoder{
1183 Val: &record,
1184 },
1185 })
1186 if err != nil {
1187 l.Error("failed to write to PDS", "err", err)
1188 cleanupKnot()
1189 rp.pages.Notice(w, "repo", "Failed to announce repository creation.")
1190 return
1191 }
1192
1193 aturi := atresp.Uri
1194 l = l.With("aturi", aturi)
1195 l.Info("wrote to PDS")
1196
1197 tx, err := rp.db.BeginTx(r.Context(), nil)
1198 if err != nil {
1199 l.Info("txn failed", "err", err)
1200 rp.pages.Notice(w, "repo", "Failed to save repository information.")
1201 return
1202 }
1203
1204 rollback := func() {
1205 err1 := tx.Rollback()
1206 err2 := rp.enforcer.E.LoadPolicy()
1207 err3 := rollbackRecord(context.Background(), aturi, atpClient)
1208
1209 if errors.Is(err1, sql.ErrTxDone) {
1210 err1 = nil
1211 }
1212
1213 if errs := errors.Join(err1, err2, err3); errs != nil {
1214 l.Error("failed to rollback changes", "errs", errs)
1215 }
1216
1217 if aturi != "" {
1218 cleanupKnot()
1219 }
1220 }
1221 defer rollback()
1222
1223 err = db.AddRepo(tx, repo)
1224 if err != nil {
1225 l.Error("failed to AddRepo", "err", err)
1226 rp.pages.Notice(w, "repo", "Failed to save repository information.")
1227 return
1228 }
1229
1230 rbacPath := repo.RepoIdentifier()
1231 err = rp.enforcer.AddRepo(user.Did, targetKnot, rbacPath)
1232 if err != nil {
1233 l.Error("failed to add ACLs", "err", err)
1234 rp.pages.Notice(w, "repo", "Failed to set up repository permissions.")
1235 return
1236 }
1237
1238 err = tx.Commit()
1239 if err != nil {
1240 l.Error("failed to commit changes", "err", err)
1241 http.Error(w, err.Error(), http.StatusInternalServerError)
1242 return
1243 }
1244
1245 err = rp.enforcer.E.SavePolicy()
1246 if err != nil {
1247 l.Error("failed to update ACLs", "err", err)
1248 http.Error(w, err.Error(), http.StatusInternalServerError)
1249 return
1250 }
1251
1252 aturi = ""
1253
1254 rp.notifier.NewRepo(r.Context(), repo)
1255 if repoDid != "" {
1256 rp.pages.HxLocation(w, fmt.Sprintf("/%s", repoDid))
1257 } else {
1258 rp.pages.HxLocation(w, fmt.Sprintf("/%s/%s", user.Did, forkName))
1259 }
1260 }
1261}
1262
1263func (rp *Repo) Stars(w http.ResponseWriter, r *http.Request) {
1264 l := rp.logger.With("handler", "Stars")
1265
1266 user := rp.oauth.GetMultiAccountUser(r)
1267 f, err := rp.repoResolver.Resolve(r)
1268 if err != nil {
1269 l.Error("failed to resolve source repo", "err", err)
1270 return
1271 }
1272
1273 starrers, err := db.GetStars(rp.db, f.RepoAt())
1274 if err != nil {
1275 l.Error("failed to fetch starrers", "err", err, "repoAt", f.RepoAt())
1276 return
1277 }
1278
1279 rp.pages.RepoStars(w, pages.RepoStarsParams{
1280 LoggedInUser: user,
1281 RepoInfo: rp.repoResolver.GetRepoInfo(r, user),
1282 Starrers: starrers,
1283 })
1284}
1285
1286// this is used to rollback changes made to the PDS
1287//
1288// it is a no-op if the provided ATURI is empty
1289func rollbackRecord(ctx context.Context, aturi string, client *atclient.APIClient) error {
1290 if aturi == "" {
1291 return nil
1292 }
1293
1294 parsed := syntax.ATURI(aturi)
1295
1296 collection := parsed.Collection().String()
1297 repo := parsed.Authority().String()
1298 rkey := parsed.RecordKey().String()
1299
1300 _, err := comatproto.RepoDeleteRecord(ctx, client, &comatproto.RepoDeleteRecord_Input{
1301 Collection: collection,
1302 Repo: repo,
1303 Rkey: rkey,
1304 })
1305 return err
1306}
1307
1308func repoCollaboratorRecord(f *models.Repo, subject string, createdAt time.Time) *tangled.RepoCollaborator {
1309 rec := &tangled.RepoCollaborator{
1310 Subject: subject,
1311 CreatedAt: createdAt.Format(time.RFC3339),
1312 }
1313 s := string(f.RepoAt())
1314 rec.Repo = &s
1315 if f.RepoDid != "" {
1316 rec.RepoDid = &f.RepoDid
1317 }
1318 return rec
1319}