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