master
c 279 lines 8.9 KB
Raw
1 /*
2 * QEMU System Emulator
3 *
4 * Copyright (c) 2003-2008 Fabrice Bellard
5 * Copyright (c) 2022 Red Hat, Inc.
6 *
7 * Permission is hereby granted, free of charge, to any person obtaining a copy
8 * of this software and associated documentation files (the "Software"), to deal
9 * in the Software without restriction, including without limitation the rights
10 * to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
11 * copies of the Software, and to permit persons to whom the Software is
12 * furnished to do so, subject to the following conditions:
13 *
14 * The above copyright notice and this permission notice shall be included in
15 * all copies or substantial portions of the Software.
16 *
17 * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
18 * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
19 * FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL
20 * THE AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
21 * LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
22 * OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN
23 * THE SOFTWARE.
24 */
25
26 #include "qemu/osdep.h"
27
28 #include "net/net.h"
29 #include "clients.h"
30 #include "qapi/error.h"
31 #include "io/net-listener.h"
32 #include "qapi/qapi-events-net.h"
33 #include "qapi/qapi-visit-sockets.h"
34 #include "qapi/clone-visitor.h"
35
36 #include "stream_data.h"
37
38 typedef struct NetStreamState {
39 NetStreamData data;
40 uint32_t reconnect_ms;
41 guint timer_tag;
42 SocketAddress *addr;
43 } NetStreamState;
44
45 static void net_stream_arm_reconnect(NetStreamState *s);
46
47 static ssize_t net_stream_receive(NetClientState *nc, const uint8_t *buf,
48 size_t size)
49 {
50 NetStreamData *d = DO_UPCAST(NetStreamData, nc, nc);
51
52 return net_stream_data_receive(d, buf, size);
53 }
54
55 static gboolean net_stream_send(QIOChannel *ioc,
56 GIOCondition condition,
57 gpointer data)
58 {
59 if (net_stream_data_send(ioc, condition, data) == G_SOURCE_REMOVE) {
60 NetStreamState *s = DO_UPCAST(NetStreamState, data, data);
61
62 qapi_event_send_netdev_stream_disconnected(s->data.nc.name);
63 net_stream_arm_reconnect(s);
64
65 return G_SOURCE_REMOVE;
66 }
67
68 return G_SOURCE_CONTINUE;
69 }
70
71 static void net_stream_cleanup(NetClientState *nc)
72 {
73 NetStreamState *s = DO_UPCAST(NetStreamState, data.nc, nc);
74 g_clear_handle_id(&s->timer_tag, g_source_remove);
75 if (s->addr) {
76 qapi_free_SocketAddress(s->addr);
77 s->addr = NULL;
78 }
79 if (s->data.ioc) {
80 if (QIO_CHANNEL_SOCKET(s->data.ioc)->fd != -1) {
81 g_clear_handle_id(&s->data.ioc_read_tag, g_source_remove);
82 g_clear_handle_id(&s->data.ioc_write_tag, g_source_remove);
83 }
84 object_unref(OBJECT(s->data.ioc));
85 s->data.ioc = NULL;
86 }
87 if (s->data.listen_ioc) {
88 if (s->data.listener) {
89 qio_net_listener_disconnect(s->data.listener);
90 object_unref(OBJECT(s->data.listener));
91 s->data.listener = NULL;
92 }
93 object_unref(OBJECT(s->data.listen_ioc));
94 s->data.listen_ioc = NULL;
95 }
96 }
97
98 static NetClientInfo net_stream_info = {
99 .type = NET_CLIENT_DRIVER_STREAM,
100 .size = sizeof(NetStreamState),
101 .receive = net_stream_receive,
102 .cleanup = net_stream_cleanup,
103 };
104
105 static void net_stream_listen(QIONetListener *listener,
106 QIOChannelSocket *cioc, gpointer data)
107 {
108 NetStreamData *d = data;
109 SocketAddress *addr;
110 char *uri;
111
112 net_stream_data_listen(listener, cioc, data);
113
114 if (cioc->localAddr.ss_family == AF_UNIX) {
115 addr = qio_channel_socket_get_local_address(cioc, NULL);
116 } else {
117 addr = qio_channel_socket_get_remote_address(cioc, NULL);
118 }
119 g_assert(addr != NULL);
120 uri = socket_uri(addr);
121 qemu_set_info_str(&d->nc, "%s", uri);
122 g_free(uri);
123 qapi_event_send_netdev_stream_connected(d->nc.name, addr);
124 qapi_free_SocketAddress(addr);
125 }
126
127 static void net_stream_server_listening(QIOTask *task, gpointer opaque)
128 {
129 NetStreamData *d = opaque;
130 QIOChannelSocket *listen_sioc = QIO_CHANNEL_SOCKET(d->listen_ioc);
131 SocketAddress *addr;
132 Error *err = NULL;
133
134 if (qio_task_propagate_error(task, &err)) {
135 qemu_set_info_str(&d->nc, "error: %s", error_get_pretty(err));
136 error_free(err);
137 return;
138 }
139
140 addr = qio_channel_socket_get_local_address(listen_sioc, NULL);
141 g_assert(addr != NULL);
142 if (!qemu_set_blocking(listen_sioc->fd, false, &err)) {
143 qemu_set_info_str(&d->nc, "error: %s", error_get_pretty(err));
144 error_free(err);
145 return;
146 }
147 qapi_free_SocketAddress(addr);
148
149 d->nc.link_down = true;
150 d->listener = qio_net_listener_new();
151
152 qemu_set_info_str(&d->nc, "listening");
153 net_socket_rs_init(&d->rs, net_stream_data_rs_finalize, false);
154 qio_net_listener_set_client_func(d->listener, d->listen, d,
155 NULL);
156 qio_net_listener_add(d->listener, listen_sioc);
157 }
158
159 static int net_stream_server_init(NetClientState *peer,
160 const char *model,
161 const char *name,
162 SocketAddress *addr,
163 Error **errp)
164 {
165 NetClientState *nc;
166 NetStreamData *d;
167 QIOChannelSocket *listen_sioc = qio_channel_socket_new();
168
169 nc = qemu_new_net_client(&net_stream_info, peer, model, name);
170 d = DO_UPCAST(NetStreamData, nc, nc);
171 d->send = net_stream_send;
172 d->listen = net_stream_listen;
173 qemu_set_info_str(&d->nc, "initializing");
174
175 d->listen_ioc = QIO_CHANNEL(listen_sioc);
176 qio_channel_socket_listen_async(listen_sioc, addr, 0,
177 net_stream_server_listening, d,
178 NULL, NULL);
179
180 return 0;
181 }
182
183 static void net_stream_client_connected(QIOTask *task, gpointer opaque)
184 {
185 NetStreamState *s = opaque;
186 NetStreamData *d = &s->data;
187 QIOChannelSocket *sioc = QIO_CHANNEL_SOCKET(d->ioc);
188 SocketAddress *addr;
189 gchar *uri;
190
191 if (net_stream_data_client_connected(task, d) == -1) {
192 net_stream_arm_reconnect(s);
193 return;
194 }
195
196 addr = qio_channel_socket_get_remote_address(sioc, NULL);
197 g_assert(addr != NULL);
198 uri = socket_uri(addr);
199 qemu_set_info_str(&d->nc, "%s", uri);
200 g_free(uri);
201 qapi_event_send_netdev_stream_connected(d->nc.name, addr);
202 qapi_free_SocketAddress(addr);
203 }
204
205 static gboolean net_stream_reconnect(gpointer data)
206 {
207 NetStreamState *s = data;
208 QIOChannelSocket *sioc;
209
210 s->timer_tag = 0;
211
212 sioc = qio_channel_socket_new();
213 s->data.ioc = QIO_CHANNEL(sioc);
214 qio_channel_socket_connect_async(sioc, s->addr,
215 net_stream_client_connected, s,
216 NULL, NULL);
217 return G_SOURCE_REMOVE;
218 }
219
220 static void net_stream_arm_reconnect(NetStreamState *s)
221 {
222 if (s->reconnect_ms && s->timer_tag == 0) {
223 qemu_set_info_str(&s->data.nc, "connecting");
224 s->timer_tag = g_timeout_add(s->reconnect_ms, net_stream_reconnect, s);
225 }
226 }
227
228 static int net_stream_client_init(NetClientState *peer,
229 const char *model,
230 const char *name,
231 SocketAddress *addr,
232 uint32_t reconnect_ms,
233 Error **errp)
234 {
235 NetStreamState *s;
236 NetClientState *nc;
237 QIOChannelSocket *sioc = qio_channel_socket_new();
238
239 nc = qemu_new_net_client(&net_stream_info, peer, model, name);
240 s = DO_UPCAST(NetStreamState, data.nc, nc);
241 qemu_set_info_str(&s->data.nc, "connecting");
242
243 s->data.ioc = QIO_CHANNEL(sioc);
244 s->data.nc.link_down = true;
245 s->data.send = net_stream_send;
246 s->data.listen = net_stream_listen;
247
248 s->reconnect_ms = reconnect_ms;
249 if (reconnect_ms) {
250 s->addr = QAPI_CLONE(SocketAddress, addr);
251 }
252 qio_channel_socket_connect_async(sioc, addr,
253 net_stream_client_connected, s,
254 NULL, NULL);
255
256 return 0;
257 }
258
259 int net_init_stream(const Netdev *netdev, const char *name,
260 NetClientState *peer, Error **errp)
261 {
262 const NetdevStreamOptions *sock;
263
264 assert(netdev->type == NET_CLIENT_DRIVER_STREAM);
265 sock = &netdev->u.stream;
266
267 if (!sock->has_server || !sock->server) {
268 return net_stream_client_init(peer, "stream", name, sock->addr,
269 sock->has_reconnect_ms ?
270 sock->reconnect_ms : 0,
271 errp);
272 }
273 if (sock->has_reconnect_ms) {
274 error_setg(errp, "'reconnect-ms' option is "
275 "incompatible with socket in server mode");
276 return -1;
277 }
278 return net_stream_server_init(peer, "stream", name, sock->addr, errp);
279 }