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