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