This repository has no description
1package appview
2
3import (
4 "context"
5 "database/sql"
6 "encoding/json"
7 "errors"
8 "fmt"
9 "io"
10 "log/slog"
11 "maps"
12 "net/http"
13 "net/url"
14 "slices"
15 "sync"
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/config"
27 "tangled.org/core/appview/db"
28 "tangled.org/core/appview/models"
29 "tangled.org/core/appview/serververify"
30 "tangled.org/core/appview/validator"
31 "tangled.org/core/idresolver"
32 "tangled.org/core/orm"
33 "tangled.org/core/rbac"
34)
35
36type Ingester struct {
37 Db db.DbWrapper
38 Enforcer *rbac.Enforcer
39 IdResolver *idresolver.Resolver
40 Config *config.Config
41 Logger *slog.Logger
42 Validator *validator.Validator
43}
44
45type processFunc func(ctx context.Context, e *jmodels.Event) error
46
47func (i *Ingester) Ingest() processFunc {
48 return func(ctx context.Context, e *jmodels.Event) error {
49 var err error
50
51 l := i.Logger.With("kind", e.Kind)
52 switch e.Kind {
53 case jmodels.EventKindAccount:
54 if !e.Account.Active && *e.Account.Status == "deactivated" {
55 err = i.IdResolver.InvalidateIdent(ctx, e.Account.Did)
56 }
57 case jmodels.EventKindIdentity:
58 err = i.IdResolver.InvalidateIdent(ctx, e.Identity.Did)
59 case jmodels.EventKindCommit:
60 switch e.Commit.Collection {
61 case tangled.GraphFollowNSID:
62 err = i.ingestFollow(e)
63 case tangled.FeedStarNSID:
64 err = i.ingestStar(e)
65 case tangled.PublicKeyNSID:
66 err = i.ingestPublicKey(e)
67 case tangled.RepoArtifactNSID:
68 err = i.ingestArtifact(e)
69 case tangled.ActorProfileNSID:
70 err = i.ingestProfile(ctx, e)
71 case tangled.SpindleMemberNSID:
72 err = i.ingestSpindleMember(ctx, e)
73 case tangled.SpindleNSID:
74 err = i.ingestSpindle(ctx, e)
75 case tangled.KnotMemberNSID:
76 err = i.ingestKnotMember(e)
77 case tangled.KnotNSID:
78 err = i.ingestKnot(e)
79 case tangled.StringNSID:
80 err = i.ingestString(e)
81 case tangled.RepoIssueNSID:
82 err = i.ingestIssue(ctx, e)
83 case tangled.RepoPullNSID:
84 err = i.ingestPull(ctx, e)
85 case tangled.RepoIssueCommentNSID:
86 err = i.ingestIssueComment(e)
87 case tangled.LabelDefinitionNSID:
88 err = i.ingestLabelDefinition(e)
89 case tangled.LabelOpNSID:
90 err = i.ingestLabelOp(e)
91 }
92 l = i.Logger.With("nsid", e.Commit.Collection)
93 }
94
95 if err != nil {
96 l.Warn("failed to ingest record, skipping", "err", err)
97 }
98
99 lastTimeUs := e.TimeUS + 1
100 if saveErr := i.Db.SaveLastTimeUs(lastTimeUs); saveErr != nil {
101 l.Error("failed to save cursor", "err", saveErr)
102 }
103
104 return nil
105 }
106}
107
108func (i *Ingester) ingestStar(e *jmodels.Event) error {
109 var err error
110 did := e.Did
111
112 l := i.Logger.With("handler", "ingestStar")
113 l = l.With("nsid", e.Commit.Collection)
114
115 switch e.Commit.Operation {
116 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate:
117 var subjectUri syntax.ATURI
118
119 raw := json.RawMessage(e.Commit.Record)
120 record := tangled.FeedStar{}
121 err := json.Unmarshal(raw, &record)
122 if err != nil {
123 l.Error("invalid record", "err", err)
124 return err
125 }
126
127 star := &models.Star{
128 Did: did,
129 Rkey: e.Commit.RKey,
130 }
131
132 switch {
133 case record.SubjectDid != nil:
134 repo, repoErr := db.GetRepo(i.Db, orm.FilterEq("repo_did", *record.SubjectDid))
135 if repoErr == nil {
136 subjectUri = repo.RepoAt()
137 star.RepoAt = subjectUri
138 }
139 case record.Subject != nil:
140 subjectUri, err = syntax.ParseATURI(*record.Subject)
141 if err != nil {
142 l.Error("invalid record", "err", err)
143 return err
144 }
145 star.RepoAt = subjectUri
146 repo, repoErr := db.GetRepoByAtUri(i.Db, subjectUri.String())
147 if repoErr == nil && repo.RepoDid != "" {
148 if enqErr := db.EnqueuePdsRewrite(i.Db, did, repo.RepoDid, tangled.FeedStarNSID, e.Commit.RKey, *record.Subject); enqErr != nil {
149 l.Warn("failed to enqueue PDS rewrite for star", "err", enqErr, "did", did, "repoDid", repo.RepoDid)
150 }
151 }
152 default:
153 l.Error("star record has neither subject nor subjectDid")
154 return fmt.Errorf("star record has neither subject nor subjectDid")
155 }
156 err = db.AddStar(i.Db, star)
157 case jmodels.CommitOperationDelete:
158 err = db.DeleteStarByRkey(i.Db, did, e.Commit.RKey)
159 }
160
161 if err != nil {
162 return fmt.Errorf("failed to %s star record: %w", e.Commit.Operation, err)
163 }
164
165 return nil
166}
167
168func (i *Ingester) ingestFollow(e *jmodels.Event) error {
169 var err error
170 did := e.Did
171
172 l := i.Logger.With("handler", "ingestFollow")
173 l = l.With("nsid", e.Commit.Collection)
174
175 switch e.Commit.Operation {
176 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate:
177 raw := json.RawMessage(e.Commit.Record)
178 record := tangled.GraphFollow{}
179 err = json.Unmarshal(raw, &record)
180 if err != nil {
181 l.Error("invalid record", "err", err)
182 return err
183 }
184
185 err = db.AddFollow(i.Db, &models.Follow{
186 UserDid: did,
187 SubjectDid: record.Subject,
188 Rkey: e.Commit.RKey,
189 })
190 case jmodels.CommitOperationDelete:
191 err = db.DeleteFollowByRkey(i.Db, did, e.Commit.RKey)
192 }
193
194 if err != nil {
195 return fmt.Errorf("failed to %s follow record: %w", e.Commit.Operation, err)
196 }
197
198 return nil
199}
200
201func (i *Ingester) ingestPublicKey(e *jmodels.Event) error {
202 did := e.Did
203 var err error
204
205 l := i.Logger.With("handler", "ingestPublicKey")
206 l = l.With("nsid", e.Commit.Collection)
207
208 switch e.Commit.Operation {
209 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate:
210 l.Debug("processing add of pubkey")
211 raw := json.RawMessage(e.Commit.Record)
212 record := tangled.PublicKey{}
213 err = json.Unmarshal(raw, &record)
214 if err != nil {
215 l.Error("invalid record", "err", err)
216 return err
217 }
218
219 name := record.Name
220 key := record.Key
221 err = db.AddPublicKey(i.Db, did, name, key, e.Commit.RKey)
222 case jmodels.CommitOperationDelete:
223 l.Debug("processing delete of pubkey")
224 err = db.DeletePublicKeyByRkey(i.Db, did, e.Commit.RKey)
225 }
226
227 if err != nil {
228 return fmt.Errorf("failed to %s pubkey record: %w", e.Commit.Operation, err)
229 }
230
231 return nil
232}
233
234func (i *Ingester) ingestArtifact(e *jmodels.Event) error {
235 did := e.Did
236 var err error
237
238 l := i.Logger.With("handler", "ingestArtifact")
239 l = l.With("nsid", e.Commit.Collection)
240
241 switch e.Commit.Operation {
242 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate:
243 raw := json.RawMessage(e.Commit.Record)
244 record := tangled.RepoArtifact{}
245 err = json.Unmarshal(raw, &record)
246 if err != nil {
247 l.Error("invalid record", "err", err)
248 return err
249 }
250
251 var repo *models.Repo
252 if record.RepoDid != nil && *record.RepoDid != "" {
253 repo, err = db.GetRepoByDid(i.Db, *record.RepoDid)
254 if err != nil && !errors.Is(err, sql.ErrNoRows) {
255 return fmt.Errorf("failed to look up repo by DID %s: %w", *record.RepoDid, err)
256 }
257 }
258 if repo == nil && record.Repo != nil {
259 repoAt, parseErr := syntax.ParseATURI(*record.Repo)
260 if parseErr != nil {
261 return parseErr
262 }
263 repo, err = db.GetRepoByAtUri(i.Db, repoAt.String())
264 if err != nil {
265 return err
266 }
267 }
268 if repo == nil {
269 return fmt.Errorf("artifact record has neither valid repoDid nor repo field")
270 }
271
272 ok, err := i.Enforcer.E.Enforce(did, repo.Knot, repo.RepoIdentifier(), "repo:push")
273 if err != nil || !ok {
274 return err
275 }
276
277 repoDid := repo.RepoDid
278 if repoDid == "" && record.RepoDid != nil {
279 repoDid = *record.RepoDid
280 }
281 if repoDid != "" && (record.RepoDid == nil || *record.RepoDid == "") && record.Repo != nil {
282 if enqErr := db.EnqueuePdsRewrite(i.Db, did, repoDid, tangled.RepoArtifactNSID, e.Commit.RKey, *record.Repo); enqErr != nil {
283 l.Warn("failed to enqueue PDS rewrite for artifact", "err", enqErr, "did", did, "repoDid", repoDid)
284 }
285 }
286
287 createdAt, err := time.Parse(time.RFC3339, record.CreatedAt)
288 if err != nil {
289 createdAt = time.Now()
290 }
291
292 artifact := models.Artifact{
293 Did: did,
294 Rkey: e.Commit.RKey,
295 RepoAt: repo.RepoAt(),
296 Tag: plumbing.Hash(record.Tag),
297 CreatedAt: createdAt,
298 BlobCid: cid.Cid(record.Artifact.Ref),
299 Name: record.Name,
300 Size: uint64(record.Artifact.Size),
301 MimeType: record.Artifact.MimeType,
302 }
303
304 err = db.AddArtifact(i.Db, artifact)
305 case jmodels.CommitOperationDelete:
306 err = db.DeleteArtifact(i.Db, orm.FilterEq("did", did), orm.FilterEq("rkey", e.Commit.RKey))
307 }
308
309 if err != nil {
310 return fmt.Errorf("failed to %s artifact record: %w", e.Commit.Operation, err)
311 }
312
313 return nil
314}
315
316func (i *Ingester) ingestProfile(ctx context.Context, e *jmodels.Event) error {
317 did := e.Did
318 var err error
319
320 l := i.Logger.With("handler", "ingestProfile")
321 l = l.With("nsid", e.Commit.Collection)
322
323 if e.Commit.RKey != "self" {
324 return fmt.Errorf("ingestProfile only ingests `self` record")
325 }
326
327 switch e.Commit.Operation {
328 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate:
329 raw := json.RawMessage(e.Commit.Record)
330 record := tangled.ActorProfile{}
331 err = json.Unmarshal(raw, &record)
332 if err != nil {
333 l.Error("invalid record", "err", err)
334 return err
335 }
336
337 avatar := ""
338 if record.Avatar != nil {
339 avatar = record.Avatar.Ref.String()
340 }
341
342 description := ""
343 if record.Description != nil {
344 description = *record.Description
345 }
346
347 includeBluesky := record.Bluesky
348
349 pronouns := ""
350 if record.Pronouns != nil {
351 pronouns = *record.Pronouns
352 }
353
354 location := ""
355 if record.Location != nil {
356 location = *record.Location
357 }
358
359 var links [5]string
360 for i, l := range record.Links {
361 if i < 5 {
362 links[i] = l
363 }
364 }
365
366 var stats [2]models.VanityStat
367 for i, s := range record.Stats {
368 if i < 2 {
369 stats[i].Kind = models.ParseVanityStatKind(s)
370 }
371 }
372
373 var pinned [6]string
374 for i, r := range record.PinnedRepositories {
375 if i < 6 {
376 pinned[i] = r
377 }
378 }
379
380 var preferredHandle syntax.Handle
381 if record.PreferredHandle != nil {
382 if h, err := syntax.ParseHandle(*record.PreferredHandle); err == nil {
383 ident, identErr := i.IdResolver.ResolveIdent(ctx, did)
384 if identErr == nil && slices.Contains(ident.AlsoKnownAs, "at://"+string(h)) {
385 preferredHandle = h
386 }
387 }
388 }
389
390 profile := models.Profile{
391 Did: did,
392 Avatar: avatar,
393 Description: description,
394 IncludeBluesky: includeBluesky,
395 Location: location,
396 Links: links,
397 Stats: stats,
398 PinnedRepos: pinned,
399 Pronouns: pronouns,
400 PreferredHandle: preferredHandle,
401 }
402
403 ddb, ok := i.Db.Execer.(*db.DB)
404 if !ok {
405 return fmt.Errorf("failed to index profile record, invalid db cast")
406 }
407
408 tx, err := ddb.Begin()
409 if err != nil {
410 return fmt.Errorf("failed to start transaction")
411 }
412
413 err = db.ValidateProfile(tx, &profile)
414 if err != nil {
415 return fmt.Errorf("invalid profile record")
416 }
417
418 err = db.UpsertProfile(tx, &profile)
419 case jmodels.CommitOperationDelete:
420 err = db.DeleteArtifact(i.Db, orm.FilterEq("did", did), orm.FilterEq("rkey", e.Commit.RKey))
421 }
422
423 if err != nil {
424 return fmt.Errorf("failed to %s profile record: %w", e.Commit.Operation, err)
425 }
426
427 return nil
428}
429
430func (i *Ingester) ingestSpindleMember(ctx context.Context, e *jmodels.Event) error {
431 did := e.Did
432 var err error
433
434 l := i.Logger.With("handler", "ingestSpindleMember")
435 l = l.With("nsid", e.Commit.Collection)
436
437 switch e.Commit.Operation {
438 case jmodels.CommitOperationCreate:
439 raw := json.RawMessage(e.Commit.Record)
440 record := tangled.SpindleMember{}
441 err = json.Unmarshal(raw, &record)
442 if err != nil {
443 l.Error("invalid record", "err", err)
444 return err
445 }
446
447 // only spindle owner can invite to spindles
448 ok, err := i.Enforcer.IsSpindleInviteAllowed(did, record.Instance)
449 if err != nil || !ok {
450 return fmt.Errorf("failed to enforce permissions: %w", err)
451 }
452
453 memberId, err := i.IdResolver.ResolveIdent(ctx, record.Subject)
454 if err != nil {
455 return err
456 }
457
458 if memberId.Handle.IsInvalidHandle() {
459 return err
460 }
461
462 ddb, ok := i.Db.Execer.(*db.DB)
463 if !ok {
464 return fmt.Errorf("invalid db cast")
465 }
466
467 err = db.AddSpindleMember(ddb, models.SpindleMember{
468 Did: syntax.DID(did),
469 Rkey: e.Commit.RKey,
470 Instance: record.Instance,
471 Subject: memberId.DID,
472 })
473 if !ok {
474 return fmt.Errorf("failed to add to db: %w", err)
475 }
476
477 err = i.Enforcer.AddSpindleMember(record.Instance, memberId.DID.String())
478 if err != nil {
479 return fmt.Errorf("failed to update ACLs: %w", err)
480 }
481
482 l.Info("added spindle member")
483 case jmodels.CommitOperationDelete:
484 rkey := e.Commit.RKey
485
486 ddb, ok := i.Db.Execer.(*db.DB)
487 if !ok {
488 return fmt.Errorf("failed to index profile record, invalid db cast")
489 }
490
491 // get record from db first
492 members, err := db.GetSpindleMembers(
493 ddb,
494 orm.FilterEq("did", did),
495 orm.FilterEq("rkey", rkey),
496 )
497 if err != nil || len(members) != 1 {
498 return fmt.Errorf("failed to get member: %w, len(members) = %d", err, len(members))
499 }
500 member := members[0]
501
502 tx, err := ddb.Begin()
503 if err != nil {
504 return fmt.Errorf("failed to start txn: %w", err)
505 }
506
507 // remove record by rkey && update enforcer
508 if err = db.RemoveSpindleMember(
509 tx,
510 orm.FilterEq("did", did),
511 orm.FilterEq("rkey", rkey),
512 ); err != nil {
513 return fmt.Errorf("failed to remove from db: %w", err)
514 }
515
516 // update enforcer
517 err = i.Enforcer.RemoveSpindleMember(member.Instance, member.Subject.String())
518 if err != nil {
519 return fmt.Errorf("failed to update ACLs: %w", err)
520 }
521
522 if err = tx.Commit(); err != nil {
523 return fmt.Errorf("failed to commit txn: %w", err)
524 }
525
526 if err = i.Enforcer.E.SavePolicy(); err != nil {
527 return fmt.Errorf("failed to save ACLs: %w", err)
528 }
529
530 l.Info("removed spindle member")
531 }
532
533 return nil
534}
535
536func (i *Ingester) ingestSpindle(ctx context.Context, e *jmodels.Event) error {
537 did := e.Did
538 var err error
539
540 l := i.Logger.With("handler", "ingestSpindle")
541 l = l.With("nsid", e.Commit.Collection)
542
543 switch e.Commit.Operation {
544 case jmodels.CommitOperationCreate:
545 raw := json.RawMessage(e.Commit.Record)
546 record := tangled.Spindle{}
547 err = json.Unmarshal(raw, &record)
548 if err != nil {
549 l.Error("invalid record", "err", err)
550 return err
551 }
552
553 instance := e.Commit.RKey
554
555 ddb, ok := i.Db.Execer.(*db.DB)
556 if !ok {
557 return fmt.Errorf("failed to index profile record, invalid db cast")
558 }
559
560 err := db.AddSpindle(ddb, models.Spindle{
561 Owner: syntax.DID(did),
562 Instance: instance,
563 })
564 if err != nil {
565 l.Error("failed to add spindle to db", "err", err, "instance", instance)
566 return err
567 }
568
569 err = retry.Do(
570 func() error { return serververify.RunVerification(ctx, instance, did, i.Config.Core.Dev) },
571 retry.Attempts(5), retry.Delay(5*time.Second), retry.MaxDelay(80*time.Second),
572 retry.DelayType(retry.BackOffDelay), retry.LastErrorOnly(true),
573 )
574 if err != nil {
575 l.Error("failed to verify spindle after retries", "err", err, "instance", instance)
576 return err
577 }
578
579 _, err = serververify.MarkSpindleVerified(ddb, i.Enforcer, instance, did)
580 if err != nil {
581 return fmt.Errorf("failed to mark verified: %w", err)
582 }
583
584 return nil
585
586 case jmodels.CommitOperationDelete:
587 instance := e.Commit.RKey
588
589 ddb, ok := i.Db.Execer.(*db.DB)
590 if !ok {
591 return fmt.Errorf("failed to index profile record, invalid db cast")
592 }
593
594 // get record from db first
595 spindles, err := db.GetSpindles(
596 ctx,
597 ddb,
598 orm.FilterEq("owner", did),
599 orm.FilterEq("instance", instance),
600 )
601 if err != nil || len(spindles) != 1 {
602 return fmt.Errorf("failed to get spindles: %w, len(spindles) = %d", err, len(spindles))
603 }
604 spindle := spindles[0]
605
606 tx, err := ddb.Begin()
607 if err != nil {
608 return err
609 }
610 defer func() {
611 tx.Rollback()
612 i.Enforcer.E.LoadPolicy()
613 }()
614
615 // remove spindle members first
616 err = db.RemoveSpindleMember(
617 tx,
618 orm.FilterEq("owner", did),
619 orm.FilterEq("instance", instance),
620 )
621 if err != nil {
622 return err
623 }
624
625 err = db.DeleteSpindle(
626 tx,
627 orm.FilterEq("owner", did),
628 orm.FilterEq("instance", instance),
629 )
630 if err != nil {
631 return err
632 }
633
634 if spindle.Verified != nil {
635 err = i.Enforcer.RemoveSpindle(instance)
636 if err != nil {
637 return err
638 }
639 }
640
641 err = tx.Commit()
642 if err != nil {
643 return err
644 }
645
646 err = i.Enforcer.E.SavePolicy()
647 if err != nil {
648 return err
649 }
650 }
651
652 return nil
653}
654
655func (i *Ingester) ingestString(e *jmodels.Event) error {
656 did := e.Did
657 rkey := e.Commit.RKey
658
659 var err error
660
661 l := i.Logger.With("handler", "ingestString", "nsid", e.Commit.Collection, "did", did, "rkey", rkey)
662 l.Info("ingesting record")
663
664 ddb, ok := i.Db.Execer.(*db.DB)
665 if !ok {
666 return fmt.Errorf("failed to index string record, invalid db cast")
667 }
668
669 switch e.Commit.Operation {
670 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate:
671 raw := json.RawMessage(e.Commit.Record)
672 record := tangled.String{}
673 err = json.Unmarshal(raw, &record)
674 if err != nil {
675 l.Error("invalid record", "err", err)
676 return err
677 }
678
679 string := models.StringFromRecord(did, rkey, record)
680
681 if err = i.Validator.ValidateString(&string); err != nil {
682 l.Error("invalid record", "err", err)
683 return err
684 }
685
686 if err = db.AddString(ddb, string); err != nil {
687 l.Error("failed to add string", "err", err)
688 return err
689 }
690
691 return nil
692
693 case jmodels.CommitOperationDelete:
694 if err := db.DeleteString(
695 ddb,
696 orm.FilterEq("did", did),
697 orm.FilterEq("rkey", rkey),
698 ); err != nil {
699 l.Error("failed to delete", "err", err)
700 return fmt.Errorf("failed to delete string record: %w", err)
701 }
702
703 return nil
704 }
705
706 return nil
707}
708
709func (i *Ingester) ingestKnotMember(e *jmodels.Event) error {
710 did := e.Did
711 var err error
712
713 l := i.Logger.With("handler", "ingestKnotMember")
714 l = l.With("nsid", e.Commit.Collection)
715
716 switch e.Commit.Operation {
717 case jmodels.CommitOperationCreate:
718 raw := json.RawMessage(e.Commit.Record)
719 record := tangled.KnotMember{}
720 err = json.Unmarshal(raw, &record)
721 if err != nil {
722 l.Error("invalid record", "err", err)
723 return err
724 }
725
726 // only knot owner can invite to knots
727 ok, err := i.Enforcer.IsKnotInviteAllowed(did, record.Domain)
728 if err != nil || !ok {
729 return fmt.Errorf("failed to enforce permissions: %w", err)
730 }
731
732 memberId, err := i.IdResolver.ResolveIdent(context.Background(), record.Subject)
733 if err != nil {
734 return err
735 }
736
737 if memberId.Handle.IsInvalidHandle() {
738 return err
739 }
740
741 err = i.Enforcer.AddKnotMember(record.Domain, memberId.DID.String())
742 if err != nil {
743 return fmt.Errorf("failed to update ACLs: %w", err)
744 }
745
746 l.Info("added knot member")
747 case jmodels.CommitOperationDelete:
748 // we don't store knot members in a table (like we do for spindle)
749 // and we can't remove this just yet. possibly fixed if we switch
750 // to either:
751 // 1. a knot_members table like with spindle and store the rkey
752 // 2. use the knot host as the rkey
753 //
754 // TODO: implement member deletion
755 l.Info("skipping knot member delete", "did", did, "rkey", e.Commit.RKey)
756 }
757
758 return nil
759}
760
761func (i *Ingester) ingestKnot(e *jmodels.Event) error {
762 did := e.Did
763 var err error
764
765 l := i.Logger.With("handler", "ingestKnot")
766 l = l.With("nsid", e.Commit.Collection)
767
768 switch e.Commit.Operation {
769 case jmodels.CommitOperationCreate:
770 raw := json.RawMessage(e.Commit.Record)
771 record := tangled.Knot{}
772 err = json.Unmarshal(raw, &record)
773 if err != nil {
774 l.Error("invalid record", "err", err)
775 return err
776 }
777
778 domain := e.Commit.RKey
779
780 ddb, ok := i.Db.Execer.(*db.DB)
781 if !ok {
782 return fmt.Errorf("failed to index profile record, invalid db cast")
783 }
784
785 err := db.AddKnot(ddb, domain, did)
786 if err != nil {
787 l.Error("failed to add knot to db", "err", err, "domain", domain)
788 return err
789 }
790
791 err = retry.Do(
792 func() error {
793 return serververify.RunVerification(context.Background(), domain, did, i.Config.Core.Dev)
794 },
795 retry.Attempts(5), retry.Delay(5*time.Second), retry.MaxDelay(80*time.Second),
796 retry.DelayType(retry.BackOffDelay), retry.LastErrorOnly(true),
797 )
798 if err != nil {
799 l.Error("failed to verify knot after retries", "err", err, "domain", domain)
800 return err
801 }
802
803 err = serververify.MarkKnotVerified(ddb, i.Enforcer, domain, did)
804 if err != nil {
805 return fmt.Errorf("failed to mark verified: %w", err)
806 }
807
808 return nil
809
810 case jmodels.CommitOperationDelete:
811 domain := e.Commit.RKey
812
813 ddb, ok := i.Db.Execer.(*db.DB)
814 if !ok {
815 return fmt.Errorf("failed to index knot record, invalid db cast")
816 }
817
818 // get record from db first
819 registrations, err := db.GetRegistrations(
820 ddb,
821 orm.FilterEq("domain", domain),
822 orm.FilterEq("did", did),
823 )
824 if err != nil {
825 return fmt.Errorf("failed to get registration: %w", err)
826 }
827 if len(registrations) != 1 {
828 return fmt.Errorf("got incorrect number of registrations: %d, expected 1", len(registrations))
829 }
830 registration := registrations[0]
831
832 tx, err := ddb.Begin()
833 if err != nil {
834 return err
835 }
836 defer func() {
837 tx.Rollback()
838 i.Enforcer.E.LoadPolicy()
839 }()
840
841 err = db.DeleteKnot(
842 tx,
843 orm.FilterEq("did", did),
844 orm.FilterEq("domain", domain),
845 )
846 if err != nil {
847 return err
848 }
849
850 if registration.Registered != nil {
851 err = i.Enforcer.RemoveKnot(domain)
852 if err != nil {
853 return err
854 }
855 }
856
857 err = tx.Commit()
858 if err != nil {
859 return err
860 }
861
862 err = i.Enforcer.E.SavePolicy()
863 if err != nil {
864 return err
865 }
866 }
867
868 return nil
869}
870func (i *Ingester) ingestIssue(ctx context.Context, e *jmodels.Event) error {
871 did := e.Did
872 rkey := e.Commit.RKey
873
874 var err error
875
876 l := i.Logger.With("handler", "ingestIssue", "nsid", e.Commit.Collection, "did", did, "rkey", rkey)
877 l.Info("ingesting record")
878
879 ddb, ok := i.Db.Execer.(*db.DB)
880 if !ok {
881 return fmt.Errorf("failed to index issue record, invalid db cast")
882 }
883
884 switch e.Commit.Operation {
885 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate:
886 raw := json.RawMessage(e.Commit.Record)
887 record := tangled.RepoIssue{}
888 err = json.Unmarshal(raw, &record)
889 if err != nil {
890 l.Error("invalid record", "err", err)
891 return err
892 }
893
894 issue := models.IssueFromRecord(did, rkey, record)
895
896 if issue.RepoAt == "" {
897 return fmt.Errorf("issue record has no repo field")
898 }
899
900 if err := i.Validator.ValidateIssue(&issue); err != nil {
901 return fmt.Errorf("failed to validate issue: %w", err)
902 }
903
904 if record.Repo != nil {
905 repo, repoErr := db.GetRepoByAtUri(i.Db, *record.Repo)
906 if repoErr == nil && repo.RepoDid != "" {
907 if enqErr := db.EnqueuePdsRewrite(i.Db, did, repo.RepoDid, tangled.RepoIssueNSID, rkey, *record.Repo); enqErr != nil {
908 l.Warn("failed to enqueue PDS rewrite for issue", "err", enqErr, "did", did, "repoDid", repo.RepoDid)
909 }
910 }
911 }
912
913 tx, err := ddb.BeginTx(ctx, nil)
914 if err != nil {
915 l.Error("failed to begin transaction", "err", err)
916 return err
917 }
918 defer tx.Rollback()
919
920 err = db.PutIssue(tx, &issue)
921 if err != nil {
922 l.Error("failed to create issue", "err", err)
923 return err
924 }
925
926 err = tx.Commit()
927 if err != nil {
928 l.Error("failed to commit txn", "err", err)
929 return err
930 }
931
932 return nil
933
934 case jmodels.CommitOperationDelete:
935 tx, err := ddb.BeginTx(ctx, nil)
936 if err != nil {
937 l.Error("failed to begin transaction", "err", err)
938 return err
939 }
940 defer tx.Rollback()
941
942 if err := db.DeleteIssues(
943 tx,
944 did,
945 rkey,
946 ); err != nil {
947 l.Error("failed to delete", "err", err)
948 return fmt.Errorf("failed to delete issue record: %w", err)
949 }
950 if err := tx.Commit(); err != nil {
951 l.Error("failed to commit txn", "err", err)
952 return err
953 }
954
955 return nil
956 }
957
958 return nil
959}
960
961func (i *Ingester) ingestPull(ctx context.Context, e *jmodels.Event) error {
962 did := e.Did
963 rkey := e.Commit.RKey
964
965 var err error
966
967 l := i.Logger.With("handler", "ingestPull", "nsid", e.Commit.Collection, "did", did, "rkey", rkey)
968 l.Info("ingesting record")
969
970 ddb, ok := i.Db.Execer.(*db.DB)
971 if !ok {
972 return fmt.Errorf("failed to index pull record, invalid db cast")
973 }
974
975 switch e.Commit.Operation {
976 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate:
977 raw := json.RawMessage(e.Commit.Record)
978 record := tangled.RepoPull{}
979 err = json.Unmarshal(raw, &record)
980 if err != nil {
981 l.Error("invalid record", "err", err)
982 return err
983 }
984
985 ownerId, err := i.IdResolver.ResolveIdent(ctx, did)
986 if err != nil {
987 l.Error("failed to resolve did")
988 return err
989 }
990
991 // go through and fetch all blobs in parallel
992 readers := make([]*io.ReadCloser, len(record.Rounds))
993 var mu sync.Mutex
994
995 g, gctx := errgroup.WithContext(ctx)
996
997 for idx, b := range record.Rounds {
998 g.Go(func() error {
999 // for some reason, a blob is empty
1000 if b.PatchBlob == nil {
1001 return fmt.Errorf("missing patchBlob in round %d", idx)
1002 }
1003
1004 ownerPds := ownerId.PDSEndpoint()
1005 url, _ := url.Parse(fmt.Sprintf("%s/xrpc/com.atproto.sync.getBlob", ownerPds))
1006 q := url.Query()
1007 q.Set("cid", b.PatchBlob.Ref.String())
1008 q.Set("did", did)
1009 url.RawQuery = q.Encode()
1010
1011 req, err := http.NewRequestWithContext(gctx, http.MethodGet, url.String(), nil)
1012 if err != nil {
1013 l.Error("failed to create request")
1014 return err
1015 }
1016 req.Header.Set("Content-Type", "application/json")
1017
1018 resp, err := http.DefaultClient.Do(req)
1019 if err != nil {
1020 l.Error("failed to make request")
1021 return err
1022 }
1023
1024 mu.Lock()
1025 readers[idx] = &resp.Body
1026 mu.Unlock()
1027
1028 return nil
1029 })
1030 }
1031
1032 if err := g.Wait(); err != nil {
1033 for _, r := range readers {
1034 if r != nil && *r != nil {
1035 (*r).Close()
1036 }
1037 }
1038 return err
1039 }
1040
1041 defer func() {
1042 for _, r := range readers {
1043 if r != nil && *r != nil {
1044 (*r).Close()
1045 }
1046 }
1047 }()
1048
1049 pull := models.PullFromRecord(did, rkey, record, readers)
1050 if err := i.Validator.ValidatePull(&pull); err != nil {
1051 return fmt.Errorf("failed to validate pull: %w", err)
1052 }
1053
1054 tx, err := ddb.BeginTx(ctx, nil)
1055 if err != nil {
1056 l.Error("failed to begin transaction", "err", err)
1057 return err
1058 }
1059 defer tx.Rollback()
1060
1061 err = db.PutPull(tx, &pull)
1062 if err != nil {
1063 l.Error("failed to create pull", "err", err)
1064 return err
1065 }
1066
1067 err = tx.Commit()
1068 if err != nil {
1069 l.Error("failed to commit txn", "err", err)
1070 return err
1071 }
1072
1073 return nil
1074
1075 case jmodels.CommitOperationDelete:
1076 tx, err := ddb.BeginTx(ctx, nil)
1077 if err != nil {
1078 l.Error("failed to begin transaction", "err", err)
1079 return err
1080 }
1081 defer tx.Rollback()
1082
1083 if err := db.AbandonPulls(
1084 tx,
1085 orm.FilterEq("owner_did", did),
1086 orm.FilterEq("rkey", rkey),
1087 ); err != nil {
1088 l.Error("failed to abandon", "err", err)
1089 return fmt.Errorf("failed to abandon pull record: %w", err)
1090 }
1091 if err := tx.Commit(); err != nil {
1092 l.Error("failed to commit txn", "err", err)
1093 return err
1094 }
1095
1096 return nil
1097 }
1098
1099 return nil
1100}
1101
1102func (i *Ingester) ingestIssueComment(e *jmodels.Event) error {
1103 did := e.Did
1104 rkey := e.Commit.RKey
1105
1106 var err error
1107
1108 l := i.Logger.With("handler", "ingestIssueComment", "nsid", e.Commit.Collection, "did", did, "rkey", rkey)
1109 l.Info("ingesting record")
1110
1111 ddb, ok := i.Db.Execer.(*db.DB)
1112 if !ok {
1113 return fmt.Errorf("failed to index issue comment record, invalid db cast")
1114 }
1115
1116 switch e.Commit.Operation {
1117 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate:
1118 raw := json.RawMessage(e.Commit.Record)
1119 record := tangled.RepoIssueComment{}
1120 err = json.Unmarshal(raw, &record)
1121 if err != nil {
1122 return fmt.Errorf("invalid record: %w", err)
1123 }
1124
1125 comment, err := models.IssueCommentFromRecord(did, rkey, record)
1126 if err != nil {
1127 return fmt.Errorf("failed to parse comment from record: %w", err)
1128 }
1129
1130 if err := i.Validator.ValidateIssueComment(comment); err != nil {
1131 return fmt.Errorf("failed to validate comment: %w", err)
1132 }
1133
1134 tx, err := ddb.Begin()
1135 if err != nil {
1136 return fmt.Errorf("failed to start transaction: %w", err)
1137 }
1138 defer tx.Rollback()
1139
1140 _, err = db.AddIssueComment(tx, *comment)
1141 if err != nil {
1142 return fmt.Errorf("failed to create issue comment: %w", err)
1143 }
1144
1145 return tx.Commit()
1146
1147 case jmodels.CommitOperationDelete:
1148 if err := db.DeleteIssueComments(
1149 ddb,
1150 orm.FilterEq("did", did),
1151 orm.FilterEq("rkey", rkey),
1152 ); err != nil {
1153 return fmt.Errorf("failed to delete issue comment record: %w", err)
1154 }
1155
1156 return nil
1157 }
1158
1159 return nil
1160}
1161
1162func (i *Ingester) ingestLabelDefinition(e *jmodels.Event) error {
1163 did := e.Did
1164 rkey := e.Commit.RKey
1165
1166 var err error
1167
1168 l := i.Logger.With("handler", "ingestLabelDefinition", "nsid", e.Commit.Collection, "did", did, "rkey", rkey)
1169 l.Info("ingesting record")
1170
1171 ddb, ok := i.Db.Execer.(*db.DB)
1172 if !ok {
1173 return fmt.Errorf("failed to index label definition, invalid db cast")
1174 }
1175
1176 switch e.Commit.Operation {
1177 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate:
1178 raw := json.RawMessage(e.Commit.Record)
1179 record := tangled.LabelDefinition{}
1180 err = json.Unmarshal(raw, &record)
1181 if err != nil {
1182 return fmt.Errorf("invalid record: %w", err)
1183 }
1184
1185 def, err := models.LabelDefinitionFromRecord(did, rkey, record)
1186 if err != nil {
1187 return fmt.Errorf("failed to parse labeldef from record: %w", err)
1188 }
1189
1190 if err := i.Validator.ValidateLabelDefinition(def); err != nil {
1191 return fmt.Errorf("failed to validate labeldef: %w", err)
1192 }
1193
1194 _, err = db.AddLabelDefinition(ddb, def)
1195 if err != nil {
1196 return fmt.Errorf("failed to create labeldef: %w", err)
1197 }
1198
1199 return nil
1200
1201 case jmodels.CommitOperationDelete:
1202 if err := db.DeleteLabelDefinition(
1203 ddb,
1204 orm.FilterEq("did", did),
1205 orm.FilterEq("rkey", rkey),
1206 ); err != nil {
1207 return fmt.Errorf("failed to delete labeldef record: %w", err)
1208 }
1209
1210 return nil
1211 }
1212
1213 return nil
1214}
1215
1216func (i *Ingester) ingestLabelOp(e *jmodels.Event) error {
1217 did := e.Did
1218 rkey := e.Commit.RKey
1219
1220 var err error
1221
1222 l := i.Logger.With("handler", "ingestLabelOp", "nsid", e.Commit.Collection, "did", did, "rkey", rkey)
1223 l.Info("ingesting record")
1224
1225 ddb, ok := i.Db.Execer.(*db.DB)
1226 if !ok {
1227 return fmt.Errorf("failed to index label op, invalid db cast")
1228 }
1229
1230 switch e.Commit.Operation {
1231 case jmodels.CommitOperationCreate:
1232 raw := json.RawMessage(e.Commit.Record)
1233 record := tangled.LabelOp{}
1234 err = json.Unmarshal(raw, &record)
1235 if err != nil {
1236 return fmt.Errorf("invalid record: %w", err)
1237 }
1238
1239 subject := syntax.ATURI(record.Subject)
1240 collection := subject.Collection()
1241
1242 var repo *models.Repo
1243 switch collection {
1244 case tangled.RepoIssueNSID:
1245 i, err := db.GetIssues(ddb, orm.FilterEq("at_uri", subject))
1246 if err != nil || len(i) != 1 {
1247 return fmt.Errorf("failed to find subject: %w || subject count %d", err, len(i))
1248 }
1249 repo = i[0].Repo
1250 default:
1251 return fmt.Errorf("unsupported label subject: %s", collection)
1252 }
1253
1254 actx, err := db.NewLabelApplicationCtx(ddb, orm.FilterIn("at_uri", repo.Labels))
1255 if err != nil {
1256 return fmt.Errorf("failed to build label application ctx: %w", err)
1257 }
1258
1259 ops := models.LabelOpsFromRecord(did, rkey, record)
1260
1261 for _, o := range ops {
1262 def, ok := actx.Defs[o.OperandKey]
1263 if !ok {
1264 return fmt.Errorf("failed to find label def for key: %s, expected: %q", o.OperandKey, slices.Collect(maps.Keys(actx.Defs)))
1265 }
1266 if err := i.Validator.ValidateLabelOp(def, repo, &o); err != nil {
1267 return fmt.Errorf("failed to validate labelop: %w", err)
1268 }
1269 }
1270
1271 tx, err := ddb.Begin()
1272 if err != nil {
1273 return err
1274 }
1275 defer tx.Rollback()
1276
1277 for _, o := range ops {
1278 _, err = db.AddLabelOp(tx, &o)
1279 if err != nil {
1280 return fmt.Errorf("failed to add labelop: %w", err)
1281 }
1282 }
1283
1284 if err = tx.Commit(); err != nil {
1285 return err
1286 }
1287 }
1288
1289 return nil
1290}