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