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