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