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