pipeline.mx raw
1 // Package pipeline provides the event ingestion pipeline for the relay.
2 // Validates, verifies, checks ACL, rate-limits, handles special kinds,
3 // and stores events. Returns a Result suitable for OK envelope responses.
4 package pipeline
5
6 import (
7 "bytes"
8 "strconv"
9 "time"
10
11 "git.smesh.lol/morly/pkg/acl"
12 "git.smesh.lol/nostr/pkg/event"
13 "git.smesh.lol/nostr/pkg/filter"
14 "git.smesh.lol/nostr/pkg/kind"
15 "git.smesh.lol/nostr/pkg/tag"
16 "git.smesh.lol/morly/pkg/metrics"
17 "git.smesh.lol/morly/pkg/relay/protowire"
18 "git.smesh.lol/morly/pkg/relay/ratelimit"
19 "git.smesh.lol/morly/pkg/store"
20 "git.smesh.lol/musiquay/pkg/protocol"
21 )
22
23 // Result is the outcome of event ingestion, maps to the OK envelope.
24 type Result struct {
25 OK bool
26 Reason []byte
27 }
28
29 func accepted() (r *Result) { return &Result{OK: true} }
30 func rejected(reason string) (r *Result) { return &Result{Reason: []byte(reason)} }
31 func rejectedB(reason []byte) (r *Result) { return &Result{Reason: reason} }
32
33 // Config holds pipeline limits.
34 type Config struct {
35 MaxFuture int64 // max seconds ahead of now (default 900)
36 MaxPast int64 // max seconds behind now (0 = unlimited)
37 MaxContent int32 // max content bytes (default 70000)
38 MaxTags int32 // max tags per event (default 2000)
39 MaxTagElem int32 // max bytes per tag element (default 1024)
40 }
41
42 // DefaultConfig returns sensible defaults.
43 func DefaultConfig() (c Config) {
44 return Config{
45 MaxFuture: 900,
46 MaxContent: 70000,
47 MaxTags: 2000,
48 MaxTagElem: 1024,
49 }
50 }
51
52 // Pipeline processes incoming events through the full ingestion path.
53 type Pipeline struct {
54 store *store.Engine
55 acl acl.Checker
56 limiter *ratelimit.Limiter
57 cfg Config
58 }
59
60 // New creates a pipeline. Limiter may be nil to disable rate limiting.
61 func New(s *store.Engine, a acl.Checker, l *ratelimit.Limiter, cfg Config) (p *Pipeline) {
62 return &Pipeline{store: s, acl: a, limiter: l, cfg: cfg}
63 }
64
65 // SetACL updates the ACL checker.
66 func (p *Pipeline) SetACL(a acl.Checker) { p.acl = a }
67
68 // StageA runs schema validation + signature verification on a parsed event.
69 // Returns nil on success, or a rejection Result on failure.
70 // Used by both the synchronous path (Ingest) and async worker domains
71 // (IngestWorker) so the logic is implemented exactly once.
72 func StageA(ev *event.E) (r *Result) {
73 if len(ev.ID) != 32 {
74 return &Result{Reason: []byte("invalid: id must be 32 bytes")}
75 }
76 if len(ev.Pubkey) != 32 {
77 return &Result{Reason: []byte("invalid: pubkey must be 32 bytes")}
78 }
79 if len(ev.Sig) != 64 {
80 return &Result{Reason: []byte("invalid: sig must be 64 bytes")}
81 }
82 if !bytes.Equal(ev.ID, ev.GetIDBytes()) {
83 return &Result{Reason: []byte("invalid: id mismatch")}
84 }
85 verifyStart := metrics.Now()
86 valid, _ := ev.Verify()
87 metrics.SigVerifyNs.Observe(metrics.Since(verifyStart))
88 if !valid {
89 return &Result{Reason: []byte("invalid: bad signature")}
90 }
91 return nil
92 }
93
94 // Ingest validates, verifies, and stores an event. Full single-threaded
95 // path (Stage A + Stage B). Used when no worker pool is configured.
96 func (p *Pipeline) Ingest(ev *event.E) (r *Result) {
97 ingestStart := metrics.Now()
98 defer func() { metrics.IngestPipelineNs.Observe(metrics.Since(ingestStart)) }()
99 if res := StageA(ev); res != nil {
100 return res
101 }
102 return p.IngestPostVerify(ev)
103 }
104
105 // IngestPostVerify runs Stage B: timestamp/content limits, ACL, rate limit,
106 // expiration, replaceable resolution, save. The caller is responsible for
107 // having already run schema validation and signature verification (Stage A),
108 // typically in a worker domain so verify is parallelized.
109 func (p *Pipeline) IngestPostVerify(ev *event.E) (r *Result) {
110 ingestStart := metrics.Now()
111 defer func() { metrics.IngestPipelineNs.Observe(metrics.Since(ingestStart)) }()
112 return p.ingestStageB(ev)
113 }
114
115 func (p *Pipeline) ingestStageB(ev *event.E) (r *Result) {
116 // 2. Timestamp and content limits (checked in both sync and worker paths).
117 if msg := p.validateLimits(ev); msg != nil {
118 return rejectedB(msg)
119 }
120
121 // 2b. Musiquay payload validation. The gate is IsMusiquay, not every kind
122 // ValidateEvent can parse: the kinds this vocabulary reuses from a NIP
123 // (comment 1111, calendar 31923, listing 30402, ...) are shared with
124 // clients that never heard of the protocol, and enforcing Musiquay's
125 // stricter parser on them would turn their valid traffic away. The kinds
126 // the protocol defines itself are the relay's to police. A Musiquay kind
127 // with no parser yet returns nil below, so it is stored rather than
128 // rejected by a rule that does not exist.
129 if protocol.IsMusiquay(ev.Kind) {
130 if verr := protowire.ValidateEvent(ev); verr != nil {
131 return rejected("invalid: " | verr.Error())
132 }
133 }
134
135 // 3. ACL.
136 aclStart := metrics.Now()
137 allowed := p.acl.AllowWrite(ev.Pubkey, ev.Kind)
138 metrics.ACLAllowWriteNs.Observe(metrics.Since(aclStart))
139 if !allowed {
140 return rejected("blocked: pubkey not allowed")
141 }
142
143 // 4. Rate limit.
144 if p.limiter != nil && !p.limiter.Allow(ev.Pubkey) {
145 return rejected("rate-limited: slow down")
146 }
147
148 // 5. Expiration check (NIP-40).
149 if isExpired(ev) {
150 return rejected("invalid: event expired")
151 }
152
153 // 6. Ephemeral - accept but don't store.
154 if kind.IsEphemeral(ev.Kind) {
155 return accepted()
156 }
157
158 // 7. Replaceable kinds.
159 if kind.IsReplaceable(ev.Kind) {
160 return p.handleReplaceable(ev)
161 }
162 // Musiquay's counted kinds sit inside the parameterized-replaceable range
163 // but are regular records, not addresses: each receipt, aggregate and
164 // rollup is stored. Without this the range rule keeps one per kind and
165 // author (all of them have no `d` tag) and deletes every other.
166 if countedInReplaceableRange(ev.Kind) {
167 return p.saveOrDup(ev)
168 }
169 if kind.IsParameterizedReplaceable(ev.Kind) {
170 return p.handleParamReplaceable(ev)
171 }
172
173 // 8. Deletion (NIP-09).
174 if ev.Kind == kind.EventDeletion.K {
175 return p.handleDeletion(ev)
176 }
177
178 // 9. Store regular event.
179 return p.saveOrDup(ev)
180 }
181
182 // --- validation ---
183
184 // validateLimits checks timestamp range and content/tag size limits.
185 // Schema + sig checks are in StageA (already run before this).
186 func (p *Pipeline) validateLimits(ev *event.E) (buf []byte) {
187 if ev.CreatedAt == 0 {
188 return []byte("invalid: created_at missing")
189 }
190
191 now := time.Now().Unix()
192 if p.cfg.MaxFuture > 0 && ev.CreatedAt > now+p.cfg.MaxFuture {
193 return []byte("invalid: created_at too far in future")
194 }
195 if p.cfg.MaxPast > 0 && ev.CreatedAt < now-p.cfg.MaxPast {
196 return []byte("invalid: created_at too far in past")
197 }
198
199 if p.cfg.MaxContent > 0 && len(ev.Content) > p.cfg.MaxContent {
200 return []byte("invalid: content too large")
201 }
202
203 if ev.Tags != nil {
204 if p.cfg.MaxTags > 0 && ev.Tags.Len() > p.cfg.MaxTags {
205 return []byte("invalid: too many tags")
206 }
207 if p.cfg.MaxTagElem > 0 {
208 for _, tg := range ev.Tags.T {
209 for _, elem := range tg.T {
210 if len(elem) > p.cfg.MaxTagElem {
211 return []byte("invalid: tag element too large")
212 }
213 }
214 }
215 }
216 }
217
218 return nil
219 }
220
221 // --- expiration (NIP-40) ---
222
223 func isExpired(ev *event.E) (ok bool) {
224 if ev.Tags == nil {
225 return false
226 }
227 t := ev.Tags.GetFirst([]byte("expiration"))
228 if t == nil || t.Len() < 2 {
229 return false
230 }
231 exp, err := strconv.ParseInt(string(t.Value()), 10, 64)
232 if err != nil || exp <= 0 {
233 return false
234 }
235 return time.Now().Unix() > exp
236 }
237
238 // --- replaceable events ---
239 // Tiebreaker when timestamps are equal: keep the event with the
240 // lexicographically lower ID (deterministic on SHA-256 hashes).
241 // NIP-01 does not specify tiebreaker behavior.
242
243 func (p *Pipeline) handleReplaceable(ev *event.E) (r *Result) {
244 f := &filter.F{
245 Kinds: kind.NewS(kind.New(ev.Kind)),
246 Authors: tag.NewFromBytesSlice(ev.Pubkey),
247 }
248 limit := uint32(1)
249 f.Limit = &limit
250
251 existing, err := p.store.QueryEvents(f)
252 if err == nil && len(existing) > 0 {
253 old := existing[0]
254 if old.CreatedAt > ev.CreatedAt {
255 return rejected("duplicate: newer version exists")
256 }
257 if old.CreatedAt == ev.CreatedAt && bytes.Compare(old.ID, ev.ID) < 0 {
258 return rejected("duplicate: event with lower id exists at same timestamp")
259 }
260 p.store.DeleteEvent(old.ID)
261 }
262 return p.saveOrDup(ev)
263 }
264
265 func (p *Pipeline) handleParamReplaceable(ev *event.E) (r *Result) {
266 dVal := dTagValue(ev)
267
268 f := &filter.F{
269 Kinds: kind.NewS(kind.New(ev.Kind)),
270 Authors: tag.NewFromBytesSlice(ev.Pubkey),
271 }
272 existing, err := p.store.QueryEvents(f)
273 if err == nil {
274 for _, old := range existing {
275 if !bytes.Equal(dTagValue(old), dVal) {
276 continue
277 }
278 if old.CreatedAt > ev.CreatedAt {
279 return rejected("duplicate: newer version exists")
280 }
281 if old.CreatedAt == ev.CreatedAt && bytes.Compare(old.ID, ev.ID) < 0 {
282 return rejected("duplicate: event with lower id exists at same timestamp")
283 }
284 p.store.DeleteEvent(old.ID)
285 break
286 }
287 }
288 return p.saveOrDup(ev)
289 }
290
291 func dTagValue(ev *event.E) (buf []byte) {
292 if ev.Tags == nil {
293 return nil
294 }
295 t := ev.Tags.GetFirst([]byte("d"))
296 if t == nil || t.Len() < 2 {
297 return nil
298 }
299 return t.Value()
300 }
301
302 // --- deletion (NIP-09) ---
303
304 func (p *Pipeline) handleDeletion(ev *event.E) (r *Result) {
305 if ev.Tags == nil {
306 return p.saveOrDup(ev)
307 }
308 eTags := ev.Tags.GetAll([]byte("e"))
309 for _, et := range eTags {
310 targetID := et.ValueBinary()
311 if targetID == nil || len(targetID) != 32 {
312 continue
313 }
314 // Only delete events owned by the same pubkey.
315 tf := &filter.F{Ids: tag.NewFromBytesSlice(targetID)}
316 targets, err := p.store.QueryEvents(tf)
317 if err != nil || len(targets) == 0 {
318 continue
319 }
320 if !bytes.Equal(targets[0].Pubkey, ev.Pubkey) {
321 continue
322 }
323 p.store.DeleteEvent(targetID)
324 }
325 return p.saveOrDup(ev)
326 }
327
328 // --- storage helper ---
329
330 func (p *Pipeline) saveOrDup(ev *event.E) (r *Result) {
331 err := p.store.SaveEvent(ev)
332 if err == nil {
333 return accepted()
334 }
335 msg := err.Error()
336 if len(msg) >= 9 && msg[:9] == "duplicate" {
337 return rejected("duplicate: already have this event")
338 }
339 return rejectedB([]byte("error: ") | msg)
340 }
341