This repository has no description
0

Configure Feed

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

core / appview / state / state.go
23 kB 819 lines
1package state 2 3import ( 4 "context" 5 "database/sql" 6 "errors" 7 "fmt" 8 "log/slog" 9 "net/http" 10 "strings" 11 "time" 12 13 "tangled.org/core/api/tangled" 14 "tangled.org/core/appview" 15 "tangled.org/core/appview/bsky" 16 "tangled.org/core/appview/cache" 17 "tangled.org/core/appview/cloudflare" 18 "tangled.org/core/appview/config" 19 "tangled.org/core/appview/db" 20 "tangled.org/core/appview/email" 21 "tangled.org/core/appview/indexer" 22 "tangled.org/core/appview/mentions" 23 "tangled.org/core/appview/models" 24 "tangled.org/core/appview/notify" 25 dbnotify "tangled.org/core/appview/notify/db" 26 lognotify "tangled.org/core/appview/notify/logging" 27 phnotify "tangled.org/core/appview/notify/posthog" 28 whnotify "tangled.org/core/appview/notify/webhook" 29 "tangled.org/core/appview/oauth" 30 "tangled.org/core/appview/pages" 31 "tangled.org/core/appview/pipelines" 32 pipelinessh "tangled.org/core/appview/pipelines/ssh" 33 "tangled.org/core/appview/reporesolver" 34 "tangled.org/core/appview/repoverify" 35 "tangled.org/core/appview/validator" 36 xrpcclient "tangled.org/core/appview/xrpcclient" 37 "tangled.org/core/consts" 38 "tangled.org/core/eventconsumer" 39 "tangled.org/core/idresolver" 40 "tangled.org/core/jetstream" 41 "tangled.org/core/log" 42 tlog "tangled.org/core/log" 43 "tangled.org/core/orm" 44 "tangled.org/core/rbac" 45 46 comatproto "github.com/bluesky-social/indigo/api/atproto" 47 "github.com/bluesky-social/indigo/atproto/atclient" 48 "github.com/bluesky-social/indigo/atproto/syntax" 49 lexutil "github.com/bluesky-social/indigo/lex/util" 50 "github.com/bluesky-social/indigo/xrpc" 51 52 "github.com/go-chi/chi/v5" 53 "github.com/posthog/posthog-go" 54) 55 56type State struct { 57 db *db.DB 58 notifier notify.Notifier 59 indexer *indexer.Indexer 60 oauth *oauth.OAuth 61 enforcer *rbac.Enforcer 62 pages *pages.Pages 63 idResolver *idresolver.Resolver 64 rdb *cache.Cache 65 mentionsResolver *mentions.Resolver 66 posthog posthog.Client 67 jc *jetstream.JetstreamClient 68 config *config.Config 69 repoResolver *reporesolver.RepoResolver 70 knotstream *eventconsumer.Consumer 71 spindlestream *eventconsumer.Consumer 72 pipelineNotifier *pipelines.StatusNotifier 73 logger *slog.Logger 74 validator *validator.Validator 75 cfClient *cloudflare.Client 76} 77 78func Make(ctx context.Context, config *config.Config) (*State, error) { 79 logger := tlog.FromContext(ctx) 80 81 d, err := db.Make(ctx, config.Core.DbPath) 82 if err != nil { 83 return nil, fmt.Errorf("failed to create db: %w", err) 84 } 85 86 indexer := indexer.New(log.SubLogger(logger, "indexer"), d) 87 err = indexer.Init(ctx) 88 if err != nil { 89 return nil, fmt.Errorf("failed to create indexer: %w", err) 90 } 91 92 enforcer, err := rbac.NewEnforcer(config.Core.DbPath) 93 if err != nil { 94 return nil, fmt.Errorf("failed to create enforcer: %w", err) 95 } 96 97 res, err := idresolver.RedisResolver(config.Redis.ToURL(), config.Plc.PLCURL) 98 if err != nil { 99 logger.Error("failed to create redis resolver", "err", err) 100 res = idresolver.DefaultResolver(config.Plc.PLCURL) 101 } 102 103 var rdb *cache.Cache 104 if config.Redis.Addr != "" { 105 rdb = cache.New(config.Redis.Addr) 106 } 107 108 posthog, err := posthog.NewWithConfig(config.Posthog.ApiKey, posthog.Config{Endpoint: config.Posthog.Endpoint}) 109 if err != nil { 110 return nil, fmt.Errorf("failed to create posthog client: %w", err) 111 } 112 113 pages := pages.NewPages(config, res, d, rdb, log.SubLogger(logger, "pages")) 114 oauth, err := oauth.New(config, posthog, d, enforcer, res, log.SubLogger(logger, "oauth")) 115 if err != nil { 116 return nil, fmt.Errorf("failed to start oauth handler: %w", err) 117 } 118 validator := validator.New(d, res, enforcer) 119 120 repoResolver := reporesolver.New(config, enforcer, d, rdb) 121 122 mentionsResolver := mentions.New(config, res, d, log.SubLogger(logger, "mentionsResolver")) 123 124 jc, err := jetstream.NewJetstreamClient( 125 config.Jetstream.Endpoint, 126 "appview", 127 []string{ 128 tangled.ActorProfileNSID, 129 tangled.FeedStarNSID, 130 tangled.FeedCommentNSID, 131 tangled.GraphFollowNSID, 132 tangled.GraphVouchNSID, 133 tangled.KnotMemberNSID, 134 tangled.KnotNSID, 135 tangled.LabelDefinitionNSID, 136 tangled.LabelOpNSID, 137 tangled.PublicKeyNSID, 138 tangled.RepoArtifactNSID, 139 tangled.RepoIssueCommentNSID, 140 tangled.RepoIssueNSID, 141 tangled.RepoNSID, 142 tangled.RepoPullNSID, 143 tangled.RepoPullCommentNSID, 144 tangled.SpindleMemberNSID, 145 tangled.SpindleNSID, 146 tangled.StringNSID, 147 }, 148 nil, 149 tlog.SubLogger(logger, "jetstream"), 150 d, 151 false, 152 153 // in-memory filter is inapplicable to appview so 154 // we'll never log dids anyway. 155 false, 156 ) 157 if err != nil { 158 return nil, fmt.Errorf("failed to create jetstream client: %w", err) 159 } 160 161 if err := BackfillDefaultDefs(d, res, config.Label.DefaultLabelDefs); err != nil { 162 return nil, fmt.Errorf("failed to backfill default label defs: %w", err) 163 } 164 165 var notifiers []notify.Notifier 166 167 // Always add the database notifier 168 notifiers = append(notifiers, dbnotify.NewDatabaseNotifier(d, res)) 169 170 // Add other notifiers in production only 171 if !config.Core.Dev { 172 notifiers = append(notifiers, phnotify.NewPosthogNotifier(posthog)) 173 } 174 notifiers = append(notifiers, indexer) 175 176 notifiers = append(notifiers, whnotify.NewNotifier(d)) 177 178 notifier := notify.NewMergedNotifier(notifiers) 179 notifier = lognotify.NewLoggingNotifier(notifier, tlog.SubLogger(logger, "notify")) 180 181 ingester := appview.Ingester{ 182 Ctx: ctx, 183 Db: d, 184 Enforcer: enforcer, 185 IdResolver: res, 186 Cache: rdb, 187 Config: config, 188 Logger: log.SubLogger(logger, "ingester"), 189 Validator: validator, 190 MentionsResolver: mentionsResolver, 191 Notifier: notifier, 192 Verifier: repoverify.New(res, config.Core.Dev), 193 } 194 err = jc.StartJetstream(ctx, ingester.Ingest()) 195 if err != nil { 196 return nil, fmt.Errorf("failed to start jetstream watcher: %w", err) 197 } 198 199 go ingester.SweepPendingVerifications() 200 201 var cfClient *cloudflare.Client 202 if config.Cloudflare.ApiToken != "" { 203 cfClient, err = cloudflare.New(config) 204 if err != nil { 205 logger.Warn("failed to create cloudflare client, sites upload will be disabled", "err", err) 206 cfClient = nil 207 } 208 } 209 210 knotstream, err := Knotstream(ctx, config, d, enforcer, posthog, notifier, cfClient) 211 if err != nil { 212 return nil, fmt.Errorf("failed to start knotstream consumer: %w", err) 213 } 214 knotstream.Start(ctx) 215 216 pipelineNotifier := pipelines.NewStatusNotifier() 217 218 spindlestream, err := Spindlestream(ctx, config, d, enforcer, pipelineNotifier) 219 if err != nil { 220 return nil, fmt.Errorf("failed to start spindlestream consumer: %w", err) 221 } 222 spindlestream.Start(ctx) 223 224 state := &State{ 225 db: d, 226 notifier: notifier, 227 indexer: indexer, 228 oauth: oauth, 229 enforcer: enforcer, 230 pages: pages, 231 idResolver: res, 232 rdb: rdb, 233 mentionsResolver: mentionsResolver, 234 posthog: posthog, 235 jc: jc, 236 config: config, 237 repoResolver: repoResolver, 238 knotstream: knotstream, 239 spindlestream: spindlestream, 240 pipelineNotifier: pipelineNotifier, 241 logger: logger, 242 validator: validator, 243 cfClient: cfClient, 244 } 245 246 // fetch initial bluesky posts if configured 247 go fetchBskyPosts(ctx, res, config, d, logger) 248 249 return state, nil 250} 251 252func (s *State) Close() error { 253 // other close up logic goes here 254 return s.db.Close() 255} 256 257func (s *State) SecurityTxt(w http.ResponseWriter, r *http.Request) { 258 w.Header().Set("Content-Type", "text/plain") 259 w.Header().Set("Cache-Control", "public, max-age=86400") // one day 260 261 securityTxt := `Contact: mailto:security@tangled.org 262Preferred-Languages: en 263Canonical: https://tangled.org/.well-known/security.txt 264Expires: 2030-01-01T21:59:00.000Z 265` 266 w.Write([]byte(securityTxt)) 267} 268 269func (s *State) RobotsTxt(w http.ResponseWriter, r *http.Request) { 270 w.Header().Set("Content-Type", "text/plain") 271 w.Header().Set("Cache-Control", "public, max-age=86400") // one day 272 273 robotsTxt := `# Hello, Tanglers! 274User-agent: * 275Allow: / 276Disallow: /*/*/settings 277Disallow: /settings 278Disallow: /*/*/compare 279Disallow: /*/*/fork 280 281Crawl-delay: 1 282` 283 w.Write([]byte(robotsTxt)) 284} 285 286func (s *State) TermsOfService(w http.ResponseWriter, r *http.Request) { 287 user := s.oauth.GetMultiAccountUser(r) 288 s.pages.TermsOfService(w, pages.TermsOfServiceParams{ 289 LoggedInUser: user, 290 }) 291} 292 293func (s *State) PrivacyPolicy(w http.ResponseWriter, r *http.Request) { 294 user := s.oauth.GetMultiAccountUser(r) 295 s.pages.PrivacyPolicy(w, pages.PrivacyPolicyParams{ 296 LoggedInUser: user, 297 }) 298} 299 300func (s *State) Brand(w http.ResponseWriter, r *http.Request) { 301 user := s.oauth.GetMultiAccountUser(r) 302 s.pages.Brand(w, pages.BrandParams{ 303 LoggedInUser: user, 304 }) 305} 306 307func (s *State) UpgradeBanner(w http.ResponseWriter, r *http.Request) { 308 user := s.oauth.GetMultiAccountUser(r) 309 if user == nil { 310 return 311 } 312 313 l := s.logger.With("handler", "UpgradeBanner") 314 l = l.With("did", user.Did) 315 316 regs, err := db.GetRegistrations( 317 s.db, 318 orm.FilterEq("did", user.Did), 319 orm.FilterEq("needs_upgrade", 1), 320 ) 321 if err != nil { 322 l.Error("non-fatal: failed to get registrations", "err", err) 323 } 324 325 spindles, err := db.GetSpindles( 326 r.Context(), 327 s.db, 328 orm.FilterEq("owner", user.Did), 329 orm.FilterEq("needs_upgrade", 1), 330 ) 331 if err != nil { 332 l.Error("non-fatal: failed to get spindles", "err", err) 333 } 334 335 if regs == nil && spindles == nil { 336 return 337 } 338 339 s.pages.UpgradeBanner(w, pages.UpgradeBannerParams{ 340 Registrations: regs, 341 Spindles: spindles, 342 }) 343} 344 345func (s *State) NewsletterSignup(w http.ResponseWriter, r *http.Request) { 346 // target is echoed back from the form via hx-vals so the response span's 347 // id matches the form's hx-target. Fallback keeps the handler useful if 348 // a caller forgets to send it. 349 target := strings.TrimSpace(r.FormValue("target")) 350 if target == "" { 351 target = "home" 352 } 353 354 w.Header().Set("Content-Type", "text/html") 355 356 emailAddr := strings.TrimSpace(r.FormValue("email")) 357 if !email.IsValidEmail(emailAddr) { 358 s.pages.NewsletterResponse(w, pages.NewsletterResponseParams{ 359 Id: target, 360 Error: "Invalid email address.", 361 }) 362 return 363 } 364 365 // For logged-in users, persist the signup locally so the widget stays 366 // hidden across devices. The DB row is the render-time source of truth; 367 // Resend still owns the mailing list itself. 368 if user := s.oauth.GetMultiAccountUser(r); user != nil { 369 if err := db.UpsertNewsletterPref(s.db, user.Did, db.NewsletterStatusSubscribed, emailAddr); err != nil { 370 s.logger.Error("failed to persist newsletter preference", "did", user.Did, "err", err) 371 } 372 } 373 374 if s.config.Resend.ApiKey != "" && s.config.Resend.NewsletterSegmentId != "" { 375 go func() { 376 if err := email.AddNewsletterContact(s.config.Resend.ApiKey, s.config.Resend.NewsletterSegmentId, emailAddr); err != nil { 377 s.logger.Error("failed to add newsletter contact", "error", err) 378 } 379 }() 380 } else { 381 s.logger.Error( 382 "failed to add newsletter contact, missing resend config", 383 "isKeyPresent", s.config.Resend.ApiKey != "", 384 "isSegmentIdPresent", s.config.Resend.NewsletterSegmentId != "", 385 "emailAddr", emailAddr, 386 ) 387 } 388 389 s.pages.NewsletterResponse(w, pages.NewsletterResponseParams{Id: target}) 390} 391 392// NewsletterDismiss records that a logged-in user has dismissed the newsletter 393// widget so it stays hidden across their devices. Anonymous callers get a 204 394// with no DB write — localStorage handles the per-browser fallback. 395func (s *State) NewsletterDismiss(w http.ResponseWriter, r *http.Request) { 396 user := s.oauth.GetMultiAccountUser(r) 397 if user == nil { 398 w.WriteHeader(http.StatusNoContent) 399 return 400 } 401 402 if err := db.UpsertNewsletterPref(s.db, user.Did, db.NewsletterStatusDismissed, ""); err != nil { 403 s.logger.Error("failed to persist newsletter dismissal", "did", user.Did, "err", err) 404 } 405 w.WriteHeader(http.StatusNoContent) 406} 407 408func (s *State) Keys(w http.ResponseWriter, r *http.Request) { 409 user := chi.URLParam(r, "user") 410 user = strings.TrimPrefix(user, "@") 411 412 if user == "" { 413 w.WriteHeader(http.StatusBadRequest) 414 return 415 } 416 417 id, err := s.idResolver.ResolveIdent(r.Context(), user) 418 if err != nil { 419 w.WriteHeader(http.StatusInternalServerError) 420 return 421 } 422 423 pubKeys, err := db.GetPublicKeysForDid(s.db, id.DID.String()) 424 if err != nil { 425 s.logger.Error("failed to get public keys", "err", err) 426 http.Error(w, "failed to get public keys", http.StatusInternalServerError) 427 return 428 } 429 430 if len(pubKeys) == 0 { 431 w.WriteHeader(http.StatusNoContent) 432 return 433 } 434 435 for _, k := range pubKeys { 436 key := strings.TrimRight(k.Key, "\n") 437 fmt.Fprintln(w, key) 438 } 439} 440 441func (s *State) NewRepo(w http.ResponseWriter, r *http.Request) { 442 switch r.Method { 443 case http.MethodGet: 444 user := s.oauth.GetMultiAccountUser(r) 445 knots, err := s.enforcer.GetKnotsForUser(user.Did) 446 if err != nil { 447 s.pages.Notice(w, "repo", "Invalid user account.") 448 return 449 } 450 451 s.pages.NewRepo(w, pages.NewRepoParams{ 452 LoggedInUser: user, 453 Knots: knots, 454 }) 455 456 case http.MethodPost: 457 l := s.logger.With("handler", "NewRepo") 458 459 user := s.oauth.GetMultiAccountUser(r) 460 l = l.With("did", user.Did) 461 462 // form validation 463 domain := r.FormValue("domain") 464 if domain == "" { 465 s.pages.Notice(w, "repo", "Invalid form submission&mdash;missing knot domain.") 466 return 467 } 468 l = l.With("knot", domain) 469 470 repoName := r.FormValue("name") 471 if repoName == "" { 472 s.pages.Notice(w, "repo", "Repository name cannot be empty.") 473 return 474 } 475 476 if err := models.ValidateRepoName(repoName); err != nil { 477 s.pages.Notice(w, "repo", err.Error()) 478 return 479 } 480 repoName = models.StripGitExt(repoName) 481 rkey := strings.ToLower(repoName) 482 l = l.With("repoName", repoName, "rkey", rkey) 483 484 defaultBranch := r.FormValue("branch") 485 if defaultBranch == "" { 486 defaultBranch = "main" 487 } 488 l = l.With("defaultBranch", defaultBranch) 489 490 description := r.FormValue("description") 491 if len([]rune(description)) > 140 { 492 s.pages.Notice(w, "repo", "Description must be 140 characters or fewer.") 493 return 494 } 495 496 // ACL validation 497 ok, err := s.enforcer.E.Enforce(user.Did, domain, domain, "repo:create") 498 if err != nil || !ok { 499 l.Info("unauthorized") 500 s.pages.Notice(w, "repo", "You do not have permission to create a repo in this knot.") 501 return 502 } 503 504 // Check for existing repos 505 existingRepo, err := db.GetRepo( 506 s.db, 507 orm.FilterEq("did", user.Did), 508 orm.FilterEq("rkey", rkey), 509 ) 510 if err == nil && existingRepo != nil { 511 l.Info("repo exists") 512 s.pages.Notice(w, "repo", fmt.Sprintf("You already have a repository by this name on %s", existingRepo.Knot)) 513 return 514 } 515 516 atpClient, err := s.oauth.AuthorizedClient(r) 517 if err != nil { 518 l.Error("failed to get authorized client", "err", err) 519 s.pages.Notice(w, "repo", "Failed to authorize. Try again later.") 520 return 521 } 522 523 if rkeyOccupied(r.Context(), atpClient, user.Did, rkey) { 524 l.Info("rkey occupied by prior rename alias") 525 s.pages.Notice(w, "repo", fmt.Sprintf("The name %q still has a record on your PDS from a prior rename. Pick a different name, or delete at://%s/%s/%s first.", rkey, user.Did, tangled.RepoNSID, rkey)) 526 return 527 } 528 529 client, err := s.oauth.ServiceClient( 530 r, 531 oauth.WithService(domain), 532 oauth.WithLxm(tangled.RepoCreateNSID), 533 oauth.WithDev(s.config.Core.Dev), 534 ) 535 if err != nil { 536 l.Error("service auth failed", "err", err) 537 s.pages.Notice(w, "repo", "Failed to authenticate. Please log out and log back in again.") 538 return 539 } 540 541 input := &tangled.RepoCreate_Input{ 542 Rkey: rkey, 543 Name: rkey, 544 DefaultBranch: &defaultBranch, 545 } 546 createResp, err := tangled.RepoCreate( 547 r.Context(), 548 client, 549 input, 550 ) 551 if xrpcerr := xrpcclient.HandleXrpcErr(err); xrpcerr != nil { 552 l.Error("failed to call XRPC repo.create", "xrpcerr", xrpcerr, "err", err) 553 s.pages.Notice(w, "repo", err.Error()) 554 return 555 } 556 557 var repoDid string 558 if createResp != nil && createResp.RepoDid != nil { 559 repoDid = *createResp.RepoDid 560 } 561 if repoDid == "" { 562 l.Error("knot returned empty repo DID") 563 s.pages.Notice(w, "repo", "Knot failed to mint a repo DID. The knot may need to be upgraded.") 564 return 565 } 566 567 repo := &models.Repo{ 568 Did: user.Did, 569 Name: repoName, 570 Knot: domain, 571 Rkey: rkey, 572 Description: description, 573 Created: time.Now(), 574 Labels: s.config.Label.DefaultLabelDefs, 575 RepoDid: repoDid, 576 } 577 record := repo.AsRecord() 578 579 cleanupKnot := func() { 580 go func() { 581 delays := []time.Duration{0, 2 * time.Second, 5 * time.Second} 582 for attempt, delay := range delays { 583 time.Sleep(delay) 584 deleteClient, dErr := s.oauth.ServiceClient( 585 r, 586 oauth.WithService(domain), 587 oauth.WithLxm(tangled.RepoDeleteNSID), 588 oauth.WithDev(s.config.Core.Dev), 589 ) 590 if dErr != nil { 591 l.Error("failed to create delete client for knot cleanup", "attempt", attempt+1, "err", dErr) 592 continue 593 } 594 ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) 595 if dErr := tangled.RepoDelete(ctx, deleteClient, &tangled.RepoDelete_Input{ 596 Did: user.Did, 597 Name: rkey, 598 Rkey: rkey, 599 }); dErr != nil { 600 cancel() 601 l.Error("failed to clean up repo on knot after rollback", "attempt", attempt+1, "err", dErr) 602 continue 603 } 604 cancel() 605 l.Info("successfully cleaned up repo on knot after rollback", "attempt", attempt+1) 606 return 607 } 608 l.Error("exhausted retries for knot cleanup, repo may be orphaned", 609 "did", user.Did, "repo", repoName, "knot", domain) 610 }() 611 } 612 613 _, err = comatproto.RepoCreateRecord(r.Context(), atpClient, &comatproto.RepoCreateRecord_Input{ 614 Collection: tangled.RepoNSID, 615 Repo: user.Did, 616 Rkey: &rkey, 617 Record: &lexutil.LexiconTypeDecoder{ 618 Val: &record, 619 }, 620 }) 621 if err != nil { 622 l.Info("PDS write failed", "err", err) 623 cleanupKnot() 624 if rkeyOccupied(r.Context(), atpClient, user.Did, rkey) { 625 s.pages.Notice(w, "repo", fmt.Sprintf("You already have a repository named %q.", rkey)) 626 } else { 627 s.pages.Notice(w, "repo", "Failed to announce repository creation.") 628 } 629 return 630 } 631 632 aturi := fmt.Sprintf("at://%s/%s/%s", user.Did, tangled.RepoNSID, rkey) 633 l = l.With("aturi", aturi) 634 l.Info("wrote to PDS") 635 636 tx, err := s.db.BeginTx(r.Context(), nil) 637 if err != nil { 638 l.Info("txn failed", "err", err) 639 s.pages.Notice(w, "repo", "Failed to save repository information.") 640 return 641 } 642 643 rollback := func() { 644 err1 := tx.Rollback() 645 err2 := s.enforcer.E.LoadPolicy() 646 err3 := rollbackRecord(context.Background(), aturi, atpClient) 647 648 if errors.Is(err1, sql.ErrTxDone) { 649 err1 = nil 650 } 651 652 if errs := errors.Join(err1, err2, err3); errs != nil { 653 l.Error("failed to rollback changes", "errs", errs) 654 } 655 656 if aturi != "" { 657 cleanupKnot() 658 } 659 } 660 defer rollback() 661 662 err = db.AddRepo(tx, repo) 663 if err != nil { 664 l.Error("db write failed", "err", err) 665 s.pages.Notice(w, "repo", "Failed to save repository information.") 666 return 667 } 668 669 rbacPath := repo.RepoIdentifier() 670 err = s.enforcer.AddRepo(user.Did, domain, rbacPath) 671 if err != nil { 672 l.Error("acl setup failed", "err", err) 673 s.pages.Notice(w, "repo", "Failed to set up repository permissions.") 674 return 675 } 676 677 err = tx.Commit() 678 if err != nil { 679 l.Error("txn commit failed", "err", err) 680 http.Error(w, err.Error(), http.StatusInternalServerError) 681 return 682 } 683 684 err = s.enforcer.E.SavePolicy() 685 if err != nil { 686 l.Error("acl save failed", "err", err) 687 http.Error(w, err.Error(), http.StatusInternalServerError) 688 return 689 } 690 691 aturi = "" 692 693 s.notifier.NewRepo(r.Context(), repo) 694 switch { 695 case repoDid != "": 696 s.pages.HxLocation(w, fmt.Sprintf("/%s", repoDid)) 697 default: 698 handle := s.pages.DisplayHandle(r.Context(), user.Did) 699 s.pages.HxLocation(w, fmt.Sprintf("/%s/%s", handle, rkey)) 700 } 701 } 702} 703 704func rkeyOccupied(ctx context.Context, client *atclient.APIClient, did, rkey string) bool { 705 probeCtx, cancel := context.WithTimeout(ctx, 10*time.Second) 706 defer cancel() 707 resp, err := comatproto.RepoGetRecord(probeCtx, client, "", tangled.RepoNSID, did, rkey) 708 return err == nil && resp != nil 709} 710 711// this is used to rollback changes made to the PDS 712// 713// it is a no-op if the provided ATURI is empty 714func rollbackRecord(ctx context.Context, aturi string, client *atclient.APIClient) error { 715 if aturi == "" { 716 return nil 717 } 718 719 parsed := syntax.ATURI(aturi) 720 721 collection := parsed.Collection().String() 722 repo := parsed.Authority().String() 723 rkey := parsed.RecordKey().String() 724 725 _, err := comatproto.RepoDeleteRecord(ctx, client, &comatproto.RepoDeleteRecord_Input{ 726 Collection: collection, 727 Repo: repo, 728 Rkey: rkey, 729 }) 730 return err 731} 732 733func BackfillDefaultDefs(e db.Execer, r *idresolver.Resolver, defaults []string) error { 734 defaultLabels, err := db.GetLabelDefinitions(e, orm.FilterIn("at_uri", defaults)) 735 if err != nil { 736 return err 737 } 738 // already present 739 if len(defaultLabels) == len(defaults) { 740 return nil 741 } 742 743 labelDefs, err := models.FetchLabelDefs(r, defaults) 744 if err != nil { 745 return err 746 } 747 748 // Insert each label definition to the database 749 for _, labelDef := range labelDefs { 750 _, err = db.AddLabelDefinition(e, &labelDef) 751 if err != nil { 752 return fmt.Errorf("failed to add label definition %s: %v", labelDef.Name, err) 753 } 754 } 755 756 return nil 757} 758 759func fetchBskyPosts(ctx context.Context, res *idresolver.Resolver, config *config.Config, d *db.DB, logger *slog.Logger) { 760 resolved, err := res.ResolveIdent(context.Background(), consts.TangledDid) 761 if err != nil { 762 logger.Error("failed to resolve tangled.org DID", "err", err) 763 return 764 } 765 766 pdsEndpoint := resolved.PDSEndpoint() 767 if pdsEndpoint == "" { 768 logger.Error("no PDS endpoint found for tangled.sh DID") 769 return 770 } 771 772 session, err := oauth.CreateAppPasswordSession(res, config.Core.AppPassword, consts.TangledDid, logger) 773 if err != nil { 774 logger.Error("failed to create appassword session... skipping fetch", "err", err) 775 return 776 } 777 778 l := log.SubLogger(logger, "bluesky") 779 780 ticker := time.NewTicker(config.Bluesky.UpdateInterval) 781 defer ticker.Stop() 782 783 for { 784 // refresh session if necessary 785 if !session.IsValid() { 786 l.Debug("access token expired, refreshing session") 787 if err := session.RefreshSession(); err != nil { 788 l.Error("failed to refresh session, stopping bluesky updater", "err", err) 789 return 790 } 791 l.Debug("session refreshed") 792 } 793 794 // make client 795 client := xrpc.Client{ 796 Auth: &xrpc.AuthInfo{ 797 AccessJwt: session.AccessJwt, 798 Did: session.Did, 799 }, 800 Host: session.PdsEndpoint, 801 } 802 803 posts, _, err := bsky.FetchPosts(ctx, &client, 20, "") 804 if err != nil { 805 l.Error("failed to fetch bluesky posts", "err", err) 806 } else if err := db.InsertBlueskyPosts(d, posts); err != nil { 807 l.Error("failed to insert bluesky posts", "err", err) 808 } else { 809 l.Info("inserted bluesky posts", "count", len(posts)) 810 } 811 812 select { 813 case <-ticker.C: 814 case <-ctx.Done(): 815 l.Info("stopping bluesky updater") 816 return 817 } 818 } 819}