This repository has no description
0

Configure Feed

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

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