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