This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / appview / ingester.go
32 kB 1290 lines
1package appview 2 3import ( 4 "context" 5 "database/sql" 6 "encoding/json" 7 "errors" 8 "fmt" 9 "io" 10 "log/slog" 11 "maps" 12 "net/http" 13 "net/url" 14 "slices" 15 "sync" 16 17 "time" 18 19 "github.com/avast/retry-go/v4" 20 "github.com/bluesky-social/indigo/atproto/syntax" 21 jmodels "github.com/bluesky-social/jetstream/pkg/models" 22 "github.com/go-git/go-git/v5/plumbing" 23 "github.com/ipfs/go-cid" 24 "golang.org/x/sync/errgroup" 25 "tangled.org/core/api/tangled" 26 "tangled.org/core/appview/config" 27 "tangled.org/core/appview/db" 28 "tangled.org/core/appview/models" 29 "tangled.org/core/appview/serververify" 30 "tangled.org/core/appview/validator" 31 "tangled.org/core/idresolver" 32 "tangled.org/core/orm" 33 "tangled.org/core/rbac" 34) 35 36type Ingester struct { 37 Db db.DbWrapper 38 Enforcer *rbac.Enforcer 39 IdResolver *idresolver.Resolver 40 Config *config.Config 41 Logger *slog.Logger 42 Validator *validator.Validator 43} 44 45type processFunc func(ctx context.Context, e *jmodels.Event) error 46 47func (i *Ingester) Ingest() processFunc { 48 return func(ctx context.Context, e *jmodels.Event) error { 49 var err error 50 51 l := i.Logger.With("kind", e.Kind) 52 switch e.Kind { 53 case jmodels.EventKindAccount: 54 if !e.Account.Active && *e.Account.Status == "deactivated" { 55 err = i.IdResolver.InvalidateIdent(ctx, e.Account.Did) 56 } 57 case jmodels.EventKindIdentity: 58 err = i.IdResolver.InvalidateIdent(ctx, e.Identity.Did) 59 case jmodels.EventKindCommit: 60 switch e.Commit.Collection { 61 case tangled.GraphFollowNSID: 62 err = i.ingestFollow(e) 63 case tangled.FeedStarNSID: 64 err = i.ingestStar(e) 65 case tangled.PublicKeyNSID: 66 err = i.ingestPublicKey(e) 67 case tangled.RepoArtifactNSID: 68 err = i.ingestArtifact(e) 69 case tangled.ActorProfileNSID: 70 err = i.ingestProfile(ctx, e) 71 case tangled.SpindleMemberNSID: 72 err = i.ingestSpindleMember(ctx, e) 73 case tangled.SpindleNSID: 74 err = i.ingestSpindle(ctx, e) 75 case tangled.KnotMemberNSID: 76 err = i.ingestKnotMember(e) 77 case tangled.KnotNSID: 78 err = i.ingestKnot(e) 79 case tangled.StringNSID: 80 err = i.ingestString(e) 81 case tangled.RepoIssueNSID: 82 err = i.ingestIssue(ctx, e) 83 case tangled.RepoPullNSID: 84 err = i.ingestPull(ctx, e) 85 case tangled.RepoIssueCommentNSID: 86 err = i.ingestIssueComment(e) 87 case tangled.LabelDefinitionNSID: 88 err = i.ingestLabelDefinition(e) 89 case tangled.LabelOpNSID: 90 err = i.ingestLabelOp(e) 91 } 92 l = i.Logger.With("nsid", e.Commit.Collection) 93 } 94 95 if err != nil { 96 l.Warn("failed to ingest record, skipping", "err", err) 97 } 98 99 lastTimeUs := e.TimeUS + 1 100 if saveErr := i.Db.SaveLastTimeUs(lastTimeUs); saveErr != nil { 101 l.Error("failed to save cursor", "err", saveErr) 102 } 103 104 return nil 105 } 106} 107 108func (i *Ingester) ingestStar(e *jmodels.Event) error { 109 var err error 110 did := e.Did 111 112 l := i.Logger.With("handler", "ingestStar") 113 l = l.With("nsid", e.Commit.Collection) 114 115 switch e.Commit.Operation { 116 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 117 var subjectUri syntax.ATURI 118 119 raw := json.RawMessage(e.Commit.Record) 120 record := tangled.FeedStar{} 121 err := json.Unmarshal(raw, &record) 122 if err != nil { 123 l.Error("invalid record", "err", err) 124 return err 125 } 126 127 star := &models.Star{ 128 Did: did, 129 Rkey: e.Commit.RKey, 130 } 131 132 switch { 133 case record.SubjectDid != nil: 134 repo, repoErr := db.GetRepo(i.Db, orm.FilterEq("repo_did", *record.SubjectDid)) 135 if repoErr == nil { 136 subjectUri = repo.RepoAt() 137 star.RepoAt = subjectUri 138 } 139 case record.Subject != nil: 140 subjectUri, err = syntax.ParseATURI(*record.Subject) 141 if err != nil { 142 l.Error("invalid record", "err", err) 143 return err 144 } 145 star.RepoAt = subjectUri 146 repo, repoErr := db.GetRepoByAtUri(i.Db, subjectUri.String()) 147 if repoErr == nil && repo.RepoDid != "" { 148 if enqErr := db.EnqueuePdsRewrite(i.Db, did, repo.RepoDid, tangled.FeedStarNSID, e.Commit.RKey, *record.Subject); enqErr != nil { 149 l.Warn("failed to enqueue PDS rewrite for star", "err", enqErr, "did", did, "repoDid", repo.RepoDid) 150 } 151 } 152 default: 153 l.Error("star record has neither subject nor subjectDid") 154 return fmt.Errorf("star record has neither subject nor subjectDid") 155 } 156 err = db.AddStar(i.Db, star) 157 case jmodels.CommitOperationDelete: 158 err = db.DeleteStarByRkey(i.Db, did, e.Commit.RKey) 159 } 160 161 if err != nil { 162 return fmt.Errorf("failed to %s star record: %w", e.Commit.Operation, err) 163 } 164 165 return nil 166} 167 168func (i *Ingester) ingestFollow(e *jmodels.Event) error { 169 var err error 170 did := e.Did 171 172 l := i.Logger.With("handler", "ingestFollow") 173 l = l.With("nsid", e.Commit.Collection) 174 175 switch e.Commit.Operation { 176 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 177 raw := json.RawMessage(e.Commit.Record) 178 record := tangled.GraphFollow{} 179 err = json.Unmarshal(raw, &record) 180 if err != nil { 181 l.Error("invalid record", "err", err) 182 return err 183 } 184 185 err = db.AddFollow(i.Db, &models.Follow{ 186 UserDid: did, 187 SubjectDid: record.Subject, 188 Rkey: e.Commit.RKey, 189 }) 190 case jmodels.CommitOperationDelete: 191 err = db.DeleteFollowByRkey(i.Db, did, e.Commit.RKey) 192 } 193 194 if err != nil { 195 return fmt.Errorf("failed to %s follow record: %w", e.Commit.Operation, err) 196 } 197 198 return nil 199} 200 201func (i *Ingester) ingestPublicKey(e *jmodels.Event) error { 202 did := e.Did 203 var err error 204 205 l := i.Logger.With("handler", "ingestPublicKey") 206 l = l.With("nsid", e.Commit.Collection) 207 208 switch e.Commit.Operation { 209 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 210 l.Debug("processing add of pubkey") 211 raw := json.RawMessage(e.Commit.Record) 212 record := tangled.PublicKey{} 213 err = json.Unmarshal(raw, &record) 214 if err != nil { 215 l.Error("invalid record", "err", err) 216 return err 217 } 218 219 name := record.Name 220 key := record.Key 221 err = db.AddPublicKey(i.Db, did, name, key, e.Commit.RKey) 222 case jmodels.CommitOperationDelete: 223 l.Debug("processing delete of pubkey") 224 err = db.DeletePublicKeyByRkey(i.Db, did, e.Commit.RKey) 225 } 226 227 if err != nil { 228 return fmt.Errorf("failed to %s pubkey record: %w", e.Commit.Operation, err) 229 } 230 231 return nil 232} 233 234func (i *Ingester) ingestArtifact(e *jmodels.Event) error { 235 did := e.Did 236 var err error 237 238 l := i.Logger.With("handler", "ingestArtifact") 239 l = l.With("nsid", e.Commit.Collection) 240 241 switch e.Commit.Operation { 242 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 243 raw := json.RawMessage(e.Commit.Record) 244 record := tangled.RepoArtifact{} 245 err = json.Unmarshal(raw, &record) 246 if err != nil { 247 l.Error("invalid record", "err", err) 248 return err 249 } 250 251 var repo *models.Repo 252 if record.RepoDid != nil && *record.RepoDid != "" { 253 repo, err = db.GetRepoByDid(i.Db, *record.RepoDid) 254 if err != nil && !errors.Is(err, sql.ErrNoRows) { 255 return fmt.Errorf("failed to look up repo by DID %s: %w", *record.RepoDid, err) 256 } 257 } 258 if repo == nil && record.Repo != nil { 259 repoAt, parseErr := syntax.ParseATURI(*record.Repo) 260 if parseErr != nil { 261 return parseErr 262 } 263 repo, err = db.GetRepoByAtUri(i.Db, repoAt.String()) 264 if err != nil { 265 return err 266 } 267 } 268 if repo == nil { 269 return fmt.Errorf("artifact record has neither valid repoDid nor repo field") 270 } 271 272 ok, err := i.Enforcer.E.Enforce(did, repo.Knot, repo.RepoIdentifier(), "repo:push") 273 if err != nil || !ok { 274 return err 275 } 276 277 repoDid := repo.RepoDid 278 if repoDid == "" && record.RepoDid != nil { 279 repoDid = *record.RepoDid 280 } 281 if repoDid != "" && (record.RepoDid == nil || *record.RepoDid == "") && record.Repo != nil { 282 if enqErr := db.EnqueuePdsRewrite(i.Db, did, repoDid, tangled.RepoArtifactNSID, e.Commit.RKey, *record.Repo); enqErr != nil { 283 l.Warn("failed to enqueue PDS rewrite for artifact", "err", enqErr, "did", did, "repoDid", repoDid) 284 } 285 } 286 287 createdAt, err := time.Parse(time.RFC3339, record.CreatedAt) 288 if err != nil { 289 createdAt = time.Now() 290 } 291 292 artifact := models.Artifact{ 293 Did: did, 294 Rkey: e.Commit.RKey, 295 RepoAt: repo.RepoAt(), 296 Tag: plumbing.Hash(record.Tag), 297 CreatedAt: createdAt, 298 BlobCid: cid.Cid(record.Artifact.Ref), 299 Name: record.Name, 300 Size: uint64(record.Artifact.Size), 301 MimeType: record.Artifact.MimeType, 302 } 303 304 err = db.AddArtifact(i.Db, artifact) 305 case jmodels.CommitOperationDelete: 306 err = db.DeleteArtifact(i.Db, orm.FilterEq("did", did), orm.FilterEq("rkey", e.Commit.RKey)) 307 } 308 309 if err != nil { 310 return fmt.Errorf("failed to %s artifact record: %w", e.Commit.Operation, err) 311 } 312 313 return nil 314} 315 316func (i *Ingester) ingestProfile(ctx context.Context, e *jmodels.Event) error { 317 did := e.Did 318 var err error 319 320 l := i.Logger.With("handler", "ingestProfile") 321 l = l.With("nsid", e.Commit.Collection) 322 323 if e.Commit.RKey != "self" { 324 return fmt.Errorf("ingestProfile only ingests `self` record") 325 } 326 327 switch e.Commit.Operation { 328 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 329 raw := json.RawMessage(e.Commit.Record) 330 record := tangled.ActorProfile{} 331 err = json.Unmarshal(raw, &record) 332 if err != nil { 333 l.Error("invalid record", "err", err) 334 return err 335 } 336 337 avatar := "" 338 if record.Avatar != nil { 339 avatar = record.Avatar.Ref.String() 340 } 341 342 description := "" 343 if record.Description != nil { 344 description = *record.Description 345 } 346 347 includeBluesky := record.Bluesky 348 349 pronouns := "" 350 if record.Pronouns != nil { 351 pronouns = *record.Pronouns 352 } 353 354 location := "" 355 if record.Location != nil { 356 location = *record.Location 357 } 358 359 var links [5]string 360 for i, l := range record.Links { 361 if i < 5 { 362 links[i] = l 363 } 364 } 365 366 var stats [2]models.VanityStat 367 for i, s := range record.Stats { 368 if i < 2 { 369 stats[i].Kind = models.ParseVanityStatKind(s) 370 } 371 } 372 373 var pinned [6]string 374 for i, r := range record.PinnedRepositories { 375 if i < 6 { 376 pinned[i] = r 377 } 378 } 379 380 var preferredHandle syntax.Handle 381 if record.PreferredHandle != nil { 382 if h, err := syntax.ParseHandle(*record.PreferredHandle); err == nil { 383 ident, identErr := i.IdResolver.ResolveIdent(ctx, did) 384 if identErr == nil && slices.Contains(ident.AlsoKnownAs, "at://"+string(h)) { 385 preferredHandle = h 386 } 387 } 388 } 389 390 profile := models.Profile{ 391 Did: did, 392 Avatar: avatar, 393 Description: description, 394 IncludeBluesky: includeBluesky, 395 Location: location, 396 Links: links, 397 Stats: stats, 398 PinnedRepos: pinned, 399 Pronouns: pronouns, 400 PreferredHandle: preferredHandle, 401 } 402 403 ddb, ok := i.Db.Execer.(*db.DB) 404 if !ok { 405 return fmt.Errorf("failed to index profile record, invalid db cast") 406 } 407 408 tx, err := ddb.Begin() 409 if err != nil { 410 return fmt.Errorf("failed to start transaction") 411 } 412 413 err = db.ValidateProfile(tx, &profile) 414 if err != nil { 415 return fmt.Errorf("invalid profile record") 416 } 417 418 err = db.UpsertProfile(tx, &profile) 419 case jmodels.CommitOperationDelete: 420 err = db.DeleteArtifact(i.Db, orm.FilterEq("did", did), orm.FilterEq("rkey", e.Commit.RKey)) 421 } 422 423 if err != nil { 424 return fmt.Errorf("failed to %s profile record: %w", e.Commit.Operation, err) 425 } 426 427 return nil 428} 429 430func (i *Ingester) ingestSpindleMember(ctx context.Context, e *jmodels.Event) error { 431 did := e.Did 432 var err error 433 434 l := i.Logger.With("handler", "ingestSpindleMember") 435 l = l.With("nsid", e.Commit.Collection) 436 437 switch e.Commit.Operation { 438 case jmodels.CommitOperationCreate: 439 raw := json.RawMessage(e.Commit.Record) 440 record := tangled.SpindleMember{} 441 err = json.Unmarshal(raw, &record) 442 if err != nil { 443 l.Error("invalid record", "err", err) 444 return err 445 } 446 447 // only spindle owner can invite to spindles 448 ok, err := i.Enforcer.IsSpindleInviteAllowed(did, record.Instance) 449 if err != nil || !ok { 450 return fmt.Errorf("failed to enforce permissions: %w", err) 451 } 452 453 memberId, err := i.IdResolver.ResolveIdent(ctx, record.Subject) 454 if err != nil { 455 return err 456 } 457 458 if memberId.Handle.IsInvalidHandle() { 459 return err 460 } 461 462 ddb, ok := i.Db.Execer.(*db.DB) 463 if !ok { 464 return fmt.Errorf("invalid db cast") 465 } 466 467 err = db.AddSpindleMember(ddb, models.SpindleMember{ 468 Did: syntax.DID(did), 469 Rkey: e.Commit.RKey, 470 Instance: record.Instance, 471 Subject: memberId.DID, 472 }) 473 if !ok { 474 return fmt.Errorf("failed to add to db: %w", err) 475 } 476 477 err = i.Enforcer.AddSpindleMember(record.Instance, memberId.DID.String()) 478 if err != nil { 479 return fmt.Errorf("failed to update ACLs: %w", err) 480 } 481 482 l.Info("added spindle member") 483 case jmodels.CommitOperationDelete: 484 rkey := e.Commit.RKey 485 486 ddb, ok := i.Db.Execer.(*db.DB) 487 if !ok { 488 return fmt.Errorf("failed to index profile record, invalid db cast") 489 } 490 491 // get record from db first 492 members, err := db.GetSpindleMembers( 493 ddb, 494 orm.FilterEq("did", did), 495 orm.FilterEq("rkey", rkey), 496 ) 497 if err != nil || len(members) != 1 { 498 return fmt.Errorf("failed to get member: %w, len(members) = %d", err, len(members)) 499 } 500 member := members[0] 501 502 tx, err := ddb.Begin() 503 if err != nil { 504 return fmt.Errorf("failed to start txn: %w", err) 505 } 506 507 // remove record by rkey && update enforcer 508 if err = db.RemoveSpindleMember( 509 tx, 510 orm.FilterEq("did", did), 511 orm.FilterEq("rkey", rkey), 512 ); err != nil { 513 return fmt.Errorf("failed to remove from db: %w", err) 514 } 515 516 // update enforcer 517 err = i.Enforcer.RemoveSpindleMember(member.Instance, member.Subject.String()) 518 if err != nil { 519 return fmt.Errorf("failed to update ACLs: %w", err) 520 } 521 522 if err = tx.Commit(); err != nil { 523 return fmt.Errorf("failed to commit txn: %w", err) 524 } 525 526 if err = i.Enforcer.E.SavePolicy(); err != nil { 527 return fmt.Errorf("failed to save ACLs: %w", err) 528 } 529 530 l.Info("removed spindle member") 531 } 532 533 return nil 534} 535 536func (i *Ingester) ingestSpindle(ctx context.Context, e *jmodels.Event) error { 537 did := e.Did 538 var err error 539 540 l := i.Logger.With("handler", "ingestSpindle") 541 l = l.With("nsid", e.Commit.Collection) 542 543 switch e.Commit.Operation { 544 case jmodels.CommitOperationCreate: 545 raw := json.RawMessage(e.Commit.Record) 546 record := tangled.Spindle{} 547 err = json.Unmarshal(raw, &record) 548 if err != nil { 549 l.Error("invalid record", "err", err) 550 return err 551 } 552 553 instance := e.Commit.RKey 554 555 ddb, ok := i.Db.Execer.(*db.DB) 556 if !ok { 557 return fmt.Errorf("failed to index profile record, invalid db cast") 558 } 559 560 err := db.AddSpindle(ddb, models.Spindle{ 561 Owner: syntax.DID(did), 562 Instance: instance, 563 }) 564 if err != nil { 565 l.Error("failed to add spindle to db", "err", err, "instance", instance) 566 return err 567 } 568 569 err = retry.Do( 570 func() error { return serververify.RunVerification(ctx, instance, did, i.Config.Core.Dev) }, 571 retry.Attempts(5), retry.Delay(5*time.Second), retry.MaxDelay(80*time.Second), 572 retry.DelayType(retry.BackOffDelay), retry.LastErrorOnly(true), 573 ) 574 if err != nil { 575 l.Error("failed to verify spindle after retries", "err", err, "instance", instance) 576 return err 577 } 578 579 _, err = serververify.MarkSpindleVerified(ddb, i.Enforcer, instance, did) 580 if err != nil { 581 return fmt.Errorf("failed to mark verified: %w", err) 582 } 583 584 return nil 585 586 case jmodels.CommitOperationDelete: 587 instance := e.Commit.RKey 588 589 ddb, ok := i.Db.Execer.(*db.DB) 590 if !ok { 591 return fmt.Errorf("failed to index profile record, invalid db cast") 592 } 593 594 // get record from db first 595 spindles, err := db.GetSpindles( 596 ctx, 597 ddb, 598 orm.FilterEq("owner", did), 599 orm.FilterEq("instance", instance), 600 ) 601 if err != nil || len(spindles) != 1 { 602 return fmt.Errorf("failed to get spindles: %w, len(spindles) = %d", err, len(spindles)) 603 } 604 spindle := spindles[0] 605 606 tx, err := ddb.Begin() 607 if err != nil { 608 return err 609 } 610 defer func() { 611 tx.Rollback() 612 i.Enforcer.E.LoadPolicy() 613 }() 614 615 // remove spindle members first 616 err = db.RemoveSpindleMember( 617 tx, 618 orm.FilterEq("owner", did), 619 orm.FilterEq("instance", instance), 620 ) 621 if err != nil { 622 return err 623 } 624 625 err = db.DeleteSpindle( 626 tx, 627 orm.FilterEq("owner", did), 628 orm.FilterEq("instance", instance), 629 ) 630 if err != nil { 631 return err 632 } 633 634 if spindle.Verified != nil { 635 err = i.Enforcer.RemoveSpindle(instance) 636 if err != nil { 637 return err 638 } 639 } 640 641 err = tx.Commit() 642 if err != nil { 643 return err 644 } 645 646 err = i.Enforcer.E.SavePolicy() 647 if err != nil { 648 return err 649 } 650 } 651 652 return nil 653} 654 655func (i *Ingester) ingestString(e *jmodels.Event) error { 656 did := e.Did 657 rkey := e.Commit.RKey 658 659 var err error 660 661 l := i.Logger.With("handler", "ingestString", "nsid", e.Commit.Collection, "did", did, "rkey", rkey) 662 l.Info("ingesting record") 663 664 ddb, ok := i.Db.Execer.(*db.DB) 665 if !ok { 666 return fmt.Errorf("failed to index string record, invalid db cast") 667 } 668 669 switch e.Commit.Operation { 670 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 671 raw := json.RawMessage(e.Commit.Record) 672 record := tangled.String{} 673 err = json.Unmarshal(raw, &record) 674 if err != nil { 675 l.Error("invalid record", "err", err) 676 return err 677 } 678 679 string := models.StringFromRecord(did, rkey, record) 680 681 if err = i.Validator.ValidateString(&string); err != nil { 682 l.Error("invalid record", "err", err) 683 return err 684 } 685 686 if err = db.AddString(ddb, string); err != nil { 687 l.Error("failed to add string", "err", err) 688 return err 689 } 690 691 return nil 692 693 case jmodels.CommitOperationDelete: 694 if err := db.DeleteString( 695 ddb, 696 orm.FilterEq("did", did), 697 orm.FilterEq("rkey", rkey), 698 ); err != nil { 699 l.Error("failed to delete", "err", err) 700 return fmt.Errorf("failed to delete string record: %w", err) 701 } 702 703 return nil 704 } 705 706 return nil 707} 708 709func (i *Ingester) ingestKnotMember(e *jmodels.Event) error { 710 did := e.Did 711 var err error 712 713 l := i.Logger.With("handler", "ingestKnotMember") 714 l = l.With("nsid", e.Commit.Collection) 715 716 switch e.Commit.Operation { 717 case jmodels.CommitOperationCreate: 718 raw := json.RawMessage(e.Commit.Record) 719 record := tangled.KnotMember{} 720 err = json.Unmarshal(raw, &record) 721 if err != nil { 722 l.Error("invalid record", "err", err) 723 return err 724 } 725 726 // only knot owner can invite to knots 727 ok, err := i.Enforcer.IsKnotInviteAllowed(did, record.Domain) 728 if err != nil || !ok { 729 return fmt.Errorf("failed to enforce permissions: %w", err) 730 } 731 732 memberId, err := i.IdResolver.ResolveIdent(context.Background(), record.Subject) 733 if err != nil { 734 return err 735 } 736 737 if memberId.Handle.IsInvalidHandle() { 738 return err 739 } 740 741 err = i.Enforcer.AddKnotMember(record.Domain, memberId.DID.String()) 742 if err != nil { 743 return fmt.Errorf("failed to update ACLs: %w", err) 744 } 745 746 l.Info("added knot member") 747 case jmodels.CommitOperationDelete: 748 // we don't store knot members in a table (like we do for spindle) 749 // and we can't remove this just yet. possibly fixed if we switch 750 // to either: 751 // 1. a knot_members table like with spindle and store the rkey 752 // 2. use the knot host as the rkey 753 // 754 // TODO: implement member deletion 755 l.Info("skipping knot member delete", "did", did, "rkey", e.Commit.RKey) 756 } 757 758 return nil 759} 760 761func (i *Ingester) ingestKnot(e *jmodels.Event) error { 762 did := e.Did 763 var err error 764 765 l := i.Logger.With("handler", "ingestKnot") 766 l = l.With("nsid", e.Commit.Collection) 767 768 switch e.Commit.Operation { 769 case jmodels.CommitOperationCreate: 770 raw := json.RawMessage(e.Commit.Record) 771 record := tangled.Knot{} 772 err = json.Unmarshal(raw, &record) 773 if err != nil { 774 l.Error("invalid record", "err", err) 775 return err 776 } 777 778 domain := e.Commit.RKey 779 780 ddb, ok := i.Db.Execer.(*db.DB) 781 if !ok { 782 return fmt.Errorf("failed to index profile record, invalid db cast") 783 } 784 785 err := db.AddKnot(ddb, domain, did) 786 if err != nil { 787 l.Error("failed to add knot to db", "err", err, "domain", domain) 788 return err 789 } 790 791 err = retry.Do( 792 func() error { 793 return serververify.RunVerification(context.Background(), domain, did, i.Config.Core.Dev) 794 }, 795 retry.Attempts(5), retry.Delay(5*time.Second), retry.MaxDelay(80*time.Second), 796 retry.DelayType(retry.BackOffDelay), retry.LastErrorOnly(true), 797 ) 798 if err != nil { 799 l.Error("failed to verify knot after retries", "err", err, "domain", domain) 800 return err 801 } 802 803 err = serververify.MarkKnotVerified(ddb, i.Enforcer, domain, did) 804 if err != nil { 805 return fmt.Errorf("failed to mark verified: %w", err) 806 } 807 808 return nil 809 810 case jmodels.CommitOperationDelete: 811 domain := e.Commit.RKey 812 813 ddb, ok := i.Db.Execer.(*db.DB) 814 if !ok { 815 return fmt.Errorf("failed to index knot record, invalid db cast") 816 } 817 818 // get record from db first 819 registrations, err := db.GetRegistrations( 820 ddb, 821 orm.FilterEq("domain", domain), 822 orm.FilterEq("did", did), 823 ) 824 if err != nil { 825 return fmt.Errorf("failed to get registration: %w", err) 826 } 827 if len(registrations) != 1 { 828 return fmt.Errorf("got incorrect number of registrations: %d, expected 1", len(registrations)) 829 } 830 registration := registrations[0] 831 832 tx, err := ddb.Begin() 833 if err != nil { 834 return err 835 } 836 defer func() { 837 tx.Rollback() 838 i.Enforcer.E.LoadPolicy() 839 }() 840 841 err = db.DeleteKnot( 842 tx, 843 orm.FilterEq("did", did), 844 orm.FilterEq("domain", domain), 845 ) 846 if err != nil { 847 return err 848 } 849 850 if registration.Registered != nil { 851 err = i.Enforcer.RemoveKnot(domain) 852 if err != nil { 853 return err 854 } 855 } 856 857 err = tx.Commit() 858 if err != nil { 859 return err 860 } 861 862 err = i.Enforcer.E.SavePolicy() 863 if err != nil { 864 return err 865 } 866 } 867 868 return nil 869} 870func (i *Ingester) ingestIssue(ctx context.Context, e *jmodels.Event) error { 871 did := e.Did 872 rkey := e.Commit.RKey 873 874 var err error 875 876 l := i.Logger.With("handler", "ingestIssue", "nsid", e.Commit.Collection, "did", did, "rkey", rkey) 877 l.Info("ingesting record") 878 879 ddb, ok := i.Db.Execer.(*db.DB) 880 if !ok { 881 return fmt.Errorf("failed to index issue record, invalid db cast") 882 } 883 884 switch e.Commit.Operation { 885 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 886 raw := json.RawMessage(e.Commit.Record) 887 record := tangled.RepoIssue{} 888 err = json.Unmarshal(raw, &record) 889 if err != nil { 890 l.Error("invalid record", "err", err) 891 return err 892 } 893 894 issue := models.IssueFromRecord(did, rkey, record) 895 896 if issue.RepoAt == "" { 897 return fmt.Errorf("issue record has no repo field") 898 } 899 900 if err := i.Validator.ValidateIssue(&issue); err != nil { 901 return fmt.Errorf("failed to validate issue: %w", err) 902 } 903 904 if record.Repo != nil { 905 repo, repoErr := db.GetRepoByAtUri(i.Db, *record.Repo) 906 if repoErr == nil && repo.RepoDid != "" { 907 if enqErr := db.EnqueuePdsRewrite(i.Db, did, repo.RepoDid, tangled.RepoIssueNSID, rkey, *record.Repo); enqErr != nil { 908 l.Warn("failed to enqueue PDS rewrite for issue", "err", enqErr, "did", did, "repoDid", repo.RepoDid) 909 } 910 } 911 } 912 913 tx, err := ddb.BeginTx(ctx, nil) 914 if err != nil { 915 l.Error("failed to begin transaction", "err", err) 916 return err 917 } 918 defer tx.Rollback() 919 920 err = db.PutIssue(tx, &issue) 921 if err != nil { 922 l.Error("failed to create issue", "err", err) 923 return err 924 } 925 926 err = tx.Commit() 927 if err != nil { 928 l.Error("failed to commit txn", "err", err) 929 return err 930 } 931 932 return nil 933 934 case jmodels.CommitOperationDelete: 935 tx, err := ddb.BeginTx(ctx, nil) 936 if err != nil { 937 l.Error("failed to begin transaction", "err", err) 938 return err 939 } 940 defer tx.Rollback() 941 942 if err := db.DeleteIssues( 943 tx, 944 did, 945 rkey, 946 ); err != nil { 947 l.Error("failed to delete", "err", err) 948 return fmt.Errorf("failed to delete issue record: %w", err) 949 } 950 if err := tx.Commit(); err != nil { 951 l.Error("failed to commit txn", "err", err) 952 return err 953 } 954 955 return nil 956 } 957 958 return nil 959} 960 961func (i *Ingester) ingestPull(ctx context.Context, e *jmodels.Event) error { 962 did := e.Did 963 rkey := e.Commit.RKey 964 965 var err error 966 967 l := i.Logger.With("handler", "ingestPull", "nsid", e.Commit.Collection, "did", did, "rkey", rkey) 968 l.Info("ingesting record") 969 970 ddb, ok := i.Db.Execer.(*db.DB) 971 if !ok { 972 return fmt.Errorf("failed to index pull record, invalid db cast") 973 } 974 975 switch e.Commit.Operation { 976 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 977 raw := json.RawMessage(e.Commit.Record) 978 record := tangled.RepoPull{} 979 err = json.Unmarshal(raw, &record) 980 if err != nil { 981 l.Error("invalid record", "err", err) 982 return err 983 } 984 985 ownerId, err := i.IdResolver.ResolveIdent(ctx, did) 986 if err != nil { 987 l.Error("failed to resolve did") 988 return err 989 } 990 991 // go through and fetch all blobs in parallel 992 readers := make([]*io.ReadCloser, len(record.Rounds)) 993 var mu sync.Mutex 994 995 g, gctx := errgroup.WithContext(ctx) 996 997 for idx, b := range record.Rounds { 998 g.Go(func() error { 999 // for some reason, a blob is empty 1000 if b.PatchBlob == nil { 1001 return fmt.Errorf("missing patchBlob in round %d", idx) 1002 } 1003 1004 ownerPds := ownerId.PDSEndpoint() 1005 url, _ := url.Parse(fmt.Sprintf("%s/xrpc/com.atproto.sync.getBlob", ownerPds)) 1006 q := url.Query() 1007 q.Set("cid", b.PatchBlob.Ref.String()) 1008 q.Set("did", did) 1009 url.RawQuery = q.Encode() 1010 1011 req, err := http.NewRequestWithContext(gctx, http.MethodGet, url.String(), nil) 1012 if err != nil { 1013 l.Error("failed to create request") 1014 return err 1015 } 1016 req.Header.Set("Content-Type", "application/json") 1017 1018 resp, err := http.DefaultClient.Do(req) 1019 if err != nil { 1020 l.Error("failed to make request") 1021 return err 1022 } 1023 1024 mu.Lock() 1025 readers[idx] = &resp.Body 1026 mu.Unlock() 1027 1028 return nil 1029 }) 1030 } 1031 1032 if err := g.Wait(); err != nil { 1033 for _, r := range readers { 1034 if r != nil && *r != nil { 1035 (*r).Close() 1036 } 1037 } 1038 return err 1039 } 1040 1041 defer func() { 1042 for _, r := range readers { 1043 if r != nil && *r != nil { 1044 (*r).Close() 1045 } 1046 } 1047 }() 1048 1049 pull := models.PullFromRecord(did, rkey, record, readers) 1050 if err := i.Validator.ValidatePull(&pull); err != nil { 1051 return fmt.Errorf("failed to validate pull: %w", err) 1052 } 1053 1054 tx, err := ddb.BeginTx(ctx, nil) 1055 if err != nil { 1056 l.Error("failed to begin transaction", "err", err) 1057 return err 1058 } 1059 defer tx.Rollback() 1060 1061 err = db.PutPull(tx, &pull) 1062 if err != nil { 1063 l.Error("failed to create pull", "err", err) 1064 return err 1065 } 1066 1067 err = tx.Commit() 1068 if err != nil { 1069 l.Error("failed to commit txn", "err", err) 1070 return err 1071 } 1072 1073 return nil 1074 1075 case jmodels.CommitOperationDelete: 1076 tx, err := ddb.BeginTx(ctx, nil) 1077 if err != nil { 1078 l.Error("failed to begin transaction", "err", err) 1079 return err 1080 } 1081 defer tx.Rollback() 1082 1083 if err := db.AbandonPulls( 1084 tx, 1085 orm.FilterEq("owner_did", did), 1086 orm.FilterEq("rkey", rkey), 1087 ); err != nil { 1088 l.Error("failed to abandon", "err", err) 1089 return fmt.Errorf("failed to abandon pull record: %w", err) 1090 } 1091 if err := tx.Commit(); err != nil { 1092 l.Error("failed to commit txn", "err", err) 1093 return err 1094 } 1095 1096 return nil 1097 } 1098 1099 return nil 1100} 1101 1102func (i *Ingester) ingestIssueComment(e *jmodels.Event) error { 1103 did := e.Did 1104 rkey := e.Commit.RKey 1105 1106 var err error 1107 1108 l := i.Logger.With("handler", "ingestIssueComment", "nsid", e.Commit.Collection, "did", did, "rkey", rkey) 1109 l.Info("ingesting record") 1110 1111 ddb, ok := i.Db.Execer.(*db.DB) 1112 if !ok { 1113 return fmt.Errorf("failed to index issue comment record, invalid db cast") 1114 } 1115 1116 switch e.Commit.Operation { 1117 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 1118 raw := json.RawMessage(e.Commit.Record) 1119 record := tangled.RepoIssueComment{} 1120 err = json.Unmarshal(raw, &record) 1121 if err != nil { 1122 return fmt.Errorf("invalid record: %w", err) 1123 } 1124 1125 comment, err := models.IssueCommentFromRecord(did, rkey, record) 1126 if err != nil { 1127 return fmt.Errorf("failed to parse comment from record: %w", err) 1128 } 1129 1130 if err := i.Validator.ValidateIssueComment(comment); err != nil { 1131 return fmt.Errorf("failed to validate comment: %w", err) 1132 } 1133 1134 tx, err := ddb.Begin() 1135 if err != nil { 1136 return fmt.Errorf("failed to start transaction: %w", err) 1137 } 1138 defer tx.Rollback() 1139 1140 _, err = db.AddIssueComment(tx, *comment) 1141 if err != nil { 1142 return fmt.Errorf("failed to create issue comment: %w", err) 1143 } 1144 1145 return tx.Commit() 1146 1147 case jmodels.CommitOperationDelete: 1148 if err := db.DeleteIssueComments( 1149 ddb, 1150 orm.FilterEq("did", did), 1151 orm.FilterEq("rkey", rkey), 1152 ); err != nil { 1153 return fmt.Errorf("failed to delete issue comment record: %w", err) 1154 } 1155 1156 return nil 1157 } 1158 1159 return nil 1160} 1161 1162func (i *Ingester) ingestLabelDefinition(e *jmodels.Event) error { 1163 did := e.Did 1164 rkey := e.Commit.RKey 1165 1166 var err error 1167 1168 l := i.Logger.With("handler", "ingestLabelDefinition", "nsid", e.Commit.Collection, "did", did, "rkey", rkey) 1169 l.Info("ingesting record") 1170 1171 ddb, ok := i.Db.Execer.(*db.DB) 1172 if !ok { 1173 return fmt.Errorf("failed to index label definition, invalid db cast") 1174 } 1175 1176 switch e.Commit.Operation { 1177 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 1178 raw := json.RawMessage(e.Commit.Record) 1179 record := tangled.LabelDefinition{} 1180 err = json.Unmarshal(raw, &record) 1181 if err != nil { 1182 return fmt.Errorf("invalid record: %w", err) 1183 } 1184 1185 def, err := models.LabelDefinitionFromRecord(did, rkey, record) 1186 if err != nil { 1187 return fmt.Errorf("failed to parse labeldef from record: %w", err) 1188 } 1189 1190 if err := i.Validator.ValidateLabelDefinition(def); err != nil { 1191 return fmt.Errorf("failed to validate labeldef: %w", err) 1192 } 1193 1194 _, err = db.AddLabelDefinition(ddb, def) 1195 if err != nil { 1196 return fmt.Errorf("failed to create labeldef: %w", err) 1197 } 1198 1199 return nil 1200 1201 case jmodels.CommitOperationDelete: 1202 if err := db.DeleteLabelDefinition( 1203 ddb, 1204 orm.FilterEq("did", did), 1205 orm.FilterEq("rkey", rkey), 1206 ); err != nil { 1207 return fmt.Errorf("failed to delete labeldef record: %w", err) 1208 } 1209 1210 return nil 1211 } 1212 1213 return nil 1214} 1215 1216func (i *Ingester) ingestLabelOp(e *jmodels.Event) error { 1217 did := e.Did 1218 rkey := e.Commit.RKey 1219 1220 var err error 1221 1222 l := i.Logger.With("handler", "ingestLabelOp", "nsid", e.Commit.Collection, "did", did, "rkey", rkey) 1223 l.Info("ingesting record") 1224 1225 ddb, ok := i.Db.Execer.(*db.DB) 1226 if !ok { 1227 return fmt.Errorf("failed to index label op, invalid db cast") 1228 } 1229 1230 switch e.Commit.Operation { 1231 case jmodels.CommitOperationCreate: 1232 raw := json.RawMessage(e.Commit.Record) 1233 record := tangled.LabelOp{} 1234 err = json.Unmarshal(raw, &record) 1235 if err != nil { 1236 return fmt.Errorf("invalid record: %w", err) 1237 } 1238 1239 subject := syntax.ATURI(record.Subject) 1240 collection := subject.Collection() 1241 1242 var repo *models.Repo 1243 switch collection { 1244 case tangled.RepoIssueNSID: 1245 i, err := db.GetIssues(ddb, orm.FilterEq("at_uri", subject)) 1246 if err != nil || len(i) != 1 { 1247 return fmt.Errorf("failed to find subject: %w || subject count %d", err, len(i)) 1248 } 1249 repo = i[0].Repo 1250 default: 1251 return fmt.Errorf("unsupported label subject: %s", collection) 1252 } 1253 1254 actx, err := db.NewLabelApplicationCtx(ddb, orm.FilterIn("at_uri", repo.Labels)) 1255 if err != nil { 1256 return fmt.Errorf("failed to build label application ctx: %w", err) 1257 } 1258 1259 ops := models.LabelOpsFromRecord(did, rkey, record) 1260 1261 for _, o := range ops { 1262 def, ok := actx.Defs[o.OperandKey] 1263 if !ok { 1264 return fmt.Errorf("failed to find label def for key: %s, expected: %q", o.OperandKey, slices.Collect(maps.Keys(actx.Defs))) 1265 } 1266 if err := i.Validator.ValidateLabelOp(def, repo, &o); err != nil { 1267 return fmt.Errorf("failed to validate labelop: %w", err) 1268 } 1269 } 1270 1271 tx, err := ddb.Begin() 1272 if err != nil { 1273 return err 1274 } 1275 defer tx.Rollback() 1276 1277 for _, o := range ops { 1278 _, err = db.AddLabelOp(tx, &o) 1279 if err != nil { 1280 return fmt.Errorf("failed to add labelop: %w", err) 1281 } 1282 } 1283 1284 if err = tx.Commit(); err != nil { 1285 return err 1286 } 1287 } 1288 1289 return nil 1290}