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