| 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 | } |