This repository has no description
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—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, xe := tangled.RepoCreate(
483 r.Context(),
484 client,
485 input,
486 )
487 if err := xrpcclient.HandleXrpcErr(xe); err != nil {
488 l.Error("xrpc error", "xe", xe)
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}