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