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 2196 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 tx, err := i.Db.Begin() 602 if err != nil { 603 return fmt.Errorf("failed to start transaction: %w", err) 604 } 605 defer tx.Rollback() 606 607 err = db.ValidateProfile(tx, &profile) 608 if err != nil { 609 return fmt.Errorf("invalid profile record") 610 } 611 612 err = db.UpsertProfile(tx, &profile) 613 if err != nil { 614 return fmt.Errorf("upserting profile: %w", err) 615 } 616 617 err = tx.Commit() 618 if err != nil { 619 return fmt.Errorf("tx.Commit: %w", err) 620 } 621 if i.Cache != nil { 622 pipe := i.Cache.Pipeline() 623 didKey := fmt.Sprintf(cache.PreferredHandleByDid, did) 624 if preferredHandle != "" { 625 pipe.Set(ctx, didKey, string(preferredHandle), cache.PreferredHandleTTL) 626 pipe.Set(ctx, fmt.Sprintf(cache.PreferredHandleByHandle, string(preferredHandle)), did, cache.PreferredHandleTTL) 627 } else { 628 pipe.Del(ctx, didKey) 629 } 630 if _, execErr := pipe.Exec(ctx); execErr != nil { 631 l.Warn("failed to update preferred handle cache", "err", execErr) 632 } 633 } 634 case jmodels.CommitOperationDelete: 635 tx, beginErr := i.Db.Begin() 636 if beginErr != nil { 637 return fmt.Errorf("failed to start transaction: %w", beginErr) 638 } 639 640 priorHandle, phErr := db.GetPreferredHandle(tx, did) 641 if phErr != nil && !errors.Is(phErr, sql.ErrNoRows) { 642 l.Warn("failed to read prior preferred handle", "err", phErr) 643 } 644 645 err = db.DeleteProfile(tx, did) 646 if err == nil && i.Cache != nil { 647 pipe := i.Cache.Pipeline() 648 pipe.Del(ctx, fmt.Sprintf(cache.PreferredHandleByDid, did)) 649 if priorHandle != "" { 650 pipe.Del(ctx, fmt.Sprintf(cache.PreferredHandleByHandle, string(priorHandle))) 651 } 652 if _, execErr := pipe.Exec(ctx); execErr != nil { 653 l.Warn("failed to evict preferred handle cache", "err", execErr) 654 } 655 } 656 } 657 658 if err != nil { 659 return fmt.Errorf("failed to %s profile record: %w", e.Commit.Operation, err) 660 } 661 662 l.Info("ingested record") 663 return nil 664} 665 666func (i *Ingester) ingestSpindleMember(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { 667 did := e.Did 668 var err error 669 670 l = l.With("handler", "ingestSpindleMember") 671 672 switch e.Commit.Operation { 673 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 674 raw := json.RawMessage(e.Commit.Record) 675 record := tangled.SpindleMember{} 676 err = json.Unmarshal(raw, &record) 677 if err != nil { 678 l.Error("invalid record", "err", err) 679 return err 680 } 681 682 // only spindle owner can invite to spindles 683 ok, err := i.Enforcer.IsSpindleInviteAllowed(did, record.Instance) 684 if err != nil { 685 return fmt.Errorf("failed to check invite permission: %w", err) 686 } 687 if !ok { 688 if verifyErr := i.verifySpindle(ctx, record.Instance, did); verifyErr != nil { 689 return fmt.Errorf("invite denied and verify failed: %w", verifyErr) 690 } 691 ok, err = i.Enforcer.IsSpindleInviteAllowed(did, record.Instance) 692 if err != nil { 693 return fmt.Errorf("failed to re-check invite permission: %w", err) 694 } 695 if !ok { 696 return fmt.Errorf("invite denied for did %s on spindle %s", did, record.Instance) 697 } 698 } 699 700 memberId, err := i.IdResolver.ResolveIdent(ctx, record.Subject) 701 if err != nil { 702 return err 703 } 704 705 if memberId.Handle.IsInvalidHandle() { 706 return fmt.Errorf("invalid handle for member %s", record.Subject) 707 } 708 709 existing, err := db.GetSpindleMembers(i.Db, 710 orm.FilterEq("did", did), 711 orm.FilterEq("rkey", e.Commit.RKey), 712 ) 713 if err != nil { 714 return fmt.Errorf("failed to look up existing member: %w", err) 715 } 716 if len(existing) > 1 { 717 return fmt.Errorf("multiple spindle members with rkey %s", e.Commit.RKey) 718 } 719 720 tx, err := i.Db.Begin() 721 if err != nil { 722 return fmt.Errorf("failed to start txn: %w", err) 723 } 724 committed := false 725 defer func() { 726 if committed { 727 return 728 } 729 tx.Rollback() 730 i.Enforcer.E.LoadPolicy() 731 }() 732 733 if len(existing) == 1 { 734 prev := existing[0] 735 if prev.Instance != record.Instance || prev.Subject != memberId.DID { 736 if err = db.RemoveSpindleMember(tx, 737 orm.FilterEq("did", did), 738 orm.FilterEq("rkey", e.Commit.RKey), 739 ); err != nil { 740 return fmt.Errorf("failed to remove stale row: %w", err) 741 } 742 if err = i.Enforcer.RemoveSpindleMember(prev.Instance, prev.Subject.String()); err != nil { 743 return fmt.Errorf("failed to remove stale ACL: %w", err) 744 } 745 } 746 } 747 748 if err = db.AddSpindleMember(tx, models.SpindleMember{ 749 Did: syntax.DID(did), 750 Rkey: e.Commit.RKey, 751 Instance: record.Instance, 752 Subject: memberId.DID, 753 }); err != nil { 754 return fmt.Errorf("failed to add to db: %w", err) 755 } 756 757 if err = i.Enforcer.AddSpindleMember(record.Instance, memberId.DID.String()); err != nil { 758 return fmt.Errorf("failed to update ACLs: %w", err) 759 } 760 761 if err = tx.Commit(); err != nil { 762 return fmt.Errorf("failed to commit txn: %w", err) 763 } 764 765 if err = i.Enforcer.E.SavePolicy(); err != nil { 766 return fmt.Errorf("failed to save ACLs: %w", err) 767 } 768 committed = true 769 770 l.Info("upserted spindle member") 771 case jmodels.CommitOperationDelete: 772 rkey := e.Commit.RKey 773 774 // get record from db first 775 members, err := db.GetSpindleMembers( 776 i.Db, 777 orm.FilterEq("did", did), 778 orm.FilterEq("rkey", rkey), 779 ) 780 if err != nil || len(members) != 1 { 781 return fmt.Errorf("failed to get member: %w, len(members) = %d", err, len(members)) 782 } 783 member := members[0] 784 785 tx, err := i.Db.Begin() 786 if err != nil { 787 return fmt.Errorf("failed to start txn: %w", err) 788 } 789 committed := false 790 defer func() { 791 if committed { 792 return 793 } 794 tx.Rollback() 795 i.Enforcer.E.LoadPolicy() 796 }() 797 798 // remove record by rkey && update enforcer 799 if err = db.RemoveSpindleMember( 800 tx, 801 orm.FilterEq("did", did), 802 orm.FilterEq("rkey", rkey), 803 ); err != nil { 804 return fmt.Errorf("failed to remove from db: %w", err) 805 } 806 807 // update enforcer 808 err = i.Enforcer.RemoveSpindleMember(member.Instance, member.Subject.String()) 809 if err != nil { 810 return fmt.Errorf("failed to update ACLs: %w", err) 811 } 812 813 if err = tx.Commit(); err != nil { 814 return fmt.Errorf("failed to commit txn: %w", err) 815 } 816 817 if err = i.Enforcer.E.SavePolicy(); err != nil { 818 return fmt.Errorf("failed to save ACLs: %w", err) 819 } 820 committed = true 821 822 l.Info("removed spindle member") 823 } 824 825 return nil 826} 827 828func (i *Ingester) ingestSpindle(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { 829 did := e.Did 830 var err error 831 832 l = l.With("handler", "ingestSpindle") 833 834 switch e.Commit.Operation { 835 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 836 raw := json.RawMessage(e.Commit.Record) 837 record := tangled.Spindle{} 838 err = json.Unmarshal(raw, &record) 839 if err != nil { 840 l.Error("invalid record", "err", err) 841 return err 842 } 843 844 instance := e.Commit.RKey 845 846 err := db.AddSpindle(i.Db, models.Spindle{ 847 Owner: syntax.DID(did), 848 Instance: instance, 849 }) 850 if err != nil { 851 l.Error("failed to add spindle to db", "err", err, "instance", instance) 852 return err 853 } 854 855 if err := i.verifySpindle(ctx, instance, did); err != nil { 856 l.Warn("failed to verify spindle", "instance", instance, "did", did, "err", err) 857 } 858 859 l.Info("ingested record", "instance", instance) 860 return nil 861 862 case jmodels.CommitOperationDelete: 863 instance := e.Commit.RKey 864 865 // get record from db first 866 spindles, err := db.GetSpindles( 867 ctx, 868 i.Db, 869 orm.FilterEq("owner", did), 870 orm.FilterEq("instance", instance), 871 ) 872 if err != nil || len(spindles) != 1 { 873 return fmt.Errorf("failed to get spindles: %w, len(spindles) = %d", err, len(spindles)) 874 } 875 spindle := spindles[0] 876 877 tx, err := i.Db.Begin() 878 if err != nil { 879 return fmt.Errorf("failed to start txn: %w", err) 880 } 881 defer func() { 882 tx.Rollback() 883 i.Enforcer.E.LoadPolicy() 884 }() 885 886 // remove spindle members first 887 err = db.RemoveSpindleMember( 888 tx, 889 orm.FilterEq("owner", did), 890 orm.FilterEq("instance", instance), 891 ) 892 if err != nil { 893 return fmt.Errorf("failed to remove spindle members: %w", err) 894 } 895 896 err = db.DeleteSpindle( 897 tx, 898 orm.FilterEq("owner", did), 899 orm.FilterEq("instance", instance), 900 ) 901 if err != nil { 902 return fmt.Errorf("failed to delete spindle: %w", err) 903 } 904 905 if spindle.Verified != nil { 906 err = i.Enforcer.RemoveSpindle(instance) 907 if err != nil { 908 return fmt.Errorf("failed to remove spindle from enforcer: %w", err) 909 } 910 } 911 912 err = tx.Commit() 913 if err != nil { 914 return fmt.Errorf("failed to commit txn: %w", err) 915 } 916 917 err = i.Enforcer.E.SavePolicy() 918 if err != nil { 919 return fmt.Errorf("failed to save ACLs: %w", err) 920 } 921 922 l.Info("ingested record", "instance", instance) 923 } 924 925 return nil 926} 927 928func (i *Ingester) ingestString(e *jmodels.Event, l *slog.Logger) error { 929 did := e.Did 930 rkey := e.Commit.RKey 931 932 var err error 933 934 l = l.With("handler", "ingestString") 935 936 switch e.Commit.Operation { 937 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 938 raw := json.RawMessage(e.Commit.Record) 939 record := tangled.String{} 940 err = json.Unmarshal(raw, &record) 941 if err != nil { 942 l.Error("invalid record", "err", err) 943 return err 944 } 945 946 string := models.StringFromRecord(did, rkey, record) 947 948 if err = string.Validate(); err != nil { 949 l.Error("invalid record", "err", err) 950 return err 951 } 952 953 if err = db.AddString(i.Db, string); err != nil { 954 l.Error("failed to add string", "err", err) 955 return err 956 } 957 958 l.Info("ingested record") 959 return nil 960 961 case jmodels.CommitOperationDelete: 962 if err := db.DeleteString( 963 i.Db, 964 orm.FilterEq("did", did), 965 orm.FilterEq("rkey", rkey), 966 ); err != nil { 967 l.Error("failed to delete", "err", err) 968 return fmt.Errorf("failed to delete string record: %w", err) 969 } 970 971 l.Info("ingested record") 972 return nil 973 } 974 975 return nil 976} 977 978func (i *Ingester) ingestKnotMember(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { 979 did := e.Did 980 var err error 981 982 l = l.With("handler", "ingestKnotMember") 983 984 switch e.Commit.Operation { 985 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 986 raw := json.RawMessage(e.Commit.Record) 987 record := tangled.KnotMember{} 988 err = json.Unmarshal(raw, &record) 989 if err != nil { 990 l.Error("invalid record", "err", err) 991 return err 992 } 993 994 // only knot owner can invite to knots 995 ok, err := i.Enforcer.IsKnotInviteAllowed(did, record.Domain) 996 if err != nil { 997 return fmt.Errorf("failed to check invite permission: %w", err) 998 } 999 if !ok { 1000 if verifyErr := i.verifyKnot(ctx, record.Domain, did); verifyErr != nil { 1001 return fmt.Errorf("invite denied and verify failed: %w", verifyErr) 1002 } 1003 ok, err = i.Enforcer.IsKnotInviteAllowed(did, record.Domain) 1004 if err != nil { 1005 return fmt.Errorf("failed to re-check invite permission: %w", err) 1006 } 1007 if !ok { 1008 return fmt.Errorf("invite denied for did %s on knot %s", did, record.Domain) 1009 } 1010 } 1011 1012 memberId, err := i.IdResolver.ResolveIdent(ctx, record.Subject) 1013 if err != nil { 1014 return err 1015 } 1016 1017 if memberId.Handle.IsInvalidHandle() { 1018 return fmt.Errorf("invalid handle for member %s", record.Subject) 1019 } 1020 1021 existing, err := db.GetKnotMembers(i.Db, 1022 orm.FilterEq("did", did), 1023 orm.FilterEq("rkey", e.Commit.RKey), 1024 ) 1025 if err != nil { 1026 return fmt.Errorf("failed to look up existing member: %w", err) 1027 } 1028 if len(existing) > 1 { 1029 return fmt.Errorf("multiple knot members with rkey %s", e.Commit.RKey) 1030 } 1031 1032 tx, err := i.Db.Begin() 1033 if err != nil { 1034 return fmt.Errorf("failed to start txn: %w", err) 1035 } 1036 committed := false 1037 defer func() { 1038 if committed { 1039 return 1040 } 1041 tx.Rollback() 1042 i.Enforcer.E.LoadPolicy() 1043 }() 1044 1045 if len(existing) == 1 { 1046 prev := existing[0] 1047 if prev.Domain != record.Domain || prev.Subject != memberId.DID { 1048 if err = db.RemoveKnotMember(tx, 1049 orm.FilterEq("did", did), 1050 orm.FilterEq("rkey", e.Commit.RKey), 1051 ); err != nil { 1052 return fmt.Errorf("failed to remove stale row: %w", err) 1053 } 1054 if err = i.Enforcer.RemoveKnotMember(prev.Domain, prev.Subject.String()); err != nil { 1055 return fmt.Errorf("failed to remove stale ACL: %w", err) 1056 } 1057 } 1058 } 1059 1060 if err = db.AddKnotMember(tx, models.KnotMember{ 1061 Did: syntax.DID(did), 1062 Rkey: e.Commit.RKey, 1063 Domain: record.Domain, 1064 Subject: memberId.DID, 1065 }); err != nil { 1066 return fmt.Errorf("failed to add to db: %w", err) 1067 } 1068 1069 if err = i.Enforcer.AddKnotMember(record.Domain, memberId.DID.String()); err != nil { 1070 return fmt.Errorf("failed to update ACLs: %w", err) 1071 } 1072 1073 if err = tx.Commit(); err != nil { 1074 return fmt.Errorf("failed to commit txn: %w", err) 1075 } 1076 1077 if err = i.Enforcer.E.SavePolicy(); err != nil { 1078 return fmt.Errorf("failed to save ACLs: %w", err) 1079 } 1080 committed = true 1081 1082 l.Info("upserted knot member") 1083 case jmodels.CommitOperationDelete: 1084 rkey := e.Commit.RKey 1085 1086 members, err := db.GetKnotMembers( 1087 i.Db, 1088 orm.FilterEq("did", did), 1089 orm.FilterEq("rkey", rkey), 1090 ) 1091 if err != nil { 1092 return fmt.Errorf("failed to look up knot member with rkey %s: %w", rkey, err) 1093 } 1094 if len(members) == 0 { 1095 l.Info("knot member already removed", "rkey", rkey) 1096 return nil 1097 } 1098 if len(members) > 1 { 1099 return fmt.Errorf("multiple knot members with rkey %s", rkey) 1100 } 1101 member := members[0] 1102 1103 tx, err := i.Db.Begin() 1104 if err != nil { 1105 return fmt.Errorf("failed to start txn: %w", err) 1106 } 1107 committed := false 1108 defer func() { 1109 if committed { 1110 return 1111 } 1112 tx.Rollback() 1113 i.Enforcer.E.LoadPolicy() 1114 }() 1115 1116 if err = db.RemoveKnotMember( 1117 tx, 1118 orm.FilterEq("did", did), 1119 orm.FilterEq("rkey", rkey), 1120 ); err != nil { 1121 return fmt.Errorf("failed to remove from db: %w", err) 1122 } 1123 1124 if err = i.Enforcer.RemoveKnotMember(member.Domain, member.Subject.String()); err != nil { 1125 return fmt.Errorf("failed to update ACLs: %w", err) 1126 } 1127 1128 if err = tx.Commit(); err != nil { 1129 return fmt.Errorf("failed to commit txn: %w", err) 1130 } 1131 1132 if err = i.Enforcer.E.SavePolicy(); err != nil { 1133 return fmt.Errorf("failed to save ACLs: %w", err) 1134 } 1135 committed = true 1136 1137 l.Info("removed knot member") 1138 } 1139 1140 return nil 1141} 1142 1143func (i *Ingester) ingestKnot(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { 1144 did := e.Did 1145 var err error 1146 1147 l = l.With("handler", "ingestKnot") 1148 1149 switch e.Commit.Operation { 1150 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 1151 raw := json.RawMessage(e.Commit.Record) 1152 record := tangled.Knot{} 1153 err = json.Unmarshal(raw, &record) 1154 if err != nil { 1155 l.Error("invalid record", "err", err) 1156 return err 1157 } 1158 1159 domain := e.Commit.RKey 1160 1161 err := db.AddKnot(i.Db, domain, did) 1162 if err != nil { 1163 l.Error("failed to add knot to db", "err", err, "domain", domain) 1164 return err 1165 } 1166 1167 if err := i.verifyKnot(ctx, domain, did); err != nil { 1168 l.Warn("failed to verify knot", "domain", domain, "did", did, "err", err) 1169 } 1170 1171 l.Info("ingested record", "domain", domain) 1172 return nil 1173 1174 case jmodels.CommitOperationDelete: 1175 domain := e.Commit.RKey 1176 1177 // get record from db first 1178 registrations, err := db.GetRegistrations( 1179 i.Db, 1180 orm.FilterEq("domain", domain), 1181 orm.FilterEq("did", did), 1182 ) 1183 if err != nil { 1184 return fmt.Errorf("failed to get registration: %w", err) 1185 } 1186 if len(registrations) != 1 { 1187 return fmt.Errorf("got incorrect number of registrations: %d, expected 1", len(registrations)) 1188 } 1189 registration := registrations[0] 1190 1191 tx, err := i.Db.Begin() 1192 if err != nil { 1193 return fmt.Errorf("failed to start txn: %w", err) 1194 } 1195 defer func() { 1196 tx.Rollback() 1197 i.Enforcer.E.LoadPolicy() 1198 }() 1199 1200 err = db.RemoveKnotMember( 1201 tx, 1202 orm.FilterEq("did", did), 1203 orm.FilterEq("domain", domain), 1204 ) 1205 if err != nil { 1206 return fmt.Errorf("failed to remove knot members: %w", err) 1207 } 1208 1209 err = db.DeleteKnot( 1210 tx, 1211 orm.FilterEq("did", did), 1212 orm.FilterEq("domain", domain), 1213 ) 1214 if err != nil { 1215 return fmt.Errorf("failed to delete knot: %w", err) 1216 } 1217 1218 err = db.RemoveReposByKnot(tx, domain) 1219 if err != nil { 1220 return fmt.Errorf("failed to remove repos by knot: %w", err) 1221 } 1222 1223 if registration.Registered != nil { 1224 err = i.Enforcer.RemoveKnot(domain) 1225 if err != nil { 1226 return fmt.Errorf("failed to remove knot from enforcer: %w", err) 1227 } 1228 } 1229 1230 err = tx.Commit() 1231 if err != nil { 1232 return fmt.Errorf("failed to commit txn: %w", err) 1233 } 1234 1235 err = i.Enforcer.E.SavePolicy() 1236 if err != nil { 1237 return fmt.Errorf("failed to save ACLs: %w", err) 1238 } 1239 1240 l.Info("ingested record", "domain", domain) 1241 } 1242 1243 return nil 1244} 1245 1246const ( 1247 verifyAttempts = 4 1248 verifyMinDelay = 1 * time.Second 1249 verifyMaxDelay = 5 * time.Second 1250) 1251 1252func (i *Ingester) verifyKnot(ctx context.Context, domain, did string) error { 1253 regs, err := db.GetRegistrations(i.Db, 1254 orm.FilterEq("domain", domain), 1255 orm.FilterEq("did", did), 1256 ) 1257 if err != nil { 1258 return fmt.Errorf("look up registration: %w", err) 1259 } 1260 if len(regs) != 1 { 1261 return fmt.Errorf("no registration for %s by %s", domain, did) 1262 } 1263 if regs[0].Registered != nil { 1264 return nil 1265 } 1266 1267 err = retry.Do( 1268 func() error { return serververify.RunVerification(ctx, domain, did, i.Config.Core.Dev) }, 1269 retry.Context(ctx), 1270 retry.Attempts(verifyAttempts), 1271 retry.Delay(verifyMinDelay), 1272 retry.MaxDelay(verifyMaxDelay), 1273 retry.DelayType(retry.BackOffDelay), 1274 retry.LastErrorOnly(true), 1275 ) 1276 if err != nil { 1277 return fmt.Errorf("verify: %w", err) 1278 } 1279 return serververify.MarkKnotVerified(i.Db, i.Enforcer, domain, did) 1280} 1281 1282func (i *Ingester) verifySpindle(ctx context.Context, instance, did string) error { 1283 spindles, err := db.GetSpindles(ctx, i.Db, 1284 orm.FilterEq("instance", instance), 1285 orm.FilterEq("owner", did), 1286 ) 1287 if err != nil { 1288 return fmt.Errorf("look up spindle: %w", err) 1289 } 1290 if len(spindles) != 1 { 1291 return fmt.Errorf("no spindle for %s by %s", instance, did) 1292 } 1293 if spindles[0].Verified != nil { 1294 return nil 1295 } 1296 1297 err = retry.Do( 1298 func() error { return serververify.RunVerification(ctx, instance, did, i.Config.Core.Dev) }, 1299 retry.Context(ctx), 1300 retry.Attempts(verifyAttempts), 1301 retry.Delay(verifyMinDelay), 1302 retry.MaxDelay(verifyMaxDelay), 1303 retry.DelayType(retry.BackOffDelay), 1304 retry.LastErrorOnly(true), 1305 ) 1306 if err != nil { 1307 return fmt.Errorf("verify: %w", err) 1308 } 1309 _, err = serververify.MarkSpindleVerified(i.Db, i.Enforcer, instance, did) 1310 return err 1311} 1312 1313const sweepConcurrency = 4 1314 1315func (i *Ingester) SweepPendingVerifications() { 1316 l := i.Logger.With("handler", "SweepPendingVerifications") 1317 1318 var g errgroup.Group 1319 g.SetLimit(sweepConcurrency) 1320 1321 regs, err := db.GetRegistrations(i.Db, orm.FilterIs("registered", nil)) 1322 if err != nil { 1323 l.Error("failed to list unverified knots", "err", err) 1324 } else { 1325 for _, reg := range regs { 1326 g.Go(func() error { 1327 if err := i.verifyKnot(i.Ctx, reg.Domain, reg.ByDid); err != nil { 1328 l.Warn("verify knot failed", "domain", reg.Domain, "did", reg.ByDid, "err", err) 1329 } 1330 return nil 1331 }) 1332 } 1333 } 1334 1335 spindles, err := db.GetSpindles(i.Ctx, i.Db, orm.FilterIs("verified", nil)) 1336 if err != nil { 1337 l.Error("failed to list unverified spindles", "err", err) 1338 g.Wait() 1339 return 1340 } 1341 for _, s := range spindles { 1342 g.Go(func() error { 1343 if err := i.verifySpindle(i.Ctx, s.Instance, s.Owner.String()); err != nil { 1344 l.Warn("verify spindle failed", "instance", s.Instance, "owner", s.Owner, "err", err) 1345 } 1346 return nil 1347 }) 1348 } 1349 g.Wait() 1350} 1351 1352func (i *Ingester) ingestIssue(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { 1353 did := e.Did 1354 rkey := e.Commit.RKey 1355 1356 var err error 1357 1358 l = l.With("handler", "ingestIssue") 1359 1360 switch e.Commit.Operation { 1361 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 1362 raw := json.RawMessage(e.Commit.Record) 1363 record := tangled.RepoIssue{} 1364 err = json.Unmarshal(raw, &record) 1365 if err != nil { 1366 l.Error("invalid record", "err", err) 1367 return err 1368 } 1369 1370 issue := models.IssueFromRecord(did, rkey, record) 1371 1372 if issue.RepoDid == "" { 1373 return fmt.Errorf("issue record has no repo field") 1374 } 1375 if _, err := syntax.ParseDID(string(issue.RepoDid)); err != nil { 1376 return fmt.Errorf("issue record repo field is not a valid DID: %w", err) 1377 } 1378 1379 if err := issue.Validate(); err != nil { 1380 return fmt.Errorf("failed to validate issue: %w", err) 1381 } 1382 1383 if record.Repo != "" && !strings.HasPrefix(record.Repo, "did:") { 1384 repo, repoErr := db.GetRepoByAtUri(i.Db, record.Repo) 1385 if repoErr == nil && repo.RepoDid != "" { 1386 if enqErr := db.EnqueuePdsRecordMigration(ctx, i.Db, "add-repo-did", syntax.DID(did), syntax.NSID(tangled.RepoIssueNSID), syntax.RecordKey(e.Commit.RKey)); enqErr != nil { 1387 l.Warn("failed to enqueue PDS rewrite for issue", "err", enqErr, "did", did, "repoDid", repo.RepoDid) 1388 } 1389 } 1390 } 1391 1392 tx, err := i.Db.BeginTx(ctx, nil) 1393 if err != nil { 1394 l.Error("failed to begin transaction", "err", err) 1395 return err 1396 } 1397 defer tx.Rollback() 1398 1399 err = db.PutIssue(tx, &issue) 1400 if err != nil { 1401 l.Error("failed to create issue", "err", err) 1402 return err 1403 } 1404 1405 if err := db.ResolveIssueState(tx, issue.AtUri()); err != nil { 1406 l.Error("failed to resolve issue state", "err", err) 1407 return err 1408 } 1409 1410 err = tx.Commit() 1411 if err != nil { 1412 l.Error("failed to commit txn", "err", err) 1413 return err 1414 } 1415 1416 i.drainPendingState(ctx, issue.AtUri(), issueStateSpec, l) 1417 1418 l.Info("ingested record") 1419 return nil 1420 1421 case jmodels.CommitOperationDelete: 1422 tx, err := i.Db.BeginTx(ctx, nil) 1423 if err != nil { 1424 l.Error("failed to begin transaction", "err", err) 1425 return err 1426 } 1427 defer tx.Rollback() 1428 1429 if err := db.DeleteIssues( 1430 tx, 1431 did, 1432 rkey, 1433 ); err != nil { 1434 l.Error("failed to delete", "err", err) 1435 return fmt.Errorf("failed to delete issue record: %w", err) 1436 } 1437 if err := tx.Commit(); err != nil { 1438 l.Error("failed to commit txn", "err", err) 1439 return err 1440 } 1441 1442 l.Info("ingested record") 1443 return nil 1444 } 1445 1446 return nil 1447} 1448 1449func (i *Ingester) ingestPull(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { 1450 did := e.Did 1451 rkey := e.Commit.RKey 1452 1453 var err error 1454 1455 l = l.With("handler", "ingestPull") 1456 1457 switch e.Commit.Operation { 1458 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 1459 raw := json.RawMessage(e.Commit.Record) 1460 record := tangled.RepoPull{} 1461 err = json.Unmarshal(raw, &record) 1462 if err != nil { 1463 l.Error("invalid record", "err", err) 1464 return err 1465 } 1466 1467 ownerId, err := i.IdResolver.ResolveIdent(ctx, did) 1468 if err != nil { 1469 l.Error("failed to resolve did", "err", err) 1470 return err 1471 } 1472 1473 // go through and fetch all blobs in parallel 1474 readers := make([]*io.ReadCloser, len(record.Rounds)) 1475 var mu sync.Mutex 1476 1477 g, gctx := errgroup.WithContext(ctx) 1478 1479 for idx, b := range record.Rounds { 1480 g.Go(func() error { 1481 // for some reason, a blob is empty 1482 if b.PatchBlob == nil { 1483 return fmt.Errorf("missing patchBlob in round %d", idx) 1484 } 1485 1486 ownerPds := ownerId.PDSEndpoint() 1487 url, _ := url.Parse(fmt.Sprintf("%s/xrpc/com.atproto.sync.getBlob", ownerPds)) 1488 q := url.Query() 1489 q.Set("cid", b.PatchBlob.Ref.String()) 1490 q.Set("did", did) 1491 url.RawQuery = q.Encode() 1492 1493 req, err := http.NewRequestWithContext(gctx, http.MethodGet, url.String(), nil) 1494 if err != nil { 1495 l.Error("failed to create request") 1496 return err 1497 } 1498 req.Header.Set("Content-Type", "application/json") 1499 1500 resp, err := http.DefaultClient.Do(req) 1501 if err != nil { 1502 l.Error("failed to make request") 1503 return err 1504 } 1505 1506 mu.Lock() 1507 readers[idx] = &resp.Body 1508 mu.Unlock() 1509 1510 return nil 1511 }) 1512 } 1513 1514 if err := g.Wait(); err != nil { 1515 for _, r := range readers { 1516 if r != nil && *r != nil { 1517 (*r).Close() 1518 } 1519 } 1520 return err 1521 } 1522 1523 defer func() { 1524 for _, r := range readers { 1525 if r != nil && *r != nil { 1526 (*r).Close() 1527 } 1528 } 1529 }() 1530 1531 pull, err := models.PullFromRecord(did, rkey, record, readers) 1532 if err != nil { 1533 return fmt.Errorf("failed to parse pull from record: %w", err) 1534 } 1535 if err := pull.Validate(); err != nil { 1536 return fmt.Errorf("failed to validate pull: %w", err) 1537 } 1538 if pull.DependentOn != nil { 1539 if err := func() error { 1540 dependentPull, err := db.GetPull( 1541 i.Db, 1542 orm.FilterEq("dependent_on", pull.DependentOn.String()), 1543 ) 1544 if errors.Is(err, sql.ErrNoRows) { 1545 return nil 1546 } 1547 if err != nil { 1548 return fmt.Errorf("failed to fetch pulls with same dependency: %w", err) 1549 } 1550 if dependentPull.AtUri() == pull.AtUri() { 1551 return nil 1552 } 1553 return fmt.Errorf("another pull already depends on %s, which would form a DAG, this is presently disallowed", pull.DependentOn.String()) 1554 }(); err != nil { 1555 return fmt.Errorf("failed to validate pull stack: %w", err) 1556 } 1557 } 1558 1559 tx, err := i.Db.BeginTx(ctx, nil) 1560 if err != nil { 1561 l.Error("failed to begin transaction", "err", err) 1562 return err 1563 } 1564 defer tx.Rollback() 1565 1566 err = db.PutPull(tx, pull) 1567 if err != nil { 1568 l.Error("failed to create pull", "err", err) 1569 return err 1570 } 1571 1572 if err := db.ResolvePullStatus(tx, pull.AtUri()); err != nil { 1573 l.Error("failed to resolve pull status", "err", err) 1574 return err 1575 } 1576 1577 err = tx.Commit() 1578 if err != nil { 1579 l.Error("failed to commit txn", "err", err) 1580 return err 1581 } 1582 1583 i.drainPendingState(ctx, pull.AtUri(), pullStatusSpec, l) 1584 1585 l.Info("ingested record") 1586 return nil 1587 1588 case jmodels.CommitOperationDelete: 1589 tx, err := i.Db.BeginTx(ctx, nil) 1590 if err != nil { 1591 l.Error("failed to begin transaction", "err", err) 1592 return err 1593 } 1594 defer tx.Rollback() 1595 1596 if err := db.AbandonPulls( 1597 tx, 1598 orm.FilterEq("owner_did", did), 1599 orm.FilterEq("rkey", rkey), 1600 ); err != nil { 1601 l.Error("failed to abandon", "err", err) 1602 return fmt.Errorf("failed to abandon pull record: %w", err) 1603 } 1604 if err := tx.Commit(); err != nil { 1605 l.Error("failed to commit txn", "err", err) 1606 return err 1607 } 1608 1609 l.Info("ingested record") 1610 return nil 1611 } 1612 1613 return nil 1614} 1615 1616func (i *Ingester) authorizeStateRecord(ctx context.Context, repo *models.Repo, subjectAuthorDid, recordAuthorDid string, l *slog.Logger) (bool, error) { 1617 if recordAuthorDid == subjectAuthorDid { 1618 return true, nil 1619 } 1620 if recordAuthorDid == consts.TangledDid { 1621 return true, nil 1622 } 1623 1624 ok, err := i.Acl.HasRepoPermissionErr(ctx, repo, recordAuthorDid, "repo:push") 1625 if err != nil { 1626 if errors.Is(err, knotacl.ErrKnotUnreachable) { 1627 l.Warn("ingesting state record without permission check", "did", recordAuthorDid, "err", err) 1628 return true, nil 1629 } 1630 return false, err 1631 } 1632 return ok, nil 1633} 1634 1635type stateIngestSpec struct { 1636 subjectNSID string 1637 parse func(did, rkey string, raw json.RawMessage) (models.StateRecord, error) 1638 findSubject func(e db.Execer, subject syntax.ATURI) (repo *models.Repo, authorDid string, found bool, err error) 1639 put func(tx *sql.Tx, rec models.StateRecord) (syntax.ATURI, error) 1640 resolve func(tx *sql.Tx, subject syntax.ATURI) error 1641 recompute func(tx *sql.Tx, subject syntax.ATURI) error 1642 del func(tx *sql.Tx, did, rkey string) (syntax.ATURI, error) 1643} 1644 1645var issueStateSpec = stateIngestSpec{ 1646 subjectNSID: tangled.RepoIssueNSID, 1647 parse: func(did, rkey string, raw json.RawMessage) (models.StateRecord, error) { 1648 record := tangled.RepoIssueState{} 1649 if err := json.Unmarshal(raw, &record); err != nil { 1650 return models.StateRecord{}, fmt.Errorf("invalid record: %w", err) 1651 } 1652 return models.IssueStateFromRecord(did, rkey, record) 1653 }, 1654 findSubject: func(e db.Execer, subject syntax.ATURI) (*models.Repo, string, bool, error) { 1655 issues, err := db.GetIssues(e, orm.FilterEq("at_uri", subject)) 1656 if err != nil { 1657 return nil, "", false, err 1658 } 1659 if len(issues) != 1 || issues[0].Repo == nil { 1660 return nil, "", false, nil 1661 } 1662 return issues[0].Repo, issues[0].Did, true, nil 1663 }, 1664 put: db.PutIssueState, 1665 resolve: db.ResolveIssueState, 1666 recompute: db.RecomputeIssueState, 1667 del: db.DeleteIssueState, 1668} 1669 1670var pullStatusSpec = stateIngestSpec{ 1671 subjectNSID: tangled.RepoPullNSID, 1672 parse: func(did, rkey string, raw json.RawMessage) (models.StateRecord, error) { 1673 record := tangled.RepoPullStatus{} 1674 if err := json.Unmarshal(raw, &record); err != nil { 1675 return models.StateRecord{}, fmt.Errorf("invalid record: %w", err) 1676 } 1677 return models.PullStatusFromRecord(did, rkey, record) 1678 }, 1679 findSubject: func(e db.Execer, subject syntax.ATURI) (*models.Repo, string, bool, error) { 1680 pulls, err := db.GetPulls(e, orm.FilterEq("at_uri", subject)) 1681 if err != nil { 1682 return nil, "", false, err 1683 } 1684 if len(pulls) != 1 || pulls[0].Repo == nil { 1685 return nil, "", false, nil 1686 } 1687 return pulls[0].Repo, pulls[0].OwnerDid, true, nil 1688 }, 1689 put: db.PutPullStatus, 1690 resolve: db.ResolvePullStatus, 1691 recompute: db.RecomputePullStatus, 1692 del: db.DeletePullStatus, 1693} 1694 1695func (i *Ingester) ingestState(ctx context.Context, e *jmodels.Event, l *slog.Logger, spec stateIngestSpec) error { 1696 did := e.Did 1697 rkey := e.Commit.RKey 1698 nsid := e.Commit.Collection 1699 1700 l = l.With("handler", "ingestState", "nsid", nsid) 1701 1702 switch e.Commit.Operation { 1703 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 1704 return i.applyStateRecord(ctx, did, rkey, nsid, e.Commit.Record, spec, l) 1705 case jmodels.CommitOperationDelete: 1706 return i.deleteStateRecord(ctx, did, rkey, nsid, spec, l) 1707 } 1708 1709 return nil 1710} 1711 1712func (i *Ingester) applyStateRecord(ctx context.Context, did, rkey, nsid string, raw []byte, spec stateIngestSpec, l *slog.Logger) error { 1713 rec, err := spec.parse(did, rkey, json.RawMessage(raw)) 1714 if err != nil { 1715 return err 1716 } 1717 if string(rec.Subject.Collection()) != spec.subjectNSID { 1718 return fmt.Errorf("state subject is not %s: %s", spec.subjectNSID, rec.Subject) 1719 } 1720 1721 repo, authorDid, found, err := spec.findSubject(i.Db, rec.Subject) 1722 if err != nil { 1723 return fmt.Errorf("failed to look up state subject: %w", err) 1724 } 1725 if !found { 1726 return i.parkStateRecord(ctx, did, rkey, nsid, rec.Subject, raw, l) 1727 } 1728 1729 authorized, err := i.authorizeStateRecord(ctx, repo, authorDid, did, l) 1730 if err != nil { 1731 return i.parkStateRecord(ctx, did, rkey, nsid, rec.Subject, raw, l) 1732 } 1733 1734 tx, err := i.Db.BeginTx(ctx, nil) 1735 if err != nil { 1736 return err 1737 } 1738 defer tx.Rollback() 1739 1740 if err := db.UnparkStateRecord(tx, did, rkey, nsid); err != nil { 1741 return fmt.Errorf("failed to unpark state record: %w", err) 1742 } 1743 1744 if !authorized { 1745 if err := tx.Commit(); err != nil { 1746 return err 1747 } 1748 l.Warn("dropped unauthorized state record", "did", did, "rkey", rkey, "subject", rec.Subject) 1749 return nil 1750 } 1751 1752 priorSubject, err := spec.put(tx, rec) 1753 if err != nil { 1754 return fmt.Errorf("failed to put state record: %w", err) 1755 } 1756 if err := spec.resolve(tx, rec.Subject); err != nil { 1757 return fmt.Errorf("failed to resolve state: %w", err) 1758 } 1759 if priorSubject != "" { 1760 if err := spec.recompute(tx, priorSubject); err != nil { 1761 return fmt.Errorf("failed to recompute prior subject state: %w", err) 1762 } 1763 } 1764 1765 if err := tx.Commit(); err != nil { 1766 return err 1767 } 1768 1769 l.Info("ingested record") 1770 return nil 1771} 1772 1773func (i *Ingester) deleteStateRecord(ctx context.Context, did, rkey, nsid string, spec stateIngestSpec, l *slog.Logger) error { 1774 tx, err := i.Db.BeginTx(ctx, nil) 1775 if err != nil { 1776 return err 1777 } 1778 defer tx.Rollback() 1779 1780 if err := db.UnparkStateRecord(tx, did, rkey, nsid); err != nil { 1781 return fmt.Errorf("failed to unpark state record: %w", err) 1782 } 1783 subject, err := spec.del(tx, did, rkey) 1784 if err != nil { 1785 return fmt.Errorf("failed to delete state record: %w", err) 1786 } 1787 if subject != "" { 1788 if err := spec.recompute(tx, subject); err != nil { 1789 return fmt.Errorf("failed to recompute state: %w", err) 1790 } 1791 } 1792 1793 if err := tx.Commit(); err != nil { 1794 return err 1795 } 1796 1797 l.Info("ingested record") 1798 return nil 1799} 1800 1801func (i *Ingester) parkStateRecord(ctx context.Context, did, rkey, nsid string, subject syntax.ATURI, raw []byte, l *slog.Logger) error { 1802 tx, err := i.Db.BeginTx(ctx, nil) 1803 if err != nil { 1804 return err 1805 } 1806 defer tx.Rollback() 1807 1808 if err := db.ParkStateRecord(tx, db.PendingStateRecord{ 1809 Did: did, 1810 Rkey: rkey, 1811 Nsid: nsid, 1812 Subject: subject, 1813 Record: raw, 1814 }); err != nil { 1815 return fmt.Errorf("failed to park state record: %w", err) 1816 } 1817 1818 if err := tx.Commit(); err != nil { 1819 return err 1820 } 1821 1822 l.Info("parked state record for retry", "subject", subject) 1823 return nil 1824} 1825 1826func (i *Ingester) drainPendingState(ctx context.Context, subject syntax.ATURI, spec stateIngestSpec, l *slog.Logger) { 1827 pending, err := db.PendingStateRecordsForSubject(i.Db, subject) 1828 if err != nil { 1829 l.Error("failed to load pending state records", "err", err, "subject", subject) 1830 return 1831 } 1832 for _, p := range pending { 1833 if err := i.applyStateRecord(ctx, p.Did, p.Rkey, p.Nsid, p.Record, spec, l); err != nil { 1834 l.Error("failed to drain pending state record", "err", err, "did", p.Did, "rkey", p.Rkey) 1835 } 1836 } 1837} 1838 1839const ( 1840 pendingStateReconcileInterval = time.Hour 1841 pendingStateRecordTTL = 7 * 24 * time.Hour 1842) 1843 1844func stateSpecForSubject(subject syntax.ATURI) (stateIngestSpec, bool) { 1845 switch string(subject.Collection()) { 1846 case tangled.RepoIssueNSID: 1847 return issueStateSpec, true 1848 case tangled.RepoPullNSID: 1849 return pullStatusSpec, true 1850 default: 1851 return stateIngestSpec{}, false 1852 } 1853} 1854 1855func (i *Ingester) StartPendingStateReconciler() { 1856 i.ReconcilePendingState() 1857 1858 ticker := time.NewTicker(pendingStateReconcileInterval) 1859 defer ticker.Stop() 1860 for { 1861 select { 1862 case <-i.Ctx.Done(): 1863 return 1864 case <-ticker.C: 1865 i.ReconcilePendingState() 1866 } 1867 } 1868} 1869 1870func (i *Ingester) ReconcilePendingState() { 1871 l := i.Logger.With("handler", "reconcilePendingState") 1872 1873 subjects, err := db.DistinctPendingStateSubjects(i.Db) 1874 if err != nil { 1875 l.Error("failed to list pending state subjects", "err", err) 1876 } 1877 for _, subject := range subjects { 1878 spec, ok := stateSpecForSubject(subject) 1879 if !ok { 1880 continue 1881 } 1882 i.drainPendingState(i.Ctx, subject, spec, l) 1883 } 1884 1885 cutoff := time.Now().Add(-pendingStateRecordTTL).UTC().Format(time.RFC3339) 1886 evicted, err := db.EvictStalePendingStateRecords(i.Db, cutoff) 1887 if err != nil { 1888 l.Error("failed to evict stale pending state records", "err", err) 1889 return 1890 } 1891 if evicted > 0 { 1892 l.Warn("evicted stale pending state records", "count", evicted, "olderThan", cutoff) 1893 } 1894} 1895 1896// ingestIssueComment ingests legacy sh.tangled.repo.issue.comment deletions 1897func (i *Ingester) ingestIssueComment(e *jmodels.Event, l *slog.Logger) error { 1898 l = l.With("handler", "ingestIssueComment") 1899 1900 switch e.Commit.Operation { 1901 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 1902 // no-op. sh.tangled.repo.issue.comment is deprecated 1903 1904 case jmodels.CommitOperationDelete: 1905 if err := db.PurgeComments( 1906 i.Db, 1907 orm.FilterEq("did", e.Did), 1908 orm.FilterEq("collection", e.Commit.Collection), 1909 orm.FilterEq("rkey", e.Commit.RKey), 1910 ); err != nil { 1911 return fmt.Errorf("failed to delete comment record: %w", err) 1912 } 1913 } 1914 1915 l.Info("ingested record") 1916 return nil 1917} 1918 1919// ingestPullComment ingests legacy sh.tangled.repo.pull.comment deletions 1920func (i *Ingester) ingestPullComment(e *jmodels.Event, l *slog.Logger) error { 1921 l = l.With("handler", "ingestPullComment") 1922 1923 switch e.Commit.Operation { 1924 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 1925 // no-op. sh.tangled.repo.pull.comment is deprecated 1926 1927 case jmodels.CommitOperationDelete: 1928 if err := db.PurgeComments( 1929 i.Db, 1930 orm.FilterEq("did", e.Did), 1931 orm.FilterEq("collection", e.Commit.Collection), 1932 orm.FilterEq("rkey", e.Commit.RKey), 1933 ); err != nil { 1934 return fmt.Errorf("failed to delete comment record: %w", err) 1935 } 1936 } 1937 1938 l.Info("ingested record") 1939 return nil 1940} 1941 1942func (i *Ingester) ingestComment(e *jmodels.Event, l *slog.Logger) error { 1943 did := e.Did 1944 rkey := e.Commit.RKey 1945 cid := e.Commit.CID 1946 1947 var err error 1948 1949 l = l.With("handler", "ingestComment") 1950 1951 ctx := context.Background() 1952 1953 switch e.Commit.Operation { 1954 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 1955 raw := json.RawMessage(e.Commit.Record) 1956 record := tangled.FeedComment{} 1957 err = json.Unmarshal(raw, &record) 1958 if err != nil { 1959 return fmt.Errorf("invalid record: %w", err) 1960 } 1961 1962 comment, err := models.CommentFromRecord(syntax.DID(did), syntax.RecordKey(rkey), syntax.CID(cid), record) 1963 if err != nil { 1964 return fmt.Errorf("failed to parse comment from record: %w", err) 1965 } 1966 1967 if err := comment.Validate(); err != nil { 1968 return fmt.Errorf("failed to validate comment: %w", err) 1969 } 1970 1971 var references []syntax.ATURI 1972 if comment.Body.Original != nil { 1973 _, references = i.MentionsResolver.Resolve(ctx, *comment.Body.Original) 1974 } 1975 1976 tx, err := i.Db.Begin() 1977 if err != nil { 1978 return fmt.Errorf("failed to start transaction: %w", err) 1979 } 1980 defer tx.Rollback() 1981 1982 _, err = db.PutComment(tx, comment, references) 1983 if err != nil { 1984 return fmt.Errorf("failed to create comment: %w", err) 1985 } 1986 1987 if err := tx.Commit(); err != nil { 1988 return err 1989 } 1990 1991 case jmodels.CommitOperationDelete: 1992 if err := db.DeleteComments( 1993 i.Db, 1994 orm.FilterEq("did", did), 1995 orm.FilterEq("collection", e.Commit.Collection), 1996 orm.FilterEq("rkey", rkey), 1997 ); err != nil { 1998 return fmt.Errorf("failed to delete comment record: %w", err) 1999 } 2000 } 2001 2002 l.Info("ingested record") 2003 return nil 2004} 2005 2006func (i *Ingester) ingestReaction(e *jmodels.Event, l *slog.Logger) error { 2007 did := e.Did 2008 rkey := e.Commit.RKey 2009 2010 l = l.With("handler", "ingestReaction") 2011 2012 switch e.Commit.Operation { 2013 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 2014 raw := json.RawMessage(e.Commit.Record) 2015 record := tangled.FeedReaction{} 2016 if err := json.Unmarshal(raw, &record); err != nil { 2017 return fmt.Errorf("invalid record: %w", err) 2018 } 2019 2020 subjectUri, err := syntax.ParseATURI(record.Subject) 2021 if err != nil { 2022 return fmt.Errorf("invalid reaction subject %q: %w", record.Subject, err) 2023 } 2024 subjectUri = models.NormalizeReactionSubject(subjectUri) 2025 2026 kind, ok := models.ParseReactionKind(record.Reaction) 2027 if !ok { 2028 return fmt.Errorf("invalid reaction kind: %q", record.Reaction) 2029 } 2030 2031 created, parseErr := time.Parse(time.RFC3339, record.CreatedAt) 2032 if parseErr != nil { 2033 created = time.Now() 2034 } 2035 2036 reaction := models.Reaction{ 2037 ReactedByDid: did, 2038 Rkey: rkey, 2039 ThreadAt: subjectUri, 2040 Kind: kind, 2041 Created: created, 2042 } 2043 if err := db.UpsertReaction(i.Db, reaction); err != nil { 2044 return fmt.Errorf("failed to upsert reaction: %w", err) 2045 } 2046 2047 case jmodels.CommitOperationDelete: 2048 if err := db.DeleteReactionByRkey(i.Db, did, rkey); err != nil { 2049 return fmt.Errorf("failed to delete reaction record: %w", err) 2050 } 2051 } 2052 2053 l.Info("ingested record") 2054 return nil 2055} 2056 2057func (i *Ingester) ingestLabelDefinition(e *jmodels.Event, l *slog.Logger) error { 2058 did := e.Did 2059 rkey := e.Commit.RKey 2060 2061 var err error 2062 2063 l = l.With("handler", "ingestLabelDefinition") 2064 2065 switch e.Commit.Operation { 2066 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 2067 raw := json.RawMessage(e.Commit.Record) 2068 record := tangled.LabelDefinition{} 2069 err = json.Unmarshal(raw, &record) 2070 if err != nil { 2071 return fmt.Errorf("invalid record: %w", err) 2072 } 2073 2074 def, err := models.LabelDefinitionFromRecord(did, rkey, record) 2075 if err != nil { 2076 return fmt.Errorf("failed to parse labeldef from record: %w", err) 2077 } 2078 2079 if err := def.Validate(); err != nil { 2080 return fmt.Errorf("failed to validate labeldef: %w", err) 2081 } 2082 2083 _, err = db.AddLabelDefinition(i.Db, def) 2084 if err != nil { 2085 return fmt.Errorf("failed to create labeldef: %w", err) 2086 } 2087 2088 l.Info("ingested record") 2089 return nil 2090 2091 case jmodels.CommitOperationDelete: 2092 if err := db.DeleteLabelDefinition( 2093 i.Db, 2094 orm.FilterEq("did", did), 2095 orm.FilterEq("rkey", rkey), 2096 ); err != nil { 2097 return fmt.Errorf("failed to delete labeldef record: %w", err) 2098 } 2099 2100 l.Info("ingested record") 2101 return nil 2102 } 2103 2104 return nil 2105} 2106 2107func (i *Ingester) ingestLabelOp(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { 2108 did := e.Did 2109 rkey := e.Commit.RKey 2110 2111 var err error 2112 2113 l = l.With("handler", "ingestLabelOp") 2114 2115 switch e.Commit.Operation { 2116 case jmodels.CommitOperationCreate: 2117 raw := json.RawMessage(e.Commit.Record) 2118 record := tangled.LabelOp{} 2119 err = json.Unmarshal(raw, &record) 2120 if err != nil { 2121 return fmt.Errorf("invalid record: %w", err) 2122 } 2123 2124 subject := syntax.ATURI(record.Subject) 2125 collection := subject.Collection() 2126 2127 var repo *models.Repo 2128 switch collection { 2129 case tangled.RepoIssueNSID: 2130 i, err := db.GetIssues(i.Db, orm.FilterEq("at_uri", subject)) 2131 if err != nil || len(i) != 1 { 2132 return fmt.Errorf("failed to find subject: %w || subject count %d", err, len(i)) 2133 } 2134 repo = i[0].Repo 2135 case tangled.RepoPullNSID: 2136 p, err := db.GetPulls(i.Db, orm.FilterEq("at_uri", subject)) 2137 if err != nil || len(p) != 1 { 2138 return fmt.Errorf("failed to find subject: %w || subject count %d", err, len(p)) 2139 } 2140 repo = p[0].Repo 2141 default: 2142 return fmt.Errorf("unsupported label subject: %s", collection) 2143 } 2144 2145 actx, err := db.NewLabelApplicationCtx(i.Db, orm.FilterIn("at_uri", repo.Labels)) 2146 if err != nil { 2147 return fmt.Errorf("failed to build label application ctx: %w", err) 2148 } 2149 2150 ops := models.LabelOpsFromRecord(did, rkey, record) 2151 2152 for _, o := range ops { 2153 def, ok := actx.Defs[o.OperandKey] 2154 if !ok { 2155 return fmt.Errorf("failed to find label def for key: %s, expected: %q", o.OperandKey, slices.Collect(maps.Keys(actx.Defs))) 2156 } 2157 // validate permissions: only collaborators can apply labels currently 2158 // 2159 // TODO: introduce a repo:triage permission 2160 allowed, permErr := i.Acl.HasRepoPermissionErr(ctx, repo, o.Did, "repo:push") 2161 if permErr != nil { 2162 if !errors.Is(permErr, knotacl.ErrKnotUnreachable) { 2163 return fmt.Errorf("enforcing permission: %w", permErr) 2164 } 2165 l.Warn("ingesting labelop without permission check", "did", o.Did, "err", permErr) 2166 } else if !allowed { 2167 return fmt.Errorf("unauthorized label operation") 2168 } 2169 2170 if err := def.ValidateOperandValue(&o); err != nil { 2171 return fmt.Errorf("failed to validate labelop: %w", err) 2172 } 2173 } 2174 2175 tx, err := i.Db.Begin() 2176 if err != nil { 2177 return err 2178 } 2179 defer tx.Rollback() 2180 2181 for _, o := range ops { 2182 _, err = db.AddLabelOp(tx, &o) 2183 if err != nil { 2184 return fmt.Errorf("failed to add labelop: %w", err) 2185 } 2186 } 2187 2188 if err = tx.Commit(); err != nil { 2189 return err 2190 } 2191 2192 l.Info("ingested record") 2193 } 2194 2195 return nil 2196}