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