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 emaildispatch "tangled.org/core/appview/notify/email"
32 whnotify "tangled.org/core/appview/notify/webhook"
33 "tangled.org/core/appview/oauth"
34 "tangled.org/core/appview/pages"
35 pipelinessh "tangled.org/core/appview/pipelines/ssh"
36 "tangled.org/core/appview/reporesolver"
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 "tangled.org/core/repoverify"
46 xrpcclient "tangled.org/core/xrpc/xrpcclient"
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 logger *slog.Logger
75 cfClient *cloudflare.Client
76 codesearch *codesearch.CodeSearch
77}
78
79func Make(ctx context.Context, config *config.Config) (*State, error) {
80 logger := tlog.FromContext(ctx)
81
82 d, err := db.Make(ctx, config.Core.DbPath)
83 if err != nil {
84 return nil, fmt.Errorf("failed to create db: %w", err)
85 }
86
87 indexer := indexer.New(log.SubLogger(logger, "indexer"), d)
88 err = indexer.Init(ctx)
89 if err != nil {
90 return nil, fmt.Errorf("failed to create indexer: %w", err)
91 }
92
93 enforcer, err := rbac.NewEnforcer(config.Core.DbPath)
94 if err != nil {
95 return nil, fmt.Errorf("failed to create enforcer: %w", err)
96 }
97
98 res, err := idresolver.RedisResolver(config.Redis.ToURL(), config.Plc.PLCURL)
99 if err != nil {
100 logger.Error("failed to create redis resolver", "err", err)
101 res = idresolver.DefaultResolver(config.Plc.PLCURL)
102 }
103
104 var rdb *cache.Cache
105 if config.Redis.Addr != "" {
106 rdb = cache.New(config.Redis.Addr)
107 }
108
109 posthog, err := posthog.NewWithConfig(config.Posthog.ApiKey, posthog.Config{Endpoint: config.Posthog.Endpoint})
110 if err != nil {
111 return nil, fmt.Errorf("failed to create posthog client: %w", err)
112 }
113
114 pages := pages.NewPages(config, res, d, rdb, log.SubLogger(logger, "pages"))
115 knotcompat.UseNativeLatch(knotacl.NewLatch(d, log.SubLogger(logger, "knotacl-latch")))
116 aclService := knotacl.NewService(enforcer, d, config.Core.Dev, log.SubLogger(logger, "knotacl"))
117 oauth, err := oauth.New(config, posthog, d, enforcer, aclService, res, log.SubLogger(logger, "oauth"))
118 if err != nil {
119 return nil, fmt.Errorf("failed to start oauth handler: %w", err)
120 }
121 repoResolver := reporesolver.New(config, aclService, d, rdb)
122
123 mentionsResolver := mentions.New(config, res, d, log.SubLogger(logger, "mentionsResolver"))
124
125 jc, err := jetstream.NewJetstreamClient(
126 config.Jetstream.Endpoint,
127 "appview",
128 []string{
129 tangled.ActorProfileNSID,
130 tangled.FeedStarNSID,
131 tangled.FeedReactionNSID,
132 tangled.FeedCommentNSID,
133 tangled.GraphFollowNSID,
134 tangled.GraphVouchNSID,
135 tangled.KnotMemberNSID,
136 tangled.KnotNSID,
137 tangled.LabelDefinitionNSID,
138 tangled.LabelOpNSID,
139 tangled.PublicKeyNSID,
140 tangled.RepoArtifactNSID,
141 tangled.RepoIssueCommentNSID,
142 tangled.RepoIssueNSID,
143 tangled.RepoIssueStateNSID,
144 tangled.RepoNSID,
145 tangled.RepoPullNSID,
146 tangled.RepoPullStatusNSID,
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 go ingester.StartPendingStateReconciler()
205
206 if config.Resend.ApiKey != "" {
207 dispatcher := emaildispatch.NewDispatcher(d, config.Resend, config.Core.BaseUrl(), res, tlog.SubLogger(logger, "email-dispatcher"), config.Core.Dev)
208 go dispatcher.Start(ctx)
209 }
210
211 var cfClient *cloudflare.Client
212 if config.Cloudflare.ApiToken != "" {
213 cfClient, err = cloudflare.New(config)
214 if err != nil {
215 logger.Warn("failed to create cloudflare client, sites upload will be disabled", "err", err)
216 cfClient = nil
217 }
218 }
219
220 knotstream, err := Knotstream(ctx, config, d, aclService, enforcer, posthog, notifier, cfClient)
221 if err != nil {
222 return nil, fmt.Errorf("failed to start knotstream consumer: %w", err)
223 }
224 knotstream.Start(ctx)
225
226 state := &State{
227 db: d,
228 notifier: notifier,
229 indexer: indexer,
230 oauth: oauth,
231 enforcer: enforcer,
232 pages: pages,
233 idResolver: res,
234 rdb: rdb,
235 mentionsResolver: mentionsResolver,
236 posthog: posthog,
237 jc: jc,
238 config: config,
239 repoResolver: repoResolver,
240 aclService: aclService,
241 knotstream: knotstream,
242 logger: logger,
243 cfClient: cfClient,
244 codesearch: &codesearch.CodeSearch{Host: config.CodeSearch.ZoektUrl},
245 }
246
247 // fetch initial bluesky posts if configured
248 go fetchBskyPosts(ctx, res, config, d, logger)
249
250 return state, nil
251}
252
253func (s *State) Close() error {
254 // other close up logic goes here
255 return s.db.Close()
256}
257
258func (s *State) NewSSHServer() *pipelinessh.Server {
259 return pipelinessh.New(s.db, s.config, log.SubLogger(s.logger, "pipelinessh"))
260}
261
262func (s *State) SecurityTxt(w http.ResponseWriter, r *http.Request) {
263 w.Header().Set("Content-Type", "text/plain")
264 w.Header().Set("Cache-Control", "public, max-age=86400") // one day
265
266 securityTxt := `Contact: mailto:security@tangled.org
267Preferred-Languages: en
268Canonical: https://tangled.org/.well-known/security.txt
269Expires: 2030-01-01T21:59:00.000Z
270`
271 w.Write([]byte(securityTxt))
272}
273
274func (s *State) RobotsTxt(w http.ResponseWriter, r *http.Request) {
275 w.Header().Set("Content-Type", "text/plain")
276 w.Header().Set("Cache-Control", "public, max-age=86400") // one day
277
278 robotsTxt := `# Hello, Tanglers!
279User-agent: *
280Allow: /
281Disallow: /*/*/settings
282Disallow: /settings
283Disallow: /*/*/compare
284Disallow: /*/*/fork
285Disallow: /search
286Disallow: /*/*/search
287
288Crawl-delay: 1
289`
290 w.Write([]byte(robotsTxt))
291}
292
293func (s *State) TermsOfService(w http.ResponseWriter, r *http.Request) {
294 s.pages.TermsOfService(w, pages.TermsOfServiceParams{
295 BaseParams: pages.BaseParamsFromContext(r.Context()),
296 })
297}
298
299func (s *State) PrivacyPolicy(w http.ResponseWriter, r *http.Request) {
300 s.pages.PrivacyPolicy(w, pages.PrivacyPolicyParams{
301 BaseParams: pages.BaseParamsFromContext(r.Context()),
302 })
303}
304
305func (s *State) Brand(w http.ResponseWriter, r *http.Request) {
306 s.pages.Brand(w, pages.BrandParams{
307 BaseParams: pages.BaseParamsFromContext(r.Context()),
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 user = strings.TrimSuffix(user, ".keys")
416
417 if user == "" {
418 w.WriteHeader(http.StatusBadRequest)
419 return
420 }
421
422 id, err := s.idResolver.ResolveIdent(r.Context(), user)
423 if err != nil {
424 w.WriteHeader(http.StatusInternalServerError)
425 return
426 }
427
428 pubKeys, err := db.GetPublicKeysForDid(s.db, id.DID.String())
429 if err != nil {
430 s.logger.Error("failed to get public keys", "err", err)
431 http.Error(w, "failed to get public keys", http.StatusInternalServerError)
432 return
433 }
434
435 if len(pubKeys) == 0 {
436 w.WriteHeader(http.StatusNoContent)
437 return
438 }
439
440 for _, k := range pubKeys {
441 key := strings.TrimRight(k.Key, "\n")
442 fmt.Fprintln(w, key)
443 }
444}
445
446func (s *State) NewRepo(w http.ResponseWriter, r *http.Request) {
447 switch r.Method {
448 case http.MethodGet:
449 user := s.oauth.GetMultiAccountUser(r)
450 knots := s.aclService.KnotsForUser(r.Context(), user.Did)
451
452 s.pages.NewRepo(w, pages.NewRepoParams{
453 BaseParams: pages.BaseParamsFromContext(r.Context()),
454 Knots: knots,
455 })
456
457 case http.MethodPost:
458 l := s.logger.With("handler", "NewRepo")
459
460 user := s.oauth.GetMultiAccountUser(r)
461 l = l.With("did", user.Did)
462
463 // form validation
464 domain := r.FormValue("domain")
465 if domain == "" {
466 s.pages.Notice(w, "repo", "Invalid form submission—missing knot domain.")
467 return
468 }
469 l = l.With("knot", domain)
470
471 repoName := r.FormValue("name")
472 if repoName == "" {
473 s.pages.Notice(w, "repo", "Repository name cannot be empty.")
474 return
475 }
476
477 if err := models.ValidateRepoName(repoName); err != nil {
478 s.pages.Notice(w, "repo", err.Error())
479 return
480 }
481 repoName = models.StripGitExt(repoName)
482 rkey := strings.ToLower(repoName)
483 l = l.With("repoName", repoName, "rkey", rkey)
484
485 defaultBranch := r.FormValue("branch")
486 if defaultBranch == "" {
487 defaultBranch = "main"
488 }
489 l = l.With("defaultBranch", defaultBranch)
490
491 description := r.FormValue("description")
492 if len([]rune(description)) > 140 {
493 s.pages.Notice(w, "repo", "Description must be 140 characters or fewer.")
494 return
495 }
496
497 // ACL validation
498 if !s.aclService.IsRepoCreateAllowed(r.Context(), domain, user.Did) {
499 l.Info("unauthorized")
500 s.pages.Notice(w, "repo", "You do not have permission to create a repo in this knot.")
501 return
502 }
503
504 // Check for existing repos
505 existingRepo, err := db.GetRepo(
506 s.db,
507 orm.FilterEq("did", user.Did),
508 orm.FilterEq("rkey", rkey),
509 )
510 if err == nil && existingRepo != nil {
511 l.Info("repo exists")
512 s.pages.Notice(w, "repo", fmt.Sprintf("You already have a repository by this name on %s", existingRepo.Knot))
513 return
514 }
515
516 atpClient, err := s.oauth.AuthorizedClient(r)
517 if err != nil {
518 l.Error("failed to get authorized client", "err", err)
519 s.pages.Notice(w, "repo", "Failed to authorize. Try again later.")
520 return
521 }
522
523 if rkeyOccupied(r.Context(), atpClient, user.Did, rkey) {
524 l.Info("rkey occupied by prior rename alias")
525 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))
526 return
527 }
528
529 client, err := s.oauth.ServiceClient(
530 r,
531 oauth.WithService(domain),
532 oauth.WithLxm(tangled.RepoCreateNSID),
533 oauth.WithDev(s.config.Core.Dev),
534 )
535 if err != nil {
536 l.Error("service auth failed", "err", err)
537 s.pages.Notice(w, "repo", "Failed to authenticate. Please log out and log back in again.")
538 return
539 }
540
541 input := &tangled.RepoCreate_Input{
542 Rkey: rkey,
543 Name: rkey,
544 DefaultBranch: &defaultBranch,
545 }
546 createResp, err := tangled.RepoCreate(
547 r.Context(),
548 client,
549 input,
550 )
551 if xrpcerr := xrpcclient.HandleXrpcErr(err); xrpcerr != nil {
552 l.Error("failed to call XRPC repo.create", "xrpcerr", xrpcerr, "err", err)
553 s.pages.Notice(w, "repo", err.Error())
554 return
555 }
556
557 var repoDid string
558 if createResp != nil && createResp.RepoDid != nil {
559 repoDid = *createResp.RepoDid
560 }
561 if repoDid == "" {
562 l.Error("knot returned empty repo DID")
563 s.pages.Notice(w, "repo", "Knot failed to mint a repo DID. The knot may need to be upgraded.")
564 return
565 }
566
567 repo := &models.Repo{
568 Did: user.Did,
569 Name: repoName,
570 Knot: domain,
571 Rkey: rkey,
572 Description: description,
573 Created: time.Now(),
574 Labels: s.config.Label.DefaultLabelDefs,
575 RepoDid: repoDid,
576 }
577 record := repo.AsRecord()
578
579 cleanupKnot := func() {
580 go func() {
581 delays := []time.Duration{0, 2 * time.Second, 5 * time.Second}
582 for attempt, delay := range delays {
583 time.Sleep(delay)
584 deleteClient, dErr := s.oauth.ServiceClient(
585 r,
586 oauth.WithService(domain),
587 oauth.WithLxm(tangled.RepoDeleteNSID),
588 oauth.WithDev(s.config.Core.Dev),
589 )
590 if dErr != nil {
591 l.Error("failed to create delete client for knot cleanup", "attempt", attempt+1, "err", dErr)
592 continue
593 }
594 ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
595 if dErr := tangled.RepoDelete(ctx, deleteClient, &tangled.RepoDelete_Input{
596 Did: user.Did,
597 Name: rkey,
598 Rkey: rkey,
599 }); dErr != nil {
600 cancel()
601 l.Error("failed to clean up repo on knot after rollback", "attempt", attempt+1, "err", dErr)
602 continue
603 }
604 cancel()
605 l.Info("successfully cleaned up repo on knot after rollback", "attempt", attempt+1)
606 return
607 }
608 l.Error("exhausted retries for knot cleanup, repo may be orphaned",
609 "did", user.Did, "repo", repoName, "knot", domain)
610 }()
611 }
612
613 _, err = comatproto.RepoCreateRecord(r.Context(), atpClient, &comatproto.RepoCreateRecord_Input{
614 Collection: tangled.RepoNSID,
615 Repo: user.Did,
616 Rkey: &rkey,
617 Record: &lexutil.LexiconTypeDecoder{
618 Val: &record,
619 },
620 })
621 if err != nil {
622 l.Info("PDS write failed", "err", err)
623 cleanupKnot()
624 if rkeyOccupied(r.Context(), atpClient, user.Did, rkey) {
625 s.pages.Notice(w, "repo", fmt.Sprintf("You already have a repository named %q.", rkey))
626 } else {
627 s.pages.Notice(w, "repo", "Failed to announce repository creation.")
628 }
629 return
630 }
631
632 aturi := fmt.Sprintf("at://%s/%s/%s", user.Did, tangled.RepoNSID, rkey)
633 l = l.With("aturi", aturi)
634 l.Info("wrote to PDS")
635
636 tx, err := s.db.BeginTx(r.Context(), nil)
637 if err != nil {
638 l.Info("txn failed", "err", err)
639 s.pages.Notice(w, "repo", "Failed to save repository information.")
640 return
641 }
642
643 rollback := func() {
644 err1 := tx.Rollback()
645 err2 := s.enforcer.E.LoadPolicy()
646 err3 := rollbackRecord(context.Background(), aturi, atpClient)
647
648 if errors.Is(err1, sql.ErrTxDone) {
649 err1 = nil
650 }
651
652 if errs := errors.Join(err1, err2, err3); errs != nil {
653 l.Error("failed to rollback changes", "errs", errs)
654 }
655
656 if aturi != "" {
657 cleanupKnot()
658 }
659 }
660 defer rollback()
661
662 err = db.AddRepo(tx, repo)
663 if err != nil {
664 l.Error("db write failed", "err", err)
665 s.pages.Notice(w, "repo", "Failed to save repository information.")
666 return
667 }
668
669 rbacPath := repo.RepoIdentifier()
670 err = s.enforcer.AddRepo(user.Did, domain, rbacPath)
671 if err != nil {
672 l.Error("acl setup failed", "err", err)
673 s.pages.Notice(w, "repo", "Failed to set up repository permissions.")
674 return
675 }
676
677 err = tx.Commit()
678 if err != nil {
679 l.Error("txn commit failed", "err", err)
680 http.Error(w, err.Error(), http.StatusInternalServerError)
681 return
682 }
683
684 err = s.enforcer.E.SavePolicy()
685 if err != nil {
686 l.Error("acl save failed", "err", err)
687 http.Error(w, err.Error(), http.StatusInternalServerError)
688 return
689 }
690
691 aturi = ""
692
693 s.notifier.NewRepo(r.Context(), repo)
694 switch {
695 case repoDid != "":
696 s.pages.HxLocation(w, fmt.Sprintf("/%s", repoDid))
697 default:
698 handle := s.pages.DisplayHandle(r.Context(), user.Did)
699 s.pages.HxLocation(w, fmt.Sprintf("/%s/%s", handle, rkey))
700 }
701 }
702}
703
704func rkeyOccupied(ctx context.Context, client *atclient.APIClient, did, rkey string) bool {
705 probeCtx, cancel := context.WithTimeout(ctx, 10*time.Second)
706 defer cancel()
707 resp, err := comatproto.RepoGetRecord(probeCtx, client, "", tangled.RepoNSID, did, rkey)
708 return err == nil && resp != nil
709}
710
711// this is used to rollback changes made to the PDS
712//
713// it is a no-op if the provided ATURI is empty
714func rollbackRecord(ctx context.Context, aturi string, client *atclient.APIClient) error {
715 if aturi == "" {
716 return nil
717 }
718
719 parsed := syntax.ATURI(aturi)
720
721 collection := parsed.Collection().String()
722 repo := parsed.Authority().String()
723 rkey := parsed.RecordKey().String()
724
725 _, err := comatproto.RepoDeleteRecord(ctx, client, &comatproto.RepoDeleteRecord_Input{
726 Collection: collection,
727 Repo: repo,
728 Rkey: rkey,
729 })
730 return err
731}
732
733func BackfillDefaultDefs(e db.Execer, r *idresolver.Resolver, defaults []string) error {
734 defaultLabels, err := db.GetLabelDefinitions(e, orm.FilterIn("at_uri", defaults))
735 if err != nil {
736 return err
737 }
738 // already present
739 if len(defaultLabels) == len(defaults) {
740 return nil
741 }
742
743 labelDefs, err := models.FetchLabelDefs(r, defaults)
744 if err != nil {
745 return err
746 }
747
748 // Insert each label definition to the database
749 for _, labelDef := range labelDefs {
750 _, err = db.AddLabelDefinition(e, &labelDef)
751 if err != nil {
752 return fmt.Errorf("failed to add label definition %s: %v", labelDef.Name, err)
753 }
754 }
755
756 return nil
757}
758
759func fetchBskyPosts(ctx context.Context, res *idresolver.Resolver, config *config.Config, d *db.DB, logger *slog.Logger) {
760 resolved, err := res.ResolveIdent(context.Background(), consts.TangledDid)
761 if err != nil {
762 logger.Error("failed to resolve tangled.org DID", "err", err)
763 return
764 }
765
766 pdsEndpoint := resolved.PDSEndpoint()
767 if pdsEndpoint == "" {
768 logger.Error("no PDS endpoint found for tangled.sh DID")
769 return
770 }
771
772 session, err := oauth.CreateAppPasswordSession(res, config.Core.AppPassword, consts.TangledDid, logger)
773 if err != nil {
774 logger.Error("failed to create appassword session... skipping fetch", "err", err)
775 return
776 }
777
778 l := log.SubLogger(logger, "bluesky")
779
780 ticker := time.NewTicker(config.Bluesky.UpdateInterval)
781 defer ticker.Stop()
782
783 for {
784 // refresh session if necessary
785 if !session.IsValid() {
786 l.Debug("access token expired, refreshing session")
787 if err := session.RefreshSession(); err != nil {
788 l.Error("failed to refresh session, stopping bluesky updater", "err", err)
789 return
790 }
791 l.Debug("session refreshed")
792 }
793
794 // make client
795 client := xrpc.Client{
796 Auth: &xrpc.AuthInfo{
797 AccessJwt: session.AccessJwt,
798 Did: session.Did,
799 },
800 Host: session.PdsEndpoint,
801 }
802
803 posts, _, err := bsky.FetchPosts(ctx, &client, 20, "")
804 if err != nil {
805 l.Error("failed to fetch bluesky posts", "err", err)
806 } else if err := db.InsertBlueskyPosts(d, posts); err != nil {
807 l.Error("failed to insert bluesky posts", "err", err)
808 } else {
809 l.Info("inserted bluesky posts", "count", len(posts))
810 }
811
812 select {
813 case <-ticker.C:
814 case <-ctx.Done():
815 l.Info("stopping bluesky updater")
816 return
817 }
818 }
819}