diff --git a/qkc/slave/frame.go b/qkc/slave/frame.go index 7cc6189221..063c35d2ee 100644 --- a/qkc/slave/frame.go +++ b/qkc/slave/frame.go @@ -1,56 +1,33 @@ -// Package slave: binary frame codec compatible with pyquarkchain protocol.py. +// Copyright 2026-2027, QuarkChain. // -// Wire format (per-frame): +// Package slave implements a binary frame codec compatible with pyquarkchain. // -// [4B payload_len] [metaSize B metadata] [1B opcode] [8B rpc_id] [payload_len bytes] +// Wire format (per frame): // -// This matches Python's Connection.read_metadata_and_raw_data / -// Connection.write_raw_data exactly: +// [4B payload_len] [metaSize B metadata] [1B opcode] [8B rpc_id] [payload bytes] // -// Python: protocol.py lines 285-308 -// -// Metadata sizes (matching Python Metadata subclasses): -// -// ClusterMetadata.get_byte_size() = 12 (branch 4B + cluster_peer_id 8B) -// Used for master ↔ slave traffic. -// -// Metadata.get_byte_size() = 0 (base class) -// Used for slave ↔ slave traffic (SlaveConnection). -// -// P2PMetadata.get_byte_size() = 4 (branch 4B) -// Used for inter-cluster P2P. Handled by Python master only; -// Go slave never sends or receives P2PMetadata frames. -// -// payload_len definition (matching Python line 305): -// -// cmd_length_bytes = (len(raw_data) - 8 - 1).to_bytes(4, "big") -// -// That is, payload_len = len(raw_data) - 9, where raw_data is -// [1B opcode][8B rpc_id][N bytes payload]. +// payload_len is the length of payload bytes only. +// metadata size depends on metaSize parameter: +// - 12 bytes: ClusterMetadata (branch uint32 + cluster_peer_id uint64) for master↔slave +// - 0 bytes: no metadata for slave↔slave traffic package slave import ( - "bufio" "encoding/binary" "errors" "fmt" "io" ) -// Metadata is the 12-byte frame header carrying routing information. -// Matches Python's ClusterMetadata: -// -// class ClusterMetadata(Metadata): -// FIELDS = [("branch", Branch), ("cluster_peer_id", uint64)] -// @staticmethod -// def get_byte_size(): return 12 +// Metadata is the routing header for ClusterMetadata traffic. +// Wire representation: 12 bytes (4B branch + 8B cluster_peer_id). type Metadata struct { - Branch uint32 // shard identifier (Python: Branch = uint32) - ClusterPeerID uint64 // 0 = cluster RPC (master commands), ≠0 = specific external peer + Branch uint32 + ClusterPeerID uint64 } // Frame is a complete protocol frame. -// raw_data layout on wire: [1B opcode][8B rpc_id][N bytes payload] +// Wire layout after metadata: [1B opcode][8B rpc_id][N bytes payload] type Frame struct { Meta Metadata Opcode byte @@ -59,53 +36,53 @@ type Frame struct { } const ( - metaSize = 12 // ClusterMetadata.get_byte_size() = 12 (branch 4B + cluster_peer_id 8B) - opcodeSize = 1 - rpcIDSize = 8 - frameHeader = 4 // payload_len prefix - totalOverhead = frameHeader + metaSize + opcodeSize + rpcIDSize // 4+12+1+8 = 25 + metaSize = 12 // ClusterMetadata.get_byte_size() = 12 (branch 4B + cluster_peer_id 8B) + opcodeSize = 1 + rpcIDSize = 8 + frameHeader = 4 // payload_len prefix ) -// ReadFrame reads a frame with 12-byte ClusterMetadata. -// -// Matches Python's Connection.read_metadata_and_raw_data (protocol.py lines 285-300) -// with metadata_class = ClusterMetadata (get_byte_size() == 12). -// -// Used for master ↔ slave traffic. -// For slave ↔ slave traffic, use ReadFrameNoMeta (0-byte metadata). +// ReadFrame reads a frame with 12-byte metadata and no payload-length cap. func ReadFrame(r io.Reader) (*Frame, error) { - return readFrameWithMetaSize(r, metaSize) + return readFrame(r, metaSize, 0) } -// ReadFrameNoMeta reads a frame with 0-byte Metadata. -// -// Matches Python's SlaveConnection which uses metadata_class = Metadata -// (get_byte_size() == 0). -// -// Used for slave ↔ slave direct TCP traffic. +// ReadFrameWithMaxPayload reads a frame with 12-byte metadata and rejects frames +// whose payload_len exceeds maxPayloadLen before allocation. +func ReadFrameWithMaxPayload(r io.Reader, maxPayloadLen uint32) (*Frame, error) { + if maxPayloadLen == 0 { + return nil, errors.New("maxPayloadLen must be greater than zero") + } + return readFrame(r, metaSize, maxPayloadLen) +} + +// ReadFrameNoMeta reads a frame with 0-byte metadata and no payload-length cap. func ReadFrameNoMeta(r io.Reader) (*Frame, error) { - return readFrameWithMetaSize(r, 0) + return readFrame(r, 0, 0) } -// readFrameWithMetaSize is the underlying frame reader. -// -// Wire layout (matching Python read_metadata_and_raw_data): -// -// size_bytes = await read_fully(4) → 4B payload_len (big-endian) -// metadata_bytes = await read_fully(metaSize) → metaSize B metadata -// raw_data_without_size = await read_fully(1+8+size) → opcode + rpc_id + payload -func readFrameWithMetaSize(r io.Reader, metaSize int) (*Frame, error) { +// ReadFrameNoMetaWithMaxPayload reads a frame with 0-byte metadata and rejects +// frames whose payload_len exceeds maxPayloadLen before allocation. +func ReadFrameNoMetaWithMaxPayload(r io.Reader, maxPayloadLen uint32) (*Frame, error) { + if maxPayloadLen == 0 { + return nil, errors.New("maxPayloadLen must be greater than zero") + } + return readFrame(r, 0, maxPayloadLen) +} + +// readFrame is the underlying frame reader. +func readFrame(r io.Reader, metaSize int, maxPayloadLen uint32) (*Frame, error) { // 1. Read 4-byte big-endian payload length - // Python: size_bytes = await self.__read_fully(4, allow_eof=True) - // size = int.from_bytes(size_bytes, byteorder="big") var payloadLen uint32 if err := binary.Read(r, binary.BigEndian, &payloadLen); err != nil { return nil, fmt.Errorf("reading frame length: %w", err) } + if maxPayloadLen != 0 && payloadLen > maxPayloadLen { + return nil, fmt.Errorf("frame payload too large: %d > %d", payloadLen, maxPayloadLen) + } + // 2. Read metadata (size depends on metadata_class) - // Python: metadata_bytes = await self.__read_fully(self.metadata_class.get_byte_size()) - // metadata = self.metadata_class.deserialize(metadata_bytes) var meta Metadata if metaSize > 0 { if metaSize != 12 { @@ -122,7 +99,6 @@ func readFrameWithMetaSize(r io.Reader, metaSize int) (*Frame, error) { } // 3. Read raw_data_without_size: [1B opcode][8B rpc_id][N bytes payload] - // Python: raw_data_without_size = await self.__read_fully(1 + 8 + size) bodySize := opcodeSize + rpcIDSize + int(payloadLen) body := make([]byte, bodySize) if _, err := io.ReadFull(r, body); err != nil { @@ -137,47 +113,26 @@ func readFrameWithMetaSize(r io.Reader, metaSize int) (*Frame, error) { }, nil } -// WriteFrame serializes f with 12-byte ClusterMetadata and writes it to w. -// -// Matches Python's Connection.write_raw_data (protocol.py lines 302-308) -// with metadata_class = ClusterMetadata. -// -// Used for master ↔ slave traffic. +// WriteFrame serializes f with 12-byte metadata and writes it to w. func WriteFrame(w io.Writer, f *Frame) error { return writeFrameWithMetaSize(w, f, metaSize) } -// WriteFrameNoMeta writes a frame with 0-byte Metadata. -// -// Matches Python's SlaveConnection which uses metadata_class = Metadata -// (get_byte_size() == 0). -// -// Used for slave ↔ slave traffic. +// WriteFrameNoMeta writes a frame with 0-byte metadata. func WriteFrameNoMeta(w io.Writer, f *Frame) error { return writeFrameWithMetaSize(w, f, 0) } -// writeFrameWithMetaSize serializes f with the given metadata size and writes -// it to w. -// -// Wire layout (matching Python write_raw_data, protocol.py lines 302-308): -// -// cmd_length_bytes = (len(raw_data) - 8 - 1).to_bytes(4, "big") -// self.writer.write(cmd_length_bytes) → 4B payload_len -// self.writer.write(metadata.serialize()) → metaSize B -// self.writer.write(raw_data) → [1B opcode][8B rpc_id][payload] +// writeFrameWithMetaSize serializes f with the given metadata size and writes it to w. func writeFrameWithMetaSize(w io.Writer, f *Frame, metaSize int) error { payloadLen := uint32(len(f.Payload)) if int(payloadLen) != len(f.Payload) { return errors.New("payload too large") } - // Build the buffer in one write (Python writes in 3 chunks, but the wire - // bytes are identical). total := frameHeader + metaSize + opcodeSize + rpcIDSize + int(payloadLen) buf := make([]byte, total) - // Frame length (payload only): matches Python's cmd_length_bytes binary.BigEndian.PutUint32(buf[0:frameHeader], payloadLen) // Metadata @@ -189,7 +144,6 @@ func writeFrameWithMetaSize(w io.Writer, f *Frame, metaSize int) error { // Opcode buf[frameHeader+metaSize] = f.Opcode - // RPC ID (big-endian, matches Python: rpc_id.to_bytes(8, "big")) binary.BigEndian.PutUint64(buf[frameHeader+metaSize+opcodeSize:frameHeader+metaSize+opcodeSize+rpcIDSize], f.RPCID) // Payload @@ -217,33 +171,3 @@ func UnmarshalMetadata(b []byte) (Metadata, error) { ClusterPeerID: binary.BigEndian.Uint64(b[4:12]), }, nil } - -// ── Convenience wrappers (used in tests) ───────────────────────────────────── - -// ReadFrameFromReader wraps r in a bufio.Reader. -func ReadFrameFromReader(r io.Reader) (*Frame, error) { - return ReadFrame(bufio.NewReader(r)) -} - -// WriteFrameToWriter wraps w with a bufio.Writer and flushes. -func WriteFrameToWriter(w io.Writer, frame *Frame) error { - bw := bufio.NewWriter(w) - if err := WriteFrame(bw, frame); err != nil { - return err - } - return bw.Flush() -} - -// ReadFrameNoMetaFromReader wraps r for ReadFrameNoMeta. -func ReadFrameNoMetaFromReader(r io.Reader) (*Frame, error) { - return ReadFrameNoMeta(bufio.NewReader(r)) -} - -// WriteFrameNoMetaToWriter wraps w for WriteFrameNoMeta. -func WriteFrameNoMetaToWriter(w io.Writer, frame *Frame) error { - bw := bufio.NewWriter(w) - if err := WriteFrameNoMeta(bw, frame); err != nil { - return err - } - return bw.Flush() -} diff --git a/qkc/slave/frame_test.go b/qkc/slave/frame_test.go index dceca81790..cfc6879928 100644 --- a/qkc/slave/frame_test.go +++ b/qkc/slave/frame_test.go @@ -1,218 +1,115 @@ +// Copyright 2026-2027, QuarkChain. + package slave import ( "bytes" "encoding/binary" "encoding/hex" - "strings" "testing" ) -// ============================================================================= -// Python compatibility reference vectors -// -// Each vector below is the exact wire bytes a Python peer would send/receive -// for a given frame. Source: qkc/quarkchain/protocol.py -// - Connection.write_raw_data (lines 302-308) -// - Connection.read_metadata_and_raw_data (lines 285-300) -// -// Wire layout: [4B payload_len] [metaSize B metadata] [1B opcode] [8B rpc_id] [payload] -// -// metadata_class: -// ClusterMetadata (12B) for master↔slave traffic -// Metadata (0B) for slave↔slave traffic -// ============================================================================= +var ( + // Test vectors serialized by pyquarkchain serializers using synthetic test values. + // + // Payload content (Ping/Pong commands): + // - id = "id" (ASCII, 2 bytes) + // - full_shard_id_list = [1, 2] + // - root_tip = None (Ping only) + // - opcode = 0x81 (PING), 0x82 (PONG) from ClusterOp (CLUSTER_OP_BASE=128) + pythonClusterPing = "0000001300000001000000000000303981000000000000000100000002696400000002000000010000000200" + pythonClusterPong = "00000012000000010000000000003039820000000000000001000000026964000000020000000100000002" + pythonNoMetaPing = "0000001381000000000000000100000002696400000002000000010000000200" + pythonNoMetaPong = "00000012820000000000000001000000026964000000020000000100000002" +) -// pingMasterWire: meta=(branch=0, peer=0), opcode=0x81 (PING), rpc_id=1, payload=empty -// -// Equivalent Python: write_raw_command(op=ClusterOp.PING, cmd_data=b"", rpc_id=1, metadata=ClusterMetadata(0, 0)) -// payload_len = 0 -// 00000000 | 00000000 0000000000000000 | 81 | 0000000000000001 -var pingMasterWire = "00000000" + "00000000" + "0000000000000000" + "81" + "0000000000000001" - -// pongMasterWire: meta=(branch=0, peer=0), opcode=0x82 (PONG), rpc_id=1, payload=empty -// -// Equivalent Python: write_raw_command(op=ClusterOp.PONG, cmd_data=b"", rpc_id=1, metadata=ClusterMetadata(0, 0)) -// 00000000 | 00000000 0000000000000000 | 82 | 0000000000000001 -var pongMasterWire = "00000000" + "00000000" + "0000000000000000" + "82" + "0000000000000001" - -// peerNewBlockWire: meta=(branch=1, peer=12345), opcode=0x01 (NEW_MINOR_BLOCK_HEADER_LIST), rpc_id=0, payload=12B -// -// cluster_peer_id=12345=0x3039, rpc_id=0 (NON-RPC fire-and-forget) -// payload = 02 00 00 00 (list len=2) + a1b2c3d4 + e5f60718 -// 0000000c | 00000001 0000000000003039 | 01 | 0000000000000000 | 02000000a1b2c3d4e5f60718 -var peerNewBlockWire = "0000000c" + - "00000001" + "0000000000003039" + - "01" + "0000000000000000" + - "02000000a1b2c3d4e5f60718" - -// xshardWire: 0-byte metadata, opcode=0x93 (ADD_XSHARD_TX_LIST_REQUEST), rpc_id=42, payload=56B -// -// Used for slave↔slave direct TCP (Python SlaveConnection, metadata_class=Metadata) -// payload_len = 56 = 0x38 -// payload = 01 00 00 00 (list len=1) + ff*32 + 00*20 -// 00000038 | (no meta) | 93 | 000000000000002a | 01000000 + ff×32 + 00×20 -var xshardWire = "00000038" + - "93" + "000000000000002a" + - "01000000" + - strings.Repeat("ff", 32) + - strings.Repeat("00", 20) - -// largePayloadWire: meta=(branch=3, peer=0), opcode=0x10, rpc_id=7, payload=10000×0xAB -// -// payload_len = 10000 = 0x2710 -// 00002710 | 00000003 0000000000000000 | 10 | 0000000000000007 | ab×10000 -var largePayloadWire = "00002710" + - "00000003" + "0000000000000000" + - "10" + "0000000000000007" + - strings.Repeat("ab", 10000) - -// ============================================================================= -// ReadFrame tests — parse Python-generated wire bytes -// ============================================================================= - -func TestPythonRead_Ping(t *testing.T) { - wire, _ := hex.DecodeString(pingMasterWire) +func TestPythonVectors_ClusterPing(t *testing.T) { + wire := mustPythonVectorBytes(t, pythonClusterPing) f, err := ReadFrame(bytes.NewReader(wire)) if err != nil { t.Fatalf("ReadFrame: %v", err) } - if f.Opcode != 0x81 { - t.Errorf("Opcode: got 0x%02x, want 0x81", f.Opcode) + if f.Opcode != 0x81 || f.RPCID != 1 { + t.Fatalf("unexpected header: opcode=0x%02x rpcid=%d", f.Opcode, f.RPCID) } - if f.RPCID != 1 { - t.Errorf("RPCID: got %d, want 1", f.RPCID) + if f.Meta.Branch != 1 || f.Meta.ClusterPeerID != 12345 { + t.Fatalf("unexpected meta: %+v", f.Meta) } - if f.Meta != (Metadata{Branch: 0, ClusterPeerID: 0}) { - t.Errorf("Meta: got %+v", f.Meta) + var out bytes.Buffer + if err := WriteFrame(&out, f); err != nil { + t.Fatalf("WriteFrame: %v", err) } - if len(f.Payload) != 0 { - t.Errorf("Payload len: got %d, want 0", len(f.Payload)) + if !bytes.Equal(out.Bytes(), wire) { + t.Fatalf("frame mismatch:\n go %x\n py %x", out.Bytes(), wire) } } -func TestPythonRead_Pong(t *testing.T) { - wire, _ := hex.DecodeString(pongMasterWire) +func TestPythonVectors_ClusterPong(t *testing.T) { + wire := mustPythonVectorBytes(t, pythonClusterPong) f, err := ReadFrame(bytes.NewReader(wire)) if err != nil { t.Fatalf("ReadFrame: %v", err) } - if f.Opcode != 0x82 { - t.Errorf("Opcode: got 0x%02x, want 0x82", f.Opcode) + if f.Opcode != 0x82 || f.RPCID != 1 { + t.Fatalf("unexpected header: opcode=0x%02x rpcid=%d", f.Opcode, f.RPCID) } - if f.RPCID != 1 { - t.Errorf("RPCID: got %d, want 1", f.RPCID) + if f.Meta.Branch != 1 || f.Meta.ClusterPeerID != 12345 { + t.Fatalf("unexpected meta: %+v", f.Meta) + } + var out bytes.Buffer + if err := WriteFrame(&out, f); err != nil { + t.Fatalf("WriteFrame: %v", err) + } + if !bytes.Equal(out.Bytes(), wire) { + t.Fatalf("frame mismatch:\n go %x\n py %x", out.Bytes(), wire) } } -func TestPythonRead_PeerNewBlock(t *testing.T) { - wire, _ := hex.DecodeString(peerNewBlockWire) - f, err := ReadFrame(bytes.NewReader(wire)) - if err != nil { - t.Fatalf("ReadFrame: %v", err) - } - if f.Opcode != 0x01 { - t.Errorf("Opcode: got 0x%02x, want 0x01", f.Opcode) - } - if f.RPCID != 0 { - t.Errorf("RPCID: got %d, want 0 (non-RPC)", f.RPCID) - } - if f.Meta.Branch != 1 { - t.Errorf("Branch: got %d, want 1", f.Meta.Branch) - } - if f.Meta.ClusterPeerID != 12345 { - t.Errorf("ClusterPeerID: got %d, want 12345", f.Meta.ClusterPeerID) - } -} - -func TestPythonRead_Xshard(t *testing.T) { - wire, _ := hex.DecodeString(xshardWire) +func TestPythonVectors_NoMetaPing(t *testing.T) { + wire := mustPythonVectorBytes(t, pythonNoMetaPing) f, err := ReadFrameNoMeta(bytes.NewReader(wire)) if err != nil { t.Fatalf("ReadFrameNoMeta: %v", err) } - if f.Opcode != 0x93 { - t.Errorf("Opcode: got 0x%02x, want 0x93", f.Opcode) + if f.Opcode != 0x81 || f.RPCID != 1 { + t.Fatalf("unexpected header: opcode=0x%02x rpcid=%d", f.Opcode, f.RPCID) } - if f.RPCID != 42 { - t.Errorf("RPCID: got %d, want 42", f.RPCID) + var out bytes.Buffer + if err := WriteFrameNoMeta(&out, f); err != nil { + t.Fatalf("WriteFrameNoMeta: %v", err) } - if len(f.Payload) != 56 { - t.Errorf("Payload len: got %d, want 56", len(f.Payload)) + if !bytes.Equal(out.Bytes(), wire) { + t.Fatalf("frame mismatch:\n go %x\n py %x", out.Bytes(), wire) } } -func TestPythonRead_LargePayload(t *testing.T) { - wire, _ := hex.DecodeString(largePayloadWire) - f, err := ReadFrame(bytes.NewReader(wire)) +func TestPythonVectors_NoMetaPong(t *testing.T) { + wire := mustPythonVectorBytes(t, pythonNoMetaPong) + f, err := ReadFrameNoMeta(bytes.NewReader(wire)) if err != nil { - t.Fatalf("ReadFrame: %v", err) + t.Fatalf("ReadFrameNoMeta: %v", err) } - if f.Opcode != 0x10 { - t.Errorf("Opcode: got 0x%02x, want 0x10", f.Opcode) + if f.Opcode != 0x82 || f.RPCID != 1 { + t.Fatalf("unexpected header: opcode=0x%02x rpcid=%d", f.Opcode, f.RPCID) } - if f.RPCID != 7 { - t.Errorf("RPCID: got %d, want 7", f.RPCID) + var out bytes.Buffer + if err := WriteFrameNoMeta(&out, f); err != nil { + t.Fatalf("WriteFrameNoMeta: %v", err) } - if f.Meta.Branch != 3 { - t.Errorf("Branch: got %d, want 3", f.Meta.Branch) - } - if len(f.Payload) != 10000 { - t.Errorf("Payload len: got %d, want 10000", len(f.Payload)) - } - for i, b := range f.Payload { - if b != 0xAB { - t.Errorf("Payload[%d]: got 0x%02x, want 0xAB", i, b) - break - } + if !bytes.Equal(out.Bytes(), wire) { + t.Fatalf("frame mismatch:\n go %x\n py %x", out.Bytes(), wire) } } -// ============================================================================= -// WriteFrame tests — Go's wire output must match Python's byte-for-byte. -// This is the strongest compatibility test: any wire-format drift is caught. -// ============================================================================= - -func TestPythonWrite_Ping(t *testing.T) { - want, _ := hex.DecodeString(pingMasterWire) - got := writeFrameForTest(Metadata{Branch: 0, ClusterPeerID: 0}, 0x81, 1, nil) - if !bytes.Equal(got, want) { - t.Errorf("WriteFrame mismatch:\n got %x\n want %x", got, want) +func mustPythonVectorBytes(t *testing.T, hexStr string) []byte { + t.Helper() + b, err := hex.DecodeString(hexStr) + if err != nil { + t.Fatalf("decode python hex: %v", err) } + return b } -func TestPythonWrite_PeerNewBlock(t *testing.T) { - want, _ := hex.DecodeString(peerNewBlockWire) - payload := []byte{0x02, 0x00, 0x00, 0x00, 0xa1, 0xb2, 0xc3, 0xd4, 0xe5, 0xf6, 0x07, 0x18} - got := writeFrameForTest(Metadata{Branch: 1, ClusterPeerID: 12345}, 0x01, 0, payload) - if !bytes.Equal(got, want) { - t.Errorf("WriteFrame mismatch:\n got %x\n want %x", got, want) - } -} - -func TestPythonWrite_Xshard(t *testing.T) { - want, _ := hex.DecodeString(xshardWire) - payload := append([]byte{0x01, 0x00, 0x00, 0x00}, bytes.Repeat([]byte{0xff}, 32)...) - payload = append(payload, bytes.Repeat([]byte{0x00}, 20)...) - got := writeFrameNoMetaForTest(0x93, 42, payload) - if !bytes.Equal(got, want) { - t.Errorf("WriteFrameNoMeta mismatch:\n got %x\n want %x", got, want) - } -} - -func TestPythonWrite_LargePayload(t *testing.T) { - want, _ := hex.DecodeString(largePayloadWire) - payload := bytes.Repeat([]byte{0xAB}, 10000) - got := writeFrameForTest(Metadata{Branch: 3, ClusterPeerID: 0}, 0x10, 7, payload) - if !bytes.Equal(got, want) { - t.Errorf("WriteFrame mismatch (large):\n got %d bytes\n want %d bytes", len(got), len(want)) - } -} - -// ============================================================================= -// Write+Read round-trip — defensive tests independent of the Python reference -// ============================================================================= - func TestRoundTrip_Meta(t *testing.T) { cases := []struct { name string @@ -255,11 +152,6 @@ func TestRoundTrip_NoMeta(t *testing.T) { } } -// ============================================================================= -// Wire-format layout — verify byte-level structure with hand-computed expected -// values, independent of any Python reference. -// ============================================================================= - func TestWireFormatLayout(t *testing.T) { f := &Frame{ Meta: Metadata{Branch: 1, ClusterPeerID: 0x1122334455667788}, @@ -269,35 +161,27 @@ func TestWireFormatLayout(t *testing.T) { } wire := writeFrameForTest(f.Meta, f.Opcode, f.RPCID, f.Payload) - // payload_len = 3 if got := binary.BigEndian.Uint32(wire[0:4]); got != 3 { t.Errorf("payload_len: got %d, want 3", got) } - // branch = 1 if got := binary.BigEndian.Uint32(wire[4:8]); got != 1 { t.Errorf("branch: got %d, want 1", got) } - // cluster_peer_id = 0x1122334455667788 if got := binary.BigEndian.Uint64(wire[8:16]); got != 0x1122334455667788 { t.Errorf("cluster_peer_id: got 0x%x", got) } - // opcode = 0x42 if wire[16] != 0x42 { t.Errorf("opcode: got 0x%02x, want 0x42", wire[16]) } - // rpc_id = 0xDEADBEEFCAFEBABE if got := binary.BigEndian.Uint64(wire[17:25]); got != 0xDEADBEEFCAFEBABE { t.Errorf("rpc_id: got 0x%x", got) } - // payload if !bytes.Equal(wire[25:28], []byte{0xAA, 0xBB, 0xCC}) { t.Errorf("payload: got %x", wire[25:28]) } } func TestMultiFrameStream(t *testing.T) { - // 3 consecutive frames on a single stream, matching Python's back-to-back - // write_raw_data() calls on the same TCP connection. frames := []*Frame{ {Meta: Metadata{Branch: 0, ClusterPeerID: 0}, Opcode: 0x81, RPCID: 0, Payload: []byte("ping")}, {Meta: Metadata{Branch: 2, ClusterPeerID: 999}, Opcode: 0x05, RPCID: 100, Payload: []byte("block_data")}, @@ -323,10 +207,6 @@ func TestMultiFrameStream(t *testing.T) { } } -// ============================================================================= -// Error handling -// ============================================================================= - func TestReadFrame_EOF(t *testing.T) { if _, err := ReadFrame(bytes.NewReader(nil)); err == nil { t.Error("expected error on empty stream") @@ -334,7 +214,6 @@ func TestReadFrame_EOF(t *testing.T) { } func TestReadFrame_Truncated(t *testing.T) { - // payload_len says 100 bytes, but we only give 4 bytes hdr := make([]byte, 4) binary.BigEndian.PutUint32(hdr, 100) if _, err := ReadFrame(bytes.NewReader(hdr)); err == nil { @@ -342,11 +221,36 @@ func TestReadFrame_Truncated(t *testing.T) { } } -// ============================================================================= -// Helpers -// ============================================================================= +func TestReadFrame_PayloadLimit(t *testing.T) { + hdr := make([]byte, 4) + binary.BigEndian.PutUint32(hdr, 100) + + if _, err := ReadFrameWithMaxPayload(bytes.NewReader(hdr), 64); err == nil { + t.Fatal("expected error when payload_len exceeds limit") + } + + if _, err := ReadFrameWithMaxPayload(bytes.NewReader(hdr), 0); err == nil { + t.Fatal("expected error when maxPayloadLen is zero") + } + + if _, err := ReadFrame(bytes.NewReader(hdr)); err == nil { + t.Error("expected error on truncated frame without payload limit") + } +} + +func TestReadFrameNoMeta_PayloadLimit(t *testing.T) { + hdr := make([]byte, 4) + binary.BigEndian.PutUint32(hdr, 32) + + if _, err := ReadFrameNoMetaWithMaxPayload(bytes.NewReader(hdr), 16); err == nil { + t.Fatal("expected error when payload_len exceeds limit") + } + + if _, err := ReadFrameNoMetaWithMaxPayload(bytes.NewReader(hdr), 0); err == nil { + t.Fatal("expected error when maxPayloadLen is zero") + } +} -// writeFrameForTest is a thin wrapper that returns the wire bytes directly. func writeFrameForTest(meta Metadata, opcode byte, rpcID uint64, payload []byte) []byte { var buf bytes.Buffer _ = WriteFrame(&buf, &Frame{Meta: meta, Opcode: opcode, RPCID: rpcID, Payload: payload})