This repository has no description
0

Configure Feed

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

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