This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

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