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