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
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/ogcard"
24 "tangled.org/core/appview/pages"
25 "tangled.org/core/appview/reporesolver"
26 "tangled.org/core/appview/validator"
27 xrpcclient "tangled.org/core/appview/xrpcclient"
28 "tangled.org/core/eventconsumer"
29 "tangled.org/core/idresolver"
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 securejoin "github.com/cyphar/filepath-securejoin"
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 ogcardClient *ogcard.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 ogcardClient: ogcard.NewClient(config.Ogcard.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.Active.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.Active.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.Active.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.Active.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 err = db.SubscribeLabel(tx, &models.RepoLabel{
322 RepoAt: f.RepoAt(),
323 LabelAt: label.AtUri(),
324 })
325
326 err = tx.Commit()
327 if err != nil {
328 fail("Failed to add label.", err)
329 return
330 }
331
332 // clear aturi when everything is successful
333 aturi = ""
334
335 rp.pages.HxRefresh(w)
336}
337
338func (rp *Repo) DeleteLabelDef(w http.ResponseWriter, r *http.Request) {
339 user := rp.oauth.GetMultiAccountUser(r)
340 l := rp.logger.With("handler", "DeleteLabel")
341 l = l.With("did", user.Active.Did)
342
343 f, err := rp.repoResolver.Resolve(r)
344 if err != nil {
345 l.Error("failed to get repo and knot", "err", err)
346 return
347 }
348
349 errorId := "label-operation"
350 fail := func(msg string, err error) {
351 l.Error(msg, "err", err)
352 rp.pages.Notice(w, errorId, msg)
353 }
354
355 // get form values
356 labelId := r.FormValue("label-id")
357
358 label, err := db.GetLabelDefinition(rp.db, orm.FilterEq("id", labelId))
359 if err != nil {
360 fail("Failed to find label definition.", err)
361 return
362 }
363
364 client, err := rp.oauth.AuthorizedClient(r)
365 if err != nil {
366 fail(err.Error(), err)
367 return
368 }
369
370 // delete label record from PDS
371 _, err = comatproto.RepoDeleteRecord(r.Context(), client, &comatproto.RepoDeleteRecord_Input{
372 Collection: tangled.LabelDefinitionNSID,
373 Repo: label.Did,
374 Rkey: label.Rkey,
375 })
376 if err != nil {
377 fail("Failed to delete label record from PDS.", err)
378 return
379 }
380
381 // update repo record to remove the label reference
382 newRepo := *f
383 var updated []string
384 removedAt := label.AtUri().String()
385 for _, l := range newRepo.Labels {
386 if l != removedAt {
387 updated = append(updated, l)
388 }
389 }
390 newRepo.Labels = updated
391 repoRecord := newRepo.AsRecord()
392
393 ex, err := comatproto.RepoGetRecord(r.Context(), client, "", tangled.RepoNSID, newRepo.Did, newRepo.Rkey)
394 if err != nil {
395 fail("Failed to update labels, no record found on PDS.", err)
396 return
397 }
398 _, err = comatproto.RepoPutRecord(r.Context(), client, &comatproto.RepoPutRecord_Input{
399 Collection: tangled.RepoNSID,
400 Repo: newRepo.Did,
401 Rkey: newRepo.Rkey,
402 SwapRecord: ex.Cid,
403 Record: &lexutil.LexiconTypeDecoder{
404 Val: &repoRecord,
405 },
406 })
407 if err != nil {
408 fail("Failed to update repo record.", err)
409 return
410 }
411
412 // transaction for DB changes
413 tx, err := rp.db.BeginTx(r.Context(), nil)
414 if err != nil {
415 fail("Failed to delete label.", err)
416 return
417 }
418 defer tx.Rollback()
419
420 err = db.UnsubscribeLabel(
421 tx,
422 orm.FilterEq("repo_at", f.RepoAt()),
423 orm.FilterEq("label_at", removedAt),
424 )
425 if err != nil {
426 fail("Failed to unsubscribe label.", err)
427 return
428 }
429
430 err = db.DeleteLabelDefinition(tx, orm.FilterEq("id", label.Id))
431 if err != nil {
432 fail("Failed to delete label definition.", err)
433 return
434 }
435
436 err = tx.Commit()
437 if err != nil {
438 fail("Failed to delete label.", err)
439 return
440 }
441
442 // everything succeeded
443 rp.pages.HxRefresh(w)
444}
445
446func (rp *Repo) SubscribeLabel(w http.ResponseWriter, r *http.Request) {
447 user := rp.oauth.GetMultiAccountUser(r)
448 l := rp.logger.With("handler", "SubscribeLabel")
449 l = l.With("did", user.Active.Did)
450
451 f, err := rp.repoResolver.Resolve(r)
452 if err != nil {
453 l.Error("failed to get repo and knot", "err", err)
454 return
455 }
456
457 if err := r.ParseForm(); err != nil {
458 l.Error("invalid form", "err", err)
459 return
460 }
461
462 errorId := "default-label-operation"
463 fail := func(msg string, err error) {
464 l.Error(msg, "err", err)
465 rp.pages.Notice(w, errorId, msg)
466 }
467
468 labelAts := r.Form["label"]
469 _, err = db.GetLabelDefinitions(rp.db, orm.FilterIn("at_uri", labelAts))
470 if err != nil {
471 fail("Failed to subscribe to label.", err)
472 return
473 }
474
475 newRepo := *f
476 newRepo.Labels = append(newRepo.Labels, labelAts...)
477
478 // dedup
479 slices.Sort(newRepo.Labels)
480 newRepo.Labels = slices.Compact(newRepo.Labels)
481
482 repoRecord := newRepo.AsRecord()
483
484 client, err := rp.oauth.AuthorizedClient(r)
485 if err != nil {
486 fail(err.Error(), err)
487 return
488 }
489
490 ex, err := comatproto.RepoGetRecord(r.Context(), client, "", tangled.RepoNSID, f.Did, f.Rkey)
491 if err != nil {
492 fail("Failed to update labels, no record found on PDS.", err)
493 return
494 }
495 _, err = comatproto.RepoPutRecord(r.Context(), client, &comatproto.RepoPutRecord_Input{
496 Collection: tangled.RepoNSID,
497 Repo: newRepo.Did,
498 Rkey: newRepo.Rkey,
499 SwapRecord: ex.Cid,
500 Record: &lexutil.LexiconTypeDecoder{
501 Val: &repoRecord,
502 },
503 })
504
505 tx, err := rp.db.Begin()
506 if err != nil {
507 fail("Failed to subscribe to label.", err)
508 return
509 }
510 defer tx.Rollback()
511
512 for _, l := range labelAts {
513 err = db.SubscribeLabel(tx, &models.RepoLabel{
514 RepoAt: f.RepoAt(),
515 LabelAt: syntax.ATURI(l),
516 })
517 if err != nil {
518 fail("Failed to subscribe to label.", err)
519 return
520 }
521 }
522
523 if err := tx.Commit(); err != nil {
524 fail("Failed to subscribe to label.", err)
525 return
526 }
527
528 // everything succeeded
529 rp.pages.HxRefresh(w)
530}
531
532func (rp *Repo) UnsubscribeLabel(w http.ResponseWriter, r *http.Request) {
533 user := rp.oauth.GetMultiAccountUser(r)
534 l := rp.logger.With("handler", "UnsubscribeLabel")
535 l = l.With("did", user.Active.Did)
536
537 f, err := rp.repoResolver.Resolve(r)
538 if err != nil {
539 l.Error("failed to get repo and knot", "err", err)
540 return
541 }
542
543 if err := r.ParseForm(); err != nil {
544 l.Error("invalid form", "err", err)
545 return
546 }
547
548 errorId := "default-label-operation"
549 fail := func(msg string, err error) {
550 l.Error(msg, "err", err)
551 rp.pages.Notice(w, errorId, msg)
552 }
553
554 labelAts := r.Form["label"]
555 _, err = db.GetLabelDefinitions(rp.db, orm.FilterIn("at_uri", labelAts))
556 if err != nil {
557 fail("Failed to unsubscribe to label.", err)
558 return
559 }
560
561 // update repo record to remove the label reference
562 newRepo := *f
563 var updated []string
564 for _, l := range newRepo.Labels {
565 if !slices.Contains(labelAts, l) {
566 updated = append(updated, l)
567 }
568 }
569 newRepo.Labels = updated
570 repoRecord := newRepo.AsRecord()
571
572 client, err := rp.oauth.AuthorizedClient(r)
573 if err != nil {
574 fail(err.Error(), err)
575 return
576 }
577
578 ex, err := comatproto.RepoGetRecord(r.Context(), client, "", tangled.RepoNSID, f.Did, f.Rkey)
579 if err != nil {
580 fail("Failed to update labels, no record found on PDS.", err)
581 return
582 }
583 _, err = comatproto.RepoPutRecord(r.Context(), client, &comatproto.RepoPutRecord_Input{
584 Collection: tangled.RepoNSID,
585 Repo: newRepo.Did,
586 Rkey: newRepo.Rkey,
587 SwapRecord: ex.Cid,
588 Record: &lexutil.LexiconTypeDecoder{
589 Val: &repoRecord,
590 },
591 })
592
593 err = db.UnsubscribeLabel(
594 rp.db,
595 orm.FilterEq("repo_at", f.RepoAt()),
596 orm.FilterIn("label_at", labelAts),
597 )
598 if err != nil {
599 fail("Failed to unsubscribe label.", err)
600 return
601 }
602
603 // everything succeeded
604 rp.pages.HxRefresh(w)
605}
606
607func (rp *Repo) LabelPanel(w http.ResponseWriter, r *http.Request) {
608 l := rp.logger.With("handler", "LabelPanel")
609
610 f, err := rp.repoResolver.Resolve(r)
611 if err != nil {
612 l.Error("failed to get repo and knot", "err", err)
613 return
614 }
615
616 subjectStr := r.FormValue("subject")
617 subject, err := syntax.ParseATURI(subjectStr)
618 if err != nil {
619 l.Error("failed to get repo and knot", "err", err)
620 return
621 }
622
623 labelDefs, err := db.GetLabelDefinitions(
624 rp.db,
625 orm.FilterIn("at_uri", f.Labels),
626 orm.FilterContains("scope", subject.Collection().String()),
627 )
628 if err != nil {
629 l.Error("failed to fetch label defs", "err", err)
630 return
631 }
632
633 defs := make(map[string]*models.LabelDefinition)
634 for _, l := range labelDefs {
635 defs[l.AtUri().String()] = &l
636 }
637
638 states, err := db.GetLabels(rp.db, orm.FilterEq("subject", subject))
639 if err != nil {
640 l.Error("failed to build label state", "err", err)
641 return
642 }
643 state := states[subject]
644
645 user := rp.oauth.GetMultiAccountUser(r)
646 rp.pages.LabelPanel(w, pages.LabelPanelParams{
647 LoggedInUser: user,
648 RepoInfo: rp.repoResolver.GetRepoInfo(r, user),
649 Defs: defs,
650 Subject: subject.String(),
651 State: state,
652 })
653}
654
655func (rp *Repo) EditLabelPanel(w http.ResponseWriter, r *http.Request) {
656 l := rp.logger.With("handler", "EditLabelPanel")
657
658 f, err := rp.repoResolver.Resolve(r)
659 if err != nil {
660 l.Error("failed to get repo and knot", "err", err)
661 return
662 }
663
664 subjectStr := r.FormValue("subject")
665 subject, err := syntax.ParseATURI(subjectStr)
666 if err != nil {
667 l.Error("failed to get repo and knot", "err", err)
668 return
669 }
670
671 labelDefs, err := db.GetLabelDefinitions(
672 rp.db,
673 orm.FilterIn("at_uri", f.Labels),
674 orm.FilterContains("scope", subject.Collection().String()),
675 )
676 if err != nil {
677 l.Error("failed to fetch labels", "err", err)
678 return
679 }
680
681 defs := make(map[string]*models.LabelDefinition)
682 for _, l := range labelDefs {
683 defs[l.AtUri().String()] = &l
684 }
685
686 states, err := db.GetLabels(rp.db, orm.FilterEq("subject", subject))
687 if err != nil {
688 l.Error("failed to build label state", "err", err)
689 return
690 }
691 state := states[subject]
692
693 user := rp.oauth.GetMultiAccountUser(r)
694 rp.pages.EditLabelPanel(w, pages.EditLabelPanelParams{
695 LoggedInUser: user,
696 RepoInfo: rp.repoResolver.GetRepoInfo(r, user),
697 Defs: defs,
698 Subject: subject.String(),
699 State: state,
700 })
701}
702
703func (rp *Repo) AddCollaborator(w http.ResponseWriter, r *http.Request) {
704 user := rp.oauth.GetMultiAccountUser(r)
705 l := rp.logger.With("handler", "AddCollaborator")
706 l = l.With("did", user.Active.Did)
707
708 f, err := rp.repoResolver.Resolve(r)
709 if err != nil {
710 l.Error("failed to get repo and knot", "err", err)
711 return
712 }
713
714 errorId := "add-collaborator-error"
715 fail := func(msg string, err error) {
716 l.Error(msg, "err", err)
717 rp.pages.Notice(w, errorId, msg)
718 }
719
720 collaborator := r.FormValue("collaborator")
721 if collaborator == "" {
722 fail("Invalid form.", nil)
723 return
724 }
725
726 // remove a single leading `@`, to make @handle work with ResolveIdent
727 collaborator = strings.TrimPrefix(collaborator, "@")
728
729 collaboratorIdent, err := rp.idResolver.ResolveIdent(r.Context(), collaborator)
730 if err != nil {
731 fail(fmt.Sprintf("'%s' is not a valid DID/handle.", collaborator), err)
732 return
733 }
734
735 if collaboratorIdent.DID.String() == user.Active.Did {
736 fail("You seem to be adding yourself as a collaborator.", nil)
737 return
738 }
739 l = l.With("collaborator", collaboratorIdent.Handle)
740 l = l.With("knot", f.Knot)
741
742 // announce this relation into the firehose, store into owners' pds
743 client, err := rp.oauth.AuthorizedClient(r)
744 if err != nil {
745 fail("Failed to write to PDS.", err)
746 return
747 }
748
749 // emit a record
750 currentUser := rp.oauth.GetMultiAccountUser(r)
751 rkey := tid.TID()
752 createdAt := time.Now()
753 resp, err := comatproto.RepoPutRecord(r.Context(), client, &comatproto.RepoPutRecord_Input{
754 Collection: tangled.RepoCollaboratorNSID,
755 Repo: currentUser.Active.Did,
756 Rkey: rkey,
757 Record: &lexutil.LexiconTypeDecoder{
758 Val: &tangled.RepoCollaborator{
759 Subject: collaboratorIdent.DID.String(),
760 Repo: string(f.RepoAt()),
761 CreatedAt: createdAt.Format(time.RFC3339),
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.DidSlashRepo())
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.Active.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.Active.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 err := xrpcclient.HandleXrpcErr(err); err != nil {
883 rp.pages.Notice(w, noticeId, err.Error())
884 return
885 }
886 l.Info("deleted repo from knot")
887
888 tx, err := rp.db.BeginTx(r.Context(), nil)
889 if err != nil {
890 l.Error("failed to start tx")
891 w.Write(fmt.Append(nil, "failed to add collaborator: ", err))
892 return
893 }
894 defer func() {
895 tx.Rollback()
896 err = rp.enforcer.E.LoadPolicy()
897 if err != nil {
898 l.Error("failed to rollback policies")
899 }
900 }()
901
902 // remove collaborator RBAC
903 repoCollaborators, err := rp.enforcer.E.GetImplicitUsersForResourceByDomain(f.DidSlashRepo(), f.Knot)
904 if err != nil {
905 rp.pages.Notice(w, noticeId, "Failed to remove collaborators")
906 return
907 }
908 for _, c := range repoCollaborators {
909 did := c[0]
910 rp.enforcer.RemoveCollaborator(did, f.Knot, f.DidSlashRepo())
911 }
912 l.Info("removed collaborators")
913
914 // remove repo RBAC
915 err = rp.enforcer.RemoveRepo(f.Did, f.Knot, f.DidSlashRepo())
916 if err != nil {
917 rp.pages.Notice(w, noticeId, "Failed to update RBAC rules")
918 return
919 }
920
921 // remove repo from db
922 err = db.RemoveRepo(tx, f.Did, f.Name)
923 if err != nil {
924 rp.pages.Notice(w, noticeId, "Failed to update appview")
925 return
926 }
927 l.Info("removed repo from db")
928
929 err = tx.Commit()
930 if err != nil {
931 l.Error("failed to commit changes", "err", err)
932 http.Error(w, err.Error(), http.StatusInternalServerError)
933 return
934 }
935
936 err = rp.enforcer.E.SavePolicy()
937 if err != nil {
938 l.Error("failed to update ACLs", "err", err)
939 http.Error(w, err.Error(), http.StatusInternalServerError)
940 return
941 }
942
943 rp.pages.HxRedirect(w, fmt.Sprintf("/%s", f.Did))
944}
945
946func (rp *Repo) SyncRepoFork(w http.ResponseWriter, r *http.Request) {
947 l := rp.logger.With("handler", "SyncRepoFork")
948
949 ref := chi.URLParam(r, "ref")
950 ref, _ = url.PathUnescape(ref)
951
952 user := rp.oauth.GetMultiAccountUser(r)
953 f, err := rp.repoResolver.Resolve(r)
954 if err != nil {
955 l.Error("failed to resolve source repo", "err", err)
956 return
957 }
958
959 switch r.Method {
960 case http.MethodPost:
961 client, err := rp.oauth.ServiceClient(
962 r,
963 oauth.WithService(f.Knot),
964 oauth.WithLxm(tangled.RepoForkSyncNSID),
965 oauth.WithDev(rp.config.Core.Dev),
966 )
967 if err != nil {
968 rp.pages.Notice(w, "repo", "Failed to connect to knot server.")
969 return
970 }
971
972 if f.Source == "" {
973 rp.pages.Notice(w, "repo", "This repository is not a fork.")
974 return
975 }
976
977 err = tangled.RepoForkSync(
978 r.Context(),
979 client,
980 &tangled.RepoForkSync_Input{
981 Did: user.Active.Did,
982 Name: f.Name,
983 Source: f.Source,
984 Branch: ref,
985 },
986 )
987 if err := xrpcclient.HandleXrpcErr(err); err != nil {
988 rp.pages.Notice(w, "repo", err.Error())
989 return
990 }
991
992 rp.pages.HxRefresh(w)
993 return
994 }
995}
996
997func (rp *Repo) ForkRepo(w http.ResponseWriter, r *http.Request) {
998 l := rp.logger.With("handler", "ForkRepo")
999
1000 user := rp.oauth.GetMultiAccountUser(r)
1001 f, err := rp.repoResolver.Resolve(r)
1002 if err != nil {
1003 l.Error("failed to resolve source repo", "err", err)
1004 return
1005 }
1006
1007 switch r.Method {
1008 case http.MethodGet:
1009 user := rp.oauth.GetMultiAccountUser(r)
1010 knots, err := rp.enforcer.GetKnotsForUser(user.Active.Did)
1011 if err != nil {
1012 rp.pages.Notice(w, "repo", "Invalid user account.")
1013 return
1014 }
1015
1016 rp.pages.ForkRepo(w, pages.ForkRepoParams{
1017 LoggedInUser: user,
1018 Knots: knots,
1019 RepoInfo: rp.repoResolver.GetRepoInfo(r, user),
1020 })
1021
1022 case http.MethodPost:
1023 l := rp.logger.With("handler", "ForkRepo")
1024
1025 targetKnot := r.FormValue("knot")
1026 if targetKnot == "" {
1027 rp.pages.Notice(w, "repo", "Invalid form submission—missing knot domain.")
1028 return
1029 }
1030 l = l.With("targetKnot", targetKnot)
1031
1032 ok, err := rp.enforcer.E.Enforce(user.Active.Did, targetKnot, targetKnot, "repo:create")
1033 if err != nil || !ok {
1034 rp.pages.Notice(w, "repo", "You do not have permission to create a repo in this knot.")
1035 return
1036 }
1037
1038 // choose a name for a fork
1039 forkName := r.FormValue("repo_name")
1040 if forkName == "" {
1041 rp.pages.Notice(w, "repo", "Repository name cannot be empty.")
1042 return
1043 }
1044
1045 // this check is *only* to see if the forked repo name already exists
1046 // in the user's account.
1047 existingRepo, err := db.GetRepo(
1048 rp.db,
1049 orm.FilterEq("did", user.Active.Did),
1050 orm.FilterEq("name", forkName),
1051 )
1052 if err != nil {
1053 if !errors.Is(err, sql.ErrNoRows) {
1054 l.Error("error fetching existing repo from db", "err", err)
1055 rp.pages.Notice(w, "repo", "Failed to fork this repository. Try again later.")
1056 return
1057 }
1058 } else if existingRepo != nil {
1059 // repo with this name already exists
1060 rp.pages.Notice(w, "repo", "A repository with this name already exists.")
1061 return
1062 }
1063 l = l.With("forkName", forkName)
1064
1065 uri := "https"
1066 if rp.config.Core.Dev {
1067 uri = "http"
1068 }
1069
1070 forkSourceUrl := fmt.Sprintf("%s://%s/%s/%s", uri, f.Knot, f.Did, f.Name)
1071 l = l.With("cloneUrl", forkSourceUrl)
1072
1073 sourceAt := f.RepoAt().String()
1074
1075 // create an atproto record for this fork
1076 rkey := tid.TID()
1077 repo := &models.Repo{
1078 Did: user.Active.Did,
1079 Name: forkName,
1080 Knot: targetKnot,
1081 Rkey: rkey,
1082 Source: sourceAt,
1083 Description: f.Description,
1084 Created: time.Now(),
1085 Labels: rp.config.Label.DefaultLabelDefs,
1086 }
1087 record := repo.AsRecord()
1088
1089 atpClient, err := rp.oauth.AuthorizedClient(r)
1090 if err != nil {
1091 l.Error("failed to create xrpcclient", "err", err)
1092 rp.pages.Notice(w, "repo", "Failed to fork repository.")
1093 return
1094 }
1095
1096 atresp, err := comatproto.RepoPutRecord(r.Context(), atpClient, &comatproto.RepoPutRecord_Input{
1097 Collection: tangled.RepoNSID,
1098 Repo: user.Active.Did,
1099 Rkey: rkey,
1100 Record: &lexutil.LexiconTypeDecoder{
1101 Val: &record,
1102 },
1103 })
1104 if err != nil {
1105 l.Error("failed to write to PDS", "err", err)
1106 rp.pages.Notice(w, "repo", "Failed to announce repository creation.")
1107 return
1108 }
1109
1110 aturi := atresp.Uri
1111 l = l.With("aturi", aturi)
1112 l.Info("wrote to PDS")
1113
1114 tx, err := rp.db.BeginTx(r.Context(), nil)
1115 if err != nil {
1116 l.Info("txn failed", "err", err)
1117 rp.pages.Notice(w, "repo", "Failed to save repository information.")
1118 return
1119 }
1120
1121 // The rollback function reverts a few things on failure:
1122 // - the pending txn
1123 // - the ACLs
1124 // - the atproto record created
1125 rollback := func() {
1126 err1 := tx.Rollback()
1127 err2 := rp.enforcer.E.LoadPolicy()
1128 err3 := rollbackRecord(context.Background(), aturi, atpClient)
1129
1130 // ignore txn complete errors, this is okay
1131 if errors.Is(err1, sql.ErrTxDone) {
1132 err1 = nil
1133 }
1134
1135 if errs := errors.Join(err1, err2, err3); errs != nil {
1136 l.Error("failed to rollback changes", "errs", errs)
1137 return
1138 }
1139 }
1140 defer rollback()
1141
1142 // TODO: this could coordinate better with the knot to recieve a clone status
1143 client, err := rp.oauth.ServiceClient(
1144 r,
1145 oauth.WithService(targetKnot),
1146 oauth.WithLxm(tangled.RepoCreateNSID),
1147 oauth.WithDev(rp.config.Core.Dev),
1148 oauth.WithTimeout(time.Second*20), // big repos take time to clone
1149 )
1150 if err != nil {
1151 l.Error("could not create service client", "err", err)
1152 rp.pages.Notice(w, "repo", "Failed to connect to knot server.")
1153 return
1154 }
1155
1156 err = tangled.RepoCreate(
1157 r.Context(),
1158 client,
1159 &tangled.RepoCreate_Input{
1160 Rkey: rkey,
1161 Source: &forkSourceUrl,
1162 },
1163 )
1164 if err := xrpcclient.HandleXrpcErr(err); err != nil {
1165 rp.pages.Notice(w, "repo", err.Error())
1166 return
1167 }
1168
1169 err = db.AddRepo(tx, repo)
1170 if err != nil {
1171 l.Error("failed to AddRepo", "err", err)
1172 rp.pages.Notice(w, "repo", "Failed to save repository information.")
1173 return
1174 }
1175
1176 // acls
1177 p, _ := securejoin.SecureJoin(user.Active.Did, forkName)
1178 err = rp.enforcer.AddRepo(user.Active.Did, targetKnot, p)
1179 if err != nil {
1180 l.Error("failed to add ACLs", "err", err)
1181 rp.pages.Notice(w, "repo", "Failed to set up repository permissions.")
1182 return
1183 }
1184
1185 err = tx.Commit()
1186 if err != nil {
1187 l.Error("failed to commit changes", "err", err)
1188 http.Error(w, err.Error(), http.StatusInternalServerError)
1189 return
1190 }
1191
1192 err = rp.enforcer.E.SavePolicy()
1193 if err != nil {
1194 l.Error("failed to update ACLs", "err", err)
1195 http.Error(w, err.Error(), http.StatusInternalServerError)
1196 return
1197 }
1198
1199 // reset the ATURI because the transaction completed successfully
1200 aturi = ""
1201
1202 rp.notifier.NewRepo(r.Context(), repo)
1203 rp.pages.HxLocation(w, fmt.Sprintf("/%s/%s", user.Active.Did, forkName))
1204 }
1205}
1206
1207// this is used to rollback changes made to the PDS
1208//
1209// it is a no-op if the provided ATURI is empty
1210func rollbackRecord(ctx context.Context, aturi string, client *atclient.APIClient) error {
1211 if aturi == "" {
1212 return nil
1213 }
1214
1215 parsed := syntax.ATURI(aturi)
1216
1217 collection := parsed.Collection().String()
1218 repo := parsed.Authority().String()
1219 rkey := parsed.RecordKey().String()
1220
1221 _, err := comatproto.RepoDeleteRecord(ctx, client, &comatproto.RepoDeleteRecord_Input{
1222 Collection: collection,
1223 Repo: repo,
1224 Rkey: rkey,
1225 })
1226 return err
1227}