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