This repository has no description
0

Configure Feed

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

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