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