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