Ensure no pthread cancellation can occur during an ALSA-related call.

This commit is contained in:
Mike Brady
2018-12-03 12:44:57 +00:00
9 changed files with 434 additions and 284 deletions
+17 -10
View File
@@ -70,7 +70,7 @@ audio_output audio_alsa = {
.rate_info = &get_rate_information,
.mute = NULL, // a function will be provided if it can, and is allowed to, do hardware mute
.volume = NULL, // a function will be provided if it can do hardware volume
.parameters = &parameters};
.parameters = NULL}; // a function will be provided if it can do hardware volume
static pthread_mutex_t alsa_mutex = PTHREAD_MUTEX_INITIALIZER;
@@ -136,6 +136,8 @@ static void help(void) {
void set_alsa_out_dev(char *dev) { alsa_out_dev = dev; }
//assuming pthread cancellation is disabled
int open_mixer() {
int response = 0;
if (hardware_mixer) {
@@ -179,6 +181,7 @@ int open_mixer() {
return response;
}
//assuming pthread cancellation is disabled
void close_mixer() {
if (alsa_mix_handle) {
snd_mixer_close(alsa_mix_handle);
@@ -186,6 +189,7 @@ void close_mixer() {
}
}
//assuming pthread cancellation is disabled
void do_snd_mixer_selem_set_playback_dB_all(snd_mixer_elem_t *mix_elem, double vol) {
if (snd_mixer_selem_set_playback_dB_all(mix_elem, vol, 0) != 0) {
debug(1, "Can't set playback volume accurately to %f dB.", vol);
@@ -196,6 +200,7 @@ void do_snd_mixer_selem_set_playback_dB_all(snd_mixer_elem_t *mix_elem, double v
}
static int init(int argc, char **argv) {
// for debugging
snd_output_stdio_attach(&output, stdout, 0);
// debug(2,"audio_alsa init called.");
@@ -406,7 +411,7 @@ static int init(int argc, char **argv) {
warn("Invalid audio argument: \"%s\" -- ignored", argv[optind]);
}
debug(1, "alsa output device name is \"%s\".", alsa_out_dev);
debug(1, "alsa: output device name is \"%s\".", alsa_out_dev);
if (hardware_mixer) {
int oldState;
@@ -459,8 +464,8 @@ static int init(int argc, char **argv) {
snd_ctl_elem_id_set_name(elem_id, alsa_mix_ctrl);
if (snd_ctl_get_dB_range(ctl, elem_id, &alsa_mix_mindb, &alsa_mix_maxdb) == 0) {
debug(1, "Volume control \"%s\" has dB volume from %f to %f.", alsa_mix_ctrl,
(1.0 * alsa_mix_mindb) / 100.0, (1.0 * alsa_mix_maxdb) / 100.0);
debug(1, "alsa: hardware mixer \"%s\" selected, with dB volume from %f to %f.",
alsa_mix_ctrl, (1.0 * alsa_mix_mindb) / 100.0, (1.0 * alsa_mix_maxdb) / 100.0);
has_softvol = 1;
audio_alsa.volume =
&volume; // insert the volume function now we know it can do dB stuff
@@ -499,7 +504,7 @@ static int init(int argc, char **argv) {
pthread_cleanup_pop(0);
pthread_setcancelstate(oldState, NULL);
} else {
// debug(1, "Has no mixer and thus no hardware mute.");
debug(1, "alsa: no hardware mixer selected.");
}
alsa_mix_handle = NULL;
@@ -511,6 +516,7 @@ static void deinit(void) {
stop();
}
//assuming pthread cancellation is disabled
int actual_open_alsa_device(void) {
// the alsa mutex is already acquired when this is called
const snd_pcm_uframes_t minimal_buffer_headroom =
@@ -906,6 +912,7 @@ static void start(int i_sample_rate, int i_sample_format) {
measurement_data_is_valid = 0;
}
//assuming pthread cancellation is disabled
int my_snd_pcm_delay(snd_pcm_t *pcm, snd_pcm_sframes_t *delayp) {
int ret;
snd_pcm_status_t *alsa_snd_pcm_status;
@@ -998,7 +1005,7 @@ int delay(long *the_delay) {
}
}
}
debug_mutex_unlock(&alsa_mutex, 3);
debug_mutex_unlock(&alsa_mutex, 0);
pthread_cleanup_pop(0);
// here, occasionally pretend there's a problem with pcm_get_delay()
// if ((random() % 100000) < 3) // keep it pretty rare
@@ -1041,7 +1048,7 @@ static int play(void *buf, int samples) {
pthread_cleanup_pop(0); // release the mutex
}
if (ret == 0) {
pthread_cleanup_debug_mutex_lock(&alsa_mutex, 10000, 1);
pthread_cleanup_debug_mutex_lock(&alsa_mutex, 10000, 0);
// snd_pcm_sframes_t current_delay = 0;
int err, err2;
if (snd_pcm_state(alsa_handle) == SND_PCM_STATE_XRUN) {
@@ -1061,6 +1068,7 @@ static int play(void *buf, int samples) {
if (samples == 0)
debug(1, "empty buffer being passed to pcm_writei -- skipping it");
if ((samples != 0) && (buf != NULL)) {
debug(3, "write %d frames.", samples);
err = alsa_pcm_write(alsa_handle, buf, samples);
if (err < 0) {
frame_index = 0;
@@ -1125,7 +1133,7 @@ static int play(void *buf, int samples) {
frame_index = 0;
measurement_data_is_valid = 0;
}
debug_mutex_unlock(&alsa_mutex, 3);
debug_mutex_unlock(&alsa_mutex, 0);
pthread_cleanup_pop(0); // release the mutex
}
pthread_setcancelstate(oldState, NULL);
@@ -1292,7 +1300,6 @@ void do_mute(int mute_state_requested) {
close_mixer();
}
}
mute_request_pending = 0;
mute_request_pending = 0;
pthread_setcancelstate(oldState, NULL);
}
+15 -3
View File
@@ -94,9 +94,21 @@ volatile int debuglev = 0;
sigset_t pselect_sigset;
int usleep_uncancellable(useconds_t usec) {
int response;
int oldState;
pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, &oldState);
response = usleep_uncancellable(usec);
pthread_setcancelstate(oldState, NULL);
return response;
}
static uint16_t UDPPortIndex = 0;
void resetFreeUDPPort() { UDPPortIndex = 0; }
void resetFreeUDPPort() {
debug(3, "Resetting UDP Port Suggestion to %u", config.udp_port_base);
UDPPortIndex = 0;
}
uint16_t nextFreeUDPPort() {
if (UDPPortIndex == 0)
@@ -1130,10 +1142,10 @@ int sps_pthread_mutex_timedlock(pthread_mutex_t *mutex, useconds_t dally_time,
// this is not pthread_cancellation safe because is contains a cancellation point
int oldState;
pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, &oldState);
useconds_t time_to_wait = dally_time;
int time_to_wait = dally_time;
int r = pthread_mutex_trylock(mutex);
while ((r == EBUSY) && (time_to_wait > 0)) {
useconds_t st = time_to_wait;
int st = time_to_wait;
if (st > 1000)
st = 1000;
sps_nanosleep(0, st * 1000); // this contains a cancellation point
+2
View File
@@ -331,4 +331,6 @@ char *get_version_string(); // mallocs a string space -- remember to free it aft
void sps_nanosleep(const time_t sec,
const long nanosec); // waits for this time, even through interruptions
int usleep_uncancellable(useconds_t usec);
#endif // _COMMON_H
+1 -1
View File
@@ -2,7 +2,7 @@
# Process this file with autoconf to produce a configure script.
AC_PREREQ([2.50])
AC_INIT([shairport-sync], [3.3d29], [mikebrady@eircom.net])
AC_INIT([shairport-sync], [3.3d32], [mikebrady@eircom.net])
AM_INIT_AUTOMAKE
AC_CONFIG_SRCDIR([shairport.c])
AC_CONFIG_HEADERS([config.h])
+34 -33
View File
@@ -442,7 +442,7 @@ void player_put_packet(seq_t seqno, uint32_t actual_timestamp, uint8_t *data, in
debug_mutex_unlock(&conn->flush_mutex, 3);
}
debug_mutex_lock(&conn->ab_mutex, 30000, 1);
debug_mutex_lock(&conn->ab_mutex, 30000, 0);
conn->packet_count++;
conn->packet_count_since_flush++;
conn->time_of_last_audio_packet = get_absolute_time_in_fp();
@@ -614,7 +614,7 @@ void player_put_packet(seq_t seqno, uint32_t actual_timestamp, uint8_t *data, in
}
}
}
debug_mutex_unlock(&conn->ab_mutex, 3);
debug_mutex_unlock(&conn->ab_mutex, 0);
}
int32_t rand_in_range(int32_t exclusive_range_limit) {
@@ -758,7 +758,7 @@ static inline void process_sample(int32_t sample, char **outp, enum sps_format_t
void buffer_get_frame_cleanup_handler(void *arg) {
rtsp_conn_info *conn = (rtsp_conn_info *)arg;
debug_mutex_unlock(&conn->ab_mutex, 3);
debug_mutex_unlock(&conn->ab_mutex, 0);
}
// get the next frame, when available. return 0 if underrun/stream reset.
@@ -769,7 +769,7 @@ static abuf_t *buffer_get_frame(rtsp_conn_info *conn) {
abuf_t *curframe = NULL;
int notified_buffer_empty = 0; // diagnostic only
debug_mutex_lock(&conn->ab_mutex, 30000, 1);
debug_mutex_lock(&conn->ab_mutex, 30000, 0);
int wait;
long dac_delay = 0; // long because alsa returns a long
@@ -779,27 +779,8 @@ static abuf_t *buffer_get_frame(rtsp_conn_info *conn) {
do {
// get the time
local_time_now = get_absolute_time_in_fp(); // type okay
debug(3, "buffer_get_frame is iterating");
// debug(3, "buffer_get_frame is iterating");
// if config.timeout (default 120) seconds have elapsed since the last audio packet was
// received, then we should stop.
// config.timeout of zero means don't check..., but iTunes may be confused by a long gap
// followed by a resumption...
if ((conn->time_of_last_audio_packet != 0) && (conn->stop == 0) &&
(config.dont_check_timeout == 0)) {
uint64_t ct = config.timeout; // go from int to 64-bit int
// if (conn->packet_count>500) { //for testing -- about 4 seconds of play first
if ((local_time_now > conn->time_of_last_audio_packet) &&
(local_time_now - conn->time_of_last_audio_packet >= ct << 32)) {
debug(1, "As Yeats almost said, \"Too long a silence / can make a stone of the heart\" "
"from RTSP conversation %d.",
conn->connection_number);
conn->stop = 1;
pthread_cancel(conn->thread);
// pthread_kill(conn->thread, SIGUSR1);
}
}
int rco = get_requested_connection_state_to_output();
if (conn->connection_state_to_output != rco) {
@@ -815,12 +796,12 @@ static abuf_t *buffer_get_frame(rtsp_conn_info *conn) {
if (config.output->is_running)
if (config.output->is_running() != 0) { // if the back end isn't running for any reason
debug(3, "not running");
debug_mutex_lock(&conn->flush_mutex, 1000, 1);
debug_mutex_lock(&conn->flush_mutex, 1000, 0);
conn->flush_requested = 1;
debug_mutex_unlock(&conn->flush_mutex, 3);
debug_mutex_unlock(&conn->flush_mutex, 0);
}
debug_mutex_lock(&conn->flush_mutex, 1000, 1);
debug_mutex_lock(&conn->flush_mutex, 1000, 0);
if (conn->flush_requested == 1) {
if (config.output->flush)
config.output->flush(); // no cancellation points
@@ -830,7 +811,7 @@ static abuf_t *buffer_get_frame(rtsp_conn_info *conn) {
conn->time_since_play_started = 0;
conn->flush_requested = 0;
}
debug_mutex_unlock(&conn->flush_mutex, 3);
debug_mutex_unlock(&conn->flush_mutex, 0);
if (conn->ab_synced) {
curframe = conn->audio_buffer + BUFIDX(conn->ab_read);
@@ -1415,6 +1396,8 @@ void player_thread_initial_cleanup_handler(__attribute__((unused)) void *arg) {
void player_thread_cleanup_handler(void *arg) {
rtsp_conn_info *conn = (rtsp_conn_info *)arg;
int oldState;
pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, &oldState);
debug(3, "Connection %d: player thread main loop exit via player_thread_cleanup_handler.",
conn->connection_number);
@@ -1484,6 +1467,7 @@ void player_thread_cleanup_handler(void *arg) {
clear_reference_timestamp(conn);
conn->rtp_running = 0;
pthread_setcancelstate(oldState, NULL);
}
void *player_thread_func(void *arg) {
@@ -1744,7 +1728,9 @@ void *player_thread_func(void *arg) {
pthread_cleanup_push(player_thread_cleanup_handler, arg); // undo what's been done so far
// stop looking elsewhere for DACP stuff
// stop looking elsewhere for DACP stuff
int oldState;
pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, &oldState);
#ifdef CONFIG_DACP_CLIENT
// debug(1, "Set dacp server info");
@@ -1761,8 +1747,9 @@ void *player_thread_func(void *arg) {
// almost certainly, this has pthread cancellation points in it -- beware
conn->dapo_private_storage = mdns_dacp_monitor(conn->dacp_id);
#endif
pthread_setcancelstate(oldState, NULL);
// set the default volume to whaterver it was before, as stored in the config airplay_volume
// set the default volume to whatever it was before, as stored in the config airplay_volume
debug(2, "Set initial volume to %f.", config.airplay_volume);
player_volume(config.airplay_volume, conn); // will contain a cancellation point if asked to wait
@@ -2275,6 +2262,14 @@ void *player_thread_func(void *arg) {
inframe->given_timestamp = 0;
inframe->sequence_number = 0;
// update the watchdog
if ((config.dont_check_timeout == 0) && (config.timeout != 0)) {
uint64_t time_now = get_absolute_time_in_fp();
debug_mutex_lock(&conn->watchdog_mutex, 1000, 0);
conn->watchdog_bark_time = time_now;
debug_mutex_unlock(&conn->watchdog_mutex, 0);
}
// debug(1,"Sync error %lld frames. Amount to stuff %d." ,sync_error,amount_to_stuff);
// new stats calculation. We want a running average of sync error, drift, adjustment,
@@ -2783,8 +2778,14 @@ int player_stop(rtsp_conn_info *conn) {
debug(2, "player_thread cancel...");
pthread_cancel(*conn->player_thread);
debug(2, "player_thread join...");
pthread_join(*conn->player_thread, NULL);
debug(2, "player_thread joined.");
if (pthread_join(*conn->player_thread, NULL) == -1) {
char errorstring[1024];
strerror_r(errno, (char *)errorstring, sizeof(errorstring));
debug(1, "Connection %d: error %d joining player thread: \"%s\".", conn->connection_number,
errno, (char *)errorstring);
} else {
debug(2, "player_thread joined.");
}
free(conn->player_thread);
conn->player_thread = NULL;
#ifdef CONFIG_METADATA
@@ -2799,4 +2800,4 @@ int player_stop(rtsp_conn_info *conn) {
// debuglev = dl;
return -1;
}
}
}
+4 -1
View File
@@ -85,7 +85,9 @@ typedef struct {
int stop;
int running;
time_t playstart;
pthread_t thread, timer_requester, rtp_audio_thread, rtp_control_thread, rtp_timing_thread;
pthread_t thread, timer_requester, rtp_audio_thread, rtp_control_thread, rtp_timing_thread,
player_watchdog_thread;
int64_t watchdog_bark_time;
// pthread_t *ptp;
// buffers to delete on exit
@@ -213,6 +215,7 @@ typedef struct {
// request in flight at the same time
pthread_mutex_t reference_time_mutex;
pthread_mutex_t watchdog_mutex;
double local_to_remote_time_gradient; // if no drift, this would be exactly 1.0; likely it's
// slightly above or below.
+34 -38
View File
@@ -97,17 +97,8 @@ uint64_t local_to_remote_time_difference_now(rtsp_conn_info *conn) {
return conn->local_to_remote_time_difference + (uint64_t)(drift * (uint64_t)0x100000000);
}
void rtp_audio_receiver_cleanup_handler(void *arg) {
int oldState;
pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, &oldState);
debug(3, "Audio Receiver Cleanup.");
rtsp_conn_info *conn = (rtsp_conn_info *)arg;
debug(3, "shutdown audio socket.");
shutdown(conn->audio_socket, SHUT_RDWR);
debug(3, "close audio socket.");
close(conn->audio_socket);
debug(3, "Connection %d: Audio Receiver Cleanup.", conn->connection_number);
pthread_setcancelstate(oldState, NULL);
void rtp_audio_receiver_cleanup_handler(__attribute__((unused)) void *arg) {
debug(3, "Audio Receiver Cleanup Done.");
}
void *rtp_audio_receiver(void *arg) {
@@ -241,17 +232,8 @@ void *rtp_audio_receiver(void *arg) {
pthread_exit(NULL);
}
void rtp_control_handler_cleanup_handler(void *arg) {
int oldState;
pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, &oldState);
debug(3, "Control Receiver Cleanup.");
rtsp_conn_info *conn = (rtsp_conn_info *)arg;
debug(3, "shutdown control socket.");
shutdown(conn->control_socket, SHUT_RDWR);
debug(3, "close control socket.");
close(conn->control_socket);
debug(3, "Connection %d: Control Receiver Cleanup.", conn->connection_number);
pthread_setcancelstate(oldState, NULL);
void rtp_control_handler_cleanup_handler(__attribute__((unused)) void *arg) {
debug(3, "Control Receiver Cleanup Done.");
}
void *rtp_control_receiver(void *arg) {
@@ -571,17 +553,15 @@ void *rtp_timing_sender(void *arg) {
}
void rtp_timing_receiver_cleanup_handler(void *arg) {
int oldState;
pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, &oldState);
debug(3, "Timing Receiver Cleanup.");
rtsp_conn_info *conn = (rtsp_conn_info *)arg;
debug(3, "Cancel Timing Requester.");
pthread_cancel(conn->timer_requester);
int oldState;
pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, &oldState);
debug(3, "Join Timing Requester.");
pthread_join(conn->timer_requester, NULL);
debug(3, "shutdown timing socket.");
shutdown(conn->timing_socket, SHUT_RDWR);
debug(3, "close timing socket.");
close(conn->timing_socket);
debug(3, "Connection %d: Timing Receiver Cleanup.", conn->connection_number);
debug(3, "Timing Receiver Cleanup Successful.");
pthread_setcancelstate(oldState, NULL);
}
@@ -873,6 +853,17 @@ static uint16_t bind_port(int ip_family, const char *self_ip_address, uint32_t s
int local_socket = socket(ip_family, SOCK_DGRAM, IPPROTO_UDP);
if (local_socket == -1)
die("Could not allocate a socket.");
/*
int val = 1;
ret = setsockopt(local_socket, SOL_SOCKET, SO_REUSEADDR, &val, sizeof(val));
if (ret < 0) {
char errorstring[1024];
strerror_r(errno, (char *)errorstring, sizeof(errorstring));
debug(1, "Error %d: \"%s\". Couldn't set SO_REUSEADDR");
}
*/
SOCKADDR myaddr;
int tryCount = 0;
uint16_t desired_port;
@@ -905,9 +896,14 @@ static uint16_t bind_port(int ip_family, const char *self_ip_address, uint32_t s
if (ret < 0) {
close(local_socket);
die("error: could not bind a UDP port! Check the udp_port_range is large enough -- it must be "
char errorstring[1024];
strerror_r(errno, (char *)errorstring, sizeof(errorstring));
die("error %d: \"%s\". Could not bind a UDP port! Check the udp_port_range is large enough -- "
"it must be "
"at least 3, and 10 or more is suggested -- or "
"check for restrictive firewall settings or a bad router!");
"check for restrictive firewall settings or a bad router! UDP base is %u, range is %u and "
"current suggestion is %u.",
errno, errorstring, config.udp_port_base, config.udp_port_range, desired_port);
}
uint16_t sport;
@@ -1055,7 +1051,7 @@ void rtp_setup(SOCKADDR *local, SOCKADDR *remote, uint16_t cport, uint16_t tport
void get_reference_timestamp_stuff(uint32_t *timestamp, uint64_t *timestamp_time,
uint64_t *remote_timestamp_time, rtsp_conn_info *conn) {
// types okay
debug_mutex_lock(&conn->reference_time_mutex, 1000, 1);
debug_mutex_lock(&conn->reference_time_mutex, 1000, 0);
*timestamp = conn->reference_timestamp;
*remote_timestamp_time = conn->remote_reference_timestamp_time;
*timestamp_time =
@@ -1063,7 +1059,7 @@ void get_reference_timestamp_stuff(uint32_t *timestamp, uint64_t *timestamp_time
// if ((*timestamp == 0) && (*timestamp_time == 0)) {
// debug(1,"Reference timestamp is invalid.");
//}
debug_mutex_unlock(&conn->reference_time_mutex, 3);
debug_mutex_unlock(&conn->reference_time_mutex, 0);
}
void clear_reference_timestamp(rtsp_conn_info *conn) {
@@ -1102,7 +1098,7 @@ int sanitised_source_rate_information(uint32_t *frames, uint64_t *time, rtsp_con
double calculated_frame_rate = ((1.0 * local_frames) / local_time) * one_fp;
if (((calculated_frame_rate / conn->input_rate) > 1.002) ||
((calculated_frame_rate / conn->input_rate) < 0.998)) {
debug(1, "input frame rate out of bounds at %.2f fps.", calculated_frame_rate);
debug(2, "input frame rate out of bounds at %.2f fps.", calculated_frame_rate);
result = 1;
} else {
*frames = local_frames;
@@ -1118,7 +1114,7 @@ int sanitised_source_rate_information(uint32_t *frames, uint64_t *time, rtsp_con
// the reference timestamps are denominated in terms of the input rate
int frame_to_local_time(uint32_t timestamp, uint64_t *time, rtsp_conn_info *conn) {
debug_mutex_lock(&conn->reference_time_mutex, 1000, 1);
debug_mutex_lock(&conn->reference_time_mutex, 1000, 0);
int result = 0;
uint64_t time_difference;
uint32_t frame_difference;
@@ -1150,12 +1146,12 @@ int frame_to_local_time(uint32_t timestamp, uint64_t *time, rtsp_conn_info *conn
// on the specified fps.
}
*time = remote_time_of_timestamp - local_to_remote_time_difference_now(conn);
debug_mutex_unlock(&conn->reference_time_mutex, 3);
debug_mutex_unlock(&conn->reference_time_mutex, 0);
return result;
}
int local_time_to_frame(uint64_t time, uint32_t *frame, rtsp_conn_info *conn) {
debug_mutex_lock(&conn->reference_time_mutex, 1000, 1);
debug_mutex_lock(&conn->reference_time_mutex, 1000, 0);
int result = 0;
uint64_t time_difference;
@@ -1186,7 +1182,7 @@ int local_time_to_frame(uint64_t time, uint32_t *frame, rtsp_conn_info *conn) {
// debug(1,"Frame interval is %" PRId64 " frames.",-frame_interval);
*frame = (conn->reference_timestamp - frame_interval);
}
debug_mutex_unlock(&conn->reference_time_mutex, 3);
debug_mutex_unlock(&conn->reference_time_mutex, 0);
return result;
}
+319 -191
View File
@@ -170,6 +170,7 @@ int send_ssnc_metadata(uint32_t code, char *data, uint32_t length, int block) {
}
void pc_queue_cleanup_handler(void *arg) {
// debug(1, "pc_queue_cleanup_handler called.");
pc_queue *the_queue = (pc_queue *)arg;
int rc = pthread_mutex_unlock(&the_queue->pc_queue_lock);
if (rc)
@@ -253,6 +254,48 @@ int pc_queue_get_item(pc_queue *the_queue, void *the_stuff) {
#endif
int have_player(rtsp_conn_info *conn) {
int response = 0;
debug_mutex_lock(&playing_conn_lock, 1000000, 3);
if (playing_conn == conn) // this connection definitely has the play lock
response = 1;
debug_mutex_unlock(&playing_conn_lock, 3);
return response;
}
void player_watchdog_thread_cleanup_handler(void *arg) {
rtsp_conn_info *conn = (rtsp_conn_info *)arg;
debug(2, "Connection %d: Watchdog Exit.", conn->connection_number);
}
void *player_watchdog_thread_code(void *arg) {
pthread_cleanup_push(player_watchdog_thread_cleanup_handler, arg);
rtsp_conn_info *conn = (rtsp_conn_info *)arg;
do {
usleep(2000000); // check every two seconds
debug(3, "Connection %d: Check the thread is doing something...", conn->connection_number);
if ((config.dont_check_timeout == 0) && (config.timeout != 0)) {
debug_mutex_lock(&conn->watchdog_mutex, 1000, 0);
uint64_t last_watchdog_bark_time = conn->watchdog_bark_time;
debug_mutex_unlock(&conn->watchdog_mutex, 0);
if (last_watchdog_bark_time != 0) {
uint64_t time_since_last_bark = (get_absolute_time_in_fp() - last_watchdog_bark_time) >> 32;
uint64_t ct = config.timeout; // go from int to 64-bit int
if (time_since_last_bark >= ct) {
debug(1, "Connection %d: As Yeats almost said, \"Too long a silence / can make a stone "
"of the heart\".",
conn->connection_number);
conn->stop = 1;
pthread_cancel(conn->thread);
}
}
}
} while (1);
pthread_cleanup_pop(0); // should never happen
pthread_exit(NULL);
}
void ask_other_rtsp_conversation_threads_to_stop(pthread_t except_this_thread);
void rtsp_request_shutdown_stream(void) {
@@ -425,13 +468,9 @@ static void debug_print_msg_content(int level, rtsp_message *msg) {
void msg_free(rtsp_message *msg) {
if (msg) {
int rc = pthread_mutex_lock(&reference_counter_lock);
if (rc)
debug(1, "Error %d locking reference counter lock during msg_free()", rc);
debug_mutex_lock(&reference_counter_lock, 1000, 3);
msg->referenceCount--;
rc = pthread_mutex_unlock(&reference_counter_lock);
if (rc)
debug(1, "Error %d unlocking reference counter lock during msg_free()", rc);
debug_mutex_unlock(&reference_counter_lock, 3);
if (msg->referenceCount == 0) {
unsigned int i;
for (i = 0; i < msg->nheaders; i++) {
@@ -517,18 +556,21 @@ enum rtsp_read_request_response rtsp_read_request(rtsp_conn_info *conn, rtsp_mes
int msg_size = -1;
while (msg_size < 0) {
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);
/*
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;
goto shutdown;
}
nread = read(conn->fd, buf + inbuf, buflen - inbuf);
if (nread == 0) {
@@ -541,10 +583,12 @@ enum rtsp_read_request_response rtsp_read_request(rtsp_conn_info *conn, rtsp_mes
if (nread < 0) {
if (errno == EINTR)
continue;
char errorstring[1024];
strerror_r(errno, (char *)errorstring, sizeof(errorstring));
debug(1, "Connection %d: rtsp_read_request_response_read_error %d: \"%s\".",
conn->connection_number, errno, (char *)errorstring);
if (errno != ECONNRESET) {
char errorstring[1024];
strerror_r(errno, (char *)errorstring, sizeof(errorstring));
debug(1, "Connection %d: rtsp_read_request_response_read_error %d: \"%s\".",
conn->connection_number, errno, (char *)errorstring);
}
reply = rtsp_read_request_response_read_error;
goto shutdown;
}
@@ -600,6 +644,8 @@ enum rtsp_read_request_response rtsp_read_request(rtsp_conn_info *conn, rtsp_mes
warning_message_sent = 1;
}
}
/*
fd_set readfds;
FD_ZERO(&readfds);
FD_SET(conn->fd, &readfds);
@@ -607,6 +653,8 @@ enum rtsp_read_request_response rtsp_read_request(rtsp_conn_info *conn, rtsp_mes
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;
@@ -696,50 +744,62 @@ int msg_write_response(int fd, rtsp_message *resp) {
debug(1, "Attempted to write overlong RTSP packet 3");
return -3;
}
if (write(fd, pkt, p - pkt) != p - pkt) {
debug(1, "Error writing an RTSP packet -- requested bytes not fully written.");
ssize_t reply = write(fd, pkt, p - pkt);
if (reply == -1) {
char errorstring[1024];
strerror_r(errno, (char *)errorstring, sizeof(errorstring));
debug(1, "msg_write_response error %d: \"%s\".", errno, (char *)errorstring);
return -4;
}
if (reply != p - pkt) {
debug(1, "msg_write_response error -- requested bytes not fully written.");
return -5;
}
return 0;
}
void handle_record(rtsp_conn_info *conn, rtsp_message *req, rtsp_message *resp) {
debug(2, "Connection %d: RECORD", conn->connection_number);
if (have_player(conn)) {
if (conn->player_thread)
warn("Connection %d: RECORD: Duplicate RECORD message -- ignored", conn->connection_number);
else
player_play(conn); // the thread better be 0
if (conn->player_thread)
warn("Duplicate RECORD message -- ignored");
else
player_play(conn); // the thread better be 0
resp->respcode = 200;
// I think this is for telling the client what the absolute minimum latency
// actually is,
// and when the client specifies a latency, it should be added to this figure.
resp->respcode = 200;
// I think this is for telling the client what the absolute minimum latency
// actually is,
// and when the client specifies a latency, it should be added to this figure.
// Thus, [the old version of] AirPlay's latency figure of 77175, when added to 11025 gives you
// exactly 88200
// and iTunes' latency figure of 88553, when added to 11025 gives you 99578,
// pretty close to the 99400 we guessed.
// Thus, [the old version of] AirPlay's latency figure of 77175, when added to 11025 gives you
// exactly 88200
// and iTunes' latency figure of 88553, when added to 11025 gives you 99578,
// pretty close to the 99400 we guessed.
msg_add_header(resp, "Audio-Latency", "11025");
msg_add_header(resp, "Audio-Latency", "11025");
char *p;
uint32_t rtptime = 0;
char *hdr = msg_get_header(req, "RTP-Info");
char *p;
uint32_t rtptime = 0;
char *hdr = msg_get_header(req, "RTP-Info");
if (hdr) {
// debug(1,"FLUSH message received: \"%s\".",hdr);
// get the rtp timestamp
p = strstr(hdr, "rtptime=");
if (p) {
p = strchr(p, '=');
if (hdr) {
// debug(1,"FLUSH message received: \"%s\".",hdr);
// get the rtp timestamp
p = strstr(hdr, "rtptime=");
if (p) {
rtptime = uatoi(p + 1); // unsigned integer -- up to 2^32-1
// rtptime--;
// debug(1,"RTSP Flush Requested by handle_record: %u.",rtptime);
player_flush(rtptime, conn);
p = strchr(p, '=');
if (p) {
rtptime = uatoi(p + 1); // unsigned integer -- up to 2^32-1
// rtptime--;
// debug(1,"RTSP Flush Requested by handle_record: %u.",rtptime);
player_flush(rtptime, conn);
}
}
}
} else {
warn("Connection %d RECORD received without having the player (no ANNOUNCE?)",
conn->connection_number);
resp->respcode = 451;
}
}
@@ -755,152 +815,170 @@ void handle_options(rtsp_conn_info *conn, __attribute__((unused)) rtsp_message *
void handle_teardown(rtsp_conn_info *conn, __attribute__((unused)) rtsp_message *req,
rtsp_message *resp) {
debug(2, "Connection %d: TEARDOWN", conn->connection_number);
// if (!rtsp_playing())
// debug(1, "This RTSP connection thread (%d) doesn't think it's playing, but "
// "it's sending a response to teardown anyway",conn->connection_number);
resp->respcode = 200;
msg_add_header(resp, "Connection", "close");
debug(3,
if (have_player(conn)) {
resp->respcode = 200;
msg_add_header(resp, "Connection", "close");
debug(
3,
"TEARDOWN: synchronously terminating the player thread of RTSP conversation thread %d (2).",
conn->connection_number);
player_stop(conn);
debug(3, "TEARDOWN: successful termination of playing thread of RTSP conversation thread %d.",
conn->connection_number);
player_stop(conn);
debug(3, "TEARDOWN: successful termination of playing thread of RTSP conversation thread %d.",
conn->connection_number);
} else {
warn("Connection %d TEARDOWN received without having the player (no ANNOUNCE?)",
conn->connection_number);
resp->respcode = 451;
}
}
void handle_flush(rtsp_conn_info *conn, rtsp_message *req, rtsp_message *resp) {
debug(3, "Connection %d: FLUSH", conn->connection_number);
// if (!rtsp_playing())
// debug(1, "This RTSP conversation thread (%d) doesn't think it's playing, but "
// "it's sending a response to flush anyway",conn->connection_number);
char *p = NULL;
uint32_t rtptime = 0;
char *hdr = msg_get_header(req, "RTP-Info");
if (have_player(conn)) {
char *p = NULL;
uint32_t rtptime = 0;
char *hdr = msg_get_header(req, "RTP-Info");
if (hdr) {
// debug(1,"FLUSH message received: \"%s\".",hdr);
// get the rtp timestamp
p = strstr(hdr, "rtptime=");
if (p) {
p = strchr(p, '=');
if (p)
rtptime = uatoi(p + 1); // unsigned integer -- up to 2^32-1
if (hdr) {
// debug(1,"FLUSH message received: \"%s\".",hdr);
// get the rtp timestamp
p = strstr(hdr, "rtptime=");
if (p) {
p = strchr(p, '=');
if (p)
rtptime = uatoi(p + 1); // unsigned integer -- up to 2^32-1
}
}
}
// debug(1,"RTSP Flush Requested: %u.",rtptime);
#ifdef CONFIG_METADATA
if (p)
send_metadata('ssnc', 'flsr', p + 1, strlen(p + 1), req, 1);
else
send_metadata('ssnc', 'flsr', NULL, 0, NULL, 0);
if (p)
send_metadata('ssnc', 'flsr', p + 1, strlen(p + 1), req, 1);
else
send_metadata('ssnc', 'flsr', NULL, 0, NULL, 0);
#endif
player_flush(rtptime, conn); // will not crash even it there is no player thread.
resp->respcode = 200;
player_flush(rtptime, conn); // will not crash even it there is no player thread.
resp->respcode = 200;
} else {
warn("Connection %d FLUSH received without having the player (no ANNOUNCE?)",
conn->connection_number);
resp->respcode = 451;
}
}
void handle_setup(rtsp_conn_info *conn, rtsp_message *req, rtsp_message *resp) {
debug(3, "Connection %d: SETUP", conn->connection_number);
uint16_t cport, tport;
char *ar = msg_get_header(req, "Active-Remote");
if (ar) {
debug(2, "Connection %d: SETUP -- Active-Remote string seen: \"%s\".", conn->connection_number,
ar);
// get the active remote
char *p;
conn->dacp_active_remote = strtoul(ar, &p, 10);
resp->respcode = 451; // invalid arguments -- expect them
if (have_player(conn)) {
uint16_t cport, tport;
char *ar = msg_get_header(req, "Active-Remote");
if (ar) {
debug(2, "Connection %d: SETUP -- Active-Remote string seen: \"%s\".",
conn->connection_number, ar);
// get the active remote
char *p;
conn->dacp_active_remote = strtoul(ar, &p, 10);
#ifdef CONFIG_METADATA
send_metadata('ssnc', 'acre', ar, strlen(ar), req, 1);
send_metadata('ssnc', 'acre', ar, strlen(ar), req, 1);
#endif
} else {
debug(2, "Connection %d: SETUP -- Note: no Active-Remote information the SETUP Record.",
conn->connection_number);
conn->dacp_active_remote = 0;
}
ar = msg_get_header(req, "DACP-ID");
if (ar) {
debug(2, "Connection %d: SETUP -- DACP-ID string seen: \"%s\".", conn->connection_number, ar);
if (conn->dacp_id) // this is in case SETUP was previously called
free(conn->dacp_id);
conn->dacp_id = strdup(ar);
#ifdef CONFIG_METADATA
send_metadata('ssnc', 'daid', ar, strlen(ar), req, 1);
#endif
} else {
debug(2, "Connection %d: SETUP doesn't include DACP-ID string information.",
conn->connection_number);
if (conn->dacp_id) // this is in case SETUP was previously called
free(conn->dacp_id);
conn->dacp_id = NULL;
}
char *hdr = msg_get_header(req, "Transport");
if (!hdr) {
debug(1, "Connection %d: SETUP doesn't contain a Transport header.", conn->connection_number);
goto error;
}
char *p;
p = strstr(hdr, "control_port=");
if (!p) {
debug(1, "Connection %d: SETUP doesn't specify a control_port.", conn->connection_number);
goto error;
}
p = strchr(p, '=') + 1;
cport = atoi(p);
p = strstr(hdr, "timing_port=");
if (!p) {
debug(1, "Connection %d: SETUP doesn't specify a timing_port.", conn->connection_number);
goto error;
}
p = strchr(p, '=') + 1;
tport = atoi(p);
if (conn->rtp_running) {
if ((conn->remote_control_port != cport) || (conn->remote_timing_port != tport)) {
warn("Connection %d: Duplicate SETUP message with different control (old %u, new %u) or "
"timing (old %u, new "
"%u) ports! This is probably fatal!",
conn->connection_number, conn->remote_control_port, cport, conn->remote_timing_port,
tport);
} else {
warn("Connection %d: Duplicate SETUP message with the same control (%u) and timing (%u) "
"ports. This is "
"probably not fatal.",
conn->connection_number, conn->remote_control_port, conn->remote_timing_port);
debug(2, "Connection %d: SETUP -- Note: no Active-Remote information the SETUP Record.",
conn->connection_number);
conn->dacp_active_remote = 0;
}
ar = msg_get_header(req, "DACP-ID");
if (ar) {
debug(2, "Connection %d: SETUP -- DACP-ID string seen: \"%s\".", conn->connection_number, ar);
if (conn->dacp_id) // this is in case SETUP was previously called
free(conn->dacp_id);
conn->dacp_id = strdup(ar);
#ifdef CONFIG_METADATA
send_metadata('ssnc', 'daid', ar, strlen(ar), req, 1);
#endif
} else {
debug(2, "Connection %d: SETUP doesn't include DACP-ID string information.",
conn->connection_number);
if (conn->dacp_id) // this is in case SETUP was previously called
free(conn->dacp_id);
conn->dacp_id = NULL;
}
char *hdr = msg_get_header(req, "Transport");
if (hdr) {
char *p;
p = strstr(hdr, "control_port=");
if (p) {
p = strchr(p, '=') + 1;
cport = atoi(p);
p = strstr(hdr, "timing_port=");
if (p) {
p = strchr(p, '=') + 1;
tport = atoi(p);
if (conn->rtp_running) {
if ((conn->remote_control_port != cport) || (conn->remote_timing_port != tport)) {
warn("Connection %d: Duplicate SETUP message with different control (old %u, new %u) "
"or "
"timing (old %u, new "
"%u) ports! This is probably fatal!",
conn->connection_number, conn->remote_control_port, cport,
conn->remote_timing_port, tport);
} else {
warn("Connection %d: Duplicate SETUP message with the same control (%u) and timing "
"(%u) "
"ports. This is "
"probably not fatal.",
conn->connection_number, conn->remote_control_port, conn->remote_timing_port);
}
} else {
rtp_setup(&conn->local, &conn->remote, cport, tport, conn);
}
if (conn->local_audio_port != 0) {
char resphdr[256] = "";
snprintf(resphdr, sizeof(resphdr),
"RTP/AVP/"
"UDP;unicast;interleaved=0-1;mode=record;control_port=%d;"
"timing_port=%d;server_"
"port=%d",
conn->local_control_port, conn->local_timing_port, conn->local_audio_port);
msg_add_header(resp, "Transport", resphdr);
msg_add_header(resp, "Session", "1");
resp->respcode = 200; // it all worked out okay
debug(1, "Connection %d: SETUP with UDP ports Control: %d, Timing: %d and Audio: %d.",
conn->connection_number, conn->local_control_port, conn->local_timing_port,
conn->local_audio_port);
} else {
debug(1, "Connection %d: SETUP seems to specify a null audio port.",
conn->connection_number);
}
} else {
debug(1, "Connection %d: SETUP doesn't specify a timing_port.", conn->connection_number);
}
} else {
debug(1, "Connection %d: SETUP doesn't specify a control_port.", conn->connection_number);
}
} else {
debug(1, "Connection %d: SETUP doesn't contain a Transport header.", conn->connection_number);
}
if (resp->respcode != 200) {
debug(1, "Connection %d: SETUP error -- releasing the player lock.", conn->connection_number);
debug_mutex_lock(&playing_conn_lock, 1000000, 3);
if (playing_conn == conn) // if we have the player
playing_conn = NULL; // let it go
debug_mutex_unlock(&playing_conn_lock, 3);
}
} else {
rtp_setup(&conn->local, &conn->remote, cport, tport, conn);
warn("Connection %d SETUP received without having the player (no ANNOUNCE?)",
conn->connection_number);
}
if (conn->local_audio_port == 0) {
debug(1, "Connection %d: SETUP seems to specify a null audio port.", conn->connection_number);
goto error;
}
char resphdr[256] = "";
snprintf(resphdr, sizeof(resphdr), "RTP/AVP/"
"UDP;unicast;interleaved=0-1;mode=record;control_port=%d;"
"timing_port=%d;server_"
"port=%d",
conn->local_control_port, conn->local_timing_port, conn->local_audio_port);
msg_add_header(resp, "Transport", resphdr);
msg_add_header(resp, "Session", "1");
resp->respcode = 200;
return;
error:
warn("Connection %d: SETUP -- Error in setup request -- unlocking play lock.",
conn->connection_number);
debug_mutex_lock(&playing_conn_lock, 1000000, 3);
playing_conn = NULL; // this connection definitely doesn't have the play lock
debug_mutex_unlock(&playing_conn_lock, 3);
resp->respcode = 451; // invalid arguments
}
/*
@@ -1522,7 +1600,7 @@ static void handle_set_parameter(rtsp_conn_info *conn, rtsp_message *req, rtsp_m
}
static void handle_announce(rtsp_conn_info *conn, rtsp_message *req, rtsp_message *resp) {
debug(2, "Connection %d: ANNOUNCE", conn->connection_number);
debug(3, "Connection %d: ANNOUNCE", conn->connection_number);
int have_the_player = 0;
int should_wait = 0; // this will be true if you're trying to break in to the current session
@@ -1539,8 +1617,8 @@ static void handle_announce(rtsp_conn_info *conn, rtsp_message *req, rtsp_messag
have_the_player = 1;
warn("Duplicate ANNOUNCE, by the look of it!");
} else if (playing_conn->stop) {
debug(2, "Connection %d: ANNOUNCE: already shutting down; waiting for it...",
playing_conn->connection_number);
debug(1, "Connection %d ANNOUNCE is waiting for connection %d to shut down.",
conn->connection_number, playing_conn->connection_number);
should_wait = 1;
} else if (config.allow_session_interruption == 1) {
debug(2, "Connection %d: ANNOUNCE: asking playing connection %d to shut down.",
@@ -1553,7 +1631,7 @@ static void handle_announce(rtsp_conn_info *conn, rtsp_message *req, rtsp_messag
debug_mutex_unlock(&playing_conn_lock, 3);
if (should_wait) {
useconds_t time_remaining = 3000000;
int time_remaining = 3000000; // must be signed, as it could go negative...
while ((time_remaining > 0) && (have_the_player == 0)) {
debug_mutex_lock(&playing_conn_lock, 1000000, 3); // get it
@@ -1752,7 +1830,8 @@ out:
debug(1, "Connection %d: Error in handling ANNOUNCE. Unlocking the play lock.",
conn->connection_number);
debug_mutex_lock(&playing_conn_lock, 1000000, 3); // get it
playing_conn = NULL;
if (playing_conn == conn) // if we managed to acquire it
playing_conn = NULL; // let it go
debug_mutex_unlock(&playing_conn_lock, 3);
}
}
@@ -1993,14 +2072,26 @@ authenticate:
void rtsp_conversation_thread_cleanup_function(void *arg) {
rtsp_conn_info *conn = (rtsp_conn_info *)arg;
// debug(1, "Connection %d: rtsp_conversation_thread_func_cleanup_function called.",
// conn->connection_number);
int oldState;
pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, &oldState);
debug(3, "Connection %d: rtsp_conversation_thread_func_cleanup_function called.",
conn->connection_number);
if (conn->player_thread)
player_stop(conn);
debug(3, "Closing timing, control and audio sockets...");
if (conn->control_socket)
close(conn->control_socket);
if (conn->timing_socket)
close(conn->timing_socket);
if (conn->audio_socket)
close(conn->audio_socket);
if (conn->fd > 0) {
// debug(1, "Connection %d: closing fd %d.",
// conn->connection_number,conn->fd);
debug(3, "Connection %d: closing fd %d.", conn->connection_number, conn->fd);
close(conn->fd);
debug(3, "Connection %d: closed fd %d.", conn->connection_number, conn->fd);
}
if (conn->auth_nonce) {
free(conn->auth_nonce);
@@ -2025,25 +2116,38 @@ void rtsp_conversation_thread_cleanup_function(void *arg) {
if (rc)
debug(1, "Connection %d: error %d destroying flush_mutex.", conn->connection_number, rc);
debug(2, "Cancel watchdog thread.");
pthread_cancel(conn->player_watchdog_thread);
debug(2, "Join watchdog thread.");
pthread_join(conn->player_watchdog_thread, NULL);
debug(2, "Delete watchdog mutex.");
pthread_mutex_destroy(&conn->watchdog_mutex);
debug(3, "Connection %d: Checking play lock.", conn->connection_number);
debug_mutex_lock(&playing_conn_lock, 1000000, 3); // get it
if (playing_conn == conn) {
debug(3, "Connection %d: Unlocking play lock.", conn->connection_number);
playing_conn = NULL;
if (playing_conn == conn) { // if it's ours
debug(2, "Connection %d: Unlocking play lock.", conn->connection_number);
playing_conn = NULL; // let it go
}
debug_mutex_unlock(&playing_conn_lock, 3);
debug(2, "Connection %d: RTSP thread terminated.", conn->connection_number);
conn->running = 0;
pthread_setcancelstate(oldState, NULL);
}
void msg_cleanup_function(void *arg) {
// debug(1, "msg_cleanup_function called.");
debug(3, "msg_cleanup_function called.");
msg_free((rtsp_message *)arg);
}
static void *rtsp_conversation_thread_func(void *pconn) {
rtsp_conn_info *conn = pconn;
// create the watchdog mutex and start the watchdog thread;
pthread_mutex_init(&conn->watchdog_mutex, NULL);
pthread_create(&conn->player_watchdog_thread, NULL, &player_watchdog_thread_code, (void *)conn);
int rc = pthread_mutex_init(&conn->flush_mutex, NULL);
if (rc)
die("Connection %d: error %d initialising flush_mutex.", conn->connection_number, rc);
@@ -2064,6 +2168,9 @@ static void *rtsp_conversation_thread_func(void *pconn) {
die("Connection %d: error %d initialising flow control condition variable.",
conn->connection_number, rc);
// nothing before this is cancellable
pthread_cleanup_push(rtsp_conversation_thread_cleanup_function, (void *)conn);
rtp_initialise(conn);
char *hdr = NULL;
@@ -2072,7 +2179,6 @@ static void *rtsp_conversation_thread_func(void *pconn) {
int rtsp_read_request_attempt_count = 1; // 1 means exit immediately
rtsp_message *req, *resp;
pthread_cleanup_push(rtsp_conversation_thread_cleanup_function, (void *)conn);
while (conn->stop == 0) {
int debug_level = 3; // for printing the request and response
reply = rtsp_read_request(conn, &req);
@@ -2131,6 +2237,8 @@ static void *rtsp_conversation_thread_func(void *pconn) {
}
debug(debug_level, "Connection %d: RTSP Response:", conn->connection_number);
debug_print_msg_headers(debug_level, resp);
/*
fd_set writefds;
FD_ZERO(&writefds);
FD_SET(conn->fd, &writefds);
@@ -2138,13 +2246,23 @@ static void *rtsp_conversation_thread_func(void *pconn) {
memory_barrier();
} while (conn->stop == 0 &&
pselect(conn->fd + 1, NULL, &writefds, NULL, NULL, &pselect_sigset) <= 0);
*/
if (conn->stop == 0) {
int err = msg_write_response(conn->fd, resp);
if (err) {
debug(1, "Connection %d: Unable to write an RTSP message response. Terminating the "
"connection.",
conn->connection_number);
struct linger so_linger;
so_linger.l_onoff = 1; // "true"
so_linger.l_linger = 0;
err = setsockopt(conn->fd, SOL_SOCKET, SO_LINGER, &so_linger, sizeof so_linger);
if (err)
debug(1, "Could not set the RTSP socket to abort due to a write error on closing.");
conn->stop = 1;
// if (debuglev >= 1)
// debuglev = 3; // see what happens next
}
}
pthread_cleanup_pop(1);
@@ -2159,7 +2277,14 @@ static void *rtsp_conversation_thread_func(void *pconn) {
rtsp_read_request_attempt_count--;
if (rtsp_read_request_attempt_count == 0) {
tstop = 1;
// if ((reply == rtsp_read_request_response_read_error) && (debuglev != 0))
if (reply == rtsp_read_request_response_read_error) {
struct linger so_linger;
so_linger.l_onoff = 1; // "true"
so_linger.l_linger = 0;
int err = setsockopt(conn->fd, SOL_SOCKET, SO_LINGER, &so_linger, sizeof so_linger);
if (err)
debug(1, "Could not set the RTSP socket to abort due to a read error on closing.");
}
// debuglev = 3; // see what happens next
} else {
if (reply == rtsp_read_request_response_channel_closed)
@@ -2208,15 +2333,20 @@ static const char *format_address(struct sockaddr *fsa) {
*/
void rtsp_listen_loop_cleanup_handler(__attribute__((unused)) void *arg) {
int oldState;
pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, &oldState);
debug(1, "rtsp_listen_loop_cleanup_handler called.");
cancel_all_RTSP_threads();
int *sockfd = (int *)arg;
mdns_unregister();
if (sockfd)
free(sockfd);
pthread_setcancelstate(oldState, NULL);
}
void rtsp_listen_loop(void) {
int oldState;
pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, &oldState);
struct addrinfo hints, *info, *p;
char portstr[6];
int *sockfd = NULL;
@@ -2315,9 +2445,7 @@ void rtsp_listen_loop(void) {
mdns_register();
// printf("Listening for connections.");
// shairport_startup_complete();
pthread_setcancelstate(oldState, NULL);
int acceptfd;
struct timeval tv;
pthread_cleanup_push(rtsp_listen_loop_cleanup_handler, (void *)sockfd);
@@ -2407,7 +2535,7 @@ 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);
// fcntl(conn->fd, F_SETFL, O_NONBLOCK);
ret = pthread_create(&conn->thread, NULL, rtsp_conversation_thread_func,
conn); // also acts as a memory barrier
@@ -2419,6 +2547,6 @@ void rtsp_listen_loop(void) {
}
} while (1);
pthread_cleanup_pop(0);
pthread_cleanup_pop(1); // should never happen
debug(1, "Oops -- fell out of the RTSP select loop");
}
+8 -7
View File
@@ -31,11 +31,11 @@
#include <libconfig.h>
#include <libgen.h>
#include <memory.h>
#include <net/if.h>
#include <popt.h>
#include <stdio.h>
#include <stdlib.h>
#include <sys/socket.h>
#include <net/if.h>
#include <sys/stat.h>
#include <sys/types.h>
#include <sys/wait.h>
@@ -1483,13 +1483,18 @@ int main(int argc, char **argv) {
char *version_dbs = get_version_string();
if (version_dbs) {
debug(1, "Version: \"%s\"", version_dbs);
debug(1, "software version: \"%s\"", version_dbs);
free(version_dbs);
} else {
debug(1, "Can't print the version information!");
debug(1, "can't print the version information!");
}
/* Print out options */
debug(1, "log verbosity is %d.", debuglev);
debug(1, "disable resend requests is %s.", config.disable_resend_requests ? "on" : "off");
debug(1, "diagnostic_drop_packet_fraction is %f. A value of 0.0 means no packets will be dropped "
"deliberately.",
config.diagnostic_drop_packet_fraction);
debug(1, "statistics_requester status is %d.", config.statistics_requested);
debug(1, "daemon status is %d.", config.daemonise);
debug(1, "deamon pid file path is \"%s\".", pid_file_proc());
@@ -1563,10 +1568,6 @@ int main(int argc, char **argv) {
#endif
debug(1, "loudness is %d.", config.loudness);
debug(1, "loudness reference level is %f", config.loudness_reference_volume_db);
debug(1, "disable resend requests is %s.", config.disable_resend_requests ? "on" : "off");
debug(1, "diagnostic_drop_packet_fraction is %f. A value of 0.0 means no packets will be dropped "
"deliberately.",
config.diagnostic_drop_packet_fraction);
uint8_t ap_md5[16];