stream and proxy code

This commit is contained in:
Ailin Nemui 2026-01-25 22:37:22 +01:00
commit 571599f673
9 changed files with 357 additions and 2 deletions

View file

@ -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 */

View file

@ -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);

View file

@ -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);

View file

@ -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)
{

View file

@ -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,

View file

@ -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 */

View file

@ -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);

View file

@ -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)

View file

@ -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);
}
}