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