Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
147 changes: 106 additions & 41 deletions erts/emulator/nifs/unix/unix_socket_syncio.c
Original file line number Diff line number Diff line change
Expand Up @@ -335,6 +335,45 @@ typedef struct {
} ESSIOControl;


#ifdef HAVE_RECVMMSG
/* Grow-to-fit scratch block for essio_recvmmsg, kept per scheduler thread. */
typedef struct {
char* base;
size_t capacity;
unsigned int laid_vlen; /* dimensions the block is currently set up for */
size_t laid_bufSz;
size_t laid_ctrlSz;
unsigned int used; /* leading slots the previous call mutated */
} ESSIOMMsgPool;

static ErlNifTSDKey esock_mmsg_pool_key;

/* Return this thread's scratch block, grown (grow-only) to hold 'need' bytes. */
static ESSIOMMsgPool* essio_recvmmsg_pool(const size_t need)
{
ESSIOMMsgPool* pool = enif_tsd_get(esock_mmsg_pool_key);
if (pool == NULL) {
ESOCK_ASSERT( (pool = MALLOC(sizeof(ESSIOMMsgPool))) != NULL );
pool->base = NULL;
pool->capacity = 0;
pool->laid_vlen = 0;
pool->laid_bufSz = 0;
pool->laid_ctrlSz = 0;
pool->used = 0;
enif_tsd_set(esock_mmsg_pool_key, pool);
}
if (pool->capacity < need) {
char* nbase = (pool->base == NULL) ?
(char*) MALLOC(need) : (char*) REALLOC(pool->base, need);
ESOCK_ASSERT( nbase != NULL );
pool->base = nbase;
pool->capacity = need;
}
return pool;
}
#endif /* HAVE_RECVMMSG */


/* ======================================================================== *
* Function Forwards *
* ======================================================================== *
Expand Down Expand Up @@ -1112,6 +1151,11 @@ int essio_init(unsigned int numThreads,

essio_sctp_init();

#ifdef HAVE_RECVMMSG
ESOCK_ASSERT( enif_tsd_key_create("esock_mmsg_pool",
&esock_mmsg_pool_key) == 0 );
#endif

return ESOCK_IO_OK;
}

Expand Down Expand Up @@ -4125,7 +4169,8 @@ ERL_NIF_TERM essio_recvmmsg(ErlNifEnv* env,
struct iovec* recvIovecs = NULL;
ErlNifBinary* bufs = NULL;
ErlNifBinary* ctrls = NULL;
char* heapPool = NULL;
ERL_NIF_TERM* elems = NULL;
ESSIOMMsgPool* pool;
SOCKLEN_T addrLen = sizeof(ESockAddress);
size_t bufSz = (bufLen != 0 ? bufLen : descP->rBufSz);
size_t ctrlSz = (ctrlLen != 0 ? ctrlLen : descP->rCtrlSz);
Expand Down Expand Up @@ -4158,32 +4203,61 @@ ERL_NIF_TERM essio_recvmmsg(ErlNifEnv* env,
size_t addrs_sz = vlen * sizeof(ESockAddress);
size_t mmsghdrs_sz = vlen * sizeof(struct mmsghdr);
size_t iovecs_sz = vlen * sizeof(struct iovec);
size_t elems_sz = vlen * sizeof(ERL_NIF_TERM);
size_t meta_sz = bufs_sz + ctrls_sz + addrs_sz + mmsghdrs_sz + iovecs_sz + elems_sz;
size_t bufdata_sz = vlen * bufSz;
size_t ctrldata_sz = vlen * ctrlSz;
size_t total_sz = bufs_sz + ctrls_sz + addrs_sz + mmsghdrs_sz + iovecs_sz + bufdata_sz + ctrldata_sz;
ESOCK_ASSERT((heapPool = (char*) MALLOC(total_sz)) != NULL );
sys_memzero(heapPool, bufs_sz + ctrls_sz + addrs_sz);
bufs = (ErlNifBinary*) (heapPool);
ctrls = (ErlNifBinary*) (heapPool + bufs_sz);
addrs = (ESockAddress*) (heapPool + bufs_sz + ctrls_sz);
recvMmsghdrs = (struct mmsghdr*) (heapPool + bufs_sz + ctrls_sz + addrs_sz);
recvIovecs = (struct iovec*) (heapPool + bufs_sz + ctrls_sz + addrs_sz + mmsghdrs_sz);
recvBufs = (unsigned char*) (heapPool + bufs_sz + ctrls_sz + addrs_sz + mmsghdrs_sz + iovecs_sz);
recvCtrl = (unsigned char*) (heapPool + bufs_sz + ctrls_sz + addrs_sz + mmsghdrs_sz + iovecs_sz + bufdata_sz);
}

/* Set up mmsghdr structures to point into raw memory blocks */
for (i = 0; i < vlen; i++) {
recvIovecs[i].iov_base = recvBufs + (i * bufSz);
recvIovecs[i].iov_len = bufSz;
recvMmsghdrs[i].msg_hdr.msg_name = &addrs[i];
recvMmsghdrs[i].msg_hdr.msg_namelen = addrLen;
recvMmsghdrs[i].msg_hdr.msg_iov = &recvIovecs[i];
recvMmsghdrs[i].msg_hdr.msg_iovlen = 1;
recvMmsghdrs[i].msg_hdr.msg_control = recvCtrl + (i * ctrlSz);
recvMmsghdrs[i].msg_hdr.msg_controllen = ctrlSz;
recvMmsghdrs[i].msg_hdr.msg_flags = 0;
recvMmsghdrs[i].msg_len = 0;
size_t total_sz = meta_sz + bufdata_sz + ctrldata_sz;
char* p;

pool = essio_recvmmsg_pool(total_sz);
p = pool->base;

bufs = (ErlNifBinary*) (p);
ctrls = (ErlNifBinary*) (p + bufs_sz);
addrs = (ESockAddress*) (p + bufs_sz + ctrls_sz);
recvMmsghdrs = (struct mmsghdr*) (p + bufs_sz + ctrls_sz + addrs_sz);
recvIovecs = (struct iovec*) (p + bufs_sz + ctrls_sz + addrs_sz + mmsghdrs_sz);
elems = (ERL_NIF_TERM*) (p + bufs_sz + ctrls_sz + addrs_sz + mmsghdrs_sz + iovecs_sz);
recvBufs = (unsigned char*) (p + meta_sz);
recvCtrl = (unsigned char*) (p + meta_sz + bufdata_sz);

/* Full setup on (re)layout; otherwise only restore the fields the
* kernel and post-processing mutate, for the previously-used slots.
* A grow only happens when total_sz (hence the dimensions) changed,
* so the layout check below already covers it. 's' keeps 'i' at 0
* until the allocation loop below. */
if (vlen != pool->laid_vlen ||
bufSz != pool->laid_bufSz ||
ctrlSz != pool->laid_ctrlSz) {
unsigned int s;
for (s = 0; s < vlen; s++) {
recvIovecs[s].iov_base = recvBufs + (s * bufSz);
recvIovecs[s].iov_len = bufSz;
recvMmsghdrs[s].msg_hdr.msg_name = &addrs[s];
recvMmsghdrs[s].msg_hdr.msg_namelen = addrLen;
recvMmsghdrs[s].msg_hdr.msg_iov = &recvIovecs[s];
recvMmsghdrs[s].msg_hdr.msg_iovlen = 1;
recvMmsghdrs[s].msg_hdr.msg_control = recvCtrl + (s * ctrlSz);
recvMmsghdrs[s].msg_hdr.msg_controllen = ctrlSz;
recvMmsghdrs[s].msg_hdr.msg_flags = 0;
recvMmsghdrs[s].msg_len = 0;
}
pool->laid_vlen = vlen;
pool->laid_bufSz = bufSz;
pool->laid_ctrlSz = ctrlSz;
pool->used = 0;
} else {
unsigned int s;
const unsigned int used = pool->used;
for (s = 0; s < used; s++) {
recvMmsghdrs[s].msg_hdr.msg_namelen = addrLen;
recvMmsghdrs[s].msg_hdr.msg_control = recvCtrl + (s * ctrlSz);
recvMmsghdrs[s].msg_hdr.msg_controllen = ctrlSz;
recvMmsghdrs[s].msg_hdr.msg_flags = 0;
recvMmsghdrs[s].msg_len = 0;
}
}
}

ESOCK_CNT_INC(env, descP, sockRef, esock_atom_read_tries, &descP->readTries, 1);
Expand All @@ -4194,6 +4268,8 @@ ERL_NIF_TERM essio_recvmmsg(ErlNifEnv* env,
readResult = sock_recvmmsg(descP->sock, recvMmsghdrs, vlen, flags, NULL);
saveErrno = ESOCK_IS_ERROR(readResult) ? sock_errno() : 0;

pool->used = (readResult > 0) ? (unsigned int) readResult : 0;

if (readResult == 0) {
ret = esock_make_ok2(env, MKEL(env));
goto cleanup;
Expand All @@ -4208,8 +4284,6 @@ ERL_NIF_TERM essio_recvmmsg(ErlNifEnv* env,
*/
{
size_t totalBytes = 0;
ERL_NIF_TERM* elems;
ESOCK_ASSERT( (elems = MALLOC(readResult * sizeof(ERL_NIF_TERM))) != NULL );

for (i = 0; i < (unsigned int) readResult; i++) {
ErlNifBinary bin;
Expand All @@ -4223,16 +4297,14 @@ ERL_NIF_TERM essio_recvmmsg(ErlNifEnv* env,
if (msgLen > bufSz)
msgLen = bufSz;

ESOCK_ASSERT( ALLOC_BIN(bufSz, &bufs[i]) );
ESOCK_ASSERT( ALLOC_BIN(msgLen, &bufs[i]) );
sys_memcpy(bufs[i].data, recvBufs + (i * bufSz), msgLen);
bufs[i].size = bufSz;

ESOCK_ASSERT( ALLOC_BIN(ctrlSz, &ctrls[i]) );
ctrlLen = (recvMmsghdrs[i].msg_hdr.msg_controllen < ctrlSz)
? recvMmsghdrs[i].msg_hdr.msg_controllen
: ctrlSz;
ESOCK_ASSERT( ALLOC_BIN(ctrlLen, &ctrls[i]) );
sys_memcpy(ctrls[i].data, recvCtrl + (i * ctrlSz), ctrlLen);
ctrls[i].size = ctrlSz;

recvMmsghdrs[i].msg_hdr.msg_control = ctrls[i].data;

Expand All @@ -4248,7 +4320,6 @@ ERL_NIF_TERM essio_recvmmsg(ErlNifEnv* env,
}

resultList = enif_make_list_from_array(env, elems, readResult);
enif_free(elems);

/* Update packet and byte counters */
ESOCK_CNT_INC(env, descP, sockRef, esock_atom_read_pkg, &descP->readPkgCnt, readResult);
Expand All @@ -4268,13 +4339,9 @@ ERL_NIF_TERM essio_recvmmsg(ErlNifEnv* env,
}

cleanup:
/* Free ErlNifBinary structures only for received messages.
* Note: recv_create_bin may have transferred ownership (set data = NULL),
* in which case FREE_BIN is a no-op. We only free binaries we still own.
* When exiting early from the allocation loop, i is at the index where
* allocation failed, so we free indices 0 to i-1. On successful completion,
* i equals readResult, so we free indices 0 to readResult-1.
*/
/* Free the binaries we still own; recv_create_bin may have handed some off
* (data == NULL -> FREE_BIN is a no-op). 'i' bounds countToFree to slots
* the allocation loop populated (0 on the empty/error paths). */
{
unsigned int countToFree = (i < (unsigned int) readResult) ? i : (unsigned int) readResult;
if (countToFree > 0) {
Expand All @@ -4289,8 +4356,6 @@ ERL_NIF_TERM essio_recvmmsg(ErlNifEnv* env,
}
}

if (heapPool) enif_free(heapPool);

return ret;
}
#else /* HAVE_RECVMMSG */
Expand Down
113 changes: 112 additions & 1 deletion lib/kernel/test/socket_SUITE.erl
Original file line number Diff line number Diff line change
Expand Up @@ -158,6 +158,8 @@
sendmmsg_invalid_msg_format/1,
recvmmsg_dirty_scheduler_udp4/1,
sendmmsg_dirty_scheduler_udp4/1,
recvmmsg_pool_reuse_udp4/1,
recvmmsg_ctrl_udp4/1,

%% Socket IOCTL simple
ioctl_simple1/1,
Expand Down Expand Up @@ -393,7 +395,9 @@ batch_cases() ->
sendmmsg_with_addresses_udp4,
sendmmsg_invalid_msg_format,
recvmmsg_dirty_scheduler_udp4,
sendmmsg_dirty_scheduler_udp4
sendmmsg_dirty_scheduler_udp4,
recvmmsg_pool_reuse_udp4,
recvmmsg_ctrl_udp4
].

ioctl_cases() ->
Expand Down Expand Up @@ -15811,6 +15815,113 @@ sendmmsg_dirty_scheduler_udp4(_Config) when is_list(_Config) ->
).


%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%
%% Many recvmmsg calls on one socket, varying VLen/BufSz and received
%% count, to exercise the recvmmsg scratch-pool reuse/grow/reset paths.
%%
recvmmsg_pool_reuse_udp4(_Config) when is_list(_Config) ->
?TT(?SECS(60)),
tc_try(
recvmmsg_pool_reuse_udp4,
fun() ->
has_support_ipv4(),
has_recvmmsg_support()
end,
fun() ->
{ok, S1} = socket:open(inet, dgram, udp),
{ok, S2} = socket:open(inet, dgram, udp),
{ok, Addr} = inet:getaddr("localhost", inet),
ok = socket:bind(S1, #{family => inet, addr => Addr, port => 0}),
{ok, #{port := LocalPort}} = socket:sockname(S1),
ok = socket:connect(S2,
#{family => inet, addr => Addr,
port => LocalPort}),
%% Varying VLen/BufSz (repeated) -> grow and pure-reuse paths.
Rounds = [{5, 64}, {50, 2048}, {3, 512}, {120, 256},
{10, 4096}, {1, 8}, {80, 1024}],
lists:foreach(
fun({VLen, BufSz}) ->
recvmmsg_pool_reuse_round(S1, S2, VLen, BufSz)
end,
Rounds ++ Rounds),
%% Fixed layout, varying received count -> incremental reset,
%% incl. large-VLen/few-received.
lists:foreach(
fun(Count) ->
recvmmsg_pool_reuse_count(S1, S2, 64, 512, Count)
end,
[64, 1, 30, 64, 5, 1, 40, 64]),
ok = socket:close(S1),
ok = socket:close(S2),
ok
end
).

recvmmsg_pool_reuse_round(S1, S2, VLen, BufSz) ->
recvmmsg_pool_reuse_count(S1, S2, VLen, BufSz, VLen).

%% Send Count datagrams, receive with a recvmmsg of capacity VLen (Count =< VLen).
recvmmsg_pool_reuse_count(S1, S2, VLen, BufSz, Count) ->
Expected = [list_to_binary(io_lib:format("r~p_~p_m~p", [VLen, Count, N]))
|| N <- lists:seq(1, Count)],
lists:foreach(fun(D) -> ok = socket:send(S2, D) end, Expected),
{ok, Received} = socket:recvmmsg(S1, VLen, BufSz, 0, [], infinity),
Count = length(Received),
ReceivedData = [Data || Msg <- Received, [Data] <- [maps:get(iov, Msg)]],
true = lists:sort(ReceivedData) =:= lists:sort(Expected),
ok.


%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%
%% recvmmsg with ancillary data: enable ip pktinfo on the receiver so each
%% datagram carries a control message, and verify recvmmsg decodes a
%% non-empty ctrl per message (exercises the ctrl output path with actual
%% cmsg bytes, not just the data-only empty-ctrl case).
%%
recvmmsg_ctrl_udp4(_Config) when is_list(_Config) ->
?TT(?SECS(10)),
tc_try(
recvmmsg_ctrl_udp4,
fun() ->
has_support_ipv4(),
has_recvmmsg_support()
end,
fun() ->
{ok, S1} = socket:open(inet, dgram, udp),
{ok, S2} = socket:open(inet, dgram, udp),
{ok, Addr} = inet:getaddr("localhost", inet),
ok = socket:bind(S1, #{family => inet, addr => Addr, port => 0}),
{ok, #{port := LocalPort}} = socket:sockname(S1),
ok = socket:connect(S2,
#{family => inet, addr => Addr,
port => LocalPort}),
case socket:setopt(S1, ip, pktinfo, true) of
{error, _} ->
_ = socket:close(S1),
_ = socket:close(S2),
skip("ip pktinfo not supported");
ok ->
N = 5,
lists:foreach(
fun(I) -> ok = socket:send(S2, integer_to_binary(I)) end,
lists:seq(1, N)),
{ok, Received} = socket:recvmmsg(S1, 10, 0, 0, [], infinity),
N = length(Received),
lists:foreach(
fun(#{ctrl := Ctrl}) ->
true = lists:any(
fun(#{level := ip, type := pktinfo}) -> true;
(_) -> false
end, Ctrl)
end, Received),
ok = socket:close(S1),
ok = socket:close(S2),
ok
end
end
).


%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%
%% Helper function to check if recvmmsg is supported
%%
Expand Down
Loading