Monorepo for Tangled
0

Configure Feed

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

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 Repo: f.RepoDidPtr(), 1387 Source: f.Source, 1388 Branch: ref, 1389 }, 1390 ) 1391 if xrpcerr := xrpcclient.HandleXrpcErr(err); xrpcerr != nil { 1392 l.Error("failed to call XRPC repo.forkSync", "xrpcerr", xrpcerr, "err", err) 1393 rp.pages.Notice(w, "repo", err.Error()) 1394 return 1395 } 1396 1397 rp.pages.HxRefresh(w) 1398 return 1399 } 1400} 1401 1402func (rp *Repo) ForkRepo(w http.ResponseWriter, r *http.Request) { 1403 l := rp.logger.With("handler", "ForkRepo") 1404 1405 user := rp.oauth.GetMultiAccountUser(r) 1406 f, err := rp.repoResolver.Resolve(r) 1407 if err != nil { 1408 l.Error("failed to resolve source repo", "err", err) 1409 return 1410 } 1411 1412 switch r.Method { 1413 case http.MethodGet: 1414 user := rp.oauth.GetMultiAccountUser(r) 1415 knots := rp.acl.KnotsForUser(r.Context(), user.Did) 1416 1417 rp.pages.ForkRepo(w, pages.ForkRepoParams{ 1418 BaseParams: pages.BaseParamsFromContext(r.Context()), 1419 Knots: knots, 1420 RepoInfo: rp.repoResolver.GetRepoInfo(r, user), 1421 }) 1422 1423 case http.MethodPost: 1424 l := rp.logger.With("handler", "ForkRepo") 1425 1426 targetKnot := r.FormValue("knot") 1427 if targetKnot == "" { 1428 rp.pages.Notice(w, "repo", "Invalid form submission&mdash;missing knot domain.") 1429 return 1430 } 1431 l = l.With("targetKnot", targetKnot) 1432 1433 if !rp.acl.IsRepoCreateAllowed(r.Context(), targetKnot, user.Did) { 1434 rp.pages.Notice(w, "repo", "You do not have permission to create a repo in this knot.") 1435 return 1436 } 1437 1438 // choose a name for a fork 1439 forkName := strings.ToLower(r.FormValue("repo_name")) 1440 if forkName == "" { 1441 rp.pages.Notice(w, "repo", "Repository name cannot be empty.") 1442 return 1443 } 1444 1445 // this check is *only* to see if the forked repo name already exists 1446 // in the user's account. 1447 existingRepo, err := db.GetRepo( 1448 rp.db, 1449 orm.FilterEq("did", user.Did), 1450 orm.FilterEq("name", forkName), 1451 ) 1452 if err != nil { 1453 if !errors.Is(err, sql.ErrNoRows) { 1454 l.Error("error fetching existing repo from db", "err", err) 1455 rp.pages.Notice(w, "repo", "Failed to fork this repository. Try again later.") 1456 return 1457 } 1458 } else if existingRepo != nil { 1459 // repo with this name already exists 1460 rp.pages.Notice(w, "repo", "A repository with this name already exists.") 1461 return 1462 } 1463 l = l.With("forkName", forkName) 1464 1465 uri := "https" 1466 if rp.config.Core.Dev { 1467 uri = "http" 1468 } 1469 1470 forkSourceUrl := fmt.Sprintf("%s://%s/%s", uri, f.Knot, f.RepoIdentifier()) 1471 l = l.With("cloneUrl", forkSourceUrl) 1472 1473 rkey := strings.ToLower(forkName) 1474 1475 // TODO: this could coordinate better with the knot to receive a clone status 1476 client, err := rp.oauth.ServiceClient( 1477 r, 1478 oauth.WithService(targetKnot), 1479 oauth.WithLxm(tangled.RepoCreateNSID), 1480 oauth.WithDev(rp.config.Core.Dev), 1481 oauth.WithTimeout(time.Second*20), 1482 ) 1483 if err != nil { 1484 l.Error("could not create service client", "err", err) 1485 rp.pages.Notice(w, "repo", "Failed to connect to knot server.") 1486 return 1487 } 1488 1489 forkInput := &tangled.RepoCreate_Input{ 1490 Rkey: rkey, 1491 Name: rkey, 1492 Source: &forkSourceUrl, 1493 } 1494 createResp, err := tangled.RepoCreate( 1495 r.Context(), 1496 client, 1497 forkInput, 1498 ) 1499 if xrpcerr := xrpcclient.HandleXrpcErr(err); xrpcerr != nil { 1500 l.Error("failed to call XRPC repo.create", "xrpcerr", xrpcerr, "err", err) 1501 rp.pages.Notice(w, "repo", xrpcerr.Error()) 1502 return 1503 } 1504 1505 var repoDid string 1506 if createResp != nil && createResp.RepoDid != nil { 1507 repoDid = *createResp.RepoDid 1508 } 1509 if repoDid == "" { 1510 l.Error("knot returned empty repo DID for fork") 1511 rp.pages.Notice(w, "repo", "Knot failed to mint a repo DID. The knot may need to be upgraded.") 1512 return 1513 } 1514 1515 forkSource := f.RepoAt().String() 1516 if f.RepoDid != "" { 1517 forkSource = f.RepoDid 1518 } 1519 1520 forkDescription := r.Form.Get("description") 1521 1522 repo := &models.Repo{ 1523 Did: user.Did, 1524 Name: rkey, 1525 Knot: targetKnot, 1526 Rkey: rkey, 1527 Source: forkSource, 1528 Description: forkDescription, 1529 Created: time.Now(), 1530 Labels: rp.config.Label.DefaultLabelDefs, 1531 RepoDid: repoDid, 1532 } 1533 record := repo.AsRecord() 1534 1535 cleanupKnot := func() { 1536 go func() { 1537 delays := []time.Duration{0, 2 * time.Second, 5 * time.Second} 1538 for attempt, delay := range delays { 1539 time.Sleep(delay) 1540 deleteClient, dErr := rp.oauth.ServiceClient( 1541 r, 1542 oauth.WithService(targetKnot), 1543 oauth.WithLxm(tangled.RepoDeleteNSID), 1544 oauth.WithDev(rp.config.Core.Dev), 1545 ) 1546 if dErr != nil { 1547 l.Error("failed to create delete client for knot cleanup", "attempt", attempt+1, "err", dErr) 1548 continue 1549 } 1550 ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) 1551 if dErr := tangled.RepoDelete(ctx, deleteClient, &tangled.RepoDelete_Input{ 1552 Did: user.Did, 1553 Name: forkName, 1554 Rkey: rkey, 1555 }); dErr != nil { 1556 cancel() 1557 l.Error("failed to clean up fork on knot after rollback", "attempt", attempt+1, "err", dErr) 1558 continue 1559 } 1560 cancel() 1561 l.Info("successfully cleaned up fork on knot after rollback", "attempt", attempt+1) 1562 return 1563 } 1564 l.Error("exhausted retries for knot cleanup, fork may be orphaned", 1565 "did", user.Did, "fork", forkName, "knot", targetKnot) 1566 }() 1567 } 1568 1569 atpClient, err := rp.oauth.AuthorizedClient(r) 1570 if err != nil { 1571 l.Error("failed to create xrpcclient", "err", err) 1572 cleanupKnot() 1573 rp.pages.Notice(w, "repo", "Failed to fork repository.") 1574 return 1575 } 1576 1577 atresp, err := comatproto.RepoPutRecord(r.Context(), atpClient, &comatproto.RepoPutRecord_Input{ 1578 Collection: tangled.RepoNSID, 1579 Repo: user.Did, 1580 Rkey: rkey, 1581 Record: &lexutil.LexiconTypeDecoder{ 1582 Val: &record, 1583 }, 1584 }) 1585 if err != nil { 1586 l.Error("failed to write to PDS", "err", err) 1587 cleanupKnot() 1588 rp.pages.Notice(w, "repo", "Failed to announce repository creation.") 1589 return 1590 } 1591 1592 aturi := atresp.Uri 1593 l = l.With("aturi", aturi) 1594 l.Info("wrote to PDS") 1595 1596 tx, err := rp.db.BeginTx(r.Context(), nil) 1597 if err != nil { 1598 l.Info("txn failed", "err", err) 1599 rp.pages.Notice(w, "repo", "Failed to save repository information.") 1600 return 1601 } 1602 1603 rollback := func() { 1604 err1 := tx.Rollback() 1605 err2 := rp.enforcer.E.LoadPolicy() 1606 err3 := rollbackRecord(context.Background(), aturi, atpClient) 1607 1608 if errors.Is(err1, sql.ErrTxDone) { 1609 err1 = nil 1610 } 1611 1612 if errs := errors.Join(err1, err2, err3); errs != nil { 1613 l.Error("failed to rollback changes", "errs", errs) 1614 } 1615 1616 if aturi != "" { 1617 cleanupKnot() 1618 } 1619 } 1620 defer rollback() 1621 1622 err = db.AddRepo(tx, repo) 1623 if err != nil { 1624 l.Error("failed to AddRepo", "err", err) 1625 rp.pages.Notice(w, "repo", "Failed to save repository information.") 1626 return 1627 } 1628 1629 rbacPath := repo.RepoIdentifier() 1630 err = rp.enforcer.AddRepo(user.Did, targetKnot, rbacPath) 1631 if err != nil { 1632 l.Error("failed to add ACLs", "err", err) 1633 rp.pages.Notice(w, "repo", "Failed to set up repository permissions.") 1634 return 1635 } 1636 1637 err = tx.Commit() 1638 if err != nil { 1639 l.Error("failed to commit changes", "err", err) 1640 http.Error(w, err.Error(), http.StatusInternalServerError) 1641 return 1642 } 1643 1644 err = rp.enforcer.E.SavePolicy() 1645 if err != nil { 1646 l.Error("failed to update ACLs", "err", err) 1647 http.Error(w, err.Error(), http.StatusInternalServerError) 1648 return 1649 } 1650 1651 aturi = "" 1652 1653 rp.notifier.NewRepo(r.Context(), repo) 1654 if repoDid != "" { 1655 rp.pages.HxLocation(w, fmt.Sprintf("/%s", repoDid)) 1656 } else { 1657 rp.pages.HxLocation(w, fmt.Sprintf("/%s/%s", user.Did, forkName)) 1658 } 1659 } 1660} 1661 1662func (rp *Repo) Stars(w http.ResponseWriter, r *http.Request) { 1663 l := rp.logger.With("handler", "Stars") 1664 1665 user := rp.oauth.GetMultiAccountUser(r) 1666 f, err := rp.repoResolver.Resolve(r) 1667 if err != nil { 1668 l.Error("failed to resolve source repo", "err", err) 1669 return 1670 } 1671 1672 page := pagination.FromContext(r.Context()) 1673 if page.Limit > 30 || page.Limit <= 0 { 1674 page.Limit = 30 1675 } 1676 1677 starrers, err := db.GetStars(rp.db, string(f.RepoDid), page) 1678 if err != nil { 1679 l.Error("failed to fetch starrers", "err", err, "repoDid", f.RepoDid) 1680 return 1681 } 1682 1683 totalCount, err := db.GetStarCount(rp.db, models.StarSubjectRepo, string(f.RepoDid)) 1684 if err != nil { 1685 l.Error("failed to fetch star count", "err", err, "repoDid", f.RepoDid) 1686 return 1687 } 1688 1689 rp.pages.RepoStars(w, pages.RepoStarsParams{ 1690 BaseParams: pages.BaseParamsFromContext(r.Context()), 1691 RepoInfo: rp.repoResolver.GetRepoInfo(r, user), 1692 Starrers: starrers, 1693 Page: page, 1694 TotalCount: totalCount, 1695 }) 1696} 1697 1698func (rp *Repo) Forks(w http.ResponseWriter, r *http.Request) { 1699 l := rp.logger.With("handler", "Forks") 1700 1701 user := rp.oauth.GetMultiAccountUser(r) 1702 f, err := rp.repoResolver.Resolve(r) 1703 if err != nil { 1704 l.Error("failed to resolve source repo", "err", err) 1705 return 1706 } 1707 1708 var forks []models.Repo 1709 totalCount := 0 1710 page := pagination.FromContext(r.Context()) 1711 if f.RepoDid != "" { 1712 forks, err = db.GetReposPaginated(rp.db, page, orm.FilterEq("source", f.RepoDid)) 1713 if err != nil { 1714 l.Error("failed to fetch forks", "err", err, "repoAt", f.RepoAt()) 1715 return 1716 } 1717 1718 totalCount, err = db.GetForkCount(rp.db, f.RepoDid) 1719 if err != nil { 1720 l.Error("failed to fetch fork count", "err", err, "repoAt", f.RepoAt()) 1721 return 1722 } 1723 } 1724 1725 err = rp.pages.RepoForks(w, pages.RepoForksParams{ 1726 BaseParams: pages.BaseParamsFromContext(r.Context()), 1727 RepoInfo: rp.repoResolver.GetRepoInfo(r, user), 1728 Forks: forks, 1729 Page: page, 1730 TotalCount: totalCount, 1731 }) 1732 if err != nil { 1733 l.Error("failed to render page", "err", err) 1734 } 1735} 1736 1737// this is used to rollback changes made to the PDS 1738// 1739// it is a no-op if the provided ATURI is empty 1740func rollbackRecord(ctx context.Context, aturi string, client *atclient.APIClient) error { 1741 if aturi == "" { 1742 return nil 1743 } 1744 1745 parsed := syntax.ATURI(aturi) 1746 1747 collection := parsed.Collection().String() 1748 repo := parsed.Authority().String() 1749 rkey := parsed.RecordKey().String() 1750 1751 _, err := comatproto.RepoDeleteRecord(ctx, client, &comatproto.RepoDeleteRecord_Input{ 1752 Collection: collection, 1753 Repo: repo, 1754 Rkey: rkey, 1755 }) 1756 return err 1757} 1758 1759func repoCollaboratorRecord(f *models.Repo, subject string, createdAt time.Time) *tangled.RepoCollaborator { 1760 return &tangled.RepoCollaborator{ 1761 Subject: subject, 1762 CreatedAt: createdAt.Format(time.RFC3339), 1763 Repo: f.RepoDid, 1764 } 1765}