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