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