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