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