Monorepo for Tangled
tangled.org
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—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}