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