From 571599f673784f53e8913de7ced5a45521f48f82 Mon Sep 17 00:00:00 2001 From: Ailin Nemui Date: Sun, 25 Jan 2026 22:37:22 +0100 Subject: [PATCH] stream and proxy code --- src/core/net-disconnect.c | 59 +++++++++++++ src/core/net-sendbuffer.c | 42 +++++++++ src/core/net-sendbuffer.h | 4 + src/core/network.c | 61 +++++++++++++ src/core/network.h | 6 ++ src/core/server-connect-rec.h | 2 + src/core/servers.c | 158 +++++++++++++++++++++++++++++++++- src/irc/core/irc-servers.c | 7 +- src/irc/core/irc.c | 20 +++++ 9 files changed, 357 insertions(+), 2 deletions(-) diff --git a/src/core/net-disconnect.c b/src/core/net-disconnect.c index ccdb8f7f..5941b27d 100644 --- a/src/core/net-disconnect.c +++ b/src/core/net-disconnect.c @@ -33,6 +33,7 @@ typedef struct { time_t created; GIOChannel *channel; + GIOStream *stream; int tag; } NET_DISCONNECT_REC; @@ -44,10 +45,13 @@ void net_disconnect_any(NET_DISCONNECT_REC *rec) { if (rec->channel != NULL) net_disconnect_channel(rec->channel); + else + net_disconnect_stream(rec->stream); } static void net_disconnect_remove(NET_DISCONNECT_REC *rec) { + g_warning("disconnect finished"); disconnects = g_slist_remove(disconnects, rec); g_source_remove(rec->tag); @@ -73,6 +77,35 @@ static void sig_disconnect(NET_DISCONNECT_REC *rec) } while (ret == sizeof(buf) && count < 18); } +// static gboolean sig_disconnect_source(GObject *pollable_stream, NET_DISCONNECT_REC *rec) +static gboolean sig_disconnect_source(GSocket *socket, GIOCondition cond, NET_DISCONNECT_REC *rec) +{ + char buf[512]; + int count, ret; + + g_warning("sig_disconnect_source, condition: %d", cond); + + if (cond & (G_IO_ERR | G_IO_HUP)) { + net_disconnect_remove(rec); + return FALSE; + } + + /* check if there's any data waiting in socket. read max. 9kB so + if server just keeps sending us stuff we won't get stuck */ + count = 0; + do { + ret = net_receive_stream(rec->stream, buf, sizeof(buf)); + if (ret == -1) { + /* socket was closed */ + net_disconnect_remove(rec); + return FALSE; + } + count++; + } while (ret == sizeof(buf) && count < 18); + + return TRUE; +} + static int sig_timeout_disconnect(void) { NET_DISCONNECT_REC *rec; @@ -109,6 +142,21 @@ void net_disconnect_later(NET_SENDBUF_REC *handle) rec->channel = handle->channel; rec->tag = i_input_add(rec->channel, I_INPUT_READ, (GInputFunction) sig_disconnect, rec); + } else if (handle->stream != NULL) { + // GInputStream *in; + GSocket *socket; + GSource *source; + + socket = g_socket_connection_get_socket((GSocketConnection *) handle->stream); + source = g_socket_create_source(socket, G_IO_IN | G_IO_HUP, NULL); + // in = g_io_stream_get_input_stream((GIOStream *) rec->stream); + // source = g_pollable_input_stream_create_source(G_POLLABLE_INPUT_STREAM(in), + // NULL); + g_source_set_callback(source, G_SOURCE_FUNC(sig_disconnect_source), rec, NULL); + rec->stream = handle->stream; + rec->tag = g_source_attach(source, NULL); + g_warning("net_disconnect_rec.tag: %d", rec->tag); + // TODO } if (timeout_tag == -1) { timeout_tag = g_timeout_add(10000, (GSourceFunc) @@ -155,6 +203,17 @@ void net_disconnect_deinit(void) /* data coming .. check if we can close the handle */ sig_disconnect(rec); } + } else if (rec->stream != NULL) { + // TODO! timeouts + /* GInputStream *iin; */ + /* GPollableInputStream *in; */ + + /* iin = g_io_stream_get_input_stream(rec->stream); */ + /* in = G_POLLABLE_INPUT_STREAM(iin); */ + /* if (g_pollable_input_stream_is_readable(in)) { */ + /* (void)sig_disconnect_source(g_socket_connection_get_socket((GSocketConnection + * *)rec->stream), G_IO_IN, rec); */ + /* } */ } else if (first) { /* Display the text when we have already waited for a while */ diff --git a/src/core/net-sendbuffer.c b/src/core/net-sendbuffer.c index 384bcdb6..62e40070 100644 --- a/src/core/net-sendbuffer.c +++ b/src/core/net-sendbuffer.c @@ -41,6 +41,21 @@ NET_SENDBUF_REC *net_sendbuffer_create_channel(GIOChannel *channel, int bufsize) return rec; } +NET_SENDBUF_REC *net_sendbuffer_create_stream(GIOStream *stream, int bufsize) +{ + NET_SENDBUF_REC *rec; + + g_return_val_if_fail(stream != NULL, NULL); + + rec = g_new0(NET_SENDBUF_REC, 1); + rec->send_tag = -1; + rec->stream = stream; + rec->bufsize = bufsize > 0 ? bufsize : DEFAULT_BUFFER_SIZE; + rec->def_bufsize = rec->bufsize; + + return rec; +} + /* Destroy the buffer. `close' specifies if socket handle should be closed. */ void net_sendbuffer_destroy(NET_SENDBUF_REC *rec, int close) { @@ -48,8 +63,12 @@ void net_sendbuffer_destroy(NET_SENDBUF_REC *rec, int close) if (close) { if (rec->channel != NULL) net_disconnect_channel(rec->channel); + else + net_disconnect_stream(rec->stream); } if (rec->readbuffer != NULL) line_split_free(rec->readbuffer); + if (rec->stream != NULL) + g_object_unref(rec->stream); g_free_not_null(rec->buffer); g_free(rec); } @@ -86,6 +105,14 @@ static void sig_sendbuffer(NET_SENDBUF_REC *rec) rec->send_tag = -1; } +static gboolean sig_sendbuffer_source(GObject *pollable_stream, NET_SENDBUF_REC *rec) +{ + g_warning("sig_sendbuffer_source"); + + sig_sendbuffer(rec); + return FALSE; +} + /* Add `data' to transmit buffer - return FALSE if buffer is full */ static int buffer_add(NET_SENDBUF_REC *rec, const void *data, int size) { @@ -125,6 +152,8 @@ int net_sendbuffer_send(NET_SENDBUF_REC *rec, const void *data, int size) /* nothing in buffer - transmit immediately */ if (rec->channel != NULL) ret = net_transmit_channel(rec->channel, data, size); + else if (rec->stream != NULL) + ret = net_transmit_stream(rec->stream, data, size); if (ret < 0) return -1; size -= ret; data = ((const char *) data) + ret; @@ -138,6 +167,17 @@ int net_sendbuffer_send(NET_SENDBUF_REC *rec, const void *data, int size) if (rec->channel != NULL) { rec->send_tag = i_input_add(rec->channel, I_INPUT_WRITE, (GInputFunction) sig_sendbuffer, rec); + } else if (rec->stream != NULL) { + GOutputStream *out; + GSource *source; + + out = g_io_stream_get_output_stream(rec->stream); + source = g_pollable_output_stream_create_source( + G_POLLABLE_OUTPUT_STREAM(out), NULL); + g_source_set_callback(source, G_SOURCE_FUNC(sig_sendbuffer_source), rec, + NULL); + rec->send_tag = g_source_attach(source, NULL); + g_warning("send_tag: %d", rec->send_tag); } } @@ -152,6 +192,8 @@ int net_sendbuffer_receive_line(NET_SENDBUF_REC *rec, char **str, int read_socke if (read_socket) { if (rec->channel != NULL) recvlen = net_receive_channel(rec->channel, tmpbuf, sizeof(tmpbuf)); + else if (rec->stream != NULL) + recvlen = net_receive_stream(rec->stream, tmpbuf, sizeof(tmpbuf)); } return line_split(tmpbuf, recvlen, str, &rec->readbuffer); diff --git a/src/core/net-sendbuffer.h b/src/core/net-sendbuffer.h index 5e7fbe96..61da1697 100644 --- a/src/core/net-sendbuffer.h +++ b/src/core/net-sendbuffer.h @@ -8,6 +8,7 @@ struct _NET_SENDBUF_REC { GIOChannel *channel; + GIOStream *stream; LINEBUF_REC *readbuffer; /* receive buffer */ int send_tag; @@ -21,6 +22,9 @@ struct _NET_SENDBUF_REC { /* Create new buffer - if `bufsize' is zero or less, DEFAULT_BUFFER_SIZE is used */ NET_SENDBUF_REC *net_sendbuffer_create_channel(GIOChannel *channel, int bufsize); +/* Create new buffer - if `bufsize' is zero or less, DEFAULT_BUFFER_SIZE + is used */ +NET_SENDBUF_REC *net_sendbuffer_create_stream(GIOStream *stream, int bufsize); /* Destroy the buffer. `close' specifies if socket handle should be closed. */ void net_sendbuffer_destroy(NET_SENDBUF_REC *rec, int close); diff --git a/src/core/network.c b/src/core/network.c index a5ebdb88..38ce1255 100644 --- a/src/core/network.c +++ b/src/core/network.c @@ -241,6 +241,14 @@ void net_disconnect_channel(GIOChannel *channel) g_io_channel_unref(channel); } +void net_disconnect_stream(GIOStream *stream) +{ + g_return_if_fail(stream != NULL); + + g_io_stream_close(stream, NULL, NULL); + g_object_unref(stream); +} + /* Listen for connections on a socket. if `my_ip' is NULL, listen in any address. */ GIOChannel *net_listen_channel(IPADDR *my_ip, int *port) @@ -338,6 +346,59 @@ int net_receive_channel(GIOChannel *channel, char *buf, int len) return ret; } +int net_receive_stream(GIOStream *stream, char *buf, int len) +{ + GInputStream *iin; + GPollableInputStream *in; + GError *error; + gsize ret; + + g_return_val_if_fail(stream != NULL, -1); + g_return_val_if_fail(buf != NULL, -1); + + error = NULL; + iin = g_io_stream_get_input_stream(stream); + in = G_POLLABLE_INPUT_STREAM(iin); + + ret = g_pollable_input_stream_read_nonblocking(in, buf, len, NULL, &error); + if (error == NULL) { + if (ret == 0) { + g_warning("net_receive returned has_pending:%d, is_closed:%d, connection " + "is_closed:%d, %lu [%s]", + g_input_stream_has_pending(iin), g_input_stream_is_closed(iin), + g_io_stream_is_closed(stream), ret, buf); + return -1; + } + return ret; + } else if (error->code == G_IO_ERROR_WOULD_BLOCK) { + // g_warning("net_receive would block: %s, is_connected:%d", error->message, + // g_socket_connection_is_connected((GSocketConnection *)stream)); + return 0; + } else { + g_warning("net_receive failed: %d:%s", error->code, error->message); + return -1; + } +} + +int net_transmit_stream(GIOStream *stream, const char *data, int len) +{ + GPollableOutputStream *out; + gsize ret; + GError *err = NULL; + + g_return_val_if_fail(stream != NULL, -1); + g_return_val_if_fail(data != NULL, -1); + + out = G_POLLABLE_OUTPUT_STREAM(g_io_stream_get_output_stream(stream)); + ret = g_pollable_output_stream_write_nonblocking(out, data, len, NULL, &err); + if (err == NULL || err->code == G_IO_ERROR_WOULD_BLOCK) { + return ret; + } else { + g_warning("net_transmit: %d:%s", err->code, err->message); + return -1; + } +} + /* Transmit data, return number of bytes sent, -1 = error */ int net_transmit_channel(GIOChannel *channel, const char *data, int len) { diff --git a/src/core/network.h b/src/core/network.h index 14e6133e..c0fed718 100644 --- a/src/core/network.h +++ b/src/core/network.h @@ -68,6 +68,8 @@ GIOChannel *net_connect_ip_channel(IPADDR *ip, int port, IPADDR *my_ip); GIOChannel *net_connect_unix_channel(const char *path); /* Disconnect socket */ void net_disconnect_channel(GIOChannel *channel); +/* Disconnect socket */ +void net_disconnect_stream(GIOStream *stream); /* Listen for connections on a socket */ GIOChannel *net_listen_channel(IPADDR *my_ip, int *port); @@ -76,8 +78,12 @@ GIOChannel *net_accept_channel(GIOChannel *channel, IPADDR *addr, int *port); /* Read data from socket, return number of bytes read, -1 = error */ int net_receive_channel(GIOChannel *channel, char *buf, int len); +/* Read data from stream, return number of bytes read, -1 = error */ +int net_receive_stream(GIOStream *stream, char *buf, int len); /* Transmit data, return number of bytes sent, -1 = error */ int net_transmit_channel(GIOChannel *channel, const char *data, int len); +/* Transmit data, return number of bytes sent, -1 = error */ +int net_transmit_stream(GIOStream *stream, const char *data, int len); /* Get the first IP address for host, both IPv4 and IPv6 if possible. */ int net_gethostbyname_first_ips(const char *addr, GResolverNameLookupFlags flags, IPADDR *ip4, diff --git a/src/core/server-connect-rec.h b/src/core/server-connect-rec.h index c11bae54..ebb2ba44 100644 --- a/src/core/server-connect-rec.h +++ b/src/core/server-connect-rec.h @@ -9,6 +9,7 @@ int refcount; char *proxy; int proxy_port; char *proxy_string, *proxy_string_after, *proxy_password; +GProxyResolver *proxy_resolver; unsigned short family; /* 0 = don't care, AF_INET or AF_INET6 */ unsigned short chosen_family; /* family actually chosen during name resolution */ @@ -33,6 +34,7 @@ char *tls_capath; char *tls_ciphers; char *tls_pinned_cert; char *tls_pinned_pubkey; +GSocketConnectable *tls_identity; GIOChannel *connect_channel; /* connect using this handle */ diff --git a/src/core/servers.c b/src/core/servers.c index 7993df46..f0a40176 100644 --- a/src/core/servers.c +++ b/src/core/servers.c @@ -164,6 +164,27 @@ static void server_connect_callback_init_channel(SERVER_REC *server, GIOChannel server_connect_finished(server); } +static gboolean server_connect_callback_init_source(GSocket *socket, GIOCondition cond, + SERVER_REC *server) +{ + g_return_val_if_fail(IS_SERVER(server), FALSE); + + if (cond & (G_IO_ERR | G_IO_HUP)) { + server->connection_lost = TRUE; + server->connrec->last_failed = server->connrec->last_connected; + server_connect_failed(server, "Connection broken"); // TODO: g_strerror(error) + return FALSE; + } + + lookup_servers = g_slist_remove(lookup_servers, server); + g_source_remove(server->connect_tag); + server->connect_tag = -1; + + server_connect_finished(server); + + return FALSE; +} + static void server_connect_callback_init_ssl_channel(SERVER_REC *server, GIOChannel *channel) { int error; @@ -195,13 +216,96 @@ static void server_connect_callback_init_ssl_channel(SERVER_REC *server, GIOChan server_connect_finished(server); } +void server_real_connect_connected(GSocketClient *client, GAsyncResult *res, SERVER_REC *server) +{ + GIOStream *stream; + // int fd; + // GIOChannel *channel; + GError *error; + error = NULL; + stream = (GIOStream *) g_socket_client_connect_finish(client, res, &error); + g_object_unref(client); + if (error != NULL) { + const char *errormsg = error->message; + if (errormsg == NULL) + errormsg = "Connect failed"; + if (error->code == G_IO_ERROR_CONNECTION_REFUSED) { +#if 0 + // TODO + if (own_ip != NULL) { + /* show the IP which is causing the error */ + net_ip2host(own_ip, ipaddr); + errmsg2 = g_strconcat(errmsg, ": ", ipaddr, NULL); + } +#endif + server->no_reconnect = TRUE; + } +#if 0 + // TODO + if (server->connrec->use_tls && errno == ENOSYS) + server->no_reconnect = TRUE; +#endif + + server->connection_lost = TRUE; + if (server->connrec->ipaddr != NULL) { + server->connrec->last_failed = server->connrec->last_connected; + } + server_connect_failed(server, errormsg); + } else { + GSocket *socket; + GSource *source; + + server->connrec->last_failed = 0; + g_object_ref(stream); // TODO: unref + g_assert(server->handle == NULL); + server->handle = net_sendbuffer_create_stream(stream, 0); + socket = g_socket_connection_get_socket((GSocketConnection *) stream); + source = g_socket_create_source(socket, G_IO_IN | G_IO_OUT, NULL); + g_source_set_callback(source, G_SOURCE_FUNC(server_connect_callback_init_source), + server, NULL); + server->connect_tag = g_source_attach(source, NULL); + } +} + +static void server_real_connect_event(GSocketClient *client, GSocketClientEvent event, + GSocketConnectable *connectable, GIOStream *stream, + SERVER_REC *server) +{ + if (event == G_SOCKET_CLIENT_TLS_HANDSHAKING) { + server->connrec->tls_identity = + g_network_address_new(server->connrec->address, server->connrec->port); + g_tls_client_connection_set_server_identity((GTlsClientConnection *) stream, + server->connrec->tls_identity); + if (server->connrec->tls_cert != NULL) { + GTlsCertificate *cert; + GError *err; + char *scert; + scert = convert_home(server->connrec->tls_cert); + err = NULL; + cert = g_tls_certificate_new_from_file(scert, &err); + if (err != NULL) { + g_warning("Tls Certificate: %s", err->message); + } else { + g_tls_connection_set_certificate((GTlsConnection *) stream, cert); + } + if (cert != NULL) { + g_object_unref(cert); + } + g_free(scert); + } + } +} + static void server_real_connect(SERVER_REC *server, IPADDR *ip, const char *unix_socket, const char *host) { + /* GIOChannel *channel; const char *errmsg; char *errmsg2; IPADDR *own_ip = NULL; + */ + GSocketClient *client; char ipaddr[MAX_IP_LEN]; int port = 0; @@ -218,6 +322,7 @@ static void server_real_connect(SERVER_REC *server, IPADDR *ip, const char *unix if (server->connrec->no_connect) return; +/* if (ip != NULL) { own_ip = IPADDR_IS_V6(ip) ? server->connrec->own_ip6 : server->connrec->own_ip4; port = server->connrec->proxy != NULL ? @@ -225,8 +330,44 @@ static void server_real_connect(SERVER_REC *server, IPADDR *ip, const char *unix channel = net_connect_ip_channel(ip, port, own_ip); } else { channel = net_connect_unix_channel(unix_socket); +*/ + + client = g_socket_client_new(); + // server->connect_cancellable = g_cancellable_new(); + if (ip != NULL || host != NULL) { + // own_ip = IPADDR_IS_V6(ip) ? server->connrec->own_ip6 : server->connrec->own_ip4; + port = /* server->connrec->proxy != NULL ? + server->connrec->proxy_port : */ + server->connrec->port; + if (server->connrec->use_tls) + g_socket_client_set_tls(client, TRUE); + if (server->connrec->proxy != NULL) { + GProxyResolver *res; + char *addr; + addr = g_strdup_printf("socks://%s:%d", server->connrec->proxy, + server->connrec->proxy_port); + res = g_simple_proxy_resolver_new(addr, NULL); + // g_free(addr); // TODO: free + g_warning("configuring %s as proxy", addr); + g_socket_client_set_proxy_resolver(client, res); + g_socket_client_set_enable_proxy(client, TRUE); + server->connrec->proxy_resolver = res; + } + g_signal_connect(client, "event", G_CALLBACK(server_real_connect_event), server); + g_socket_client_connect_to_host_async( + client, ip ? ipaddr : host, port, server->connect_cancellable, + (GAsyncReadyCallback) server_real_connect_connected, server); + return; + } else { + GSocketAddress *addr; + addr = g_unix_socket_address_new(unix_socket); + g_socket_client_connect_async( + client, (GSocketConnectable *) addr, server->connect_cancellable, + (GAsyncReadyCallback) server_real_connect_connected, server); + return; } +#if 0 if (server->connrec->use_tls && channel != NULL) { server->handle = net_sendbuffer_create_channel(channel, 0); channel = net_start_ssl_channel(server); @@ -270,6 +411,7 @@ static void server_real_connect(SERVER_REC *server, IPADDR *ip, const char *unix channel, I_INPUT_WRITE | I_INPUT_READ, (GInputFunction) server_connect_callback_init_channel, server); } +#endif } static int server_start_connect_resolve(SERVER_REC *server); @@ -280,6 +422,11 @@ static void server_connect_use_resolved(SERVER_REC *server) const char *errormsg; RESOLVED_IP_REC *iprec = server->connrec->resolved_host; + if (server->connrec->proxy != NULL) { + server_real_connect(server, NULL, NULL, server->connrec->address); + errormsg = NULL; + return; + } if (iprec->error != NULL) { /* error */ ip = NULL; @@ -368,8 +515,12 @@ static int server_start_connect_resolve(SERVER_REC *server) const char *connect_address; GResolverNameLookupFlags net_gethostbyname_flags; + if (server->connrec->proxy != NULL) + return TRUE; + connect_address = - server->connrec->proxy != NULL ? server->connrec->proxy : server->connrec->address; + /* server->connrec->proxy != NULL ? server->connrec->proxy : */ server->connrec + ->address; net_gethostbyname_flags = G_RESOLVER_NAME_LOOKUP_FLAGS_DEFAULT; if (server->connrec->family == AF_INET) { net_gethostbyname_flags = G_RESOLVER_NAME_LOOKUP_FLAGS_IPV4_ONLY; @@ -441,6 +592,7 @@ int server_start_connect(SERVER_REC *server) server->rawlog = rawlog_create(); + // TODO connect_connection if (server->connrec->connect_channel != NULL) { /* already connected */ GIOChannel *channel = server->connrec->connect_channel; @@ -661,7 +813,11 @@ void server_connect_unref(SERVER_CONNECT_REC *conn) g_free_not_null(conn->proxy_string); g_free_not_null(conn->proxy_string_after); g_free_not_null(conn->proxy_password); + if (conn->proxy_resolver != NULL) + g_object_unref(conn->proxy_resolver); + if (conn->tls_identity != NULL) + g_object_unref(conn->tls_identity); g_free_not_null(conn->ipaddr); g_free_not_null(conn->tag); g_free_not_null(conn->address); diff --git a/src/irc/core/irc-servers.c b/src/irc/core/irc-servers.c index a369707f..fabf9b29 100644 --- a/src/irc/core/irc-servers.c +++ b/src/irc/core/irc-servers.c @@ -387,6 +387,8 @@ static void event_starttls(IRC_SERVER_REC *server, const char *data) char *str; line_split("", -1, &str, &server->handle->readbuffer); } +#if 0 + // TODO ssl_channel = net_start_ssl_channel((SERVER_REC *) server); if (ssl_channel != NULL) { g_source_remove(server->readtag); @@ -394,8 +396,11 @@ static void event_starttls(IRC_SERVER_REC *server, const char *data) server->handle->channel = ssl_channel; init_ssl_loop_channel(server, server->handle->channel); } else { - g_warning("net_start_ssl failed"); +#endif + g_warning("net_start_ssl failed"); +#if 0 } +#endif } static void event_registerfirst(IRC_SERVER_REC *server, const char *data) diff --git a/src/irc/core/irc.c b/src/irc/core/irc.c index b44aa0c9..c3b17ad8 100644 --- a/src/irc/core/irc.c +++ b/src/irc/core/irc.c @@ -572,6 +572,16 @@ static void irc_parse_incoming(SERVER_REC *server) server_unref(server); } +static gboolean irc_parse_incoming_source(GObject *pollable_stream, SERVER_REC *server) +{ + g_return_val_if_fail(server != NULL, FALSE); + // g_warning("irc_parse_incoming_source"); + + irc_parse_incoming(server); + + return TRUE; +} + static void irc_init_server(IRC_SERVER_REC *server) { g_return_if_fail(server != NULL); @@ -582,6 +592,16 @@ static void irc_init_server(IRC_SERVER_REC *server) if (server->handle->channel != NULL) { server->readtag = i_input_add(net_sendbuffer_channel(server->handle), I_INPUT_READ, (GInputFunction) irc_parse_incoming, server); + } else { + GInputStream *in; + GSource *source; + + in = g_io_stream_get_input_stream(server->handle->stream); + source = g_pollable_input_stream_create_source(G_POLLABLE_INPUT_STREAM(in), NULL); + g_source_set_callback(source, G_SOURCE_FUNC(irc_parse_incoming_source), server, + NULL); + server->readtag = g_source_attach(source, NULL); + g_warning("readtag: %d", server->readtag); } }