Monorepo for Tangled tangled.org
1

Configure Feed

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

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