This repository has no description
0

Configure Feed

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

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