Monorepo for Tangled tangled.org
1

Configure Feed

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

core / appview / repo / repo.go
30 kB 1227 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/ogcard" 24 "tangled.org/core/appview/pages" 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/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 securejoin "github.com/cyphar/filepath-securejoin" 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 ogcardClient *ogcard.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 ogcardClient: ogcard.NewClient(config.Ogcard.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 err = db.SubscribeLabel(tx, &models.RepoLabel{ 322 RepoAt: f.RepoAt(), 323 LabelAt: label.AtUri(), 324 }) 325 326 err = tx.Commit() 327 if err != nil { 328 fail("Failed to add label.", err) 329 return 330 } 331 332 // clear aturi when everything is successful 333 aturi = "" 334 335 rp.pages.HxRefresh(w) 336} 337 338func (rp *Repo) DeleteLabelDef(w http.ResponseWriter, r *http.Request) { 339 user := rp.oauth.GetMultiAccountUser(r) 340 l := rp.logger.With("handler", "DeleteLabel") 341 l = l.With("did", user.Active.Did) 342 343 f, err := rp.repoResolver.Resolve(r) 344 if err != nil { 345 l.Error("failed to get repo and knot", "err", err) 346 return 347 } 348 349 errorId := "label-operation" 350 fail := func(msg string, err error) { 351 l.Error(msg, "err", err) 352 rp.pages.Notice(w, errorId, msg) 353 } 354 355 // get form values 356 labelId := r.FormValue("label-id") 357 358 label, err := db.GetLabelDefinition(rp.db, orm.FilterEq("id", labelId)) 359 if err != nil { 360 fail("Failed to find label definition.", err) 361 return 362 } 363 364 client, err := rp.oauth.AuthorizedClient(r) 365 if err != nil { 366 fail(err.Error(), err) 367 return 368 } 369 370 // delete label record from PDS 371 _, err = comatproto.RepoDeleteRecord(r.Context(), client, &comatproto.RepoDeleteRecord_Input{ 372 Collection: tangled.LabelDefinitionNSID, 373 Repo: label.Did, 374 Rkey: label.Rkey, 375 }) 376 if err != nil { 377 fail("Failed to delete label record from PDS.", err) 378 return 379 } 380 381 // update repo record to remove the label reference 382 newRepo := *f 383 var updated []string 384 removedAt := label.AtUri().String() 385 for _, l := range newRepo.Labels { 386 if l != removedAt { 387 updated = append(updated, l) 388 } 389 } 390 newRepo.Labels = updated 391 repoRecord := newRepo.AsRecord() 392 393 ex, err := comatproto.RepoGetRecord(r.Context(), client, "", tangled.RepoNSID, newRepo.Did, newRepo.Rkey) 394 if err != nil { 395 fail("Failed to update labels, no record found on PDS.", err) 396 return 397 } 398 _, err = comatproto.RepoPutRecord(r.Context(), client, &comatproto.RepoPutRecord_Input{ 399 Collection: tangled.RepoNSID, 400 Repo: newRepo.Did, 401 Rkey: newRepo.Rkey, 402 SwapRecord: ex.Cid, 403 Record: &lexutil.LexiconTypeDecoder{ 404 Val: &repoRecord, 405 }, 406 }) 407 if err != nil { 408 fail("Failed to update repo record.", err) 409 return 410 } 411 412 // transaction for DB changes 413 tx, err := rp.db.BeginTx(r.Context(), nil) 414 if err != nil { 415 fail("Failed to delete label.", err) 416 return 417 } 418 defer tx.Rollback() 419 420 err = db.UnsubscribeLabel( 421 tx, 422 orm.FilterEq("repo_at", f.RepoAt()), 423 orm.FilterEq("label_at", removedAt), 424 ) 425 if err != nil { 426 fail("Failed to unsubscribe label.", err) 427 return 428 } 429 430 err = db.DeleteLabelDefinition(tx, orm.FilterEq("id", label.Id)) 431 if err != nil { 432 fail("Failed to delete label definition.", err) 433 return 434 } 435 436 err = tx.Commit() 437 if err != nil { 438 fail("Failed to delete label.", err) 439 return 440 } 441 442 // everything succeeded 443 rp.pages.HxRefresh(w) 444} 445 446func (rp *Repo) SubscribeLabel(w http.ResponseWriter, r *http.Request) { 447 user := rp.oauth.GetMultiAccountUser(r) 448 l := rp.logger.With("handler", "SubscribeLabel") 449 l = l.With("did", user.Active.Did) 450 451 f, err := rp.repoResolver.Resolve(r) 452 if err != nil { 453 l.Error("failed to get repo and knot", "err", err) 454 return 455 } 456 457 if err := r.ParseForm(); err != nil { 458 l.Error("invalid form", "err", err) 459 return 460 } 461 462 errorId := "default-label-operation" 463 fail := func(msg string, err error) { 464 l.Error(msg, "err", err) 465 rp.pages.Notice(w, errorId, msg) 466 } 467 468 labelAts := r.Form["label"] 469 _, err = db.GetLabelDefinitions(rp.db, orm.FilterIn("at_uri", labelAts)) 470 if err != nil { 471 fail("Failed to subscribe to label.", err) 472 return 473 } 474 475 newRepo := *f 476 newRepo.Labels = append(newRepo.Labels, labelAts...) 477 478 // dedup 479 slices.Sort(newRepo.Labels) 480 newRepo.Labels = slices.Compact(newRepo.Labels) 481 482 repoRecord := newRepo.AsRecord() 483 484 client, err := rp.oauth.AuthorizedClient(r) 485 if err != nil { 486 fail(err.Error(), err) 487 return 488 } 489 490 ex, err := comatproto.RepoGetRecord(r.Context(), client, "", tangled.RepoNSID, f.Did, f.Rkey) 491 if err != nil { 492 fail("Failed to update labels, no record found on PDS.", err) 493 return 494 } 495 _, err = comatproto.RepoPutRecord(r.Context(), client, &comatproto.RepoPutRecord_Input{ 496 Collection: tangled.RepoNSID, 497 Repo: newRepo.Did, 498 Rkey: newRepo.Rkey, 499 SwapRecord: ex.Cid, 500 Record: &lexutil.LexiconTypeDecoder{ 501 Val: &repoRecord, 502 }, 503 }) 504 505 tx, err := rp.db.Begin() 506 if err != nil { 507 fail("Failed to subscribe to label.", err) 508 return 509 } 510 defer tx.Rollback() 511 512 for _, l := range labelAts { 513 err = db.SubscribeLabel(tx, &models.RepoLabel{ 514 RepoAt: f.RepoAt(), 515 LabelAt: syntax.ATURI(l), 516 }) 517 if err != nil { 518 fail("Failed to subscribe to label.", err) 519 return 520 } 521 } 522 523 if err := tx.Commit(); err != nil { 524 fail("Failed to subscribe to label.", err) 525 return 526 } 527 528 // everything succeeded 529 rp.pages.HxRefresh(w) 530} 531 532func (rp *Repo) UnsubscribeLabel(w http.ResponseWriter, r *http.Request) { 533 user := rp.oauth.GetMultiAccountUser(r) 534 l := rp.logger.With("handler", "UnsubscribeLabel") 535 l = l.With("did", user.Active.Did) 536 537 f, err := rp.repoResolver.Resolve(r) 538 if err != nil { 539 l.Error("failed to get repo and knot", "err", err) 540 return 541 } 542 543 if err := r.ParseForm(); err != nil { 544 l.Error("invalid form", "err", err) 545 return 546 } 547 548 errorId := "default-label-operation" 549 fail := func(msg string, err error) { 550 l.Error(msg, "err", err) 551 rp.pages.Notice(w, errorId, msg) 552 } 553 554 labelAts := r.Form["label"] 555 _, err = db.GetLabelDefinitions(rp.db, orm.FilterIn("at_uri", labelAts)) 556 if err != nil { 557 fail("Failed to unsubscribe to label.", err) 558 return 559 } 560 561 // update repo record to remove the label reference 562 newRepo := *f 563 var updated []string 564 for _, l := range newRepo.Labels { 565 if !slices.Contains(labelAts, l) { 566 updated = append(updated, l) 567 } 568 } 569 newRepo.Labels = updated 570 repoRecord := newRepo.AsRecord() 571 572 client, err := rp.oauth.AuthorizedClient(r) 573 if err != nil { 574 fail(err.Error(), err) 575 return 576 } 577 578 ex, err := comatproto.RepoGetRecord(r.Context(), client, "", tangled.RepoNSID, f.Did, f.Rkey) 579 if err != nil { 580 fail("Failed to update labels, no record found on PDS.", err) 581 return 582 } 583 _, err = comatproto.RepoPutRecord(r.Context(), client, &comatproto.RepoPutRecord_Input{ 584 Collection: tangled.RepoNSID, 585 Repo: newRepo.Did, 586 Rkey: newRepo.Rkey, 587 SwapRecord: ex.Cid, 588 Record: &lexutil.LexiconTypeDecoder{ 589 Val: &repoRecord, 590 }, 591 }) 592 593 err = db.UnsubscribeLabel( 594 rp.db, 595 orm.FilterEq("repo_at", f.RepoAt()), 596 orm.FilterIn("label_at", labelAts), 597 ) 598 if err != nil { 599 fail("Failed to unsubscribe label.", err) 600 return 601 } 602 603 // everything succeeded 604 rp.pages.HxRefresh(w) 605} 606 607func (rp *Repo) LabelPanel(w http.ResponseWriter, r *http.Request) { 608 l := rp.logger.With("handler", "LabelPanel") 609 610 f, err := rp.repoResolver.Resolve(r) 611 if err != nil { 612 l.Error("failed to get repo and knot", "err", err) 613 return 614 } 615 616 subjectStr := r.FormValue("subject") 617 subject, err := syntax.ParseATURI(subjectStr) 618 if err != nil { 619 l.Error("failed to get repo and knot", "err", err) 620 return 621 } 622 623 labelDefs, err := db.GetLabelDefinitions( 624 rp.db, 625 orm.FilterIn("at_uri", f.Labels), 626 orm.FilterContains("scope", subject.Collection().String()), 627 ) 628 if err != nil { 629 l.Error("failed to fetch label defs", "err", err) 630 return 631 } 632 633 defs := make(map[string]*models.LabelDefinition) 634 for _, l := range labelDefs { 635 defs[l.AtUri().String()] = &l 636 } 637 638 states, err := db.GetLabels(rp.db, orm.FilterEq("subject", subject)) 639 if err != nil { 640 l.Error("failed to build label state", "err", err) 641 return 642 } 643 state := states[subject] 644 645 user := rp.oauth.GetMultiAccountUser(r) 646 rp.pages.LabelPanel(w, pages.LabelPanelParams{ 647 LoggedInUser: user, 648 RepoInfo: rp.repoResolver.GetRepoInfo(r, user), 649 Defs: defs, 650 Subject: subject.String(), 651 State: state, 652 }) 653} 654 655func (rp *Repo) EditLabelPanel(w http.ResponseWriter, r *http.Request) { 656 l := rp.logger.With("handler", "EditLabelPanel") 657 658 f, err := rp.repoResolver.Resolve(r) 659 if err != nil { 660 l.Error("failed to get repo and knot", "err", err) 661 return 662 } 663 664 subjectStr := r.FormValue("subject") 665 subject, err := syntax.ParseATURI(subjectStr) 666 if err != nil { 667 l.Error("failed to get repo and knot", "err", err) 668 return 669 } 670 671 labelDefs, err := db.GetLabelDefinitions( 672 rp.db, 673 orm.FilterIn("at_uri", f.Labels), 674 orm.FilterContains("scope", subject.Collection().String()), 675 ) 676 if err != nil { 677 l.Error("failed to fetch labels", "err", err) 678 return 679 } 680 681 defs := make(map[string]*models.LabelDefinition) 682 for _, l := range labelDefs { 683 defs[l.AtUri().String()] = &l 684 } 685 686 states, err := db.GetLabels(rp.db, orm.FilterEq("subject", subject)) 687 if err != nil { 688 l.Error("failed to build label state", "err", err) 689 return 690 } 691 state := states[subject] 692 693 user := rp.oauth.GetMultiAccountUser(r) 694 rp.pages.EditLabelPanel(w, pages.EditLabelPanelParams{ 695 LoggedInUser: user, 696 RepoInfo: rp.repoResolver.GetRepoInfo(r, user), 697 Defs: defs, 698 Subject: subject.String(), 699 State: state, 700 }) 701} 702 703func (rp *Repo) AddCollaborator(w http.ResponseWriter, r *http.Request) { 704 user := rp.oauth.GetMultiAccountUser(r) 705 l := rp.logger.With("handler", "AddCollaborator") 706 l = l.With("did", user.Active.Did) 707 708 f, err := rp.repoResolver.Resolve(r) 709 if err != nil { 710 l.Error("failed to get repo and knot", "err", err) 711 return 712 } 713 714 errorId := "add-collaborator-error" 715 fail := func(msg string, err error) { 716 l.Error(msg, "err", err) 717 rp.pages.Notice(w, errorId, msg) 718 } 719 720 collaborator := r.FormValue("collaborator") 721 if collaborator == "" { 722 fail("Invalid form.", nil) 723 return 724 } 725 726 // remove a single leading `@`, to make @handle work with ResolveIdent 727 collaborator = strings.TrimPrefix(collaborator, "@") 728 729 collaboratorIdent, err := rp.idResolver.ResolveIdent(r.Context(), collaborator) 730 if err != nil { 731 fail(fmt.Sprintf("'%s' is not a valid DID/handle.", collaborator), err) 732 return 733 } 734 735 if collaboratorIdent.DID.String() == user.Active.Did { 736 fail("You seem to be adding yourself as a collaborator.", nil) 737 return 738 } 739 l = l.With("collaborator", collaboratorIdent.Handle) 740 l = l.With("knot", f.Knot) 741 742 // announce this relation into the firehose, store into owners' pds 743 client, err := rp.oauth.AuthorizedClient(r) 744 if err != nil { 745 fail("Failed to write to PDS.", err) 746 return 747 } 748 749 // emit a record 750 currentUser := rp.oauth.GetMultiAccountUser(r) 751 rkey := tid.TID() 752 createdAt := time.Now() 753 resp, err := comatproto.RepoPutRecord(r.Context(), client, &comatproto.RepoPutRecord_Input{ 754 Collection: tangled.RepoCollaboratorNSID, 755 Repo: currentUser.Active.Did, 756 Rkey: rkey, 757 Record: &lexutil.LexiconTypeDecoder{ 758 Val: &tangled.RepoCollaborator{ 759 Subject: collaboratorIdent.DID.String(), 760 Repo: string(f.RepoAt()), 761 CreatedAt: createdAt.Format(time.RFC3339), 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.DidSlashRepo()) 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 err := xrpcclient.HandleXrpcErr(err); err != nil { 883 rp.pages.Notice(w, noticeId, err.Error()) 884 return 885 } 886 l.Info("deleted repo from knot") 887 888 tx, err := rp.db.BeginTx(r.Context(), nil) 889 if err != nil { 890 l.Error("failed to start tx") 891 w.Write(fmt.Append(nil, "failed to add collaborator: ", err)) 892 return 893 } 894 defer func() { 895 tx.Rollback() 896 err = rp.enforcer.E.LoadPolicy() 897 if err != nil { 898 l.Error("failed to rollback policies") 899 } 900 }() 901 902 // remove collaborator RBAC 903 repoCollaborators, err := rp.enforcer.E.GetImplicitUsersForResourceByDomain(f.DidSlashRepo(), f.Knot) 904 if err != nil { 905 rp.pages.Notice(w, noticeId, "Failed to remove collaborators") 906 return 907 } 908 for _, c := range repoCollaborators { 909 did := c[0] 910 rp.enforcer.RemoveCollaborator(did, f.Knot, f.DidSlashRepo()) 911 } 912 l.Info("removed collaborators") 913 914 // remove repo RBAC 915 err = rp.enforcer.RemoveRepo(f.Did, f.Knot, f.DidSlashRepo()) 916 if err != nil { 917 rp.pages.Notice(w, noticeId, "Failed to update RBAC rules") 918 return 919 } 920 921 // remove repo from db 922 err = db.RemoveRepo(tx, f.Did, f.Name) 923 if err != nil { 924 rp.pages.Notice(w, noticeId, "Failed to update appview") 925 return 926 } 927 l.Info("removed repo from db") 928 929 err = tx.Commit() 930 if err != nil { 931 l.Error("failed to commit changes", "err", err) 932 http.Error(w, err.Error(), http.StatusInternalServerError) 933 return 934 } 935 936 err = rp.enforcer.E.SavePolicy() 937 if err != nil { 938 l.Error("failed to update ACLs", "err", err) 939 http.Error(w, err.Error(), http.StatusInternalServerError) 940 return 941 } 942 943 rp.pages.HxRedirect(w, fmt.Sprintf("/%s", f.Did)) 944} 945 946func (rp *Repo) SyncRepoFork(w http.ResponseWriter, r *http.Request) { 947 l := rp.logger.With("handler", "SyncRepoFork") 948 949 ref := chi.URLParam(r, "ref") 950 ref, _ = url.PathUnescape(ref) 951 952 user := rp.oauth.GetMultiAccountUser(r) 953 f, err := rp.repoResolver.Resolve(r) 954 if err != nil { 955 l.Error("failed to resolve source repo", "err", err) 956 return 957 } 958 959 switch r.Method { 960 case http.MethodPost: 961 client, err := rp.oauth.ServiceClient( 962 r, 963 oauth.WithService(f.Knot), 964 oauth.WithLxm(tangled.RepoForkSyncNSID), 965 oauth.WithDev(rp.config.Core.Dev), 966 ) 967 if err != nil { 968 rp.pages.Notice(w, "repo", "Failed to connect to knot server.") 969 return 970 } 971 972 if f.Source == "" { 973 rp.pages.Notice(w, "repo", "This repository is not a fork.") 974 return 975 } 976 977 err = tangled.RepoForkSync( 978 r.Context(), 979 client, 980 &tangled.RepoForkSync_Input{ 981 Did: user.Active.Did, 982 Name: f.Name, 983 Source: f.Source, 984 Branch: ref, 985 }, 986 ) 987 if err := xrpcclient.HandleXrpcErr(err); err != nil { 988 rp.pages.Notice(w, "repo", err.Error()) 989 return 990 } 991 992 rp.pages.HxRefresh(w) 993 return 994 } 995} 996 997func (rp *Repo) ForkRepo(w http.ResponseWriter, r *http.Request) { 998 l := rp.logger.With("handler", "ForkRepo") 999 1000 user := rp.oauth.GetMultiAccountUser(r) 1001 f, err := rp.repoResolver.Resolve(r) 1002 if err != nil { 1003 l.Error("failed to resolve source repo", "err", err) 1004 return 1005 } 1006 1007 switch r.Method { 1008 case http.MethodGet: 1009 user := rp.oauth.GetMultiAccountUser(r) 1010 knots, err := rp.enforcer.GetKnotsForUser(user.Active.Did) 1011 if err != nil { 1012 rp.pages.Notice(w, "repo", "Invalid user account.") 1013 return 1014 } 1015 1016 rp.pages.ForkRepo(w, pages.ForkRepoParams{ 1017 LoggedInUser: user, 1018 Knots: knots, 1019 RepoInfo: rp.repoResolver.GetRepoInfo(r, user), 1020 }) 1021 1022 case http.MethodPost: 1023 l := rp.logger.With("handler", "ForkRepo") 1024 1025 targetKnot := r.FormValue("knot") 1026 if targetKnot == "" { 1027 rp.pages.Notice(w, "repo", "Invalid form submission—missing knot domain.") 1028 return 1029 } 1030 l = l.With("targetKnot", targetKnot) 1031 1032 ok, err := rp.enforcer.E.Enforce(user.Active.Did, targetKnot, targetKnot, "repo:create") 1033 if err != nil || !ok { 1034 rp.pages.Notice(w, "repo", "You do not have permission to create a repo in this knot.") 1035 return 1036 } 1037 1038 // choose a name for a fork 1039 forkName := r.FormValue("repo_name") 1040 if forkName == "" { 1041 rp.pages.Notice(w, "repo", "Repository name cannot be empty.") 1042 return 1043 } 1044 1045 // this check is *only* to see if the forked repo name already exists 1046 // in the user's account. 1047 existingRepo, err := db.GetRepo( 1048 rp.db, 1049 orm.FilterEq("did", user.Active.Did), 1050 orm.FilterEq("name", forkName), 1051 ) 1052 if err != nil { 1053 if !errors.Is(err, sql.ErrNoRows) { 1054 l.Error("error fetching existing repo from db", "err", err) 1055 rp.pages.Notice(w, "repo", "Failed to fork this repository. Try again later.") 1056 return 1057 } 1058 } else if existingRepo != nil { 1059 // repo with this name already exists 1060 rp.pages.Notice(w, "repo", "A repository with this name already exists.") 1061 return 1062 } 1063 l = l.With("forkName", forkName) 1064 1065 uri := "https" 1066 if rp.config.Core.Dev { 1067 uri = "http" 1068 } 1069 1070 forkSourceUrl := fmt.Sprintf("%s://%s/%s/%s", uri, f.Knot, f.Did, f.Name) 1071 l = l.With("cloneUrl", forkSourceUrl) 1072 1073 sourceAt := f.RepoAt().String() 1074 1075 // create an atproto record for this fork 1076 rkey := tid.TID() 1077 repo := &models.Repo{ 1078 Did: user.Active.Did, 1079 Name: forkName, 1080 Knot: targetKnot, 1081 Rkey: rkey, 1082 Source: sourceAt, 1083 Description: f.Description, 1084 Created: time.Now(), 1085 Labels: rp.config.Label.DefaultLabelDefs, 1086 } 1087 record := repo.AsRecord() 1088 1089 atpClient, err := rp.oauth.AuthorizedClient(r) 1090 if err != nil { 1091 l.Error("failed to create xrpcclient", "err", err) 1092 rp.pages.Notice(w, "repo", "Failed to fork repository.") 1093 return 1094 } 1095 1096 atresp, err := comatproto.RepoPutRecord(r.Context(), atpClient, &comatproto.RepoPutRecord_Input{ 1097 Collection: tangled.RepoNSID, 1098 Repo: user.Active.Did, 1099 Rkey: rkey, 1100 Record: &lexutil.LexiconTypeDecoder{ 1101 Val: &record, 1102 }, 1103 }) 1104 if err != nil { 1105 l.Error("failed to write to PDS", "err", err) 1106 rp.pages.Notice(w, "repo", "Failed to announce repository creation.") 1107 return 1108 } 1109 1110 aturi := atresp.Uri 1111 l = l.With("aturi", aturi) 1112 l.Info("wrote to PDS") 1113 1114 tx, err := rp.db.BeginTx(r.Context(), nil) 1115 if err != nil { 1116 l.Info("txn failed", "err", err) 1117 rp.pages.Notice(w, "repo", "Failed to save repository information.") 1118 return 1119 } 1120 1121 // The rollback function reverts a few things on failure: 1122 // - the pending txn 1123 // - the ACLs 1124 // - the atproto record created 1125 rollback := func() { 1126 err1 := tx.Rollback() 1127 err2 := rp.enforcer.E.LoadPolicy() 1128 err3 := rollbackRecord(context.Background(), aturi, atpClient) 1129 1130 // ignore txn complete errors, this is okay 1131 if errors.Is(err1, sql.ErrTxDone) { 1132 err1 = nil 1133 } 1134 1135 if errs := errors.Join(err1, err2, err3); errs != nil { 1136 l.Error("failed to rollback changes", "errs", errs) 1137 return 1138 } 1139 } 1140 defer rollback() 1141 1142 // TODO: this could coordinate better with the knot to recieve a clone status 1143 client, err := rp.oauth.ServiceClient( 1144 r, 1145 oauth.WithService(targetKnot), 1146 oauth.WithLxm(tangled.RepoCreateNSID), 1147 oauth.WithDev(rp.config.Core.Dev), 1148 oauth.WithTimeout(time.Second*20), // big repos take time to clone 1149 ) 1150 if err != nil { 1151 l.Error("could not create service client", "err", err) 1152 rp.pages.Notice(w, "repo", "Failed to connect to knot server.") 1153 return 1154 } 1155 1156 err = tangled.RepoCreate( 1157 r.Context(), 1158 client, 1159 &tangled.RepoCreate_Input{ 1160 Rkey: rkey, 1161 Source: &forkSourceUrl, 1162 }, 1163 ) 1164 if err := xrpcclient.HandleXrpcErr(err); err != nil { 1165 rp.pages.Notice(w, "repo", err.Error()) 1166 return 1167 } 1168 1169 err = db.AddRepo(tx, repo) 1170 if err != nil { 1171 l.Error("failed to AddRepo", "err", err) 1172 rp.pages.Notice(w, "repo", "Failed to save repository information.") 1173 return 1174 } 1175 1176 // acls 1177 p, _ := securejoin.SecureJoin(user.Active.Did, forkName) 1178 err = rp.enforcer.AddRepo(user.Active.Did, targetKnot, p) 1179 if err != nil { 1180 l.Error("failed to add ACLs", "err", err) 1181 rp.pages.Notice(w, "repo", "Failed to set up repository permissions.") 1182 return 1183 } 1184 1185 err = tx.Commit() 1186 if err != nil { 1187 l.Error("failed to commit changes", "err", err) 1188 http.Error(w, err.Error(), http.StatusInternalServerError) 1189 return 1190 } 1191 1192 err = rp.enforcer.E.SavePolicy() 1193 if err != nil { 1194 l.Error("failed to update ACLs", "err", err) 1195 http.Error(w, err.Error(), http.StatusInternalServerError) 1196 return 1197 } 1198 1199 // reset the ATURI because the transaction completed successfully 1200 aturi = "" 1201 1202 rp.notifier.NewRepo(r.Context(), repo) 1203 rp.pages.HxLocation(w, fmt.Sprintf("/%s/%s", user.Active.Did, forkName)) 1204 } 1205} 1206 1207// this is used to rollback changes made to the PDS 1208// 1209// it is a no-op if the provided ATURI is empty 1210func rollbackRecord(ctx context.Context, aturi string, client *atclient.APIClient) error { 1211 if aturi == "" { 1212 return nil 1213 } 1214 1215 parsed := syntax.ATURI(aturi) 1216 1217 collection := parsed.Collection().String() 1218 repo := parsed.Authority().String() 1219 rkey := parsed.RecordKey().String() 1220 1221 _, err := comatproto.RepoDeleteRecord(ctx, client, &comatproto.RepoDeleteRecord_Input{ 1222 Collection: collection, 1223 Repo: repo, 1224 Rkey: rkey, 1225 }) 1226 return err 1227}