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