Monorepo for Tangled tangled.org
1

Configure Feed

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

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