Monorepo for Tangled tangled.org
2

Configure Feed

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

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