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