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