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