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