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