This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

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