This repository has no description
0

Configure Feed

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

core / appview / ingester.go
49 kB 1896 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/knotacl" 31 "tangled.org/core/appview/mentions" 32 "tangled.org/core/appview/models" 33 "tangled.org/core/appview/notify" 34 "tangled.org/core/appview/repoverify" 35 "tangled.org/core/appview/serververify" 36 "tangled.org/core/idresolver" 37 "tangled.org/core/orm" 38 "tangled.org/core/rbac" 39) 40 41type Ingester struct { 42 Ctx context.Context 43 Db *db.DB 44 Enforcer *rbac.Enforcer 45 Acl *knotacl.Service 46 IdResolver *idresolver.Resolver 47 Cache *cache.Cache 48 Config *config.Config 49 Logger *slog.Logger 50 MentionsResolver *mentions.Resolver 51 Notifier notify.Notifier 52 Verifier repoverify.Verifier 53} 54 55type processFunc func(ctx context.Context, e *jmodels.Event) error 56 57func (i *Ingester) Ingest() processFunc { 58 return func(ctx context.Context, e *jmodels.Event) error { 59 var err error 60 61 l := i.Logger.With("kind", e.Kind) 62 switch e.Kind { 63 case jmodels.EventKindAccount: 64 // TODO: sync account state to db 65 if e.Account.Active { 66 break 67 } 68 // TODO: revoke sessions by DID 69 if *e.Account.Status == "deactivated" { 70 err = i.IdResolver.InvalidateIdent(ctx, e.Account.Did) 71 } 72 case jmodels.EventKindIdentity: 73 err = i.IdResolver.InvalidateIdent(ctx, e.Identity.Did) 74 case jmodels.EventKindCommit: 75 l = l.With( 76 "nsid", e.Commit.Collection, 77 "did", e.Did, 78 "rkey", e.Commit.RKey, 79 "op", e.Commit.Operation, 80 ) 81 switch e.Commit.Collection { 82 case tangled.GraphFollowNSID: 83 err = i.ingestFollow(e, l) 84 case tangled.GraphVouchNSID: 85 err = i.ingestVouch(ctx, e, l) 86 case tangled.FeedStarNSID: 87 err = i.ingestStar(ctx, e, l) 88 case tangled.FeedReactionNSID: 89 err = i.ingestReaction(e, l) 90 case tangled.PublicKeyNSID: 91 err = i.ingestPublicKey(e, l) 92 case tangled.RepoArtifactNSID: 93 err = i.ingestArtifact(ctx, e, l) 94 case tangled.ActorProfileNSID: 95 err = i.ingestProfile(ctx, e, l) 96 case tangled.SpindleMemberNSID: 97 err = i.ingestSpindleMember(ctx, e, l) 98 case tangled.SpindleNSID: 99 err = i.ingestSpindle(ctx, e, l) 100 case tangled.KnotMemberNSID: 101 err = i.ingestKnotMember(ctx, e, l) 102 case tangled.KnotNSID: 103 err = i.ingestKnot(ctx, e, l) 104 case tangled.StringNSID: 105 err = i.ingestString(e, l) 106 case tangled.RepoIssueNSID: 107 err = i.ingestIssue(ctx, e, l) 108 case tangled.RepoPullNSID: 109 err = i.ingestPull(ctx, e, l) 110 case tangled.FeedCommentNSID: 111 err = i.ingestComment(e, l) 112 case tangled.RepoIssueCommentNSID: 113 err = i.ingestIssueComment(e, l) 114 case tangled.RepoPullCommentNSID: 115 err = i.ingestPullComment(e, l) 116 case tangled.LabelDefinitionNSID: 117 err = i.ingestLabelDefinition(e, l) 118 case tangled.LabelOpNSID: 119 err = i.ingestLabelOp(ctx, e, l) 120 case tangled.RepoNSID: 121 err = i.ingestRepo(ctx, e, l) 122 } 123 } 124 125 if err != nil { 126 l.Warn("failed to ingest record, skipping", "err", err) 127 } 128 129 lastTimeUs := e.TimeUS + 1 130 if saveErr := i.Db.SaveLastTimeUs(lastTimeUs); saveErr != nil { 131 l.Error("failed to save cursor", "err", saveErr) 132 } 133 134 return nil 135 } 136} 137 138func (i *Ingester) resolveRepoRef(ref string) (*models.Repo, error) { 139 if strings.HasPrefix(ref, "did:") { 140 return db.GetRepoByDid(i.Db, ref) 141 } 142 return db.GetRepoByAtUri(i.Db, ref) 143} 144 145func (i *Ingester) resolveOldFormatStar(raw json.RawMessage, star *models.Star, l *slog.Logger) (bool, error) { 146 var legacy struct { 147 Subject *string `json:"subject"` 148 SubjectDid *string `json:"subjectDid"` 149 } 150 if err := json.Unmarshal(raw, &legacy); err != nil { 151 return false, err 152 } 153 154 switch { 155 case legacy.SubjectDid != nil: 156 repo, err := i.resolveRepoRef(*legacy.SubjectDid) 157 if err != nil { 158 l.Warn("skipping old-format star for unknown repo", "subjectDid", *legacy.SubjectDid) 159 return false, nil 160 } 161 star.SubjectType = models.StarSubjectRepo 162 star.Subject = repo.RepoDid 163 return true, nil 164 165 case legacy.Subject != nil: 166 uri, err := syntax.ParseATURI(*legacy.Subject) 167 if err != nil { 168 return false, fmt.Errorf("invalid old-format star subject: %w", err) 169 } 170 switch uri.Collection().String() { 171 case tangled.RepoNSID: 172 repo, err := db.GetRepoByAtUri(i.Db, uri.String()) 173 if err != nil { 174 l.Warn("skipping old-format star for unknown repo", "subject", *legacy.Subject) 175 return false, nil 176 } 177 star.SubjectType = models.StarSubjectRepo 178 star.Subject = repo.RepoDid 179 return true, nil 180 default: 181 star.SubjectType = models.StarSubjectString 182 star.Subject = *legacy.Subject 183 return true, nil 184 } 185 186 default: 187 return false, fmt.Errorf("old-format star has neither subject nor subjectDid") 188 } 189} 190 191func (i *Ingester) ingestStar(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { 192 var err error 193 did := e.Did 194 195 l = l.With("handler", "ingestStar") 196 197 switch e.Commit.Operation { 198 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 199 raw := json.RawMessage(e.Commit.Record) 200 record := tangled.FeedStar{} 201 unmarshalErr := json.Unmarshal(raw, &record) 202 203 star := models.Star{ 204 Did: did, 205 Rkey: e.Commit.RKey, 206 } 207 208 switch { 209 case unmarshalErr != nil: 210 resolved, resolveErr := i.resolveOldFormatStar(raw, &star, l) 211 if resolveErr != nil { 212 l.Error("invalid record", "newFmtErr", unmarshalErr, "oldFmtErr", resolveErr) 213 return unmarshalErr 214 } 215 if !resolved { 216 return nil 217 } 218 219 case record.Subject == nil: 220 return fmt.Errorf("star record has nil subject") 221 222 case record.Subject.FeedStar_Repo != nil: 223 repo, repoErr := i.resolveRepoRef(record.Subject.FeedStar_Repo.Did) 224 if repoErr != nil { 225 l.Warn("skipping star for unknown repo", "did", record.Subject.FeedStar_Repo.Did) 226 return nil 227 } 228 star.SubjectType = models.StarSubjectRepo 229 star.Subject = repo.RepoDid 230 231 case record.Subject.FeedStar_String != nil: 232 star.SubjectType = models.StarSubjectString 233 star.Subject = record.Subject.FeedStar_String.Uri 234 235 default: 236 return fmt.Errorf("star record has empty subject union") 237 } 238 239 err = db.UpsertStar(i.Db, star) 240 case jmodels.CommitOperationDelete: 241 err = db.DeleteStarByRkey(i.Db, did, e.Commit.RKey) 242 } 243 244 if err != nil { 245 return fmt.Errorf("failed to %s star record: %w", e.Commit.Operation, err) 246 } 247 248 l.Info("ingested record") 249 return nil 250} 251 252func (i *Ingester) ingestFollow(e *jmodels.Event, l *slog.Logger) error { 253 var err error 254 did := e.Did 255 256 l = l.With("handler", "ingestFollow") 257 258 switch e.Commit.Operation { 259 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 260 raw := json.RawMessage(e.Commit.Record) 261 record := tangled.GraphFollow{} 262 err = json.Unmarshal(raw, &record) 263 if err != nil { 264 l.Error("invalid record", "err", err) 265 return err 266 } 267 268 err = db.UpsertFollow(i.Db, models.Follow{ 269 UserDid: did, 270 SubjectDid: record.Subject, 271 Rkey: e.Commit.RKey, 272 }) 273 case jmodels.CommitOperationDelete: 274 err = db.DeleteFollowByRkey(i.Db, did, e.Commit.RKey) 275 } 276 277 if err != nil { 278 return fmt.Errorf("failed to %s follow record: %w", e.Commit.Operation, err) 279 } 280 281 l.Info("ingested record") 282 return nil 283} 284 285func (i *Ingester) ingestVouch(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { 286 var err error 287 did := e.Did 288 289 l = l.With("handler", "ingestVouch") 290 291 switch e.Commit.Operation { 292 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 293 raw := json.RawMessage(e.Commit.Record) 294 record := tangled.GraphVouch{} 295 err = json.Unmarshal(raw, &record) 296 if err != nil { 297 l.Error("invalid record", "err", err) 298 return err 299 } 300 301 // rkey is the subject_did being vouched for/denounced 302 subjectDID := e.Commit.RKey 303 304 _, err = syntax.ParseDID(subjectDID) 305 if err != nil { 306 l.Error("invalid subject_did in rkey", "err", err, "rkey", subjectDID) 307 return fmt.Errorf("invalid subject_did: %w", err) 308 } 309 310 if did == subjectDID { 311 l.Warn("attempted self-vouch", "did", did) 312 return fmt.Errorf("cannot vouch for self") 313 } 314 315 subjectId, err := i.IdResolver.ResolveIdent(ctx, subjectDID) 316 if err != nil { 317 return err 318 } 319 320 if subjectId.Handle.IsInvalidHandle() { 321 return err 322 } 323 324 kind, err := models.ParseVouchKind(record.Kind) 325 if err != nil { 326 l.Error("invalid kind", "kind", kind) 327 return fmt.Errorf("invalid kind: %s", kind) 328 } 329 330 recordCid, err := cid.Parse(e.Commit.CID) 331 if err != nil { 332 l.Error("invalid cid", "err", err, "cid", e.Commit.CID) 333 return fmt.Errorf("invalid cid: %w", err) 334 } 335 336 var evidences []syntax.ATURI 337 for _, raw := range record.Evidences { 338 uri, parseErr := syntax.ParseATURI(raw) 339 if parseErr != nil { 340 l.Warn("invalid evidence AT-URI, skipping", "uri", raw, "err", parseErr) 341 continue 342 } 343 evidences = append(evidences, uri) 344 } 345 346 tx, txErr := i.Db.Begin() 347 if txErr != nil { 348 return fmt.Errorf("failed to start transaction: %w", txErr) 349 } 350 351 addErr := db.AddVouch(tx, &models.Vouch{ 352 Did: syntax.DID(did), 353 SubjectDid: subjectId.DID, 354 Cid: recordCid, 355 Kind: kind, 356 Reason: record.Reason, 357 Evidences: evidences, 358 }) 359 if addErr != nil { 360 tx.Rollback() 361 err = addErr 362 } else { 363 err = tx.Commit() 364 } 365 366 case jmodels.CommitOperationDelete: 367 err = db.DeleteVouchByRkey(i.Db, did, e.Commit.RKey) 368 } 369 370 if err != nil { 371 return fmt.Errorf("failed to %s vouch record: %w", e.Commit.Operation, err) 372 } 373 374 l.Info("ingested record") 375 return nil 376} 377 378func (i *Ingester) ingestPublicKey(e *jmodels.Event, l *slog.Logger) error { 379 did := e.Did 380 var err error 381 382 l = l.With("handler", "ingestPublicKey") 383 384 switch e.Commit.Operation { 385 case jmodels.CommitOperationCreate: 386 l.Debug("processing add of pubkey") 387 raw := json.RawMessage(e.Commit.Record) 388 record := tangled.PublicKey{} 389 err = json.Unmarshal(raw, &record) 390 if err != nil { 391 l.Error("invalid record", "err", err) 392 return err 393 } 394 395 name := record.Name 396 key := record.Key 397 err = db.AddPublicKey(i.Db, did, name, key, e.Commit.RKey) 398 case jmodels.CommitOperationUpdate: 399 l.Debug("processing update of pubkey") 400 raw := json.RawMessage(e.Commit.Record) 401 record := tangled.PublicKey{} 402 err = json.Unmarshal(raw, &record) 403 if err != nil { 404 l.Error("invalid record", "err", err) 405 return err 406 } 407 408 name := record.Name 409 key := record.Key 410 err = db.UpdatePublicKey(i.Db, did, name, key, e.Commit.RKey) 411 case jmodels.CommitOperationDelete: 412 l.Debug("processing delete of pubkey") 413 err = db.DeletePublicKeyByRkey(i.Db, did, e.Commit.RKey) 414 } 415 416 if err != nil { 417 return fmt.Errorf("failed to %s pubkey record: %w", e.Commit.Operation, err) 418 } 419 420 l.Info("ingested record") 421 return nil 422} 423 424func (i *Ingester) ingestArtifact(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { 425 did := e.Did 426 var err error 427 428 l = l.With("handler", "ingestArtifact") 429 430 switch e.Commit.Operation { 431 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 432 raw := json.RawMessage(e.Commit.Record) 433 record := tangled.RepoArtifact{} 434 err = json.Unmarshal(raw, &record) 435 if err != nil { 436 l.Error("invalid record", "err", err) 437 return err 438 } 439 440 var repo *models.Repo 441 if record.RepoDid != nil && *record.RepoDid != "" { 442 repo, err = db.GetRepoByDid(i.Db, *record.RepoDid) 443 if err != nil && !errors.Is(err, sql.ErrNoRows) { 444 return fmt.Errorf("failed to look up repo by DID %s: %w", *record.RepoDid, err) 445 } 446 } 447 if repo == nil && record.Repo != nil { 448 repoAt, parseErr := syntax.ParseATURI(*record.Repo) 449 if parseErr != nil { 450 return parseErr 451 } 452 repo, err = db.GetRepoByAtUri(i.Db, repoAt.String()) 453 if err != nil { 454 return err 455 } 456 } 457 if repo == nil { 458 return fmt.Errorf("artifact record has neither valid repoDid nor repo field") 459 } 460 461 allowed, permErr := i.Acl.HasRepoPermissionErr(ctx, repo, did, "repo:push") 462 if permErr != nil { 463 l.Warn("ingesting artifact without permission check", "did", did, "repo", repo.RepoIdentifier(), "err", permErr) 464 } else if !allowed { 465 l.Info("skipping unauthorized artifact", "did", did, "repo", repo.RepoIdentifier()) 466 return nil 467 } 468 469 repoDid := repo.RepoDid 470 if repoDid == "" && record.RepoDid != nil { 471 repoDid = *record.RepoDid 472 } 473 if repoDid != "" && (record.RepoDid == nil || *record.RepoDid == "") && record.Repo != nil { 474 if enqErr := db.EnqueuePdsRecordMigration(ctx, i.Db, "add-repo-did", syntax.DID(did), syntax.NSID(tangled.RepoArtifactNSID), syntax.RecordKey(e.Commit.RKey)); enqErr != nil { 475 l.Warn("failed to enqueue PDS rewrite for artifact", "err", enqErr, "did", did, "repoDid", repoDid) 476 } 477 } 478 479 createdAt, parseErr := time.Parse(time.RFC3339, record.CreatedAt) 480 if parseErr != nil { 481 createdAt = time.Now() 482 } 483 484 artifact := models.Artifact{ 485 Did: did, 486 Rkey: e.Commit.RKey, 487 RepoDid: syntax.DID(repo.RepoDid), 488 Tag: plumbing.Hash(record.Tag), 489 CreatedAt: createdAt, 490 BlobCid: cid.Cid(record.Artifact.Ref), 491 Name: record.Name, 492 Size: uint64(record.Artifact.Size), 493 MimeType: record.Artifact.MimeType, 494 } 495 496 err = db.AddArtifact(i.Db, artifact) 497 case jmodels.CommitOperationDelete: 498 err = db.DeleteArtifact(i.Db, orm.FilterEq("did", did), orm.FilterEq("rkey", e.Commit.RKey)) 499 } 500 501 if err != nil { 502 return fmt.Errorf("failed to %s artifact record: %w", e.Commit.Operation, err) 503 } 504 505 l.Info("ingested record") 506 return nil 507} 508 509func (i *Ingester) ingestProfile(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { 510 did := e.Did 511 var err error 512 513 l = l.With("handler", "ingestProfile") 514 515 if e.Commit.RKey != "self" { 516 return fmt.Errorf("ingestProfile only ingests `self` record") 517 } 518 519 switch e.Commit.Operation { 520 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 521 raw := json.RawMessage(e.Commit.Record) 522 record := tangled.ActorProfile{} 523 err = json.Unmarshal(raw, &record) 524 if err != nil { 525 l.Error("invalid record", "err", err) 526 return err 527 } 528 529 avatar := "" 530 if record.Avatar != nil { 531 avatar = record.Avatar.Ref.String() 532 } 533 534 description := "" 535 if record.Description != nil { 536 description = *record.Description 537 } 538 539 includeBluesky := record.Bluesky 540 541 pronouns := "" 542 if record.Pronouns != nil { 543 pronouns = *record.Pronouns 544 } 545 546 location := "" 547 if record.Location != nil { 548 location = *record.Location 549 } 550 551 var links [5]string 552 for i, l := range record.Links { 553 if i < 5 { 554 links[i] = l 555 } 556 } 557 558 var stats [2]models.VanityStat 559 for i, s := range record.Stats { 560 if i < 2 { 561 stats[i].Kind = models.ParseVanityStatKind(s) 562 } 563 } 564 565 var pinned [6]string 566 for i, r := range record.PinnedRepositories { 567 if i < 6 { 568 pinned[i] = r 569 } 570 } 571 572 var preferredHandle syntax.Handle 573 if record.PreferredHandle != nil { 574 if h, err := syntax.ParseHandle(*record.PreferredHandle); err == nil { 575 ident, identErr := i.IdResolver.ResolveIdent(ctx, did) 576 if identErr == nil && slices.Contains(ident.AlsoKnownAs, "at://"+string(h)) { 577 preferredHandle = h 578 } 579 } 580 } 581 582 profile := models.Profile{ 583 Did: did, 584 Avatar: avatar, 585 Description: description, 586 IncludeBluesky: includeBluesky, 587 Location: location, 588 Links: links, 589 Stats: stats, 590 PinnedRepos: pinned, 591 Pronouns: pronouns, 592 PreferredHandle: preferredHandle, 593 } 594 595 tx, err := i.Db.Begin() 596 if err != nil { 597 return fmt.Errorf("failed to start transaction: %w", err) 598 } 599 defer tx.Rollback() 600 601 err = db.ValidateProfile(tx, &profile) 602 if err != nil { 603 return fmt.Errorf("invalid profile record") 604 } 605 606 err = db.UpsertProfile(tx, &profile) 607 if err != nil { 608 return fmt.Errorf("upserting profile: %w", err) 609 } 610 611 err = tx.Commit() 612 if err != nil { 613 return fmt.Errorf("tx.Commit: %w", err) 614 } 615 if i.Cache != nil { 616 pipe := i.Cache.Pipeline() 617 didKey := fmt.Sprintf(cache.PreferredHandleByDid, did) 618 if preferredHandle != "" { 619 pipe.Set(ctx, didKey, string(preferredHandle), cache.PreferredHandleTTL) 620 pipe.Set(ctx, fmt.Sprintf(cache.PreferredHandleByHandle, string(preferredHandle)), did, cache.PreferredHandleTTL) 621 } else { 622 pipe.Del(ctx, didKey) 623 } 624 if _, execErr := pipe.Exec(ctx); execErr != nil { 625 l.Warn("failed to update preferred handle cache", "err", execErr) 626 } 627 } 628 case jmodels.CommitOperationDelete: 629 tx, beginErr := i.Db.Begin() 630 if beginErr != nil { 631 return fmt.Errorf("failed to start transaction: %w", beginErr) 632 } 633 634 priorHandle, phErr := db.GetPreferredHandle(tx, did) 635 if phErr != nil && !errors.Is(phErr, sql.ErrNoRows) { 636 l.Warn("failed to read prior preferred handle", "err", phErr) 637 } 638 639 err = db.DeleteProfile(tx, did) 640 if err == nil && i.Cache != nil { 641 pipe := i.Cache.Pipeline() 642 pipe.Del(ctx, fmt.Sprintf(cache.PreferredHandleByDid, did)) 643 if priorHandle != "" { 644 pipe.Del(ctx, fmt.Sprintf(cache.PreferredHandleByHandle, string(priorHandle))) 645 } 646 if _, execErr := pipe.Exec(ctx); execErr != nil { 647 l.Warn("failed to evict preferred handle cache", "err", execErr) 648 } 649 } 650 } 651 652 if err != nil { 653 return fmt.Errorf("failed to %s profile record: %w", e.Commit.Operation, err) 654 } 655 656 l.Info("ingested record") 657 return nil 658} 659 660func (i *Ingester) ingestSpindleMember(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { 661 did := e.Did 662 var err error 663 664 l = l.With("handler", "ingestSpindleMember") 665 666 switch e.Commit.Operation { 667 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 668 raw := json.RawMessage(e.Commit.Record) 669 record := tangled.SpindleMember{} 670 err = json.Unmarshal(raw, &record) 671 if err != nil { 672 l.Error("invalid record", "err", err) 673 return err 674 } 675 676 // only spindle owner can invite to spindles 677 ok, err := i.Enforcer.IsSpindleInviteAllowed(did, record.Instance) 678 if err != nil { 679 return fmt.Errorf("failed to check invite permission: %w", err) 680 } 681 if !ok { 682 if verifyErr := i.verifySpindle(ctx, record.Instance, did); verifyErr != nil { 683 return fmt.Errorf("invite denied and verify failed: %w", verifyErr) 684 } 685 ok, err = i.Enforcer.IsSpindleInviteAllowed(did, record.Instance) 686 if err != nil { 687 return fmt.Errorf("failed to re-check invite permission: %w", err) 688 } 689 if !ok { 690 return fmt.Errorf("invite denied for did %s on spindle %s", did, record.Instance) 691 } 692 } 693 694 memberId, err := i.IdResolver.ResolveIdent(ctx, record.Subject) 695 if err != nil { 696 return err 697 } 698 699 if memberId.Handle.IsInvalidHandle() { 700 return fmt.Errorf("invalid handle for member %s", record.Subject) 701 } 702 703 existing, err := db.GetSpindleMembers(i.Db, 704 orm.FilterEq("did", did), 705 orm.FilterEq("rkey", e.Commit.RKey), 706 ) 707 if err != nil { 708 return fmt.Errorf("failed to look up existing member: %w", err) 709 } 710 if len(existing) > 1 { 711 return fmt.Errorf("multiple spindle members with rkey %s", e.Commit.RKey) 712 } 713 714 tx, err := i.Db.Begin() 715 if err != nil { 716 return fmt.Errorf("failed to start txn: %w", err) 717 } 718 committed := false 719 defer func() { 720 if committed { 721 return 722 } 723 tx.Rollback() 724 i.Enforcer.E.LoadPolicy() 725 }() 726 727 if len(existing) == 1 { 728 prev := existing[0] 729 if prev.Instance != record.Instance || prev.Subject != memberId.DID { 730 if err = db.RemoveSpindleMember(tx, 731 orm.FilterEq("did", did), 732 orm.FilterEq("rkey", e.Commit.RKey), 733 ); err != nil { 734 return fmt.Errorf("failed to remove stale row: %w", err) 735 } 736 if err = i.Enforcer.RemoveSpindleMember(prev.Instance, prev.Subject.String()); err != nil { 737 return fmt.Errorf("failed to remove stale ACL: %w", err) 738 } 739 } 740 } 741 742 if err = db.AddSpindleMember(tx, models.SpindleMember{ 743 Did: syntax.DID(did), 744 Rkey: e.Commit.RKey, 745 Instance: record.Instance, 746 Subject: memberId.DID, 747 }); err != nil { 748 return fmt.Errorf("failed to add to db: %w", err) 749 } 750 751 if err = i.Enforcer.AddSpindleMember(record.Instance, memberId.DID.String()); err != nil { 752 return fmt.Errorf("failed to update ACLs: %w", err) 753 } 754 755 if err = tx.Commit(); err != nil { 756 return fmt.Errorf("failed to commit txn: %w", err) 757 } 758 759 if err = i.Enforcer.E.SavePolicy(); err != nil { 760 return fmt.Errorf("failed to save ACLs: %w", err) 761 } 762 committed = true 763 764 l.Info("upserted spindle member") 765 case jmodels.CommitOperationDelete: 766 rkey := e.Commit.RKey 767 768 // get record from db first 769 members, err := db.GetSpindleMembers( 770 i.Db, 771 orm.FilterEq("did", did), 772 orm.FilterEq("rkey", rkey), 773 ) 774 if err != nil || len(members) != 1 { 775 return fmt.Errorf("failed to get member: %w, len(members) = %d", err, len(members)) 776 } 777 member := members[0] 778 779 tx, err := i.Db.Begin() 780 if err != nil { 781 return fmt.Errorf("failed to start txn: %w", err) 782 } 783 committed := false 784 defer func() { 785 if committed { 786 return 787 } 788 tx.Rollback() 789 i.Enforcer.E.LoadPolicy() 790 }() 791 792 // remove record by rkey && update enforcer 793 if err = db.RemoveSpindleMember( 794 tx, 795 orm.FilterEq("did", did), 796 orm.FilterEq("rkey", rkey), 797 ); err != nil { 798 return fmt.Errorf("failed to remove from db: %w", err) 799 } 800 801 // update enforcer 802 err = i.Enforcer.RemoveSpindleMember(member.Instance, member.Subject.String()) 803 if err != nil { 804 return fmt.Errorf("failed to update ACLs: %w", err) 805 } 806 807 if err = tx.Commit(); err != nil { 808 return fmt.Errorf("failed to commit txn: %w", err) 809 } 810 811 if err = i.Enforcer.E.SavePolicy(); err != nil { 812 return fmt.Errorf("failed to save ACLs: %w", err) 813 } 814 committed = true 815 816 l.Info("removed spindle member") 817 } 818 819 return nil 820} 821 822func (i *Ingester) ingestSpindle(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { 823 did := e.Did 824 var err error 825 826 l = l.With("handler", "ingestSpindle") 827 828 switch e.Commit.Operation { 829 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 830 raw := json.RawMessage(e.Commit.Record) 831 record := tangled.Spindle{} 832 err = json.Unmarshal(raw, &record) 833 if err != nil { 834 l.Error("invalid record", "err", err) 835 return err 836 } 837 838 instance := e.Commit.RKey 839 840 err := db.AddSpindle(i.Db, models.Spindle{ 841 Owner: syntax.DID(did), 842 Instance: instance, 843 }) 844 if err != nil { 845 l.Error("failed to add spindle to db", "err", err, "instance", instance) 846 return err 847 } 848 849 if err := i.verifySpindle(ctx, instance, did); err != nil { 850 l.Warn("failed to verify spindle", "instance", instance, "did", did, "err", err) 851 } 852 853 l.Info("ingested record", "instance", instance) 854 return nil 855 856 case jmodels.CommitOperationDelete: 857 instance := e.Commit.RKey 858 859 // get record from db first 860 spindles, err := db.GetSpindles( 861 ctx, 862 i.Db, 863 orm.FilterEq("owner", did), 864 orm.FilterEq("instance", instance), 865 ) 866 if err != nil || len(spindles) != 1 { 867 return fmt.Errorf("failed to get spindles: %w, len(spindles) = %d", err, len(spindles)) 868 } 869 spindle := spindles[0] 870 871 tx, err := i.Db.Begin() 872 if err != nil { 873 return fmt.Errorf("failed to start txn: %w", err) 874 } 875 defer func() { 876 tx.Rollback() 877 i.Enforcer.E.LoadPolicy() 878 }() 879 880 // remove spindle members first 881 err = db.RemoveSpindleMember( 882 tx, 883 orm.FilterEq("owner", did), 884 orm.FilterEq("instance", instance), 885 ) 886 if err != nil { 887 return fmt.Errorf("failed to remove spindle members: %w", err) 888 } 889 890 err = db.DeleteSpindle( 891 tx, 892 orm.FilterEq("owner", did), 893 orm.FilterEq("instance", instance), 894 ) 895 if err != nil { 896 return fmt.Errorf("failed to delete spindle: %w", err) 897 } 898 899 if spindle.Verified != nil { 900 err = i.Enforcer.RemoveSpindle(instance) 901 if err != nil { 902 return fmt.Errorf("failed to remove spindle from enforcer: %w", err) 903 } 904 } 905 906 err = tx.Commit() 907 if err != nil { 908 return fmt.Errorf("failed to commit txn: %w", err) 909 } 910 911 err = i.Enforcer.E.SavePolicy() 912 if err != nil { 913 return fmt.Errorf("failed to save ACLs: %w", err) 914 } 915 916 l.Info("ingested record", "instance", instance) 917 } 918 919 return nil 920} 921 922func (i *Ingester) ingestString(e *jmodels.Event, l *slog.Logger) error { 923 did := e.Did 924 rkey := e.Commit.RKey 925 926 var err error 927 928 l = l.With("handler", "ingestString") 929 930 switch e.Commit.Operation { 931 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 932 raw := json.RawMessage(e.Commit.Record) 933 record := tangled.String{} 934 err = json.Unmarshal(raw, &record) 935 if err != nil { 936 l.Error("invalid record", "err", err) 937 return err 938 } 939 940 string := models.StringFromRecord(did, rkey, record) 941 942 if err = string.Validate(); err != nil { 943 l.Error("invalid record", "err", err) 944 return err 945 } 946 947 if err = db.AddString(i.Db, string); err != nil { 948 l.Error("failed to add string", "err", err) 949 return err 950 } 951 952 l.Info("ingested record") 953 return nil 954 955 case jmodels.CommitOperationDelete: 956 if err := db.DeleteString( 957 i.Db, 958 orm.FilterEq("did", did), 959 orm.FilterEq("rkey", rkey), 960 ); err != nil { 961 l.Error("failed to delete", "err", err) 962 return fmt.Errorf("failed to delete string record: %w", err) 963 } 964 965 l.Info("ingested record") 966 return nil 967 } 968 969 return nil 970} 971 972func (i *Ingester) ingestKnotMember(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { 973 did := e.Did 974 var err error 975 976 l = l.With("handler", "ingestKnotMember") 977 978 switch e.Commit.Operation { 979 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 980 raw := json.RawMessage(e.Commit.Record) 981 record := tangled.KnotMember{} 982 err = json.Unmarshal(raw, &record) 983 if err != nil { 984 l.Error("invalid record", "err", err) 985 return err 986 } 987 988 // only knot owner can invite to knots 989 ok, err := i.Enforcer.IsKnotInviteAllowed(did, record.Domain) 990 if err != nil { 991 return fmt.Errorf("failed to check invite permission: %w", err) 992 } 993 if !ok { 994 if verifyErr := i.verifyKnot(ctx, record.Domain, did); verifyErr != nil { 995 return fmt.Errorf("invite denied and verify failed: %w", verifyErr) 996 } 997 ok, err = i.Enforcer.IsKnotInviteAllowed(did, record.Domain) 998 if err != nil { 999 return fmt.Errorf("failed to re-check invite permission: %w", err) 1000 } 1001 if !ok { 1002 return fmt.Errorf("invite denied for did %s on knot %s", did, record.Domain) 1003 } 1004 } 1005 1006 memberId, err := i.IdResolver.ResolveIdent(ctx, record.Subject) 1007 if err != nil { 1008 return err 1009 } 1010 1011 if memberId.Handle.IsInvalidHandle() { 1012 return fmt.Errorf("invalid handle for member %s", record.Subject) 1013 } 1014 1015 existing, err := db.GetKnotMembers(i.Db, 1016 orm.FilterEq("did", did), 1017 orm.FilterEq("rkey", e.Commit.RKey), 1018 ) 1019 if err != nil { 1020 return fmt.Errorf("failed to look up existing member: %w", err) 1021 } 1022 if len(existing) > 1 { 1023 return fmt.Errorf("multiple knot members with rkey %s", e.Commit.RKey) 1024 } 1025 1026 tx, err := i.Db.Begin() 1027 if err != nil { 1028 return fmt.Errorf("failed to start txn: %w", err) 1029 } 1030 committed := false 1031 defer func() { 1032 if committed { 1033 return 1034 } 1035 tx.Rollback() 1036 i.Enforcer.E.LoadPolicy() 1037 }() 1038 1039 if len(existing) == 1 { 1040 prev := existing[0] 1041 if prev.Domain != record.Domain || prev.Subject != memberId.DID { 1042 if err = db.RemoveKnotMember(tx, 1043 orm.FilterEq("did", did), 1044 orm.FilterEq("rkey", e.Commit.RKey), 1045 ); err != nil { 1046 return fmt.Errorf("failed to remove stale row: %w", err) 1047 } 1048 if err = i.Enforcer.RemoveKnotMember(prev.Domain, prev.Subject.String()); err != nil { 1049 return fmt.Errorf("failed to remove stale ACL: %w", err) 1050 } 1051 } 1052 } 1053 1054 if err = db.AddKnotMember(tx, models.KnotMember{ 1055 Did: syntax.DID(did), 1056 Rkey: e.Commit.RKey, 1057 Domain: record.Domain, 1058 Subject: memberId.DID, 1059 }); err != nil { 1060 return fmt.Errorf("failed to add to db: %w", err) 1061 } 1062 1063 if err = i.Enforcer.AddKnotMember(record.Domain, memberId.DID.String()); err != nil { 1064 return fmt.Errorf("failed to update ACLs: %w", err) 1065 } 1066 1067 if err = tx.Commit(); err != nil { 1068 return fmt.Errorf("failed to commit txn: %w", err) 1069 } 1070 1071 if err = i.Enforcer.E.SavePolicy(); err != nil { 1072 return fmt.Errorf("failed to save ACLs: %w", err) 1073 } 1074 committed = true 1075 1076 l.Info("upserted knot member") 1077 case jmodels.CommitOperationDelete: 1078 rkey := e.Commit.RKey 1079 1080 members, err := db.GetKnotMembers( 1081 i.Db, 1082 orm.FilterEq("did", did), 1083 orm.FilterEq("rkey", rkey), 1084 ) 1085 if err != nil { 1086 return fmt.Errorf("failed to look up knot member with rkey %s: %w", rkey, err) 1087 } 1088 if len(members) == 0 { 1089 l.Info("knot member already removed", "rkey", rkey) 1090 return nil 1091 } 1092 if len(members) > 1 { 1093 return fmt.Errorf("multiple knot members with rkey %s", rkey) 1094 } 1095 member := members[0] 1096 1097 tx, err := i.Db.Begin() 1098 if err != nil { 1099 return fmt.Errorf("failed to start txn: %w", err) 1100 } 1101 committed := false 1102 defer func() { 1103 if committed { 1104 return 1105 } 1106 tx.Rollback() 1107 i.Enforcer.E.LoadPolicy() 1108 }() 1109 1110 if err = db.RemoveKnotMember( 1111 tx, 1112 orm.FilterEq("did", did), 1113 orm.FilterEq("rkey", rkey), 1114 ); err != nil { 1115 return fmt.Errorf("failed to remove from db: %w", err) 1116 } 1117 1118 if err = i.Enforcer.RemoveKnotMember(member.Domain, member.Subject.String()); err != nil { 1119 return fmt.Errorf("failed to update ACLs: %w", err) 1120 } 1121 1122 if err = tx.Commit(); err != nil { 1123 return fmt.Errorf("failed to commit txn: %w", err) 1124 } 1125 1126 if err = i.Enforcer.E.SavePolicy(); err != nil { 1127 return fmt.Errorf("failed to save ACLs: %w", err) 1128 } 1129 committed = true 1130 1131 l.Info("removed knot member") 1132 } 1133 1134 return nil 1135} 1136 1137func (i *Ingester) ingestKnot(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { 1138 did := e.Did 1139 var err error 1140 1141 l = l.With("handler", "ingestKnot") 1142 1143 switch e.Commit.Operation { 1144 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 1145 raw := json.RawMessage(e.Commit.Record) 1146 record := tangled.Knot{} 1147 err = json.Unmarshal(raw, &record) 1148 if err != nil { 1149 l.Error("invalid record", "err", err) 1150 return err 1151 } 1152 1153 domain := e.Commit.RKey 1154 1155 err := db.AddKnot(i.Db, domain, did) 1156 if err != nil { 1157 l.Error("failed to add knot to db", "err", err, "domain", domain) 1158 return err 1159 } 1160 1161 if err := i.verifyKnot(ctx, domain, did); err != nil { 1162 l.Warn("failed to verify knot", "domain", domain, "did", did, "err", err) 1163 } 1164 1165 l.Info("ingested record", "domain", domain) 1166 return nil 1167 1168 case jmodels.CommitOperationDelete: 1169 domain := e.Commit.RKey 1170 1171 // get record from db first 1172 registrations, err := db.GetRegistrations( 1173 i.Db, 1174 orm.FilterEq("domain", domain), 1175 orm.FilterEq("did", did), 1176 ) 1177 if err != nil { 1178 return fmt.Errorf("failed to get registration: %w", err) 1179 } 1180 if len(registrations) != 1 { 1181 return fmt.Errorf("got incorrect number of registrations: %d, expected 1", len(registrations)) 1182 } 1183 registration := registrations[0] 1184 1185 tx, err := i.Db.Begin() 1186 if err != nil { 1187 return fmt.Errorf("failed to start txn: %w", err) 1188 } 1189 defer func() { 1190 tx.Rollback() 1191 i.Enforcer.E.LoadPolicy() 1192 }() 1193 1194 err = db.RemoveKnotMember( 1195 tx, 1196 orm.FilterEq("did", did), 1197 orm.FilterEq("domain", domain), 1198 ) 1199 if err != nil { 1200 return fmt.Errorf("failed to remove knot members: %w", err) 1201 } 1202 1203 err = db.DeleteKnot( 1204 tx, 1205 orm.FilterEq("did", did), 1206 orm.FilterEq("domain", domain), 1207 ) 1208 if err != nil { 1209 return fmt.Errorf("failed to delete knot: %w", err) 1210 } 1211 1212 err = db.RemoveReposByKnot(tx, domain) 1213 if err != nil { 1214 return fmt.Errorf("failed to remove repos by knot: %w", err) 1215 } 1216 1217 if registration.Registered != nil { 1218 err = i.Enforcer.RemoveKnot(domain) 1219 if err != nil { 1220 return fmt.Errorf("failed to remove knot from enforcer: %w", err) 1221 } 1222 } 1223 1224 err = tx.Commit() 1225 if err != nil { 1226 return fmt.Errorf("failed to commit txn: %w", err) 1227 } 1228 1229 err = i.Enforcer.E.SavePolicy() 1230 if err != nil { 1231 return fmt.Errorf("failed to save ACLs: %w", err) 1232 } 1233 1234 l.Info("ingested record", "domain", domain) 1235 } 1236 1237 return nil 1238} 1239 1240const ( 1241 verifyAttempts = 4 1242 verifyMinDelay = 1 * time.Second 1243 verifyMaxDelay = 5 * time.Second 1244) 1245 1246func (i *Ingester) verifyKnot(ctx context.Context, domain, did string) error { 1247 regs, err := db.GetRegistrations(i.Db, 1248 orm.FilterEq("domain", domain), 1249 orm.FilterEq("did", did), 1250 ) 1251 if err != nil { 1252 return fmt.Errorf("look up registration: %w", err) 1253 } 1254 if len(regs) != 1 { 1255 return fmt.Errorf("no registration for %s by %s", domain, did) 1256 } 1257 if regs[0].Registered != nil { 1258 return nil 1259 } 1260 1261 err = retry.Do( 1262 func() error { return serververify.RunVerification(ctx, domain, did, i.Config.Core.Dev) }, 1263 retry.Context(ctx), 1264 retry.Attempts(verifyAttempts), 1265 retry.Delay(verifyMinDelay), 1266 retry.MaxDelay(verifyMaxDelay), 1267 retry.DelayType(retry.BackOffDelay), 1268 retry.LastErrorOnly(true), 1269 ) 1270 if err != nil { 1271 return fmt.Errorf("verify: %w", err) 1272 } 1273 return serververify.MarkKnotVerified(i.Db, i.Enforcer, domain, did) 1274} 1275 1276func (i *Ingester) verifySpindle(ctx context.Context, instance, did string) error { 1277 spindles, err := db.GetSpindles(ctx, i.Db, 1278 orm.FilterEq("instance", instance), 1279 orm.FilterEq("owner", did), 1280 ) 1281 if err != nil { 1282 return fmt.Errorf("look up spindle: %w", err) 1283 } 1284 if len(spindles) != 1 { 1285 return fmt.Errorf("no spindle for %s by %s", instance, did) 1286 } 1287 if spindles[0].Verified != nil { 1288 return nil 1289 } 1290 1291 err = retry.Do( 1292 func() error { return serververify.RunVerification(ctx, instance, did, i.Config.Core.Dev) }, 1293 retry.Context(ctx), 1294 retry.Attempts(verifyAttempts), 1295 retry.Delay(verifyMinDelay), 1296 retry.MaxDelay(verifyMaxDelay), 1297 retry.DelayType(retry.BackOffDelay), 1298 retry.LastErrorOnly(true), 1299 ) 1300 if err != nil { 1301 return fmt.Errorf("verify: %w", err) 1302 } 1303 _, err = serververify.MarkSpindleVerified(i.Db, i.Enforcer, instance, did) 1304 return err 1305} 1306 1307const sweepConcurrency = 4 1308 1309func (i *Ingester) SweepPendingVerifications() { 1310 l := i.Logger.With("handler", "SweepPendingVerifications") 1311 1312 var g errgroup.Group 1313 g.SetLimit(sweepConcurrency) 1314 1315 regs, err := db.GetRegistrations(i.Db, orm.FilterIs("registered", nil)) 1316 if err != nil { 1317 l.Error("failed to list unverified knots", "err", err) 1318 } else { 1319 for _, reg := range regs { 1320 g.Go(func() error { 1321 if err := i.verifyKnot(i.Ctx, reg.Domain, reg.ByDid); err != nil { 1322 l.Warn("verify knot failed", "domain", reg.Domain, "did", reg.ByDid, "err", err) 1323 } 1324 return nil 1325 }) 1326 } 1327 } 1328 1329 spindles, err := db.GetSpindles(i.Ctx, i.Db, orm.FilterIs("verified", nil)) 1330 if err != nil { 1331 l.Error("failed to list unverified spindles", "err", err) 1332 g.Wait() 1333 return 1334 } 1335 for _, s := range spindles { 1336 g.Go(func() error { 1337 if err := i.verifySpindle(i.Ctx, s.Instance, s.Owner.String()); err != nil { 1338 l.Warn("verify spindle failed", "instance", s.Instance, "owner", s.Owner, "err", err) 1339 } 1340 return nil 1341 }) 1342 } 1343 g.Wait() 1344} 1345 1346func (i *Ingester) ingestIssue(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { 1347 did := e.Did 1348 rkey := e.Commit.RKey 1349 1350 var err error 1351 1352 l = l.With("handler", "ingestIssue") 1353 1354 switch e.Commit.Operation { 1355 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 1356 raw := json.RawMessage(e.Commit.Record) 1357 record := tangled.RepoIssue{} 1358 err = json.Unmarshal(raw, &record) 1359 if err != nil { 1360 l.Error("invalid record", "err", err) 1361 return err 1362 } 1363 1364 issue := models.IssueFromRecord(did, rkey, record) 1365 1366 if issue.RepoDid == "" { 1367 return fmt.Errorf("issue record has no repo field") 1368 } 1369 if _, err := syntax.ParseDID(string(issue.RepoDid)); err != nil { 1370 return fmt.Errorf("issue record repo field is not a valid DID: %w", err) 1371 } 1372 1373 if err := issue.Validate(); err != nil { 1374 return fmt.Errorf("failed to validate issue: %w", err) 1375 } 1376 1377 if record.Repo != "" && !strings.HasPrefix(record.Repo, "did:") { 1378 repo, repoErr := db.GetRepoByAtUri(i.Db, record.Repo) 1379 if repoErr == nil && repo.RepoDid != "" { 1380 if enqErr := db.EnqueuePdsRecordMigration(ctx, i.Db, "add-repo-did", syntax.DID(did), syntax.NSID(tangled.RepoIssueNSID), syntax.RecordKey(e.Commit.RKey)); enqErr != nil { 1381 l.Warn("failed to enqueue PDS rewrite for issue", "err", enqErr, "did", did, "repoDid", repo.RepoDid) 1382 } 1383 } 1384 } 1385 1386 tx, err := i.Db.BeginTx(ctx, nil) 1387 if err != nil { 1388 l.Error("failed to begin transaction", "err", err) 1389 return err 1390 } 1391 defer tx.Rollback() 1392 1393 err = db.PutIssue(tx, &issue) 1394 if err != nil { 1395 l.Error("failed to create issue", "err", err) 1396 return err 1397 } 1398 1399 err = tx.Commit() 1400 if err != nil { 1401 l.Error("failed to commit txn", "err", err) 1402 return err 1403 } 1404 1405 l.Info("ingested record") 1406 return nil 1407 1408 case jmodels.CommitOperationDelete: 1409 tx, err := i.Db.BeginTx(ctx, nil) 1410 if err != nil { 1411 l.Error("failed to begin transaction", "err", err) 1412 return err 1413 } 1414 defer tx.Rollback() 1415 1416 if err := db.DeleteIssues( 1417 tx, 1418 did, 1419 rkey, 1420 ); err != nil { 1421 l.Error("failed to delete", "err", err) 1422 return fmt.Errorf("failed to delete issue record: %w", err) 1423 } 1424 if err := tx.Commit(); err != nil { 1425 l.Error("failed to commit txn", "err", err) 1426 return err 1427 } 1428 1429 l.Info("ingested record") 1430 return nil 1431 } 1432 1433 return nil 1434} 1435 1436func (i *Ingester) ingestPull(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { 1437 did := e.Did 1438 rkey := e.Commit.RKey 1439 1440 var err error 1441 1442 l = l.With("handler", "ingestPull") 1443 1444 switch e.Commit.Operation { 1445 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 1446 raw := json.RawMessage(e.Commit.Record) 1447 record := tangled.RepoPull{} 1448 err = json.Unmarshal(raw, &record) 1449 if err != nil { 1450 l.Error("invalid record", "err", err) 1451 return err 1452 } 1453 1454 ownerId, err := i.IdResolver.ResolveIdent(ctx, did) 1455 if err != nil { 1456 l.Error("failed to resolve did", "err", err) 1457 return err 1458 } 1459 1460 // go through and fetch all blobs in parallel 1461 readers := make([]*io.ReadCloser, len(record.Rounds)) 1462 var mu sync.Mutex 1463 1464 g, gctx := errgroup.WithContext(ctx) 1465 1466 for idx, b := range record.Rounds { 1467 g.Go(func() error { 1468 // for some reason, a blob is empty 1469 if b.PatchBlob == nil { 1470 return fmt.Errorf("missing patchBlob in round %d", idx) 1471 } 1472 1473 ownerPds := ownerId.PDSEndpoint() 1474 url, _ := url.Parse(fmt.Sprintf("%s/xrpc/com.atproto.sync.getBlob", ownerPds)) 1475 q := url.Query() 1476 q.Set("cid", b.PatchBlob.Ref.String()) 1477 q.Set("did", did) 1478 url.RawQuery = q.Encode() 1479 1480 req, err := http.NewRequestWithContext(gctx, http.MethodGet, url.String(), nil) 1481 if err != nil { 1482 l.Error("failed to create request") 1483 return err 1484 } 1485 req.Header.Set("Content-Type", "application/json") 1486 1487 resp, err := http.DefaultClient.Do(req) 1488 if err != nil { 1489 l.Error("failed to make request") 1490 return err 1491 } 1492 1493 mu.Lock() 1494 readers[idx] = &resp.Body 1495 mu.Unlock() 1496 1497 return nil 1498 }) 1499 } 1500 1501 if err := g.Wait(); err != nil { 1502 for _, r := range readers { 1503 if r != nil && *r != nil { 1504 (*r).Close() 1505 } 1506 } 1507 return err 1508 } 1509 1510 defer func() { 1511 for _, r := range readers { 1512 if r != nil && *r != nil { 1513 (*r).Close() 1514 } 1515 } 1516 }() 1517 1518 pull, err := models.PullFromRecord(did, rkey, record, readers) 1519 if err != nil { 1520 return fmt.Errorf("failed to parse pull from record: %w", err) 1521 } 1522 if err := pull.Validate(); err != nil { 1523 return fmt.Errorf("failed to validate pull: %w", err) 1524 } 1525 if pull.DependentOn != nil { 1526 if err := func() error { 1527 dependentPull, err := db.GetPull( 1528 i.Db, 1529 orm.FilterEq("dependent_on", pull.DependentOn.String()), 1530 ) 1531 if errors.Is(err, sql.ErrNoRows) { 1532 return nil 1533 } 1534 if err != nil { 1535 return fmt.Errorf("failed to fetch pulls with same dependency: %w", err) 1536 } 1537 if dependentPull.AtUri() == pull.AtUri() { 1538 return nil 1539 } 1540 return fmt.Errorf("another pull already depends on %s, which would form a DAG, this is presently disallowed", pull.DependentOn.String()) 1541 }(); err != nil { 1542 return fmt.Errorf("failed to validate pull stack: %w", err) 1543 } 1544 } 1545 1546 tx, err := i.Db.BeginTx(ctx, nil) 1547 if err != nil { 1548 l.Error("failed to begin transaction", "err", err) 1549 return err 1550 } 1551 defer tx.Rollback() 1552 1553 err = db.PutPull(tx, pull) 1554 if err != nil { 1555 l.Error("failed to create pull", "err", err) 1556 return err 1557 } 1558 1559 err = tx.Commit() 1560 if err != nil { 1561 l.Error("failed to commit txn", "err", err) 1562 return err 1563 } 1564 1565 l.Info("ingested record") 1566 return nil 1567 1568 case jmodels.CommitOperationDelete: 1569 tx, err := i.Db.BeginTx(ctx, nil) 1570 if err != nil { 1571 l.Error("failed to begin transaction", "err", err) 1572 return err 1573 } 1574 defer tx.Rollback() 1575 1576 if err := db.AbandonPulls( 1577 tx, 1578 orm.FilterEq("owner_did", did), 1579 orm.FilterEq("rkey", rkey), 1580 ); err != nil { 1581 l.Error("failed to abandon", "err", err) 1582 return fmt.Errorf("failed to abandon pull record: %w", err) 1583 } 1584 if err := tx.Commit(); err != nil { 1585 l.Error("failed to commit txn", "err", err) 1586 return err 1587 } 1588 1589 l.Info("ingested record") 1590 return nil 1591 } 1592 1593 return nil 1594} 1595 1596// ingestIssueComment ingests legacy sh.tangled.repo.issue.comment deletions 1597func (i *Ingester) ingestIssueComment(e *jmodels.Event, l *slog.Logger) error { 1598 l = l.With("handler", "ingestIssueComment") 1599 1600 switch e.Commit.Operation { 1601 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 1602 // no-op. sh.tangled.repo.issue.comment is deprecated 1603 1604 case jmodels.CommitOperationDelete: 1605 if err := db.PurgeComments( 1606 i.Db, 1607 orm.FilterEq("did", e.Did), 1608 orm.FilterEq("collection", e.Commit.Collection), 1609 orm.FilterEq("rkey", e.Commit.RKey), 1610 ); err != nil { 1611 return fmt.Errorf("failed to delete comment record: %w", err) 1612 } 1613 } 1614 1615 l.Info("ingested record") 1616 return nil 1617} 1618 1619// ingestPullComment ingests legacy sh.tangled.repo.pull.comment deletions 1620func (i *Ingester) ingestPullComment(e *jmodels.Event, l *slog.Logger) error { 1621 l = l.With("handler", "ingestPullComment") 1622 1623 switch e.Commit.Operation { 1624 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 1625 // no-op. sh.tangled.repo.pull.comment is deprecated 1626 1627 case jmodels.CommitOperationDelete: 1628 if err := db.PurgeComments( 1629 i.Db, 1630 orm.FilterEq("did", e.Did), 1631 orm.FilterEq("collection", e.Commit.Collection), 1632 orm.FilterEq("rkey", e.Commit.RKey), 1633 ); err != nil { 1634 return fmt.Errorf("failed to delete comment record: %w", err) 1635 } 1636 } 1637 1638 l.Info("ingested record") 1639 return nil 1640} 1641 1642func (i *Ingester) ingestComment(e *jmodels.Event, l *slog.Logger) error { 1643 did := e.Did 1644 rkey := e.Commit.RKey 1645 cid := e.Commit.CID 1646 1647 var err error 1648 1649 l = l.With("handler", "ingestComment") 1650 1651 ctx := context.Background() 1652 1653 switch e.Commit.Operation { 1654 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 1655 raw := json.RawMessage(e.Commit.Record) 1656 record := tangled.FeedComment{} 1657 err = json.Unmarshal(raw, &record) 1658 if err != nil { 1659 return fmt.Errorf("invalid record: %w", err) 1660 } 1661 1662 comment, err := models.CommentFromRecord(syntax.DID(did), syntax.RecordKey(rkey), syntax.CID(cid), record) 1663 if err != nil { 1664 return fmt.Errorf("failed to parse comment from record: %w", err) 1665 } 1666 1667 if err := comment.Validate(); err != nil { 1668 return fmt.Errorf("failed to validate comment: %w", err) 1669 } 1670 1671 var references []syntax.ATURI 1672 if comment.Body.Original != nil { 1673 _, references = i.MentionsResolver.Resolve(ctx, *comment.Body.Original) 1674 } 1675 1676 tx, err := i.Db.Begin() 1677 if err != nil { 1678 return fmt.Errorf("failed to start transaction: %w", err) 1679 } 1680 defer tx.Rollback() 1681 1682 _, err = db.PutComment(tx, comment, references) 1683 if err != nil { 1684 return fmt.Errorf("failed to create comment: %w", err) 1685 } 1686 1687 if err := tx.Commit(); err != nil { 1688 return err 1689 } 1690 1691 case jmodels.CommitOperationDelete: 1692 if err := db.DeleteComments( 1693 i.Db, 1694 orm.FilterEq("did", did), 1695 orm.FilterEq("collection", e.Commit.Collection), 1696 orm.FilterEq("rkey", rkey), 1697 ); err != nil { 1698 return fmt.Errorf("failed to delete comment record: %w", err) 1699 } 1700 } 1701 1702 l.Info("ingested record") 1703 return nil 1704} 1705 1706func (i *Ingester) ingestReaction(e *jmodels.Event, l *slog.Logger) error { 1707 did := e.Did 1708 rkey := e.Commit.RKey 1709 1710 l = l.With("handler", "ingestReaction") 1711 1712 switch e.Commit.Operation { 1713 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 1714 raw := json.RawMessage(e.Commit.Record) 1715 record := tangled.FeedReaction{} 1716 if err := json.Unmarshal(raw, &record); err != nil { 1717 return fmt.Errorf("invalid record: %w", err) 1718 } 1719 1720 subjectUri, err := syntax.ParseATURI(record.Subject) 1721 if err != nil { 1722 return fmt.Errorf("invalid reaction subject %q: %w", record.Subject, err) 1723 } 1724 subjectUri = models.NormalizeReactionSubject(subjectUri) 1725 1726 kind, ok := models.ParseReactionKind(record.Reaction) 1727 if !ok { 1728 return fmt.Errorf("invalid reaction kind: %q", record.Reaction) 1729 } 1730 1731 created, parseErr := time.Parse(time.RFC3339, record.CreatedAt) 1732 if parseErr != nil { 1733 created = time.Now() 1734 } 1735 1736 reaction := models.Reaction{ 1737 ReactedByDid: did, 1738 Rkey: rkey, 1739 ThreadAt: subjectUri, 1740 Kind: kind, 1741 Created: created, 1742 } 1743 if err := db.UpsertReaction(i.Db, reaction); err != nil { 1744 return fmt.Errorf("failed to upsert reaction: %w", err) 1745 } 1746 1747 case jmodels.CommitOperationDelete: 1748 if err := db.DeleteReactionByRkey(i.Db, did, rkey); err != nil { 1749 return fmt.Errorf("failed to delete reaction record: %w", err) 1750 } 1751 } 1752 1753 l.Info("ingested record") 1754 return nil 1755} 1756 1757func (i *Ingester) ingestLabelDefinition(e *jmodels.Event, l *slog.Logger) error { 1758 did := e.Did 1759 rkey := e.Commit.RKey 1760 1761 var err error 1762 1763 l = l.With("handler", "ingestLabelDefinition") 1764 1765 switch e.Commit.Operation { 1766 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 1767 raw := json.RawMessage(e.Commit.Record) 1768 record := tangled.LabelDefinition{} 1769 err = json.Unmarshal(raw, &record) 1770 if err != nil { 1771 return fmt.Errorf("invalid record: %w", err) 1772 } 1773 1774 def, err := models.LabelDefinitionFromRecord(did, rkey, record) 1775 if err != nil { 1776 return fmt.Errorf("failed to parse labeldef from record: %w", err) 1777 } 1778 1779 if err := def.Validate(); err != nil { 1780 return fmt.Errorf("failed to validate labeldef: %w", err) 1781 } 1782 1783 _, err = db.AddLabelDefinition(i.Db, def) 1784 if err != nil { 1785 return fmt.Errorf("failed to create labeldef: %w", err) 1786 } 1787 1788 l.Info("ingested record") 1789 return nil 1790 1791 case jmodels.CommitOperationDelete: 1792 if err := db.DeleteLabelDefinition( 1793 i.Db, 1794 orm.FilterEq("did", did), 1795 orm.FilterEq("rkey", rkey), 1796 ); err != nil { 1797 return fmt.Errorf("failed to delete labeldef record: %w", err) 1798 } 1799 1800 l.Info("ingested record") 1801 return nil 1802 } 1803 1804 return nil 1805} 1806 1807func (i *Ingester) ingestLabelOp(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { 1808 did := e.Did 1809 rkey := e.Commit.RKey 1810 1811 var err error 1812 1813 l = l.With("handler", "ingestLabelOp") 1814 1815 switch e.Commit.Operation { 1816 case jmodels.CommitOperationCreate: 1817 raw := json.RawMessage(e.Commit.Record) 1818 record := tangled.LabelOp{} 1819 err = json.Unmarshal(raw, &record) 1820 if err != nil { 1821 return fmt.Errorf("invalid record: %w", err) 1822 } 1823 1824 subject := syntax.ATURI(record.Subject) 1825 collection := subject.Collection() 1826 1827 var repo *models.Repo 1828 switch collection { 1829 case tangled.RepoIssueNSID: 1830 i, err := db.GetIssues(i.Db, orm.FilterEq("at_uri", subject)) 1831 if err != nil || len(i) != 1 { 1832 return fmt.Errorf("failed to find subject: %w || subject count %d", err, len(i)) 1833 } 1834 repo = i[0].Repo 1835 case tangled.RepoPullNSID: 1836 p, err := db.GetPulls(i.Db, orm.FilterEq("at_uri", subject)) 1837 if err != nil || len(p) != 1 { 1838 return fmt.Errorf("failed to find subject: %w || subject count %d", err, len(p)) 1839 } 1840 repo = p[0].Repo 1841 default: 1842 return fmt.Errorf("unsupported label subject: %s", collection) 1843 } 1844 1845 actx, err := db.NewLabelApplicationCtx(i.Db, orm.FilterIn("at_uri", repo.Labels)) 1846 if err != nil { 1847 return fmt.Errorf("failed to build label application ctx: %w", err) 1848 } 1849 1850 ops := models.LabelOpsFromRecord(did, rkey, record) 1851 1852 for _, o := range ops { 1853 def, ok := actx.Defs[o.OperandKey] 1854 if !ok { 1855 return fmt.Errorf("failed to find label def for key: %s, expected: %q", o.OperandKey, slices.Collect(maps.Keys(actx.Defs))) 1856 } 1857 // validate permissions: only collaborators can apply labels currently 1858 // 1859 // TODO: introduce a repo:triage permission 1860 allowed, permErr := i.Acl.HasRepoPermissionErr(ctx, repo, o.Did, "repo:push") 1861 if permErr != nil { 1862 if !errors.Is(permErr, knotacl.ErrKnotUnreachable) { 1863 return fmt.Errorf("enforcing permission: %w", permErr) 1864 } 1865 l.Warn("ingesting labelop without permission check", "did", o.Did, "err", permErr) 1866 } else if !allowed { 1867 return fmt.Errorf("unauthorized label operation") 1868 } 1869 1870 if err := def.ValidateOperandValue(&o); err != nil { 1871 return fmt.Errorf("failed to validate labelop: %w", err) 1872 } 1873 } 1874 1875 tx, err := i.Db.Begin() 1876 if err != nil { 1877 return err 1878 } 1879 defer tx.Rollback() 1880 1881 for _, o := range ops { 1882 _, err = db.AddLabelOp(tx, &o) 1883 if err != nil { 1884 return fmt.Errorf("failed to add labelop: %w", err) 1885 } 1886 } 1887 1888 if err = tx.Commit(); err != nil { 1889 return err 1890 } 1891 1892 l.Info("ingested record") 1893 } 1894 1895 return nil 1896}