master
c 188 lines 4.95 KB
Raw
1 /*
2 * net stream generic functions
3 *
4 * Copyright Red Hat
5 *
6 * SPDX-License-Identifier: GPL-2.0-or-later
7 */
8
9 #include "qemu/osdep.h"
10 #include "qemu/iov.h"
11 #include "qapi/error.h"
12 #include "net/net.h"
13 #include "io/channel.h"
14 #include "io/net-listener.h"
15 #include "qemu/sockets.h"
16
17 #include "stream_data.h"
18
19 static gboolean net_stream_data_writable(QIOChannel *ioc,
20 GIOCondition condition, gpointer data)
21 {
22 NetStreamData *d = data;
23
24 d->ioc_write_tag = 0;
25
26 qemu_flush_queued_packets(&d->nc);
27
28 return G_SOURCE_REMOVE;
29 }
30
31 ssize_t net_stream_data_receive(NetStreamData *d, const uint8_t *buf,
32 size_t size)
33 {
34 uint32_t len = htonl(size);
35 struct iovec iov[] = {
36 {
37 .iov_base = &len,
38 .iov_len = sizeof(len),
39 }, {
40 .iov_base = (void *)buf,
41 .iov_len = size,
42 },
43 };
44 struct iovec local_iov[2];
45 unsigned int nlocal_iov;
46 size_t remaining;
47 ssize_t ret;
48
49 remaining = iov_size(iov, 2) - d->send_index;
50 nlocal_iov = iov_copy(local_iov, 2, iov, 2, d->send_index, remaining);
51 ret = qio_channel_writev(d->ioc, local_iov, nlocal_iov, NULL);
52 if (ret == QIO_CHANNEL_ERR_BLOCK) {
53 ret = 0; /* handled further down */
54 }
55 if (ret == -1) {
56 d->send_index = 0;
57 return -errno;
58 }
59 if (ret < (ssize_t)remaining) {
60 d->send_index += ret;
61 d->ioc_write_tag = qio_channel_add_watch(d->ioc, G_IO_OUT,
62 net_stream_data_writable, d,
63 NULL);
64 return 0;
65 }
66 d->send_index = 0;
67 return size;
68 }
69
70 static void net_stream_data_send_completed(NetClientState *nc, ssize_t len)
71 {
72 NetStreamData *d = DO_UPCAST(NetStreamData, nc, nc);
73
74 if (!d->ioc_read_tag) {
75 d->ioc_read_tag = qio_channel_add_watch(d->ioc, G_IO_IN, d->send, d,
76 NULL);
77 }
78 }
79
80 void net_stream_data_rs_finalize(SocketReadState *rs)
81 {
82 NetStreamData *d = container_of(rs, NetStreamData, rs);
83
84 if (qemu_send_packet_async(&d->nc, rs->buf,
85 rs->packet_len,
86 net_stream_data_send_completed) == 0) {
87 g_clear_handle_id(&d->ioc_read_tag, g_source_remove);
88 }
89 }
90
91 gboolean net_stream_data_send(QIOChannel *ioc, GIOCondition condition,
92 NetStreamData *d)
93 {
94 int size;
95 int ret;
96 QEMU_UNINITIALIZED char buf1[NET_BUFSIZE];
97 const char *buf;
98
99 size = qio_channel_read(d->ioc, buf1, sizeof(buf1), NULL);
100 if (size < 0) {
101 if (errno != EWOULDBLOCK) {
102 goto eoc;
103 }
104 } else if (size == 0) {
105 /* end of connection */
106 eoc:
107 d->ioc_read_tag = 0;
108 if (d->ioc_write_tag) {
109 g_source_remove(d->ioc_write_tag);
110 d->ioc_write_tag = 0;
111 }
112 if (d->listener) {
113 qemu_set_info_str(&d->nc, "listening");
114 qio_net_listener_set_client_func(d->listener,
115 d->listen, d, NULL);
116 }
117 object_unref(OBJECT(d->ioc));
118 d->ioc = NULL;
119
120 net_socket_rs_init(&d->rs, net_stream_data_rs_finalize, false);
121 d->nc.link_down = true;
122
123 return G_SOURCE_REMOVE;
124 }
125 buf = buf1;
126
127 ret = net_fill_rstate(&d->rs, (const uint8_t *)buf, size);
128
129 if (ret == -1) {
130 goto eoc;
131 }
132
133 return G_SOURCE_CONTINUE;
134 }
135
136 void net_stream_data_listen(QIONetListener *listener,
137 QIOChannelSocket *cioc,
138 NetStreamData *d)
139 {
140 object_ref(OBJECT(cioc));
141
142 qio_net_listener_set_client_func(d->listener, NULL, d, NULL);
143
144 d->ioc = QIO_CHANNEL(cioc);
145 qio_channel_set_name(d->ioc, "stream-server");
146 d->nc.link_down = false;
147
148 d->ioc_read_tag = qio_channel_add_watch(d->ioc, G_IO_IN, d->send, d, NULL);
149 }
150
151 int net_stream_data_client_connected(QIOTask *task, NetStreamData *d)
152 {
153 QIOChannelSocket *sioc = QIO_CHANNEL_SOCKET(d->ioc);
154 SocketAddress *addr;
155 Error *err = NULL;
156
157 if (qio_task_propagate_error(task, &err)) {
158 qemu_set_info_str(&d->nc, "error: %s", error_get_pretty(err));
159 error_free(err);
160 goto error;
161 }
162
163 addr = qio_channel_socket_get_remote_address(sioc, NULL);
164 g_assert(addr != NULL);
165
166 if (!qemu_set_blocking(sioc->fd, false, &err)) {
167 qemu_set_info_str(&d->nc, "error: %s", error_get_pretty(err));
168 error_free(err);
169 qapi_free_SocketAddress(addr);
170 goto error;
171 }
172 qapi_free_SocketAddress(addr);
173
174 net_socket_rs_init(&d->rs, net_stream_data_rs_finalize, false);
175
176 /* Disable Nagle algorithm on TCP sockets to reduce latency */
177 qio_channel_set_delay(d->ioc, false);
178
179 d->ioc_read_tag = qio_channel_add_watch(d->ioc, G_IO_IN, d->send, d, NULL);
180 d->nc.link_down = false;
181
182 return 0;
183 error:
184 object_unref(OBJECT(d->ioc));
185 d->ioc = NULL;
186
187 return -1;
188 }