Monorepo for Tangled tangled.org
1

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/codesearch" 19 "tangled.org/core/appview/config" 20 "tangled.org/core/appview/db" 21 "tangled.org/core/appview/email" 22 "tangled.org/core/appview/indexer" 23 "tangled.org/core/appview/knotacl" 24 "tangled.org/core/appview/knotcompat" 25 "tangled.org/core/appview/mentions" 26 "tangled.org/core/appview/models" 27 "tangled.org/core/appview/notify" 28 dbnotify "tangled.org/core/appview/notify/db" 29 lognotify "tangled.org/core/appview/notify/logging" 30 phnotify "tangled.org/core/appview/notify/posthog" 31 whnotify "tangled.org/core/appview/notify/webhook" 32 "tangled.org/core/appview/oauth" 33 "tangled.org/core/appview/pages" 34 pipelinessh "tangled.org/core/appview/pipelines/ssh" 35 "tangled.org/core/appview/reporesolver" 36 "tangled.org/core/consts" 37 "tangled.org/core/eventconsumer" 38 "tangled.org/core/idresolver" 39 "tangled.org/core/jetstream" 40 "tangled.org/core/log" 41 tlog "tangled.org/core/log" 42 "tangled.org/core/orm" 43 "tangled.org/core/rbac" 44 "tangled.org/core/repoverify" 45 xrpcclient "tangled.org/core/xrpc/xrpcclient" 46 47 comatproto "github.com/bluesky-social/indigo/api/atproto" 48 "github.com/bluesky-social/indigo/atproto/atclient" 49 "github.com/bluesky-social/indigo/atproto/syntax" 50 lexutil "github.com/bluesky-social/indigo/lex/util" 51 "github.com/bluesky-social/indigo/xrpc" 52 53 "github.com/go-chi/chi/v5" 54 "github.com/posthog/posthog-go" 55) 56 57type State struct { 58 db *db.DB 59 notifier notify.Notifier 60 indexer *indexer.Indexer 61 oauth *oauth.OAuth 62 enforcer *rbac.Enforcer 63 pages *pages.Pages 64 idResolver *idresolver.Resolver 65 rdb *cache.Cache 66 mentionsResolver *mentions.Resolver 67 posthog posthog.Client 68 jc *jetstream.JetstreamClient 69 config *config.Config 70 repoResolver *reporesolver.RepoResolver 71 aclService *knotacl.Service 72 knotstream *eventconsumer.Consumer 73 logger *slog.Logger 74 cfClient *cloudflare.Client 75 codesearch *codesearch.CodeSearch 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 knotcompat.UseNativeLatch(knotacl.NewLatch(d, log.SubLogger(logger, "knotacl-latch"))) 115 aclService := knotacl.NewService(enforcer, d, config.Core.Dev, log.SubLogger(logger, "knotacl")) 116 oauth, err := oauth.New(config, posthog, d, enforcer, aclService, res, log.SubLogger(logger, "oauth")) 117 if err != nil { 118 return nil, fmt.Errorf("failed to start oauth handler: %w", err) 119 } 120 repoResolver := reporesolver.New(config, aclService, 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.FeedReactionNSID, 131 tangled.FeedCommentNSID, 132 tangled.GraphFollowNSID, 133 tangled.GraphVouchNSID, 134 tangled.KnotMemberNSID, 135 tangled.KnotNSID, 136 tangled.LabelDefinitionNSID, 137 tangled.LabelOpNSID, 138 tangled.PublicKeyNSID, 139 tangled.RepoArtifactNSID, 140 tangled.RepoIssueCommentNSID, 141 tangled.RepoIssueNSID, 142 tangled.RepoIssueStateNSID, 143 tangled.RepoNSID, 144 tangled.RepoPullNSID, 145 tangled.RepoPullStatusNSID, 146 tangled.RepoPullCommentNSID, 147 tangled.SpindleMemberNSID, 148 tangled.SpindleNSID, 149 tangled.StringNSID, 150 }, 151 nil, 152 tlog.SubLogger(logger, "jetstream"), 153 d, 154 false, 155 156 // in-memory filter is inapplicable to appview so 157 // we'll never log dids anyway. 158 false, 159 ) 160 if err != nil { 161 return nil, fmt.Errorf("failed to create jetstream client: %w", err) 162 } 163 164 if err := BackfillDefaultDefs(d, res, config.Label.DefaultLabelDefs); err != nil { 165 return nil, fmt.Errorf("failed to backfill default label defs: %w", err) 166 } 167 168 var notifiers []notify.Notifier 169 170 // Always add the database notifier 171 notifiers = append(notifiers, dbnotify.NewDatabaseNotifier(d, res)) 172 173 // Add other notifiers in production only 174 if !config.Core.Dev { 175 notifiers = append(notifiers, phnotify.NewPosthogNotifier(posthog)) 176 } 177 notifiers = append(notifiers, indexer) 178 179 notifiers = append(notifiers, whnotify.NewNotifier(d)) 180 181 notifier := notify.NewMergedNotifier(notifiers) 182 notifier = lognotify.NewLoggingNotifier(notifier, tlog.SubLogger(logger, "notify")) 183 184 ingester := appview.Ingester{ 185 Ctx: ctx, 186 Db: d, 187 Enforcer: enforcer, 188 Acl: aclService, 189 IdResolver: res, 190 Cache: rdb, 191 Config: config, 192 Logger: log.SubLogger(logger, "ingester"), 193 MentionsResolver: mentionsResolver, 194 Notifier: notifier, 195 Verifier: repoverify.New(res, config.Core.Dev), 196 } 197 err = jc.StartJetstream(ctx, ingester.Ingest()) 198 if err != nil { 199 return nil, fmt.Errorf("failed to start jetstream watcher: %w", err) 200 } 201 202 go ingester.SweepPendingVerifications() 203 go ingester.StartPendingStateReconciler() 204 205 var cfClient *cloudflare.Client 206 if config.Cloudflare.ApiToken != "" { 207 cfClient, err = cloudflare.New(config) 208 if err != nil { 209 logger.Warn("failed to create cloudflare client, sites upload will be disabled", "err", err) 210 cfClient = nil 211 } 212 } 213 214 knotstream, err := Knotstream(ctx, config, d, aclService, enforcer, posthog, notifier, cfClient) 215 if err != nil { 216 return nil, fmt.Errorf("failed to start knotstream consumer: %w", err) 217 } 218 knotstream.Start(ctx) 219 220 state := &State{ 221 db: d, 222 notifier: notifier, 223 indexer: indexer, 224 oauth: oauth, 225 enforcer: enforcer, 226 pages: pages, 227 idResolver: res, 228 rdb: rdb, 229 mentionsResolver: mentionsResolver, 230 posthog: posthog, 231 jc: jc, 232 config: config, 233 repoResolver: repoResolver, 234 aclService: aclService, 235 knotstream: knotstream, 236 logger: logger, 237 cfClient: cfClient, 238 codesearch: &codesearch.CodeSearch{Host: config.CodeSearch.ZoektUrl}, 239 } 240 241 // fetch initial bluesky posts if configured 242 go fetchBskyPosts(ctx, res, config, d, logger) 243 244 return state, nil 245} 246 247func (s *State) Close() error { 248 // other close up logic goes here 249 return s.db.Close() 250} 251 252func (s *State) NewSSHServer() *pipelinessh.Server { 253 return pipelinessh.New(s.db, s.config, log.SubLogger(s.logger, "pipelinessh")) 254} 255 256func (s *State) SecurityTxt(w http.ResponseWriter, r *http.Request) { 257 w.Header().Set("Content-Type", "text/plain") 258 w.Header().Set("Cache-Control", "public, max-age=86400") // one day 259 260 securityTxt := `Contact: mailto:security@tangled.org 261Preferred-Languages: en 262Canonical: https://tangled.org/.well-known/security.txt 263Expires: 2030-01-01T21:59:00.000Z 264` 265 w.Write([]byte(securityTxt)) 266} 267 268func (s *State) RobotsTxt(w http.ResponseWriter, r *http.Request) { 269 w.Header().Set("Content-Type", "text/plain") 270 w.Header().Set("Cache-Control", "public, max-age=86400") // one day 271 272 robotsTxt := `# Hello, Tanglers! 273User-agent: * 274Allow: / 275Disallow: /*/*/settings 276Disallow: /settings 277Disallow: /*/*/compare 278Disallow: /*/*/fork 279Disallow: /search 280Disallow: /*/*/search 281 282Crawl-delay: 1 283` 284 w.Write([]byte(robotsTxt)) 285} 286 287func (s *State) TermsOfService(w http.ResponseWriter, r *http.Request) { 288 s.pages.TermsOfService(w, pages.TermsOfServiceParams{ 289 BaseParams: pages.BaseParamsFromContext(r.Context()), 290 }) 291} 292 293func (s *State) PrivacyPolicy(w http.ResponseWriter, r *http.Request) { 294 s.pages.PrivacyPolicy(w, pages.PrivacyPolicyParams{ 295 BaseParams: pages.BaseParamsFromContext(r.Context()), 296 }) 297} 298 299func (s *State) Brand(w http.ResponseWriter, r *http.Request) { 300 s.pages.Brand(w, pages.BrandParams{ 301 BaseParams: pages.BaseParamsFromContext(r.Context()), 302 }) 303} 304 305func (s *State) UpgradeBanner(w http.ResponseWriter, r *http.Request) { 306 user := s.oauth.GetMultiAccountUser(r) 307 if user == nil { 308 return 309 } 310 311 l := s.logger.With("handler", "UpgradeBanner") 312 l = l.With("did", user.Did) 313 314 regs, err := db.GetRegistrations( 315 s.db, 316 orm.FilterEq("did", user.Did), 317 orm.FilterEq("needs_upgrade", 1), 318 ) 319 if err != nil { 320 l.Error("non-fatal: failed to get registrations", "err", err) 321 } 322 323 spindles, err := db.GetSpindles( 324 r.Context(), 325 s.db, 326 orm.FilterEq("owner", user.Did), 327 orm.FilterEq("needs_upgrade", 1), 328 ) 329 if err != nil { 330 l.Error("non-fatal: failed to get spindles", "err", err) 331 } 332 333 if regs == nil && spindles == nil { 334 return 335 } 336 337 s.pages.UpgradeBanner(w, pages.UpgradeBannerParams{ 338 Registrations: regs, 339 Spindles: spindles, 340 }) 341} 342 343func (s *State) NewsletterSignup(w http.ResponseWriter, r *http.Request) { 344 // target is echoed back from the form via hx-vals so the response span's 345 // id matches the form's hx-target. Fallback keeps the handler useful if 346 // a caller forgets to send it. 347 target := strings.TrimSpace(r.FormValue("target")) 348 if target == "" { 349 target = "home" 350 } 351 352 w.Header().Set("Content-Type", "text/html") 353 354 emailAddr := strings.TrimSpace(r.FormValue("email")) 355 if !email.IsValidEmail(emailAddr) { 356 s.pages.NewsletterResponse(w, pages.NewsletterResponseParams{ 357 Id: target, 358 Error: "Invalid email address.", 359 }) 360 return 361 } 362 363 // For logged-in users, persist the signup locally so the widget stays 364 // hidden across devices. The DB row is the render-time source of truth; 365 // Resend still owns the mailing list itself. 366 if user := s.oauth.GetMultiAccountUser(r); user != nil { 367 if err := db.UpsertNewsletterPref(s.db, user.Did, db.NewsletterStatusSubscribed, emailAddr); err != nil { 368 s.logger.Error("failed to persist newsletter preference", "did", user.Did, "err", err) 369 } 370 } 371 372 if s.config.Resend.ApiKey != "" && s.config.Resend.NewsletterSegmentId != "" { 373 go func() { 374 if err := email.AddNewsletterContact(s.config.Resend.ApiKey, s.config.Resend.NewsletterSegmentId, emailAddr); err != nil { 375 s.logger.Error("failed to add newsletter contact", "error", err) 376 } 377 }() 378 } else { 379 s.logger.Error( 380 "failed to add newsletter contact, missing resend config", 381 "isKeyPresent", s.config.Resend.ApiKey != "", 382 "isSegmentIdPresent", s.config.Resend.NewsletterSegmentId != "", 383 "emailAddr", emailAddr, 384 ) 385 } 386 387 s.pages.NewsletterResponse(w, pages.NewsletterResponseParams{Id: target}) 388} 389 390// NewsletterDismiss records that a logged-in user has dismissed the newsletter 391// widget so it stays hidden across their devices. Anonymous callers get a 204 392// with no DB write — localStorage handles the per-browser fallback. 393func (s *State) NewsletterDismiss(w http.ResponseWriter, r *http.Request) { 394 user := s.oauth.GetMultiAccountUser(r) 395 if user == nil { 396 w.WriteHeader(http.StatusNoContent) 397 return 398 } 399 400 if err := db.UpsertNewsletterPref(s.db, user.Did, db.NewsletterStatusDismissed, ""); err != nil { 401 s.logger.Error("failed to persist newsletter dismissal", "did", user.Did, "err", err) 402 } 403 w.WriteHeader(http.StatusNoContent) 404} 405 406func (s *State) Keys(w http.ResponseWriter, r *http.Request) { 407 user := chi.URLParam(r, "user") 408 user = strings.TrimPrefix(user, "@") 409 user = strings.TrimSuffix(user, ".keys") 410 411 if user == "" { 412 w.WriteHeader(http.StatusBadRequest) 413 return 414 } 415 416 id, err := s.idResolver.ResolveIdent(r.Context(), user) 417 if err != nil { 418 w.WriteHeader(http.StatusInternalServerError) 419 return 420 } 421 422 pubKeys, err := db.GetPublicKeysForDid(s.db, id.DID.String()) 423 if err != nil { 424 s.logger.Error("failed to get public keys", "err", err) 425 http.Error(w, "failed to get public keys", http.StatusInternalServerError) 426 return 427 } 428 429 if len(pubKeys) == 0 { 430 w.WriteHeader(http.StatusNoContent) 431 return 432 } 433 434 for _, k := range pubKeys { 435 key := strings.TrimRight(k.Key, "\n") 436 fmt.Fprintln(w, key) 437 } 438} 439 440func (s *State) NewRepo(w http.ResponseWriter, r *http.Request) { 441 switch r.Method { 442 case http.MethodGet: 443 user := s.oauth.GetMultiAccountUser(r) 444 knots := s.aclService.KnotsForUser(r.Context(), user.Did) 445 446 s.pages.NewRepo(w, pages.NewRepoParams{ 447 BaseParams: pages.BaseParamsFromContext(r.Context()), 448 Knots: knots, 449 }) 450 451 case http.MethodPost: 452 l := s.logger.With("handler", "NewRepo") 453 454 user := s.oauth.GetMultiAccountUser(r) 455 l = l.With("did", user.Did) 456 457 // form validation 458 domain := r.FormValue("domain") 459 if domain == "" { 460 s.pages.Notice(w, "repo", "Invalid form submission&mdash;missing knot domain.") 461 return 462 } 463 l = l.With("knot", domain) 464 465 repoName := r.FormValue("name") 466 if repoName == "" { 467 s.pages.Notice(w, "repo", "Repository name cannot be empty.") 468 return 469 } 470 471 if err := models.ValidateRepoName(repoName); err != nil { 472 s.pages.Notice(w, "repo", err.Error()) 473 return 474 } 475 repoName = models.StripGitExt(repoName) 476 rkey := strings.ToLower(repoName) 477 l = l.With("repoName", repoName, "rkey", rkey) 478 479 defaultBranch := r.FormValue("branch") 480 if defaultBranch == "" { 481 defaultBranch = "main" 482 } 483 l = l.With("defaultBranch", defaultBranch) 484 485 description := r.FormValue("description") 486 if len([]rune(description)) > 140 { 487 s.pages.Notice(w, "repo", "Description must be 140 characters or fewer.") 488 return 489 } 490 491 // ACL validation 492 if !s.aclService.IsRepoCreateAllowed(r.Context(), domain, user.Did) { 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}