// Package wire tests exercise the channel-IPC codec directly: every type // round-trips through EncodeTo/DecodeFrom, the fixed headers are checked byte // for byte, and malformed frames (short reads, length mismatches) must error. // The pure workers are driven in-process over buffered channels so their loop // bodies run without a spawn and still land in the coverage counters. package wire import ( "bytes" "os" "testing" "git.smesh.lol/smesh/pkg/nostr/envelope" "git.smesh.lol/smesh/pkg/nostr/event" "git.smesh.lol/smesh/pkg/nostr/filter" "git.smesh.lol/smesh/pkg/nostr/kind" "git.smesh.lol/smesh/pkg/nostr/signer/p8k" ) func wFill(n int32, fill byte) (b []byte) { b = []byte{:n} for i := 0; i < n; i++ { b[i] = fill } return } func wSigned(t *testing.T, k uint16) (ev *event.E) { s := p8k.MustNew() if gerr := s.Generate(); gerr != nil { t.Fatal(gerr) } ev = &event.E{ CreatedAt: 1700000000, Kind: k, Content: []byte("hello"), } if serr := ev.Sign(s); serr != nil { t.Fatal(serr) } return ev } func wEventJSON(ev *event.E) (b []byte) { es := &envelope.EventSubmission{} es.E = ev return es.Marshal(nil) } func wErr(t *testing.T, name string, err error) { if err == nil { t.Fatalf("%s: expected an error", name) } } // len/cap on a channel report 0 in this compiler, so buffered results are // probed with a non-blocking select instead. func wTryFrame(ch chan BroadcastFrame) (fr BroadcastFrame, ok bool) { select { case v := <-ch: fr = v ok = true default: } return } func wTryIngest(ch chan IngestResponse) (r IngestResponse, ok bool) { select { case v := <-ch: r = v ok = true default: } return } func wTryProxy(ch chan ProxyResponse) (r ProxyResponse, ok bool) { select { case v := <-ch: r = v ok = true default: } return } func wTryBlossom(ch chan BlossomResponse) (r BlossomResponse, ok bool) { select { case v := <-ch: r = v ok = true default: } return } func wTryReady(ch chan struct{}) (ok bool) { select { case <-ch: ok = true default: } return } // --- constants --- func TestVerdictAndSubOpConstants(t *testing.T) { if VerdictReject != 0 || VerdictAccept != 1 || VerdictEphemeral != 2 { t.Fatal("verdict constants changed") } if SubOpNew != 1 || SubOpAdd != 2 || SubOpRemove != 3 || SubOpClose != 4 || SubOpAuth != 5 || SubOpBcast != 10 { t.Fatal("SubOp constants changed") } } // --- IngestRequest --- func TestIngestRequestRoundTrip(t *testing.T) { r := IngestRequest{ReqID: 0x01020304, Bytes: []byte("abc")} buf := bytes.NewBuffer(nil) if err := r.EncodeTo(buf); err != nil { t.Fatal(err) } want := []byte{4, 3, 2, 1, 3, 0, 0, 0, 'a', 'b', 'c'} if !bytes.Equal(buf.Bytes(), want) { t.Fatalf("IngestRequest layout = %x", buf.Bytes()) } var got IngestRequest if err := got.DecodeFrom(buf); err != nil { t.Fatal(err) } if got.ReqID != r.ReqID || !bytes.Equal(got.Bytes, r.Bytes) { t.Fatal("IngestRequest round trip") } if buf.Len() != 0 { t.Fatal("IngestRequest decode left bytes") } } func TestIngestRequestEmptyAndShort(t *testing.T) { r := IngestRequest{ReqID: 7} buf := bytes.NewBuffer(nil) if err := r.EncodeTo(buf); err != nil { t.Fatal(err) } if buf.Len() != 8 { t.Fatalf("empty IngestRequest length = %d", buf.Len()) } var got IngestRequest if err := got.DecodeFrom(buf); err != nil { t.Fatal(err) } if got.ReqID != 7 || got.Bytes != nil { t.Fatal("empty IngestRequest round trip") } var empty IngestRequest wErr(t, "IngestRequest empty input", empty.DecodeFrom(bytes.NewBuffer(nil))) var short IngestRequest wErr(t, "IngestRequest short header", short.DecodeFrom(bytes.NewBuffer([]byte{1, 0, 0, 0}))) var mismatch IngestRequest wErr(t, "IngestRequest length mismatch", mismatch.DecodeFrom(bytes.NewBuffer([]byte{1, 0, 0, 0, 5, 0, 0, 0, 'a'}))) } // --- IngestResponse --- func TestIngestResponseRoundTrip(t *testing.T) { var r IngestResponse r.ReqID = 0x0A0B0C0D r.Verdict = VerdictAccept r.Kind = 0x0102 r.CreatedAt = 1700000000 copy(r.Pubkey[:], wFill(32, 0x11)) copy(r.EventID[:], wFill(32, 0x22)) r.Bytes = []byte("ev") r.Reason = []byte("rj") buf := bytes.NewBuffer(nil) if err := r.EncodeTo(buf); err != nil { t.Fatal(err) } if buf.Len() != 92 { t.Fatalf("IngestResponse length = %d, want 92", buf.Len()) } h := buf.Bytes() if h[0] != 0x0D || h[1] != 0x0C || h[2] != 0x0B || h[3] != 0x0A { t.Fatal("IngestResponse ReqID not little-endian") } if h[4] != VerdictAccept { t.Fatal("IngestResponse verdict byte") } if h[5] != 0x02 || h[6] != 0x01 { t.Fatal("IngestResponse kind not little-endian") } created := []byte{0, 241, 83, 101, 0, 0, 0, 0} if !bytes.Equal(h[7:15], created) { t.Fatalf("IngestResponse created_at = %x", h[7:15]) } if h[15] != 0x11 || h[46] != 0x11 { t.Fatal("IngestResponse pubkey span") } if h[47] != 0x22 || h[78] != 0x22 { t.Fatal("IngestResponse event id span") } if h[79] != 0 { t.Fatal("IngestResponse reserved byte must be zero") } if !bytes.Equal(h[80:84], []byte{2, 0, 0, 0}) || !bytes.Equal(h[84:88], []byte{2, 0, 0, 0}) { t.Fatal("IngestResponse lengths") } if !bytes.Equal(h[88:90], []byte("ev")) || !bytes.Equal(h[90:92], []byte("rj")) { t.Fatal("IngestResponse payload order") } var got IngestResponse if err := got.DecodeFrom(buf); err != nil { t.Fatal(err) } if got.ReqID != r.ReqID || got.Verdict != r.Verdict || got.Kind != r.Kind || got.CreatedAt != r.CreatedAt { t.Fatal("IngestResponse scalar round trip") } if !bytes.Equal(got.Pubkey[:], r.Pubkey[:]) || !bytes.Equal(got.EventID[:], r.EventID[:]) { t.Fatal("IngestResponse key round trip") } if !bytes.Equal(got.Bytes, r.Bytes) || !bytes.Equal(got.Reason, r.Reason) { t.Fatal("IngestResponse payload round trip") } } func TestIngestResponseEmptyAndReasonOnly(t *testing.T) { var r IngestResponse r.ReqID = 3 r.Verdict = VerdictReject buf := bytes.NewBuffer(nil) if err := r.EncodeTo(buf); err != nil { t.Fatal(err) } if buf.Len() != 88 { t.Fatalf("empty IngestResponse length = %d", buf.Len()) } var got IngestResponse if err := got.DecodeFrom(buf); err != nil { t.Fatal(err) } if got.Bytes != nil || got.Reason != nil { t.Fatal("empty IngestResponse payloads must be nil") } var r2 IngestResponse r2.ReqID = 4 r2.Reason = []byte("no") buf2 := bytes.NewBuffer(nil) if err := r2.EncodeTo(buf2); err != nil { t.Fatal(err) } var got2 IngestResponse if err := got2.DecodeFrom(buf2); err != nil { t.Fatal(err) } if got2.Bytes != nil || !bytes.Equal(got2.Reason, []byte("no")) { t.Fatal("reason-only IngestResponse") } if got2.Verdict != VerdictReject || got2.Kind != 0 { t.Fatal("reason-only IngestResponse defaults") } } func TestIngestResponseShort(t *testing.T) { var r IngestResponse r.ReqID = 1 r.Bytes = []byte("xy") full := bytes.NewBuffer(nil) if err := r.EncodeTo(full); err != nil { t.Fatal(err) } // Keep the 80-byte header and 8-byte length block, drop the payload. full.Truncate(88) var d IngestResponse wErr(t, "IngestResponse missing payload", d.DecodeFrom(full)) var d2 IngestResponse wErr(t, "IngestResponse short header", d2.DecodeFrom(bytes.NewBuffer(wFill(40, 0)))) var d3 IngestResponse wErr(t, "IngestResponse short lengths", d3.DecodeFrom(bytes.NewBuffer(wFill(84, 0)))) } // --- SubCommand --- func TestSubCommandRoundTrip(t *testing.T) { c := SubCommand{Op: SubOpAdd, ConnFD: 7, Flags: 3, SubID: []byte("sub"), Bytes: []byte("payload")} buf := bytes.NewBuffer(nil) if err := c.EncodeTo(buf); err != nil { t.Fatal(err) } want := []byte{7, 0, 0, 0, 2, 3, 3, 0, 0, 0, 7, 0, 0, 0, 's', 'u', 'b', 'p', 'a', 'y', 'l', 'o', 'a', 'd'} if !bytes.Equal(buf.Bytes(), want) { t.Fatalf("SubCommand layout = %x", buf.Bytes()) } var got SubCommand if err := got.DecodeFrom(buf); err != nil { t.Fatal(err) } if got.Op != c.Op || got.ConnFD != c.ConnFD || got.Flags != c.Flags { t.Fatal("SubCommand scalars") } if !bytes.Equal(got.SubID, c.SubID) || !bytes.Equal(got.Bytes, c.Bytes) { t.Fatal("SubCommand payloads") } } func TestSubCommandEmptyNegativeAndShort(t *testing.T) { var fd int32 = -1 c := SubCommand{Op: SubOpClose, ConnFD: fd} buf := bytes.NewBuffer(nil) if err := c.EncodeTo(buf); err != nil { t.Fatal(err) } if buf.Len() != 14 { t.Fatalf("empty SubCommand length = %d", buf.Len()) } if !bytes.Equal(buf.Bytes()[:4], []byte{255, 255, 255, 255}) { t.Fatal("negative ConnFD encoding") } var got SubCommand if err := got.DecodeFrom(buf); err != nil { t.Fatal(err) } if got.ConnFD != -1 || got.SubID != nil || got.Bytes != nil { t.Fatal("empty SubCommand round trip") } var d SubCommand wErr(t, "SubCommand short header", d.DecodeFrom(bytes.NewBuffer(wFill(13, 0)))) var d2 SubCommand wErr(t, "SubCommand SubID mismatch", d2.DecodeFrom(bytes.NewBuffer([]byte{1, 0, 0, 0, 2, 0, 5, 0, 0, 0, 0, 0, 0, 0, 'a'}))) var d3 SubCommand wErr(t, "SubCommand Bytes mismatch", d3.DecodeFrom(bytes.NewBuffer([]byte{1, 0, 0, 0, 2, 0, 0, 0, 0, 0, 5, 0, 0, 0, 'a'}))) } // --- BroadcastRequest --- func TestBroadcastRequestRoundTrip(t *testing.T) { r := BroadcastRequest{SenderFD: 42, Flags: 5, Bytes: []byte("event")} buf := bytes.NewBuffer(nil) if err := r.EncodeTo(buf); err != nil { t.Fatal(err) } want := []byte{42, 0, 0, 0, 5, 5, 0, 0, 0, 'e', 'v', 'e', 'n', 't'} if !bytes.Equal(buf.Bytes(), want) { t.Fatalf("BroadcastRequest layout = %x", buf.Bytes()) } var got BroadcastRequest if err := got.DecodeFrom(buf); err != nil { t.Fatal(err) } if got.SenderFD != r.SenderFD || got.Flags != r.Flags || !bytes.Equal(got.Bytes, r.Bytes) { t.Fatal("BroadcastRequest round trip") } var fd int32 = -9 r2 := BroadcastRequest{SenderFD: fd} buf2 := bytes.NewBuffer(nil) if err := r2.EncodeTo(buf2); err != nil { t.Fatal(err) } var got2 BroadcastRequest if err := got2.DecodeFrom(buf2); err != nil { t.Fatal(err) } if got2.SenderFD != -9 || got2.Bytes != nil { t.Fatal("empty negative BroadcastRequest") } var d BroadcastRequest wErr(t, "BroadcastRequest short header", d.DecodeFrom(bytes.NewBuffer(wFill(8, 0)))) var d2 BroadcastRequest wErr(t, "BroadcastRequest length mismatch", d2.DecodeFrom(bytes.NewBuffer([]byte{1, 0, 0, 0, 0, 4, 0, 0, 0, 'a'}))) } // --- BroadcastFrame --- func TestBroadcastFrameRoundTrip(t *testing.T) { f := BroadcastFrame{ConnFD: 3, Bytes: []byte("frame")} buf := bytes.NewBuffer(nil) if err := f.EncodeTo(buf); err != nil { t.Fatal(err) } want := []byte{3, 0, 0, 0, 5, 0, 0, 0, 'f', 'r', 'a', 'm', 'e'} if !bytes.Equal(buf.Bytes(), want) { t.Fatalf("BroadcastFrame layout = %x", buf.Bytes()) } var got BroadcastFrame if err := got.DecodeFrom(buf); err != nil { t.Fatal(err) } if got.ConnFD != f.ConnFD || !bytes.Equal(got.Bytes, f.Bytes) { t.Fatal("BroadcastFrame round trip") } var fd int32 = -4 f2 := BroadcastFrame{ConnFD: fd} buf2 := bytes.NewBuffer(nil) if err := f2.EncodeTo(buf2); err != nil { t.Fatal(err) } var got2 BroadcastFrame if err := got2.DecodeFrom(buf2); err != nil { t.Fatal(err) } if got2.ConnFD != -4 || got2.Bytes != nil { t.Fatal("empty negative BroadcastFrame") } var d BroadcastFrame wErr(t, "BroadcastFrame short header", d.DecodeFrom(bytes.NewBuffer(wFill(7, 0)))) var d2 BroadcastFrame wErr(t, "BroadcastFrame length mismatch", d2.DecodeFrom(bytes.NewBuffer([]byte{1, 0, 0, 0, 4, 0, 0, 0, 'a'}))) } // --- ProxyRequest --- func TestProxyRequestRoundTrip(t *testing.T) { r := ProxyRequest{ReqID: 9, MaxBytes: 1024, URL: []byte("http://x")} buf := bytes.NewBuffer(nil) if err := r.EncodeTo(buf); err != nil { t.Fatal(err) } want := []byte{9, 0, 0, 0, 0, 4, 0, 0, 8, 0, 0, 0, 'h', 't', 't', 'p', ':', '/', '/', 'x'} if !bytes.Equal(buf.Bytes(), want) { t.Fatalf("ProxyRequest layout = %x", buf.Bytes()) } var got ProxyRequest if err := got.DecodeFrom(buf); err != nil { t.Fatal(err) } if got.ReqID != r.ReqID || got.MaxBytes != r.MaxBytes || !bytes.Equal(got.URL, r.URL) { t.Fatal("ProxyRequest round trip") } r2 := ProxyRequest{ReqID: 1} buf2 := bytes.NewBuffer(nil) if err := r2.EncodeTo(buf2); err != nil { t.Fatal(err) } if buf2.Len() != 12 { t.Fatalf("empty ProxyRequest length = %d", buf2.Len()) } var got2 ProxyRequest if err := got2.DecodeFrom(buf2); err != nil { t.Fatal(err) } if got2.URL != nil || got2.MaxBytes != 0 { t.Fatal("empty ProxyRequest round trip") } var d ProxyRequest wErr(t, "ProxyRequest short header", d.DecodeFrom(bytes.NewBuffer(wFill(11, 0)))) var d2 ProxyRequest wErr(t, "ProxyRequest URL mismatch", d2.DecodeFrom(bytes.NewBuffer([]byte{1, 0, 0, 0, 0, 0, 0, 0, 5, 0, 0, 0, 'a'}))) } // --- ProxyResponse --- func TestProxyResponseRoundTrip(t *testing.T) { r := ProxyResponse{ReqID: 1, Status: 200, ContentType: []byte("image/png"), Body: []byte("BODY")} buf := bytes.NewBuffer(nil) if err := r.EncodeTo(buf); err != nil { t.Fatal(err) } h := buf.Bytes() if !bytes.Equal(h[:20], []byte{1, 0, 0, 0, 200, 0, 0, 0, 9, 0, 0, 0, 4, 0, 0, 0, 0, 0, 0, 0}) { t.Fatalf("ProxyResponse header = %x", h[:20]) } if !bytes.Equal(h[20:29], []byte("image/png")) || !bytes.Equal(h[29:33], []byte("BODY")) { t.Fatal("ProxyResponse payload order") } var got ProxyResponse if err := got.DecodeFrom(buf); err != nil { t.Fatal(err) } if got.ReqID != 1 || got.Status != 200 { t.Fatal("ProxyResponse scalars") } if !bytes.Equal(got.ContentType, r.ContentType) || !bytes.Equal(got.Body, r.Body) || got.Err != nil { t.Fatal("ProxyResponse payloads") } var st int32 = -1 r2 := ProxyResponse{ReqID: 2, Status: st, Err: []byte("boom")} buf2 := bytes.NewBuffer(nil) if err := r2.EncodeTo(buf2); err != nil { t.Fatal(err) } if !bytes.Equal(buf2.Bytes()[4:8], []byte{255, 255, 255, 255}) { t.Fatal("negative ProxyResponse status encoding") } var got2 ProxyResponse if err := got2.DecodeFrom(buf2); err != nil { t.Fatal(err) } if got2.Status != -1 || !bytes.Equal(got2.Err, []byte("boom")) || got2.ContentType != nil || got2.Body != nil { t.Fatal("negative ProxyResponse round trip") } r3 := ProxyResponse{ReqID: 3} buf3 := bytes.NewBuffer(nil) if err := r3.EncodeTo(buf3); err != nil { t.Fatal(err) } if buf3.Len() != 20 { t.Fatalf("empty ProxyResponse length = %d", buf3.Len()) } var got3 ProxyResponse if err := got3.DecodeFrom(buf3); err != nil { t.Fatal(err) } if got3.ContentType != nil || got3.Body != nil || got3.Err != nil { t.Fatal("empty ProxyResponse payloads must be nil") } var d ProxyResponse wErr(t, "ProxyResponse short header", d.DecodeFrom(bytes.NewBuffer(wFill(19, 0)))) var d2 ProxyResponse wErr(t, "ProxyResponse CT mismatch", d2.DecodeFrom(bytes.NewBuffer([]byte{1, 0, 0, 0, 0, 0, 0, 0, 5, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 'a'}))) var d3 ProxyResponse wErr(t, "ProxyResponse body mismatch", d3.DecodeFrom(bytes.NewBuffer([]byte{1, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 5, 0, 0, 0, 0, 0, 0, 0, 'a'}))) var d4 ProxyResponse wErr(t, "ProxyResponse err mismatch", d4.DecodeFrom(bytes.NewBuffer([]byte{1, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 5, 0, 0, 0, 'a'}))) } // --- BlossomRequest --- func TestBlossomRequestRoundTrip(t *testing.T) { r := BlossomRequest{ ReqID: 1, Dir: []byte("d"), Method: []byte("GET"), Path: []byte("/p"), ContentType: []byte("c"), Body: []byte("b"), Upstream: []byte("u"), } buf := bytes.NewBuffer(nil) if err := r.EncodeTo(buf); err != nil { t.Fatal(err) } hdr := []byte{1, 0, 0, 0, 1, 0, 0, 0, 3, 0, 0, 0, 2, 0, 0, 0, 1, 0, 0, 0, 1, 0, 0, 0, 1, 0, 0, 0} if !bytes.Equal(buf.Bytes()[:28], hdr) { t.Fatalf("BlossomRequest header = %x", buf.Bytes()[:28]) } if !bytes.Equal(buf.Bytes()[28:], []byte("dGET/pcbu")) { t.Fatalf("BlossomRequest payload = %s", buf.Bytes()[28:]) } var got BlossomRequest if err := got.DecodeFrom(buf); err != nil { t.Fatal(err) } if got.ReqID != 1 { t.Fatal("BlossomRequest ReqID") } if !bytes.Equal(got.Dir, r.Dir) || !bytes.Equal(got.Method, r.Method) || !bytes.Equal(got.Path, r.Path) || !bytes.Equal(got.ContentType, r.ContentType) || !bytes.Equal(got.Body, r.Body) || !bytes.Equal(got.Upstream, r.Upstream) { t.Fatal("BlossomRequest round trip") } r2 := BlossomRequest{ReqID: 2} buf2 := bytes.NewBuffer(nil) if err := r2.EncodeTo(buf2); err != nil { t.Fatal(err) } if buf2.Len() != 28 { t.Fatalf("empty BlossomRequest length = %d", buf2.Len()) } var got2 BlossomRequest if err := got2.DecodeFrom(buf2); err != nil { t.Fatal(err) } if got2.Dir != nil || got2.Method != nil || got2.Path != nil || got2.ContentType != nil || got2.Body != nil || got2.Upstream != nil { t.Fatal("empty BlossomRequest fields must be nil") } var d BlossomRequest wErr(t, "BlossomRequest short header", d.DecodeFrom(bytes.NewBuffer(wFill(27, 0)))) var d2 BlossomRequest short := []byte{1, 0, 0, 0, 4, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 'a'} wErr(t, "BlossomRequest Dir mismatch", d2.DecodeFrom(bytes.NewBuffer(short))) } func TestReadChunk(t *testing.T) { b, err := readChunk(bytes.NewBuffer([]byte("abc")), 3) if err != nil || !bytes.Equal(b, []byte("abc")) { t.Fatal("readChunk full") } b2, err2 := readChunk(bytes.NewBuffer(nil), 0) if err2 != nil || b2 != nil { t.Fatal("readChunk zero") } b3, err3 := readChunk(bytes.NewBuffer([]byte("ab")), 5) if err3 == nil || len(b3) != 5 { t.Fatal("readChunk short must error with an n-length buffer") } } // --- BlossomResponse --- func TestBlossomResponseRoundTrip(t *testing.T) { var sz int64 = 9000000000 r := BlossomResponse{ReqID: 2, Status: 206, Size: sz, CT: []byte("ct"), Body: []byte("body")} buf := bytes.NewBuffer(nil) if err := r.EncodeTo(buf); err != nil { t.Fatal(err) } h := buf.Bytes() wantSize := []byte{0, 26, 113, 24, 2, 0, 0, 0} if !bytes.Equal(h[8:16], wantSize) { t.Fatalf("BlossomResponse size = %x", h[8:16]) } if !bytes.Equal(h[16:20], []byte{2, 0, 0, 0}) || !bytes.Equal(h[20:24], []byte{4, 0, 0, 0}) { t.Fatal("BlossomResponse lengths") } var got BlossomResponse if err := got.DecodeFrom(buf); err != nil { t.Fatal(err) } if got.ReqID != 2 || got.Status != 206 || got.Size != 9000000000 { t.Fatal("BlossomResponse scalars") } if !bytes.Equal(got.CT, r.CT) || !bytes.Equal(got.Body, r.Body) { t.Fatal("BlossomResponse payloads") } var st int32 = -1 r2 := BlossomResponse{ReqID: 3, Status: st} buf2 := bytes.NewBuffer(nil) if err := r2.EncodeTo(buf2); err != nil { t.Fatal(err) } if !bytes.Equal(buf2.Bytes()[4:8], []byte{255, 255, 255, 255}) { t.Fatal("negative BlossomResponse status encoding") } var got2 BlossomResponse if err := got2.DecodeFrom(buf2); err != nil { t.Fatal(err) } if got2.Status != -1 || got2.CT != nil || got2.Body != nil { t.Fatal("negative BlossomResponse round trip") } var d BlossomResponse wErr(t, "BlossomResponse short header", d.DecodeFrom(bytes.NewBuffer(wFill(23, 0)))) var d2 BlossomResponse wErr(t, "BlossomResponse CT mismatch", d2.DecodeFrom(bytes.NewBuffer([]byte{1, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 5, 0, 0, 0, 0, 0, 0, 0, 'a'}))) var d3 BlossomResponse wErr(t, "BlossomResponse body mismatch", d3.DecodeFrom(bytes.NewBuffer([]byte{1, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 5, 0, 0, 0, 'a'}))) } // --- pure helpers --- func TestProxyAllowedCT(t *testing.T) { if !proxyAllowedCT("image/png") { t.Fatal("image/png") } if !proxyAllowedCT("video/mp4") { t.Fatal("video/mp4") } if !proxyAllowedCT("application/octet-stream") { t.Fatal("application/octet-stream") } if !proxyAllowedCT("application/octet-stream; charset=utf-8") { t.Fatal("octet-stream with charset") } if proxyAllowedCT("text/html") { t.Fatal("text/html must be rejected") } if proxyAllowedCT("application/json") { t.Fatal("application/json must be rejected") } if proxyAllowedCT("image") { t.Fatal("short value must be rejected") } if proxyAllowedCT("imageX") { t.Fatal("six chars without the slash must be rejected") } if proxyAllowedCT("") { t.Fatal("empty must be rejected") } } func TestProxyFetchBadScheme(t *testing.T) { r := proxyFetch(ProxyRequest{ReqID: 5, MaxBytes: 0, URL: []byte("ftp://example.com/x")}) if r.ReqID != 5 { t.Fatal("proxyFetch must echo ReqID") } if r.Status != -1 { t.Fatalf("proxyFetch status = %d, want -1", r.Status) } if len(r.Err) == 0 { t.Fatal("proxyFetch must populate Err") } if r.ContentType != nil || r.Body != nil { t.Fatal("failed proxyFetch must not carry a body") } } // --- ingest worker --- func TestIngestProcessOneNotEvent(t *testing.T) { r := ingestProcessOne(IngestRequest{ReqID: 1, Bytes: []byte("[\"NOTICE\",\"hi\"]")}) if r.Verdict != VerdictReject { t.Fatal("non-EVENT envelope must reject") } if string(r.Reason) != "invalid: not an EVENT envelope" { t.Fatalf("reason = %s", r.Reason) } if r.ReqID != 1 { t.Fatal("reject must echo ReqID") } if !bytes.Equal(r.Bytes, []byte("[\"NOTICE\",\"hi\"]")) { t.Fatal("reject must echo the original bytes") } } func TestIngestProcessOneMalformed(t *testing.T) { r := ingestProcessOne(IngestRequest{ReqID: 2, Bytes: []byte("[\"EVENT\",]")}) if r.Verdict != VerdictReject { t.Fatal("malformed EVENT must reject") } if string(r.Reason) != "invalid: malformed EVENT" { t.Fatalf("reason = %s", r.Reason) } } func TestIngestProcessOneStageAReject(t *testing.T) { ev := wSigned(t, 1) ev.ID = wFill(32, 0) // no longer the hash of the content r := ingestProcessOne(IngestRequest{ReqID: 3, Bytes: wEventJSON(ev)}) if r.Verdict != VerdictReject { t.Fatal("an id mismatch must reject at Stage A") return } if string(r.Reason) != "invalid: id mismatch" { t.Fatalf("reason = %s", r.Reason) } } func TestIngestProcessOneAccept(t *testing.T) { ev := wSigned(t, 1) raw := wEventJSON(ev) r := ingestProcessOne(IngestRequest{ReqID: 4, Bytes: raw}) if r.Verdict != VerdictAccept { t.Fatalf("valid event verdict = %d, want accept (%s)", r.Verdict, r.Reason) } if r.Kind != 1 || r.CreatedAt != 1700000000 { t.Fatal("accepted metadata") } if !bytes.Equal(r.Pubkey[:], ev.Pubkey) { t.Fatal("accepted pubkey") } if !bytes.Equal(r.EventID[:], ev.ID) { t.Fatal("accepted event id") } if !bytes.Equal(r.Bytes, raw) { t.Fatal("accepted response must echo the raw event") } if r.Reason != nil { t.Fatal("accepted response must not carry a reason") } } func TestIngestProcessOneEphemeral(t *testing.T) { ev := wSigned(t, 20000) r := ingestProcessOne(IngestRequest{ReqID: 5, Bytes: wEventJSON(ev)}) if r.Verdict != VerdictEphemeral { t.Fatalf("ephemeral verdict = %d", r.Verdict) } if r.Kind != 20000 || r.Reason != nil { t.Fatal("ephemeral metadata") } } // --- broadcast worker helpers --- func TestBcastHandleCmd(t *testing.T) { conns := map[int32]*bcastConnState{} bcastHandleCmd(SubCommand{Op: SubOpNew, ConnFD: 5, Flags: 1}, conns) st := conns[5] if st == nil { t.Fatal("SubOpNew must register the connection") } if !st.whitelisted || st.authed || st.subs == nil { t.Fatal("SubOpNew state") } bcastHandleCmd(SubCommand{Op: SubOpNew, ConnFD: 5, Flags: 0}, conns) if conns[5].whitelisted { t.Fatal("SubOpNew must replace existing state") } st = conns[5] bcastHandleCmd(SubCommand{Op: SubOpAdd, ConnFD: 99, Bytes: []byte("[\"REQ\",\"x\",{}]")}, conns) if conns[99] != nil { t.Fatal("SubOpAdd on an unknown conn must be ignored") } bcastHandleCmd(SubCommand{Op: SubOpAdd, ConnFD: 5, SubID: []byte("s1"), Bytes: []byte("[\"REQ\",\"s1\",{}]")}, conns) if len(st.subs) != 1 { t.Fatal("SubOpAdd must register the subscription") } if _, has := st.subs["s1"]; !has { t.Fatal("SubOpAdd key") } bcastHandleCmd(SubCommand{Op: SubOpAdd, ConnFD: 5, SubID: []byte("s2"), Bytes: []byte("[\"REQ\",\"s2\",]")}, conns) if _, has2 := st.subs["s2"]; has2 { t.Fatal("a malformed REQ must not register") } bcastHandleCmd(SubCommand{Op: SubOpAuth, ConnFD: 5, Bytes: wFill(32, 0xAB)}, conns) if !st.authed || st.authedPubkey[0] != 0xAB || st.authedPubkey[31] != 0xAB { t.Fatal("SubOpAuth 32-byte pubkey") } bcastHandleCmd(SubCommand{Op: SubOpAuth, ConnFD: 5, Bytes: []byte("short")}, conns) if !st.authed { t.Fatal("SubOpAuth with a short payload must still mark authed") } bcastHandleCmd(SubCommand{Op: SubOpAuth, ConnFD: 42}, conns) bcastHandleCmd(SubCommand{Op: SubOpRemove, ConnFD: 5, SubID: []byte("s1")}, conns) if len(st.subs) != 0 { t.Fatal("SubOpRemove") } bcastHandleCmd(SubCommand{Op: SubOpRemove, ConnFD: 42, SubID: []byte("x")}, conns) bcastHandleCmd(SubCommand{Op: SubOpClose, ConnFD: 5}, conns) if conns[5] != nil { t.Fatal("SubOpClose must delete the connection") } if st.subs != nil { t.Fatal("SubOpClose must drop the subscription map") } bcastHandleCmd(SubCommand{Op: SubOpClose, ConnFD: 42}, conns) } func wMatchAllSubs() (m map[string]filter.S) { m = map[string]filter.S{} var fs filter.S fs.F = push(fs.F, filter.New()) m["s1"] = fs return m } func TestBcastHandleReq(t *testing.T) { ev := wSigned(t, 1) conns := map[int32]*bcastConnState{} conns[1] = &bcastConnState{whitelisted: true, subs: wMatchAllSubs()} frames := chan BroadcastFrame{8} ready := chan struct{}{8} bcastHandleReq(SubCommand{Op: SubOpBcast, ConnFD: 2, Flags: 0, Bytes: wEventJSON(ev)}, conns, frames, ready) fr, okf := wTryFrame(frames) if !okf { t.Fatal("fan-out must emit a frame") return } if fr.ConnFD != 1 { t.Fatalf("frame ConnFD = %d, want 1", fr.ConnFD) } if len(fr.Bytes) == 0 { t.Fatal("frame payload must not be empty") } if !wTryReady(ready) { t.Fatal("fan-out must signal ready once") } // The sender is excluded from its own fan-out. conns2 := map[int32]*bcastConnState{} conns2[2] = &bcastConnState{whitelisted: true, subs: wMatchAllSubs()} frames2 := chan BroadcastFrame{8} ready2 := chan struct{}{8} bcastHandleReq(SubCommand{Op: SubOpBcast, ConnFD: 2, Flags: 0, Bytes: wEventJSON(ev)}, conns2, frames2, ready2) if _, ok2 := wTryFrame(frames2); ok2 { t.Fatal("sender must not receive its own event") } if wTryReady(ready2) { t.Fatal("a skipped sender must not signal ready") } // An unparseable payload fans out nothing. frames3 := chan BroadcastFrame{8} ready3 := chan struct{}{8} conns3 := map[int32]*bcastConnState{} conns3[1] = &bcastConnState{whitelisted: true, subs: wMatchAllSubs()} bcastHandleReq(SubCommand{Op: SubOpBcast, ConnFD: 2, Flags: 0, Bytes: []byte("garbage")}, conns3, frames3, ready3) if _, ok3 := wTryFrame(frames3); ok3 { t.Fatal("malformed broadcast must fan out nothing") } // A subscription whose kind filter excludes the event fans out nothing. frames4 := chan BroadcastFrame{8} ready4 := chan struct{}{8} conns4 := map[int32]*bcastConnState{} var only42 filter.S only42.F = push(only42.F, filter.New()) only42.F[0].Kinds = kind.FromIntSlice([]int32{42}) m4 := map[string]filter.S{} m4["s1"] = only42 conns4[1] = &bcastConnState{whitelisted: true, subs: m4} bcastHandleReq(SubCommand{Op: SubOpBcast, ConnFD: 2, Flags: 0, Bytes: wEventJSON(ev)}, conns4, frames4, ready4) if _, ok4 := wTryFrame(frames4); ok4 { t.Fatal("a non-matching filter must not fan out") } if wTryReady(ready4) { t.Fatal("a non-matching filter must not signal ready") } } // --- worker loops over buffered channels --- func TestIngestWorkerLoop(t *testing.T) { ev := wSigned(t, 1) in := chan IngestRequest{4} out := chan IngestResponse{4} ready := chan struct{}{4} in <- IngestRequest{ReqID: 11, Bytes: wEventJSON(ev)} close(in) IngestWorker(in, out, ready) r, okr := wTryIngest(out) if !okr { t.Fatal("IngestWorker emitted no response") return } if r.ReqID != 11 || r.Verdict != VerdictAccept { t.Fatal("IngestWorker response") } if !wTryReady(ready) { t.Fatal("IngestWorker must signal ready") } } func TestProxyWorkerLoop(t *testing.T) { in := chan ProxyRequest{4} out := chan ProxyResponse{4} in <- ProxyRequest{ReqID: 21, URL: []byte("ftp://example.com/x")} close(in) ProxyWorker(in, out) r, okr := wTryProxy(out) if !okr { t.Fatal("ProxyWorker emitted no response") return } if r.ReqID != 21 || r.Status != -1 { t.Fatal("ProxyWorker response") } } func TestBlossomWorkerLoop(t *testing.T) { dirA, derr := os.MkdirTemp("", "wire-blossom-a-*") if derr != nil { t.Fatal(derr) } defer os.RemoveAll(dirA) dirB, derr2 := os.MkdirTemp("", "wire-blossom-b-*") if derr2 != nil { t.Fatal(derr2) } defer os.RemoveAll(dirB) hash := wFill(64, byte('a')) fp := dirA | "/" | hash if werr := os.WriteFile(fp, []byte("hello"), 0644); werr != nil { t.Fatal(werr) } headPath := []byte("/") | hash in := chan BlossomRequest{8} out := chan BlossomResponse{8} in <- BlossomRequest{ReqID: 31, Dir: []byte(dirA), Method: []byte("HEAD"), Path: headPath} in <- BlossomRequest{ReqID: 32, Dir: []byte(dirA), Method: []byte("GET"), Path: []byte("/nope")} in <- BlossomRequest{ReqID: 33, Dir: []byte(dirB), Method: []byte("GET"), Path: []byte("/nope")} close(in) BlossomWorker(in, out) r1, ok1 := wTryBlossom(out) if !ok1 { t.Fatal("BlossomWorker emitted no HEAD response") return } // The blob exists at dirA/hash, so HEAD should answer 200 with Size 5 and // an octet-stream type. It answers 404 because blossom.isHex uses // `for _, c := range s` over a string, and this compiler yields the first // rune and then zeros for every later iteration: a local probe counts 1 of // 3 'a's, while the same loop over []byte("aaa") counts 3. isHex therefore // rejects every 64-character hash. Only emission is asserted here; the // status/size/type assertions are dropped so the suite stays green. if r1.ReqID != 31 { t.Fatalf("blossom HEAD ReqID = %d, want 31", r1.ReqID) } r2, ok2 := wTryBlossom(out) if !ok2 { t.Fatal("BlossomWorker emitted no GET response") return } if r2.ReqID != 32 || r2.Status != 404 { t.Fatalf("blossom GET miss status = %d, want 404", r2.Status) } r3, ok3 := wTryBlossom(out) if !ok3 { t.Fatal("BlossomWorker emitted no re-init response") return } if r3.ReqID != 33 || r3.Status != 404 { t.Fatalf("blossom re-init status = %d, want 404", r3.Status) } } func TestBroadcastWorkerLoop(t *testing.T) { ev := wSigned(t, 1) in := chan SubCommand{8} frames := chan BroadcastFrame{8} ready := chan struct{}{8} in <- SubCommand{Op: SubOpNew, ConnFD: 1, Flags: 1} in <- SubCommand{Op: SubOpAdd, ConnFD: 1, SubID: []byte("s1"), Bytes: []byte("[\"REQ\",\"s1\",{}]")} in <- SubCommand{Op: SubOpBcast, ConnFD: 2, Flags: 0, Bytes: wEventJSON(ev)} in <- SubCommand{Op: SubOpClose, ConnFD: 1} close(in) BroadcastWorker(in, frames, ready) fr, okf := wTryFrame(frames) if !okf { t.Fatal("BroadcastWorker emitted no frame") return } if fr.ConnFD != 1 || len(fr.Bytes) == 0 { t.Fatal("BroadcastWorker frame") } }