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