pipeline_test.mx raw
1 package pipeline
2
3 import (
4 "bytes"
5 "os"
6 "strconv"
7 "testing"
8 "time"
9
10 "git.smesh.lol/morly/pkg/acl"
11 "git.smesh.lol/nostr/pkg/event"
12 "git.smesh.lol/nostr/pkg/filter"
13 "git.smesh.lol/nostr/pkg/kind"
14 "git.smesh.lol/nostr/pkg/signer/p8k"
15 "git.smesh.lol/nostr/pkg/tag"
16 "git.smesh.lol/morly/pkg/relay/ratelimit"
17 "git.smesh.lol/morly/pkg/store"
18 "git.smesh.lol/musiquay/pkg/protocol"
19 )
20
21 // tOpenStore opens a store in a fresh directory. The caller owns both:
22 // defer os.RemoveAll(dir) and defer eng.Close().
23 func tOpenStore(t *testing.T) (eng *store.Engine, dir string) {
24 t.Helper()
25 dir, derr := os.MkdirTemp("", "moxie-pipe-test")
26 if derr != nil {
27 t.Fatal(derr)
28 }
29 eng, oerr := store.Open(dir)
30 if oerr != nil {
31 t.Fatal(oerr)
32 }
33 return eng, dir
34 }
35
36 func tNewSigner(t *testing.T) (s *p8k.Signer) {
37 t.Helper()
38 s = p8k.MustNew()
39 if gerr := s.Generate(); gerr != nil {
40 t.Fatal(gerr)
41 }
42 return s
43 }
44
45 func tSetup(t *testing.T) (p *Pipeline, s *p8k.Signer, eng *store.Engine, dir string) {
46 t.Helper()
47 eng, dir = tOpenStore(t)
48 s = tNewSigner(t)
49 p = New(eng, &acl.Open{}, nil, DefaultConfig())
50 return p, s, eng, dir
51 }
52
53 func tSetupCfg(t *testing.T, cfg Config) (p *Pipeline, s *p8k.Signer, eng *store.Engine, dir string) {
54 t.Helper()
55 eng, dir = tOpenStore(t)
56 s = tNewSigner(t)
57 p = New(eng, &acl.Open{}, nil, cfg)
58 return p, s, eng, dir
59 }
60
61 func tEvent(t *testing.T, s *p8k.Signer, k uint16, ts int64, tags *tag.S, content string) (ev *event.E) {
62 t.Helper()
63 ev = &event.E{
64 CreatedAt: ts,
65 Kind: k,
66 Tags: tags,
67 Content: []byte(content),
68 }
69 if serr := ev.Sign(s); serr != nil {
70 t.Fatal(serr)
71 }
72 return ev
73 }
74
75 func tTag(k, v string) (tt *tag.T) {
76 return tag.NewFromBytesSlice([]byte(k), []byte(v))
77 }
78
79 // tReason flattens a Result for comparison; a nil Result (accepted) is "".
80 func tReason(r *Result) (s string) {
81 if r == nil {
82 return ""
83 }
84 return string(r.Reason)
85 }
86
87 // tCount returns the total number of events visible in the store.
88 func tCount(t *testing.T, eng *store.Engine) (n int32) {
89 t.Helper()
90 evs, qerr := eng.QueryEvents(&filter.F{})
91 if qerr != nil {
92 t.Fatalf("query: %s", qerr.Error())
93 }
94 return int32(len(evs))
95 }
96
97 // --- config ---
98
99 func TestDefaultConfig(t *testing.T) {
100 c := DefaultConfig()
101 if c.MaxFuture != 900 {
102 t.Errorf("MaxFuture = %d, want 900", c.MaxFuture)
103 }
104 if c.MaxPast != 0 {
105 t.Errorf("MaxPast = %d, want 0 (unlimited)", c.MaxPast)
106 }
107 if c.MaxContent != 70000 {
108 t.Errorf("MaxContent = %d, want 70000", c.MaxContent)
109 }
110 if c.MaxTags != 2000 {
111 t.Errorf("MaxTags = %d, want 2000", c.MaxTags)
112 }
113 if c.MaxTagElem != 1024 {
114 t.Errorf("MaxTagElem = %d, want 1024", c.MaxTagElem)
115 }
116 }
117
118 // --- Stage A: schema + signature ---
119
120 func TestStageAValid(t *testing.T) {
121 s := tNewSigner(t)
122 ev := tEvent(t, s, 1, time.Now().Unix(), nil, "hello")
123 if r := StageA(ev); r != nil {
124 t.Fatalf("valid event rejected: %s", r.Reason)
125 }
126 }
127
128 func TestStageABadIDLength(t *testing.T) {
129 s := tNewSigner(t)
130 short := tEvent(t, s, 1, time.Now().Unix(), nil, "short id")
131 short.ID = short.ID[:16]
132 r1 := StageA(short)
133 if tReason(r1) != "invalid: id must be 32 bytes" {
134 t.Errorf("short id: got %q", tReason(r1))
135 }
136
137 long := tEvent(t, s, 1, time.Now().Unix(), nil, "long id")
138 long.ID = []byte{:33}
139 r2 := StageA(long)
140 if tReason(r2) != "invalid: id must be 32 bytes" {
141 t.Errorf("long id: got %q", tReason(r2))
142 }
143 }
144
145 func TestStageABadPubkeyLength(t *testing.T) {
146 s := tNewSigner(t)
147 short := tEvent(t, s, 1, time.Now().Unix(), nil, "short pubkey")
148 short.Pubkey = short.Pubkey[:16]
149 r1 := StageA(short)
150 if tReason(r1) != "invalid: pubkey must be 32 bytes" {
151 t.Errorf("short pubkey: got %q", tReason(r1))
152 }
153
154 long := tEvent(t, s, 1, time.Now().Unix(), nil, "long pubkey")
155 long.Pubkey = []byte{:33}
156 r2 := StageA(long)
157 if tReason(r2) != "invalid: pubkey must be 32 bytes" {
158 t.Errorf("long pubkey: got %q", tReason(r2))
159 }
160 }
161
162 func TestStageABadSigLength(t *testing.T) {
163 s := tNewSigner(t)
164 short := tEvent(t, s, 1, time.Now().Unix(), nil, "short sig")
165 short.Sig = short.Sig[:32]
166 r1 := StageA(short)
167 if tReason(r1) != "invalid: sig must be 64 bytes" {
168 t.Errorf("short sig: got %q", tReason(r1))
169 }
170
171 long := tEvent(t, s, 1, time.Now().Unix(), nil, "long sig")
172 long.Sig = []byte{:65}
173 r2 := StageA(long)
174 if tReason(r2) != "invalid: sig must be 64 bytes" {
175 t.Errorf("long sig: got %q", tReason(r2))
176 }
177 }
178
179 func TestStageAIDMismatch(t *testing.T) {
180 s := tNewSigner(t)
181 ev := tEvent(t, s, 1, time.Now().Unix(), nil, "hello")
182 ev.ID[0] ^= 0xFF
183 r := StageA(ev)
184 if tReason(r) != "invalid: id mismatch" {
185 t.Errorf("mismatched id: got %q", tReason(r))
186 }
187 }
188
189 func TestStageABadSignature(t *testing.T) {
190 s := tNewSigner(t)
191 ev := tEvent(t, s, 1, time.Now().Unix(), nil, "hello")
192 ev.Sig[0] ^= 0xFF
193 r := StageA(ev)
194 if tReason(r) != "invalid: bad signature" {
195 t.Errorf("bad signature: got %q", tReason(r))
196 }
197 }
198
199 // --- Stage B: limits ---
200
201 func TestIngestValid(t *testing.T) {
202 p, s, eng, dir := tSetup(t)
203 defer os.RemoveAll(dir)
204 defer eng.Close()
205
206 ev := tEvent(t, s, 1, time.Now().Unix(), nil, "hello")
207 r := p.Ingest(ev)
208 if !r.OK {
209 t.Fatalf("expected OK, got %q", tReason(r))
210 }
211 if n := tCount(t, eng); n != 1 {
212 t.Fatalf("expected 1 stored event, got %d", n)
213 }
214 }
215
216 func TestIngestDuplicate(t *testing.T) {
217 p, s, eng, dir := tSetup(t)
218 defer os.RemoveAll(dir)
219 defer eng.Close()
220
221 ev := tEvent(t, s, 1, time.Now().Unix(), nil, "hello")
222 r1 := p.Ingest(ev)
223 if !r1.OK {
224 t.Fatalf("first ingest failed: %q", tReason(r1))
225 }
226 r2 := p.Ingest(ev)
227 if r2.OK {
228 t.Fatal("expected duplicate rejection")
229 }
230 if tReason(r2) != "duplicate: already have this event" {
231 t.Fatalf("unexpected reason: %q", tReason(r2))
232 }
233 if n := tCount(t, eng); n != 1 {
234 t.Fatalf("duplicate must not add a second event, got %d", n)
235 }
236 }
237
238 func TestIngestPostVerify(t *testing.T) {
239 p, s, eng, dir := tSetup(t)
240 defer os.RemoveAll(dir)
241 defer eng.Close()
242
243 now := time.Now().Unix()
244 ev1 := tEvent(t, s, 1, now, nil, "post-verify")
245 r1 := p.IngestPostVerify(ev1)
246 if !r1.OK {
247 t.Fatalf("Stage B rejected a valid event: %q", tReason(r1))
248 }
249 if n := tCount(t, eng); n != 1 {
250 t.Fatalf("expected 1 stored event, got %d", n)
251 }
252
253 // Stage B alone still enforces the timestamp window.
254 ev2 := tEvent(t, s, 1, now+3600, nil, "too future")
255 r2 := p.IngestPostVerify(ev2)
256 if r2.OK {
257 t.Error("Stage B accepted a future timestamp")
258 }
259 if tReason(r2) != "invalid: created_at too far in future" {
260 t.Errorf("unexpected reason: %q", tReason(r2))
261 }
262 }
263
264 func TestIngestCreatedAtMissing(t *testing.T) {
265 p, s, eng, dir := tSetup(t)
266 defer os.RemoveAll(dir)
267 defer eng.Close()
268
269 ev := tEvent(t, s, 1, 0, nil, "no timestamp")
270 r := p.Ingest(ev)
271 if r.OK {
272 t.Error("created_at=0 should be rejected")
273 }
274 if tReason(r) != "invalid: created_at missing" {
275 t.Errorf("unexpected reason: %q", tReason(r))
276 }
277 if n := tCount(t, eng); n != 0 {
278 t.Errorf("rejected event must not be stored, got %d", n)
279 }
280 }
281
282 func TestIngestFutureTimestamp(t *testing.T) {
283 p, s, eng, dir := tSetup(t)
284 defer os.RemoveAll(dir)
285 defer eng.Close()
286
287 now := time.Now().Unix()
288 future := tEvent(t, s, 1, now+3600, nil, "future")
289 r1 := p.Ingest(future)
290 if r1.OK {
291 t.Error("event too far in future should be rejected")
292 }
293 if tReason(r1) != "invalid: created_at too far in future" {
294 t.Errorf("unexpected reason: %q", tReason(r1))
295 }
296
297 soon := tEvent(t, s, 1, now+100, nil, "soon")
298 r2 := p.Ingest(soon)
299 if !r2.OK {
300 t.Fatalf("event within MaxFuture should be accepted: %q", tReason(r2))
301 }
302 }
303
304 func TestIngestPastTimestamp(t *testing.T) {
305 cfg := DefaultConfig()
306 cfg.MaxPast = 60
307 p, s, eng, dir := tSetupCfg(t, cfg)
308 defer os.RemoveAll(dir)
309 defer eng.Close()
310
311 now := time.Now().Unix()
312 old := tEvent(t, s, 1, now-600, nil, "old")
313 r1 := p.Ingest(old)
314 if r1.OK {
315 t.Error("event too far in past should be rejected")
316 }
317 if tReason(r1) != "invalid: created_at too far in past" {
318 t.Errorf("unexpected reason: %q", tReason(r1))
319 }
320
321 recent := tEvent(t, s, 1, now-10, nil, "recent")
322 r2 := p.Ingest(recent)
323 if !r2.OK {
324 t.Fatalf("event within MaxPast should be accepted: %q", tReason(r2))
325 }
326 }
327
328 func TestIngestContentLimit(t *testing.T) {
329 cfg := DefaultConfig()
330 cfg.MaxContent = 4
331 p, s, eng, dir := tSetupCfg(t, cfg)
332 defer os.RemoveAll(dir)
333 defer eng.Close()
334
335 over := tEvent(t, s, 1, time.Now().Unix(), nil, "hello")
336 r1 := p.Ingest(over)
337 if r1.OK {
338 t.Error("content over MaxContent should be rejected")
339 }
340 if tReason(r1) != "invalid: content too large" {
341 t.Errorf("unexpected reason: %q", tReason(r1))
342 }
343
344 // Exactly MaxContent is allowed (the check is strictly greater).
345 exact := tEvent(t, s, 1, time.Now().Unix(), nil, "abcd")
346 r2 := p.Ingest(exact)
347 if !r2.OK {
348 t.Fatalf("content at MaxContent should be accepted: %q", tReason(r2))
349 }
350 }
351
352 func TestIngestTooManyTags(t *testing.T) {
353 cfg := DefaultConfig()
354 cfg.MaxTags = 2
355 p, s, eng, dir := tSetupCfg(t, cfg)
356 defer os.RemoveAll(dir)
357 defer eng.Close()
358
359 many := tag.NewS(tTag("t", "a"), tTag("t", "b"), tTag("t", "c"))
360 over := tEvent(t, s, 1, time.Now().Unix(), many, "x")
361 r1 := p.Ingest(over)
362 if r1.OK {
363 t.Error("more tags than MaxTags should be rejected")
364 }
365 if tReason(r1) != "invalid: too many tags" {
366 t.Errorf("unexpected reason: %q", tReason(r1))
367 }
368
369 ok := tag.NewS(tTag("t", "a"), tTag("t", "b"))
370 exact := tEvent(t, s, 1, time.Now().Unix(), ok, "x")
371 r2 := p.Ingest(exact)
372 if !r2.OK {
373 t.Fatalf("tag count at MaxTags should be accepted: %q", tReason(r2))
374 }
375 }
376
377 func TestIngestTagElementLimit(t *testing.T) {
378 cfg := DefaultConfig()
379 cfg.MaxTagElem = 4
380 p, s, eng, dir := tSetupCfg(t, cfg)
381 defer os.RemoveAll(dir)
382 defer eng.Close()
383
384 over := tEvent(t, s, 1, time.Now().Unix(), tag.NewS(tTag("t", "abcde")), "x")
385 r1 := p.Ingest(over)
386 if r1.OK {
387 t.Error("tag element over MaxTagElem should be rejected")
388 }
389 if tReason(r1) != "invalid: tag element too large" {
390 t.Errorf("unexpected reason: %q", tReason(r1))
391 }
392
393 exact := tEvent(t, s, 1, time.Now().Unix(), tag.NewS(tTag("t", "abcd")), "x")
394 r2 := p.Ingest(exact)
395 if !r2.OK {
396 t.Fatalf("tag element at MaxTagElem should be accepted: %q", tReason(r2))
397 }
398 }
399
400 // --- Stage B: special kinds ---
401
402 func TestIngestEphemeral(t *testing.T) {
403 p, s, eng, dir := tSetup(t)
404 defer os.RemoveAll(dir)
405 defer eng.Close()
406
407 // Kind 20001 is in the ephemeral range (20000-29999).
408 eph := tEvent(t, s, 20001, time.Now().Unix(), nil, "ephemeral")
409 r1 := p.Ingest(eph)
410 if !r1.OK {
411 t.Fatalf("ephemeral should be accepted: %q", tReason(r1))
412 }
413 if n1 := tCount(t, eng); n1 != 0 {
414 t.Errorf("ephemeral event must not be stored, got %d stored", n1)
415 }
416
417 regular := tEvent(t, s, 1, time.Now().Unix(), nil, "regular")
418 r2 := p.Ingest(regular)
419 if !r2.OK {
420 t.Fatalf("regular event failed: %q", tReason(r2))
421 }
422 if n2 := tCount(t, eng); n2 != 1 {
423 t.Errorf("regular event should be stored, got %d", n2)
424 }
425 }
426
427 func TestIngestReplaceable(t *testing.T) {
428 p, s, eng, dir := tSetup(t)
429 defer os.RemoveAll(dir)
430 defer eng.Close()
431
432 now := time.Now().Unix()
433 ev1 := tEvent(t, s, 0, now, nil, "old")
434 r1 := p.Ingest(ev1)
435 if !r1.OK {
436 t.Fatalf("first replaceable rejected: %q", tReason(r1))
437 }
438
439 ev2 := tEvent(t, s, 0, now+10, nil, "new")
440 r2 := p.Ingest(ev2)
441 if !r2.OK {
442 t.Fatalf("newer replaceable rejected: %q", tReason(r2))
443 }
444 if n1 := tCount(t, eng); n1 != 1 {
445 t.Fatalf("replaceable set should hold 1 event, got %d", n1)
446 }
447
448 evs, qerr := eng.QueryEvents(&filter.F{Kinds: kind.NewS(kind.New(uint16(0)))})
449 if qerr != nil {
450 t.Fatalf("query: %s", qerr.Error())
451 }
452 if len(evs) != 1 || string(evs[0].Content) != "new" {
453 t.Error("newer replaceable did not replace the older one")
454 }
455
456 ev3 := tEvent(t, s, 0, now-10, nil, "ancient")
457 r3 := p.Ingest(ev3)
458 if r3.OK {
459 t.Error("older replaceable should be rejected")
460 }
461 if tReason(r3) != "duplicate: newer version exists" {
462 t.Errorf("unexpected reason: %q", tReason(r3))
463 }
464 }
465
466 func TestIngestReplaceableTiebreak(t *testing.T) {
467 p, s, eng, dir := tSetup(t)
468 defer os.RemoveAll(dir)
469 defer eng.Close()
470
471 now := time.Now().Unix()
472 ev1 := tEvent(t, s, 0, now, nil, "first")
473 r1 := p.Ingest(ev1)
474 if !r1.OK {
475 t.Fatalf("initial ingest rejected: %q", tReason(r1))
476 }
477
478 ev2 := tEvent(t, s, 0, now, nil, "second")
479 r2 := p.Ingest(ev2)
480
481 // At equal timestamps the lexicographically lower ID wins.
482 lowerWins := bytes.Compare(ev1.ID, ev2.ID) < 0
483 if lowerWins && r2.OK {
484 t.Error("higher id at the same timestamp should be rejected")
485 }
486 if !lowerWins && !r2.OK {
487 t.Errorf("lower id at the same timestamp should replace: %q", tReason(r2))
488 }
489 if n1 := tCount(t, eng); n1 != 1 {
490 t.Fatalf("tiebreak should leave exactly 1 event, got %d", n1)
491 }
492
493 evs, qerr := eng.QueryEvents(&filter.F{Kinds: kind.NewS(kind.New(uint16(0)))})
494 if qerr != nil {
495 t.Fatalf("query: %s", qerr.Error())
496 }
497 want := "second"
498 if lowerWins {
499 want = "first"
500 }
501 if len(evs) != 1 || string(evs[0].Content) != want {
502 t.Errorf("tiebreak kept the wrong event, want %q", want)
503 }
504 }
505
506 func TestIngestParamReplaceable(t *testing.T) {
507 p, s, eng, dir := tSetup(t)
508 defer os.RemoveAll(dir)
509 defer eng.Close()
510
511 now := time.Now().Unix()
512 ev1 := tEvent(t, s, 30023, now, tag.NewS(tTag("d", "article-1")), "v1")
513 r1 := p.Ingest(ev1)
514 if !r1.OK {
515 t.Fatalf("first param replaceable rejected: %q", tReason(r1))
516 }
517
518 ev2 := tEvent(t, s, 30023, now+10, tag.NewS(tTag("d", "article-1")), "v2")
519 r2 := p.Ingest(ev2)
520 if !r2.OK {
521 t.Fatalf("newer param replaceable rejected: %q", tReason(r2))
522 }
523 if n1 := tCount(t, eng); n1 != 1 {
524 t.Fatalf("same d-tag set should hold 1 event, got %d", n1)
525 }
526
527 evs, qerr := eng.QueryEvents(&filter.F{Kinds: kind.NewS(kind.New(uint16(30023)))})
528 if qerr != nil {
529 t.Fatalf("query: %s", qerr.Error())
530 }
531 if len(evs) != 1 || string(evs[0].Content) != "v2" {
532 t.Error("newer d-tag version did not replace the older one")
533 }
534
535 ev3 := tEvent(t, s, 30023, now-10, tag.NewS(tTag("d", "article-1")), "v0")
536 r3 := p.Ingest(ev3)
537 if r3.OK {
538 t.Error("older same-d event should be rejected")
539 }
540 if tReason(r3) != "duplicate: newer version exists" {
541 t.Errorf("unexpected reason: %q", tReason(r3))
542 }
543
544 ev4 := tEvent(t, s, 30023, now, tag.NewS(tTag("d", "article-2")), "other")
545 r4 := p.Ingest(ev4)
546 if !r4.OK {
547 t.Fatalf("different d-tag should be independent: %q", tReason(r4))
548 }
549 if n2 := tCount(t, eng); n2 != 2 {
550 t.Fatalf("independent d-tag should be stored alongside, got %d", n2)
551 }
552 }
553
554 func TestIngestParamReplaceableNoDTag(t *testing.T) {
555 p, s, eng, dir := tSetup(t)
556 defer os.RemoveAll(dir)
557 defer eng.Close()
558
559 now := time.Now().Unix()
560 // Tags present, but no "d" tag: dTagValue is nil.
561 ev1 := tEvent(t, s, 30023, now, tag.NewS(tTag("x", "y")), "a")
562 r1 := p.Ingest(ev1)
563 if !r1.OK {
564 t.Fatalf("param replaceable without d-tag rejected: %q", tReason(r1))
565 }
566
567 // No tags at all: dTagValue is nil too, so this replaces the first.
568 ev2 := tEvent(t, s, 30023, now+10, nil, "b")
569 r2 := p.Ingest(ev2)
570 if !r2.OK {
571 t.Fatalf("param replaceable with nil tags rejected: %q", tReason(r2))
572 }
573 if n1 := tCount(t, eng); n1 != 1 {
574 t.Fatalf("nil d-tag events should collapse to 1, got %d", n1)
575 }
576 }
577
578 // countedRecordTags builds the smallest payload protocol.Validate accepts for
579 // each counted Musiquay kind. The storage rule under test is about keeping
580 // three records per kind, so the records have to be ones the protocol layer
581 // admits; a bare x tag is not.
582 func countedRecordTags(k uint16) (ts *tag.S) {
583 pub := "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"
584 blobHash := "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"
585 asset := "32210:" | pub | ":track-1"
586 switch k {
587 case 32218:
588 return tag.NewS(tTag("a", asset), tTag("p", pub), tTag("x", blobHash), tTag("mode", "download"), tTag("amount_sats", "1"))
589 case 32219, 32220:
590 return tag.NewS(tTag("a", asset), tTag("start", "1"), tTag("end", "2"))
591 }
592 return tag.NewS()
593 }
594
595 // Musiquay's counted kinds (32218 delivery receipt, 32219 host aggregate,
596 // 32220 publisher rollup) sit inside NIP-01's parameterized-replaceable range
597 // but are regular records: every one is stored. Before the policy hook the
598 // range rule treated them all as d="" and the second receipt silently deleted
599 // the first.
600 func TestIngestCountedKindsInReplaceableRange(t *testing.T) {
601 p, s, eng, dir := tSetup(t)
602 defer os.RemoveAll(dir)
603 defer eng.Close()
604
605 now := time.Now().Unix()
606 for _, k := range []uint16{32218, 32219, 32220} {
607 for i := int32(0); i < 3; i++ {
608 ev := tEvent(t, s, k, now+int64(i), countedRecordTags(k), "rec")
609 r := p.Ingest(ev)
610 if !r.OK {
611 t.Fatalf("kind %d record %d rejected: %q", k, i, tReason(r))
612 }
613 }
614 evs, qerr := eng.QueryEvents(&filter.F{Kinds: kind.NewS(kind.New(k))})
615 if qerr != nil {
616 t.Fatalf("query kind %d: %s", k, qerr.Error())
617 }
618 if len(evs) != 3 {
619 t.Fatalf("kind %d should keep 3 records, got %d", k, len(evs))
620 }
621 }
622
623 // An addressable kind in the same range still replaces on its d tag.
624 a1 := tEvent(t, s, 32210, now, tag.NewS(tTag("d", "track-1")), "v1")
625 if r := p.Ingest(a1); !r.OK {
626 t.Fatalf("addressable track rejected: %q", tReason(r))
627 }
628 a2 := tEvent(t, s, 32210, now+5, tag.NewS(tTag("d", "track-1")), "v2")
629 if r := p.Ingest(a2); !r.OK {
630 t.Fatalf("addressable track update rejected: %q", tReason(r))
631 }
632 aevs, aqerr := eng.QueryEvents(&filter.F{Kinds: kind.NewS(kind.New(uint16(32210)))})
633 if aqerr != nil {
634 t.Fatalf("query 32210: %s", aqerr.Error())
635 }
636 if len(aevs) != 1 || string(aevs[0].Content) != "v2" {
637 t.Fatalf("32210 should still be parameterized-replaceable, got %d events", len(aevs))
638 }
639 }
640
641 func TestIngestExpired(t *testing.T) {
642 p, s, eng, dir := tSetup(t)
643 defer os.RemoveAll(dir)
644 defer eng.Close()
645
646 now := time.Now().Unix()
647
648 // NIP-40 enforcement goes through strconv.ParseInt. ParseUint computes
649 // maxVal := uint64(1)<<uint32(bitSize) - 1, which Go defines as 0 for a
650 // 64-bit shift; before that rule was materialized maxVal read 0, every
651 // digit tripped the range check and ParseInt swallowed the unsigned range
652 // error. Assert the parse first, then that a past expiration is rejected.
653 one, parseErr := strconv.ParseInt([]byte("1"), 10, 64)
654 if parseErr != nil || one != 1 {
655 t.Fatalf("ParseInt(1,10,64) = (%d, %v), want (1, nil)", one, parseErr)
656 }
657 expired := tEvent(t, s, 1, now-60, tag.NewS(tTag("expiration", "1000000000")), "expired")
658 r1 := p.Ingest(expired)
659 if r1.OK {
660 t.Error("expired event should be rejected")
661 }
662 if tReason(r1) != "invalid: event expired" {
663 t.Errorf("unexpected reason: %q", tReason(r1))
664 }
665
666 // Future expiration: accepted.
667 future := tEvent(t, s, 1, now, tag.NewS(tTag("expiration", "9999999999")), "valid")
668 r2 := p.Ingest(future)
669 if !r2.OK {
670 t.Fatalf("future expiration should be accepted: %q", tReason(r2))
671 }
672
673 // Non-numeric expiration is ignored, not treated as expired.
674 malformed := tEvent(t, s, 1, now, tag.NewS(tTag("expiration", "not-a-number")), "ok")
675 r3 := p.Ingest(malformed)
676 if !r3.OK {
677 t.Fatalf("non-numeric expiration should be ignored: %q", tReason(r3))
678 }
679
680 // Zero expiration is ignored.
681 zero := tEvent(t, s, 1, now, tag.NewS(tTag("expiration", "0")), "ok")
682 r4 := p.Ingest(zero)
683 if !r4.OK {
684 t.Fatalf("zero expiration should be ignored: %q", tReason(r4))
685 }
686
687 // A key-only expiration tag (Len < 2) is ignored.
688 short := tEvent(t, s, 1, now, tag.NewS(tag.NewFromBytesSlice([]byte("expiration"))), "ok")
689 r5 := p.Ingest(short)
690 if !r5.OK {
691 t.Fatalf("key-only expiration should be ignored: %q", tReason(r5))
692 }
693 }
694
695 func TestIngestDeletion(t *testing.T) {
696 p, s, eng, dir := tSetup(t)
697 defer os.RemoveAll(dir)
698 defer eng.Close()
699
700 now := time.Now().Unix()
701 target := tEvent(t, s, 1, now, nil, "to be deleted")
702 r1 := p.Ingest(target)
703 if !r1.OK {
704 t.Fatalf("target ingest failed: %q", tReason(r1))
705 }
706
707 // Binary-encoded e-tag: 32 raw bytes plus the zero terminator that
708 // tag.isBinaryEncoded requires (the JSON parse path produces this shape).
709 eVal := []byte{:33}
710 copy(eVal, target.ID)
711 delTags := tag.NewS(tag.NewFromBytesSlice([]byte("e"), eVal))
712 del := tEvent(t, s, 5, now+1, delTags, "")
713 r2 := p.Ingest(del)
714 if !r2.OK {
715 t.Fatalf("deletion ingest failed: %q", tReason(r2))
716 }
717
718 if n1 := tCount(t, eng); n1 != 1 {
719 t.Fatalf("only the deletion event should remain, got %d", n1)
720 }
721 // Verify by scanning, not by filter.Ids: the by-ID path uses
722 // store.getEventSerial -> File.Get, and File.Get ignores Delete
723 // tombstones, so a filter.Ids query still returns the deleted target.
724 // That is a store defect (reported), not pipeline behaviour.
725 evs, qerr := eng.QueryEvents(&filter.F{})
726 if qerr != nil {
727 t.Fatalf("query: %s", qerr.Error())
728 }
729 for _, ev := range evs {
730 if bytes.Equal(ev.ID, target.ID) {
731 t.Error("target event was not deleted")
732 }
733 }
734 }
735
736 func TestIngestDeletionOtherAuthor(t *testing.T) {
737 p, s, eng, dir := tSetup(t)
738 defer os.RemoveAll(dir)
739 defer eng.Close()
740
741 other := tNewSigner(t)
742 now := time.Now().Unix()
743 target := tEvent(t, s, 1, now, nil, "not yours")
744 r1 := p.Ingest(target)
745 if !r1.OK {
746 t.Fatalf("target ingest failed: %q", tReason(r1))
747 }
748
749 // A deletion signed by a different key is still a valid event, but it
750 // must not remove the target.
751 eVal := []byte{:33}
752 copy(eVal, target.ID)
753 delTags := tag.NewS(tag.NewFromBytesSlice([]byte("e"), eVal))
754 del := tEvent(t, other, 5, now+1, delTags, "")
755 r2 := p.Ingest(del)
756 if !r2.OK {
757 t.Fatalf("foreign deletion should still be accepted: %q", tReason(r2))
758 }
759
760 if n1 := tCount(t, eng); n1 != 2 {
761 t.Fatalf("foreign deletion should be stored next to the target, got %d", n1)
762 }
763 // Scan-based: the by-ID path cannot be trusted after a delete (see above).
764 evs, qerr := eng.QueryEvents(&filter.F{})
765 if qerr != nil {
766 t.Fatalf("query: %s", qerr.Error())
767 }
768 found := false
769 for _, ev := range evs {
770 if bytes.Equal(ev.ID, target.ID) {
771 found = true
772 }
773 }
774 if !found {
775 t.Error("deletion must not remove another author's event")
776 }
777 }
778
779 func TestIngestDeletionNoTargets(t *testing.T) {
780 p, s, eng, dir := tSetup(t)
781 defer os.RemoveAll(dir)
782 defer eng.Close()
783
784 now := time.Now().Unix()
785
786 // A deletion with no tags is stored as-is.
787 plain := tEvent(t, s, 5, now, nil, "")
788 r1 := p.Ingest(plain)
789 if !r1.OK {
790 t.Fatalf("tagless deletion should be accepted: %q", tReason(r1))
791 }
792
793 // A well-formed e-tag for an unknown id is a no-op.
794 missing := []byte{:33}
795 unknown := tag.NewS(tag.NewFromBytesSlice([]byte("e"), missing))
796 r2 := p.Ingest(tEvent(t, s, 5, now+1, unknown, ""))
797 if !r2.OK {
798 t.Fatalf("deletion of an unknown id should be accepted: %q", tReason(r2))
799 }
800
801 // A malformed (non-32-byte) e-tag value is skipped.
802 short := tag.NewS(tag.NewFromBytesSlice([]byte("e"), []byte("nope")))
803 r3 := p.Ingest(tEvent(t, s, 5, now+2, short, ""))
804 if !r3.OK {
805 t.Fatalf("malformed e-tag deletion should be accepted: %q", tReason(r3))
806 }
807
808 if n1 := tCount(t, eng); n1 != 3 {
809 t.Fatalf("expected 3 deletion events, got %d", n1)
810 }
811 }
812
813 // --- Stage B: ACL and rate limit ---
814
815 func TestIngestACL(t *testing.T) {
816 eng, dir := tOpenStore(t)
817 defer os.RemoveAll(dir)
818 defer eng.Close()
819
820 s := tNewSigner(t)
821 other := tNewSigner(t)
822 now := time.Now().Unix()
823
824 // Matching: the writer's own key is whitelisted.
825 wl := &acl.Whitelist{Pubkeys: [][]byte{s.Pub()}}
826 p := New(eng, wl, nil, DefaultConfig())
827 r1 := p.Ingest(tEvent(t, s, 1, now, nil, "allowed"))
828 if !r1.OK {
829 t.Fatalf("whitelisted pubkey rejected: %q", tReason(r1))
830 }
831
832 // Non-matching: only another key is whitelisted.
833 wl2 := &acl.Whitelist{Pubkeys: [][]byte{other.Pub()}}
834 p.SetACL(wl2)
835 r2 := p.Ingest(tEvent(t, s, 1, now+1, nil, "blocked"))
836 if r2.OK {
837 t.Error("non-whitelisted pubkey accepted")
838 }
839 if tReason(r2) != "blocked: pubkey not allowed" {
840 t.Errorf("unexpected reason: %q", tReason(r2))
841 }
842
843 // ReadOnly rejects every write.
844 p.SetACL(&acl.ReadOnly{})
845 r3 := p.Ingest(tEvent(t, s, 1, now+2, nil, "readonly"))
846 if r3.OK {
847 t.Error("ReadOnly accepted a write")
848 }
849 }
850
851 func TestIngestRateLimit(t *testing.T) {
852 eng, dir := tOpenStore(t)
853 defer os.RemoveAll(dir)
854 defer eng.Close()
855
856 s := tNewSigner(t)
857 lim := ratelimit.New(1.0, 2)
858 p := New(eng, &acl.Open{}, lim, DefaultConfig())
859 now := time.Now().Unix()
860
861 r1 := p.Ingest(tEvent(t, s, 1, now, nil, "m1"))
862 if !r1.OK {
863 t.Fatalf("first event should pass: %q", tReason(r1))
864 }
865 r2 := p.Ingest(tEvent(t, s, 1, now+1, nil, "m2"))
866 if !r2.OK {
867 t.Fatalf("second event should pass: %q", tReason(r2))
868 }
869 r3 := p.Ingest(tEvent(t, s, 1, now+2, nil, "m3"))
870 if r3.OK {
871 t.Error("third event should be rate-limited")
872 }
873 if tReason(r3) != "rate-limited: slow down" {
874 t.Errorf("unexpected reason: %q", tReason(r3))
875 }
876 }
877
878 // --- Stage B: Musiquay payload validation ---
879
880 func TestIngestMusiquayInvalidPayload(t *testing.T) {
881 p, s, eng, dir := tSetup(t)
882 defer os.RemoveAll(dir)
883 defer eng.Close()
884
885 // music.mx's Track.Validate rejects an empty d tag; a track with only a
886 // title has no d tag and must be turned away with an invalid: reason.
887 ev := tEvent(t, s, protocol.KindTrack, time.Now().Unix(), tag.NewS(tTag("title", "No identity")), "")
888 r := p.Ingest(ev)
889 if r.OK {
890 t.Fatal("track with an empty d tag was accepted")
891 }
892 reason := tReason(r)
893 if len(reason) < 9 || reason[:9] != "invalid: " {
894 t.Fatalf("reason = %q, want an invalid: prefix", reason)
895 }
896 if n := tCount(t, eng); n != 0 {
897 t.Fatalf("rejected track was stored: %d", n)
898 }
899 }
900
901 func TestIngestMusiquayValidPayload(t *testing.T) {
902 p, s, eng, dir := tSetup(t)
903 defer os.RemoveAll(dir)
904 defer eng.Close()
905
906 ev := tEvent(t, s, protocol.KindTrack, time.Now().Unix(), tag.NewS(tTag("d", "isrc-1"), tTag("title", "Song")), "")
907 r := p.Ingest(ev)
908 if !r.OK {
909 t.Fatalf("valid track rejected: %q", tReason(r))
910 }
911 if n := tCount(t, eng); n != 1 {
912 t.Fatalf("expected 1 stored event, got %d", n)
913 }
914 }
915