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