This repository has no description
0

Configure Feed

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

core / appview / ingester.go
61 kB 2259 lines
1package appview 2 3import ( 4 "context" 5 "database/sql" 6 "encoding/json" 7 "errors" 8 "fmt" 9 "log/slog" 10 "net/http" 11 "net/url" 12 "slices" 13 "strings" 14 "sync" 15 16 "time" 17 18 "github.com/avast/retry-go/v4" 19 "github.com/bluesky-social/indigo/atproto/syntax" 20 jmodels "github.com/bluesky-social/jetstream/pkg/models" 21 "github.com/go-git/go-git/v5/plumbing" 22 "github.com/ipfs/go-cid" 23 "golang.org/x/sync/errgroup" 24 "tangled.org/core/api/tangled" 25 "tangled.org/core/appview/cache" 26 "tangled.org/core/appview/config" 27 "tangled.org/core/appview/db" 28 "tangled.org/core/appview/knotacl" 29 "tangled.org/core/appview/mentions" 30 "tangled.org/core/appview/models" 31 "tangled.org/core/appview/notify" 32 "tangled.org/core/appview/serververify" 33 "tangled.org/core/consts" 34 "tangled.org/core/idresolver" 35 "tangled.org/core/orm" 36 "tangled.org/core/rbac" 37 "tangled.org/core/repoverify" 38) 39 40type RepoPermissionChecker interface { 41 HasRepoPermissionErr(ctx context.Context, repo *models.Repo, userDid, perm string) (bool, error) 42} 43 44type Ingester struct { 45 Ctx context.Context 46 Db *db.DB 47 Enforcer *rbac.Enforcer 48 Acl RepoPermissionChecker 49 IdResolver *idresolver.Resolver 50 Cache *cache.Cache 51 Config *config.Config 52 Logger *slog.Logger 53 MentionsResolver *mentions.Resolver 54 Notifier notify.Notifier 55 Verifier repoverify.Verifier 56} 57 58type processFunc func(ctx context.Context, e *jmodels.Event) error 59 60func (i *Ingester) Ingest() processFunc { 61 return func(ctx context.Context, e *jmodels.Event) error { 62 var err error 63 64 l := i.Logger.With("kind", e.Kind) 65 switch e.Kind { 66 case jmodels.EventKindAccount: 67 // TODO: sync account state to db 68 if e.Account.Active { 69 break 70 } 71 // TODO: revoke sessions by DID 72 if *e.Account.Status == "deactivated" { 73 err = i.IdResolver.InvalidateIdent(ctx, e.Account.Did) 74 } 75 case jmodels.EventKindIdentity: 76 err = i.IdResolver.InvalidateIdent(ctx, e.Identity.Did) 77 case jmodels.EventKindCommit: 78 l = l.With( 79 "nsid", e.Commit.Collection, 80 "did", e.Did, 81 "rkey", e.Commit.RKey, 82 "op", e.Commit.Operation, 83 ) 84 switch e.Commit.Collection { 85 case tangled.GraphFollowNSID: 86 err = i.ingestFollow(e, l) 87 case tangled.GraphVouchNSID: 88 err = i.ingestVouch(ctx, e, l) 89 case tangled.FeedStarNSID: 90 err = i.ingestStar(ctx, e, l) 91 case tangled.FeedReactionNSID: 92 err = i.ingestReaction(e, l) 93 case tangled.PublicKeyNSID: 94 err = i.ingestPublicKey(e, l) 95 case tangled.RepoArtifactNSID: 96 err = i.ingestArtifact(ctx, e, l) 97 case tangled.ActorProfileNSID: 98 err = i.ingestProfile(ctx, e, l) 99 case tangled.SpindleMemberNSID: 100 err = i.ingestSpindleMember(ctx, e, l) 101 case tangled.SpindleNSID: 102 err = i.ingestSpindle(ctx, e, l) 103 case tangled.KnotMemberNSID: 104 err = i.ingestKnotMember(ctx, e, l) 105 case tangled.KnotNSID: 106 err = i.ingestKnot(ctx, e, l) 107 case tangled.StringNSID: 108 err = i.ingestString(e, l) 109 case tangled.RepoIssueNSID: 110 err = i.ingestIssue(ctx, e, l) 111 case tangled.RepoIssueStateNSID: 112 err = i.ingestState(ctx, e, l, issueStateSpec) 113 case tangled.RepoPullNSID: 114 err = i.ingestPull(ctx, e, l) 115 case tangled.RepoPullStatusNSID: 116 err = i.ingestState(ctx, e, l, pullStatusSpec) 117 case tangled.FeedCommentNSID: 118 err = i.ingestComment(e, l) 119 case tangled.RepoIssueCommentNSID: 120 err = i.ingestIssueComment(e, l) 121 case tangled.RepoPullCommentNSID: 122 err = i.ingestPullComment(e, l) 123 case tangled.LabelDefinitionNSID: 124 err = i.ingestLabelDefinition(e, l) 125 case tangled.LabelOpNSID: 126 err = i.ingestLabelOp(ctx, e, l) 127 case tangled.RepoNSID: 128 err = i.ingestRepo(ctx, e, l) 129 } 130 } 131 132 if err != nil { 133 l.Warn("failed to ingest record, skipping", "err", err) 134 } 135 136 return nil 137 } 138} 139 140func (i *Ingester) resolveRepoRef(ref string) (*models.Repo, error) { 141 if strings.HasPrefix(ref, "did:") { 142 return db.GetRepoByDid(i.Db, ref) 143 } 144 return db.GetRepoByAtUri(i.Db, ref) 145} 146 147func (i *Ingester) resolveOldFormatStar(raw json.RawMessage, star *models.Star, l *slog.Logger) (bool, error) { 148 var legacy struct { 149 Subject *string `json:"subject"` 150 SubjectDid *string `json:"subjectDid"` 151 } 152 if err := json.Unmarshal(raw, &legacy); err != nil { 153 return false, err 154 } 155 156 switch { 157 case legacy.SubjectDid != nil: 158 repo, err := i.resolveRepoRef(*legacy.SubjectDid) 159 if err != nil { 160 l.Warn("skipping old-format star for unknown repo", "subjectDid", *legacy.SubjectDid) 161 return false, nil 162 } 163 star.SubjectType = models.StarSubjectRepo 164 star.Subject = repo.RepoDid 165 return true, nil 166 167 case legacy.Subject != nil: 168 uri, err := syntax.ParseATURI(*legacy.Subject) 169 if err != nil { 170 return false, fmt.Errorf("invalid old-format star subject: %w", err) 171 } 172 switch uri.Collection().String() { 173 case tangled.RepoNSID: 174 repo, err := db.GetRepoByAtUri(i.Db, uri.String()) 175 if err != nil { 176 l.Warn("skipping old-format star for unknown repo", "subject", *legacy.Subject) 177 return false, nil 178 } 179 star.SubjectType = models.StarSubjectRepo 180 star.Subject = repo.RepoDid 181 return true, nil 182 default: 183 star.SubjectType = models.StarSubjectString 184 star.Subject = *legacy.Subject 185 return true, nil 186 } 187 188 default: 189 return false, fmt.Errorf("old-format star has neither subject nor subjectDid") 190 } 191} 192 193func (i *Ingester) ingestStar(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { 194 var err error 195 did := e.Did 196 197 l = l.With("handler", "ingestStar") 198 199 switch e.Commit.Operation { 200 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 201 raw := json.RawMessage(e.Commit.Record) 202 record := tangled.FeedStar{} 203 unmarshalErr := json.Unmarshal(raw, &record) 204 205 createdAt, parseErr := time.Parse(time.RFC3339, record.CreatedAt) 206 if parseErr != nil { 207 createdAt = time.Now() 208 } 209 210 star := models.Star{ 211 Did: did, 212 Created: createdAt, 213 } 214 215 switch { 216 case unmarshalErr != nil: 217 resolved, resolveErr := i.resolveOldFormatStar(raw, &star, l) 218 if resolveErr != nil { 219 l.Error("invalid record", "newFmtErr", unmarshalErr, "oldFmtErr", resolveErr) 220 return unmarshalErr 221 } 222 if !resolved { 223 return nil 224 } 225 226 case record.Subject == nil: 227 return fmt.Errorf("star record has nil subject") 228 229 case record.Subject.FeedStar_Repo != nil: 230 repo, repoErr := i.resolveRepoRef(record.Subject.FeedStar_Repo.Did) 231 if repoErr != nil { 232 l.Warn("skipping star for unknown repo", "did", record.Subject.FeedStar_Repo.Did) 233 return nil 234 } 235 star.SubjectType = models.StarSubjectRepo 236 star.Subject = repo.RepoDid 237 238 case record.Subject.FeedStar_String != nil: 239 star.SubjectType = models.StarSubjectString 240 star.Subject = record.Subject.FeedStar_String.Uri 241 242 default: 243 return fmt.Errorf("star record has empty subject union") 244 } 245 246 err = db.UpsertStar(i.Db, e.Commit.RKey, star) 247 case jmodels.CommitOperationDelete: 248 err = db.DeleteStarByRkey(i.Db, did, e.Commit.RKey) 249 } 250 251 if err != nil { 252 return fmt.Errorf("failed to %s star record: %w", e.Commit.Operation, err) 253 } 254 l.Info("processed star", "operation", e.Commit.Operation, "rkey", e.Commit.RKey) 255 256 l.Info("ingested record") 257 return nil 258} 259 260func (i *Ingester) ingestFollow(e *jmodels.Event, l *slog.Logger) error { 261 var err error 262 did := e.Did 263 264 l = l.With("handler", "ingestFollow") 265 266 switch e.Commit.Operation { 267 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 268 raw := json.RawMessage(e.Commit.Record) 269 record := tangled.GraphFollow{} 270 err = json.Unmarshal(raw, &record) 271 if err != nil { 272 l.Error("invalid record", "err", err) 273 return err 274 } 275 _, err := syntax.ParseDID(record.Subject) 276 if err != nil { 277 l.Error("invalid record. subject is invalid DID", "err", err) 278 return err 279 } 280 281 followedAt, err := time.Parse(time.RFC3339, record.CreatedAt) 282 if err != nil { 283 err = fmt.Errorf("createdAt is invalid datetime: %w", err) 284 l.Error("invalid record", "err", err) 285 return err 286 } 287 288 err = db.UpsertFollow(i.Db, e.Commit.RKey, models.Follow{ 289 UserDid: did, 290 SubjectDid: record.Subject, 291 FollowedAt: followedAt, 292 }) 293 case jmodels.CommitOperationDelete: 294 err = db.DeleteFollowByRkey(i.Db, did, e.Commit.RKey) 295 } 296 297 if err != nil { 298 return fmt.Errorf("failed to %s follow record: %w", e.Commit.Operation, err) 299 } 300 l.Info("processed follow", "operation", e.Commit.Operation, "rkey", e.Commit.RKey) 301 302 l.Info("ingested record") 303 return nil 304} 305 306func (i *Ingester) ingestVouch(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { 307 var err error 308 did := e.Did 309 310 l = l.With("handler", "ingestVouch") 311 312 switch e.Commit.Operation { 313 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 314 raw := json.RawMessage(e.Commit.Record) 315 record := tangled.GraphVouch{} 316 err = json.Unmarshal(raw, &record) 317 if err != nil { 318 l.Error("invalid record", "err", err) 319 return err 320 } 321 322 // rkey is the subject_did being vouched for/denounced 323 subjectDID := e.Commit.RKey 324 325 _, err = syntax.ParseDID(subjectDID) 326 if err != nil { 327 l.Error("invalid subject_did in rkey", "err", err, "rkey", subjectDID) 328 return fmt.Errorf("invalid subject_did: %w", err) 329 } 330 331 if did == subjectDID { 332 l.Warn("attempted self-vouch", "did", did) 333 return fmt.Errorf("cannot vouch for self") 334 } 335 336 subjectId, err := i.IdResolver.ResolveIdent(ctx, subjectDID) 337 if err != nil { 338 return err 339 } 340 341 if subjectId.Handle.IsInvalidHandle() { 342 return err 343 } 344 345 kind, err := models.ParseVouchKind(record.Kind) 346 if err != nil { 347 l.Error("invalid kind", "kind", kind) 348 return fmt.Errorf("invalid kind: %s", kind) 349 } 350 351 recordCid, err := cid.Parse(e.Commit.CID) 352 if err != nil { 353 l.Error("invalid cid", "err", err, "cid", e.Commit.CID) 354 return fmt.Errorf("invalid cid: %w", err) 355 } 356 357 var evidences []syntax.ATURI 358 for _, raw := range record.Evidences { 359 uri, parseErr := syntax.ParseATURI(raw) 360 if parseErr != nil { 361 l.Warn("invalid evidence AT-URI, skipping", "uri", raw, "err", parseErr) 362 continue 363 } 364 evidences = append(evidences, uri) 365 } 366 367 tx, txErr := i.Db.Begin() 368 if txErr != nil { 369 return fmt.Errorf("failed to start transaction: %w", txErr) 370 } 371 372 addErr := db.AddVouch(tx, &models.Vouch{ 373 Did: syntax.DID(did), 374 SubjectDid: subjectId.DID, 375 Cid: recordCid, 376 Kind: kind, 377 Reason: record.Reason, 378 Evidences: evidences, 379 }) 380 if addErr != nil { 381 tx.Rollback() 382 err = addErr 383 } else { 384 err = tx.Commit() 385 } 386 387 case jmodels.CommitOperationDelete: 388 err = db.DeleteVouchByRkey(i.Db, did, e.Commit.RKey) 389 } 390 391 if err != nil { 392 return fmt.Errorf("failed to %s vouch record: %w", e.Commit.Operation, err) 393 } 394 395 l.Info("ingested record") 396 return nil 397} 398 399func (i *Ingester) ingestPublicKey(e *jmodels.Event, l *slog.Logger) error { 400 did := e.Did 401 var err error 402 403 l = l.With("handler", "ingestPublicKey") 404 405 switch e.Commit.Operation { 406 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 407 l.Debug("processing add of pubkey") 408 raw := json.RawMessage(e.Commit.Record) 409 record := tangled.PublicKey{} 410 err = json.Unmarshal(raw, &record) 411 if err != nil { 412 l.Error("invalid record", "err", err) 413 return err 414 } 415 pubKey, err := models.PublicKeyFromRecord(syntax.DID(did), syntax.RecordKey(e.Commit.RKey), record) 416 if err != nil { 417 l.Error("invalid record", "err", err) 418 return err 419 } 420 if err := pubKey.Validate(); err != nil { 421 l.Error("invalid record", "err", err) 422 return err 423 } 424 425 err = db.UpsertPublicKey(i.Db, pubKey) 426 case jmodels.CommitOperationDelete: 427 l.Debug("processing delete of pubkey") 428 err = db.DeletePublicKeyByRkey(i.Db, did, e.Commit.RKey) 429 } 430 431 if err != nil { 432 return fmt.Errorf("failed to %s pubkey record: %w", e.Commit.Operation, err) 433 } 434 l.Info("processed pubkey", "operation", e.Commit.Operation, "rkey", e.Commit.RKey) 435 436 l.Info("ingested record") 437 return nil 438} 439 440func (i *Ingester) ingestArtifact(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { 441 did := e.Did 442 var err error 443 444 l = l.With("handler", "ingestArtifact") 445 446 switch e.Commit.Operation { 447 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 448 raw := json.RawMessage(e.Commit.Record) 449 record := tangled.RepoArtifact{} 450 err = json.Unmarshal(raw, &record) 451 if err != nil { 452 l.Error("invalid record", "err", err) 453 return err 454 } 455 456 var repo *models.Repo 457 if record.RepoDid != nil && *record.RepoDid != "" { 458 repo, err = db.GetRepoByDid(i.Db, *record.RepoDid) 459 if err != nil && !errors.Is(err, sql.ErrNoRows) { 460 return fmt.Errorf("failed to look up repo by DID %s: %w", *record.RepoDid, err) 461 } 462 } 463 if repo == nil && record.Repo != nil { 464 repoAt, parseErr := syntax.ParseATURI(*record.Repo) 465 if parseErr != nil { 466 return parseErr 467 } 468 repo, err = db.GetRepoByAtUri(i.Db, repoAt.String()) 469 if err != nil { 470 return err 471 } 472 } 473 if repo == nil { 474 return fmt.Errorf("artifact record has neither valid repoDid nor repo field") 475 } 476 477 allowed, permErr := i.Acl.HasRepoPermissionErr(ctx, repo, did, "repo:push") 478 if permErr != nil { 479 l.Warn("ingesting artifact without permission check", "did", did, "repo", repo.RepoIdentifier(), "err", permErr) 480 } else if !allowed { 481 l.Info("skipping unauthorized artifact", "did", did, "repo", repo.RepoIdentifier()) 482 return nil 483 } 484 485 repoDid := repo.RepoDid 486 if repoDid == "" && record.RepoDid != nil { 487 repoDid = *record.RepoDid 488 } 489 if repoDid != "" && (record.RepoDid == nil || *record.RepoDid == "") && record.Repo != nil { 490 if enqErr := db.EnqueuePdsRecordMigration(ctx, i.Db, "add-repo-did", syntax.DID(did), syntax.NSID(tangled.RepoArtifactNSID), syntax.RecordKey(e.Commit.RKey)); enqErr != nil { 491 l.Warn("failed to enqueue PDS rewrite for artifact", "err", enqErr, "did", did, "repoDid", repoDid) 492 } 493 } 494 495 createdAt, parseErr := time.Parse(time.RFC3339, record.CreatedAt) 496 if parseErr != nil { 497 createdAt = time.Now() 498 } 499 500 artifact := models.Artifact{ 501 Did: did, 502 Rkey: e.Commit.RKey, 503 RepoDid: syntax.DID(repo.RepoDid), 504 Tag: plumbing.Hash(record.Tag), 505 CreatedAt: createdAt, 506 BlobCid: cid.Cid(record.Artifact.Ref), 507 Name: record.Name, 508 Size: uint64(record.Artifact.Size), 509 MimeType: record.Artifact.MimeType, 510 } 511 512 err = db.AddArtifact(i.Db, artifact) 513 case jmodels.CommitOperationDelete: 514 err = db.DeleteArtifact(i.Db, orm.FilterEq("did", did), orm.FilterEq("rkey", e.Commit.RKey)) 515 } 516 517 if err != nil { 518 return fmt.Errorf("failed to %s artifact record: %w", e.Commit.Operation, err) 519 } 520 521 l.Info("ingested record") 522 return nil 523} 524 525func (i *Ingester) ingestProfile(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { 526 did := e.Did 527 var err error 528 529 l = l.With("handler", "ingestProfile") 530 531 if e.Commit.RKey != "self" { 532 return fmt.Errorf("ingestProfile only ingests `self` record") 533 } 534 535 switch e.Commit.Operation { 536 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 537 raw := json.RawMessage(e.Commit.Record) 538 record := tangled.ActorProfile{} 539 err = json.Unmarshal(raw, &record) 540 if err != nil { 541 l.Error("invalid record", "err", err) 542 return err 543 } 544 545 avatar := "" 546 if record.Avatar != nil { 547 avatar = record.Avatar.Ref.String() 548 } 549 550 description := "" 551 if record.Description != nil { 552 description = *record.Description 553 } 554 555 includeBluesky := record.Bluesky 556 557 pronouns := "" 558 if record.Pronouns != nil { 559 pronouns = *record.Pronouns 560 } 561 562 location := "" 563 if record.Location != nil { 564 location = *record.Location 565 } 566 567 var links [5]string 568 for i, l := range record.Links { 569 if i < 5 { 570 links[i] = l 571 } 572 } 573 574 var stats [2]models.VanityStat 575 for i, s := range record.Stats { 576 if i < 2 { 577 stats[i].Kind = models.ParseVanityStatKind(s) 578 } 579 } 580 581 var pinned [6]string 582 for i, r := range record.PinnedRepositories { 583 if i < 6 { 584 pinned[i] = r 585 } 586 } 587 588 var preferredHandle syntax.Handle 589 if record.PreferredHandle != nil { 590 if h, err := syntax.ParseHandle(*record.PreferredHandle); err == nil { 591 ident, identErr := i.IdResolver.ResolveIdent(ctx, did) 592 if identErr == nil && slices.Contains(ident.AlsoKnownAs, "at://"+string(h)) { 593 preferredHandle = h 594 } 595 } 596 } 597 598 profile := models.Profile{ 599 Did: did, 600 Avatar: avatar, 601 Description: description, 602 IncludeBluesky: includeBluesky, 603 Location: location, 604 Links: links, 605 Stats: stats, 606 PinnedRepos: pinned, 607 Pronouns: pronouns, 608 PreferredHandle: preferredHandle, 609 } 610 611 err = db.ValidateProfile(i.Db, &profile) 612 if err != nil { 613 return fmt.Errorf("invalid profile record: %w", err) 614 } 615 616 err = db.UpsertProfile(i.Db, &profile) 617 if err != nil { 618 return fmt.Errorf("upserting profile: %w", err) 619 } 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 l.Error("attempt to delete repos by knot", "knot", domain) 1219 // err = db.RemoveReposByKnot(tx, domain) 1220 // if err != nil { 1221 // return fmt.Errorf("failed to remove repos by knot: %w", err) 1222 // } 1223 1224 if registration.Registered != nil { 1225 err = i.Enforcer.RemoveKnot(domain) 1226 if err != nil { 1227 return fmt.Errorf("failed to remove knot from enforcer: %w", err) 1228 } 1229 } 1230 1231 err = tx.Commit() 1232 if err != nil { 1233 return fmt.Errorf("failed to commit txn: %w", err) 1234 } 1235 1236 err = i.Enforcer.E.SavePolicy() 1237 if err != nil { 1238 return fmt.Errorf("failed to save ACLs: %w", err) 1239 } 1240 1241 l.Info("ingested record", "domain", domain) 1242 } 1243 1244 return nil 1245} 1246 1247const ( 1248 verifyAttempts = 4 1249 verifyMinDelay = 1 * time.Second 1250 verifyMaxDelay = 5 * time.Second 1251) 1252 1253func (i *Ingester) verifyKnot(ctx context.Context, domain, did string) error { 1254 regs, err := db.GetRegistrations(i.Db, 1255 orm.FilterEq("domain", domain), 1256 orm.FilterEq("did", did), 1257 ) 1258 if err != nil { 1259 return fmt.Errorf("look up registration: %w", err) 1260 } 1261 if len(regs) != 1 { 1262 return fmt.Errorf("no registration for %s by %s", domain, did) 1263 } 1264 if regs[0].Registered != nil { 1265 return nil 1266 } 1267 1268 err = retry.Do( 1269 func() error { return serververify.RunVerification(ctx, domain, did, i.Config.Core.Dev) }, 1270 retry.Context(ctx), 1271 retry.Attempts(verifyAttempts), 1272 retry.Delay(verifyMinDelay), 1273 retry.MaxDelay(verifyMaxDelay), 1274 retry.DelayType(retry.BackOffDelay), 1275 retry.LastErrorOnly(true), 1276 ) 1277 if err != nil { 1278 return fmt.Errorf("verify: %w", err) 1279 } 1280 return serververify.MarkKnotVerified(i.Db, i.Enforcer, domain, did) 1281} 1282 1283func (i *Ingester) verifySpindle(ctx context.Context, instance, did string) error { 1284 spindles, err := db.GetSpindles(ctx, i.Db, 1285 orm.FilterEq("instance", instance), 1286 orm.FilterEq("owner", did), 1287 ) 1288 if err != nil { 1289 return fmt.Errorf("look up spindle: %w", err) 1290 } 1291 if len(spindles) != 1 { 1292 return fmt.Errorf("no spindle for %s by %s", instance, did) 1293 } 1294 if spindles[0].Verified != nil { 1295 return nil 1296 } 1297 1298 err = retry.Do( 1299 func() error { return serververify.RunVerification(ctx, instance, did, i.Config.Core.Dev) }, 1300 retry.Context(ctx), 1301 retry.Attempts(verifyAttempts), 1302 retry.Delay(verifyMinDelay), 1303 retry.MaxDelay(verifyMaxDelay), 1304 retry.DelayType(retry.BackOffDelay), 1305 retry.LastErrorOnly(true), 1306 ) 1307 if err != nil { 1308 return fmt.Errorf("verify: %w", err) 1309 } 1310 _, err = serververify.MarkSpindleVerified(i.Db, i.Enforcer, instance, did) 1311 return err 1312} 1313 1314const sweepConcurrency = 4 1315 1316func (i *Ingester) SweepPendingVerifications() { 1317 l := i.Logger.With("handler", "SweepPendingVerifications") 1318 1319 var g errgroup.Group 1320 g.SetLimit(sweepConcurrency) 1321 1322 regs, err := db.GetRegistrations(i.Db, orm.FilterIs("registered", nil)) 1323 if err != nil { 1324 l.Error("failed to list unverified knots", "err", err) 1325 } else { 1326 for _, reg := range regs { 1327 g.Go(func() error { 1328 if err := i.verifyKnot(i.Ctx, reg.Domain, reg.ByDid); err != nil { 1329 l.Warn("verify knot failed", "domain", reg.Domain, "did", reg.ByDid, "err", err) 1330 } 1331 return nil 1332 }) 1333 } 1334 } 1335 1336 spindles, err := db.GetSpindles(i.Ctx, i.Db, orm.FilterIs("verified", nil)) 1337 if err != nil { 1338 l.Error("failed to list unverified spindles", "err", err) 1339 g.Wait() 1340 return 1341 } 1342 for _, s := range spindles { 1343 g.Go(func() error { 1344 if err := i.verifySpindle(i.Ctx, s.Instance, s.Owner.String()); err != nil { 1345 l.Warn("verify spindle failed", "instance", s.Instance, "owner", s.Owner, "err", err) 1346 } 1347 return nil 1348 }) 1349 } 1350 g.Wait() 1351} 1352 1353func (i *Ingester) ingestIssue(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { 1354 did := e.Did 1355 rkey := e.Commit.RKey 1356 1357 var err error 1358 1359 l = l.With("handler", "ingestIssue") 1360 1361 switch e.Commit.Operation { 1362 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 1363 raw := json.RawMessage(e.Commit.Record) 1364 record := tangled.RepoIssue{} 1365 err = json.Unmarshal(raw, &record) 1366 if err != nil { 1367 l.Error("invalid record", "err", err) 1368 return err 1369 } 1370 1371 issue := models.IssueFromRecord(did, rkey, record) 1372 1373 if issue.RepoDid == "" { 1374 return fmt.Errorf("issue record has no repo field") 1375 } 1376 if _, err := syntax.ParseDID(string(issue.RepoDid)); err != nil { 1377 return fmt.Errorf("issue record repo field is not a valid DID: %w", err) 1378 } 1379 1380 if err := issue.Validate(); err != nil { 1381 return fmt.Errorf("failed to validate issue: %w", err) 1382 } 1383 1384 if record.Repo != "" && !strings.HasPrefix(record.Repo, "did:") { 1385 repo, repoErr := db.GetRepoByAtUri(i.Db, record.Repo) 1386 if repoErr == nil && repo.RepoDid != "" { 1387 if enqErr := db.EnqueuePdsRecordMigration(ctx, i.Db, "add-repo-did", syntax.DID(did), syntax.NSID(tangled.RepoIssueNSID), syntax.RecordKey(e.Commit.RKey)); enqErr != nil { 1388 l.Warn("failed to enqueue PDS rewrite for issue", "err", enqErr, "did", did, "repoDid", repo.RepoDid) 1389 } 1390 } 1391 } 1392 1393 tx, err := i.Db.BeginTx(ctx, nil) 1394 if err != nil { 1395 l.Error("failed to begin transaction", "err", err) 1396 return err 1397 } 1398 defer tx.Rollback() 1399 1400 err = db.PutIssue(tx, &issue) 1401 if err != nil { 1402 l.Error("failed to create issue", "err", err) 1403 return err 1404 } 1405 1406 if err := db.ResolveIssueState(tx, issue.AtUri()); err != nil { 1407 l.Error("failed to resolve issue state", "err", err) 1408 return err 1409 } 1410 1411 err = tx.Commit() 1412 if err != nil { 1413 l.Error("failed to commit txn", "err", err) 1414 return err 1415 } 1416 1417 i.drainPendingState(ctx, issue.AtUri(), l) 1418 1419 l.Info("ingested record") 1420 return nil 1421 1422 case jmodels.CommitOperationDelete: 1423 tx, err := i.Db.BeginTx(ctx, nil) 1424 if err != nil { 1425 l.Error("failed to begin transaction", "err", err) 1426 return err 1427 } 1428 defer tx.Rollback() 1429 1430 if err := db.DeleteIssues( 1431 tx, 1432 did, 1433 rkey, 1434 ); err != nil { 1435 l.Error("failed to delete", "err", err) 1436 return fmt.Errorf("failed to delete issue record: %w", err) 1437 } 1438 if err := tx.Commit(); err != nil { 1439 l.Error("failed to commit txn", "err", err) 1440 return err 1441 } 1442 1443 l.Info("ingested record") 1444 return nil 1445 } 1446 1447 return nil 1448} 1449 1450func (i *Ingester) ingestPull(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { 1451 did := syntax.DID(e.Did) 1452 rkey := syntax.RecordKey(e.Commit.RKey) 1453 cid := syntax.CID(e.Commit.CID) 1454 1455 var err error 1456 1457 l = l.With("handler", "ingestPull") 1458 1459 switch e.Commit.Operation { 1460 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 1461 raw := json.RawMessage(e.Commit.Record) 1462 record := tangled.RepoPull{} 1463 err = json.Unmarshal(raw, &record) 1464 if err != nil { 1465 l.Error("invalid record", "err", err) 1466 return err 1467 } 1468 1469 versions, err := func() ([]models.PullVersion, error) { 1470 if len(record.Versions) > 0 { 1471 versions := make([]models.PullVersion, len(record.Versions)) 1472 var err error 1473 for i, v := range record.Versions { 1474 versions[i], err = models.PullVersionFromRecord(i, v) 1475 if err != nil { 1476 return nil, fmt.Errorf("versions[%d]: %w", i, err) 1477 } 1478 } 1479 return versions, nil 1480 } 1481 1482 if len(record.Rounds) == 0 { 1483 return nil, nil 1484 } 1485 1486 ownerId, err := i.IdResolver.Directory().LookupDID(ctx, did) 1487 if err != nil { 1488 return nil, fmt.Errorf("failed to resolve did: %w", err) 1489 } 1490 1491 // go through and fetch all blobs in parallel 1492 versions := make([]models.PullVersion, len(record.Rounds)) 1493 var mu sync.Mutex 1494 1495 g, gctx := errgroup.WithContext(ctx) 1496 1497 for idx, b := range record.Rounds { 1498 g.Go(func() error { 1499 // for some reason, a blob is empty 1500 if b.PatchBlob == nil { 1501 return fmt.Errorf("missing patchBlob in round %d", idx) 1502 } 1503 1504 ownerPds := ownerId.PDSEndpoint() 1505 url, _ := url.Parse(fmt.Sprintf("%s/xrpc/com.atproto.sync.getBlob", ownerPds)) 1506 q := url.Query() 1507 q.Set("cid", b.PatchBlob.Ref.String()) 1508 q.Set("did", did.String()) 1509 url.RawQuery = q.Encode() 1510 1511 req, err := http.NewRequestWithContext(gctx, http.MethodGet, url.String(), nil) 1512 if err != nil { 1513 return fmt.Errorf("versions[%d]: failed to create request: %w", idx, err) 1514 } 1515 req.Header.Set("Content-Type", "application/json") 1516 1517 resp, err := http.DefaultClient.Do(req) 1518 if err != nil { 1519 return fmt.Errorf("versions[%d]: failed to make request: %w", idx, err) 1520 } 1521 defer resp.Body.Close() 1522 1523 version, err := models.PullVersionFromLegacy(idx, b, resp.Body) 1524 if err != nil { 1525 return fmt.Errorf("versions[%d]: %w", idx, err) 1526 } 1527 1528 mu.Lock() 1529 versions[idx] = version 1530 mu.Unlock() 1531 1532 return nil 1533 }) 1534 } 1535 1536 if err := g.Wait(); err != nil { 1537 return nil, err 1538 } 1539 1540 return versions, nil 1541 }() 1542 if err != nil { 1543 return fmt.Errorf("failed pares pull versions: %w", err) 1544 } 1545 1546 pull, err := models.PullFromRecord(did, rkey, cid, record, versions) 1547 if err != nil { 1548 return fmt.Errorf("failed to parse pull from record: %w", err) 1549 } 1550 if err := pull.Validate(); err != nil { 1551 return fmt.Errorf("failed to validate pull: %w", err) 1552 } 1553 1554 var references []syntax.ATURI 1555 if pull.Body != "" { 1556 _, references = i.MentionsResolver.Resolve(ctx, pull.Body) 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(ctx, tx, pull, references) 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(), 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(ctx context.Context, 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(ctx context.Context, 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(ctx context.Context, e db.Execer, subject syntax.ATURI) (*models.Repo, string, bool, error) { 1680 pull, err := db.GetPull(ctx, e, orm.FilterEq("at_uri", subject)) 1681 if err != nil { 1682 return nil, "", false, err 1683 } 1684 if pull.Repo == nil { 1685 return nil, "", false, nil 1686 } 1687 return pull.Repo, pull.OwnerDid.String(), 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(ctx, 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 record for retry", "subject", subject, "nsid", nsid) 1823 return nil 1824} 1825 1826func (i *Ingester) drainPendingState(ctx context.Context, subject syntax.ATURI, l *slog.Logger) { 1827 pending, err := db.PendingStateRecordsForSubject(i.Db, subject) 1828 if err != nil { 1829 l.Error("failed to load pending records", "err", err, "subject", subject) 1830 return 1831 } 1832 for _, p := range pending { 1833 if err := i.reapplyPendingRecord(ctx, p, l); err != nil { 1834 l.Error("failed to drain pending record", "err", err, "did", p.Did, "rkey", p.Rkey, "nsid", p.Nsid) 1835 } 1836 } 1837} 1838 1839func (i *Ingester) drainPendingLabelOps(l *slog.Logger) { 1840 subjects, err := db.PendingStateSubjectsForNsid(i.Db, tangled.LabelOpNSID) 1841 if err != nil { 1842 l.Error("failed to list pending label op subjects", "err", err) 1843 return 1844 } 1845 for _, subject := range subjects { 1846 i.drainPendingState(i.Ctx, subject, l) 1847 } 1848} 1849 1850func (i *Ingester) reapplyPendingRecord(ctx context.Context, p db.PendingStateRecord, l *slog.Logger) error { 1851 switch p.Nsid { 1852 case tangled.RepoIssueStateNSID: 1853 return i.applyStateRecord(ctx, p.Did, p.Rkey, p.Nsid, p.Record, issueStateSpec, l) 1854 case tangled.RepoPullStatusNSID: 1855 return i.applyStateRecord(ctx, p.Did, p.Rkey, p.Nsid, p.Record, pullStatusSpec, l) 1856 case tangled.LabelOpNSID: 1857 return i.applyLabelOpRecord(ctx, p.Did, p.Rkey, p.Record, l) 1858 default: 1859 return fmt.Errorf("no reapply handler for parked nsid: %s", p.Nsid) 1860 } 1861} 1862 1863const ( 1864 pendingStateReconcileInterval = time.Hour 1865 pendingStateRecordTTL = 7 * 24 * time.Hour 1866) 1867 1868func (i *Ingester) StartPendingStateReconciler() { 1869 i.ReconcilePendingState() 1870 1871 ticker := time.NewTicker(pendingStateReconcileInterval) 1872 defer ticker.Stop() 1873 for { 1874 select { 1875 case <-i.Ctx.Done(): 1876 return 1877 case <-ticker.C: 1878 i.ReconcilePendingState() 1879 } 1880 } 1881} 1882 1883func (i *Ingester) ReconcilePendingState() { 1884 l := i.Logger.With("handler", "reconcilePendingState") 1885 1886 subjects, err := db.DistinctPendingStateSubjects(i.Db) 1887 if err != nil { 1888 l.Error("failed to list pending state subjects", "err", err) 1889 } 1890 for _, subject := range subjects { 1891 i.drainPendingState(i.Ctx, subject, l) 1892 } 1893 1894 cutoff := time.Now().Add(-pendingStateRecordTTL).UTC().Format(time.RFC3339) 1895 evicted, err := db.EvictStalePendingStateRecords(i.Db, cutoff) 1896 if err != nil { 1897 l.Error("failed to evict stale pending state records", "err", err) 1898 return 1899 } 1900 if evicted > 0 { 1901 l.Warn("evicted stale pending state records", "count", evicted, "olderThan", cutoff) 1902 } 1903} 1904 1905// ingestIssueComment ingests legacy sh.tangled.repo.issue.comment deletions 1906func (i *Ingester) ingestIssueComment(e *jmodels.Event, l *slog.Logger) error { 1907 l = l.With("handler", "ingestIssueComment") 1908 1909 switch e.Commit.Operation { 1910 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 1911 // no-op. sh.tangled.repo.issue.comment is deprecated 1912 1913 case jmodels.CommitOperationDelete: 1914 if err := db.PurgeComments( 1915 i.Db, 1916 orm.FilterEq("did", e.Did), 1917 orm.FilterEq("collection", e.Commit.Collection), 1918 orm.FilterEq("rkey", e.Commit.RKey), 1919 ); err != nil { 1920 return fmt.Errorf("failed to delete comment record: %w", err) 1921 } 1922 } 1923 1924 l.Info("ingested record") 1925 return nil 1926} 1927 1928// ingestPullComment ingests legacy sh.tangled.repo.pull.comment deletions 1929func (i *Ingester) ingestPullComment(e *jmodels.Event, l *slog.Logger) error { 1930 l = l.With("handler", "ingestPullComment") 1931 1932 switch e.Commit.Operation { 1933 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 1934 // no-op. sh.tangled.repo.pull.comment is deprecated 1935 1936 case jmodels.CommitOperationDelete: 1937 if err := db.PurgeComments( 1938 i.Db, 1939 orm.FilterEq("did", e.Did), 1940 orm.FilterEq("collection", e.Commit.Collection), 1941 orm.FilterEq("rkey", e.Commit.RKey), 1942 ); err != nil { 1943 return fmt.Errorf("failed to delete comment record: %w", err) 1944 } 1945 } 1946 1947 l.Info("ingested record") 1948 return nil 1949} 1950 1951func (i *Ingester) ingestComment(e *jmodels.Event, l *slog.Logger) error { 1952 did := e.Did 1953 rkey := e.Commit.RKey 1954 cid := e.Commit.CID 1955 1956 var err error 1957 1958 l = l.With("handler", "ingestComment") 1959 1960 ctx := context.Background() 1961 1962 switch e.Commit.Operation { 1963 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 1964 raw := json.RawMessage(e.Commit.Record) 1965 record := tangled.FeedComment{} 1966 err = json.Unmarshal(raw, &record) 1967 if err != nil { 1968 return fmt.Errorf("invalid record: %w", err) 1969 } 1970 1971 comment, err := models.CommentFromRecord(syntax.DID(did), syntax.RecordKey(rkey), syntax.CID(cid), record) 1972 if err != nil { 1973 return fmt.Errorf("failed to parse comment from record: %w", err) 1974 } 1975 1976 if err := comment.Validate(); err != nil { 1977 return fmt.Errorf("failed to validate comment: %w", err) 1978 } 1979 1980 var references []syntax.ATURI 1981 if comment.Body.Original != nil { 1982 _, references = i.MentionsResolver.Resolve(ctx, *comment.Body.Original) 1983 } 1984 1985 tx, err := i.Db.Begin() 1986 if err != nil { 1987 return fmt.Errorf("failed to start transaction: %w", err) 1988 } 1989 defer tx.Rollback() 1990 1991 _, err = db.PutComment(tx, comment, references) 1992 if err != nil { 1993 return fmt.Errorf("failed to create comment: %w", err) 1994 } 1995 1996 if err := tx.Commit(); err != nil { 1997 return err 1998 } 1999 2000 case jmodels.CommitOperationDelete: 2001 if err := db.DeleteComments( 2002 i.Db, 2003 orm.FilterEq("did", did), 2004 orm.FilterEq("collection", e.Commit.Collection), 2005 orm.FilterEq("rkey", rkey), 2006 ); err != nil { 2007 return fmt.Errorf("failed to delete comment record: %w", err) 2008 } 2009 } 2010 2011 l.Info("ingested record") 2012 return nil 2013} 2014 2015func (i *Ingester) ingestReaction(e *jmodels.Event, l *slog.Logger) error { 2016 did := e.Did 2017 rkey := e.Commit.RKey 2018 2019 l = l.With("handler", "ingestReaction") 2020 2021 switch e.Commit.Operation { 2022 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 2023 raw := json.RawMessage(e.Commit.Record) 2024 record := tangled.FeedReaction{} 2025 if err := json.Unmarshal(raw, &record); err != nil { 2026 return fmt.Errorf("invalid record: %w", err) 2027 } 2028 2029 subjectUri, err := syntax.ParseATURI(record.Subject) 2030 if err != nil { 2031 return fmt.Errorf("invalid reaction subject %q: %w", record.Subject, err) 2032 } 2033 subjectUri = models.NormalizeReactionSubject(subjectUri) 2034 2035 kind, ok := models.ParseReactionKind(record.Reaction) 2036 if !ok { 2037 return fmt.Errorf("invalid reaction kind: %q", record.Reaction) 2038 } 2039 2040 created, parseErr := time.Parse(time.RFC3339, record.CreatedAt) 2041 if parseErr != nil { 2042 created = time.Now() 2043 } 2044 2045 reaction := models.Reaction{ 2046 ReactedByDid: did, 2047 Rkey: rkey, 2048 ThreadAt: subjectUri, 2049 Kind: kind, 2050 Created: created, 2051 } 2052 if err := db.UpsertReaction(i.Db, reaction); err != nil { 2053 return fmt.Errorf("failed to upsert reaction: %w", err) 2054 } 2055 2056 case jmodels.CommitOperationDelete: 2057 if err := db.DeleteReactionByRkey(i.Db, did, rkey); err != nil { 2058 return fmt.Errorf("failed to delete reaction record: %w", err) 2059 } 2060 } 2061 2062 l.Info("ingested record") 2063 return nil 2064} 2065 2066func (i *Ingester) ingestLabelDefinition(e *jmodels.Event, l *slog.Logger) error { 2067 did := e.Did 2068 rkey := e.Commit.RKey 2069 2070 var err error 2071 2072 l = l.With("handler", "ingestLabelDefinition") 2073 2074 switch e.Commit.Operation { 2075 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 2076 raw := json.RawMessage(e.Commit.Record) 2077 record := tangled.LabelDefinition{} 2078 err = json.Unmarshal(raw, &record) 2079 if err != nil { 2080 return fmt.Errorf("invalid record: %w", err) 2081 } 2082 2083 def, err := models.LabelDefinitionFromRecord(did, rkey, record) 2084 if err != nil { 2085 return fmt.Errorf("failed to parse labeldef from record: %w", err) 2086 } 2087 2088 if err := def.Validate(); err != nil { 2089 return fmt.Errorf("failed to validate labeldef: %w", err) 2090 } 2091 2092 _, err = db.AddLabelDefinition(i.Db, def) 2093 if err != nil { 2094 return fmt.Errorf("failed to create labeldef: %w", err) 2095 } 2096 2097 if e.Commit.Operation == jmodels.CommitOperationCreate { 2098 i.drainPendingLabelOps(l) 2099 } 2100 2101 l.Info("ingested record") 2102 return nil 2103 2104 case jmodels.CommitOperationDelete: 2105 if err := db.DeleteLabelDefinition( 2106 i.Db, 2107 orm.FilterEq("did", did), 2108 orm.FilterEq("rkey", rkey), 2109 ); err != nil { 2110 return fmt.Errorf("failed to delete labeldef record: %w", err) 2111 } 2112 2113 l.Info("ingested record") 2114 return nil 2115 } 2116 2117 return nil 2118} 2119 2120func (i *Ingester) ingestLabelOp(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { 2121 l = l.With("handler", "ingestLabelOp") 2122 did := e.Did 2123 rkey := e.Commit.RKey 2124 2125 switch e.Commit.Operation { 2126 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 2127 return i.applyLabelOpRecord(ctx, did, rkey, e.Commit.Record, l) 2128 case jmodels.CommitOperationDelete: 2129 return i.deleteLabelOpRecord(ctx, did, rkey, l) 2130 } 2131 2132 return nil 2133} 2134 2135func (i *Ingester) findLabelSubjectRepo(ctx context.Context, subject syntax.ATURI) (*models.Repo, bool, error) { 2136 var spec stateIngestSpec 2137 switch subject.Collection() { 2138 case tangled.RepoIssueNSID: 2139 spec = issueStateSpec 2140 case tangled.RepoPullNSID: 2141 spec = pullStatusSpec 2142 default: 2143 return nil, false, fmt.Errorf("unsupported label subject: %s", subject.Collection()) 2144 } 2145 repo, _, found, err := spec.findSubject(ctx, i.Db, subject) 2146 return repo, found, err 2147} 2148 2149func (i *Ingester) applyLabelOpRecord(ctx context.Context, did, rkey string, raw []byte, l *slog.Logger) error { 2150 record := tangled.LabelOp{} 2151 if err := json.Unmarshal(raw, &record); err != nil { 2152 return fmt.Errorf("invalid record: %w", err) 2153 } 2154 2155 subject := syntax.ATURI(record.Subject) 2156 park := func() error { 2157 return i.parkStateRecord(ctx, did, rkey, tangled.LabelOpNSID, subject, raw, l) 2158 } 2159 2160 repo, found, err := i.findLabelSubjectRepo(ctx, subject) 2161 if err != nil { 2162 return err 2163 } 2164 if !found { 2165 return park() 2166 } 2167 2168 // validate permissions: only collaborators can apply labels currently 2169 // 2170 // TODO: introduce a repo:triage permission 2171 allowed, permErr := i.Acl.HasRepoPermissionErr(ctx, repo, did, "repo:push") 2172 if permErr != nil { 2173 if !errors.Is(permErr, knotacl.ErrKnotUnreachable) { 2174 return park() 2175 } 2176 l.Warn("ingesting labelop without permission check", "did", did, "err", permErr) 2177 allowed = true 2178 } 2179 2180 if !allowed { 2181 if err := i.unparkLabelOp(ctx, did, rkey); err != nil { 2182 return err 2183 } 2184 l.Warn("dropped unauthorized label op", "did", did, "rkey", rkey, "subject", subject) 2185 return nil 2186 } 2187 2188 ops := models.LabelOpsFromRecord(did, rkey, record) 2189 2190 operandKeys := make([]string, 0, len(ops)) 2191 for idx := range ops { 2192 operandKeys = append(operandKeys, ops[idx].OperandKey) 2193 } 2194 2195 actx, err := db.NewLabelApplicationCtx(i.Db, orm.FilterIn("at_uri", operandKeys)) 2196 if err != nil { 2197 return fmt.Errorf("failed to build label application ctx: %w", err) 2198 } 2199 2200 for idx := range ops { 2201 def, ok := actx.Defs[ops[idx].OperandKey] 2202 if !ok { 2203 return park() 2204 } 2205 if err := def.ValidateOperandValue(&ops[idx]); err != nil { 2206 return fmt.Errorf("failed to validate labelop: %w", err) 2207 } 2208 } 2209 2210 if err := i.materializeLabelOps(ctx, did, rkey, ops); err != nil { 2211 return err 2212 } 2213 2214 l.Info("ingested record") 2215 return nil 2216} 2217 2218func (i *Ingester) unparkLabelOp(ctx context.Context, did, rkey string) error { 2219 tx, err := i.Db.BeginTx(ctx, nil) 2220 if err != nil { 2221 return err 2222 } 2223 defer tx.Rollback() 2224 2225 if err := db.UnparkStateRecord(tx, did, rkey, tangled.LabelOpNSID); err != nil { 2226 return fmt.Errorf("failed to unpark label op: %w", err) 2227 } 2228 2229 return tx.Commit() 2230} 2231 2232func (i *Ingester) materializeLabelOps(ctx context.Context, did, rkey string, ops []models.LabelOp) error { 2233 tx, err := i.Db.BeginTx(ctx, nil) 2234 if err != nil { 2235 return err 2236 } 2237 defer tx.Rollback() 2238 2239 if err := db.UnparkStateRecord(tx, did, rkey, tangled.LabelOpNSID); err != nil { 2240 return fmt.Errorf("failed to unpark label op: %w", err) 2241 } 2242 if err := db.DeleteLabelOps(tx, orm.FilterEq("did", did), orm.FilterEq("rkey", rkey)); err != nil { 2243 return fmt.Errorf("failed to clear prior label ops: %w", err) 2244 } 2245 for idx := range ops { 2246 if _, err := db.AddLabelOp(tx, &ops[idx]); err != nil { 2247 return fmt.Errorf("failed to add labelop: %w", err) 2248 } 2249 } 2250 return tx.Commit() 2251} 2252 2253func (i *Ingester) deleteLabelOpRecord(ctx context.Context, did, rkey string, l *slog.Logger) error { 2254 if err := i.materializeLabelOps(ctx, did, rkey, nil); err != nil { 2255 return err 2256 } 2257 l.Info("ingested record") 2258 return nil 2259}