package pipeline import ( "bytes" "os" "strconv" "testing" "time" "git.smesh.lol/morly/pkg/acl" "git.smesh.lol/nostr/pkg/event" "git.smesh.lol/nostr/pkg/filter" "git.smesh.lol/nostr/pkg/kind" "git.smesh.lol/nostr/pkg/signer/p8k" "git.smesh.lol/nostr/pkg/tag" "git.smesh.lol/morly/pkg/relay/ratelimit" "git.smesh.lol/morly/pkg/store" "git.smesh.lol/musiquay/pkg/protocol" ) // tOpenStore opens a store in a fresh directory. The caller owns both: // defer os.RemoveAll(dir) and defer eng.Close(). func tOpenStore(t *testing.T) (eng *store.Engine, dir string) { t.Helper() dir, derr := os.MkdirTemp("", "moxie-pipe-test") if derr != nil { t.Fatal(derr) } eng, oerr := store.Open(dir) if oerr != nil { t.Fatal(oerr) } return eng, dir } func tNewSigner(t *testing.T) (s *p8k.Signer) { t.Helper() s = p8k.MustNew() if gerr := s.Generate(); gerr != nil { t.Fatal(gerr) } return s } func tSetup(t *testing.T) (p *Pipeline, s *p8k.Signer, eng *store.Engine, dir string) { t.Helper() eng, dir = tOpenStore(t) s = tNewSigner(t) p = New(eng, &acl.Open{}, nil, DefaultConfig()) return p, s, eng, dir } func tSetupCfg(t *testing.T, cfg Config) (p *Pipeline, s *p8k.Signer, eng *store.Engine, dir string) { t.Helper() eng, dir = tOpenStore(t) s = tNewSigner(t) p = New(eng, &acl.Open{}, nil, cfg) return p, s, eng, dir } func tEvent(t *testing.T, s *p8k.Signer, k uint16, ts int64, tags *tag.S, content string) (ev *event.E) { t.Helper() ev = &event.E{ CreatedAt: ts, Kind: k, Tags: tags, Content: []byte(content), } if serr := ev.Sign(s); serr != nil { t.Fatal(serr) } return ev } func tTag(k, v string) (tt *tag.T) { return tag.NewFromBytesSlice([]byte(k), []byte(v)) } // tReason flattens a Result for comparison; a nil Result (accepted) is "". func tReason(r *Result) (s string) { if r == nil { return "" } return string(r.Reason) } // tCount returns the total number of events visible in the store. func tCount(t *testing.T, eng *store.Engine) (n int32) { t.Helper() evs, qerr := eng.QueryEvents(&filter.F{}) if qerr != nil { t.Fatalf("query: %s", qerr.Error()) } return int32(len(evs)) } // --- config --- func TestDefaultConfig(t *testing.T) { c := DefaultConfig() if c.MaxFuture != 900 { t.Errorf("MaxFuture = %d, want 900", c.MaxFuture) } if c.MaxPast != 0 { t.Errorf("MaxPast = %d, want 0 (unlimited)", c.MaxPast) } if c.MaxContent != 70000 { t.Errorf("MaxContent = %d, want 70000", c.MaxContent) } if c.MaxTags != 2000 { t.Errorf("MaxTags = %d, want 2000", c.MaxTags) } if c.MaxTagElem != 1024 { t.Errorf("MaxTagElem = %d, want 1024", c.MaxTagElem) } } // --- Stage A: schema + signature --- func TestStageAValid(t *testing.T) { s := tNewSigner(t) ev := tEvent(t, s, 1, time.Now().Unix(), nil, "hello") if r := StageA(ev); r != nil { t.Fatalf("valid event rejected: %s", r.Reason) } } func TestStageABadIDLength(t *testing.T) { s := tNewSigner(t) short := tEvent(t, s, 1, time.Now().Unix(), nil, "short id") short.ID = short.ID[:16] r1 := StageA(short) if tReason(r1) != "invalid: id must be 32 bytes" { t.Errorf("short id: got %q", tReason(r1)) } long := tEvent(t, s, 1, time.Now().Unix(), nil, "long id") long.ID = []byte{:33} r2 := StageA(long) if tReason(r2) != "invalid: id must be 32 bytes" { t.Errorf("long id: got %q", tReason(r2)) } } func TestStageABadPubkeyLength(t *testing.T) { s := tNewSigner(t) short := tEvent(t, s, 1, time.Now().Unix(), nil, "short pubkey") short.Pubkey = short.Pubkey[:16] r1 := StageA(short) if tReason(r1) != "invalid: pubkey must be 32 bytes" { t.Errorf("short pubkey: got %q", tReason(r1)) } long := tEvent(t, s, 1, time.Now().Unix(), nil, "long pubkey") long.Pubkey = []byte{:33} r2 := StageA(long) if tReason(r2) != "invalid: pubkey must be 32 bytes" { t.Errorf("long pubkey: got %q", tReason(r2)) } } func TestStageABadSigLength(t *testing.T) { s := tNewSigner(t) short := tEvent(t, s, 1, time.Now().Unix(), nil, "short sig") short.Sig = short.Sig[:32] r1 := StageA(short) if tReason(r1) != "invalid: sig must be 64 bytes" { t.Errorf("short sig: got %q", tReason(r1)) } long := tEvent(t, s, 1, time.Now().Unix(), nil, "long sig") long.Sig = []byte{:65} r2 := StageA(long) if tReason(r2) != "invalid: sig must be 64 bytes" { t.Errorf("long sig: got %q", tReason(r2)) } } func TestStageAIDMismatch(t *testing.T) { s := tNewSigner(t) ev := tEvent(t, s, 1, time.Now().Unix(), nil, "hello") ev.ID[0] ^= 0xFF r := StageA(ev) if tReason(r) != "invalid: id mismatch" { t.Errorf("mismatched id: got %q", tReason(r)) } } func TestStageABadSignature(t *testing.T) { s := tNewSigner(t) ev := tEvent(t, s, 1, time.Now().Unix(), nil, "hello") ev.Sig[0] ^= 0xFF r := StageA(ev) if tReason(r) != "invalid: bad signature" { t.Errorf("bad signature: got %q", tReason(r)) } } // --- Stage B: limits --- func TestIngestValid(t *testing.T) { p, s, eng, dir := tSetup(t) defer os.RemoveAll(dir) defer eng.Close() ev := tEvent(t, s, 1, time.Now().Unix(), nil, "hello") r := p.Ingest(ev) if !r.OK { t.Fatalf("expected OK, got %q", tReason(r)) } if n := tCount(t, eng); n != 1 { t.Fatalf("expected 1 stored event, got %d", n) } } func TestIngestDuplicate(t *testing.T) { p, s, eng, dir := tSetup(t) defer os.RemoveAll(dir) defer eng.Close() ev := tEvent(t, s, 1, time.Now().Unix(), nil, "hello") r1 := p.Ingest(ev) if !r1.OK { t.Fatalf("first ingest failed: %q", tReason(r1)) } r2 := p.Ingest(ev) if r2.OK { t.Fatal("expected duplicate rejection") } if tReason(r2) != "duplicate: already have this event" { t.Fatalf("unexpected reason: %q", tReason(r2)) } if n := tCount(t, eng); n != 1 { t.Fatalf("duplicate must not add a second event, got %d", n) } } func TestIngestPostVerify(t *testing.T) { p, s, eng, dir := tSetup(t) defer os.RemoveAll(dir) defer eng.Close() now := time.Now().Unix() ev1 := tEvent(t, s, 1, now, nil, "post-verify") r1 := p.IngestPostVerify(ev1) if !r1.OK { t.Fatalf("Stage B rejected a valid event: %q", tReason(r1)) } if n := tCount(t, eng); n != 1 { t.Fatalf("expected 1 stored event, got %d", n) } // Stage B alone still enforces the timestamp window. ev2 := tEvent(t, s, 1, now+3600, nil, "too future") r2 := p.IngestPostVerify(ev2) if r2.OK { t.Error("Stage B accepted a future timestamp") } if tReason(r2) != "invalid: created_at too far in future" { t.Errorf("unexpected reason: %q", tReason(r2)) } } func TestIngestCreatedAtMissing(t *testing.T) { p, s, eng, dir := tSetup(t) defer os.RemoveAll(dir) defer eng.Close() ev := tEvent(t, s, 1, 0, nil, "no timestamp") r := p.Ingest(ev) if r.OK { t.Error("created_at=0 should be rejected") } if tReason(r) != "invalid: created_at missing" { t.Errorf("unexpected reason: %q", tReason(r)) } if n := tCount(t, eng); n != 0 { t.Errorf("rejected event must not be stored, got %d", n) } } func TestIngestFutureTimestamp(t *testing.T) { p, s, eng, dir := tSetup(t) defer os.RemoveAll(dir) defer eng.Close() now := time.Now().Unix() future := tEvent(t, s, 1, now+3600, nil, "future") r1 := p.Ingest(future) if r1.OK { t.Error("event too far in future should be rejected") } if tReason(r1) != "invalid: created_at too far in future" { t.Errorf("unexpected reason: %q", tReason(r1)) } soon := tEvent(t, s, 1, now+100, nil, "soon") r2 := p.Ingest(soon) if !r2.OK { t.Fatalf("event within MaxFuture should be accepted: %q", tReason(r2)) } } func TestIngestPastTimestamp(t *testing.T) { cfg := DefaultConfig() cfg.MaxPast = 60 p, s, eng, dir := tSetupCfg(t, cfg) defer os.RemoveAll(dir) defer eng.Close() now := time.Now().Unix() old := tEvent(t, s, 1, now-600, nil, "old") r1 := p.Ingest(old) if r1.OK { t.Error("event too far in past should be rejected") } if tReason(r1) != "invalid: created_at too far in past" { t.Errorf("unexpected reason: %q", tReason(r1)) } recent := tEvent(t, s, 1, now-10, nil, "recent") r2 := p.Ingest(recent) if !r2.OK { t.Fatalf("event within MaxPast should be accepted: %q", tReason(r2)) } } func TestIngestContentLimit(t *testing.T) { cfg := DefaultConfig() cfg.MaxContent = 4 p, s, eng, dir := tSetupCfg(t, cfg) defer os.RemoveAll(dir) defer eng.Close() over := tEvent(t, s, 1, time.Now().Unix(), nil, "hello") r1 := p.Ingest(over) if r1.OK { t.Error("content over MaxContent should be rejected") } if tReason(r1) != "invalid: content too large" { t.Errorf("unexpected reason: %q", tReason(r1)) } // Exactly MaxContent is allowed (the check is strictly greater). exact := tEvent(t, s, 1, time.Now().Unix(), nil, "abcd") r2 := p.Ingest(exact) if !r2.OK { t.Fatalf("content at MaxContent should be accepted: %q", tReason(r2)) } } func TestIngestTooManyTags(t *testing.T) { cfg := DefaultConfig() cfg.MaxTags = 2 p, s, eng, dir := tSetupCfg(t, cfg) defer os.RemoveAll(dir) defer eng.Close() many := tag.NewS(tTag("t", "a"), tTag("t", "b"), tTag("t", "c")) over := tEvent(t, s, 1, time.Now().Unix(), many, "x") r1 := p.Ingest(over) if r1.OK { t.Error("more tags than MaxTags should be rejected") } if tReason(r1) != "invalid: too many tags" { t.Errorf("unexpected reason: %q", tReason(r1)) } ok := tag.NewS(tTag("t", "a"), tTag("t", "b")) exact := tEvent(t, s, 1, time.Now().Unix(), ok, "x") r2 := p.Ingest(exact) if !r2.OK { t.Fatalf("tag count at MaxTags should be accepted: %q", tReason(r2)) } } func TestIngestTagElementLimit(t *testing.T) { cfg := DefaultConfig() cfg.MaxTagElem = 4 p, s, eng, dir := tSetupCfg(t, cfg) defer os.RemoveAll(dir) defer eng.Close() over := tEvent(t, s, 1, time.Now().Unix(), tag.NewS(tTag("t", "abcde")), "x") r1 := p.Ingest(over) if r1.OK { t.Error("tag element over MaxTagElem should be rejected") } if tReason(r1) != "invalid: tag element too large" { t.Errorf("unexpected reason: %q", tReason(r1)) } exact := tEvent(t, s, 1, time.Now().Unix(), tag.NewS(tTag("t", "abcd")), "x") r2 := p.Ingest(exact) if !r2.OK { t.Fatalf("tag element at MaxTagElem should be accepted: %q", tReason(r2)) } } // --- Stage B: special kinds --- func TestIngestEphemeral(t *testing.T) { p, s, eng, dir := tSetup(t) defer os.RemoveAll(dir) defer eng.Close() // Kind 20001 is in the ephemeral range (20000-29999). eph := tEvent(t, s, 20001, time.Now().Unix(), nil, "ephemeral") r1 := p.Ingest(eph) if !r1.OK { t.Fatalf("ephemeral should be accepted: %q", tReason(r1)) } if n1 := tCount(t, eng); n1 != 0 { t.Errorf("ephemeral event must not be stored, got %d stored", n1) } regular := tEvent(t, s, 1, time.Now().Unix(), nil, "regular") r2 := p.Ingest(regular) if !r2.OK { t.Fatalf("regular event failed: %q", tReason(r2)) } if n2 := tCount(t, eng); n2 != 1 { t.Errorf("regular event should be stored, got %d", n2) } } func TestIngestReplaceable(t *testing.T) { p, s, eng, dir := tSetup(t) defer os.RemoveAll(dir) defer eng.Close() now := time.Now().Unix() ev1 := tEvent(t, s, 0, now, nil, "old") r1 := p.Ingest(ev1) if !r1.OK { t.Fatalf("first replaceable rejected: %q", tReason(r1)) } ev2 := tEvent(t, s, 0, now+10, nil, "new") r2 := p.Ingest(ev2) if !r2.OK { t.Fatalf("newer replaceable rejected: %q", tReason(r2)) } if n1 := tCount(t, eng); n1 != 1 { t.Fatalf("replaceable set should hold 1 event, got %d", n1) } evs, qerr := eng.QueryEvents(&filter.F{Kinds: kind.NewS(kind.New(uint16(0)))}) if qerr != nil { t.Fatalf("query: %s", qerr.Error()) } if len(evs) != 1 || string(evs[0].Content) != "new" { t.Error("newer replaceable did not replace the older one") } ev3 := tEvent(t, s, 0, now-10, nil, "ancient") r3 := p.Ingest(ev3) if r3.OK { t.Error("older replaceable should be rejected") } if tReason(r3) != "duplicate: newer version exists" { t.Errorf("unexpected reason: %q", tReason(r3)) } } func TestIngestReplaceableTiebreak(t *testing.T) { p, s, eng, dir := tSetup(t) defer os.RemoveAll(dir) defer eng.Close() now := time.Now().Unix() ev1 := tEvent(t, s, 0, now, nil, "first") r1 := p.Ingest(ev1) if !r1.OK { t.Fatalf("initial ingest rejected: %q", tReason(r1)) } ev2 := tEvent(t, s, 0, now, nil, "second") r2 := p.Ingest(ev2) // At equal timestamps the lexicographically lower ID wins. lowerWins := bytes.Compare(ev1.ID, ev2.ID) < 0 if lowerWins && r2.OK { t.Error("higher id at the same timestamp should be rejected") } if !lowerWins && !r2.OK { t.Errorf("lower id at the same timestamp should replace: %q", tReason(r2)) } if n1 := tCount(t, eng); n1 != 1 { t.Fatalf("tiebreak should leave exactly 1 event, got %d", n1) } evs, qerr := eng.QueryEvents(&filter.F{Kinds: kind.NewS(kind.New(uint16(0)))}) if qerr != nil { t.Fatalf("query: %s", qerr.Error()) } want := "second" if lowerWins { want = "first" } if len(evs) != 1 || string(evs[0].Content) != want { t.Errorf("tiebreak kept the wrong event, want %q", want) } } func TestIngestParamReplaceable(t *testing.T) { p, s, eng, dir := tSetup(t) defer os.RemoveAll(dir) defer eng.Close() now := time.Now().Unix() ev1 := tEvent(t, s, 30023, now, tag.NewS(tTag("d", "article-1")), "v1") r1 := p.Ingest(ev1) if !r1.OK { t.Fatalf("first param replaceable rejected: %q", tReason(r1)) } ev2 := tEvent(t, s, 30023, now+10, tag.NewS(tTag("d", "article-1")), "v2") r2 := p.Ingest(ev2) if !r2.OK { t.Fatalf("newer param replaceable rejected: %q", tReason(r2)) } if n1 := tCount(t, eng); n1 != 1 { t.Fatalf("same d-tag set should hold 1 event, got %d", n1) } evs, qerr := eng.QueryEvents(&filter.F{Kinds: kind.NewS(kind.New(uint16(30023)))}) if qerr != nil { t.Fatalf("query: %s", qerr.Error()) } if len(evs) != 1 || string(evs[0].Content) != "v2" { t.Error("newer d-tag version did not replace the older one") } ev3 := tEvent(t, s, 30023, now-10, tag.NewS(tTag("d", "article-1")), "v0") r3 := p.Ingest(ev3) if r3.OK { t.Error("older same-d event should be rejected") } if tReason(r3) != "duplicate: newer version exists" { t.Errorf("unexpected reason: %q", tReason(r3)) } ev4 := tEvent(t, s, 30023, now, tag.NewS(tTag("d", "article-2")), "other") r4 := p.Ingest(ev4) if !r4.OK { t.Fatalf("different d-tag should be independent: %q", tReason(r4)) } if n2 := tCount(t, eng); n2 != 2 { t.Fatalf("independent d-tag should be stored alongside, got %d", n2) } } func TestIngestParamReplaceableNoDTag(t *testing.T) { p, s, eng, dir := tSetup(t) defer os.RemoveAll(dir) defer eng.Close() now := time.Now().Unix() // Tags present, but no "d" tag: dTagValue is nil. ev1 := tEvent(t, s, 30023, now, tag.NewS(tTag("x", "y")), "a") r1 := p.Ingest(ev1) if !r1.OK { t.Fatalf("param replaceable without d-tag rejected: %q", tReason(r1)) } // No tags at all: dTagValue is nil too, so this replaces the first. ev2 := tEvent(t, s, 30023, now+10, nil, "b") r2 := p.Ingest(ev2) if !r2.OK { t.Fatalf("param replaceable with nil tags rejected: %q", tReason(r2)) } if n1 := tCount(t, eng); n1 != 1 { t.Fatalf("nil d-tag events should collapse to 1, got %d", n1) } } // countedRecordTags builds the smallest payload protocol.Validate accepts for // each counted Musiquay kind. The storage rule under test is about keeping // three records per kind, so the records have to be ones the protocol layer // admits; a bare x tag is not. func countedRecordTags(k uint16) (ts *tag.S) { pub := "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" blobHash := "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb" asset := "32210:" | pub | ":track-1" switch k { case 32218: return tag.NewS(tTag("a", asset), tTag("p", pub), tTag("x", blobHash), tTag("mode", "download"), tTag("amount_sats", "1")) case 32219, 32220: return tag.NewS(tTag("a", asset), tTag("start", "1"), tTag("end", "2")) } return tag.NewS() } // Musiquay's counted kinds (32218 delivery receipt, 32219 host aggregate, // 32220 publisher rollup) sit inside NIP-01's parameterized-replaceable range // but are regular records: every one is stored. Before the policy hook the // range rule treated them all as d="" and the second receipt silently deleted // the first. func TestIngestCountedKindsInReplaceableRange(t *testing.T) { p, s, eng, dir := tSetup(t) defer os.RemoveAll(dir) defer eng.Close() now := time.Now().Unix() for _, k := range []uint16{32218, 32219, 32220} { for i := int32(0); i < 3; i++ { ev := tEvent(t, s, k, now+int64(i), countedRecordTags(k), "rec") r := p.Ingest(ev) if !r.OK { t.Fatalf("kind %d record %d rejected: %q", k, i, tReason(r)) } } evs, qerr := eng.QueryEvents(&filter.F{Kinds: kind.NewS(kind.New(k))}) if qerr != nil { t.Fatalf("query kind %d: %s", k, qerr.Error()) } if len(evs) != 3 { t.Fatalf("kind %d should keep 3 records, got %d", k, len(evs)) } } // An addressable kind in the same range still replaces on its d tag. a1 := tEvent(t, s, 32210, now, tag.NewS(tTag("d", "track-1")), "v1") if r := p.Ingest(a1); !r.OK { t.Fatalf("addressable track rejected: %q", tReason(r)) } a2 := tEvent(t, s, 32210, now+5, tag.NewS(tTag("d", "track-1")), "v2") if r := p.Ingest(a2); !r.OK { t.Fatalf("addressable track update rejected: %q", tReason(r)) } aevs, aqerr := eng.QueryEvents(&filter.F{Kinds: kind.NewS(kind.New(uint16(32210)))}) if aqerr != nil { t.Fatalf("query 32210: %s", aqerr.Error()) } if len(aevs) != 1 || string(aevs[0].Content) != "v2" { t.Fatalf("32210 should still be parameterized-replaceable, got %d events", len(aevs)) } } func TestIngestExpired(t *testing.T) { p, s, eng, dir := tSetup(t) defer os.RemoveAll(dir) defer eng.Close() now := time.Now().Unix() // NIP-40 enforcement goes through strconv.ParseInt. ParseUint computes // maxVal := uint64(1)< File.Get, and File.Get ignores Delete // tombstones, so a filter.Ids query still returns the deleted target. // That is a store defect (reported), not pipeline behaviour. evs, qerr := eng.QueryEvents(&filter.F{}) if qerr != nil { t.Fatalf("query: %s", qerr.Error()) } for _, ev := range evs { if bytes.Equal(ev.ID, target.ID) { t.Error("target event was not deleted") } } } func TestIngestDeletionOtherAuthor(t *testing.T) { p, s, eng, dir := tSetup(t) defer os.RemoveAll(dir) defer eng.Close() other := tNewSigner(t) now := time.Now().Unix() target := tEvent(t, s, 1, now, nil, "not yours") r1 := p.Ingest(target) if !r1.OK { t.Fatalf("target ingest failed: %q", tReason(r1)) } // A deletion signed by a different key is still a valid event, but it // must not remove the target. eVal := []byte{:33} copy(eVal, target.ID) delTags := tag.NewS(tag.NewFromBytesSlice([]byte("e"), eVal)) del := tEvent(t, other, 5, now+1, delTags, "") r2 := p.Ingest(del) if !r2.OK { t.Fatalf("foreign deletion should still be accepted: %q", tReason(r2)) } if n1 := tCount(t, eng); n1 != 2 { t.Fatalf("foreign deletion should be stored next to the target, got %d", n1) } // Scan-based: the by-ID path cannot be trusted after a delete (see above). evs, qerr := eng.QueryEvents(&filter.F{}) if qerr != nil { t.Fatalf("query: %s", qerr.Error()) } found := false for _, ev := range evs { if bytes.Equal(ev.ID, target.ID) { found = true } } if !found { t.Error("deletion must not remove another author's event") } } func TestIngestDeletionNoTargets(t *testing.T) { p, s, eng, dir := tSetup(t) defer os.RemoveAll(dir) defer eng.Close() now := time.Now().Unix() // A deletion with no tags is stored as-is. plain := tEvent(t, s, 5, now, nil, "") r1 := p.Ingest(plain) if !r1.OK { t.Fatalf("tagless deletion should be accepted: %q", tReason(r1)) } // A well-formed e-tag for an unknown id is a no-op. missing := []byte{:33} unknown := tag.NewS(tag.NewFromBytesSlice([]byte("e"), missing)) r2 := p.Ingest(tEvent(t, s, 5, now+1, unknown, "")) if !r2.OK { t.Fatalf("deletion of an unknown id should be accepted: %q", tReason(r2)) } // A malformed (non-32-byte) e-tag value is skipped. short := tag.NewS(tag.NewFromBytesSlice([]byte("e"), []byte("nope"))) r3 := p.Ingest(tEvent(t, s, 5, now+2, short, "")) if !r3.OK { t.Fatalf("malformed e-tag deletion should be accepted: %q", tReason(r3)) } if n1 := tCount(t, eng); n1 != 3 { t.Fatalf("expected 3 deletion events, got %d", n1) } } // --- Stage B: ACL and rate limit --- func TestIngestACL(t *testing.T) { eng, dir := tOpenStore(t) defer os.RemoveAll(dir) defer eng.Close() s := tNewSigner(t) other := tNewSigner(t) now := time.Now().Unix() // Matching: the writer's own key is whitelisted. wl := &acl.Whitelist{Pubkeys: [][]byte{s.Pub()}} p := New(eng, wl, nil, DefaultConfig()) r1 := p.Ingest(tEvent(t, s, 1, now, nil, "allowed")) if !r1.OK { t.Fatalf("whitelisted pubkey rejected: %q", tReason(r1)) } // Non-matching: only another key is whitelisted. wl2 := &acl.Whitelist{Pubkeys: [][]byte{other.Pub()}} p.SetACL(wl2) r2 := p.Ingest(tEvent(t, s, 1, now+1, nil, "blocked")) if r2.OK { t.Error("non-whitelisted pubkey accepted") } if tReason(r2) != "blocked: pubkey not allowed" { t.Errorf("unexpected reason: %q", tReason(r2)) } // ReadOnly rejects every write. p.SetACL(&acl.ReadOnly{}) r3 := p.Ingest(tEvent(t, s, 1, now+2, nil, "readonly")) if r3.OK { t.Error("ReadOnly accepted a write") } } func TestIngestRateLimit(t *testing.T) { eng, dir := tOpenStore(t) defer os.RemoveAll(dir) defer eng.Close() s := tNewSigner(t) lim := ratelimit.New(1.0, 2) p := New(eng, &acl.Open{}, lim, DefaultConfig()) now := time.Now().Unix() r1 := p.Ingest(tEvent(t, s, 1, now, nil, "m1")) if !r1.OK { t.Fatalf("first event should pass: %q", tReason(r1)) } r2 := p.Ingest(tEvent(t, s, 1, now+1, nil, "m2")) if !r2.OK { t.Fatalf("second event should pass: %q", tReason(r2)) } r3 := p.Ingest(tEvent(t, s, 1, now+2, nil, "m3")) if r3.OK { t.Error("third event should be rate-limited") } if tReason(r3) != "rate-limited: slow down" { t.Errorf("unexpected reason: %q", tReason(r3)) } } // --- Stage B: Musiquay payload validation --- func TestIngestMusiquayInvalidPayload(t *testing.T) { p, s, eng, dir := tSetup(t) defer os.RemoveAll(dir) defer eng.Close() // music.mx's Track.Validate rejects an empty d tag; a track with only a // title has no d tag and must be turned away with an invalid: reason. ev := tEvent(t, s, protocol.KindTrack, time.Now().Unix(), tag.NewS(tTag("title", "No identity")), "") r := p.Ingest(ev) if r.OK { t.Fatal("track with an empty d tag was accepted") } reason := tReason(r) if len(reason) < 9 || reason[:9] != "invalid: " { t.Fatalf("reason = %q, want an invalid: prefix", reason) } if n := tCount(t, eng); n != 0 { t.Fatalf("rejected track was stored: %d", n) } } func TestIngestMusiquayValidPayload(t *testing.T) { p, s, eng, dir := tSetup(t) defer os.RemoveAll(dir) defer eng.Close() ev := tEvent(t, s, protocol.KindTrack, time.Now().Unix(), tag.NewS(tTag("d", "isrc-1"), tTag("title", "Song")), "") r := p.Ingest(ev) if !r.OK { t.Fatalf("valid track rejected: %q", tReason(r)) } if n := tCount(t, eng); n != 1 { t.Fatalf("expected 1 stored event, got %d", n) } }