From a389b8f2f9c55b76fe5a28c3f1fea798ed37aa98 Mon Sep 17 00:00:00 2001 From: Andrew Opalach Date: Mon, 13 Apr 2026 16:11:04 -0400 Subject: Allow packet_stream_reconnect() to cut off a connecting stream. Cleanup timer. Signed-off-by: Andrew Opalach --- src/multiplex.c | 6 +++--- src/packet_stream.c | 19 +++++++++++-------- src/packet_stream.h | 2 +- src/rpc.c | 3 ++- src/rpc2.h | 4 ++-- src/timer.c | 15 ++------------- src/timer.h | 3 --- src/util/thread/thread_linux.c | 4 ++-- 8 files changed, 23 insertions(+), 33 deletions(-) (limited to 'src') diff --git a/src/multiplex.c b/src/multiplex.c index 2de4621..0911540 100644 --- a/src/multiplex.c +++ b/src/multiplex.c @@ -108,9 +108,9 @@ static void direct_connect(struct nn_packet_stream *client, u8 id) multiplex_direct_global.connection_callback(multiplex_direct_global.userdata, id, server); server->direct = client; client->direct = server; - // Connected happens on the server first. - nn_packet_stream_connected(server); - nn_packet_stream_connected(client); + // Explicitly signal connected on the server first. + nn_packet_stream_set_connected(server); + nn_packet_stream_set_connected(client); } void nn_multiplex_direct_connect(struct nn_packet_stream *client, u8 id) diff --git a/src/packet_stream.c b/src/packet_stream.c index 7868242..0a45b22 100644 --- a/src/packet_stream.c +++ b/src/packet_stream.c @@ -104,7 +104,7 @@ static void stream_read_callback(struct ev_loop *loop, ev_io *w, s32 revents) } } -static bool stream_connected(struct nn_packet_stream *stream) +static bool set_stream_connected(struct nn_packet_stream *stream) { stream->connect = PACKET_STREAM_CONNECTED; al_assert(stream->corked); @@ -134,7 +134,7 @@ static void stream_write_callback(struct ev_loop *loop, ev_io *w, s32 revents) // Try again on next POLLOUT, not sure if this can actually happen. return; } - if (stream_connected(stream)) { + if (set_stream_connected(stream)) { ev_io_start(stream->loop->ev, &stream->revent); } } @@ -201,7 +201,7 @@ void nn_packet_stream_from_socket(struct nn_packet_stream *stream, struct nn_eve s32 fd = nn_socket_get_fd(&stream->sock); ev_io_init_n(&stream->revent, stream_read_callback, fd, EV_READ); ev_io_init_n(&stream->wevent, stream_write_callback, fd, EV_WRITE); - if (stream_connected(stream)) { + if (set_stream_connected(stream)) { ev_io_start(stream->loop->ev, &stream->revent); } } @@ -234,15 +234,20 @@ void nn_packet_stream_connect(struct nn_packet_stream *stream, struct nn_event_l do_connect_internal(stream, addr, port); } +// Keep in mind this stream will have had it's state reset in stop_internal(). void nn_packet_stream_reconnect(struct nn_packet_stream *stream, str *addr, u16 port) { + if (stream->connect == PACKET_STREAM_CONNECTING) { + stop_write_internal(stream); + stop_internal(stream); + } al_assert(stream->connect == PACKET_STREAM_DISCONNECTED); do_connect_internal(stream, addr, port); } -bool nn_packet_stream_connected(struct nn_packet_stream *stream) +bool nn_packet_stream_set_connected(struct nn_packet_stream *stream) { - return stream_connected(stream); + return set_stream_connected(stream); } void nn_packet_stream_cork(struct nn_packet_stream *stream, bool cork) @@ -314,10 +319,8 @@ void nn_packet_stream_return_packets(struct nn_packet_stream *stream, struct nn_ void nn_packet_stream_disconnect(struct nn_packet_stream *stream) { + al_assert(stream->connect != PACKET_STREAM_DISCONNECTED); al_assert(stream->connect != PACKET_STREAM_DISCONNECTING); - // If reusing this stream, keep in mind it will have been set back - // to a default state in stop_internal(). - if (stream->connect == PACKET_STREAM_DISCONNECTED) return; if (stream->direct) { struct nn_packet_stream *direct = stream->direct; direct->connect = PACKET_STREAM_DISCONNECTED; diff --git a/src/packet_stream.h b/src/packet_stream.h index 5531d00..d60af16 100644 --- a/src/packet_stream.h +++ b/src/packet_stream.h @@ -43,7 +43,7 @@ void nn_packet_stream_from_socket(struct nn_packet_stream *stream, struct nn_eve void nn_packet_stream_connect(struct nn_packet_stream *stream, struct nn_event_loop *loop, u8 id, u8 type, str *addr, u16 port); void nn_packet_stream_reconnect(struct nn_packet_stream *stream, str *addr, u16 port); -bool nn_packet_stream_connected(struct nn_packet_stream *stream); +bool nn_packet_stream_set_connected(struct nn_packet_stream *stream); void nn_packet_stream_cork(struct nn_packet_stream *stream, bool cork); bool nn_packet_stream_send_packet(struct nn_packet_stream *stream, struct nn_packet *packet); void nn_packet_stream_discard_queue(struct nn_packet_stream *stream); diff --git a/src/rpc.c b/src/rpc.c index 1b00e63..7938341 100644 --- a/src/rpc.c +++ b/src/rpc.c @@ -125,10 +125,11 @@ void nn_rpc_connect(struct nn_rpc *rpc, u8 id, u8 type, str *addr, u16 port) nn_packet_stream_connect(rpc->conn->stream, rpc->loop, id, type, addr, port); } -void nn_rpc_reconnect(struct nn_rpc *rpc, str *addr, u16 port) +struct nn_rpc_connection *nn_rpc_reconnect(struct nn_rpc *rpc, str *addr, u16 port) { al_assert(rpc->conn); nn_packet_stream_reconnect(rpc->conn->stream, addr, port); + return rpc->conn; } // @TODO: Pack opcode into a u32, reduce ID by 8 bits (make struct with bitmask) diff --git a/src/rpc2.h b/src/rpc2.h index e6c2098..bed8822 100644 --- a/src/rpc2.h +++ b/src/rpc2.h @@ -45,11 +45,11 @@ void nn_rpc_add_command(struct nn_rpc *rpc, struct nn_rpc_command *command); void nn_rpc_add_stream(struct nn_rpc *rpc, struct nn_packet_stream *stream); struct nn_rpc_connection *nn_rpc_prepare_client(struct nn_rpc *rpc); void nn_rpc_connect(struct nn_rpc *rpc, u8 id, u8 type, str *addr, u16 port); -void nn_rpc_reconnect(struct nn_rpc *rpc, str *addr, u16 port); + struct nn_rpc_connection *nn_rpc_reconnect(struct nn_rpc *rpc, str *addr, u16 port); struct nn_packet *nn_rpc_get_packet(struct nn_rpc *rpc, s8 op); void nn_rpc_free(struct nn_rpc *rpc); -// The order of commands per connection will be respected. +// The order of commands per-connection will be respected. void nn_rpc_connection_command(struct nn_rpc_connection *conn, struct nn_packet *packet, void (*callback)(void *, struct nn_rpc_connection *conn, struct nn_packet *), void *userdata); void nn_rpc_conn_flush(struct nn_rpc_connection *conn); diff --git a/src/timer.c b/src/timer.c index c0759cc..1d14a38 100644 --- a/src/timer.c +++ b/src/timer.c @@ -5,7 +5,7 @@ static void timer_callback(struct ev_loop *loop, ev_timer *w, s32 revents) (void)loop; struct nn_timer *timer = (struct nn_timer *)w->data; (void)revents; - if (!timer->disabled) timer->callback(timer->userdata, timer); + timer->callback(timer->userdata, timer); } void nn_timer_init(struct nn_timer *timer, struct nn_event_loop *loop, @@ -15,7 +15,6 @@ void nn_timer_init(struct nn_timer *timer, struct nn_event_loop *loop, timer->callback = callback; timer->userdata = userdata; timer->timer.data = timer; - timer->disabled = false; ev_init_n(&timer->timer, timer_callback); } @@ -26,17 +25,7 @@ void nn_timer_set_repeat(struct nn_timer *timer, nn_os_tstamp repeat) void nn_timer_again(struct nn_timer *timer) { - if (!timer->disabled) ev_timer_again(timer->loop->ev, &timer->timer); -} - -void nn_timer_disable(struct nn_timer *timer) -{ - timer->disabled = true; -} - -void nn_timer_enable(struct nn_timer *timer) -{ - timer->disabled = false; + ev_timer_again(timer->loop->ev, &timer->timer); } void nn_timer_stop(struct nn_timer *timer) diff --git a/src/timer.h b/src/timer.h index 15b467f..29107df 100644 --- a/src/timer.h +++ b/src/timer.h @@ -7,7 +7,6 @@ struct nn_timer { ev_timer timer; struct nn_event_loop *loop; - bool disabled; void (*callback)(void *, struct nn_timer *); void *userdata; }; @@ -16,6 +15,4 @@ void nn_timer_init(struct nn_timer *timer, struct nn_event_loop *loop, void (*callback)(void *, struct nn_timer *), void *userdata); void nn_timer_set_repeat(struct nn_timer *timer, nn_os_tstamp repeat); void nn_timer_again(struct nn_timer *timer); -void nn_timer_disable(struct nn_timer *timer); -void nn_timer_enable(struct nn_timer *timer); void nn_timer_stop(struct nn_timer *timer); diff --git a/src/util/thread/thread_linux.c b/src/util/thread/thread_linux.c index cd80d06..c9edce5 100644 --- a/src/util/thread/thread_linux.c +++ b/src/util/thread/thread_linux.c @@ -91,9 +91,9 @@ void nn_cond_init(struct nn_cond *cond) void nn_cond_wait(struct nn_cond *cond, struct nn_mutex *mutex) { cond->condition = false; - while (!cond->condition) { + do { pthread_cond_wait(&cond->cond, &mutex->mutex); - } + } while (!cond->condition); } bool nn_cond_is_waiting(struct nn_cond *cond) -- cgit v1.2.3-101-g0448