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
289Disallow: /search
290Disallow: /*/*/search
291
292Crawl-delay: 1
293`
294 w.Write([]byte(robotsTxt))
295}
296
297func (s *State) TermsOfService(w http.ResponseWriter, r *http.Request) {
298 s.pages.TermsOfService(w, pages.TermsOfServiceParams{
299 BaseParams: pages.BaseParamsFromContext(r.Context()),
300 })
301}
302
303func (s *State) PrivacyPolicy(w http.ResponseWriter, r *http.Request) {
304 s.pages.PrivacyPolicy(w, pages.PrivacyPolicyParams{
305 BaseParams: pages.BaseParamsFromContext(r.Context()),
306 })
307}
308
309func (s *State) Brand(w http.ResponseWriter, r *http.Request) {
310 s.pages.Brand(w, pages.BrandParams{
311 BaseParams: pages.BaseParamsFromContext(r.Context()),
312 })
313}
314
315func (s *State) UpgradeBanner(w http.ResponseWriter, r *http.Request) {
316 user := s.oauth.GetMultiAccountUser(r)
317 if user == nil {
318 return
319 }
320
321 l := s.logger.With("handler", "UpgradeBanner")
322 l = l.With("did", user.Did)
323
324 regs, err := db.GetRegistrations(
325 s.db,
326 orm.FilterEq("did", user.Did),
327 orm.FilterEq("needs_upgrade", 1),
328 )
329 if err != nil {
330 l.Error("non-fatal: failed to get registrations", "err", err)
331 }
332
333 spindles, err := db.GetSpindles(
334 r.Context(),
335 s.db,
336 orm.FilterEq("owner", user.Did),
337 orm.FilterEq("needs_upgrade", 1),
338 )
339 if err != nil {
340 l.Error("non-fatal: failed to get spindles", "err", err)
341 }
342
343 if regs == nil && spindles == nil {
344 return
345 }
346
347 s.pages.UpgradeBanner(w, pages.UpgradeBannerParams{
348 Registrations: regs,
349 Spindles: spindles,
350 })
351}
352
353func (s *State) NewsletterSignup(w http.ResponseWriter, r *http.Request) {
354 // target is echoed back from the form via hx-vals so the response span's
355 // id matches the form's hx-target. Fallback keeps the handler useful if
356 // a caller forgets to send it.
357 target := strings.TrimSpace(r.FormValue("target"))
358 if target == "" {
359 target = "home"
360 }
361
362 w.Header().Set("Content-Type", "text/html")
363
364 emailAddr := strings.TrimSpace(r.FormValue("email"))
365 if !email.IsValidEmail(emailAddr) {
366 s.pages.NewsletterResponse(w, pages.NewsletterResponseParams{
367 Id: target,
368 Error: "Invalid email address.",
369 })
370 return
371 }
372
373 // For logged-in users, persist the signup locally so the widget stays
374 // hidden across devices. The DB row is the render-time source of truth;
375 // Resend still owns the mailing list itself.
376 if user := s.oauth.GetMultiAccountUser(r); user != nil {
377 if err := db.UpsertNewsletterPref(s.db, user.Did, db.NewsletterStatusSubscribed, emailAddr); err != nil {
378 s.logger.Error("failed to persist newsletter preference", "did", user.Did, "err", err)
379 }
380 }
381
382 if s.config.Resend.ApiKey != "" && s.config.Resend.NewsletterSegmentId != "" {
383 go func() {
384 if err := email.AddNewsletterContact(s.config.Resend.ApiKey, s.config.Resend.NewsletterSegmentId, emailAddr); err != nil {
385 s.logger.Error("failed to add newsletter contact", "error", err)
386 }
387 }()
388 } else {
389 s.logger.Error(
390 "failed to add newsletter contact, missing resend config",
391 "isKeyPresent", s.config.Resend.ApiKey != "",
392 "isSegmentIdPresent", s.config.Resend.NewsletterSegmentId != "",
393 "emailAddr", emailAddr,
394 )
395 }
396
397 s.pages.NewsletterResponse(w, pages.NewsletterResponseParams{Id: target})
398}
399
400// NewsletterDismiss records that a logged-in user has dismissed the newsletter
401// widget so it stays hidden across their devices. Anonymous callers get a 204
402// with no DB write — localStorage handles the per-browser fallback.
403func (s *State) NewsletterDismiss(w http.ResponseWriter, r *http.Request) {
404 user := s.oauth.GetMultiAccountUser(r)
405 if user == nil {
406 w.WriteHeader(http.StatusNoContent)
407 return
408 }
409
410 if err := db.UpsertNewsletterPref(s.db, user.Did, db.NewsletterStatusDismissed, ""); err != nil {
411 s.logger.Error("failed to persist newsletter dismissal", "did", user.Did, "err", err)
412 }
413 w.WriteHeader(http.StatusNoContent)
414}
415
416func (s *State) Keys(w http.ResponseWriter, r *http.Request) {
417 user := chi.URLParam(r, "user")
418 user = strings.TrimPrefix(user, "@")
419 user = strings.TrimSuffix(user, ".keys")
420
421 if user == "" {
422 w.WriteHeader(http.StatusBadRequest)
423 return
424 }
425
426 id, err := s.idResolver.ResolveIdent(r.Context(), user)
427 if err != nil {
428 w.WriteHeader(http.StatusInternalServerError)
429 return
430 }
431
432 pubKeys, err := db.GetPublicKeysForDid(s.db, id.DID.String())
433 if err != nil {
434 s.logger.Error("failed to get public keys", "err", err)
435 http.Error(w, "failed to get public keys", http.StatusInternalServerError)
436 return
437 }
438
439 if len(pubKeys) == 0 {
440 w.WriteHeader(http.StatusNoContent)
441 return
442 }
443
444 for _, k := range pubKeys {
445 key := strings.TrimRight(k.Key, "\n")
446 fmt.Fprintln(w, key)
447 }
448}
449
450func (s *State) NewRepo(w http.ResponseWriter, r *http.Request) {
451 switch r.Method {
452 case http.MethodGet:
453 user := s.oauth.GetMultiAccountUser(r)
454 knots := s.aclService.KnotsForUser(r.Context(), user.Did)
455
456 s.pages.NewRepo(w, pages.NewRepoParams{
457 BaseParams: pages.BaseParamsFromContext(r.Context()),
458 Knots: knots,
459 })
460
461 case http.MethodPost:
462 l := s.logger.With("handler", "NewRepo")
463
464 user := s.oauth.GetMultiAccountUser(r)
465 l = l.With("did", user.Did)
466
467 // form validation
468 domain := r.FormValue("domain")
469 if domain == "" {
470 s.pages.Notice(w, "repo", "Invalid form submission—missing knot domain.")
471 return
472 }
473 l = l.With("knot", domain)
474
475 repoName := r.FormValue("name")
476 if repoName == "" {
477 s.pages.Notice(w, "repo", "Repository name cannot be empty.")
478 return
479 }
480
481 if err := models.ValidateRepoName(repoName); err != nil {
482 s.pages.Notice(w, "repo", err.Error())
483 return
484 }
485 repoName = models.StripGitExt(repoName)
486 rkey := strings.ToLower(repoName)
487 l = l.With("repoName", repoName, "rkey", rkey)
488
489 defaultBranch := r.FormValue("branch")
490 if defaultBranch == "" {
491 defaultBranch = "main"
492 }
493 l = l.With("defaultBranch", defaultBranch)
494
495 description := r.FormValue("description")
496 if len([]rune(description)) > 140 {
497 s.pages.Notice(w, "repo", "Description must be 140 characters or fewer.")
498 return
499 }
500
501 // ACL validation
502 if !s.aclService.IsRepoCreateAllowed(r.Context(), domain, user.Did) {
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}