This repository has no description
0

Configure Feed

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

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