Monorepo for Tangled tangled.org
2

Configure Feed

Select the types of activity you want to include in your feed.

core / knotserver / ingester.go
4.0 kB 143 lines
1package knotserver 2 3import ( 4 "context" 5 "encoding/json" 6 "fmt" 7 "strings" 8 9 "github.com/bluesky-social/indigo/atproto/syntax" 10 jmodels "github.com/bluesky-social/jetstream/pkg/models" 11 "tangled.org/core/api/tangled" 12 "tangled.org/core/knotserver/db" 13 knotxrpc "tangled.org/core/knotserver/xrpc" 14 "tangled.org/core/log" 15) 16 17func (h *Knot) processPublicKey(ctx context.Context, event *jmodels.Event) error { 18 l := log.FromContext(ctx).With("handler", "processPublicKey", "did", event.Did, "rkey", event.Commit.RKey) 19 did := syntax.DID(event.Did) 20 rkey := syntax.RecordKey(event.Commit.RKey) 21 22 switch event.Commit.Operation { 23 case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: 24 var record tangled.PublicKey 25 if err := json.Unmarshal(json.RawMessage(event.Commit.Record), &record); err != nil { 26 return fmt.Errorf("failed to unmarshal record: %w", err) 27 } 28 29 pk := db.PublicKey{ 30 Did: did, 31 Rkey: rkey, 32 PublicKey: record, 33 } 34 if err := h.db.UpsertPublicKey(pk); err != nil { 35 return fmt.Errorf("failed to upsert public key: %w", err) 36 } 37 l.Info("upserted public key from firehose") 38 case jmodels.CommitOperationDelete: 39 if err := h.db.DeletePublicKeyByRkey(did, rkey); err != nil { 40 return fmt.Errorf("failed to delete public key: %w", err) 41 } 42 l.Info("deleted public key from firehose") 43 } 44 45 return nil 46} 47 48// returns a repo path on disk if present, and error if not 49type targetRepo struct { 50 RepoPath string 51 OwnerDid string 52 RepoName string 53 RepoDid string 54 DefaultBranch string // default branch 55} 56 57func (h *Knot) processRepo(ctx context.Context, event *jmodels.Event) error { 58 l := log.FromContext(ctx).With("handler", "processRepo", "did", event.Did, "rkey", event.Commit.RKey) 59 60 rkey := strings.TrimSuffix(strings.TrimSpace(event.Commit.RKey), ".git") 61 if rkey == "" { 62 return nil 63 } 64 65 if event.Commit.Operation == jmodels.CommitOperationDelete { 66 return nil 67 } 68 69 if event.Commit.Operation != jmodels.CommitOperationCreate && event.Commit.Operation != jmodels.CommitOperationUpdate { 70 return nil 71 } 72 73 raw := json.RawMessage(event.Commit.Record) 74 var record tangled.Repo 75 if err := json.Unmarshal(raw, &record); err != nil { 76 return fmt.Errorf("failed to unmarshal repo record: %w", err) 77 } 78 79 if record.Knot != h.c.Server.Hostname { 80 return nil 81 } 82 if record.RepoDid == nil || *record.RepoDid == "" { 83 l.Info("skipping repo event without repoDid") 84 return nil 85 } 86 repoDid := *record.RepoDid 87 88 if err := knotxrpc.ValidateRepoName(rkey); err != nil { 89 l.Warn("skipping repo event with invalid rkey", "repoDid", repoDid, "rkey", rkey, "err", err) 90 return nil 91 } 92 93 ownerDid, _, lookupErr := h.db.GetRepoKeyOwner(repoDid) 94 if lookupErr != nil { 95 l.Info("skipping repo event for unknown repoDid", "repoDid", repoDid) 96 return nil 97 } 98 if ownerDid != event.Did { 99 l.Warn("repo event author does not own repoDid", "repoDid", repoDid, "author", event.Did) 100 return nil 101 } 102 103 alias := db.RepoAlias{ 104 OwnerDid: event.Did, 105 Rkey: rkey, 106 RepoDid: repoDid, 107 Rev: event.Commit.Rev, 108 } 109 if err := h.db.UpsertRepoAlias(alias); err != nil { 110 l.Warn("failed to upsert repo alias", "err", err) 111 return nil 112 } 113 114 l.Info("recorded repo alias", "repoDid", repoDid, "rkey", rkey, "rev", event.Commit.Rev) 115 return nil 116} 117 118func (h *Knot) processMessages(ctx context.Context, event *jmodels.Event) error { 119 var err error 120 switch event.Kind { 121 case jmodels.EventKindIdentity: 122 err = h.resolver.InvalidateIdent(ctx, event.Did) 123 case jmodels.EventKindCommit: 124 switch event.Commit.Collection { 125 case tangled.PublicKeyNSID: 126 err = h.processPublicKey(ctx, event) 127 case tangled.RepoNSID: 128 err = h.processRepo(ctx, event) 129 } 130 default: 131 return nil 132 } 133 134 if err != nil { 135 args := []any{"kind", event.Kind, "err", err} 136 if event.Kind == jmodels.EventKindCommit { 137 args = append(args, "nsid", event.Commit.Collection, "did", event.Did, "rkey", event.Commit.RKey) 138 } 139 h.l.Warn("failed to process event, skipping", args...) 140 } 141 142 return nil 143}