Monorepo for Tangled tangled.org
1

Configure Feed

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

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