fix comment

This commit is contained in:
iteye 2026-07-02 14:57:26 +08:00
parent 58b66d610e
commit 43d0e0e580
2 changed files with 142 additions and 314 deletions

View file

@ -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 / // [4B payload_len] [metaSize B metadata] [1B opcode] [8B rpc_id] [payload bytes]
// Connection.write_raw_data exactly:
// //
// Python: protocol.py lines 285-308 // payload_len is the length of payload bytes only.
// // metadata size depends on metaSize parameter:
// Metadata sizes (matching Python Metadata subclasses): // - 12 bytes: ClusterMetadata (branch uint32 + cluster_peer_id uint64) for master↔slave
// // - 0 bytes: no metadata for slave↔slave traffic
// 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 package slave
import ( import (
"bufio"
"encoding/binary" "encoding/binary"
"errors" "errors"
"fmt" "fmt"
"io" "io"
) )
// Metadata is the 12-byte frame header carrying routing information. // Metadata is the routing header for ClusterMetadata traffic.
// Matches Python's ClusterMetadata: // Wire representation: 12 bytes (4B branch + 8B cluster_peer_id).
//
// class ClusterMetadata(Metadata):
// FIELDS = [("branch", Branch), ("cluster_peer_id", uint64)]
// @staticmethod
// def get_byte_size(): return 12
type Metadata struct { type Metadata struct {
Branch uint32 // shard identifier (Python: Branch = uint32) Branch uint32
ClusterPeerID uint64 // 0 = cluster RPC (master commands), ≠0 = specific external peer ClusterPeerID uint64
} }
// Frame is a complete protocol frame. // 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 { type Frame struct {
Meta Metadata Meta Metadata
Opcode byte Opcode byte
@ -59,53 +36,53 @@ type Frame struct {
} }
const ( const (
metaSize = 12 // ClusterMetadata.get_byte_size() = 12 (branch 4B + cluster_peer_id 8B) metaSize = 12 // ClusterMetadata.get_byte_size() = 12 (branch 4B + cluster_peer_id 8B)
opcodeSize = 1 opcodeSize = 1
rpcIDSize = 8 rpcIDSize = 8
frameHeader = 4 // payload_len prefix frameHeader = 4 // payload_len prefix
totalOverhead = frameHeader + metaSize + opcodeSize + rpcIDSize // 4+12+1+8 = 25
) )
// ReadFrame reads a frame with 12-byte ClusterMetadata. // ReadFrame reads a frame with 12-byte metadata and no payload-length cap.
//
// 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) { func ReadFrame(r io.Reader) (*Frame, error) {
return readFrameWithMetaSize(r, metaSize) return readFrame(r, metaSize, 0)
} }
// ReadFrameNoMeta reads a frame with 0-byte Metadata. // ReadFrameWithMaxPayload reads a frame with 12-byte metadata and rejects frames
// // whose payload_len exceeds maxPayloadLen before allocation.
// Matches Python's SlaveConnection which uses metadata_class = Metadata func ReadFrameWithMaxPayload(r io.Reader, maxPayloadLen uint32) (*Frame, error) {
// (get_byte_size() == 0). if maxPayloadLen == 0 {
// return nil, errors.New("maxPayloadLen must be greater than zero")
// Used for slave ↔ slave direct TCP traffic. }
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) { func ReadFrameNoMeta(r io.Reader) (*Frame, error) {
return readFrameWithMetaSize(r, 0) return readFrame(r, 0, 0)
} }
// readFrameWithMetaSize is the underlying frame reader. // ReadFrameNoMetaWithMaxPayload reads a frame with 0-byte metadata and rejects
// // frames whose payload_len exceeds maxPayloadLen before allocation.
// Wire layout (matching Python read_metadata_and_raw_data): func ReadFrameNoMetaWithMaxPayload(r io.Reader, maxPayloadLen uint32) (*Frame, error) {
// if maxPayloadLen == 0 {
// size_bytes = await read_fully(4) → 4B payload_len (big-endian) return nil, errors.New("maxPayloadLen must be greater than zero")
// metadata_bytes = await read_fully(metaSize) → metaSize B metadata }
// raw_data_without_size = await read_fully(1+8+size) → opcode + rpc_id + payload return readFrame(r, 0, maxPayloadLen)
func readFrameWithMetaSize(r io.Reader, metaSize int) (*Frame, error) { }
// 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 // 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 var payloadLen uint32
if err := binary.Read(r, binary.BigEndian, &payloadLen); err != nil { if err := binary.Read(r, binary.BigEndian, &payloadLen); err != nil {
return nil, fmt.Errorf("reading frame length: %w", err) 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) // 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 var meta Metadata
if metaSize > 0 { if metaSize > 0 {
if metaSize != 12 { 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] // 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) bodySize := opcodeSize + rpcIDSize + int(payloadLen)
body := make([]byte, bodySize) body := make([]byte, bodySize)
if _, err := io.ReadFull(r, body); err != nil { if _, err := io.ReadFull(r, body); err != nil {
@ -137,47 +113,26 @@ func readFrameWithMetaSize(r io.Reader, metaSize int) (*Frame, error) {
}, nil }, nil
} }
// WriteFrame serializes f with 12-byte ClusterMetadata and writes it to w. // WriteFrame serializes f with 12-byte metadata 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 { func WriteFrame(w io.Writer, f *Frame) error {
return writeFrameWithMetaSize(w, f, metaSize) return writeFrameWithMetaSize(w, f, metaSize)
} }
// WriteFrameNoMeta writes a frame with 0-byte Metadata. // 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 { func WriteFrameNoMeta(w io.Writer, f *Frame) error {
return writeFrameWithMetaSize(w, f, 0) return writeFrameWithMetaSize(w, f, 0)
} }
// writeFrameWithMetaSize serializes f with the given metadata size and writes // writeFrameWithMetaSize serializes f with the given metadata size and writes it to w.
// 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 { func writeFrameWithMetaSize(w io.Writer, f *Frame, metaSize int) error {
payloadLen := uint32(len(f.Payload)) payloadLen := uint32(len(f.Payload))
if int(payloadLen) != len(f.Payload) { if int(payloadLen) != len(f.Payload) {
return errors.New("payload too large") 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) total := frameHeader + metaSize + opcodeSize + rpcIDSize + int(payloadLen)
buf := make([]byte, total) buf := make([]byte, total)
// Frame length (payload only): matches Python's cmd_length_bytes
binary.BigEndian.PutUint32(buf[0:frameHeader], payloadLen) binary.BigEndian.PutUint32(buf[0:frameHeader], payloadLen)
// Metadata // Metadata
@ -189,7 +144,6 @@ func writeFrameWithMetaSize(w io.Writer, f *Frame, metaSize int) error {
// Opcode // Opcode
buf[frameHeader+metaSize] = f.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) binary.BigEndian.PutUint64(buf[frameHeader+metaSize+opcodeSize:frameHeader+metaSize+opcodeSize+rpcIDSize], f.RPCID)
// Payload // Payload
@ -217,33 +171,3 @@ func UnmarshalMetadata(b []byte) (Metadata, error) {
ClusterPeerID: binary.BigEndian.Uint64(b[4:12]), ClusterPeerID: binary.BigEndian.Uint64(b[4:12]),
}, nil }, 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()
}

View file

@ -1,218 +1,115 @@
// Copyright 2026-2027, QuarkChain.
package slave package slave
import ( import (
"bytes" "bytes"
"encoding/binary" "encoding/binary"
"encoding/hex" "encoding/hex"
"strings"
"testing" "testing"
) )
// ============================================================================= var (
// Python compatibility reference vectors // Test vectors serialized by pyquarkchain serializers using synthetic test values.
// //
// Each vector below is the exact wire bytes a Python peer would send/receive // Payload content (Ping/Pong commands):
// for a given frame. Source: qkc/quarkchain/protocol.py // - id = "id" (ASCII, 2 bytes)
// - Connection.write_raw_data (lines 302-308) // - full_shard_id_list = [1, 2]
// - Connection.read_metadata_and_raw_data (lines 285-300) // - root_tip = None (Ping only)
// // - opcode = 0x81 (PING), 0x82 (PONG) from ClusterOp (CLUSTER_OP_BASE=128)
// Wire layout: [4B payload_len] [metaSize B metadata] [1B opcode] [8B rpc_id] [payload] pythonClusterPing = "0000001300000001000000000000303981000000000000000100000002696400000002000000010000000200"
// pythonClusterPong = "00000012000000010000000000003039820000000000000001000000026964000000020000000100000002"
// metadata_class: pythonNoMetaPing = "0000001381000000000000000100000002696400000002000000010000000200"
// ClusterMetadata (12B) for master↔slave traffic pythonNoMetaPong = "00000012820000000000000001000000026964000000020000000100000002"
// Metadata (0B) for slave↔slave traffic )
// =============================================================================
// pingMasterWire: meta=(branch=0, peer=0), opcode=0x81 (PING), rpc_id=1, payload=empty func TestPythonVectors_ClusterPing(t *testing.T) {
// wire := mustPythonVectorBytes(t, pythonClusterPing)
// 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)) f, err := ReadFrame(bytes.NewReader(wire))
if err != nil { if err != nil {
t.Fatalf("ReadFrame: %v", err) t.Fatalf("ReadFrame: %v", err)
} }
if f.Opcode != 0x81 { if f.Opcode != 0x81 || f.RPCID != 1 {
t.Errorf("Opcode: got 0x%02x, want 0x81", f.Opcode) t.Fatalf("unexpected header: opcode=0x%02x rpcid=%d", f.Opcode, f.RPCID)
} }
if f.RPCID != 1 { if f.Meta.Branch != 1 || f.Meta.ClusterPeerID != 12345 {
t.Errorf("RPCID: got %d, want 1", f.RPCID) t.Fatalf("unexpected meta: %+v", f.Meta)
} }
if f.Meta != (Metadata{Branch: 0, ClusterPeerID: 0}) { var out bytes.Buffer
t.Errorf("Meta: got %+v", f.Meta) if err := WriteFrame(&out, f); err != nil {
t.Fatalf("WriteFrame: %v", err)
} }
if len(f.Payload) != 0 { if !bytes.Equal(out.Bytes(), wire) {
t.Errorf("Payload len: got %d, want 0", len(f.Payload)) t.Fatalf("frame mismatch:\n go %x\n py %x", out.Bytes(), wire)
} }
} }
func TestPythonRead_Pong(t *testing.T) { func TestPythonVectors_ClusterPong(t *testing.T) {
wire, _ := hex.DecodeString(pongMasterWire) wire := mustPythonVectorBytes(t, pythonClusterPong)
f, err := ReadFrame(bytes.NewReader(wire)) f, err := ReadFrame(bytes.NewReader(wire))
if err != nil { if err != nil {
t.Fatalf("ReadFrame: %v", err) t.Fatalf("ReadFrame: %v", err)
} }
if f.Opcode != 0x82 { if f.Opcode != 0x82 || f.RPCID != 1 {
t.Errorf("Opcode: got 0x%02x, want 0x82", f.Opcode) t.Fatalf("unexpected header: opcode=0x%02x rpcid=%d", f.Opcode, f.RPCID)
} }
if f.RPCID != 1 { if f.Meta.Branch != 1 || f.Meta.ClusterPeerID != 12345 {
t.Errorf("RPCID: got %d, want 1", f.RPCID) 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) { func TestPythonVectors_NoMetaPing(t *testing.T) {
wire, _ := hex.DecodeString(peerNewBlockWire) wire := mustPythonVectorBytes(t, pythonNoMetaPing)
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)) f, err := ReadFrameNoMeta(bytes.NewReader(wire))
if err != nil { if err != nil {
t.Fatalf("ReadFrameNoMeta: %v", err) t.Fatalf("ReadFrameNoMeta: %v", err)
} }
if f.Opcode != 0x93 { if f.Opcode != 0x81 || f.RPCID != 1 {
t.Errorf("Opcode: got 0x%02x, want 0x93", f.Opcode) t.Fatalf("unexpected header: opcode=0x%02x rpcid=%d", f.Opcode, f.RPCID)
} }
if f.RPCID != 42 { var out bytes.Buffer
t.Errorf("RPCID: got %d, want 42", f.RPCID) if err := WriteFrameNoMeta(&out, f); err != nil {
t.Fatalf("WriteFrameNoMeta: %v", err)
} }
if len(f.Payload) != 56 { if !bytes.Equal(out.Bytes(), wire) {
t.Errorf("Payload len: got %d, want 56", len(f.Payload)) t.Fatalf("frame mismatch:\n go %x\n py %x", out.Bytes(), wire)
} }
} }
func TestPythonRead_LargePayload(t *testing.T) { func TestPythonVectors_NoMetaPong(t *testing.T) {
wire, _ := hex.DecodeString(largePayloadWire) wire := mustPythonVectorBytes(t, pythonNoMetaPong)
f, err := ReadFrame(bytes.NewReader(wire)) f, err := ReadFrameNoMeta(bytes.NewReader(wire))
if err != nil { if err != nil {
t.Fatalf("ReadFrame: %v", err) t.Fatalf("ReadFrameNoMeta: %v", err)
} }
if f.Opcode != 0x10 { if f.Opcode != 0x82 || f.RPCID != 1 {
t.Errorf("Opcode: got 0x%02x, want 0x10", f.Opcode) t.Fatalf("unexpected header: opcode=0x%02x rpcid=%d", f.Opcode, f.RPCID)
} }
if f.RPCID != 7 { var out bytes.Buffer
t.Errorf("RPCID: got %d, want 7", f.RPCID) if err := WriteFrameNoMeta(&out, f); err != nil {
t.Fatalf("WriteFrameNoMeta: %v", err)
} }
if f.Meta.Branch != 3 { if !bytes.Equal(out.Bytes(), wire) {
t.Errorf("Branch: got %d, want 3", f.Meta.Branch) t.Fatalf("frame mismatch:\n go %x\n py %x", out.Bytes(), wire)
}
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
}
} }
} }
// ============================================================================= func mustPythonVectorBytes(t *testing.T, hexStr string) []byte {
// WriteFrame tests — Go's wire output must match Python's byte-for-byte. t.Helper()
// This is the strongest compatibility test: any wire-format drift is caught. b, err := hex.DecodeString(hexStr)
// ============================================================================= if err != nil {
t.Fatalf("decode python hex: %v", err)
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)
} }
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) { func TestRoundTrip_Meta(t *testing.T) {
cases := []struct { cases := []struct {
name string 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) { func TestWireFormatLayout(t *testing.T) {
f := &Frame{ f := &Frame{
Meta: Metadata{Branch: 1, ClusterPeerID: 0x1122334455667788}, 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) wire := writeFrameForTest(f.Meta, f.Opcode, f.RPCID, f.Payload)
// payload_len = 3
if got := binary.BigEndian.Uint32(wire[0:4]); got != 3 { if got := binary.BigEndian.Uint32(wire[0:4]); got != 3 {
t.Errorf("payload_len: got %d, want 3", got) t.Errorf("payload_len: got %d, want 3", got)
} }
// branch = 1
if got := binary.BigEndian.Uint32(wire[4:8]); got != 1 { if got := binary.BigEndian.Uint32(wire[4:8]); got != 1 {
t.Errorf("branch: got %d, want 1", got) t.Errorf("branch: got %d, want 1", got)
} }
// cluster_peer_id = 0x1122334455667788
if got := binary.BigEndian.Uint64(wire[8:16]); got != 0x1122334455667788 { if got := binary.BigEndian.Uint64(wire[8:16]); got != 0x1122334455667788 {
t.Errorf("cluster_peer_id: got 0x%x", got) t.Errorf("cluster_peer_id: got 0x%x", got)
} }
// opcode = 0x42
if wire[16] != 0x42 { if wire[16] != 0x42 {
t.Errorf("opcode: got 0x%02x, want 0x42", wire[16]) t.Errorf("opcode: got 0x%02x, want 0x42", wire[16])
} }
// rpc_id = 0xDEADBEEFCAFEBABE
if got := binary.BigEndian.Uint64(wire[17:25]); got != 0xDEADBEEFCAFEBABE { if got := binary.BigEndian.Uint64(wire[17:25]); got != 0xDEADBEEFCAFEBABE {
t.Errorf("rpc_id: got 0x%x", got) t.Errorf("rpc_id: got 0x%x", got)
} }
// payload
if !bytes.Equal(wire[25:28], []byte{0xAA, 0xBB, 0xCC}) { if !bytes.Equal(wire[25:28], []byte{0xAA, 0xBB, 0xCC}) {
t.Errorf("payload: got %x", wire[25:28]) t.Errorf("payload: got %x", wire[25:28])
} }
} }
func TestMultiFrameStream(t *testing.T) { 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{ frames := []*Frame{
{Meta: Metadata{Branch: 0, ClusterPeerID: 0}, Opcode: 0x81, RPCID: 0, Payload: []byte("ping")}, {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: 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) { func TestReadFrame_EOF(t *testing.T) {
if _, err := ReadFrame(bytes.NewReader(nil)); err == nil { if _, err := ReadFrame(bytes.NewReader(nil)); err == nil {
t.Error("expected error on empty stream") t.Error("expected error on empty stream")
@ -334,7 +214,6 @@ func TestReadFrame_EOF(t *testing.T) {
} }
func TestReadFrame_Truncated(t *testing.T) { func TestReadFrame_Truncated(t *testing.T) {
// payload_len says 100 bytes, but we only give 4 bytes
hdr := make([]byte, 4) hdr := make([]byte, 4)
binary.BigEndian.PutUint32(hdr, 100) binary.BigEndian.PutUint32(hdr, 100)
if _, err := ReadFrame(bytes.NewReader(hdr)); err == nil { if _, err := ReadFrame(bytes.NewReader(hdr)); err == nil {
@ -342,11 +221,36 @@ func TestReadFrame_Truncated(t *testing.T) {
} }
} }
// ============================================================================= func TestReadFrame_PayloadLimit(t *testing.T) {
// Helpers 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 { func writeFrameForTest(meta Metadata, opcode byte, rpcID uint64, payload []byte) []byte {
var buf bytes.Buffer var buf bytes.Buffer
_ = WriteFrame(&buf, &Frame{Meta: meta, Opcode: opcode, RPCID: rpcID, Payload: payload}) _ = WriteFrame(&buf, &Frame{Meta: meta, Opcode: opcode, RPCID: rpcID, Payload: payload})