This repository has no description
0

Configure Feed

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

core / appview / knots / knots.go
18 kB 706 lines
1package knots 2 3import ( 4 "context" 5 "errors" 6 "fmt" 7 "log/slog" 8 "net/http" 9 "slices" 10 "strings" 11 "time" 12 13 "github.com/go-chi/chi/v5" 14 "tangled.org/core/api/tangled" 15 "tangled.org/core/appview/config" 16 "tangled.org/core/appview/db" 17 "tangled.org/core/appview/middleware" 18 "tangled.org/core/appview/models" 19 "tangled.org/core/appview/oauth" 20 "tangled.org/core/appview/pages" 21 "tangled.org/core/appview/serververify" 22 "tangled.org/core/appview/xrpcclient" 23 "tangled.org/core/eventconsumer" 24 "tangled.org/core/idresolver" 25 "tangled.org/core/orm" 26 "tangled.org/core/rbac" 27 "tangled.org/core/tid" 28 29 comatproto "github.com/bluesky-social/indigo/api/atproto" 30 "github.com/bluesky-social/indigo/atproto/atclient" 31 lexutil "github.com/bluesky-social/indigo/lex/util" 32) 33 34type Knots struct { 35 Db *db.DB 36 OAuth *oauth.OAuth 37 Pages *pages.Pages 38 Config *config.Config 39 Enforcer *rbac.Enforcer 40 IdResolver *idresolver.Resolver 41 Logger *slog.Logger 42 Knotstream *eventconsumer.Consumer 43} 44 45func (k *Knots) Router() http.Handler { 46 r := chi.NewRouter() 47 48 r.With(middleware.AuthMiddleware(k.OAuth)).Get("/", k.knots) 49 r.With(middleware.AuthMiddleware(k.OAuth)).Post("/register", k.register) 50 51 r.With(middleware.AuthMiddleware(k.OAuth)).Get("/{domain}", k.dashboard) 52 r.With(middleware.AuthMiddleware(k.OAuth)).Delete("/{domain}", k.delete) 53 54 r.With(middleware.AuthMiddleware(k.OAuth)).Post("/{domain}/retry", k.retry) 55 r.With(middleware.AuthMiddleware(k.OAuth)).Post("/{domain}/add", k.addMember) 56 r.With(middleware.AuthMiddleware(k.OAuth)).Post("/{domain}/remove", k.removeMember) 57 58 return r 59} 60 61func (k *Knots) knots(w http.ResponseWriter, r *http.Request) { 62 user := k.OAuth.GetMultiAccountUser(r) 63 registrations, err := db.GetRegistrations( 64 k.Db, 65 orm.FilterEq("did", user.Did), 66 ) 67 if err != nil { 68 k.Logger.Error("failed to fetch knot registrations", "err", err) 69 w.WriteHeader(http.StatusInternalServerError) 70 return 71 } 72 73 knots := make([]pages.KnotListingParams, 0, len(registrations)) 74 for i := range registrations { 75 registration := &registrations[i] 76 count, err := db.CountRepos(k.Db, orm.FilterEq("knot", registration.Domain)) 77 if err != nil { 78 k.Logger.Error("failed to count knot repos", "err", err, "domain", registration.Domain) 79 w.WriteHeader(http.StatusInternalServerError) 80 return 81 } 82 knots = append(knots, pages.KnotListingParams{ 83 Registration: registration, 84 RepoCount: int(count), 85 }) 86 } 87 88 k.Pages.Knots(w, pages.KnotsParams{ 89 LoggedInUser: user, 90 Knots: knots, 91 }) 92} 93 94func (k *Knots) dashboard(w http.ResponseWriter, r *http.Request) { 95 l := k.Logger.With("handler", "dashboard") 96 97 user := k.OAuth.GetMultiAccountUser(r) 98 l = l.With("user", user.Did) 99 100 domain := chi.URLParam(r, "domain") 101 if domain == "" { 102 return 103 } 104 l = l.With("domain", domain) 105 106 registrations, err := db.GetRegistrations( 107 k.Db, 108 orm.FilterEq("did", user.Did), 109 orm.FilterEq("domain", domain), 110 ) 111 if err != nil { 112 l.Error("failed to get registrations", "err", err) 113 http.Error(w, "Not found", http.StatusNotFound) 114 return 115 } 116 if len(registrations) != 1 { 117 l.Error("got incorrect number of registrations", "got", len(registrations), "expected", 1) 118 return 119 } 120 registration := registrations[0] 121 122 members, err := k.Enforcer.GetUserByRole("server:member", domain) 123 if err != nil { 124 l.Error("failed to get knot members", "err", err) 125 http.Error(w, "Not found", http.StatusInternalServerError) 126 return 127 } 128 slices.Sort(members) 129 130 repos, err := db.GetRepos( 131 k.Db, 132 orm.FilterEq("knot", domain), 133 ) 134 if err != nil { 135 l.Error("failed to get knot repos", "err", err) 136 http.Error(w, "Not found", http.StatusInternalServerError) 137 return 138 } 139 140 // organize repos by did 141 repoMap := make(map[string][]models.Repo) 142 for _, r := range repos { 143 repoMap[r.Did] = append(repoMap[r.Did], r) 144 } 145 146 k.Pages.Knot(w, pages.KnotParams{ 147 LoggedInUser: user, 148 Registration: &registration, 149 Members: members, 150 Repos: repoMap, 151 IsOwner: true, 152 RepoCount: len(repos), 153 }) 154} 155 156func (k *Knots) register(w http.ResponseWriter, r *http.Request) { 157 user := k.OAuth.GetMultiAccountUser(r) 158 l := k.Logger.With("handler", "register") 159 160 noticeId := "register-error" 161 defaultErr := "Failed to register knot. Try again later." 162 fail := func() { 163 k.Pages.Notice(w, noticeId, defaultErr) 164 } 165 166 domain := r.FormValue("domain") 167 // Strip protocol, trailing slashes, and whitespace 168 // Rkey cannot contain slashes 169 domain = strings.TrimSpace(domain) 170 domain = strings.TrimPrefix(domain, "https://") 171 domain = strings.TrimPrefix(domain, "http://") 172 domain = strings.TrimSuffix(domain, "/") 173 if domain == "" { 174 k.Pages.Notice(w, noticeId, "Incomplete form.") 175 return 176 } 177 l = l.With("domain", domain) 178 l = l.With("user", user.Did) 179 180 tx, err := k.Db.Begin() 181 if err != nil { 182 l.Error("failed to start transaction", "err", err) 183 fail() 184 return 185 } 186 defer tx.Rollback() 187 188 if err := db.AddKnot(tx, domain, user.Did); err != nil { 189 l.Error("failed to insert", "err", err) 190 fail() 191 return 192 } 193 194 client, err := k.OAuth.AuthorizedClient(r) 195 if err != nil { 196 l.Error("failed to authorize client", "err", err) 197 fail() 198 return 199 } 200 201 ex, _ := comatproto.RepoGetRecord(r.Context(), client, "", tangled.KnotNSID, user.Did, domain) 202 var exCid *string 203 if ex != nil { 204 exCid = ex.Cid 205 } 206 207 _, err = comatproto.RepoPutRecord(r.Context(), client, &comatproto.RepoPutRecord_Input{ 208 Collection: tangled.KnotNSID, 209 Repo: user.Did, 210 Rkey: domain, 211 Record: &lexutil.LexiconTypeDecoder{ 212 Val: &tangled.Knot{ 213 CreatedAt: time.Now().Format(time.RFC3339), 214 }, 215 }, 216 SwapRecord: exCid, 217 }) 218 if err != nil { 219 l.Error("failed to put record", "err", err) 220 fail() 221 return 222 } 223 224 if err := tx.Commit(); err != nil { 225 l.Error("failed to commit transaction", "err", err) 226 fail() 227 return 228 } 229 230 go k.Knotstream.AddSource(r.Context(), eventconsumer.NewKnotSource(domain)) 231 232 k.Pages.HxRefresh(w) 233} 234 235func (k *Knots) delete(w http.ResponseWriter, r *http.Request) { 236 user := k.OAuth.GetMultiAccountUser(r) 237 l := k.Logger.With("handler", "delete") 238 239 noticeId := "operation-error" 240 defaultErr := "Failed to delete knot. Try again later." 241 fail := func() { 242 k.Pages.Notice(w, noticeId, defaultErr) 243 } 244 245 domain := chi.URLParam(r, "domain") 246 if domain == "" { 247 l.Error("empty domain") 248 fail() 249 return 250 } 251 252 // get record from db first 253 registrations, err := db.GetRegistrations( 254 k.Db, 255 orm.FilterEq("did", user.Did), 256 orm.FilterEq("domain", domain), 257 ) 258 if err != nil { 259 l.Error("failed to get registration", "err", err) 260 fail() 261 return 262 } 263 if len(registrations) != 1 { 264 l.Error("got incorrect number of registrations", "got", len(registrations), "expected", 1) 265 fail() 266 return 267 } 268 registration := registrations[0] 269 270 tx, err := k.Db.Begin() 271 if err != nil { 272 l.Error("failed to start txn", "err", err) 273 fail() 274 return 275 } 276 defer func() { 277 tx.Rollback() 278 k.Enforcer.E.LoadPolicy() 279 }() 280 281 err = db.DeleteKnot( 282 tx, 283 orm.FilterEq("did", user.Did), 284 orm.FilterEq("domain", domain), 285 ) 286 if err != nil { 287 l.Error("failed to delete registration", "err", err) 288 fail() 289 return 290 } 291 292 err = db.RemoveReposByKnot(tx, domain) 293 if err != nil { 294 l.Error("failed to delete repos", "err", err) 295 fail() 296 return 297 } 298 299 // delete from enforcer if it was registered 300 if registration.Registered != nil { 301 err = k.Enforcer.RemoveKnot(domain) 302 if err != nil { 303 l.Error("failed to update ACL", "err", err) 304 fail() 305 return 306 } 307 } 308 309 client, err := k.OAuth.AuthorizedClient(r) 310 if err != nil { 311 l.Error("failed to authorize client", "err", err) 312 fail() 313 return 314 } 315 316 _, err = comatproto.RepoDeleteRecord(r.Context(), client, &comatproto.RepoDeleteRecord_Input{ 317 Collection: tangled.KnotNSID, 318 Repo: user.Did, 319 Rkey: domain, 320 }) 321 if err != nil { 322 // non-fatal 323 l.Error("failed to delete record", "err", err) 324 } 325 326 err = tx.Commit() 327 if err != nil { 328 l.Error("failed to delete knot", "err", err) 329 fail() 330 return 331 } 332 333 err = k.Enforcer.E.SavePolicy() 334 if err != nil { 335 l.Error("failed to update ACL", "err", err) 336 k.Pages.HxRefresh(w) 337 return 338 } 339 340 if registration.Registered != nil { 341 remaining, rErr := db.GetRegistrations(k.Db, 342 orm.FilterEq("domain", domain), 343 orm.FilterIsNot("registered", "null"), 344 ) 345 if rErr != nil { 346 l.Warn("failed to check remaining registrations after delete", "err", rErr) 347 } else if len(remaining) == 0 { 348 go k.Knotstream.RemoveSource(eventconsumer.NewKnotSource(domain)) 349 } 350 } 351 352 shouldRedirect := r.Header.Get("shouldRedirect") 353 if shouldRedirect == "true" { 354 k.Pages.HxRedirect(w, "/knots") 355 return 356 } 357 358 w.Write([]byte{}) 359} 360 361func (k *Knots) retry(w http.ResponseWriter, r *http.Request) { 362 user := k.OAuth.GetMultiAccountUser(r) 363 l := k.Logger.With("handler", "retry") 364 365 noticeId := "operation-error" 366 defaultErr := "Failed to verify knot. Try again later." 367 fail := func() { 368 k.Pages.Notice(w, noticeId, defaultErr) 369 } 370 371 domain := chi.URLParam(r, "domain") 372 if domain == "" { 373 l.Error("empty domain") 374 fail() 375 return 376 } 377 l = l.With("domain", domain) 378 l = l.With("user", user.Did) 379 380 // get record from db first 381 registrations, err := db.GetRegistrations( 382 k.Db, 383 orm.FilterEq("did", user.Did), 384 orm.FilterEq("domain", domain), 385 ) 386 if err != nil { 387 l.Error("failed to get registration", "err", err) 388 fail() 389 return 390 } 391 if len(registrations) != 1 { 392 l.Error("got incorrect number of registrations", "got", len(registrations), "expected", 1) 393 fail() 394 return 395 } 396 registration := registrations[0] 397 398 // begin verification 399 err = serververify.RunVerification(r.Context(), domain, user.Did, k.Config.Core.Dev) 400 if err != nil { 401 l.Error("verification failed", "err", err) 402 403 if errors.Is(err, xrpcclient.ErrXrpcUnsupported) { 404 k.Pages.Notice(w, noticeId, "Failed to verify knot, XRPC queries are unsupported on this knot, consider upgrading!") 405 return 406 } 407 408 if e, ok := err.(*serververify.OwnerMismatch); ok { 409 k.Pages.Notice(w, noticeId, e.Error()) 410 return 411 } 412 413 fail() 414 return 415 } 416 417 err = serververify.MarkKnotVerified(k.Db, k.Enforcer, domain, user.Did) 418 if err != nil { 419 l.Error("failed to mark verified", "err", err) 420 k.Pages.Notice(w, noticeId, err.Error()) 421 return 422 } 423 424 // if this knot requires upgrade, then emit a record too 425 // 426 // this is part of migrating from the old knot system to the new one 427 if registration.NeedsUpgrade { 428 // re-announce by registering under same rkey 429 client, err := k.OAuth.AuthorizedClient(r) 430 if err != nil { 431 l.Error("failed to authorize client", "err", err) 432 fail() 433 return 434 } 435 436 ex, _ := comatproto.RepoGetRecord(r.Context(), client, "", tangled.KnotNSID, user.Did, domain) 437 var exCid *string 438 if ex != nil { 439 exCid = ex.Cid 440 } 441 442 // ignore the error here 443 _, err = comatproto.RepoPutRecord(r.Context(), client, &comatproto.RepoPutRecord_Input{ 444 Collection: tangled.KnotNSID, 445 Repo: user.Did, 446 Rkey: domain, 447 Record: &lexutil.LexiconTypeDecoder{ 448 Val: &tangled.Knot{ 449 CreatedAt: time.Now().Format(time.RFC3339), 450 }, 451 }, 452 SwapRecord: exCid, 453 }) 454 if err != nil { 455 l.Error("non-fatal: failed to reannouce knot", "err", err) 456 } 457 } 458 459 // add this knot to knotstream 460 go k.Knotstream.AddSource( 461 r.Context(), 462 eventconsumer.NewKnotSource(domain), 463 ) 464 465 shouldRefresh := r.Header.Get("shouldRefresh") 466 if shouldRefresh == "true" { 467 k.Pages.HxRefresh(w) 468 return 469 } 470 471 // Get updated registration to show 472 registrations, err = db.GetRegistrations( 473 k.Db, 474 orm.FilterEq("did", user.Did), 475 orm.FilterEq("domain", domain), 476 ) 477 if err != nil { 478 l.Error("failed to get registration", "err", err) 479 fail() 480 return 481 } 482 if len(registrations) != 1 { 483 l.Error("got incorrect number of registrations", "got", len(registrations), "expected", 1) 484 fail() 485 return 486 } 487 updatedRegistration := registrations[0] 488 489 count, err := db.CountRepos(k.Db, orm.FilterEq("knot", domain)) 490 if err != nil { 491 l.Error("failed to count knot repos", "err", err) 492 fail() 493 return 494 } 495 496 w.Header().Set("HX-Reswap", "outerHTML") 497 k.Pages.KnotListing(w, pages.KnotListingParams{ 498 Registration: &updatedRegistration, 499 RepoCount: int(count), 500 }) 501} 502 503func (k *Knots) addMember(w http.ResponseWriter, r *http.Request) { 504 user := k.OAuth.GetMultiAccountUser(r) 505 l := k.Logger.With("handler", "addMember") 506 507 domain := chi.URLParam(r, "domain") 508 if domain == "" { 509 l.Error("empty domain") 510 http.Error(w, "Not found", http.StatusNotFound) 511 return 512 } 513 l = l.With("domain", domain) 514 l = l.With("user", user.Did) 515 516 registrations, err := db.GetRegistrations( 517 k.Db, 518 orm.FilterEq("did", user.Did), 519 orm.FilterEq("domain", domain), 520 orm.FilterIsNot("registered", "null"), 521 ) 522 if err != nil { 523 l.Error("failed to get registration", "err", err) 524 return 525 } 526 if len(registrations) != 1 { 527 l.Error("got incorrect number of registrations", "got", len(registrations), "expected", 1) 528 return 529 } 530 registration := registrations[0] 531 532 noticeId := fmt.Sprintf("add-member-error-%d", registration.Id) 533 defaultErr := "Failed to add member. Try again later." 534 fail := func() { 535 k.Pages.Notice(w, noticeId, defaultErr) 536 } 537 538 member := r.FormValue("member") 539 member = strings.TrimPrefix(member, "@") 540 if member == "" { 541 l.Error("empty member") 542 k.Pages.Notice(w, noticeId, "Failed to add member, empty form.") 543 return 544 } 545 l = l.With("member", member) 546 547 memberId, err := k.IdResolver.ResolveIdent(r.Context(), member) 548 if err != nil { 549 l.Error("failed to resolve member identity to handle", "err", err) 550 k.Pages.Notice(w, noticeId, "Failed to add member, identity resolution failed.") 551 return 552 } 553 if memberId.Handle.IsInvalidHandle() { 554 l.Error("failed to resolve member identity to handle") 555 k.Pages.Notice(w, noticeId, "Failed to add member, identity resolution failed.") 556 return 557 } 558 559 client, err := k.OAuth.AuthorizedClient(r) 560 if err != nil { 561 l.Error("failed to authorize client", "err", err) 562 fail() 563 return 564 } 565 566 rkey := tid.TID() 567 _, err = comatproto.RepoPutRecord(r.Context(), client, &comatproto.RepoPutRecord_Input{ 568 Collection: tangled.KnotMemberNSID, 569 Repo: user.Did, 570 Rkey: rkey, 571 Record: &lexutil.LexiconTypeDecoder{ 572 Val: &tangled.KnotMember{ 573 CreatedAt: time.Now().Format(time.RFC3339), 574 Domain: domain, 575 Subject: memberId.DID.String(), 576 }, 577 }, 578 }) 579 if err != nil { 580 l.Error("failed to add record to PDS", "err", err) 581 k.Pages.Notice(w, noticeId, "Failed to add record to PDS, try again later.") 582 return 583 } 584 585 k.Pages.HxRedirect(w, fmt.Sprintf("/settings/knots/%s", domain)) 586} 587 588func (k *Knots) removeMember(w http.ResponseWriter, r *http.Request) { 589 user := k.OAuth.GetMultiAccountUser(r) 590 l := k.Logger.With("handler", "removeMember") 591 592 noticeId := "operation-error" 593 defaultErr := "Failed to remove member. Try again later." 594 fail := func() { 595 k.Pages.Notice(w, noticeId, defaultErr) 596 } 597 598 domain := chi.URLParam(r, "domain") 599 if domain == "" { 600 l.Error("empty domain") 601 fail() 602 return 603 } 604 l = l.With("domain", domain) 605 l = l.With("user", user.Did) 606 607 registrations, err := db.GetRegistrations( 608 k.Db, 609 orm.FilterEq("did", user.Did), 610 orm.FilterEq("domain", domain), 611 orm.FilterIsNot("registered", "null"), 612 ) 613 if err != nil { 614 l.Error("failed to get registration", "err", err) 615 return 616 } 617 if len(registrations) != 1 { 618 l.Error("got incorrect number of registrations", "got", len(registrations), "expected", 1) 619 return 620 } 621 622 member := r.FormValue("member") 623 member = strings.TrimPrefix(member, "@") 624 if member == "" { 625 l.Error("empty member") 626 k.Pages.Notice(w, noticeId, "Failed to remove member, empty form.") 627 return 628 } 629 l = l.With("member", member) 630 631 memberId, err := k.IdResolver.ResolveIdent(r.Context(), member) 632 if err != nil { 633 l.Error("failed to resolve member identity to handle", "err", err) 634 k.Pages.Notice(w, noticeId, "Failed to remove member, identity resolution failed.") 635 return 636 } 637 638 client, err := k.OAuth.AuthorizedClient(r) 639 if err != nil { 640 l.Error("failed to authorize client", "err", err) 641 fail() 642 return 643 } 644 645 rkey, err := lookupKnotMemberRkey(r.Context(), k.Db, client, user.Did, domain, memberId.DID.String()) 646 if err != nil { 647 l.Warn("failed to look up member rkey", "err", err) 648 } 649 650 if rkey == "" { 651 l.Error("no member record found to remove") 652 fail() 653 return 654 } 655 656 _, err = comatproto.RepoDeleteRecord(r.Context(), client, &comatproto.RepoDeleteRecord_Input{ 657 Collection: tangled.KnotMemberNSID, 658 Repo: user.Did, 659 Rkey: rkey, 660 }) 661 if err != nil { 662 l.Error("failed to delete record from PDS", "err", err) 663 k.Pages.Notice(w, noticeId, "Failed to delete record from PDS, try again later.") 664 return 665 } 666 667 k.Pages.HxRefresh(w) 668} 669 670func lookupKnotMemberRkey(ctx context.Context, d *db.DB, client *atclient.APIClient, ownerDid, domain, subject string) (string, error) { 671 members, err := db.GetKnotMembers( 672 d, 673 orm.FilterEq("did", ownerDid), 674 orm.FilterEq("domain", domain), 675 orm.FilterEq("subject", subject), 676 ) 677 if err != nil { 678 return "", fmt.Errorf("db lookup: %w", err) 679 } 680 if len(members) >= 1 { 681 return members[0].Rkey, nil 682 } 683 return findKnotMemberRkey(ctx, client, ownerDid, domain, subject, "") 684} 685 686func findKnotMemberRkey(ctx context.Context, client *atclient.APIClient, repo, domain, subject, cursor string) (string, error) { 687 out, err := comatproto.RepoListRecords(ctx, client, tangled.KnotMemberNSID, cursor, 100, repo, false) 688 if err != nil { 689 return "", err 690 } 691 for _, rec := range out.Records { 692 m, ok := rec.Value.Val.(*tangled.KnotMember) 693 if !ok { 694 continue 695 } 696 if m.Domain != domain || m.Subject != subject { 697 continue 698 } 699 parts := strings.Split(rec.Uri, "/") 700 return parts[len(parts)-1], nil 701 } 702 if out.Cursor == nil || *out.Cursor == "" || *out.Cursor == cursor { 703 return "", nil 704 } 705 return findKnotMemberRkey(ctx, client, repo, domain, subject, *out.Cursor) 706}