Monorepo for Tangled tangled.org
1

Configure Feed

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

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