Skip to content

Commit 5269669

Browse files
committed
unix: add Recvmmsg for linux
Add a Recvmmsg wrapper that exposes the recvmmsg(2) system call, allowing callers to receive multiple messages from a socket in a single syscall. ps and oobs are parallel slices of per-message payload and control buffers, with vlen derived from len(ps). Each message uses a single buffer rather than a scatter-gather vector, which should be sufficient for the typical recvmmsg use case. The return values ns, oobns, recvflags, and from are parallel slices of length n (the number of messages received). This is a first step toward golang/go#45886. The wrapper is useful as-is and lets callers experiment with batch receive patterns before a poller-integrated net package API is designed.
1 parent 99666ae commit 5269669

42 files changed

Lines changed: 385 additions & 0 deletions

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

unix/linux/types.go

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -772,6 +772,8 @@ type PacketMreq C.struct_packet_mreq
772772

773773
type Msghdr C.struct_msghdr
774774

775+
type Mmsghdr C.struct_mmsghdr
776+
775777
type Cmsghdr C.struct_cmsghdr
776778

777779
type Inet4Pktinfo C.struct_in_pktinfo
@@ -826,6 +828,7 @@ const (
826828
SizeofIPv6Mreq = C.sizeof_struct_ipv6_mreq
827829
SizeofPacketMreq = C.sizeof_struct_packet_mreq
828830
SizeofMsghdr = C.sizeof_struct_msghdr
831+
SizeofMmsghdr = C.sizeof_struct_mmsghdr
829832
SizeofCmsghdr = C.sizeof_struct_cmsghdr
830833
SizeofInet4Pktinfo = C.sizeof_struct_in_pktinfo
831834
SizeofInet6Pktinfo = C.sizeof_struct_in6_pktinfo

unix/syscall_linux.go

Lines changed: 65 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1617,6 +1617,71 @@ func sendmsgN(fd int, iov []Iovec, oob []byte, ptr unsafe.Pointer, salen _Sockle
16171617
return n, nil
16181618
}
16191619

1620+
// Recvmmsg receives multiple messages from a socket using the recvmmsg system
1621+
// call. ps holds the payload buffers and oobs the optional out-of-band control
1622+
// buffers, one entry per message; vlen is derived from len(ps). Pass a nil or
1623+
// zero-length oobs to receive no control data.
1624+
//
1625+
// The results are:
1626+
// - n is the number of messages received
1627+
// - ns[i] is the number of non-control bytes read into ps[i]
1628+
// - oobns[i] is the number of control bytes read into oobs[i]; interpret with [ParseSocketControlMessage]
1629+
// - recvflags[i] is the per-message flags returned by recvmmsg
1630+
// - from[i] is the sender address for message i, or nil for connected sockets
1631+
func Recvmmsg(fd int, ps, oobs [][]byte, flags int) (n int, ns, oobns, recvflags []int, from []Sockaddr, err error) {
1632+
vlen := len(ps)
1633+
if vlen == 0 {
1634+
return 0, nil, nil, nil, nil, EINVAL
1635+
}
1636+
if len(oobs) > 0 && len(oobs) != vlen {
1637+
return 0, nil, nil, nil, nil, EINVAL
1638+
}
1639+
1640+
msghdrs := make([]Mmsghdr, vlen)
1641+
iovecs := make([]Iovec, vlen)
1642+
names := make([]byte, SizeofSockaddrAny*vlen)
1643+
1644+
for i := range vlen {
1645+
if len(ps[i]) > 0 {
1646+
iovecs[i].Base = &ps[i][0]
1647+
iovecs[i].SetLen(len(ps[i]))
1648+
msghdrs[i].Hdr.Iov = &iovecs[i]
1649+
msghdrs[i].Hdr.SetIovlen(1)
1650+
}
1651+
if i < len(oobs) && len(oobs[i]) > 0 {
1652+
msghdrs[i].Hdr.Control = &oobs[i][0]
1653+
msghdrs[i].Hdr.SetControllen(len(oobs[i]))
1654+
}
1655+
msghdrs[i].Hdr.Name = &names[i*SizeofSockaddrAny]
1656+
msghdrs[i].Hdr.Namelen = uint32(SizeofSockaddrAny)
1657+
}
1658+
1659+
n, err = recvmmsg(fd, &msghdrs[0], vlen, flags, nil)
1660+
if err != nil {
1661+
return 0, nil, nil, nil, nil, err
1662+
}
1663+
1664+
ns = make([]int, n)
1665+
oobns = make([]int, n)
1666+
recvflags = make([]int, n)
1667+
from = make([]Sockaddr, n)
1668+
1669+
for i := range n {
1670+
ns[i] = int(msghdrs[i].Len)
1671+
oobns[i] = int(msghdrs[i].Hdr.Controllen)
1672+
recvflags[i] = int(msghdrs[i].Hdr.Flags)
1673+
rsa := (*RawSockaddrAny)(unsafe.Pointer(&names[i*SizeofSockaddrAny]))
1674+
if rsa.Addr.Family != AF_UNSPEC {
1675+
from[i], err = anyToSockaddr(fd, rsa)
1676+
if err != nil {
1677+
return
1678+
}
1679+
}
1680+
}
1681+
1682+
return
1683+
}
1684+
16201685
// BindToDevice binds the socket associated with fd to device.
16211686
func BindToDevice(fd int, device string) (err error) {
16221687
return SetsockoptString(fd, SOL_SOCKET, SO_BINDTODEVICE, device)

unix/syscall_linux_386.go

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -253,6 +253,14 @@ func sendmsg(s int, msg *Msghdr, flags int) (n int, err error) {
253253
return
254254
}
255255

256+
func recvmmsg(s int, mmsg *Mmsghdr, vlen int, flags int, timeout *Timespec) (n int, err error) {
257+
n, e := socketcall(_RECVMMSG, uintptr(s), uintptr(unsafe.Pointer(mmsg)), uintptr(vlen), uintptr(flags), uintptr(unsafe.Pointer(timeout)), 0)
258+
if e != 0 {
259+
err = e
260+
}
261+
return
262+
}
263+
256264
func Listen(s int, n int) (err error) {
257265
_, e := socketcall(_LISTEN, uintptr(s), uintptr(n), 0, 0, 0, 0)
258266
if e != 0 {

unix/syscall_linux_amd64.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -72,6 +72,7 @@ func Stat(path string, stat *Stat_t) (err error) {
7272
//sys sendto(s int, buf []byte, flags int, to unsafe.Pointer, addrlen _Socklen) (err error)
7373
//sys recvmsg(s int, msg *Msghdr, flags int) (n int, err error)
7474
//sys sendmsg(s int, msg *Msghdr, flags int) (n int, err error)
75+
//sys recvmmsg(s int, mmsg *Mmsghdr, vlen int, flags int, timeout *Timespec) (n int, err error)
7576
//sys mmap(addr uintptr, length uintptr, prot int, flags int, fd int, offset int64) (xaddr uintptr, err error)
7677

7778
//sys futimesat(dirfd int, path string, times *[2]Timeval) (err error)

unix/syscall_linux_arm.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -41,6 +41,7 @@ func Seek(fd int, offset int64, whence int) (newoffset int64, err error) {
4141
//sysnb socketpair(domain int, typ int, flags int, fd *[2]int32) (err error)
4242
//sys recvmsg(s int, msg *Msghdr, flags int) (n int, err error)
4343
//sys sendmsg(s int, msg *Msghdr, flags int) (n int, err error)
44+
//sys recvmmsg(s int, mmsg *Mmsghdr, vlen int, flags int, timeout *Timespec) (n int, err error)
4445

4546
// 64-bit file system and 32-bit uid calls
4647
// (16-bit uid calls are not always supported in newer kernels)

unix/syscall_linux_arm64.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -76,6 +76,7 @@ func Ustat(dev int, ubuf *Ustat_t) (err error) {
7676
//sys sendto(s int, buf []byte, flags int, to unsafe.Pointer, addrlen _Socklen) (err error)
7777
//sys recvmsg(s int, msg *Msghdr, flags int) (n int, err error)
7878
//sys sendmsg(s int, msg *Msghdr, flags int) (n int, err error)
79+
//sys recvmmsg(s int, mmsg *Mmsghdr, vlen int, flags int, timeout *Timespec) (n int, err error)
7980
//sys mmap(addr uintptr, length uintptr, prot int, flags int, fd int, offset int64) (xaddr uintptr, err error)
8081

8182
//sysnb Gettimeofday(tv *Timeval) (err error)

unix/syscall_linux_loong64.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -108,6 +108,7 @@ func Ustat(dev int, ubuf *Ustat_t) (err error) {
108108
//sys sendto(s int, buf []byte, flags int, to unsafe.Pointer, addrlen _Socklen) (err error)
109109
//sys recvmsg(s int, msg *Msghdr, flags int) (n int, err error)
110110
//sys sendmsg(s int, msg *Msghdr, flags int) (n int, err error)
111+
//sys recvmmsg(s int, mmsg *Mmsghdr, vlen int, flags int, timeout *Timespec) (n int, err error)
111112
//sys mmap(addr uintptr, length uintptr, prot int, flags int, fd int, offset int64) (xaddr uintptr, err error)
112113

113114
//sysnb Gettimeofday(tv *Timeval) (err error)

unix/syscall_linux_mips64x.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -56,6 +56,7 @@ func Select(nfd int, r *FdSet, w *FdSet, e *FdSet, timeout *Timeval) (n int, err
5656
//sys sendto(s int, buf []byte, flags int, to unsafe.Pointer, addrlen _Socklen) (err error)
5757
//sys recvmsg(s int, msg *Msghdr, flags int) (n int, err error)
5858
//sys sendmsg(s int, msg *Msghdr, flags int) (n int, err error)
59+
//sys recvmmsg(s int, mmsg *Mmsghdr, vlen int, flags int, timeout *Timespec) (n int, err error)
5960
//sys mmap(addr uintptr, length uintptr, prot int, flags int, fd int, offset int64) (xaddr uintptr, err error)
6061

6162
//sys futimesat(dirfd int, path string, times *[2]Timeval) (err error)

unix/syscall_linux_mipsx.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -50,6 +50,7 @@ func Syscall9(trap, a1, a2, a3, a4, a5, a6, a7, a8, a9 uintptr) (r1, r2 uintptr,
5050
//sys sendto(s int, buf []byte, flags int, to unsafe.Pointer, addrlen _Socklen) (err error)
5151
//sys recvmsg(s int, msg *Msghdr, flags int) (n int, err error)
5252
//sys sendmsg(s int, msg *Msghdr, flags int) (n int, err error)
53+
//sys recvmmsg(s int, mmsg *Mmsghdr, vlen int, flags int, timeout *Timespec) (n int, err error)
5354

5455
//sys Ioperm(from int, num int, on int) (err error)
5556
//sys Iopl(level int) (err error)

unix/syscall_linux_ppc.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -53,6 +53,7 @@ import (
5353
//sys sendto(s int, buf []byte, flags int, to unsafe.Pointer, addrlen _Socklen) (err error)
5454
//sys recvmsg(s int, msg *Msghdr, flags int) (n int, err error)
5555
//sys sendmsg(s int, msg *Msghdr, flags int) (n int, err error)
56+
//sys recvmmsg(s int, mmsg *Mmsghdr, vlen int, flags int, timeout *Timespec) (n int, err error)
5657

5758
//sys futimesat(dirfd int, path string, times *[2]Timeval) (err error)
5859
//sysnb Gettimeofday(tv *Timeval) (err error)

0 commit comments

Comments
 (0)