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