Skip to content

Commit d91f82d

Browse files
committed
Add test case for IORING_CQE_F_SOCK_FULL
Signed-off-by: Jens Axboe <axboe@kernel.dk>
1 parent 9a0461d commit d91f82d

3 files changed

Lines changed: 304 additions & 0 deletions

File tree

src/include/liburing/io_uring.h

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -457,12 +457,16 @@ struct io_uring_cqe {
457457
* other provided buffer type, all completions with a
458458
* buffer passed back is automatically returned to the
459459
* application.
460+
* IORING_CQE_F_SOCK_FULL If set, the socket was full when this send or
461+
* sendmsg was attempted. Hence it had to wait for POLLOUT
462+
* before being able to complete
460463
*/
461464
#define IORING_CQE_F_BUFFER (1U << 0)
462465
#define IORING_CQE_F_MORE (1U << 1)
463466
#define IORING_CQE_F_SOCK_NONEMPTY (1U << 2)
464467
#define IORING_CQE_F_NOTIF (1U << 3)
465468
#define IORING_CQE_F_BUF_MORE (1U << 4)
469+
#define IORING_CQE_F_SOCK_FULL IORING_CQE_F_SOCK_NONEMPTY
466470

467471
#define IORING_CQE_BUFFER_SHIFT 16
468472

test/Makefile

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -200,6 +200,7 @@ test_srcs := \
200200
self.c \
201201
recvsend_bundle.c \
202202
recvsend_bundle-inc.c \
203+
send_full.c \
203204
send_recv.c \
204205
send_recvmsg.c \
205206
send-zerocopy.c \

test/send_full.c

Lines changed: 299 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,299 @@
1+
/* SPDX-License-Identifier: MIT */
2+
/*
3+
* Description: Test flagging of IORING_CQE_F_SOCK_FULL on a socket
4+
*/
5+
#include <errno.h>
6+
#include <stdio.h>
7+
#include <stdlib.h>
8+
#include <string.h>
9+
#include <unistd.h>
10+
#include <arpa/inet.h>
11+
#include <sys/types.h>
12+
#include <sys/socket.h>
13+
#include <pthread.h>
14+
15+
#include "liburing.h"
16+
#include "helpers.h"
17+
18+
static int use_port = 10202;
19+
#define HOST "127.0.0.1"
20+
21+
struct recv_data {
22+
pthread_barrier_t startup;
23+
pthread_barrier_t receives;
24+
pthread_barrier_t finish;
25+
int pollfirst;
26+
unsigned int ring_flags;
27+
int port;
28+
};
29+
30+
static int recv_prep(struct recv_data *rd, int *sock)
31+
{
32+
struct sockaddr_in saddr;
33+
int sockfd, ret, val;
34+
socklen_t len;
35+
36+
memset(&saddr, 0, sizeof(saddr));
37+
saddr.sin_family = AF_INET;
38+
saddr.sin_addr.s_addr = htonl(INADDR_ANY);
39+
saddr.sin_port = htons(rd->port);
40+
41+
sockfd = socket(AF_INET, SOCK_STREAM, 0);
42+
if (sockfd < 0) {
43+
perror("socket");
44+
return 1;
45+
}
46+
47+
val = 1;
48+
setsockopt(sockfd, SOL_SOCKET, SO_REUSEADDR, &val, sizeof(val));
49+
50+
ret = bind(sockfd, (struct sockaddr *)&saddr, sizeof(saddr));
51+
if (ret < 0) {
52+
perror("bind");
53+
goto err;
54+
}
55+
56+
ret = listen(sockfd, 1);
57+
if (ret < 0) {
58+
perror("listen");
59+
goto err;
60+
}
61+
62+
pthread_barrier_wait(&rd->startup);
63+
64+
ret = listen(sockfd, 1);
65+
len = sizeof(saddr);
66+
ret = accept(sockfd, (struct sockaddr *)&saddr, &len);
67+
if (ret < 0) {
68+
perror("accept");
69+
goto err;
70+
}
71+
72+
close(sockfd);
73+
*sock = ret;
74+
return 0;
75+
err:
76+
close(sockfd);
77+
return 1;
78+
}
79+
80+
static void *recv_fn(void *data)
81+
{
82+
struct recv_data *rd = data;
83+
int ret, sock;
84+
char *buf;
85+
86+
ret = recv_prep(rd, &sock);
87+
if (ret) {
88+
fprintf(stderr, "recv_prep failed: %d\n", ret);
89+
pthread_barrier_wait(&rd->receives);
90+
goto err;
91+
}
92+
93+
pthread_barrier_wait(&rd->receives);
94+
buf = malloc(32768);
95+
do {
96+
ret = recv(sock, buf, 32768, MSG_DONTWAIT);
97+
if (ret > 0)
98+
continue;
99+
if (ret <= 0)
100+
break;
101+
} while (1);
102+
103+
free(buf);
104+
close(sock);
105+
ret = 0;
106+
err:
107+
pthread_barrier_wait(&rd->finish);
108+
return (void *)(intptr_t)ret;
109+
}
110+
111+
static int do_send(struct recv_data *rd)
112+
{
113+
struct sockaddr_in saddr;
114+
struct io_uring ring;
115+
struct io_uring_cqe *cqe;
116+
struct io_uring_sqe *sqe;
117+
int sockfd, ret, len;
118+
socklen_t optlen;
119+
char *buf;
120+
121+
pthread_barrier_wait(&rd->startup);
122+
123+
buf = malloc(4096);
124+
memset(buf, 0x5a, 4096);
125+
126+
ret = io_uring_queue_init(1, &ring, rd->ring_flags);
127+
if (ret) {
128+
fprintf(stderr, "queue init failed: %d\n", ret);
129+
return 1;
130+
}
131+
132+
memset(&saddr, 0, sizeof(saddr));
133+
saddr.sin_family = AF_INET;
134+
saddr.sin_port = htons(rd->port);
135+
inet_pton(AF_INET, HOST, &saddr.sin_addr);
136+
137+
sockfd = socket(AF_INET, SOCK_STREAM, 0);
138+
if (sockfd < 0) {
139+
perror("socket");
140+
goto err2;
141+
}
142+
143+
ret = connect(sockfd, (struct sockaddr *)&saddr, sizeof(saddr));
144+
if (ret < 0) {
145+
perror("connect");
146+
goto err;
147+
}
148+
149+
/* Limit socket buffer send size */
150+
optlen = sizeof(len);
151+
len = 4096;
152+
ret = setsockopt(sockfd, SOL_SOCKET, SO_SNDBUF, &len, optlen);
153+
if (ret < 0) {
154+
perror("setsockopt");
155+
goto err;
156+
}
157+
158+
do {
159+
sqe = io_uring_get_sqe(&ring);
160+
io_uring_prep_send(sqe, sockfd, buf, 4096, MSG_DONTWAIT);
161+
sqe->user_data = 1;
162+
163+
ret = io_uring_submit(&ring);
164+
if (ret <= 0) {
165+
fprintf(stderr, "submit failed: %d\n", ret);
166+
goto err;
167+
}
168+
169+
ret = io_uring_wait_cqe(&ring, &cqe);
170+
if (ret < 0) {
171+
fprintf(stderr, "wait: %d\n", ret);
172+
goto err;
173+
}
174+
if (cqe->flags & IORING_CQE_F_SOCK_FULL) {
175+
fprintf(stderr, "SOCK_FULL seen prematurely!\n");
176+
goto err;
177+
}
178+
if (cqe->res < 0) {
179+
/* socket now full */
180+
if (cqe->res == -EAGAIN) {
181+
io_uring_cqe_seen(&ring, cqe);
182+
break;
183+
}
184+
fprintf(stderr, "cqe res: %d\n", cqe->res);
185+
goto err;
186+
}
187+
/* if a short write, repeat until -EAGAIN is seen */
188+
io_uring_cqe_seen(&ring, cqe);
189+
} while (1);
190+
191+
/*
192+
* submit send on known full socket. this will go through the poll
193+
* machinery, waiting for POLLOUT.
194+
*/
195+
sqe = io_uring_get_sqe(&ring);
196+
io_uring_prep_send(sqe, sockfd, buf, 4096, MSG_WAITALL);
197+
if (rd->pollfirst)
198+
sqe->ioprio = IORING_RECVSEND_POLL_FIRST;
199+
sqe->user_data = 1;
200+
201+
io_uring_submit(&ring);
202+
203+
/*
204+
* Kick receive side to start freeing up socket space
205+
*/
206+
pthread_barrier_wait(&rd->receives);
207+
208+
/*
209+
* Wait for previous send to complete. That should have
210+
* IORING_CQE_F_SOCK_FULL set, if the kernel supports this kind
211+
* of notification.
212+
*/
213+
ret = io_uring_wait_cqe(&ring, &cqe);
214+
if (ret < 0) {
215+
fprintf(stderr, "wait: %d\n", ret);
216+
goto err;
217+
}
218+
if (cqe->res != 4096) {
219+
fprintf(stdout, "Unexpected send result: %d\n", cqe->res);
220+
io_uring_cqe_seen(&ring, cqe);
221+
goto err;
222+
}
223+
if (cqe->flags & IORING_CQE_F_SOCK_FULL)
224+
fprintf(stdout, "SOCK_FULL seen, kernel supports it (f=%x, p=%d)\n", rd->ring_flags, rd->pollfirst);
225+
else
226+
fprintf(stdout, "SOCK_FULL not set (f=%x, p=%d)\n", rd->ring_flags, rd->pollfirst);
227+
io_uring_cqe_seen(&ring, cqe);
228+
229+
free(buf);
230+
close(sockfd);
231+
io_uring_queue_exit(&ring);
232+
return 0;
233+
err:
234+
close(sockfd);
235+
err2:
236+
io_uring_queue_exit(&ring);
237+
return 1;
238+
}
239+
240+
static int test(unsigned int ring_flags, int pollfirst)
241+
{
242+
pthread_t recv_thread;
243+
struct recv_data rd;
244+
int ret;
245+
void *retval;
246+
247+
pthread_barrier_init(&rd.startup, NULL, 2);
248+
pthread_barrier_init(&rd.receives, NULL, 2);
249+
pthread_barrier_init(&rd.finish, NULL, 2);
250+
rd.pollfirst = pollfirst;
251+
rd.ring_flags = ring_flags;
252+
rd.port = use_port++;
253+
254+
ret = pthread_create(&recv_thread, NULL, recv_fn, &rd);
255+
if (ret) {
256+
fprintf(stderr, "Thread create failed: %d\n", ret);
257+
return 1;
258+
}
259+
260+
do_send(&rd);
261+
pthread_barrier_wait(&rd.finish);
262+
pthread_join(recv_thread, &retval);
263+
return (intptr_t)retval;
264+
}
265+
266+
int main(int argc, char *argv[])
267+
{
268+
int ret;
269+
270+
if (argc > 1)
271+
return T_EXIT_SKIP;
272+
273+
ret = test(0, 0);
274+
if (ret) {
275+
fprintf(stderr, "test 0 0 failed\n");
276+
return ret;
277+
}
278+
279+
ret = test(0, 1);
280+
if (ret) {
281+
fprintf(stderr, "test 0 1 failed\n");
282+
return ret;
283+
}
284+
285+
ret = test(IORING_SETUP_DEFER_TASKRUN|IORING_SETUP_SINGLE_ISSUER, 0);
286+
if (ret) {
287+
fprintf(stderr, "test defer 0 failed\n");
288+
return ret;
289+
}
290+
291+
ret = test(IORING_SETUP_DEFER_TASKRUN|IORING_SETUP_SINGLE_ISSUER, 1);
292+
if (ret) {
293+
fprintf(stderr, "test defer 1 failed\n");
294+
return ret;
295+
}
296+
297+
298+
return T_EXIT_PASS;
299+
}

0 commit comments

Comments
 (0)