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