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