diff --git a/erts/emulator/nifs/unix/unix_socket_syncio.c b/erts/emulator/nifs/unix/unix_socket_syncio.c index 3049d287f775..0da7833c4b3b 100644 --- a/erts/emulator/nifs/unix/unix_socket_syncio.c +++ b/erts/emulator/nifs/unix/unix_socket_syncio.c @@ -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 * * ======================================================================== * @@ -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; } @@ -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); @@ -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); @@ -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; @@ -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; @@ -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; @@ -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); @@ -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) { @@ -4289,8 +4356,6 @@ ERL_NIF_TERM essio_recvmmsg(ErlNifEnv* env, } } - if (heapPool) enif_free(heapPool); - return ret; } #else /* HAVE_RECVMMSG */ diff --git a/lib/kernel/test/socket_SUITE.erl b/lib/kernel/test/socket_SUITE.erl index 4fa05a6e583d..c2a5d3e14846 100644 --- a/lib/kernel/test/socket_SUITE.erl +++ b/lib/kernel/test/socket_SUITE.erl @@ -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, @@ -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() -> @@ -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 %%