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