From 58b66d610ed63dc1fbe24c2b0d3a74f6ab435e1c Mon Sep 17 00:00:00 2001 From: iteye Date: Wed, 1 Jul 2026 18:33:16 +0800 Subject: [PATCH 1/4] qkc/slave: add binary frame codec --- qkc/slave/frame.go | 249 +++++++++++++++++++++++++++ qkc/slave/frame_test.go | 360 ++++++++++++++++++++++++++++++++++++++++ 2 files changed, 609 insertions(+) create mode 100644 qkc/slave/frame.go create mode 100644 qkc/slave/frame_test.go diff --git a/qkc/slave/frame.go b/qkc/slave/frame.go new file mode 100644 index 0000000000..7cc6189221 --- /dev/null +++ b/qkc/slave/frame.go @@ -0,0 +1,249 @@ +// Package slave: binary frame codec compatible with pyquarkchain protocol.py. +// +// Wire format (per-frame): +// +// [4B payload_len] [metaSize B metadata] [1B opcode] [8B rpc_id] [payload_len bytes] +// +// This matches Python's Connection.read_metadata_and_raw_data / +// Connection.write_raw_data exactly: +// +// 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]. +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 +type Metadata struct { + Branch uint32 // shard identifier (Python: Branch = uint32) + ClusterPeerID uint64 // 0 = cluster RPC (master commands), ≠0 = specific external peer +} + +// Frame is a complete protocol frame. +// raw_data layout on wire: [1B opcode][8B rpc_id][N bytes payload] +type Frame struct { + Meta Metadata + Opcode byte + RPCID uint64 + Payload []byte +} + +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 +) + +// 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). +func ReadFrame(r io.Reader) (*Frame, error) { + return readFrameWithMetaSize(r, metaSize) +} + +// 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. +func ReadFrameNoMeta(r io.Reader) (*Frame, error) { + return readFrameWithMetaSize(r, 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) { + // 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) + } + + // 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 { + return nil, fmt.Errorf("unsupported metaSize %d (only 0 or 12 supported)", metaSize) + } + metaBuf := make([]byte, metaSize) + if _, err := io.ReadFull(r, metaBuf); err != nil { + return nil, fmt.Errorf("reading metadata: %w", err) + } + meta = Metadata{ + Branch: binary.BigEndian.Uint32(metaBuf[0:4]), + ClusterPeerID: binary.BigEndian.Uint64(metaBuf[4:12]), + } + } + + // 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 { + return nil, fmt.Errorf("reading frame body (payload_len=%d): %w", payloadLen, err) + } + + return &Frame{ + Meta: meta, + Opcode: body[0], + RPCID: binary.BigEndian.Uint64(body[1:9]), + Payload: body[9:], + }, 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. +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. +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] +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 + if metaSize > 0 { + binary.BigEndian.PutUint32(buf[frameHeader:frameHeader+4], f.Meta.Branch) + binary.BigEndian.PutUint64(buf[frameHeader+4:frameHeader+metaSize], f.Meta.ClusterPeerID) + } + + // 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 + copy(buf[frameHeader+metaSize+opcodeSize+rpcIDSize:], f.Payload) + + _, err := w.Write(buf) + return err +} + +// MarshalMetadata serializes Metadata into its 12-byte wire representation. +func MarshalMetadata(m Metadata) []byte { + buf := make([]byte, metaSize) + binary.BigEndian.PutUint32(buf[0:4], m.Branch) + binary.BigEndian.PutUint64(buf[4:12], m.ClusterPeerID) + return buf +} + +// UnmarshalMetadata deserializes a 12-byte wire representation into Metadata. +func UnmarshalMetadata(b []byte) (Metadata, error) { + if len(b) != metaSize { + return Metadata{}, fmt.Errorf("metadata must be %d bytes, got %d", metaSize, len(b)) + } + return Metadata{ + Branch: binary.BigEndian.Uint32(b[0:4]), + 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 new file mode 100644 index 0000000000..dceca81790 --- /dev/null +++ b/qkc/slave/frame_test.go @@ -0,0 +1,360 @@ +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 +// ============================================================================= + +// 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) + 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.RPCID != 1 { + t.Errorf("RPCID: got %d, want 1", f.RPCID) + } + if f.Meta != (Metadata{Branch: 0, ClusterPeerID: 0}) { + t.Errorf("Meta: got %+v", f.Meta) + } + if len(f.Payload) != 0 { + t.Errorf("Payload len: got %d, want 0", len(f.Payload)) + } +} + +func TestPythonRead_Pong(t *testing.T) { + wire, _ := hex.DecodeString(pongMasterWire) + 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.RPCID != 1 { + t.Errorf("RPCID: got %d, want 1", f.RPCID) + } +} + +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) + 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.RPCID != 42 { + t.Errorf("RPCID: got %d, want 42", f.RPCID) + } + if len(f.Payload) != 56 { + t.Errorf("Payload len: got %d, want 56", len(f.Payload)) + } +} + +func TestPythonRead_LargePayload(t *testing.T) { + wire, _ := hex.DecodeString(largePayloadWire) + f, err := ReadFrame(bytes.NewReader(wire)) + if err != nil { + t.Fatalf("ReadFrame: %v", err) + } + if f.Opcode != 0x10 { + t.Errorf("Opcode: got 0x%02x, want 0x10", f.Opcode) + } + if f.RPCID != 7 { + t.Errorf("RPCID: got %d, want 7", f.RPCID) + } + 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 + } + } +} + +// ============================================================================= +// 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 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 + f *Frame + }{ + {"empty", &Frame{Opcode: 1, RPCID: 0, Payload: nil}}, + {"with_meta", &Frame{Meta: Metadata{Branch: 5, ClusterPeerID: 999}, Opcode: 0x10, RPCID: 7, Payload: []byte("hello")}}, + {"large_rpc_id", &Frame{Opcode: 0xC4, RPCID: 0xFFFFFFFFFFFFFFFF, Payload: []byte("x")}}, + {"zero_meta", &Frame{Meta: Metadata{}, Opcode: 0x81, RPCID: 1, Payload: []byte{}}}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + wire := writeFrameForTest(tc.f.Meta, tc.f.Opcode, tc.f.RPCID, tc.f.Payload) + got, err := ReadFrame(bytes.NewReader(wire)) + if err != nil { + t.Fatalf("ReadFrame: %v", err) + } + if got.Opcode != tc.f.Opcode || got.RPCID != tc.f.RPCID || got.Meta != tc.f.Meta { + t.Errorf("mismatch: got %+v, want %+v", got, tc.f) + } + if !bytes.Equal(got.Payload, tc.f.Payload) { + t.Errorf("payload mismatch") + } + }) + } +} + +func TestRoundTrip_NoMeta(t *testing.T) { + original := &Frame{Opcode: 0x93, RPCID: 99, Payload: []byte("xshard-data")} + wire := writeFrameNoMetaForTest(original.Opcode, original.RPCID, original.Payload) + got, err := ReadFrameNoMeta(bytes.NewReader(wire)) + if err != nil { + t.Fatalf("ReadFrameNoMeta: %v", err) + } + if got.Opcode != original.Opcode || got.RPCID != original.RPCID { + t.Errorf("mismatch: got %+v", got) + } + if !bytes.Equal(got.Payload, original.Payload) { + t.Errorf("payload mismatch") + } +} + +// ============================================================================= +// 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}, + Opcode: 0x42, + RPCID: 0xDEADBEEFCAFEBABE, + Payload: []byte{0xAA, 0xBB, 0xCC}, + } + 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")}, + {Meta: Metadata{Branch: 1, ClusterPeerID: 0}, Opcode: 0x03, RPCID: 200, Payload: nil}, + } + var stream bytes.Buffer + for _, f := range frames { + stream.Write(writeFrameForTest(f.Meta, f.Opcode, f.RPCID, f.Payload)) + } + + reader := bytes.NewReader(stream.Bytes()) + for i, want := range frames { + got, err := ReadFrame(reader) + if err != nil { + t.Fatalf("frame %d: %v", i, err) + } + if got.Opcode != want.Opcode || got.RPCID != want.RPCID || got.Meta != want.Meta { + t.Errorf("frame %d mismatch: got %+v, want %+v", i, got, want) + } + if !bytes.Equal(got.Payload, want.Payload) { + t.Errorf("frame %d payload mismatch", i) + } + } +} + +// ============================================================================= +// Error handling +// ============================================================================= + +func TestReadFrame_EOF(t *testing.T) { + if _, err := ReadFrame(bytes.NewReader(nil)); err == nil { + t.Error("expected error on empty stream") + } +} + +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 { + t.Error("expected error on truncated frame") + } +} + +// ============================================================================= +// Helpers +// ============================================================================= + +// 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}) + return buf.Bytes() +} + +func writeFrameNoMetaForTest(opcode byte, rpcID uint64, payload []byte) []byte { + var buf bytes.Buffer + _ = WriteFrameNoMeta(&buf, &Frame{Opcode: opcode, RPCID: rpcID, Payload: payload}) + return buf.Bytes() +} From 43d0e0e5802c55e6c7cf596ceb3c9555765f709d Mon Sep 17 00:00:00 2001 From: iteye Date: Thu, 2 Jul 2026 14:57:26 +0800 Subject: [PATCH 2/4] fix comment --- qkc/slave/frame.go | 172 +++++++----------------- qkc/slave/frame_test.go | 284 +++++++++++++--------------------------- 2 files changed, 142 insertions(+), 314 deletions(-) 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}) From edf7f7b473b0abe4393ad805a38aff4988194d0e Mon Sep 17 00:00:00 2001 From: iteye Date: Thu, 2 Jul 2026 15:15:56 +0800 Subject: [PATCH 3/4] change cluster lib --- qkc/{slave => cluster/wire}/frame.go | 46 +++++------------------ qkc/{slave => cluster/wire}/frame_test.go | 44 ++++++++++++++++++---- qkc/cluster/wire/metadata.go | 38 +++++++++++++++++++ 3 files changed, 84 insertions(+), 44 deletions(-) rename qkc/{slave => cluster/wire}/frame.go (76%) rename qkc/{slave => cluster/wire}/frame_test.go (83%) create mode 100644 qkc/cluster/wire/metadata.go diff --git a/qkc/slave/frame.go b/qkc/cluster/wire/frame.go similarity index 76% rename from qkc/slave/frame.go rename to qkc/cluster/wire/frame.go index 063c35d2ee..8f06338efb 100644 --- a/qkc/slave/frame.go +++ b/qkc/cluster/wire/frame.go @@ -1,6 +1,6 @@ // Copyright 2026-2027, QuarkChain. -// -// Package slave implements a binary frame codec compatible with pyquarkchain. + +// Package wire implements a binary frame codec compatible with pyquarkchain. // // Wire format (per frame): // @@ -10,7 +10,7 @@ // 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 +package wire import ( "encoding/binary" @@ -19,17 +19,10 @@ import ( "io" ) -// Metadata is the routing header for ClusterMetadata traffic. -// Wire representation: 12 bytes (4B branch + 8B cluster_peer_id). -type Metadata struct { - Branch uint32 - ClusterPeerID uint64 -} - // Frame is a complete protocol frame. // Wire layout after metadata: [1B opcode][8B rpc_id][N bytes payload] type Frame struct { - Meta Metadata + Meta ClusterMetadata Opcode byte RPCID uint64 Payload []byte @@ -42,12 +35,12 @@ const ( frameHeader = 4 // payload_len prefix ) -// ReadFrame reads a frame with 12-byte metadata and no payload-length cap. +// ReadFrame reads a frame with 12-byte ClusterMetadata and no payload-length cap. func ReadFrame(r io.Reader) (*Frame, error) { return readFrame(r, metaSize, 0) } -// ReadFrameWithMaxPayload reads a frame with 12-byte metadata and rejects frames +// ReadFrameWithMaxPayload reads a frame with 12-byte ClusterMetadata and rejects frames // whose payload_len exceeds maxPayloadLen before allocation. func ReadFrameWithMaxPayload(r io.Reader, maxPayloadLen uint32) (*Frame, error) { if maxPayloadLen == 0 { @@ -83,16 +76,16 @@ func readFrame(r io.Reader, metaSize int, maxPayloadLen uint32) (*Frame, error) } // 2. Read metadata (size depends on metadata_class) - var meta Metadata + var meta ClusterMetadata if metaSize > 0 { if metaSize != 12 { - return nil, fmt.Errorf("unsupported metaSize %d (only 0 or 12 supported)", metaSize) + return nil, fmt.Errorf("unsupported metaSize %d (only 12 supported)", metaSize) } metaBuf := make([]byte, metaSize) if _, err := io.ReadFull(r, metaBuf); err != nil { return nil, fmt.Errorf("reading metadata: %w", err) } - meta = Metadata{ + meta = ClusterMetadata{ Branch: binary.BigEndian.Uint32(metaBuf[0:4]), ClusterPeerID: binary.BigEndian.Uint64(metaBuf[4:12]), } @@ -113,7 +106,7 @@ func readFrame(r io.Reader, metaSize int, maxPayloadLen uint32) (*Frame, error) }, nil } -// WriteFrame serializes f with 12-byte metadata and writes it to w. +// WriteFrame serializes f with 12-byte ClusterMetadata and writes it to w. func WriteFrame(w io.Writer, f *Frame) error { return writeFrameWithMetaSize(w, f, metaSize) } @@ -152,22 +145,3 @@ func writeFrameWithMetaSize(w io.Writer, f *Frame, metaSize int) error { _, err := w.Write(buf) return err } - -// MarshalMetadata serializes Metadata into its 12-byte wire representation. -func MarshalMetadata(m Metadata) []byte { - buf := make([]byte, metaSize) - binary.BigEndian.PutUint32(buf[0:4], m.Branch) - binary.BigEndian.PutUint64(buf[4:12], m.ClusterPeerID) - return buf -} - -// UnmarshalMetadata deserializes a 12-byte wire representation into Metadata. -func UnmarshalMetadata(b []byte) (Metadata, error) { - if len(b) != metaSize { - return Metadata{}, fmt.Errorf("metadata must be %d bytes, got %d", metaSize, len(b)) - } - return Metadata{ - Branch: binary.BigEndian.Uint32(b[0:4]), - ClusterPeerID: binary.BigEndian.Uint64(b[4:12]), - }, nil -} diff --git a/qkc/slave/frame_test.go b/qkc/cluster/wire/frame_test.go similarity index 83% rename from qkc/slave/frame_test.go rename to qkc/cluster/wire/frame_test.go index cfc6879928..5019de1b19 100644 --- a/qkc/slave/frame_test.go +++ b/qkc/cluster/wire/frame_test.go @@ -1,6 +1,6 @@ // Copyright 2026-2027, QuarkChain. -package slave +package wire import ( "bytes" @@ -11,6 +11,7 @@ import ( var ( // Test vectors serialized by pyquarkchain serializers using synthetic test values. + // NOT real production network data. // // Payload content (Ping/Pong commands): // - id = "id" (ASCII, 2 bytes) @@ -116,9 +117,9 @@ func TestRoundTrip_Meta(t *testing.T) { f *Frame }{ {"empty", &Frame{Opcode: 1, RPCID: 0, Payload: nil}}, - {"with_meta", &Frame{Meta: Metadata{Branch: 5, ClusterPeerID: 999}, Opcode: 0x10, RPCID: 7, Payload: []byte("hello")}}, + {"with_meta", &Frame{Meta: ClusterMetadata{Branch: 5, ClusterPeerID: 999}, Opcode: 0x10, RPCID: 7, Payload: []byte("hello")}}, {"large_rpc_id", &Frame{Opcode: 0xC4, RPCID: 0xFFFFFFFFFFFFFFFF, Payload: []byte("x")}}, - {"zero_meta", &Frame{Meta: Metadata{}, Opcode: 0x81, RPCID: 1, Payload: []byte{}}}, + {"zero_meta", &Frame{Meta: ClusterMetadata{}, Opcode: 0x81, RPCID: 1, Payload: []byte{}}}, } for _, tc := range cases { t.Run(tc.name, func(t *testing.T) { @@ -154,7 +155,7 @@ func TestRoundTrip_NoMeta(t *testing.T) { func TestWireFormatLayout(t *testing.T) { f := &Frame{ - Meta: Metadata{Branch: 1, ClusterPeerID: 0x1122334455667788}, + Meta: ClusterMetadata{Branch: 1, ClusterPeerID: 0x1122334455667788}, Opcode: 0x42, RPCID: 0xDEADBEEFCAFEBABE, Payload: []byte{0xAA, 0xBB, 0xCC}, @@ -183,9 +184,9 @@ func TestWireFormatLayout(t *testing.T) { func TestMultiFrameStream(t *testing.T) { 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")}, - {Meta: Metadata{Branch: 1, ClusterPeerID: 0}, Opcode: 0x03, RPCID: 200, Payload: nil}, + {Meta: ClusterMetadata{Branch: 0, ClusterPeerID: 0}, Opcode: 0x81, RPCID: 0, Payload: []byte("ping")}, + {Meta: ClusterMetadata{Branch: 2, ClusterPeerID: 999}, Opcode: 0x05, RPCID: 100, Payload: []byte("block_data")}, + {Meta: ClusterMetadata{Branch: 1, ClusterPeerID: 0}, Opcode: 0x03, RPCID: 200, Payload: nil}, } var stream bytes.Buffer for _, f := range frames { @@ -251,7 +252,34 @@ func TestReadFrameNoMeta_PayloadLimit(t *testing.T) { } } -func writeFrameForTest(meta Metadata, opcode byte, rpcID uint64, payload []byte) []byte { +func TestClusterMetadata_RoundTrip(t *testing.T) { + cases := []ClusterMetadata{ + {Branch: 0, ClusterPeerID: 0}, + {Branch: 1, ClusterPeerID: 12345}, + {Branch: 0xFFFFFFFF, ClusterPeerID: 0xFFFFFFFFFFFFFFFF}, + } + for _, m := range cases { + wire := MarshalClusterMetadata(m) + got, err := UnmarshalClusterMetadata(wire) + if err != nil { + t.Fatalf("UnmarshalClusterMetadata(%+v): %v", m, err) + } + if got != m { + t.Errorf("round-trip mismatch: got %+v, want %+v", got, m) + } + } +} + +func TestUnmarshalClusterMetadata_InvalidLength(t *testing.T) { + for _, n := range []int{0, 4, 8, 11, 13, 16} { + b := make([]byte, n) + if _, err := UnmarshalClusterMetadata(b); err == nil { + t.Errorf("expected error for %d-byte input", n) + } + } +} + +func writeFrameForTest(meta ClusterMetadata, opcode byte, rpcID uint64, payload []byte) []byte { var buf bytes.Buffer _ = WriteFrame(&buf, &Frame{Meta: meta, Opcode: opcode, RPCID: rpcID, Payload: payload}) return buf.Bytes() diff --git a/qkc/cluster/wire/metadata.go b/qkc/cluster/wire/metadata.go new file mode 100644 index 0000000000..89a188375a --- /dev/null +++ b/qkc/cluster/wire/metadata.go @@ -0,0 +1,38 @@ +// Copyright 2026-2027, QuarkChain. + +package wire + +import ( + "encoding/binary" + "fmt" +) + +// ClusterMetadata is the routing header for intra-cluster (master↔slave) traffic. +// Wire representation: 12 bytes (4B branch + 8B cluster_peer_id). +// +// Matches pyquarkchain's ClusterMetadata class (protocol.py). +// The base Metadata in pyquarkchain is the 0-byte variant used by direct +// slave-to-slave connections, handled by ReadFrameNoMeta / WriteFrameNoMeta. +type ClusterMetadata struct { + Branch uint32 + ClusterPeerID uint64 +} + +// MarshalClusterMetadata serializes ClusterMetadata into its 12-byte wire representation. +func MarshalClusterMetadata(m ClusterMetadata) []byte { + buf := make([]byte, metaSize) + binary.BigEndian.PutUint32(buf[0:4], m.Branch) + binary.BigEndian.PutUint64(buf[4:12], m.ClusterPeerID) + return buf +} + +// UnmarshalClusterMetadata deserializes a 12-byte wire representation into ClusterMetadata. +func UnmarshalClusterMetadata(b []byte) (ClusterMetadata, error) { + if len(b) != metaSize { + return ClusterMetadata{}, fmt.Errorf("metadata must be %d bytes, got %d", metaSize, len(b)) + } + return ClusterMetadata{ + Branch: binary.BigEndian.Uint32(b[0:4]), + ClusterPeerID: binary.BigEndian.Uint64(b[4:12]), + }, nil +} From e1d27be1f819f92930cba0330643a24cc714b084 Mon Sep 17 00:00:00 2001 From: iteye Date: Thu, 2 Jul 2026 17:31:07 +0800 Subject: [PATCH 4/4] merge frame function --- qkc/cluster/wire/frame.go | 28 +--- qkc/cluster/wire/frame_test.go | 226 +++++++++++++++++---------------- 2 files changed, 126 insertions(+), 128 deletions(-) diff --git a/qkc/cluster/wire/frame.go b/qkc/cluster/wire/frame.go index 8f06338efb..eb66f23564 100644 --- a/qkc/cluster/wire/frame.go +++ b/qkc/cluster/wire/frame.go @@ -35,31 +35,15 @@ const ( frameHeader = 4 // payload_len prefix ) -// ReadFrame reads a frame with 12-byte ClusterMetadata and no payload-length cap. -func ReadFrame(r io.Reader) (*Frame, error) { - return readFrame(r, metaSize, 0) -} - -// ReadFrameWithMaxPayload reads a frame with 12-byte ClusterMetadata 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") - } +// ReadFrame reads a frame with 12-byte ClusterMetadata. +// maxPayloadLen == 0 disables payload-size checking. +func ReadFrame(r io.Reader, maxPayloadLen uint32) (*Frame, error) { 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 readFrame(r, 0, 0) -} - -// 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") - } +// ReadFrameNoMeta reads a frame with 0-byte metadata. +// maxPayloadLen == 0 disables payload-size checking. +func ReadFrameNoMeta(r io.Reader, maxPayloadLen uint32) (*Frame, error) { return readFrame(r, 0, maxPayloadLen) } diff --git a/qkc/cluster/wire/frame_test.go b/qkc/cluster/wire/frame_test.go index 5019de1b19..4b75e5fcc0 100644 --- a/qkc/cluster/wire/frame_test.go +++ b/qkc/cluster/wire/frame_test.go @@ -6,6 +6,7 @@ import ( "bytes" "encoding/binary" "encoding/hex" + "io" "testing" ) @@ -19,14 +20,14 @@ var ( // - root_tip = None (Ping only) // - opcode = 0x81 (PING), 0x82 (PONG) from ClusterOp (CLUSTER_OP_BASE=128) pythonClusterPing = "0000001300000001000000000000303981000000000000000100000002696400000002000000010000000200" - pythonClusterPong = "00000012000000010000000000003039820000000000000001000000026964000000020000000100000002" pythonNoMetaPing = "0000001381000000000000000100000002696400000002000000010000000200" - pythonNoMetaPong = "00000012820000000000000001000000026964000000020000000100000002" ) -func TestPythonVectors_ClusterPing(t *testing.T) { - wire := mustPythonVectorBytes(t, pythonClusterPing) - f, err := ReadFrame(bytes.NewReader(wire)) +// ---- wire compatibility (pyquarkchain vectors) ---- + +func TestWireCompat_Meta(t *testing.T) { + wire := mustHex(t, pythonClusterPing) + f, err := ReadFrame(bytes.NewReader(wire), 0) if err != nil { t.Fatalf("ReadFrame: %v", err) } @@ -45,30 +46,9 @@ func TestPythonVectors_ClusterPing(t *testing.T) { } } -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 || f.RPCID != 1 { - t.Fatalf("unexpected header: opcode=0x%02x rpcid=%d", f.Opcode, 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 TestPythonVectors_NoMetaPing(t *testing.T) { - wire := mustPythonVectorBytes(t, pythonNoMetaPing) - f, err := ReadFrameNoMeta(bytes.NewReader(wire)) +func TestWireCompat_NoMeta(t *testing.T) { + wire := mustHex(t, pythonNoMetaPing) + f, err := ReadFrameNoMeta(bytes.NewReader(wire), 0) if err != nil { t.Fatalf("ReadFrameNoMeta: %v", err) } @@ -84,74 +64,39 @@ func TestPythonVectors_NoMetaPing(t *testing.T) { } } -func TestPythonVectors_NoMetaPong(t *testing.T) { - wire := mustPythonVectorBytes(t, pythonNoMetaPong) - f, err := ReadFrameNoMeta(bytes.NewReader(wire)) - if err != nil { - t.Fatalf("ReadFrameNoMeta: %v", err) - } - if f.Opcode != 0x82 || f.RPCID != 1 { - t.Fatalf("unexpected header: opcode=0x%02x rpcid=%d", f.Opcode, f.RPCID) - } - var out bytes.Buffer - if err := WriteFrameNoMeta(&out, f); err != nil { - t.Fatalf("WriteFrameNoMeta: %v", err) - } - if !bytes.Equal(out.Bytes(), wire) { - t.Fatalf("frame mismatch:\n go %x\n py %x", out.Bytes(), wire) - } -} +// ---- round-trip: write → read ---- -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 TestRoundTrip_Meta(t *testing.T) { +func TestRoundTrip(t *testing.T) { cases := []struct { - name string - f *Frame + name string + meta ClusterMetadata + opcode byte + rpcID uint64 + payload []byte }{ - {"empty", &Frame{Opcode: 1, RPCID: 0, Payload: nil}}, - {"with_meta", &Frame{Meta: ClusterMetadata{Branch: 5, ClusterPeerID: 999}, Opcode: 0x10, RPCID: 7, Payload: []byte("hello")}}, - {"large_rpc_id", &Frame{Opcode: 0xC4, RPCID: 0xFFFFFFFFFFFFFFFF, Payload: []byte("x")}}, - {"zero_meta", &Frame{Meta: ClusterMetadata{}, Opcode: 0x81, RPCID: 1, Payload: []byte{}}}, + {"nil_payload", ClusterMetadata{}, 0x01, 0, nil}, + {"with_meta", ClusterMetadata{Branch: 5, ClusterPeerID: 999}, 0x10, 7, []byte("hello")}, + {"max_rpc_id", ClusterMetadata{}, 0xC4, 0xFFFFFFFFFFFFFFFF, []byte("x")}, + {"empty_payload", ClusterMetadata{}, 0x81, 1, []byte{}}, } for _, tc := range cases { t.Run(tc.name, func(t *testing.T) { - wire := writeFrameForTest(tc.f.Meta, tc.f.Opcode, tc.f.RPCID, tc.f.Payload) - got, err := ReadFrame(bytes.NewReader(wire)) + wire := writeFrameForTest(tc.meta, tc.opcode, tc.rpcID, tc.payload) + got, err := ReadFrame(bytes.NewReader(wire), 0) if err != nil { t.Fatalf("ReadFrame: %v", err) } - if got.Opcode != tc.f.Opcode || got.RPCID != tc.f.RPCID || got.Meta != tc.f.Meta { - t.Errorf("mismatch: got %+v, want %+v", got, tc.f) + if got.Opcode != tc.opcode || got.RPCID != tc.rpcID || got.Meta != tc.meta { + t.Errorf("mismatch: got %+v, want Meta=%+v Opcode=%x RPCID=%d", got, tc.meta, tc.opcode, tc.rpcID) } - if !bytes.Equal(got.Payload, tc.f.Payload) { + if !bytes.Equal(got.Payload, tc.payload) { t.Errorf("payload mismatch") } }) } } -func TestRoundTrip_NoMeta(t *testing.T) { - original := &Frame{Opcode: 0x93, RPCID: 99, Payload: []byte("xshard-data")} - wire := writeFrameNoMetaForTest(original.Opcode, original.RPCID, original.Payload) - got, err := ReadFrameNoMeta(bytes.NewReader(wire)) - if err != nil { - t.Fatalf("ReadFrameNoMeta: %v", err) - } - if got.Opcode != original.Opcode || got.RPCID != original.RPCID { - t.Errorf("mismatch: got %+v", got) - } - if !bytes.Equal(got.Payload, original.Payload) { - t.Errorf("payload mismatch") - } -} +// ---- wire format layout ---- func TestWireFormatLayout(t *testing.T) { f := &Frame{ @@ -182,6 +127,8 @@ func TestWireFormatLayout(t *testing.T) { } } +// ---- multi-frame stream ---- + func TestMultiFrameStream(t *testing.T) { frames := []*Frame{ {Meta: ClusterMetadata{Branch: 0, ClusterPeerID: 0}, Opcode: 0x81, RPCID: 0, Payload: []byte("ping")}, @@ -195,7 +142,7 @@ func TestMultiFrameStream(t *testing.T) { reader := bytes.NewReader(stream.Bytes()) for i, want := range frames { - got, err := ReadFrame(reader) + got, err := ReadFrame(reader, 0) if err != nil { t.Fatalf("frame %d: %v", i, err) } @@ -208,51 +155,108 @@ func TestMultiFrameStream(t *testing.T) { } } -func TestReadFrame_EOF(t *testing.T) { - if _, err := ReadFrame(bytes.NewReader(nil)); err == nil { - t.Error("expected error on empty stream") +// ---- error paths ---- + +func TestReadErrors(t *testing.T) { + cases := []struct { + name string + r io.Reader + read func(io.Reader, uint32) (*Frame, error) + }{ + {"meta_empty_stream", bytes.NewReader(nil), ReadFrame}, + {"meta_truncated", bytes.NewReader(truncatedHeader()), ReadFrame}, + {"nometa_empty_stream", bytes.NewReader(nil), ReadFrameNoMeta}, + {"nometa_truncated", bytes.NewReader(truncatedHeader()), ReadFrameNoMeta}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + if _, err := tc.read(tc.r, 0); err == nil { + t.Error("expected error, got nil") + } + }) } } -func TestReadFrame_Truncated(t *testing.T) { +func truncatedHeader() []byte { hdr := make([]byte, 4) binary.BigEndian.PutUint32(hdr, 100) - if _, err := ReadFrame(bytes.NewReader(hdr)); err == nil { - t.Error("expected error on truncated frame") - } + return hdr } +// ---- payload size limit ---- + func TestReadFrame_PayloadLimit(t *testing.T) { - hdr := make([]byte, 4) - binary.BigEndian.PutUint32(hdr, 100) + payload := bytes.Repeat([]byte{0xAB}, 100) + frame := &Frame{Meta: ClusterMetadata{Branch: 1, ClusterPeerID: 2}, Opcode: 0x81, RPCID: 1, Payload: payload} - if _, err := ReadFrameWithMaxPayload(bytes.NewReader(hdr), 64); err == nil { - t.Fatal("expected error when payload_len exceeds limit") + // limit == 0: unbounded — full frame with 100-byte payload is accepted. + { + wire := writeFrameForTest(frame.Meta, frame.Opcode, frame.RPCID, frame.Payload) + got, err := ReadFrame(bytes.NewReader(wire), 0) + if err != nil { + t.Fatalf("limit=0 should accept full frame, got: %v", err) + } + if !bytes.Equal(got.Payload, payload) { + t.Fatalf("limit=0 payload mismatch") + } } - if _, err := ReadFrameWithMaxPayload(bytes.NewReader(hdr), 0); err == nil { - t.Fatal("expected error when maxPayloadLen is zero") + // limit == payload length: boundary, accepted. + { + wire := writeFrameForTest(frame.Meta, frame.Opcode, frame.RPCID, frame.Payload) + if _, err := ReadFrame(bytes.NewReader(wire), uint32(len(payload))); err != nil { + t.Fatalf("limit==payload_len should accept, got: %v", err) + } } - if _, err := ReadFrame(bytes.NewReader(hdr)); err == nil { - t.Error("expected error on truncated frame without payload limit") + // limit < payload length: rejected before metadata is read. + { + wire := writeFrameForTest(frame.Meta, frame.Opcode, frame.RPCID, frame.Payload) + _, err := ReadFrame(bytes.NewReader(wire), uint32(len(payload)-1)) + if err == nil { + t.Fatal("limit < payload_len should be rejected") + } } } func TestReadFrameNoMeta_PayloadLimit(t *testing.T) { - hdr := make([]byte, 4) - binary.BigEndian.PutUint32(hdr, 32) + payload := bytes.Repeat([]byte{0xCD}, 32) + frame := &Frame{Opcode: 0x82, RPCID: 7, Payload: payload} - if _, err := ReadFrameNoMetaWithMaxPayload(bytes.NewReader(hdr), 16); err == nil { - t.Fatal("expected error when payload_len exceeds limit") + // limit == 0: unbounded. + { + wire := writeFrameNoMetaForTest(frame.Opcode, frame.RPCID, frame.Payload) + got, err := ReadFrameNoMeta(bytes.NewReader(wire), 0) + if err != nil { + t.Fatalf("limit=0 should accept full frame, got: %v", err) + } + if !bytes.Equal(got.Payload, payload) { + t.Fatalf("limit=0 payload mismatch") + } } - if _, err := ReadFrameNoMetaWithMaxPayload(bytes.NewReader(hdr), 0); err == nil { - t.Fatal("expected error when maxPayloadLen is zero") + // limit == payload length: boundary, accepted. + { + wire := writeFrameNoMetaForTest(frame.Opcode, frame.RPCID, frame.Payload) + if _, err := ReadFrameNoMeta(bytes.NewReader(wire), uint32(len(payload))); err != nil { + t.Fatalf("limit==payload_len should accept, got: %v", err) + } + } + + // limit < payload length: rejected. + { + wire := writeFrameNoMetaForTest(frame.Opcode, frame.RPCID, frame.Payload) + _, err := ReadFrameNoMeta(bytes.NewReader(wire), uint32(len(payload)-1)) + if err == nil { + t.Fatal("limit < payload_len should be rejected") + } } } -func TestClusterMetadata_RoundTrip(t *testing.T) { +// ---- ClusterMetadata serialization ---- + +func TestClusterMetadata(t *testing.T) { + // Round-trip: edge cases. cases := []ClusterMetadata{ {Branch: 0, ClusterPeerID: 0}, {Branch: 1, ClusterPeerID: 12345}, @@ -268,9 +272,8 @@ func TestClusterMetadata_RoundTrip(t *testing.T) { t.Errorf("round-trip mismatch: got %+v, want %+v", got, m) } } -} -func TestUnmarshalClusterMetadata_InvalidLength(t *testing.T) { + // Invalid lengths. for _, n := range []int{0, 4, 8, 11, 13, 16} { b := make([]byte, n) if _, err := UnmarshalClusterMetadata(b); err == nil { @@ -279,6 +282,17 @@ func TestUnmarshalClusterMetadata_InvalidLength(t *testing.T) { } } +// ---- helpers ---- + +func mustHex(t *testing.T, hexStr string) []byte { + t.Helper() + b, err := hex.DecodeString(hexStr) + if err != nil { + t.Fatalf("decode hex: %v", err) + } + return b +} + func writeFrameForTest(meta ClusterMetadata, opcode byte, rpcID uint64, payload []byte) []byte { var buf bytes.Buffer _ = WriteFrame(&buf, &Frame{Meta: meta, Opcode: opcode, RPCID: rpcID, Payload: payload})