This repository has no description
0

Configure Feed

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

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