Make the first output backend in the list of backends the default and make its name the default output_name. Clang-format everything

This commit is contained in:
Mike Brady
2020-09-03 09:53:13 +01:00
parent 2c2ac1be93
commit ca56287254
12 changed files with 1016 additions and 920 deletions
+1 -1
View File
@@ -89,7 +89,7 @@ static audio_output *outputs[] = {
#endif
NULL};
audio_output *audio_get_output(char *name) {
audio_output *audio_get_output(const char *name) {
audio_output **out;
// default to the first
+1 -1
View File
@@ -54,7 +54,7 @@ typedef struct {
} audio_output;
audio_output *audio_get_output(char *name);
audio_output *audio_get_output(const char *name);
void audio_ls_outputs(void);
void parse_general_audio_options(void);
+16 -16
View File
@@ -51,33 +51,34 @@ static void start(__attribute__((unused)) int sample_rate,
// "ENXIO O_NONBLOCK | O_WRONLY is set, the named file is a FIFO, and no process has the FIFO
// open for reading."
fd = try_to_open_pipe_for_writing(pipename);
// we check that it's not a "real" error. From the "man 2 open" page:
// "ENXIO O_NONBLOCK | O_WRONLY is set, the named file is a FIFO, and no process has the FIFO
// open for reading." Which is okay.
if ((fd == -1) && (errno != ENXIO)) {
char errorstring[1024];
strerror_r(errno, (char *)errorstring, sizeof(errorstring));
debug(1, "audio_pipe start -- error %d (\"%s\") opening pipe: \"%s\".", errno,
(char *)errorstring, pipename);
warn("can not open audio pipe -- error %d (\"%s\") opening pipe: \"%s\".", errno,
(char *)errorstring, pipename);
}
fd = try_to_open_pipe_for_writing(pipename);
// we check that it's not a "real" error. From the "man 2 open" page:
// "ENXIO O_NONBLOCK | O_WRONLY is set, the named file is a FIFO, and no process has the FIFO
// open for reading." Which is okay.
if ((fd == -1) && (errno != ENXIO)) {
char errorstring[1024];
strerror_r(errno, (char *)errorstring, sizeof(errorstring));
debug(1, "audio_pipe start -- error %d (\"%s\") opening pipe: \"%s\".", errno,
(char *)errorstring, pipename);
warn("can not open audio pipe -- error %d (\"%s\") opening pipe: \"%s\".", errno,
(char *)errorstring, pipename);
}
}
static int play(void *buf, int samples) {
// if the file is not open, try to open it.
char errorstring[1024];
if (fd == -1) {
fd = try_to_open_pipe_for_writing(pipename);
fd = try_to_open_pipe_for_writing(pipename);
}
// if it's got a reader, write to it.
if (fd > 0) {
//int rc = non_blocking_write(fd, buf, samples * 4);
// int rc = non_blocking_write(fd, buf, samples * 4);
int rc = write(fd, buf, samples * 4);
if ((rc < 0) && (errno != EPIPE)) {
strerror_r(errno, (char *)errorstring, 1024);
debug(1, "audio_pip play: error %d writing to the pipe named \"%s\": \"%s\".", errno, pipename, errorstring);
debug(1, "audio_pip play: error %d writing to the pipe named \"%s\": \"%s\".", errno,
pipename, errorstring);
}
}
return 0;
@@ -85,7 +86,6 @@ static int play(void *buf, int samples) {
static void stop(void) {
// Don't close the pipe just because a play session has stopped.
}
static int init(int argc, char **argv) {
+90 -88
View File
@@ -29,6 +29,8 @@
#include "common.h"
#include <assert.h>
#include <errno.h>
#include <fcntl.h>
#include <libgen.h>
#include <memory.h>
#include <poll.h>
#include <popt.h>
@@ -40,8 +42,6 @@
#include <sys/wait.h>
#include <time.h>
#include <unistd.h>
#include <fcntl.h>
#include <libgen.h>
#ifdef COMPILE_FOR_OSX
#include <CoreServices/CoreServices.h>
@@ -141,64 +141,66 @@ void do_sps_log_to_stdout(__attribute__((unused)) int prio, const char *t, ...)
fprintf(stdout, "%s\n", s);
}
int create_log_file(const char* path) {
int fd = -1;
if (path != NULL) {
char *dirc = strdup(path);
if (dirc) {
char *dname = dirname(dirc);
// create the directory, if necessary
int result = 0;
if (dname) {
char *pdir = realpath(dname, NULL); // will return a NULL if the directory doesn't exist
if (pdir == NULL) {
mode_t oldumask = umask(000);
result = mkpath(dname, 0777);
umask(oldumask);
} else {
free(pdir);
}
if ((result == 0) || (result == -EEXIST)) {
// now open the file
fd = open(path, O_WRONLY | O_NONBLOCK | O_CREAT | O_EXCL, S_IRUSR | S_IWUSR | S_IRGRP | S_IROTH);
if ((fd == -1) && (errno == EEXIST))
fd = open(path, O_WRONLY | O_APPEND | O_NONBLOCK);
int create_log_file(const char *path) {
int fd = -1;
if (path != NULL) {
char *dirc = strdup(path);
if (dirc) {
char *dname = dirname(dirc);
// create the directory, if necessary
int result = 0;
if (dname) {
char *pdir = realpath(dname, NULL); // will return a NULL if the directory doesn't exist
if (pdir == NULL) {
mode_t oldumask = umask(000);
result = mkpath(dname, 0777);
umask(oldumask);
} else {
free(pdir);
}
if ((result == 0) || (result == -EEXIST)) {
// now open the file
fd = open(path, O_WRONLY | O_NONBLOCK | O_CREAT | O_EXCL,
S_IRUSR | S_IWUSR | S_IRGRP | S_IROTH);
if ((fd == -1) && (errno == EEXIST))
fd = open(path, O_WRONLY | O_APPEND | O_NONBLOCK);
if (fd >= 0) {
// now we switch to blocking mode
int flags = fcntl(fd, F_GETFL);
if (flags == -1) {
// strerror_r(errno, (char *)errorstring, sizeof(errorstring));
// debug(1, "create_log_file -- error %d (\"%s\") getting flags of pipe: \"%s\".", errno,
// (char *)errorstring, pathname);
} else {
flags = fcntl(fd, F_SETFL,flags & ~O_NONBLOCK);
// if (flags == -1) {
// strerror_r(errno, (char *)errorstring, sizeof(errorstring));
// debug(1, "create_log_file -- error %d (\"%s\") unsetting NONBLOCK of pipe: \"%s\".", errno,
// (char *)errorstring, pathname);
}
}
}
}
free(dirc);
}
}
return fd;
if (fd >= 0) {
// now we switch to blocking mode
int flags = fcntl(fd, F_GETFL);
if (flags == -1) {
// strerror_r(errno, (char
//*)errorstring, sizeof(errorstring)); debug(1, "create_log_file -- error %d (\"%s\")
//getting flags of pipe: \"%s\".", errno, (char *)errorstring, pathname);
} else {
flags = fcntl(fd, F_SETFL, flags & ~O_NONBLOCK);
// if (flags == -1) {
// strerror_r(errno,
//(char *)errorstring, sizeof(errorstring)); debug(1, "create_log_file -- error %d
//(\"%s\") unsetting NONBLOCK of pipe: \"%s\".", errno, (char *)errorstring,
//pathname);
}
}
}
}
free(dirc);
}
}
return fd;
}
void do_sps_log_to_fd(__attribute__((unused)) int prio, const char *t, ...) {
char s[1024];
va_list args;
va_start(args, t);
vsnprintf(s, sizeof(s), t, args);
va_end(args);
if (config.log_fd == -1)
config.log_fd = create_log_file(config.log_file_path);
if (config.log_fd >= 0) {
dprintf(config.log_fd, "%s\n", s);
char s[1024];
va_list args;
va_start(args, t);
vsnprintf(s, sizeof(s), t, args);
va_end(args);
if (config.log_fd == -1)
config.log_fd = create_log_file(config.log_file_path);
if (config.log_fd >= 0) {
dprintf(config.log_fd, "%s\n", s);
} else if (errno != ENXIO) { // maybe there is a pipe there but not hooked up
fprintf(stderr, "%s\n", s);
fprintf(stderr, "%s\n", s);
}
}
@@ -207,9 +209,9 @@ void log_to_stdout() { sps_log = do_sps_log_to_stdout; }
void log_to_file() { sps_log = do_sps_log_to_fd; }
void log_to_syslog() {
#ifdef CONFIG_LIBDAEMON
sps_log = daemon_log;
sps_log = daemon_log;
#else
sps_log = syslog;
sps_log = syslog;
#endif
}
@@ -309,8 +311,8 @@ void _die(const char *filename, const int linenumber, const char *format, ...) {
1.0 * time_since_last_debug_message / 1000000000, filename,
linenumber, " *fatal error: ");
} else {
strncpy(b, "fatal error: ", sizeof(b));
s = b+strlen(b);
strncpy(b, "fatal error: ", sizeof(b));
s = b + strlen(b);
}
va_list args;
va_start(args, format);
@@ -339,8 +341,8 @@ void _warn(const char *filename, const int linenumber, const char *format, ...)
1.0 * time_since_last_debug_message / 1000000000, filename,
linenumber, " *warning: ");
} else {
strncpy(b, "warning: ", sizeof(b));
s = b+strlen(b);
strncpy(b, "warning: ", sizeof(b));
s = b + strlen(b);
}
va_list args;
va_start(args, format);
@@ -1134,12 +1136,12 @@ uint64_t get_absolute_time_in_ns() {
return time_now_ns;
}
int try_to_open_pipe_for_writing(const char* pathname) {
// tries to open the pipe in non-blocking mode first.
// if it succeeds, it sets it to blocking.
// if not, it returns -1.
int try_to_open_pipe_for_writing(const char *pathname) {
// tries to open the pipe in non-blocking mode first.
// if it succeeds, it sets it to blocking.
// if not, it returns -1.
int fdis = open(pathname, O_WRONLY | O_NONBLOCK); // open it in non blocking mode first
int fdis = open(pathname, O_WRONLY | O_NONBLOCK); // open it in non blocking mode first
// we check that it's not a "real" error. From the "man 2 open" page:
// "ENXIO O_NONBLOCK | O_WRONLY is set, the named file is a FIFO, and no process has the FIFO
@@ -1147,24 +1149,24 @@ int try_to_open_pipe_for_writing(const char* pathname) {
// This is checked by the caller.
if (fdis >= 0) {
// now we switch to blocking mode
int flags = fcntl(fdis, F_GETFL);
if (flags == -1) {
char errorstring[1024];
strerror_r(errno, (char *)errorstring, sizeof(errorstring));
debug(1, "try_to_open_pipe -- error %d (\"%s\") getting flags of pipe: \"%s\".", errno,
(char *)errorstring, pathname);
} else {
flags = fcntl(fdis, F_SETFL,flags & ~O_NONBLOCK);
if (flags == -1) {
char errorstring[1024];
strerror_r(errno, (char *)errorstring, sizeof(errorstring));
debug(1, "try_to_open_pipe -- error %d (\"%s\") unsetting NONBLOCK of pipe: \"%s\".", errno,
(char *)errorstring, pathname);
}
}
// now we switch to blocking mode
int flags = fcntl(fdis, F_GETFL);
if (flags == -1) {
char errorstring[1024];
strerror_r(errno, (char *)errorstring, sizeof(errorstring));
debug(1, "try_to_open_pipe -- error %d (\"%s\") getting flags of pipe: \"%s\".", errno,
(char *)errorstring, pathname);
} else {
flags = fcntl(fdis, F_SETFL, flags & ~O_NONBLOCK);
if (flags == -1) {
char errorstring[1024];
strerror_r(errno, (char *)errorstring, sizeof(errorstring));
debug(1, "try_to_open_pipe -- error %d (\"%s\") unsetting NONBLOCK of pipe: \"%s\".", errno,
(char *)errorstring, pathname);
}
}
}
return fdis;
return fdis;
}
/* from
@@ -1710,11 +1712,11 @@ int string_update_with_size(char **str, int *flag, char *s, size_t len) {
}
// from https://stackoverflow.com/questions/13663617/memdup-function-in-c, with thanks
void* memdup(const void* mem, size_t size) {
void* out = malloc(size);
void *memdup(const void *mem, size_t size) {
void *out = malloc(size);
if(out != NULL)
memcpy(out, mem, size);
if (out != NULL)
memcpy(out, mem, size);
return out;
return out;
}
+10 -11
View File
@@ -182,8 +182,8 @@ typedef struct {
char *pidfile;
#endif
int log_fd; // file descriptor of the file or pipe to log stuff to.
char *log_file_path; // path to file or pipe to log to, if any
int log_fd; // file descriptor of the file or pipe to log stuff to.
char *log_file_path; // path to file or pipe to log to, if any
int logOutputLevel; // log output level
int debugger_show_elapsed_time; // in the debug message, display the time since startup
int debugger_show_relative_time; // in the debug message, display the time since the last one
@@ -285,10 +285,10 @@ typedef struct {
int jack_soxr_resample_quality;
#endif
#endif
void *gradients; // a linked list of the clock gradients discovered for all DACP IDs
// can't use IP numbers as they might be given to different devices
// can't get hold of MAC addresses.
// can't define the nvll linked list struct here
void *gradients; // a linked list of the clock gradients discovered for all DACP IDs
// can't use IP numbers as they might be given to different devices
// can't get hold of MAC addresses.
// can't define the nvll linked list struct here
} shairport_cfg;
// accessors to config for multi-thread access
@@ -303,9 +303,7 @@ void memory_barrier();
void log_to_stderr(); // call this to direct logging to stderr;
void log_to_stdout(); // call this to direct logging to stdout;
void log_to_syslog(); // call this to direct logging to the system log;
void log_to_file(); // call this to direct logging to a file or (pre-existing) pipe;
void log_to_file(); // call this to direct logging to a file or (pre-existing) pipe;
// true if Shairport Sync is supposed to be sending output to the output device, false otherwise
@@ -313,7 +311,8 @@ int get_requested_connection_state_to_output();
void set_requested_connection_state_to_output(int v);
int try_to_open_pipe_for_writing(const char* pathname); // open it without blocking if it's not hooked up
int try_to_open_pipe_for_writing(
const char *pathname); // open it without blocking if it's not hooked up
/* from
* http://coding.debuntu.org/c-implementing-str_replace-replace-all-occurrences-substring#comment-722
@@ -446,6 +445,6 @@ void malloc_cleanup(void *arg);
int string_update_with_size(char **str, int *flag, char *s, size_t len);
// from https://stackoverflow.com/questions/13663617/memdup-function-in-c, with thanks
void* memdup(const void* mem, size_t size);
void *memdup(const void *mem, size_t size);
#endif // _COMMON_H
+25 -19
View File
@@ -136,18 +136,20 @@ void _metadata_hub_modify_prolog(const char *filename, const int linenumber) {
// debug(1, "locking metadata hub for writing");
if (pthread_rwlock_trywrlock(&metadata_hub_re_lock) != 0) {
if (last_metadata_hub_modify_prolog_file)
debug(2, "Metadata_hub write lock at \"%s:%d\" is already taken at \"%s:%d\" -- must wait.", filename, linenumber, last_metadata_hub_modify_prolog_file, last_metadata_hub_modify_prolog_line);
debug(2, "Metadata_hub write lock at \"%s:%d\" is already taken at \"%s:%d\" -- must wait.",
filename, linenumber, last_metadata_hub_modify_prolog_file,
last_metadata_hub_modify_prolog_line);
else
debug(2, "Metadata_hub write lock is already taken by unknown -- must wait.");
debug(2, "Metadata_hub write lock is already taken by unknown -- must wait.");
metadata_hub_re_lock_access_is_delayed = 0;
pthread_rwlock_wrlock(&metadata_hub_re_lock);
debug(2, "Okay -- acquired the metadata_hub write lock at \"%s:%d\".", filename, linenumber);
} else {
if (last_metadata_hub_modify_prolog_file) {
free(last_metadata_hub_modify_prolog_file);
}
last_metadata_hub_modify_prolog_file = strdup(filename);
last_metadata_hub_modify_prolog_line = linenumber;
if (last_metadata_hub_modify_prolog_file) {
free(last_metadata_hub_modify_prolog_file);
}
last_metadata_hub_modify_prolog_file = strdup(filename);
last_metadata_hub_modify_prolog_line = linenumber;
// debug(3, "Metadata_hub write lock acquired.");
}
metadata_hub_re_lock_access_is_delayed = 0;
@@ -160,13 +162,16 @@ void _metadata_hub_modify_epilog(int modified, const char *filename, const int l
run_metadata_watchers();
}
if (metadata_hub_re_lock_access_is_delayed) {
if (last_metadata_hub_modify_prolog_file) {
debug(1, "Metadata_hub write lock taken at \"%s:%d\" is freed at \"%s:%d\".", last_metadata_hub_modify_prolog_file, last_metadata_hub_modify_prolog_line, filename, linenumber);
free(last_metadata_hub_modify_prolog_file);
last_metadata_hub_modify_prolog_file = NULL;
} else {
debug(1, "Metadata_hub write lock taken at an unknown place is freed at \"%s:%d\".", filename, linenumber);
}
if (last_metadata_hub_modify_prolog_file) {
debug(1, "Metadata_hub write lock taken at \"%s:%d\" is freed at \"%s:%d\".",
last_metadata_hub_modify_prolog_file, last_metadata_hub_modify_prolog_line, filename,
linenumber);
free(last_metadata_hub_modify_prolog_file);
last_metadata_hub_modify_prolog_file = NULL;
} else {
debug(1, "Metadata_hub write lock taken at an unknown place is freed at \"%s:%d\".", filename,
linenumber);
}
}
pthread_rwlock_unlock(&metadata_hub_re_lock);
// debug(3, "Metadata_hub write lock unlocked.");
@@ -506,10 +511,11 @@ void metadata_hub_process_metadata(uint32_t type, uint32_t code, char *data, uin
debug(2, "MH Picture received, length %u bytes.", length);
char uri[2048];
if ((length > 16) && (strcmp(config.cover_art_cache_dir,"")!=0)) { // if it's okay to write the file
// make this uncancellable
int oldState;
pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, &oldState); // make this un-cancellable
if ((length > 16) &&
(strcmp(config.cover_art_cache_dir, "") != 0)) { // if it's okay to write the file
// make this uncancellable
int oldState;
pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, &oldState); // make this un-cancellable
char *pathname = metadata_write_image_file(data, length);
snprintf(uri, sizeof(uri), "file://%s", pathname);
free(pathname);
@@ -524,7 +530,7 @@ void metadata_hub_process_metadata(uint32_t type, uint32_t code, char *data, uin
changed = 1;
else
changed = 0;
// pthread_cleanup_pop(0); // don't remove the lock -- it'll have been done
// pthread_cleanup_pop(0); // don't remove the lock -- it'll have been done
break;
case 'clip':
cs = strndup(data, length);
+5 -2
View File
@@ -149,7 +149,9 @@ void metadata_hub_release_track_artwork(void);
// these functions lock and unlock the read-write mutex on the metadata hub and run the watchers
// afterwards
void _metadata_hub_modify_prolog(const char *filename, const int linenumber);
void _metadata_hub_modify_epilog(int modified, const char *filename, const int linenumber); // set to true if modifications occurred, 0 otherwise
void _metadata_hub_modify_epilog(
int modified, const char *filename,
const int linenumber); // set to true if modifications occurred, 0 otherwise
/*
// these are for safe reading
@@ -158,7 +160,8 @@ void _metadata_hub_read_epilog(const char *filename, const int linenumber);
*/
#define metadata_hub_modify_prolog(void) _metadata_hub_modify_prolog(__FILE__, __LINE__)
#define metadata_hub_modify_epilog(modified) _metadata_hub_modify_epilog(modified, __FILE__, __LINE__)
#define metadata_hub_modify_epilog(modified) \
_metadata_hub_modify_epilog(modified, __FILE__, __LINE__)
#define metadata_hub_read_prolog(void) _metadata_hub_read_prolog(__FILE__, __LINE__)
#define metadata_hub_read_epilog(void) _metadata_hub_modify_epilog(__FILE__, __LINE__)
+512 -482
View File
File diff suppressed because it is too large Load Diff
+3 -1
View File
@@ -25,7 +25,9 @@
#include "audio.h"
#define time_ping_history_power_of_two 7
#define time_ping_history (1 << time_ping_history_power_of_two) // 2^7 is 128. At 1 per three seconds, approximately six minutes of records
#define time_ping_history \
(1 << time_ping_history_power_of_two) // 2^7 is 128. At 1 per three seconds, approximately six
// minutes of records
typedef struct time_ping_record {
uint64_t dispersion;
+52 -47
View File
@@ -46,9 +46,9 @@
#include <unistd.h>
struct Nvll {
char* name;
double value;
struct Nvll *next;
char *name;
double value;
struct Nvll *next;
};
typedef struct Nvll nvll;
@@ -263,8 +263,8 @@ void *rtp_control_receiver(void *arg) {
obfp += 2;
};
*obfp = 0;
// get raw timestamp information
// I think that a good way to understand these timestamps is that
// (1) the rtlt below is the timestamp of the frame that should be playing at the
@@ -275,19 +275,19 @@ void *rtp_control_receiver(void *arg) {
// Thus, (3) the latency can be calculated by subtracting the second from the
// first.
// There must be more to it -- there something missing.
// In addition, it seems that if the value of the short represented by the second
// pair of bytes in the packet is 7
// then an extra time lag is expected to be added, presumably by
// the AirPort Express.
// Best guess is that this delay is 11,025 frames.
uint32_t rtlt = nctohl(&packet[4]); // raw timestamp less latency
uint32_t rt = nctohl(&packet[16]); // raw timestamp
uint32_t fl = nctohs(&packet[2]); //
debug(1,"Sync Packet of %d bytes received: \"%s\", flags: %d, timestamps %u and %u,
giving a latency of %d frames.",plen,obf,fl,rt,rtlt,rt-rtlt);
//debug(1,"Monotonic timestamps are: %" PRId64 " and %" PRId64 "
@@ -527,22 +527,25 @@ void rtp_timing_receiver_cleanup_handler(void *arg) {
// walk down the list of DACP / gradient pairs, if any
nvll *gradients = config.gradients;
if (conn->dacp_id)
while ((gradients) && (strcasecmp((const char *)&conn->client_ip_string,gradients->name) != 0))
gradients = gradients->next;
while ((gradients) && (strcasecmp((const char *)&conn->client_ip_string, gradients->name) != 0))
gradients = gradients->next;
// if gradients comes out of this non-null, it is pointing to the DACP and it's last-known gradient
// if gradients comes out of this non-null, it is pointing to the DACP and it's last-known
// gradient
if (gradients) {
gradients->value = conn->local_to_remote_time_gradient;
// debug(1,"Updating a drift of %.2f ppm for \"%s\".", (conn->local_to_remote_time_gradient - 1.0)*1000000, gradients->name);
gradients->value = conn->local_to_remote_time_gradient;
// debug(1,"Updating a drift of %.2f ppm for \"%s\".", (conn->local_to_remote_time_gradient
// - 1.0)*1000000, gradients->name);
} else {
nvll *new_entry = (nvll*)malloc(sizeof(nvll));
if (new_entry) {
new_entry->name = strdup((const char *)&conn->client_ip_string);
new_entry->value = conn->local_to_remote_time_gradient;
new_entry->next = config.gradients;
config.gradients = new_entry;
// debug(1,"Setting a new drift of %.2f ppm for \"%s\".", (conn->local_to_remote_time_gradient - 1.0)*1000000, new_entry->name);
}
nvll *new_entry = (nvll *)malloc(sizeof(nvll));
if (new_entry) {
new_entry->name = strdup((const char *)&conn->client_ip_string);
new_entry->value = conn->local_to_remote_time_gradient;
new_entry->next = config.gradients;
config.gradients = new_entry;
// debug(1,"Setting a new drift of %.2f ppm for \"%s\".", (conn->local_to_remote_time_gradient
// - 1.0)*1000000, new_entry->name);
}
}
debug(3, "Cancel Timing Requester.");
@@ -574,13 +577,14 @@ void *rtp_timing_receiver(void *arg) {
conn->local_to_remote_time_gradient = 1.0; // initial value.
// walk down the list of DACP / gradient pairs, if any
nvll *gradients = config.gradients;
while ((gradients) && (strcasecmp((const char *)&conn->client_ip_string,gradients->name) != 0))
gradients = gradients->next;
while ((gradients) && (strcasecmp((const char *)&conn->client_ip_string, gradients->name) != 0))
gradients = gradients->next;
// if gradients comes out of this non-null, it is pointing to the IP and it's last-known gradient
if (gradients) {
conn->local_to_remote_time_gradient = gradients->value;
// debug(1,"Using a stored drift of %.2f ppm for \"%s\".", (conn->local_to_remote_time_gradient - 1.0)*1000000, gradients->name);
conn->local_to_remote_time_gradient = gradients->value;
// debug(1,"Using a stored drift of %.2f ppm for \"%s\".", (conn->local_to_remote_time_gradient
// - 1.0)*1000000, gradients->name);
}
// calculate diffusion factor
@@ -591,13 +595,12 @@ void *rtp_timing_receiver(void *arg) {
// be the nth root of diffusion_expansion_factor
// where n is the number of elements in the array
const double diffusion_expansion_factor = 10;
double log_of_multiplier = log10(diffusion_expansion_factor)/time_ping_history;
double multiplier = pow(10,log_of_multiplier);
const double diffusion_expansion_factor = 10;
double log_of_multiplier = log10(diffusion_expansion_factor) / time_ping_history;
double multiplier = pow(10, log_of_multiplier);
uint64_t dispersion_factor = (uint64_t)(multiplier * 100);
// debug(1,"dispersion factor is %" PRIu64 ".", dispersion_factor);
// uint64_t first_local_to_remote_time_difference_time;
// uint64_t l2rtd = 0;
int sequence_number = 0;
@@ -621,7 +624,8 @@ void *rtp_timing_receiver(void *arg) {
if (packet[1] == 0xd3) { // timing reply
return_time = arrival_time - conn->departure_time;
debug(3,"clock synchronisation request: return time is %8.3f milliseconds.",0.000001*return_time);
debug(3, "clock synchronisation request: return time is %8.3f milliseconds.",
0.000001 * return_time);
if (return_time < 200000000) { // must be less than 0.2 seconds
// distant_receive_time =
@@ -671,11 +675,11 @@ void *rtp_timing_receiver(void *arg) {
// conn->time_pings[cc].dispersion * pow(2.14,
// 1.0/conn->time_ping_count);
if (conn->time_pings[cc].dispersion > UINT64_MAX / dispersion_factor)
debug(1,"dispersion factor is too large at %" PRIu64 ".");
debug(1, "dispersion factor is too large at %" PRIu64 ".");
else
conn->time_pings[cc].dispersion =
(conn->time_pings[cc].dispersion * dispersion_factor) /
100; // make the dispersions 'age' by this rational factor
conn->time_pings[cc].dispersion =
(conn->time_pings[cc].dispersion * dispersion_factor) /
100; // make the dispersions 'age' by this rational factor
}
// these are used for doing a least squares calculation to get the drift
conn->time_pings[0].local_time = arrival_time;
@@ -752,7 +756,8 @@ void *rtp_timing_receiver(void *arg) {
if ((conn->time_pings[cc].chosen) &&
(conn->time_pings[cc].sequence_number >
(settling_time / 3))) { // wait for a approximate settling time
// have to scale them down so that the sum, possibly over every term in the array, doesn't overflow
// have to scale them down so that the sum, possibly over
// every term in the array, doesn't overflow
y_bar += (conn->time_pings[cc].remote_time >> time_ping_history_power_of_two);
x_bar += (conn->time_pings[cc].local_time >> time_ping_history_power_of_two);
sample_count++;
@@ -762,8 +767,6 @@ void *rtp_timing_receiver(void *arg) {
y_bar = y_bar / sample_count;
x_bar = x_bar / sample_count;
int64_t xid, yid;
double mtl, mbl;
mtl = 0;
@@ -791,19 +794,21 @@ void *rtp_timing_receiver(void *arg) {
conn->local_to_remote_time_gradient = mtl / mbl;
else {
// conn->local_to_remote_time_gradient = 1.0;
debug(1,"mbl is zero. Drift remains at %.2f ppm.", (conn->local_to_remote_time_gradient - 1.0)*1000000);
debug(1, "mbl is zero. Drift remains at %.2f ppm.",
(conn->local_to_remote_time_gradient - 1.0) * 1000000);
}
// scale the numbers back up
uint64_t ybf = y_bar << time_ping_history_power_of_two;
uint64_t xbf = x_bar << time_ping_history_power_of_two;
// scale the numbers back up
uint64_t ybf = y_bar << time_ping_history_power_of_two;
uint64_t xbf = x_bar << time_ping_history_power_of_two;
conn->local_to_remote_time_difference =
ybf - xbf; // make this the new local-to-remote-time-difference
conn->local_to_remote_time_difference_measurement_time = xbf;
conn->local_to_remote_time_difference =
ybf - xbf; // make this the new local-to-remote-time-difference
conn->local_to_remote_time_difference_measurement_time = xbf;
} else {
debug(3,"not enough samples to estimate drift -- remaining at %.2f ppm.", (conn->local_to_remote_time_gradient - 1.0)*1000000);
debug(3, "not enough samples to estimate drift -- remaining at %.2f ppm.",
(conn->local_to_remote_time_gradient - 1.0) * 1000000);
// conn->local_to_remote_time_gradient = 1.0;
}
// debug(1,"local to remote time gradient is %12.2f ppm, based on %d
+188 -158
View File
@@ -115,7 +115,7 @@ typedef struct {
pthread_cond_t pc_queue_item_added_signal;
pthread_cond_t pc_queue_item_removed_signal;
char *name;
size_t item_size; // number of bytes in each item
size_t item_size; // number of bytes in each item
uint32_t count; // number of items in the queue
uint32_t capacity; // maximum number of items
uint32_t toq; // first item to take
@@ -152,11 +152,12 @@ typedef struct {
rtsp_message *carrier;
} metadata_package;
void pc_queue_init(pc_queue *the_queue, char *items, size_t item_size, uint32_t number_of_items, const char* name) {
if (name)
debug(2, "Creating metadata queue \"%s\".", name);
else
debug(1, "Creating an unnamed metadata queue.");
void pc_queue_init(pc_queue *the_queue, char *items, size_t item_size, uint32_t number_of_items,
const char *name) {
if (name)
debug(2, "Creating metadata queue \"%s\".", name);
else
debug(1, "Creating an unnamed metadata queue.");
pthread_mutex_init(&the_queue->pc_queue_lock, NULL);
pthread_cond_init(&the_queue->pc_queue_item_added_signal, NULL);
pthread_cond_init(&the_queue->pc_queue_item_removed_signal, NULL);
@@ -167,25 +168,25 @@ void pc_queue_init(pc_queue *the_queue, char *items, size_t item_size, uint32_t
the_queue->toq = 0;
the_queue->eoq = 0;
if (name == NULL)
the_queue->name = NULL;
the_queue->name = NULL;
else
the_queue->name = strdup(name);
the_queue->name = strdup(name);
}
void pc_queue_delete(pc_queue *the_queue) {
if (the_queue->name)
debug(2, "Deleting metadata queue \"%s\".", the_queue->name);
else
debug(1, "Deleting an unnamed metadata queue.");
if (the_queue->name != NULL)
free(the_queue->name);
// debug(2, "destroying pc_queue_item_removed_signal");
if (the_queue->name)
debug(2, "Deleting metadata queue \"%s\".", the_queue->name);
else
debug(1, "Deleting an unnamed metadata queue.");
if (the_queue->name != NULL)
free(the_queue->name);
// debug(2, "destroying pc_queue_item_removed_signal");
pthread_cond_destroy(&the_queue->pc_queue_item_removed_signal);
// debug(2, "destroying pc_queue_item_added_signal");
// debug(2, "destroying pc_queue_item_added_signal");
pthread_cond_destroy(&the_queue->pc_queue_item_added_signal);
// debug(2, "destroying pc_queue_lock");
// debug(2, "destroying pc_queue_lock");
pthread_mutex_destroy(&the_queue->pc_queue_lock);
// debug(2, "destroying signals and locks done");
// debug(2, "destroying signals and locks done");
}
int send_metadata(uint32_t type, uint32_t code, char *data, uint32_t length, rtsp_message *carrier,
@@ -204,7 +205,7 @@ void pc_queue_cleanup_handler(void *arg) {
}
int pc_queue_add_item(pc_queue *the_queue, const void *the_stuff, int block) {
int response = 0;
int response = 0;
int rc;
if (the_queue) {
if (block == 0) {
@@ -219,34 +220,39 @@ int pc_queue_add_item(pc_queue *the_queue, const void *the_stuff, int block) {
// leave this out if you want this to return if the queue is already full
// irrespective of the block flag.
/*
while (the_queue->count == the_queue->capacity) {
rc = pthread_cond_wait(&the_queue->pc_queue_item_removed_signal, &the_queue->pc_queue_lock);
if (rc)
debug(1, "Error waiting for item to be removed");
}
*/
while (the_queue->count == the_queue->capacity) {
rc = pthread_cond_wait(&the_queue->pc_queue_item_removed_signal,
&the_queue->pc_queue_lock); if (rc) debug(1, "Error waiting for item to be removed");
}
*/
if (the_queue->count < the_queue->capacity) {
uint32_t i = the_queue->eoq;
void *p = the_queue->items + the_queue->item_size * i;
// void * p = &the_queue->qbase + the_queue->item_size*the_queue->eoq;
memcpy(p, the_stuff, the_queue->item_size);
uint32_t i = the_queue->eoq;
void *p = the_queue->items + the_queue->item_size * i;
// void * p = &the_queue->qbase + the_queue->item_size*the_queue->eoq;
memcpy(p, the_stuff, the_queue->item_size);
// update the pointer
i++;
if (i == the_queue->capacity)
// fold pointer if necessary
i = 0;
the_queue->eoq = i;
the_queue->count++;
//debug(2,"metadata queue+ \"%s\" %d/%d.", the_queue->name, the_queue->count, the_queue->capacity);
if (the_queue->count == the_queue->capacity)
debug(3, "metadata queue \"%s\": is now full with %d items in it!", the_queue->name, the_queue->count);
rc = pthread_cond_signal(&the_queue->pc_queue_item_added_signal);
if (rc)
debug(1, "metadata queue \"%s\": error signalling after pc_queue_add_item", the_queue->name);
// update the pointer
i++;
if (i == the_queue->capacity)
// fold pointer if necessary
i = 0;
the_queue->eoq = i;
the_queue->count++;
// debug(2,"metadata queue+ \"%s\" %d/%d.", the_queue->name, the_queue->count,
// the_queue->capacity);
if (the_queue->count == the_queue->capacity)
debug(3, "metadata queue \"%s\": is now full with %d items in it!", the_queue->name,
the_queue->count);
rc = pthread_cond_signal(&the_queue->pc_queue_item_added_signal);
if (rc)
debug(1, "metadata queue \"%s\": error signalling after pc_queue_add_item",
the_queue->name);
} else {
response = EWOULDBLOCK; // a bit arbitrary, this.
debug(3,"metadata queue \"%s\": is already full with %d items in it. Not adding this item to the queue.", the_queue->name, the_queue->count);
response = EWOULDBLOCK; // a bit arbitrary, this.
debug(3,
"metadata queue \"%s\": is already full with %d items in it. Not adding this item to "
"the queue.",
the_queue->name, the_queue->count);
}
pthread_cleanup_pop(1); // unlock the queue lock.
} else {
@@ -279,7 +285,8 @@ int pc_queue_get_item(pc_queue *the_queue, void *the_stuff) {
i = 0;
the_queue->toq = i;
the_queue->count--;
debug(3,"metadata queue- \"%s\" %d/%d.", the_queue->name, the_queue->count, the_queue->capacity);
debug(3, "metadata queue- \"%s\" %d/%d.", the_queue->name, the_queue->count,
the_queue->capacity);
rc = pthread_cond_signal(&the_queue->pc_queue_item_removed_signal);
if (rc)
debug(1, "metadata queue \"%s\": error signalling after pc_queue_get_item", the_queue->name);
@@ -317,15 +324,17 @@ void *player_watchdog_thread_code(void *arg) {
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_ns() - last_watchdog_bark_time) / 1000000000;
uint64_t time_since_last_bark =
(get_absolute_time_in_ns() - last_watchdog_bark_time) / 1000000000;
uint64_t ct = config.timeout; // go from int to 64-bit int
if (time_since_last_bark >= ct) {
conn->watchdog_barks++;
if (conn->watchdog_barks == 1) {
// debuglev = 3; // tell us everything.
debug(1, "Connection %d: As Yeats almost said, \"Too long a silence / can make a stone "
"of the heart\".",
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);
@@ -446,7 +455,8 @@ void msg_retain(rtsp_message *msg) {
debug(1, "Error %d locking reference counter lock");
if (msg > (rtsp_message *)0x00010000) {
msg->referenceCount++;
debug(3,"msg_free increment reference counter message %d to %d.", msg->index_number, msg->referenceCount);
debug(3, "msg_free increment reference counter message %d to %d.", msg->index_number,
msg->referenceCount);
// debug(1,"msg_retain -- item %d reference count %d.", msg->index_number, msg->referenceCount);
rc = pthread_mutex_unlock(&reference_counter_lock);
if (rc)
@@ -462,7 +472,7 @@ rtsp_message *msg_init(void) {
memset(msg, 0, sizeof(rtsp_message));
msg->referenceCount = 1; // from now on, any access to this must be protected with the lock
msg->index_number = msg_indexes++;
debug(3,"msg_init message %d", msg->index_number);
debug(3, "msg_init message %d", msg->index_number);
} else {
die("msg_init -- can not allocate memory for rtsp_message %d.", msg_indexes);
}
@@ -526,8 +536,9 @@ void msg_free(rtsp_message **msgh) {
if (*msgh > (rtsp_message *)0x00010000) {
rtsp_message *msg = *msgh;
msg->referenceCount--;
if (msg->referenceCount)
debug(3,"msg_free decrement reference counter message %d to %d", msg->index_number, msg->referenceCount);
if (msg->referenceCount)
debug(3, "msg_free decrement reference counter message %d to %d", msg->index_number,
msg->referenceCount);
if (msg->referenceCount == 0) {
unsigned int i;
for (i = 0; i < msg->nheaders; i++) {
@@ -542,7 +553,7 @@ void msg_free(rtsp_message **msgh) {
index = 0x10000; // ensure it doesn't fold to zero.
*msgh =
(rtsp_message *)(index); // put a version of the index number of the freed message in here
debug(3,"msg_free freed message %d", msg->index_number);
debug(3, "msg_free freed message %d", msg->index_number);
free(msg);
} else {
// debug(1,"msg_free item %d -- decrement reference to
@@ -607,7 +618,7 @@ int msg_handle_line(rtsp_message **pmsg, char *line) {
}
fail:
debug(3,"msg_handle_line fail");
debug(3, "msg_handle_line fail");
msg_free(pmsg);
*pmsg = NULL;
return 0;
@@ -650,7 +661,8 @@ enum rtsp_read_request_response rtsp_read_request(rtsp_conn_info *conn, rtsp_mes
if (errno == EINTR)
continue;
if (errno == EAGAIN) {
debug(1, "Connection %d: getting Error 11 -- EAGAIN from a blocking read!", conn->connection_number);
debug(1, "Connection %d: getting Error 11 -- EAGAIN from a blocking read!",
conn->connection_number);
continue;
}
if (errno != ECONNRESET) {
@@ -663,15 +675,15 @@ enum rtsp_read_request_response rtsp_read_request(rtsp_conn_info *conn, rtsp_mes
goto shutdown;
}
/* // this outputs the message received
{
void *pt = malloc(nread+1);
memset(pt, 0, nread+1);
memcpy(pt, buf + inbuf, nread);
debug(1, "Incoming string on port: \"%s\"",pt);
free(pt);
}
*/
/* // this outputs the message received
{
void *pt = malloc(nread+1);
memset(pt, 0, nread+1);
memcpy(pt, buf + inbuf, nread);
debug(1, "Incoming string on port: \"%s\"",pt);
free(pt);
}
*/
inbuf += nread;
@@ -680,7 +692,8 @@ enum rtsp_read_request_response rtsp_read_request(rtsp_conn_info *conn, rtsp_mes
msg_size = msg_handle_line(the_packet, buf);
if (!(*the_packet)) {
debug(1,"Connection %d: rtsp_read_request can't find an RTSP header.", conn->connection_number);
debug(1, "Connection %d: rtsp_read_request can't find an RTSP header.",
conn->connection_number);
reply = rtsp_read_request_response_bad_packet;
goto shutdown;
}
@@ -887,9 +900,10 @@ void handle_options(rtsp_conn_info *conn, __attribute__((unused)) rtsp_message *
rtsp_message *resp) {
debug(3, "Connection %d: OPTIONS", conn->connection_number);
resp->respcode = 200;
msg_add_header(resp, "Public", "ANNOUNCE, SETUP, RECORD, "
"PAUSE, FLUSH, TEARDOWN, "
"OPTIONS, GET_PARAMETER, SET_PARAMETER");
msg_add_header(resp, "Public",
"ANNOUNCE, SETUP, RECORD, "
"PAUSE, FLUSH, TEARDOWN, "
"OPTIONS, GET_PARAMETER, SET_PARAMETER");
}
void handle_teardown(rtsp_conn_info *conn, __attribute__((unused)) rtsp_message *req,
@@ -929,17 +943,17 @@ void handle_flush(rtsp_conn_info *conn, rtsp_message *req, rtsp_message *resp) {
rtptime = uatoi(p + 1); // unsigned integer -- up to 2^32-1
}
}
// debug(1,"RTSP Flush Requested: %u.",rtptime);
// debug(1,"RTSP Flush Requested: %u.",rtptime);
// the following is now done better by the player_flush routine as a 'pfls'
/*
#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);
#endif
*/
// the following is now done better by the player_flush routine as a 'pfls'
/*
#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);
#endif
*/
player_flush(rtptime, conn); // will not crash even it there is no player thread.
resp->respcode = 200;
@@ -1035,8 +1049,9 @@ void handle_setup(rtsp_conn_info *conn, rtsp_message *req, rtsp_message *resp) {
msg_add_header(resp, "Session", "1");
resp->respcode = 200; // it all worked out okay
debug(1, "Connection %d: SETUP DACP-ID \"%s\" from %s to %s with UDP ports Control: "
"%d, Timing: %d and Audio: %d.",
debug(1,
"Connection %d: SETUP DACP-ID \"%s\" from %s to %s with UDP ports Control: "
"%d, Timing: %d and Audio: %d.",
conn->connection_number, conn->dacp_id, &conn->client_ip_string,
&conn->self_ip_string, conn->local_control_port, conn->local_timing_port,
conn->local_audio_port);
@@ -1265,8 +1280,6 @@ char *base64_encode_so(const unsigned char *data, size_t input_length, char *enc
static int fd = -1;
// static int dirty = 0;
pc_queue metadata_queue;
#define metadata_queue_size 500
metadata_package metadata_queue_items[metadata_queue_size];
@@ -1294,7 +1307,6 @@ pc_queue metadata_multicast_queue;
metadata_package metadata_multicast_queue_items[metadata_queue_size];
pthread_t metadata_multicast_thread;
void metadata_create_multicast_socket(void) {
if (config.metadata_enabled == 0)
return;
@@ -1332,7 +1344,6 @@ void metadata_delete_multicast_socket(void) {
free(metadata_sockmsg);
}
void metadata_open(void) {
if (config.metadata_enabled == 0)
return;
@@ -1354,7 +1365,8 @@ static void metadata_close(void) {
}
void metadata_multicast_process(uint32_t type, uint32_t code, char *data, uint32_t length) {
// debug(1, "Process multicast metadata with type %x, code %x and length %u.", type, code, length);
// debug(1, "Process multicast metadata with type %x, code %x and length %u.", type, code,
// length);
if (metadata_sock >= 0 && length < config.metadata_sockmsglength - 8) {
char *ptr = metadata_sockmsg;
uint32_t v;
@@ -1455,7 +1467,7 @@ void metadata_process(uint32_t type, uint32_t code, char *data, uint32_t length)
debug(1, "Error encoding base64 data.");
// debug(1,"Remaining count: %d ret: %d, outbuf_size:
// %d.",remaining_count,ret,outbuf_size);
//ret = non_blocking_write(fd, outbuf, outbuf_size);
// ret = non_blocking_write(fd, outbuf, outbuf_size);
ret = write(fd, outbuf, outbuf_size);
if (ret < 0) {
// debug(1,"metadata_process error %d exit 3",ret);
@@ -1473,7 +1485,7 @@ void metadata_process(uint32_t type, uint32_t code, char *data, uint32_t length)
}
}
snprintf(thestring, 1024, "</item>\n");
//ret = non_blocking_write(fd, thestring, strlen(thestring));
// ret = non_blocking_write(fd, thestring, strlen(thestring));
ret = write(fd, thestring, strlen(thestring));
if (ret < 0) {
// debug(1,"metadata_process error %d exit 5",ret);
@@ -1508,11 +1520,12 @@ void *metadata_thread_function(__attribute__((unused)) void *ignore) {
pc_queue_get_item(&metadata_queue, &pack);
pthread_cleanup_push(metadata_pack_cleanup_function, (void *)&pack);
if (config.metadata_enabled) {
if (pack.carrier) {
debug(3, " pipe: type %x, code %x, length %u, message %d.", pack.type, pack.code, pack.length, pack.carrier->index_number);
} else {
debug(3, " pipe: type %x, code %x, length %u.", pack.type, pack.code, pack.length);
}
if (pack.carrier) {
debug(3, " pipe: type %x, code %x, length %u, message %d.", pack.type, pack.code,
pack.length, pack.carrier->index_number);
} else {
debug(3, " pipe: type %x, code %x, length %u.", pack.type, pack.code, pack.length);
}
metadata_process(pack.type, pack.code, pack.data, pack.length);
debug(3, " pipe: done.");
}
@@ -1530,8 +1543,8 @@ void metadata_multicast_thread_cleanup_function(__attribute__((unused)) void *ar
void *metadata_multicast_thread_function(__attribute__((unused)) void *ignore) {
// create a pc_queue for passing information to a threaded metadata handler
pc_queue_init(&metadata_multicast_queue, (char *)&metadata_multicast_queue_items, sizeof(metadata_package),
metadata_multicast_queue_size, "multicast");
pc_queue_init(&metadata_multicast_queue, (char *)&metadata_multicast_queue_items,
sizeof(metadata_package), metadata_multicast_queue_size, "multicast");
metadata_create_multicast_socket();
metadata_package pack;
pthread_cleanup_push(metadata_multicast_thread_cleanup_function, NULL);
@@ -1539,13 +1552,20 @@ void *metadata_multicast_thread_function(__attribute__((unused)) void *ignore) {
pc_queue_get_item(&metadata_multicast_queue, &pack);
pthread_cleanup_push(metadata_pack_cleanup_function, (void *)&pack);
if (config.metadata_enabled) {
if (pack.carrier) {
debug(3, " multicast: type %x, code %x, length %u, message %d.", pack.type, pack.code, pack.length, pack.carrier->index_number);
} else {
debug(3, " multicast: type %x, code %x, length %u.", pack.type, pack.code, pack.length);
}
if (pack.carrier) {
debug(3,
" multicast: type "
"%x, code %x, length %u, message %d.",
pack.type, pack.code, pack.length, pack.carrier->index_number);
} else {
debug(3,
" multicast: type "
"%x, code %x, length %u.",
pack.type, pack.code, pack.length);
}
metadata_multicast_process(pack.type, pack.code, pack.data, pack.length);
debug(3, " multicast: done.");
debug(3,
" multicast: done.");
}
pthread_cleanup_pop(1);
}
@@ -1553,10 +1573,8 @@ void *metadata_multicast_thread_function(__attribute__((unused)) void *ignore) {
pthread_exit(NULL);
}
#ifdef CONFIG_METADATA_HUB
void metadata_hub_close(void) {
}
void metadata_hub_close(void) {}
void metadata_hub_thread_cleanup_function(__attribute__((unused)) void *arg) {
// debug(2, "metadata_hub_thread_cleanup_function called");
@@ -1582,14 +1600,14 @@ void *metadata_hub_thread_function(__attribute__((unused)) void *ignore) {
// we check that it's not a "real" error. From the "man 2 open" page:
// "ENXIO O_NONBLOCK | O_WRONLY is set, the named file is a FIFO, and no process has the FIFO
// open for reading." Which is okay.
if ((fd == -1) && (errno != ENXIO)) {
char errorstring[1024];
strerror_r(errno, (char *)errorstring, sizeof(errorstring));
debug(1, "metadata_hub_thread_function -- error %d (\"%s\") opening pipe: \"%s\".", errno,
(char *)errorstring, path);
warn("can not open metadata pipe -- error %d (\"%s\") opening pipe: \"%s\".", errno,
(char *)errorstring, path);
}
if ((fd == -1) && (errno != ENXIO)) {
char errorstring[1024];
strerror_r(errno, (char *)errorstring, sizeof(errorstring));
debug(1, "metadata_hub_thread_function -- error %d (\"%s\") opening pipe: \"%s\".", errno,
(char *)errorstring, path);
warn("can not open metadata pipe -- error %d (\"%s\") opening pipe: \"%s\".", errno,
(char *)errorstring, path);
}
free(path);
// create a pc_queue for passing information to a threaded metadata handler
@@ -1600,13 +1618,15 @@ void *metadata_hub_thread_function(__attribute__((unused)) void *ignore) {
while (1) {
pc_queue_get_item(&metadata_hub_queue, &pack);
pthread_cleanup_push(metadata_pack_cleanup_function, (void *)&pack);
if (pack.carrier) {
debug(3, " hub: type %x, code %x, length %u, message %d.", pack.type, pack.code, pack.length, pack.carrier->index_number);
} else {
debug(3, " hub: type %x, code %x, length %u.", pack.type, pack.code, pack.length);
}
metadata_hub_process_metadata(pack.type, pack.code, pack.data, pack.length);
debug(3, " hub: done.");
if (pack.carrier) {
debug(3, " hub: type %x, code %x, length %u, message %d.", pack.type,
pack.code, pack.length, pack.carrier->index_number);
} else {
debug(3, " hub: type %x, code %x, length %u.", pack.type, pack.code,
pack.length);
}
metadata_hub_process_metadata(pack.type, pack.code, pack.data, pack.length);
debug(3, " hub: done.");
pthread_cleanup_pop(1);
}
pthread_cleanup_pop(1); // will never happen
@@ -1615,8 +1635,7 @@ void *metadata_hub_thread_function(__attribute__((unused)) void *ignore) {
#endif
#ifdef CONFIG_MQTT
void metadata_mqtt_close(void) {
}
void metadata_mqtt_close(void) {}
void metadata_mqtt_thread_cleanup_function(__attribute__((unused)) void *arg) {
// debug(2, "metadata_mqtt_thread_cleanup_function called");
@@ -1634,15 +1653,19 @@ void *metadata_mqtt_thread_function(__attribute__((unused)) void *ignore) {
while (1) {
pc_queue_get_item(&metadata_mqtt_queue, &pack);
pthread_cleanup_push(metadata_pack_cleanup_function, (void *)&pack);
if (config.mqtt_enabled) {
if (pack.carrier) {
debug(3, " mqtt: type %x, code %x, length %u, message %d.", pack.type, pack.code, pack.length, pack.carrier->index_number);
} else {
debug(3, " mqtt: type %x, code %x, length %u.", pack.type, pack.code, pack.length);
}
mqtt_process_metadata(pack.type, pack.code, pack.data, pack.length);
debug(3, " mqtt: done.");
}
if (config.mqtt_enabled) {
if (pack.carrier) {
debug(3,
" mqtt: type %x, code %x, length %u, message "
"%d.",
pack.type, pack.code, pack.length, pack.carrier->index_number);
} else {
debug(3, " mqtt: type %x, code %x, length %u.",
pack.type, pack.code, pack.length);
}
mqtt_process_metadata(pack.type, pack.code, pack.data, pack.length);
debug(3, " mqtt: done.");
}
pthread_cleanup_pop(1);
}
@@ -1700,8 +1723,8 @@ void metadata_stop(void) {
}
}
int send_metadata_to_queue(pc_queue* queue, uint32_t type, uint32_t code, char *data, uint32_t length, rtsp_message *carrier,
int block) {
int send_metadata_to_queue(pc_queue *queue, uint32_t type, uint32_t code, char *data,
uint32_t length, rtsp_message *carrier, int block) {
// parameters: type, code, pointer to data or NULL, length of data or NULL,
// the rtsp_message or
@@ -1737,20 +1760,27 @@ int send_metadata_to_queue(pc_queue* queue, uint32_t type, uint32_t code, char *
if (pack.carrier) {
msg_retain(pack.carrier);
} else {
if (data)
pack.data = memdup(data,length); // only if it's not a null
if (data)
pack.data = memdup(data, length); // only if it's not a null
}
int rc = pc_queue_add_item(queue, &pack, block);
if (rc != 0) {
if (pack.carrier) {
if (rc == EWOULDBLOCK)
debug(2, "metadata queue \"%s\" full, dropping message item: type %x, code %x, data %x, length %u, message %d.", queue->name, pack.type, pack.code, pack.data, pack.length, pack.carrier->index_number);
if (rc == EWOULDBLOCK)
debug(2,
"metadata queue \"%s\" full, dropping message item: type %x, code %x, data %x, "
"length %u, message %d.",
queue->name, pack.type, pack.code, pack.data, pack.length,
pack.carrier->index_number);
msg_free(&pack.carrier);
} else {
if (rc == EWOULDBLOCK)
debug(2, "metadata queue \"%s\" full, dropping data item: type %x, code %x, data %x, length %u.", queue->name, pack.type, pack.code, pack.data, pack.length);
if (pack.data)
free(pack.data);
if (rc == EWOULDBLOCK)
debug(
2,
"metadata queue \"%s\" full, dropping data item: type %x, code %x, data %x, length %u.",
queue->name, pack.type, pack.code, pack.data, pack.length);
if (pack.data)
free(pack.data);
}
}
return rc;
@@ -1758,23 +1788,20 @@ int send_metadata_to_queue(pc_queue* queue, uint32_t type, uint32_t code, char *
int send_metadata(uint32_t type, uint32_t code, char *data, uint32_t length, rtsp_message *carrier,
int block) {
int rc;
rc = send_metadata_to_queue(&metadata_queue, type, code, data, length, carrier, block);
int rc;
rc = send_metadata_to_queue(&metadata_queue, type, code, data, length, carrier, block);
#ifdef CONFIG_METADATA_HUB
rc = send_metadata_to_queue(&metadata_hub_queue, type, code, data, length, carrier, block);
rc = send_metadata_to_queue(&metadata_hub_queue, type, code, data, length, carrier, block);
#endif
#ifdef CONFIG_MQTT
rc = send_metadata_to_queue(&metadata_mqtt_queue, type, code, data, length, carrier, block);
rc = send_metadata_to_queue(&metadata_mqtt_queue, type, code, data, length, carrier, block);
#endif
return rc;
return rc;
}
static void handle_set_parameter_metadata(__attribute__((unused)) rtsp_conn_info *conn,
rtsp_message *req,
__attribute__((unused)) rtsp_message *resp) {
@@ -1838,7 +1865,7 @@ static void handle_set_parameter(rtsp_conn_info *conn, rtsp_message *req, rtsp_m
char *ct = msg_get_header(req, "Content-Type");
if (ct) {
// debug(2, "SET_PARAMETER Content-Type:\"%s\".", ct);
// debug(2, "SET_PARAMETER Content-Type:\"%s\".", ct);
#ifdef CONFIG_METADATA
// It seems that the rtptime of the message is used as a kind of an ID that
@@ -2136,7 +2163,7 @@ static void handle_announce(rtsp_conn_info *conn, rtsp_message *req, rtsp_messag
unsigned int i = 0;
unsigned int max_param = sizeof(conn->stream.fmtp) / sizeof(conn->stream.fmtp[0]);
char* found;
char *found;
while ((found = strsep(&pfmtp, " \t")) != NULL && i < max_param) {
conn->stream.fmtp[i++] = atoi(found);
}
@@ -2617,8 +2644,9 @@ static void *rtsp_conversation_thread_func(void *pconn) {
if (strcmp(req->method, "OPTIONS") !=
0) // the options message is very common, so don't log it until level 3
debug_level = 2;
debug(debug_level, "Connection %d: Received an RTSP Packet of type \"%s\":",
conn->connection_number, req->method),
debug(debug_level,
"Connection %d: Received an RTSP Packet of type \"%s\":", conn->connection_number,
req->method),
debug_print_msg_headers(debug_level, req);
apple_challenge(conn->fd, req, resp);
@@ -2667,8 +2695,9 @@ static void *rtsp_conversation_thread_func(void *pconn) {
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.",
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"
@@ -2721,10 +2750,11 @@ static void *rtsp_conversation_thread_func(void *pconn) {
if (reply == -1) {
char errorstring[1024];
strerror_r(errno, (char *)errorstring, sizeof(errorstring));
debug(1, "rtsp_read_request_response_bad_packet write response error %d: \"%s\".", errno, (char *)errorstring);
debug(1, "rtsp_read_request_response_bad_packet write response error %d: \"%s\".", errno,
(char *)errorstring);
} else if (reply != (ssize_t)strlen(response_text)) {
debug(1, "rtsp_read_request_response_bad_packet write %d bytes requested but %d written.", strlen(response_text),
reply);
debug(1, "rtsp_read_request_response_bad_packet write %d bytes requested but %d written.",
strlen(response_text), reply);
}
} else {
debug(1, "Connection %d: rtsp_read_request error %d, packet ignored.",
@@ -2937,7 +2967,7 @@ void rtsp_listen_loop(void) {
debug(2, "Connection %d: new connection from %s:%u to self at %s:%u.",
conn->connection_number, remote_ip4, rport, ip4, tport);
}
#ifdef AF_INET6
#ifdef AF_INET6
if (local_info->SAFAMILY == AF_INET6) {
// IPv6:
@@ -2954,7 +2984,7 @@ void rtsp_listen_loop(void) {
debug(2, "Connection %d: new connection from [%s]:%u to self at [%s]:%u.",
conn->connection_number, remote_ip6, rport, ip6, tport);
}
#endif
#endif
} else {
debug(1, "Error figuring out Shairport Sync's own IP number.");
+113 -94
View File
@@ -61,6 +61,7 @@
#endif
#include "activity_monitor.h"
#include "audio.h"
#include "common.h"
#include "rtp.h"
#include "rtsp.h"
@@ -117,6 +118,8 @@ int daemonisewithout = 0;
char configuration_file_path[4096 + 1];
char actual_configuration_file_path[4096 + 1];
char first_backend_name[256];
void print_version(void) {
char *version_string = get_version_string();
if (version_string) {
@@ -126,7 +129,6 @@ void print_version(void) {
debug(1, "Can't print version string!");
}
}
#ifdef CONFIG_SOXR
pthread_t soxr_time_check_thread;
void *soxr_time_check(__attribute__((unused)) void *arg) {
@@ -419,15 +421,15 @@ int parse_options(int argc, char **argv) {
// for unexpected circumstances
#ifdef CONFIG_METADATA
/* Get the metadata setting. */
config.metadata_enabled = 1; // if metadata support is included, then enable it by default
config.get_coverart = 1; // if metadata support is included, then enable it by default
/* Get the metadata setting. */
config.metadata_enabled = 1; // if metadata support is included, then enable it by default
config.get_coverart = 1; // if metadata support is included, then enable it by default
#endif
#ifdef CONFIG_CONVOLUTION
config.convolution_max_length = 8192;
config.convolution_max_length = 8192;
#endif
config.loudness_reference_volume_db = -20;
config.loudness_reference_volume_db = -20;
#ifdef CONFIG_METADATA_HUB
config.cover_art_cache_dir = "/tmp/shairport-sync/.cache/coverart";
@@ -682,9 +684,9 @@ int parse_options(int argc, char **argv) {
} else if (strcasecmp(str, "stderr") == 0) {
log_to_stderr();
} else {
config.log_file_path = (char *)str;
config.log_fd = -1;
log_to_file();
config.log_file_path = (char *)str;
config.log_fd = -1;
log_to_file();
}
}
/* Get the ignore_volume_control setting. */
@@ -1313,113 +1315,114 @@ const char *pid_file_proc(void) {
void exit_function() {
if (emergency_exit == 0) {
// the following is to ensure that if libdaemon has been included
// that most of this code will be skipped when the parent process is exiting
// exec
if (emergency_exit == 0) {
// the following is to ensure that if libdaemon has been included
// that most of this code will be skipped when the parent process is exiting
// exec
#ifdef CONFIG_LIBDAEMON
if ((this_is_the_daemon_process) || (config.daemonise == 0)) { // if this is the daemon process that is exiting or it's not actually deamonised at all
if ((this_is_the_daemon_process) ||
(config.daemonise == 0)) { // if this is the daemon process that is exiting or it's not
// actually deamonised at all
#endif
debug(2, "exit function called...");
/*
Actually, there is no terminate_mqtt() function.
#ifdef CONFIG_MQTT
if (config.mqtt_enabled) {
terminate_mqtt();
}
#endif
*/
debug(2, "exit function called...");
/*
Actually, there is no terminate_mqtt() function.
#ifdef CONFIG_MQTT
if (config.mqtt_enabled) {
terminate_mqtt();
}
#endif
*/
#if defined(CONFIG_DBUS_INTERFACE) || defined(CONFIG_MPRIS_INTERFACE)
/*
Actually, there is no stop_mpris_service() function.
#ifdef CONFIG_MPRIS_INTERFACE
stop_mpris_service();
#endif
*/
/*
Actually, there is no stop_mpris_service() function.
#ifdef CONFIG_MPRIS_INTERFACE
stop_mpris_service();
#endif
*/
#ifdef CONFIG_DBUS_INTERFACE
stop_dbus_service();
stop_dbus_service();
#endif
if (g_main_loop) {
debug(2, "Stopping DBUS Loop Thread");
g_main_loop_quit(g_main_loop);
pthread_join(dbus_thread, NULL);
}
if (g_main_loop) {
debug(2, "Stopping DBUS Loop Thread");
g_main_loop_quit(g_main_loop);
pthread_join(dbus_thread, NULL);
}
#endif
#ifdef CONFIG_DACP_CLIENT
debug(2, "Stopping DACP Monitor");
dacp_monitor_stop();
debug(2, "Stopping DACP Monitor");
dacp_monitor_stop();
#endif
#ifdef CONFIG_METADATA_HUB
debug(2, "Stopping metadata hub");
metadata_hub_stop();
debug(2, "Stopping metadata hub");
metadata_hub_stop();
#endif
#ifdef CONFIG_METADATA
metadata_stop(); // close down the metadata pipe
metadata_stop(); // close down the metadata pipe
#endif
activity_monitor_stop(0);
activity_monitor_stop(0);
if ((config.output) && (config.output->deinit)) {
debug(2, "Deinitialise the audio backend.");
config.output->deinit();
}
if ((config.output) && (config.output->deinit)) {
debug(2, "Deinitialise the audio backend.");
config.output->deinit();
}
#ifdef CONFIG_SOXR
// be careful -- not sure if the thread can be cancelled cleanly, so wait for it to shut down
pthread_join(soxr_time_check_thread, NULL);
// be careful -- not sure if the thread can be cancelled cleanly, so wait for it to shut down
pthread_join(soxr_time_check_thread, NULL);
#endif
if (conns)
free(conns); // make sure the connections have been deleted first
if (conns)
free(conns); // make sure the connections have been deleted first
if (config.service_name)
free(config.service_name);
if (config.service_name)
free(config.service_name);
#ifdef CONFIG_CONVOLUTION
if (config.convolution_ir_file)
free(config.convolution_ir_file);
if (config.convolution_ir_file)
free(config.convolution_ir_file);
#endif
if (config.regtype)
free(config.regtype);
if (config.regtype)
free(config.regtype);
#ifdef CONFIG_LIBDAEMON
if (this_is_the_daemon_process) {
daemon_retval_send(0);
daemon_pid_file_remove();
daemon_signal_done();
if (config.computed_piddir)
free(config.computed_piddir);
}
}
if (this_is_the_daemon_process) {
daemon_retval_send(0);
daemon_pid_file_remove();
daemon_signal_done();
if (config.computed_piddir)
free(config.computed_piddir);
}
}
#endif
if (config.cfg)
config_destroy(config.cfg);
if (config.appName)
free(config.appName);
// probably should be freeing malloc'ed memory here, including strdup-created strings...
if (config.cfg)
config_destroy(config.cfg);
if (config.appName)
free(config.appName);
// probably should be freeing malloc'ed memory here, including strdup-created strings...
#ifdef CONFIG_LIBDAEMON
if (this_is_the_daemon_process) { // this is the daemon that is exiting
debug(1,"libdaemon daemon exit");
} else {
if (config.daemonise)
debug(1,"libdaemon parent exit");
else
debug(1,"exit");
}
if (this_is_the_daemon_process) { // this is the daemon that is exiting
debug(1, "libdaemon daemon exit");
} else {
if (config.daemonise)
debug(1, "libdaemon parent exit");
else
debug(1, "exit");
}
#else
debug(1,"exit");
debug(1, "exit");
#endif
} else {
debug(1,"emergency exit");
}
} else {
debug(1, "emergency exit");
}
}
// for removing zombie script processes
@@ -1434,11 +1437,10 @@ void handle_sigchld(__attribute__((unused)) int sig) {
}
void main_thread_cleanup_handler(__attribute__((unused)) void *arg) {
debug(2,"main thread cleanup handler called");
debug(2, "main thread cleanup handler called");
exit(EXIT_SUCCESS);
}
int main(int argc, char **argv) {
/* Check if we are called with -V or --version parameter */
if (argc >= 2 && ((strcmp(argv[1], "-V") == 0) || (strcmp(argv[1], "--version") == 0))) {
@@ -1455,7 +1457,7 @@ int main(int argc, char **argv) {
#ifdef CONFIG_LIBDAEMON
pid = getpid();
#endif
config.log_fd = -1;
config.log_fd = -1;
conns = NULL; // no connections active
memset((void *)&main_thread_id, 0, sizeof(main_thread_id));
memset(&config, 0, sizeof(config)); // also clears all strings, BTW
@@ -1476,7 +1478,7 @@ int main(int argc, char **argv) {
setlogmask(LOG_UPTO(LOG_DEBUG));
openlog(NULL, 0, LOG_DAEMON);
#endif
emergency_exit = 0; // what to do or skip in the exit_function
emergency_exit = 0; // what to do or skip in the exit_function
atexit(exit_function);
// set defaults
@@ -1504,6 +1506,15 @@ int main(int argc, char **argv) {
// set non-zero / non-NULL default values here
// but note that audio back ends also have a chance to set defaults
// get the first output backend in the list and make it the default
audio_output *first_backend = audio_get_output(NULL);
if (first_backend == NULL) {
die("No audio backend found! Check your build of Shairport Sync.");
} else {
strncpy(first_backend_name, first_backend->name, sizeof(first_backend_name) - 1);
config.output_name = first_backend_name;
}
strcpy(configuration_file_path, SYSCONFDIR);
// strcat(configuration_file_path, "/shairport-sync"); // thinking about adding a special
// shairport-sync directory
@@ -1598,13 +1609,16 @@ int main(int argc, char **argv) {
if (errno == ENOENT)
daemon_log(LOG_WARNING, "Failed to kill %s daemon: PID file not found.", config.appName);
else
daemon_log(LOG_WARNING, "Failed to kill %s daemon: \"%s\", errno %u.", config.appName, strerror(errno), errno);
daemon_log(LOG_WARNING, "Failed to kill %s daemon: \"%s\", errno %u.", config.appName,
strerror(errno), errno);
} else {
// debug(1,"Successfully killed the %s daemon.", config.appName);
// debug(1,"Successfully killed the %s daemon.", config.appName);
if (daemon_pid_file_remove() == 0)
debug(2, "killed the %s daemon.", config.appName);
else
daemon_log(LOG_WARNING, "killed the %s deamon, but cannot remove old PID file: \"%s\", errno %u.", config.appName, strerror(errno), errno);
daemon_log(LOG_WARNING,
"killed the %s deamon, but cannot remove old PID file: \"%s\", errno %u.",
config.appName, strerror(errno), errno);
}
return ret < 0 ? 1 : 0;
#else
@@ -1651,14 +1665,19 @@ int main(int argc, char **argv) {
case 0:
break;
case 1:
daemon_log(LOG_ERR,
"the %s daemon failed to launch: could not close open file descriptors after forking.", config.appName);
daemon_log(
LOG_ERR,
"the %s daemon failed to launch: could not close open file descriptors after forking.",
config.appName);
break;
case 2:
daemon_log(LOG_ERR, "the %s daemon failed to launch: could not create PID file.", config.appName);
daemon_log(LOG_ERR, "the %s daemon failed to launch: could not create PID file.",
config.appName);
break;
case 3:
daemon_log(LOG_ERR, "the %s daemon failed to launch: could not create or access PID directory.", config.appName);
daemon_log(LOG_ERR,
"the %s daemon failed to launch: could not create or access PID directory.",
config.appName);
break;
default:
daemon_log(LOG_ERR, "the %s daemon failed to launch, error %i.", config.appName, ret);
@@ -1709,8 +1728,8 @@ int main(int argc, char **argv) {
#endif
debug(1, "Started!");
// stop a pipe signal from killing the program
signal(SIGPIPE, SIG_IGN);
// stop a pipe signal from killing the program
signal(SIGPIPE, SIG_IGN);
// install a zombie process reaper
// see: http://www.microhowto.info/howto/reap_zombie_processes_using_a_sigchld_handler.html