msg_zerocopy.c 19 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827
  1. /* Evaluate MSG_ZEROCOPY
  2. *
  3. * Send traffic between two processes over one of the supported
  4. * protocols and modes:
  5. *
  6. * PF_INET/PF_INET6
  7. * - SOCK_STREAM
  8. * - SOCK_DGRAM
  9. * - SOCK_DGRAM with UDP_CORK
  10. * - SOCK_RAW
  11. * - SOCK_RAW with IP_HDRINCL
  12. *
  13. * PF_PACKET
  14. * - SOCK_DGRAM
  15. * - SOCK_RAW
  16. *
  17. * PF_RDS
  18. * - SOCK_SEQPACKET
  19. *
  20. * Start this program on two connected hosts, one in send mode and
  21. * the other with option '-r' to put it in receiver mode.
  22. *
  23. * If zerocopy mode ('-z') is enabled, the sender will verify that
  24. * the kernel queues completions on the error queue for all zerocopy
  25. * transfers.
  26. */
  27. #define _GNU_SOURCE
  28. #include <arpa/inet.h>
  29. #include <error.h>
  30. #include <errno.h>
  31. #include <limits.h>
  32. #include <linux/errqueue.h>
  33. #include <linux/if_packet.h>
  34. #include <linux/ipv6.h>
  35. #include <linux/socket.h>
  36. #include <linux/sockios.h>
  37. #include <net/ethernet.h>
  38. #include <net/if.h>
  39. #include <netinet/ip.h>
  40. #include <netinet/ip6.h>
  41. #include <netinet/tcp.h>
  42. #include <netinet/udp.h>
  43. #include <poll.h>
  44. #include <sched.h>
  45. #include <stdbool.h>
  46. #include <stdio.h>
  47. #include <stdint.h>
  48. #include <stdlib.h>
  49. #include <string.h>
  50. #include <sys/ioctl.h>
  51. #include <sys/socket.h>
  52. #include <sys/stat.h>
  53. #include <sys/time.h>
  54. #include <sys/types.h>
  55. #include <sys/wait.h>
  56. #include <unistd.h>
  57. #include <linux/rds.h>
  58. #ifndef SO_EE_ORIGIN_ZEROCOPY
  59. #define SO_EE_ORIGIN_ZEROCOPY 5
  60. #endif
  61. #ifndef SO_ZEROCOPY
  62. #define SO_ZEROCOPY 60
  63. #endif
  64. #ifndef SO_EE_CODE_ZEROCOPY_COPIED
  65. #define SO_EE_CODE_ZEROCOPY_COPIED 1
  66. #endif
  67. #ifndef MSG_ZEROCOPY
  68. #define MSG_ZEROCOPY 0x4000000
  69. #endif
  70. static int cfg_cork;
  71. static bool cfg_cork_mixed;
  72. static int cfg_cpu = -1; /* default: pin to last cpu */
  73. static int cfg_expect_zerocopy = -1;
  74. static int cfg_family = PF_UNSPEC;
  75. static int cfg_ifindex = 1;
  76. static int cfg_payload_len;
  77. static int cfg_port = 8000;
  78. static bool cfg_rx;
  79. static int cfg_runtime_ms = 4200;
  80. static int cfg_verbose;
  81. static int cfg_waittime_ms = 500;
  82. static int cfg_notification_limit = 32;
  83. static bool cfg_zerocopy;
  84. static socklen_t cfg_alen;
  85. static struct sockaddr_storage cfg_dst_addr;
  86. static struct sockaddr_storage cfg_src_addr;
  87. static int exitcode;
  88. static char payload[IP_MAXPACKET];
  89. static long packets, bytes, completions, expected_completions;
  90. static uint32_t next_completion;
  91. static uint32_t sends_since_notify;
  92. static unsigned long gettimeofday_ms(void)
  93. {
  94. struct timeval tv;
  95. gettimeofday(&tv, NULL);
  96. return (tv.tv_sec * 1000) + (tv.tv_usec / 1000);
  97. }
  98. static uint16_t get_ip_csum(const uint16_t *start, int num_words)
  99. {
  100. unsigned long sum = 0;
  101. int i;
  102. for (i = 0; i < num_words; i++)
  103. sum += start[i];
  104. while (sum >> 16)
  105. sum = (sum & 0xFFFF) + (sum >> 16);
  106. return ~sum;
  107. }
  108. static int do_setcpu(int cpu)
  109. {
  110. cpu_set_t mask;
  111. CPU_ZERO(&mask);
  112. CPU_SET(cpu, &mask);
  113. if (sched_setaffinity(0, sizeof(mask), &mask))
  114. fprintf(stderr, "cpu: unable to pin, may increase variance.\n");
  115. else if (cfg_verbose)
  116. fprintf(stderr, "cpu: %u\n", cpu);
  117. return 0;
  118. }
  119. static void do_setsockopt(int fd, int level, int optname, int val)
  120. {
  121. if (setsockopt(fd, level, optname, &val, sizeof(val)))
  122. error(1, errno, "setsockopt %d.%d: %d", level, optname, val);
  123. }
  124. static int do_poll(int fd, int events)
  125. {
  126. struct pollfd pfd;
  127. int ret;
  128. pfd.events = events;
  129. pfd.revents = 0;
  130. pfd.fd = fd;
  131. ret = poll(&pfd, 1, cfg_waittime_ms);
  132. if (ret == -1)
  133. error(1, errno, "poll");
  134. return ret && (pfd.revents & events);
  135. }
  136. static int do_accept(int fd)
  137. {
  138. int fda = fd;
  139. fd = accept(fda, NULL, NULL);
  140. if (fd == -1)
  141. error(1, errno, "accept");
  142. if (close(fda))
  143. error(1, errno, "close listen sock");
  144. return fd;
  145. }
  146. static void add_zcopy_cookie(struct msghdr *msg, uint32_t cookie)
  147. {
  148. struct cmsghdr *cm;
  149. if (!msg->msg_control)
  150. error(1, errno, "NULL cookie");
  151. cm = (void *)msg->msg_control;
  152. cm->cmsg_len = CMSG_LEN(sizeof(cookie));
  153. cm->cmsg_level = SOL_RDS;
  154. cm->cmsg_type = RDS_CMSG_ZCOPY_COOKIE;
  155. memcpy(CMSG_DATA(cm), &cookie, sizeof(cookie));
  156. }
  157. static bool do_sendmsg(int fd, struct msghdr *msg, bool do_zerocopy, int domain)
  158. {
  159. int ret, len, i, flags;
  160. static uint32_t cookie;
  161. char ckbuf[CMSG_SPACE(sizeof(cookie))];
  162. len = 0;
  163. for (i = 0; i < msg->msg_iovlen; i++)
  164. len += msg->msg_iov[i].iov_len;
  165. flags = MSG_DONTWAIT;
  166. if (do_zerocopy) {
  167. flags |= MSG_ZEROCOPY;
  168. if (domain == PF_RDS) {
  169. memset(&msg->msg_control, 0, sizeof(msg->msg_control));
  170. msg->msg_controllen = CMSG_SPACE(sizeof(cookie));
  171. msg->msg_control = (struct cmsghdr *)ckbuf;
  172. add_zcopy_cookie(msg, ++cookie);
  173. }
  174. }
  175. ret = sendmsg(fd, msg, flags);
  176. if (ret == -1 && errno == EAGAIN)
  177. return false;
  178. if (ret == -1)
  179. error(1, errno, "send");
  180. if (cfg_verbose && ret != len)
  181. fprintf(stderr, "send: ret=%u != %u\n", ret, len);
  182. sends_since_notify++;
  183. if (len) {
  184. packets++;
  185. bytes += ret;
  186. if (do_zerocopy && ret)
  187. expected_completions++;
  188. }
  189. if (do_zerocopy && domain == PF_RDS) {
  190. msg->msg_control = NULL;
  191. msg->msg_controllen = 0;
  192. }
  193. return true;
  194. }
  195. static void do_sendmsg_corked(int fd, struct msghdr *msg)
  196. {
  197. bool do_zerocopy = cfg_zerocopy;
  198. int i, payload_len, extra_len;
  199. /* split up the packet. for non-multiple, make first buffer longer */
  200. payload_len = cfg_payload_len / cfg_cork;
  201. extra_len = cfg_payload_len - (cfg_cork * payload_len);
  202. do_setsockopt(fd, IPPROTO_UDP, UDP_CORK, 1);
  203. for (i = 0; i < cfg_cork; i++) {
  204. /* in mixed-frags mode, alternate zerocopy and copy frags
  205. * start with non-zerocopy, to ensure attach later works
  206. */
  207. if (cfg_cork_mixed)
  208. do_zerocopy = (i & 1);
  209. msg->msg_iov[0].iov_len = payload_len + extra_len;
  210. extra_len = 0;
  211. do_sendmsg(fd, msg, do_zerocopy,
  212. (cfg_dst_addr.ss_family == AF_INET ?
  213. PF_INET : PF_INET6));
  214. }
  215. do_setsockopt(fd, IPPROTO_UDP, UDP_CORK, 0);
  216. }
  217. static int setup_iph(struct iphdr *iph, uint16_t payload_len)
  218. {
  219. struct sockaddr_in *daddr = (void *) &cfg_dst_addr;
  220. struct sockaddr_in *saddr = (void *) &cfg_src_addr;
  221. memset(iph, 0, sizeof(*iph));
  222. iph->version = 4;
  223. iph->tos = 0;
  224. iph->ihl = 5;
  225. iph->ttl = 2;
  226. iph->saddr = saddr->sin_addr.s_addr;
  227. iph->daddr = daddr->sin_addr.s_addr;
  228. iph->protocol = IPPROTO_EGP;
  229. iph->tot_len = htons(sizeof(*iph) + payload_len);
  230. iph->check = get_ip_csum((void *) iph, iph->ihl << 1);
  231. return sizeof(*iph);
  232. }
  233. static int setup_ip6h(struct ipv6hdr *ip6h, uint16_t payload_len)
  234. {
  235. struct sockaddr_in6 *daddr = (void *) &cfg_dst_addr;
  236. struct sockaddr_in6 *saddr = (void *) &cfg_src_addr;
  237. memset(ip6h, 0, sizeof(*ip6h));
  238. ip6h->version = 6;
  239. ip6h->payload_len = htons(payload_len);
  240. ip6h->nexthdr = IPPROTO_EGP;
  241. ip6h->hop_limit = 2;
  242. ip6h->saddr = saddr->sin6_addr;
  243. ip6h->daddr = daddr->sin6_addr;
  244. return sizeof(*ip6h);
  245. }
  246. static void setup_sockaddr(int domain, const char *str_addr,
  247. struct sockaddr_storage *sockaddr)
  248. {
  249. struct sockaddr_in6 *addr6 = (void *) sockaddr;
  250. struct sockaddr_in *addr4 = (void *) sockaddr;
  251. switch (domain) {
  252. case PF_INET:
  253. memset(addr4, 0, sizeof(*addr4));
  254. addr4->sin_family = AF_INET;
  255. addr4->sin_port = htons(cfg_port);
  256. if (str_addr &&
  257. inet_pton(AF_INET, str_addr, &(addr4->sin_addr)) != 1)
  258. error(1, 0, "ipv4 parse error: %s", str_addr);
  259. break;
  260. case PF_INET6:
  261. memset(addr6, 0, sizeof(*addr6));
  262. addr6->sin6_family = AF_INET6;
  263. addr6->sin6_port = htons(cfg_port);
  264. if (str_addr &&
  265. inet_pton(AF_INET6, str_addr, &(addr6->sin6_addr)) != 1)
  266. error(1, 0, "ipv6 parse error: %s", str_addr);
  267. break;
  268. default:
  269. error(1, 0, "illegal domain");
  270. }
  271. }
  272. static int do_setup_tx(int domain, int type, int protocol)
  273. {
  274. int fd;
  275. fd = socket(domain, type, protocol);
  276. if (fd == -1)
  277. error(1, errno, "socket t");
  278. do_setsockopt(fd, SOL_SOCKET, SO_SNDBUF, 1 << 21);
  279. if (cfg_zerocopy)
  280. do_setsockopt(fd, SOL_SOCKET, SO_ZEROCOPY, 1);
  281. if (domain != PF_PACKET && domain != PF_RDS)
  282. if (connect(fd, (void *) &cfg_dst_addr, cfg_alen))
  283. error(1, errno, "connect");
  284. if (domain == PF_RDS) {
  285. if (bind(fd, (void *) &cfg_src_addr, cfg_alen))
  286. error(1, errno, "bind");
  287. }
  288. return fd;
  289. }
  290. static uint32_t do_process_zerocopy_cookies(struct rds_zcopy_cookies *ck)
  291. {
  292. int i;
  293. if (ck->num > RDS_MAX_ZCOOKIES)
  294. error(1, 0, "Returned %d cookies, max expected %d\n",
  295. ck->num, RDS_MAX_ZCOOKIES);
  296. for (i = 0; i < ck->num; i++)
  297. if (cfg_verbose >= 2)
  298. fprintf(stderr, "%d\n", ck->cookies[i]);
  299. return ck->num;
  300. }
  301. static bool do_recvmsg_completion(int fd)
  302. {
  303. char cmsgbuf[CMSG_SPACE(sizeof(struct rds_zcopy_cookies))];
  304. struct rds_zcopy_cookies *ck;
  305. struct cmsghdr *cmsg;
  306. struct msghdr msg;
  307. bool ret = false;
  308. memset(&msg, 0, sizeof(msg));
  309. msg.msg_control = cmsgbuf;
  310. msg.msg_controllen = sizeof(cmsgbuf);
  311. if (recvmsg(fd, &msg, MSG_DONTWAIT))
  312. return ret;
  313. if (msg.msg_flags & MSG_CTRUNC)
  314. error(1, errno, "recvmsg notification: truncated");
  315. for (cmsg = CMSG_FIRSTHDR(&msg); cmsg; cmsg = CMSG_NXTHDR(&msg, cmsg)) {
  316. if (cmsg->cmsg_level == SOL_RDS &&
  317. cmsg->cmsg_type == RDS_CMSG_ZCOPY_COMPLETION) {
  318. ck = (struct rds_zcopy_cookies *)CMSG_DATA(cmsg);
  319. completions += do_process_zerocopy_cookies(ck);
  320. ret = true;
  321. break;
  322. }
  323. error(0, 0, "ignoring cmsg at level %d type %d\n",
  324. cmsg->cmsg_level, cmsg->cmsg_type);
  325. }
  326. return ret;
  327. }
  328. static bool do_recv_completion(int fd, int domain)
  329. {
  330. struct sock_extended_err *serr;
  331. struct msghdr msg = {};
  332. struct cmsghdr *cm;
  333. uint32_t hi, lo, range;
  334. int ret, zerocopy;
  335. char control[100];
  336. if (domain == PF_RDS)
  337. return do_recvmsg_completion(fd);
  338. msg.msg_control = control;
  339. msg.msg_controllen = sizeof(control);
  340. ret = recvmsg(fd, &msg, MSG_ERRQUEUE);
  341. if (ret == -1 && errno == EAGAIN)
  342. return false;
  343. if (ret == -1)
  344. error(1, errno, "recvmsg notification");
  345. if (msg.msg_flags & MSG_CTRUNC)
  346. error(1, errno, "recvmsg notification: truncated");
  347. cm = CMSG_FIRSTHDR(&msg);
  348. if (!cm)
  349. error(1, 0, "cmsg: no cmsg");
  350. if (!((cm->cmsg_level == SOL_IP && cm->cmsg_type == IP_RECVERR) ||
  351. (cm->cmsg_level == SOL_IPV6 && cm->cmsg_type == IPV6_RECVERR) ||
  352. (cm->cmsg_level == SOL_PACKET && cm->cmsg_type == PACKET_TX_TIMESTAMP)))
  353. error(1, 0, "serr: wrong type: %d.%d",
  354. cm->cmsg_level, cm->cmsg_type);
  355. serr = (void *) CMSG_DATA(cm);
  356. if (serr->ee_origin != SO_EE_ORIGIN_ZEROCOPY)
  357. error(1, 0, "serr: wrong origin: %u", serr->ee_origin);
  358. if (serr->ee_errno != 0)
  359. error(1, 0, "serr: wrong error code: %u", serr->ee_errno);
  360. hi = serr->ee_data;
  361. lo = serr->ee_info;
  362. range = hi - lo + 1;
  363. /* Detect notification gaps. These should not happen often, if at all.
  364. * Gaps can occur due to drops, reordering and retransmissions.
  365. */
  366. if (cfg_verbose && lo != next_completion)
  367. fprintf(stderr, "gap: %u..%u does not append to %u\n",
  368. lo, hi, next_completion);
  369. next_completion = hi + 1;
  370. zerocopy = !(serr->ee_code & SO_EE_CODE_ZEROCOPY_COPIED);
  371. if (cfg_expect_zerocopy != -1 &&
  372. cfg_expect_zerocopy != zerocopy) {
  373. fprintf(stderr, "serr: ee_code: %u != expected %u\n",
  374. zerocopy, cfg_expect_zerocopy);
  375. exitcode = 1;
  376. /* suppress repeated messages */
  377. cfg_expect_zerocopy = zerocopy;
  378. }
  379. if (cfg_verbose >= 2)
  380. fprintf(stderr, "completed: %u (h=%u l=%u)\n",
  381. range, hi, lo);
  382. completions += range;
  383. return true;
  384. }
  385. /* Read all outstanding messages on the errqueue */
  386. static void do_recv_completions(int fd, int domain)
  387. {
  388. while (do_recv_completion(fd, domain)) {}
  389. sends_since_notify = 0;
  390. }
  391. /* Wait for all remaining completions on the errqueue */
  392. static void do_recv_remaining_completions(int fd, int domain)
  393. {
  394. int64_t tstop = gettimeofday_ms() + cfg_waittime_ms;
  395. while (completions < expected_completions &&
  396. gettimeofday_ms() < tstop) {
  397. if (do_poll(fd, domain == PF_RDS ? POLLIN : POLLERR))
  398. do_recv_completions(fd, domain);
  399. }
  400. if (completions < expected_completions)
  401. fprintf(stderr, "missing notifications: %lu < %lu\n",
  402. completions, expected_completions);
  403. }
  404. static void do_tx(int domain, int type, int protocol)
  405. {
  406. struct iovec iov[3] = { {0} };
  407. struct sockaddr_ll laddr;
  408. struct msghdr msg = {0};
  409. struct ethhdr eth;
  410. union {
  411. struct ipv6hdr ip6h;
  412. struct iphdr iph;
  413. } nh;
  414. uint64_t tstop;
  415. int fd;
  416. fd = do_setup_tx(domain, type, protocol);
  417. if (domain == PF_PACKET) {
  418. uint16_t proto = cfg_family == PF_INET ? ETH_P_IP : ETH_P_IPV6;
  419. /* sock_raw passes ll header as data */
  420. if (type == SOCK_RAW) {
  421. memset(eth.h_dest, 0x06, ETH_ALEN);
  422. memset(eth.h_source, 0x02, ETH_ALEN);
  423. eth.h_proto = htons(proto);
  424. iov[0].iov_base = &eth;
  425. iov[0].iov_len = sizeof(eth);
  426. msg.msg_iovlen++;
  427. }
  428. /* both sock_raw and sock_dgram expect name */
  429. memset(&laddr, 0, sizeof(laddr));
  430. laddr.sll_family = AF_PACKET;
  431. laddr.sll_ifindex = cfg_ifindex;
  432. laddr.sll_protocol = htons(proto);
  433. laddr.sll_halen = ETH_ALEN;
  434. memset(laddr.sll_addr, 0x06, ETH_ALEN);
  435. msg.msg_name = &laddr;
  436. msg.msg_namelen = sizeof(laddr);
  437. }
  438. /* packet and raw sockets with hdrincl must pass network header */
  439. if (domain == PF_PACKET || protocol == IPPROTO_RAW) {
  440. if (cfg_family == PF_INET)
  441. iov[1].iov_len = setup_iph(&nh.iph, cfg_payload_len);
  442. else
  443. iov[1].iov_len = setup_ip6h(&nh.ip6h, cfg_payload_len);
  444. iov[1].iov_base = (void *) &nh;
  445. msg.msg_iovlen++;
  446. }
  447. if (domain == PF_RDS) {
  448. msg.msg_name = &cfg_dst_addr;
  449. msg.msg_namelen = (cfg_dst_addr.ss_family == AF_INET ?
  450. sizeof(struct sockaddr_in) :
  451. sizeof(struct sockaddr_in6));
  452. }
  453. iov[2].iov_base = payload;
  454. iov[2].iov_len = cfg_payload_len;
  455. msg.msg_iovlen++;
  456. msg.msg_iov = &iov[3 - msg.msg_iovlen];
  457. tstop = gettimeofday_ms() + cfg_runtime_ms;
  458. do {
  459. if (cfg_cork)
  460. do_sendmsg_corked(fd, &msg);
  461. else
  462. do_sendmsg(fd, &msg, cfg_zerocopy, domain);
  463. if (cfg_zerocopy && sends_since_notify >= cfg_notification_limit)
  464. do_recv_completions(fd, domain);
  465. while (!do_poll(fd, POLLOUT)) {
  466. if (cfg_zerocopy)
  467. do_recv_completions(fd, domain);
  468. }
  469. } while (gettimeofday_ms() < tstop);
  470. if (cfg_zerocopy)
  471. do_recv_remaining_completions(fd, domain);
  472. if (close(fd))
  473. error(1, errno, "close");
  474. fprintf(stderr, "tx=%lu (%lu MB) txc=%lu zc=%c\n",
  475. packets, bytes >> 20, completions,
  476. cfg_zerocopy && cfg_expect_zerocopy == 1 ? 'y' : 'n');
  477. }
  478. static int do_setup_rx(int domain, int type, int protocol)
  479. {
  480. int fd;
  481. /* If tx over PF_PACKET, rx over PF_INET(6)/SOCK_RAW,
  482. * to recv the only copy of the packet, not a clone
  483. */
  484. if (domain == PF_PACKET)
  485. error(1, 0, "Use PF_INET/SOCK_RAW to read");
  486. if (type == SOCK_RAW && protocol == IPPROTO_RAW)
  487. error(1, 0, "IPPROTO_RAW: not supported on Rx");
  488. fd = socket(domain, type, protocol);
  489. if (fd == -1)
  490. error(1, errno, "socket r");
  491. do_setsockopt(fd, SOL_SOCKET, SO_RCVBUF, 1 << 21);
  492. do_setsockopt(fd, SOL_SOCKET, SO_RCVLOWAT, 1 << 16);
  493. do_setsockopt(fd, SOL_SOCKET, SO_REUSEPORT, 1);
  494. if (bind(fd, (void *) &cfg_dst_addr, cfg_alen))
  495. error(1, errno, "bind");
  496. if (type == SOCK_STREAM) {
  497. if (listen(fd, 1))
  498. error(1, errno, "listen");
  499. fd = do_accept(fd);
  500. }
  501. return fd;
  502. }
  503. /* Flush all outstanding bytes for the tcp receive queue */
  504. static void do_flush_tcp(int fd)
  505. {
  506. int ret;
  507. /* MSG_TRUNC flushes up to len bytes */
  508. ret = recv(fd, NULL, 1 << 21, MSG_TRUNC | MSG_DONTWAIT);
  509. if (ret == -1 && errno == EAGAIN)
  510. return;
  511. if (ret == -1)
  512. error(1, errno, "flush");
  513. if (!ret)
  514. return;
  515. packets++;
  516. bytes += ret;
  517. }
  518. /* Flush all outstanding datagrams. Verify first few bytes of each. */
  519. static void do_flush_datagram(int fd, int type)
  520. {
  521. int ret, off = 0;
  522. char buf[64];
  523. /* MSG_TRUNC will return full datagram length */
  524. ret = recv(fd, buf, sizeof(buf), MSG_DONTWAIT | MSG_TRUNC);
  525. if (ret == -1 && errno == EAGAIN)
  526. return;
  527. /* raw ipv4 return with header, raw ipv6 without */
  528. if (cfg_family == PF_INET && type == SOCK_RAW) {
  529. off += sizeof(struct iphdr);
  530. ret -= sizeof(struct iphdr);
  531. }
  532. if (ret == -1)
  533. error(1, errno, "recv");
  534. if (ret != cfg_payload_len)
  535. error(1, 0, "recv: ret=%u != %u", ret, cfg_payload_len);
  536. if (ret > sizeof(buf) - off)
  537. ret = sizeof(buf) - off;
  538. if (memcmp(buf + off, payload, ret))
  539. error(1, 0, "recv: data mismatch");
  540. packets++;
  541. bytes += cfg_payload_len;
  542. }
  543. static void do_rx(int domain, int type, int protocol)
  544. {
  545. const int cfg_receiver_wait_ms = 400;
  546. uint64_t tstop;
  547. int fd;
  548. fd = do_setup_rx(domain, type, protocol);
  549. tstop = gettimeofday_ms() + cfg_runtime_ms + cfg_receiver_wait_ms;
  550. do {
  551. if (type == SOCK_STREAM)
  552. do_flush_tcp(fd);
  553. else
  554. do_flush_datagram(fd, type);
  555. do_poll(fd, POLLIN);
  556. } while (gettimeofday_ms() < tstop);
  557. if (close(fd))
  558. error(1, errno, "close");
  559. fprintf(stderr, "rx=%lu (%lu MB)\n", packets, bytes >> 20);
  560. }
  561. static void do_test(int domain, int type, int protocol)
  562. {
  563. int i;
  564. if (cfg_cork && (domain == PF_PACKET || type != SOCK_DGRAM))
  565. error(1, 0, "can only cork udp sockets");
  566. do_setcpu(cfg_cpu);
  567. for (i = 0; i < IP_MAXPACKET; i++)
  568. payload[i] = 'a' + (i % 26);
  569. if (cfg_rx)
  570. do_rx(domain, type, protocol);
  571. else
  572. do_tx(domain, type, protocol);
  573. }
  574. static void usage(const char *filepath)
  575. {
  576. error(1, 0, "Usage: %s [options] <test>", filepath);
  577. }
  578. static void parse_opts(int argc, char **argv)
  579. {
  580. const int max_payload_len = sizeof(payload) -
  581. sizeof(struct ipv6hdr) -
  582. sizeof(struct tcphdr) -
  583. 40 /* max tcp options */;
  584. int c;
  585. char *daddr = NULL, *saddr = NULL;
  586. char *cfg_test;
  587. cfg_payload_len = max_payload_len;
  588. while ((c = getopt(argc, argv, "46c:C:D:i:l:mp:rs:S:t:vzZ:")) != -1) {
  589. switch (c) {
  590. case '4':
  591. if (cfg_family != PF_UNSPEC)
  592. error(1, 0, "Pass one of -4 or -6");
  593. cfg_family = PF_INET;
  594. cfg_alen = sizeof(struct sockaddr_in);
  595. break;
  596. case '6':
  597. if (cfg_family != PF_UNSPEC)
  598. error(1, 0, "Pass one of -4 or -6");
  599. cfg_family = PF_INET6;
  600. cfg_alen = sizeof(struct sockaddr_in6);
  601. break;
  602. case 'c':
  603. cfg_cork = strtol(optarg, NULL, 0);
  604. break;
  605. case 'C':
  606. cfg_cpu = strtol(optarg, NULL, 0);
  607. break;
  608. case 'D':
  609. daddr = optarg;
  610. break;
  611. case 'i':
  612. cfg_ifindex = if_nametoindex(optarg);
  613. if (cfg_ifindex == 0)
  614. error(1, errno, "invalid iface: %s", optarg);
  615. break;
  616. case 'l':
  617. cfg_notification_limit = strtoul(optarg, NULL, 0);
  618. break;
  619. case 'm':
  620. cfg_cork_mixed = true;
  621. break;
  622. case 'p':
  623. cfg_port = strtoul(optarg, NULL, 0);
  624. break;
  625. case 'r':
  626. cfg_rx = true;
  627. break;
  628. case 's':
  629. cfg_payload_len = strtoul(optarg, NULL, 0);
  630. break;
  631. case 'S':
  632. saddr = optarg;
  633. break;
  634. case 't':
  635. cfg_runtime_ms = 200 + strtoul(optarg, NULL, 10) * 1000;
  636. break;
  637. case 'v':
  638. cfg_verbose++;
  639. break;
  640. case 'z':
  641. cfg_zerocopy = true;
  642. break;
  643. case 'Z':
  644. cfg_expect_zerocopy = !!atoi(optarg);
  645. break;
  646. }
  647. }
  648. cfg_test = argv[argc - 1];
  649. if (strcmp(cfg_test, "rds") == 0) {
  650. if (!daddr)
  651. error(1, 0, "-D <server addr> required for PF_RDS\n");
  652. if (!cfg_rx && !saddr)
  653. error(1, 0, "-S <client addr> required for PF_RDS\n");
  654. }
  655. setup_sockaddr(cfg_family, daddr, &cfg_dst_addr);
  656. setup_sockaddr(cfg_family, saddr, &cfg_src_addr);
  657. if (cfg_payload_len > max_payload_len)
  658. error(1, 0, "-s: payload exceeds max (%d)", max_payload_len);
  659. if (cfg_cork_mixed && (!cfg_zerocopy || !cfg_cork))
  660. error(1, 0, "-m: cork_mixed requires corking and zerocopy");
  661. if (optind != argc - 1)
  662. usage(argv[0]);
  663. }
  664. int main(int argc, char **argv)
  665. {
  666. const char *cfg_test;
  667. parse_opts(argc, argv);
  668. cfg_test = argv[argc - 1];
  669. if (!strcmp(cfg_test, "packet"))
  670. do_test(PF_PACKET, SOCK_RAW, 0);
  671. else if (!strcmp(cfg_test, "packet_dgram"))
  672. do_test(PF_PACKET, SOCK_DGRAM, 0);
  673. else if (!strcmp(cfg_test, "raw"))
  674. do_test(cfg_family, SOCK_RAW, IPPROTO_EGP);
  675. else if (!strcmp(cfg_test, "raw_hdrincl"))
  676. do_test(cfg_family, SOCK_RAW, IPPROTO_RAW);
  677. else if (!strcmp(cfg_test, "tcp"))
  678. do_test(cfg_family, SOCK_STREAM, 0);
  679. else if (!strcmp(cfg_test, "udp"))
  680. do_test(cfg_family, SOCK_DGRAM, 0);
  681. else if (!strcmp(cfg_test, "rds"))
  682. do_test(PF_RDS, SOCK_SEQPACKET, 0);
  683. else
  684. error(1, 0, "unknown cfg_test %s", cfg_test);
  685. return exitcode;
  686. }