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