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