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