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