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