Test / test (push) Successful in 2m4s
Release / gates (push) Successful in 2m5s
Release / build (amd64, freebsd) (push) Successful in 1m27s
Release / build (amd64, linux) (push) Successful in 1m22s
Release / build (amd64, netbsd) (push) Successful in 1m19s
Release / build (amd64, openbsd) (push) Successful in 1m20s
Release / build (arm64, darwin) (push) Successful in 1m21s
Release / build (arm64, freebsd) (push) Successful in 1m26s
Release / build (arm64, linux) (push) Successful in 1m25s
Release / build (arm64, netbsd) (push) Successful in 1m31s
Release / build (arm64, openbsd) (push) Successful in 1m27s
Release / build (loong64, linux) (push) Successful in 1m37s
Release / build (riscv64, linux) (push) Successful in 1m21s
Release / release (push) Successful in 40s
Assisted-by: GLM 5.3 Flash
262 lines
7.3 KiB
Go
262 lines
7.3 KiB
Go
// Copyright (c) 2026 Petr Balvín <opensource@petrbalvin.org> (https://petrbalvin.org)
|
|
// SPDX-License-Identifier: MIT
|
|
|
|
// The RPC-over-RDMA version one framing of RFC 8166: the fixed header,
|
|
// the chunk lists and the message assembly. The transfer of the chunks
|
|
// themselves belongs to RDMA verbs, which no pure Go stack can reach
|
|
// without cgo; this layer speaks the framing over a stream transport
|
|
// (TCP in the tests), where every chunk list stands empty and the whole
|
|
// message rides inline, wire compatible with an RDMA peer that
|
|
// registers nothing.
|
|
|
|
package rdma
|
|
|
|
import (
|
|
"encoding/binary"
|
|
"errors"
|
|
)
|
|
|
|
// The protocol constants of RFC 8166 section 4.
|
|
const (
|
|
Version = 1
|
|
|
|
ProcMsg = 0
|
|
ProcNomsg = 1
|
|
ProcError = 4
|
|
|
|
ErrVers = 1
|
|
ErrChunk = 2
|
|
MaxSupported = 1
|
|
)
|
|
|
|
// The transport magic of the RDMA echo of the record marking layer: a
|
|
// frame on the stream transport carries one RDMA message.
|
|
const frameMagic = 0x52444d31 // "RDM1"
|
|
|
|
// ErrFrame marks a malformed or unsupported frame.
|
|
var ErrFrame = errors.New("rdma: bad frame")
|
|
|
|
// A Segment is one xdr_rdma_segment: the registered memory handle, the
|
|
// chunk length and the remote offset.
|
|
type Segment struct {
|
|
Handle uint32
|
|
Length uint32
|
|
Offset uint64
|
|
}
|
|
|
|
// A ReadChunk is one xdr_read_chunk: the position in the XDR stream
|
|
// and the segment that carries the bytes.
|
|
type ReadChunk struct {
|
|
Position uint32
|
|
Segment Segment
|
|
}
|
|
|
|
// A WriteChunk is one xdr_write_chunk: a segment list for one reply
|
|
// piece.
|
|
type WriteChunk struct {
|
|
Segments []Segment
|
|
}
|
|
|
|
// A Header is the decoded RPC-over-RDMA frame: the fixed four fields,
|
|
// the chunk lists and the inline payload after them.
|
|
type Header struct {
|
|
XID uint32
|
|
Version uint32
|
|
Credit uint32
|
|
Proc uint32
|
|
Reads []ReadChunk
|
|
Writes []WriteChunk
|
|
Reply *WriteChunk
|
|
// Payload carries the inline RPC call or reply bytes; empty for
|
|
// RDMA_NOMSG frames whose payload rides the chunks.
|
|
Payload []byte
|
|
}
|
|
|
|
// appendSegment encodes one xdr_rdma_segment.
|
|
func appendSegment(b []byte, s Segment) []byte {
|
|
b = binary.BigEndian.AppendUint32(b, s.Handle)
|
|
b = binary.BigEndian.AppendUint32(b, s.Length)
|
|
return binary.BigEndian.AppendUint64(b, s.Offset)
|
|
}
|
|
|
|
// appendReadList encodes the optional read list: entries until the
|
|
// zero handle terminator.
|
|
func appendReadList(b []byte, reads []ReadChunk) []byte {
|
|
for _, r := range reads {
|
|
b = binary.BigEndian.AppendUint32(b, r.Position)
|
|
b = appendSegment(b, r.Segment)
|
|
}
|
|
return binary.BigEndian.AppendUint32(b, 0) // terminator
|
|
}
|
|
|
|
// appendWriteList encodes the optional write list: each chunk leads
|
|
// with its segment count and the list closes with a zero count word,
|
|
// RFC 8166 section 4.3.2.
|
|
func appendWriteList(b []byte, writes []WriteChunk) []byte {
|
|
for _, w := range writes {
|
|
b = binary.BigEndian.AppendUint32(b, uint32(len(w.Segments)))
|
|
for _, s := range w.Segments {
|
|
b = appendSegment(b, s)
|
|
}
|
|
}
|
|
return binary.BigEndian.AppendUint32(b, 0)
|
|
}
|
|
|
|
// AppendFrame encodes one RPC-over-RDMA message: the fixed header, the
|
|
// empty chunk lists of an inline transfer and the payload. A nil
|
|
// payload makes an RDMA_NOMSG frame.
|
|
func AppendFrame(b []byte, h Header) []byte {
|
|
b = binary.BigEndian.AppendUint32(b, h.XID)
|
|
b = binary.BigEndian.AppendUint32(b, Version)
|
|
b = binary.BigEndian.AppendUint32(b, h.Credit)
|
|
if h.Payload == nil {
|
|
b = binary.BigEndian.AppendUint32(b, ProcNomsg)
|
|
} else {
|
|
b = binary.BigEndian.AppendUint32(b, ProcMsg)
|
|
}
|
|
b = appendReadList(b, h.Reads)
|
|
b = appendWriteList(b, h.Writes)
|
|
// The reply chunk count is always present: zero when no reply chunk
|
|
// rides the frame.
|
|
if h.Reply != nil {
|
|
b = binary.BigEndian.AppendUint32(b, uint32(len(h.Reply.Segments)))
|
|
for _, s := range h.Reply.Segments {
|
|
b = appendSegment(b, s)
|
|
}
|
|
} else {
|
|
b = binary.BigEndian.AppendUint32(b, 0)
|
|
}
|
|
if h.Payload != nil {
|
|
b = binary.BigEndian.AppendUint32(b, uint32(len(h.Payload)))
|
|
b = append(b, h.Payload...)
|
|
}
|
|
return b
|
|
}
|
|
|
|
// DecodeFrame decodes one RPC-over-RDMA frame from the payload of a
|
|
// stream frame. The chunk lists are parsed and skipped: an inline
|
|
// transport never carries registered memory.
|
|
func DecodeFrame(frame []byte) (Header, error) {
|
|
h := Header{}
|
|
if len(frame) < 16 {
|
|
return h, ErrFrame
|
|
}
|
|
h.XID = binary.BigEndian.Uint32(frame[0:])
|
|
if binary.BigEndian.Uint32(frame[4:]) != Version {
|
|
return h, ErrFrame
|
|
}
|
|
h.Credit = binary.BigEndian.Uint32(frame[8:])
|
|
h.Proc = binary.BigEndian.Uint32(frame[12:])
|
|
off := 16
|
|
|
|
// The read list: entries until a zero position.
|
|
for {
|
|
if off+4 > len(frame) {
|
|
return h, ErrFrame
|
|
}
|
|
pos := binary.BigEndian.Uint32(frame[off:])
|
|
off += 4
|
|
if pos == 0 {
|
|
break
|
|
}
|
|
if off+16 > len(frame) {
|
|
return h, ErrFrame
|
|
}
|
|
var rc ReadChunk
|
|
rc.Position = pos
|
|
rc.Segment.Handle = binary.BigEndian.Uint32(frame[off:])
|
|
rc.Segment.Length = binary.BigEndian.Uint32(frame[off+4:])
|
|
rc.Segment.Offset = binary.BigEndian.Uint64(frame[off+8:])
|
|
off += 16
|
|
h.Reads = append(h.Reads, rc)
|
|
}
|
|
|
|
// The write list: chunks each led by their segment count, until a
|
|
// zero count word closes the list. The reply chunk follows as one
|
|
// more count-plus-segments group, RFC 8166 section 4.3.2.
|
|
for {
|
|
if off+4 > len(frame) {
|
|
return h, ErrFrame
|
|
}
|
|
n := binary.BigEndian.Uint32(frame[off:])
|
|
off += 4
|
|
if n != 0 {
|
|
var wc WriteChunk
|
|
for i := uint32(0); i < n; i++ {
|
|
if off+16 > len(frame) {
|
|
return h, ErrFrame
|
|
}
|
|
var s Segment
|
|
s.Handle = binary.BigEndian.Uint32(frame[off:])
|
|
s.Length = binary.BigEndian.Uint32(frame[off+4:])
|
|
s.Offset = binary.BigEndian.Uint64(frame[off+8:])
|
|
off += 16
|
|
wc.Segments = append(wc.Segments, s)
|
|
}
|
|
h.Writes = append(h.Writes, wc)
|
|
continue
|
|
}
|
|
// The zero closes the write list; the word after it is the reply
|
|
// chunk count.
|
|
if off+4 > len(frame) {
|
|
return h, ErrFrame
|
|
}
|
|
n = binary.BigEndian.Uint32(frame[off:])
|
|
off += 4
|
|
if n != 0 {
|
|
h.Reply = &WriteChunk{}
|
|
for i := uint32(0); i < n; i++ {
|
|
if off+16 > len(frame) {
|
|
return h, ErrFrame
|
|
}
|
|
var s Segment
|
|
s.Handle = binary.BigEndian.Uint32(frame[off:])
|
|
s.Length = binary.BigEndian.Uint32(frame[off+4:])
|
|
s.Offset = binary.BigEndian.Uint64(frame[off+8:])
|
|
off += 16
|
|
h.Reply.Segments = append(h.Reply.Segments, s)
|
|
}
|
|
}
|
|
break
|
|
}
|
|
if h.Proc == ProcMsg {
|
|
if off+4 > len(frame) {
|
|
return h, ErrFrame
|
|
}
|
|
n := binary.BigEndian.Uint32(frame[off:])
|
|
off += 4
|
|
if off+int(n) > len(frame) {
|
|
return h, ErrFrame
|
|
}
|
|
h.Payload = frame[off : off+int(n)]
|
|
}
|
|
return h, nil
|
|
}
|
|
|
|
// AppendStreamFrame frames one RDMA message for a stream transport:
|
|
// the magic, the payload length and the message. A decode of a frame
|
|
// written this way answers DecodeFrame(frame).
|
|
func AppendStreamFrame(b []byte, h Header) []byte {
|
|
msg := AppendFrame(nil, h)
|
|
b = binary.BigEndian.AppendUint32(b, frameMagic)
|
|
b = binary.BigEndian.AppendUint32(b, uint32(len(msg)))
|
|
return append(b, msg...)
|
|
}
|
|
|
|
// ReadStreamFrame splits one framed message off the front of buf: it
|
|
// answers the frame payload, the bytes consumed, and whether a whole
|
|
// frame is present.
|
|
func ReadStreamFrame(buf []byte) (frame []byte, consumed int, ok bool, err error) {
|
|
if len(buf) < 8 {
|
|
return nil, 0, false, nil
|
|
}
|
|
if binary.BigEndian.Uint32(buf[0:]) != frameMagic {
|
|
return nil, 0, false, ErrFrame
|
|
}
|
|
n := binary.BigEndian.Uint32(buf[4:])
|
|
if int(n)+8 > len(buf) {
|
|
return nil, 0, false, nil
|
|
}
|
|
return buf[8 : 8+n], 8 + int(n), true, nil
|
|
}
|