diff --git a/common.c b/common.c index 36fa183c..e8cc16ac 100644 --- a/common.c +++ b/common.c @@ -92,6 +92,8 @@ shairport_cfg config; int debuglev = 0; +sigset_t pselect_sigset; + int get_requested_connection_state_to_output() { return requested_connection_state_to_output; } void set_requested_connection_state_to_output(int v) { requested_connection_state_to_output = v; } diff --git a/common.h b/common.h index 4c2584e9..20d142d4 100644 --- a/common.h +++ b/common.h @@ -4,6 +4,7 @@ #include #include #include +#include #include "audio.h" #include "config.h" @@ -219,4 +220,6 @@ void command_set_volume(double volume); void shairport_shutdown(); // void shairport_startup_complete(void); +extern sigset_t pselect_sigset; + #endif // _COMMON_H diff --git a/mdns_external.c b/mdns_external.c index 86ad208b..c4aecd11 100644 --- a/mdns_external.c +++ b/mdns_external.c @@ -28,7 +28,6 @@ #include "common.h" #include #include -#include #include #include #include diff --git a/player.c b/player.c index 92d7e575..0b3200a8 100644 --- a/player.c +++ b/player.c @@ -35,7 +35,6 @@ #include #include #include -#include #include #include #include diff --git a/rtp.c b/rtp.c index 89c9c736..ba12744b 100644 --- a/rtp.c +++ b/rtp.c @@ -32,17 +32,19 @@ #include #include #include -#include #include #include #include #include #include +#include #include "common.h" #include "player.h" #include "rtp.h" +void memory_barrier(); + void rtp_initialise(rtsp_conn_info *conn) { conn->rtp_running = 0; @@ -81,6 +83,15 @@ void *rtp_audio_receiver(void *arg) { ssize_t nread; while (conn->please_stop == 0) { + fd_set readfds; + FD_ZERO(&readfds); + FD_SET(conn->audio_socket, &readfds); + do { + memory_barrier(); + } while (conn->please_stop == 0 && pselect(conn->audio_socket + 1, &readfds, NULL, NULL, NULL, &pselect_sigset) <= 0); + if (conn->please_stop != 0) { + break; + } nread = recv(conn->audio_socket, packet, sizeof(packet), 0); uint64_t local_time_now_fp = get_absolute_time_in_fp(); @@ -173,6 +184,15 @@ void *rtp_control_receiver(void *arg) { int64_t sync_rtp_timestamp, rtp_timestamp_less_latency; ssize_t nread; while (conn->please_stop == 0) { + fd_set readfds; + FD_ZERO(&readfds); + FD_SET(conn->control_socket, &readfds); + do { + memory_barrier(); + } while (conn->please_stop == 0 && pselect(conn->control_socket + 1, &readfds, NULL, NULL, NULL, &pselect_sigset) <= 0); + if (conn->please_stop != 0) { + break; + } nread = recv(conn->control_socket, packet, sizeof(packet), 0); local_time_now = get_absolute_time_in_fp(); // clock_gettime(CLOCK_MONOTONIC,&tn); @@ -308,6 +328,15 @@ void *rtp_timing_sender(void *arg) { msgsize = sizeof(struct sockaddr_in6); } #endif + fd_set writefds; + FD_ZERO(&writefds); + FD_SET(conn->timing_socket, &writefds); + do { + memory_barrier(); + } while (conn->timing_sender_stop == 0 && pselect(conn->timing_socket + 1, NULL, &writefds, NULL, NULL, &pselect_sigset) <= 0); + if (conn->timing_sender_stop != 0) { + break; + } if (sendto(conn->timing_socket, &req, sizeof(req), 0, (struct sockaddr *)&conn->rtp_client_timing_socket, msgsize) == -1) { perror("Error sendto-ing to timing socket"); @@ -344,6 +373,15 @@ void *rtp_timing_receiver(void *arg) { uint64_t first_local_to_remote_time_difference_time; uint64_t l2rtd = 0; while (conn->please_stop == 0) { + fd_set readfds; + FD_ZERO(&readfds); + FD_SET(conn->timing_socket, &readfds); + do { + memory_barrier(); + } while (conn->please_stop == 0 && pselect(conn->timing_socket + 1, &readfds, NULL, NULL, NULL, &pselect_sigset) <= 0); + if (conn->please_stop != 0) { + break; + } nread = recv(conn->timing_socket, packet, sizeof(packet), 0); arrival_time = get_absolute_time_in_fp(); // clock_gettime(CLOCK_MONOTONIC,&att); @@ -589,6 +627,7 @@ static int bind_port(int ip_family, const char *self_ip_address, uint32_t scope_ struct sockaddr_in *sa = (struct sockaddr_in *)&local; sport = ntohs(sa->sin_port); } + fcntl(local_socket, F_SETFL, O_NONBLOCK); *sock = local_socket; return sport; diff --git a/rtsp.c b/rtsp.c index 792600a9..4e9749e0 100644 --- a/rtsp.c +++ b/rtsp.c @@ -34,7 +34,6 @@ #include #include #include -#include #include #include #include @@ -466,7 +465,12 @@ static enum rtsp_read_request_response rtsp_read_request(rtsp_conn_info *conn, int msg_size = -1; while (msg_size < 0) { - memory_barrier(); + fd_set readfds; + FD_ZERO(&readfds); + FD_SET(conn->fd, &readfds); + do { + memory_barrier(); + } while (conn->stop == 0 && pselect(conn->fd + 1, &readfds, NULL, NULL, NULL, &pselect_sigset) <= 0); if (conn->stop != 0) { debug(3, "RTSP conversation thread %d shutdown requested.", conn->connection_number); reply = rtsp_read_request_response_immediate_shutdown_requested; @@ -540,6 +544,17 @@ static enum rtsp_read_request_response rtsp_read_request(rtsp_conn_info *conn, warning_message_sent = 1; } } + fd_set readfds; + FD_ZERO(&readfds); + FD_SET(conn->fd, &readfds); + do { + memory_barrier(); + } while (conn->stop == 0 && pselect(conn->fd + 1, &readfds, NULL, NULL, NULL, &pselect_sigset) <= 0); + if (conn->stop != 0) { + debug(1, "RTSP shutdown requested."); + reply = rtsp_read_request_response_immediate_shutdown_requested; + goto shutdown; + } ssize_t read_chunk = msg_size - inbuf; if (read_chunk > max_read_chunk) read_chunk = max_read_chunk; @@ -1026,6 +1041,7 @@ void metadata_create(void) { } else { int buffer_size = METADATA_SNDBUF; setsockopt(metadata_sock, SOL_SOCKET, SO_SNDBUF, &buffer_size, sizeof(buffer_size)); + fcntl(fd, F_SETFL, O_NONBLOCK); bzero((char *)&metadata_sockaddr, sizeof(metadata_sockaddr)); metadata_sockaddr.sin_family = AF_INET; metadata_sockaddr.sin_addr.s_addr = inet_addr(config.metadata_sockaddr); @@ -1773,12 +1789,6 @@ authenticate: } static void *rtsp_conversation_thread_func(void *pconn) { - // SIGUSR1 is used to interrupt this thread if blocked for read - sigset_t set; - sigemptyset(&set); - sigaddset(&set, SIGUSR1); - pthread_sigmask(SIG_UNBLOCK, &set, NULL); - rtsp_conn_info *conn = pconn; rtp_initialise(conn); @@ -1821,7 +1831,15 @@ static void *rtsp_conversation_thread_func(void *pconn) { } debug(3, "RTSP thread %d: RTSP Response:", conn->connection_number); debug_print_msg_headers(3, resp); - msg_write_response(conn->fd, resp); + fd_set writefds; + FD_ZERO(&writefds); + FD_SET(conn->fd, &writefds); + do { + memory_barrier(); + } while (conn->stop == 0 && pselect(conn->fd + 1, NULL, &writefds, NULL, NULL, &pselect_sigset) <= 0); + if (conn->stop == 0) { + msg_write_response(conn->fd, resp); + } msg_free(req); msg_free(resp); } else { @@ -2063,6 +2081,8 @@ void rtsp_listen_loop(void) { // conn->thread = rtsp_conversation_thread; // conn->stop = 0; // record's memory has been zeroed // conn->authorized = 0; // record's memory has been zeroed + fcntl(conn->fd, F_SETFL, O_NONBLOCK); + ret = pthread_create(&conn->thread, NULL, rtsp_conversation_thread_func, conn); // also acts as a memory barrier if (ret) diff --git a/shairport.c b/shairport.c index 033bb98c..b2f68b4a 100644 --- a/shairport.c +++ b/shairport.c @@ -33,7 +33,6 @@ #include #include #include -#include #include #include #include @@ -987,6 +986,10 @@ void signal_setup(void) { sigdelset(&set, SIGUSR2); pthread_sigmask(SIG_BLOCK, &set, NULL); + // SIGUSR1 is used to interrupt a thread if blocked in pselect + pthread_sigmask(SIG_SETMASK, NULL, &pselect_sigset); + sigdelset(&pselect_sigset, SIGUSR1); + // setting this to SIG_IGN would prevent signalling any threads. struct sigaction sa; memset(&sa, 0, sizeof(sa));