forked from
tangled.org/core
Monorepo for Tangled
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}