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