Monorepo for Tangled
0

Configure Feed

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

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