This repository has no description
0

Configure Feed

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

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