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