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