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