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 821 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 289 290Crawl-delay: 1 291` 292 w.Write([]byte(robotsTxt)) 293} 294 295func (s *State) TermsOfService(w http.ResponseWriter, r *http.Request) { 296 s.pages.TermsOfService(w, pages.TermsOfServiceParams{ 297 BaseParams: pages.BaseParamsFromContext(r.Context()), 298 }) 299} 300 301func (s *State) PrivacyPolicy(w http.ResponseWriter, r *http.Request) { 302 s.pages.PrivacyPolicy(w, pages.PrivacyPolicyParams{ 303 BaseParams: pages.BaseParamsFromContext(r.Context()), 304 }) 305} 306 307func (s *State) Brand(w http.ResponseWriter, r *http.Request) { 308 s.pages.Brand(w, pages.BrandParams{ 309 BaseParams: pages.BaseParamsFromContext(r.Context()), 310 }) 311} 312 313func (s *State) UpgradeBanner(w http.ResponseWriter, r *http.Request) { 314 user := s.oauth.GetMultiAccountUser(r) 315 if user == nil { 316 return 317 } 318 319 l := s.logger.With("handler", "UpgradeBanner") 320 l = l.With("did", user.Did) 321 322 regs, err := db.GetRegistrations( 323 s.db, 324 orm.FilterEq("did", user.Did), 325 orm.FilterEq("needs_upgrade", 1), 326 ) 327 if err != nil { 328 l.Error("non-fatal: failed to get registrations", "err", err) 329 } 330 331 spindles, err := db.GetSpindles( 332 r.Context(), 333 s.db, 334 orm.FilterEq("owner", user.Did), 335 orm.FilterEq("needs_upgrade", 1), 336 ) 337 if err != nil { 338 l.Error("non-fatal: failed to get spindles", "err", err) 339 } 340 341 if regs == nil && spindles == nil { 342 return 343 } 344 345 s.pages.UpgradeBanner(w, pages.UpgradeBannerParams{ 346 Registrations: regs, 347 Spindles: spindles, 348 }) 349} 350 351func (s *State) NewsletterSignup(w http.ResponseWriter, r *http.Request) { 352 // target is echoed back from the form via hx-vals so the response span's 353 // id matches the form's hx-target. Fallback keeps the handler useful if 354 // a caller forgets to send it. 355 target := strings.TrimSpace(r.FormValue("target")) 356 if target == "" { 357 target = "home" 358 } 359 360 w.Header().Set("Content-Type", "text/html") 361 362 emailAddr := strings.TrimSpace(r.FormValue("email")) 363 if !email.IsValidEmail(emailAddr) { 364 s.pages.NewsletterResponse(w, pages.NewsletterResponseParams{ 365 Id: target, 366 Error: "Invalid email address.", 367 }) 368 return 369 } 370 371 // For logged-in users, persist the signup locally so the widget stays 372 // hidden across devices. The DB row is the render-time source of truth; 373 // Resend still owns the mailing list itself. 374 if user := s.oauth.GetMultiAccountUser(r); user != nil { 375 if err := db.UpsertNewsletterPref(s.db, user.Did, db.NewsletterStatusSubscribed, emailAddr); err != nil { 376 s.logger.Error("failed to persist newsletter preference", "did", user.Did, "err", err) 377 } 378 } 379 380 if s.config.Resend.ApiKey != "" && s.config.Resend.NewsletterSegmentId != "" { 381 go func() { 382 if err := email.AddNewsletterContact(s.config.Resend.ApiKey, s.config.Resend.NewsletterSegmentId, emailAddr); err != nil { 383 s.logger.Error("failed to add newsletter contact", "error", err) 384 } 385 }() 386 } else { 387 s.logger.Error( 388 "failed to add newsletter contact, missing resend config", 389 "isKeyPresent", s.config.Resend.ApiKey != "", 390 "isSegmentIdPresent", s.config.Resend.NewsletterSegmentId != "", 391 "emailAddr", emailAddr, 392 ) 393 } 394 395 s.pages.NewsletterResponse(w, pages.NewsletterResponseParams{Id: target}) 396} 397 398// NewsletterDismiss records that a logged-in user has dismissed the newsletter 399// widget so it stays hidden across their devices. Anonymous callers get a 204 400// with no DB write — localStorage handles the per-browser fallback. 401func (s *State) NewsletterDismiss(w http.ResponseWriter, r *http.Request) { 402 user := s.oauth.GetMultiAccountUser(r) 403 if user == nil { 404 w.WriteHeader(http.StatusNoContent) 405 return 406 } 407 408 if err := db.UpsertNewsletterPref(s.db, user.Did, db.NewsletterStatusDismissed, ""); err != nil { 409 s.logger.Error("failed to persist newsletter dismissal", "did", user.Did, "err", err) 410 } 411 w.WriteHeader(http.StatusNoContent) 412} 413 414func (s *State) Keys(w http.ResponseWriter, r *http.Request) { 415 user := chi.URLParam(r, "user") 416 user = strings.TrimPrefix(user, "@") 417 user = strings.TrimSuffix(user, ".keys") 418 419 if user == "" { 420 w.WriteHeader(http.StatusBadRequest) 421 return 422 } 423 424 id, err := s.idResolver.ResolveIdent(r.Context(), user) 425 if err != nil { 426 w.WriteHeader(http.StatusInternalServerError) 427 return 428 } 429 430 pubKeys, err := db.GetPublicKeysForDid(s.db, id.DID.String()) 431 if err != nil { 432 s.logger.Error("failed to get public keys", "err", err) 433 http.Error(w, "failed to get public keys", http.StatusInternalServerError) 434 return 435 } 436 437 if len(pubKeys) == 0 { 438 w.WriteHeader(http.StatusNoContent) 439 return 440 } 441 442 for _, k := range pubKeys { 443 key := strings.TrimRight(k.Key, "\n") 444 fmt.Fprintln(w, key) 445 } 446} 447 448func (s *State) NewRepo(w http.ResponseWriter, r *http.Request) { 449 switch r.Method { 450 case http.MethodGet: 451 user := s.oauth.GetMultiAccountUser(r) 452 knots := s.aclService.KnotsForUser(r.Context(), user.Did) 453 454 s.pages.NewRepo(w, pages.NewRepoParams{ 455 BaseParams: pages.BaseParamsFromContext(r.Context()), 456 Knots: knots, 457 }) 458 459 case http.MethodPost: 460 l := s.logger.With("handler", "NewRepo") 461 462 user := s.oauth.GetMultiAccountUser(r) 463 l = l.With("did", user.Did) 464 465 // form validation 466 domain := r.FormValue("domain") 467 if domain == "" { 468 s.pages.Notice(w, "repo", "Invalid form submission&mdash;missing knot domain.") 469 return 470 } 471 l = l.With("knot", domain) 472 473 repoName := r.FormValue("name") 474 if repoName == "" { 475 s.pages.Notice(w, "repo", "Repository name cannot be empty.") 476 return 477 } 478 479 if err := models.ValidateRepoName(repoName); err != nil { 480 s.pages.Notice(w, "repo", err.Error()) 481 return 482 } 483 repoName = models.StripGitExt(repoName) 484 rkey := strings.ToLower(repoName) 485 l = l.With("repoName", repoName, "rkey", rkey) 486 487 defaultBranch := r.FormValue("branch") 488 if defaultBranch == "" { 489 defaultBranch = "main" 490 } 491 l = l.With("defaultBranch", defaultBranch) 492 493 description := r.FormValue("description") 494 if len([]rune(description)) > 140 { 495 s.pages.Notice(w, "repo", "Description must be 140 characters or fewer.") 496 return 497 } 498 499 // ACL validation 500 if !s.aclService.IsRepoCreateAllowed(r.Context(), domain, user.Did) { 501 l.Info("unauthorized") 502 s.pages.Notice(w, "repo", "You do not have permission to create a repo in this knot.") 503 return 504 } 505 506 // Check for existing repos 507 existingRepo, err := db.GetRepo( 508 s.db, 509 orm.FilterEq("did", user.Did), 510 orm.FilterEq("rkey", rkey), 511 ) 512 if err == nil && existingRepo != nil { 513 l.Info("repo exists") 514 s.pages.Notice(w, "repo", fmt.Sprintf("You already have a repository by this name on %s", existingRepo.Knot)) 515 return 516 } 517 518 atpClient, err := s.oauth.AuthorizedClient(r) 519 if err != nil { 520 l.Error("failed to get authorized client", "err", err) 521 s.pages.Notice(w, "repo", "Failed to authorize. Try again later.") 522 return 523 } 524 525 if rkeyOccupied(r.Context(), atpClient, user.Did, rkey) { 526 l.Info("rkey occupied by prior rename alias") 527 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)) 528 return 529 } 530 531 client, err := s.oauth.ServiceClient( 532 r, 533 oauth.WithService(domain), 534 oauth.WithLxm(tangled.RepoCreateNSID), 535 oauth.WithDev(s.config.Core.Dev), 536 ) 537 if err != nil { 538 l.Error("service auth failed", "err", err) 539 s.pages.Notice(w, "repo", "Failed to authenticate. Please log out and log back in again.") 540 return 541 } 542 543 input := &tangled.RepoCreate_Input{ 544 Rkey: rkey, 545 Name: rkey, 546 DefaultBranch: &defaultBranch, 547 } 548 createResp, err := tangled.RepoCreate( 549 r.Context(), 550 client, 551 input, 552 ) 553 if xrpcerr := xrpcclient.HandleXrpcErr(err); xrpcerr != nil { 554 l.Error("failed to call XRPC repo.create", "xrpcerr", xrpcerr, "err", err) 555 s.pages.Notice(w, "repo", err.Error()) 556 return 557 } 558 559 var repoDid string 560 if createResp != nil && createResp.RepoDid != nil { 561 repoDid = *createResp.RepoDid 562 } 563 if repoDid == "" { 564 l.Error("knot returned empty repo DID") 565 s.pages.Notice(w, "repo", "Knot failed to mint a repo DID. The knot may need to be upgraded.") 566 return 567 } 568 569 repo := &models.Repo{ 570 Did: user.Did, 571 Name: repoName, 572 Knot: domain, 573 Rkey: rkey, 574 Description: description, 575 Created: time.Now(), 576 Labels: s.config.Label.DefaultLabelDefs, 577 RepoDid: repoDid, 578 } 579 record := repo.AsRecord() 580 581 cleanupKnot := func() { 582 go func() { 583 delays := []time.Duration{0, 2 * time.Second, 5 * time.Second} 584 for attempt, delay := range delays { 585 time.Sleep(delay) 586 deleteClient, dErr := s.oauth.ServiceClient( 587 r, 588 oauth.WithService(domain), 589 oauth.WithLxm(tangled.RepoDeleteNSID), 590 oauth.WithDev(s.config.Core.Dev), 591 ) 592 if dErr != nil { 593 l.Error("failed to create delete client for knot cleanup", "attempt", attempt+1, "err", dErr) 594 continue 595 } 596 ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) 597 if dErr := tangled.RepoDelete(ctx, deleteClient, &tangled.RepoDelete_Input{ 598 Did: user.Did, 599 Name: rkey, 600 Rkey: rkey, 601 }); dErr != nil { 602 cancel() 603 l.Error("failed to clean up repo on knot after rollback", "attempt", attempt+1, "err", dErr) 604 continue 605 } 606 cancel() 607 l.Info("successfully cleaned up repo on knot after rollback", "attempt", attempt+1) 608 return 609 } 610 l.Error("exhausted retries for knot cleanup, repo may be orphaned", 611 "did", user.Did, "repo", repoName, "knot", domain) 612 }() 613 } 614 615 _, err = comatproto.RepoCreateRecord(r.Context(), atpClient, &comatproto.RepoCreateRecord_Input{ 616 Collection: tangled.RepoNSID, 617 Repo: user.Did, 618 Rkey: &rkey, 619 Record: &lexutil.LexiconTypeDecoder{ 620 Val: &record, 621 }, 622 }) 623 if err != nil { 624 l.Info("PDS write failed", "err", err) 625 cleanupKnot() 626 if rkeyOccupied(r.Context(), atpClient, user.Did, rkey) { 627 s.pages.Notice(w, "repo", fmt.Sprintf("You already have a repository named %q.", rkey)) 628 } else { 629 s.pages.Notice(w, "repo", "Failed to announce repository creation.") 630 } 631 return 632 } 633 634 aturi := fmt.Sprintf("at://%s/%s/%s", user.Did, tangled.RepoNSID, rkey) 635 l = l.With("aturi", aturi) 636 l.Info("wrote to PDS") 637 638 tx, err := s.db.BeginTx(r.Context(), nil) 639 if err != nil { 640 l.Info("txn failed", "err", err) 641 s.pages.Notice(w, "repo", "Failed to save repository information.") 642 return 643 } 644 645 rollback := func() { 646 err1 := tx.Rollback() 647 err2 := s.enforcer.E.LoadPolicy() 648 err3 := rollbackRecord(context.Background(), aturi, atpClient) 649 650 if errors.Is(err1, sql.ErrTxDone) { 651 err1 = nil 652 } 653 654 if errs := errors.Join(err1, err2, err3); errs != nil { 655 l.Error("failed to rollback changes", "errs", errs) 656 } 657 658 if aturi != "" { 659 cleanupKnot() 660 } 661 } 662 defer rollback() 663 664 err = db.AddRepo(tx, repo) 665 if err != nil { 666 l.Error("db write failed", "err", err) 667 s.pages.Notice(w, "repo", "Failed to save repository information.") 668 return 669 } 670 671 rbacPath := repo.RepoIdentifier() 672 err = s.enforcer.AddRepo(user.Did, domain, rbacPath) 673 if err != nil { 674 l.Error("acl setup failed", "err", err) 675 s.pages.Notice(w, "repo", "Failed to set up repository permissions.") 676 return 677 } 678 679 err = tx.Commit() 680 if err != nil { 681 l.Error("txn commit failed", "err", err) 682 http.Error(w, err.Error(), http.StatusInternalServerError) 683 return 684 } 685 686 err = s.enforcer.E.SavePolicy() 687 if err != nil { 688 l.Error("acl save failed", "err", err) 689 http.Error(w, err.Error(), http.StatusInternalServerError) 690 return 691 } 692 693 aturi = "" 694 695 s.notifier.NewRepo(r.Context(), repo) 696 switch { 697 case repoDid != "": 698 s.pages.HxLocation(w, fmt.Sprintf("/%s", repoDid)) 699 default: 700 handle := s.pages.DisplayHandle(r.Context(), user.Did) 701 s.pages.HxLocation(w, fmt.Sprintf("/%s/%s", handle, rkey)) 702 } 703 } 704} 705 706func rkeyOccupied(ctx context.Context, client *atclient.APIClient, did, rkey string) bool { 707 probeCtx, cancel := context.WithTimeout(ctx, 10*time.Second) 708 defer cancel() 709 resp, err := comatproto.RepoGetRecord(probeCtx, client, "", tangled.RepoNSID, did, rkey) 710 return err == nil && resp != nil 711} 712 713// this is used to rollback changes made to the PDS 714// 715// it is a no-op if the provided ATURI is empty 716func rollbackRecord(ctx context.Context, aturi string, client *atclient.APIClient) error { 717 if aturi == "" { 718 return nil 719 } 720 721 parsed := syntax.ATURI(aturi) 722 723 collection := parsed.Collection().String() 724 repo := parsed.Authority().String() 725 rkey := parsed.RecordKey().String() 726 727 _, err := comatproto.RepoDeleteRecord(ctx, client, &comatproto.RepoDeleteRecord_Input{ 728 Collection: collection, 729 Repo: repo, 730 Rkey: rkey, 731 }) 732 return err 733} 734 735func BackfillDefaultDefs(e db.Execer, r *idresolver.Resolver, defaults []string) error { 736 defaultLabels, err := db.GetLabelDefinitions(e, orm.FilterIn("at_uri", defaults)) 737 if err != nil { 738 return err 739 } 740 // already present 741 if len(defaultLabels) == len(defaults) { 742 return nil 743 } 744 745 labelDefs, err := models.FetchLabelDefs(r, defaults) 746 if err != nil { 747 return err 748 } 749 750 // Insert each label definition to the database 751 for _, labelDef := range labelDefs { 752 _, err = db.AddLabelDefinition(e, &labelDef) 753 if err != nil { 754 return fmt.Errorf("failed to add label definition %s: %v", labelDef.Name, err) 755 } 756 } 757 758 return nil 759} 760 761func fetchBskyPosts(ctx context.Context, res *idresolver.Resolver, config *config.Config, d *db.DB, logger *slog.Logger) { 762 resolved, err := res.ResolveIdent(context.Background(), consts.TangledDid) 763 if err != nil { 764 logger.Error("failed to resolve tangled.org DID", "err", err) 765 return 766 } 767 768 pdsEndpoint := resolved.PDSEndpoint() 769 if pdsEndpoint == "" { 770 logger.Error("no PDS endpoint found for tangled.sh DID") 771 return 772 } 773 774 session, err := oauth.CreateAppPasswordSession(res, config.Core.AppPassword, consts.TangledDid, logger) 775 if err != nil { 776 logger.Error("failed to create appassword session... skipping fetch", "err", err) 777 return 778 } 779 780 l := log.SubLogger(logger, "bluesky") 781 782 ticker := time.NewTicker(config.Bluesky.UpdateInterval) 783 defer ticker.Stop() 784 785 for { 786 // refresh session if necessary 787 if !session.IsValid() { 788 l.Debug("access token expired, refreshing session") 789 if err := session.RefreshSession(); err != nil { 790 l.Error("failed to refresh session, stopping bluesky updater", "err", err) 791 return 792 } 793 l.Debug("session refreshed") 794 } 795 796 // make client 797 client := xrpc.Client{ 798 Auth: &xrpc.AuthInfo{ 799 AccessJwt: session.AccessJwt, 800 Did: session.Did, 801 }, 802 Host: session.PdsEndpoint, 803 } 804 805 posts, _, err := bsky.FetchPosts(ctx, &client, 20, "") 806 if err != nil { 807 l.Error("failed to fetch bluesky posts", "err", err) 808 } else if err := db.InsertBlueskyPosts(d, posts); err != nil { 809 l.Error("failed to insert bluesky posts", "err", err) 810 } else { 811 l.Info("inserted bluesky posts", "count", len(posts)) 812 } 813 814 select { 815 case <-ticker.C: 816 case <-ctx.Done(): 817 l.Info("stopping bluesky updater") 818 return 819 } 820 } 821}