This repository has no description
0

Configure Feed

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

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