From ff4ab45eba7c791484e5407fa9ef1c2f3d1317ba Mon Sep 17 00:00:00 2001 From: Thomas Ries Date: Fri, 8 Jun 2007 19:43:44 +0000 Subject: [PATCH] - Some cleanup in dejitter code --- src/Makefile.am | 6 +- src/dejitter.c | 451 +++++++++++++++++++++++++++++++++ src/dejitter.h | 52 ++++ src/rtpproxy.c | 4 +- src/rtpproxy.h | 34 +-- src/rtpproxy_relay.c | 575 +++++++------------------------------------ src/siproxd.h | 14 +- 7 files changed, 634 insertions(+), 502 deletions(-) create mode 100644 src/dejitter.c create mode 100644 src/dejitter.h diff --git a/src/Makefile.am b/src/Makefile.am index d69ac17..1df7535 100644 --- a/src/Makefile.am +++ b/src/Makefile.am @@ -1,5 +1,5 @@ # -# Copyright (C) 2002-2005 Thomas Ries +# Copyright (C) 2002-2007 Thomas Ries # # This file is part of Siproxd. # @@ -30,7 +30,7 @@ siproxd_SOURCES = siproxd.c proxy.c register.c sock.c utils.c \ sip_utils.c sip_layer.c log.c readconf.c rtpproxy.c \ rtpproxy_relay.c accessctl.c route_processing.c \ security.c auth.c fwapi.c resolve.c \ - plugin_shortdial.c + plugin_shortdial.c dejitter.c # addrcache.c # @@ -42,7 +42,7 @@ libcustom_fw_module_a_SOURCES = custom_fw_module.c noinst_HEADERS = log.h siproxd.h digcalc.h rtpproxy.h \ - fwapi.h plugins.h addrcache.h + fwapi.h plugins.h dejitter.h addrcache.h EXTRA_DIST = .buildno diff --git a/src/dejitter.c b/src/dejitter.c new file mode 100644 index 0000000..b41c446 --- /dev/null +++ b/src/dejitter.c @@ -0,0 +1,451 @@ +/* + Copyright (C) 2006-2007 Hans Carlos Hofmann , + Thomas Ries + + This file is part of Siproxd. + + Siproxd is free software; you can redistribute it and/or modify + it under the terms of the GNU General Public License as published by + the Free Software Foundation; either version 2 of the License, or + (at your option) any later version. + + Siproxd is distributed in the hope that it will be useful, + but WITHOUT ANY WARRANTY; without even the implied warrantry of + MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + GNU General Public License for more details. + + You should have received a copy of the GNU General Public License + along with Siproxd; if not, write to the Free Software + Foundation, Inc., 59 Temple Place, Suite 330, Boston, MA 02111-1307 USA +*/ +#include "config.h" + +//#include +//#include +//#include +#include +//#include +//#include + +#include +#include +//#include + + +#include +#include "siproxd.h" +#include "rtpproxy.h" +#include "log.h" + +#include "dejitter.h" + +#ifdef GPL + +/* + * static forward declarations + */ +static void add_time_values(const struct timeval *a, + const struct timeval *b, struct timeval *r); +static void sub_time_values(const struct timeval *a, const struct timeval *b, + struct timeval *r); +static int cmp_time_values(const struct timeval *a, const struct timeval *b); +static double make_double_time(const struct timeval *tv); +static void send_top_of_que(int nolock); +static void split_double_time(double d, struct timeval *tv); +static int fetch_missalign_long_network_oder(char *where); + + +/* + * RTP buffers for dejitter + */ + +/* + * table to buffer date for dejitter function + */ +#define NUMBER_OF_BUFFER (10*RTPPROXY_SIZE) +static rtp_delayed_message rtp_buffer_area[NUMBER_OF_BUFFER]; + +static rtp_delayed_message *free_memory; +static rtp_delayed_message *msg_que; + +static struct timeval minstep; + + + +/* + * Initialize RTP dejitter + */ +void dejitter_init(void) { + int i; + rtp_delayed_message *m; + + memset (&rtp_buffer_area, 0, sizeof(rtp_buffer_area)); + free_memory = NULL; + msg_que = NULL; + for (i=0,m=&rtp_buffer_area[0];inext = free_memory; + free_memory = m; + } +} + +/* + * Delayed send + */ +void dejitter_delayedsendto(int s, const void *msg, size_t len, int flags, + const struct sockaddr_in *to, + const struct timeval *tv, + const struct timeval *current_tv, + rtp_proxytable_t *errret, int nolock) { + rtp_delayed_message *m; + rtp_delayed_message **linkin; + + if (!free_memory) send_top_of_que(nolock); + + m = free_memory; + + m->socked = s; + memcpy(&(m->rtp_buff), msg, m->message_len = len); + m->flags = flags; + m->dst_addr = *to; + m->transm_time = *tv; + m->errret = errret; + + free_memory = m->next; + + if (cmp_time_values(current_tv,tv) >= 0) { + m->next = msg_que; + msg_que = m; + send_top_of_que(nolock); + } else { + linkin = &msg_que; + while ((*linkin != NULL) && + (cmp_time_values(&((*linkin)->transm_time),tv) < 0)) { + linkin = (rtp_delayed_message **)&((*linkin)->next); + } + m->next = *linkin; + *linkin = m; + } +} + +/* + * Cancel a message + */ +void dejitter_cancel(rtp_proxytable_t *dropentry) { + rtp_delayed_message **linkout; + rtp_delayed_message *m; + + linkout = &msg_que; + + while (*linkout != NULL) { + if ((*linkout)->errret == dropentry) { + m = *linkout; + *linkout = m->next; + m->next = free_memory; + free_memory = m; + } else { + linkout = (rtp_delayed_message **)&((*linkout)->next); + } + } +} + +/* + * Flush buffers + */ +void dejitter_flush(struct timeval *current_tv, int nolock) { + struct timezone tz; + + while (msg_que && + (cmp_time_values(&(msg_que->transm_time),current_tv)<=0)) { + send_top_of_que(nolock); + gettimeofday(current_tv,&tz); + } +} + +/* + * Delay of next transmission + */ +int dejitter_delay_of_next_tx(struct timeval *tv,struct timeval *current_tv) { + struct timezone tz ; + + if (msg_que) { + gettimeofday(current_tv,&tz); + sub_time_values(&(msg_que->transm_time),current_tv,tv); + if (cmp_time_values(tv,&minstep)<=0) { + *tv = minstep ; + } + return -1; + } + return 0; +} + +/* + * Initialize calculation of transmit the frame + */ +void dejitter_init_time(timecontrol_t *tc, int dejitter) { + struct timezone tz; + + minstep.tv_sec = 0; + minstep.tv_usec = 6000; + memset(tc, 0, sizeof(*tc)); + if (dejitter>0) { + gettimeofday(&(tc->starttime),&tz); + + tc->dejitter = dejitter; + tc->dejitter_d = dejitter; + split_double_time(tc->dejitter_d, &(tc->dejitter_tv)); + } +} + +/* + * Calculate transmit time + */ +void dejitter_calc_tx_time(rtp_buff_t *rtp_buff, timecontrol_t *tc, + struct timeval *input_tv, + struct timeval *ttv) { + int packet_time_code; + double currenttime; + double calculatedtime = 0; + double calculatedtime2 = 0; + struct timeval input_r_tv; + struct timeval output_r_tv; + + if (!tc || !tc->dejitter) { + *ttv = *input_tv; + return; + } + + + /* I hate this computer language ... :-/ quite confuse ! look modula */ + packet_time_code = fetch_missalign_long_network_oder(&((*rtp_buff)[4])); + +/*&&&& beware, it seems that when sending RTP events (payload type +telephone-event) the timestamp does not increment and stays the same. +The sequence number however DOES increment. This could lead to confusion when +transmitting RTP events (like DTMF). How can we handle this? Check for RTP event +and then do an "educated guess" for the to-be timestamp? +*/ + if (tc->calccount == 0) { + DEBUGC(DBCLASS_RTP, "initialise time calculatin"); + tc->starttime = *input_tv; + tc->time_code_a = packet_time_code; + } + + sub_time_values(input_tv,&(tc->starttime),&input_r_tv); + + calculatedtime = currenttime = make_double_time(&input_r_tv); + if (tc->calccount < 10) { + DEBUGC(DBCLASS_RTP, "initial data stage 1 %f usec", currenttime); + tc->received_a = currenttime / (packet_time_code - tc->time_code_a); + } else if (tc->calccount < 20) { + tc->received_a = 0.95 * tc->received_a + 0.05 * currenttime / + (packet_time_code - tc->time_code_a); + } else { + tc->received_a = 0.99 * tc->received_a + 0.01 * currenttime / + (packet_time_code - tc->time_code_a); + } + if (tc->calccount > 20) { + if (!tc->time_code_b) { + tc->time_code_b = packet_time_code; + tc->received_b = currenttime; + } else if (tc->time_code_b < packet_time_code) { + calculatedtime = tc->received_b = tc->received_b + + (packet_time_code - tc->time_code_b) * tc->received_a; + tc->time_code_b = packet_time_code; + if (tc->calccount < 28) { + tc->received_b = 0.90 * tc->received_b + 0.1 * currenttime; + } else if (tc->calccount < 300) { + tc->received_b = 0.95 * tc->received_b + 0.05 * currenttime; + } else { + tc->received_b = 0.99 * tc->received_b + 0.01 * currenttime; + } + } else { + calculatedtime = tc->received_b + + (packet_time_code - tc->time_code_b) * tc->received_a; + } + } + tc->received_c = currenttime; + tc->time_code_c = packet_time_code; + + if (tc->calccount < 30) { + /* + * But in the start phase, + * we asume every packet as not delayed. + */ + calculatedtime = currenttime; + } + + /* + ** theoretical value for F1000 Phone + */ + //calculatedtime = (tc->received_a = 125.) * packet_time_code; + + tc->calccount ++; + calculatedtime += tc->dejitter_d; + + if (calculatedtime < currenttime) { + calculatedtime = currenttime; + } else if (calculatedtime > currenttime + 2.* tc->dejitter_d) { + calculatedtime = currenttime + 2.* tc->dejitter_d; + } + + /* every 500 counts show statistics */ + if (tc->calccount % 500 == 0) { + DEBUGC(DBCLASS_RTPBABL, "currenttime = %f", currenttime); + DEBUGC(DBCLASS_RTPBABL, "packetcode = %i", packet_time_code); + DEBUGC(DBCLASS_RTPBABL, "timecodes %i, %i, %i", + tc->time_code_a, tc->time_code_b, tc->time_code_c); + DEBUGC(DBCLASS_RTPBABL, "measuredtimes %f usec, %f usec, %f usec", + tc->received_a, tc->received_b, tc->received_c); + DEBUGC(DBCLASS_RTPBABL, "p2 - p1 = (%i,%f usec)", + tc->time_code_b - tc->time_code_a, + tc->received_b - tc->received_a); + if (tc->time_code_c) { + DEBUGC(DBCLASS_RTPBABL, "p3 - p2 = (%i,%f usec)", + tc->time_code_c - tc->time_code_b, + tc->received_c - tc->received_b); + } + DEBUGC(DBCLASS_RTPBABL, "calculatedtime = %f", calculatedtime); + if (calculatedtime2) { + DEBUGC(DBCLASS_RTPBABL, "calculatedtime2 = %f", calculatedtime2); + } + DEBUGC(DBCLASS_RTPBABL, "transmtime = %f (%f)", calculatedtime / + (160. * tc->received_a) - packet_time_code / 160, + currenttime / (160. * tc->received_a) - + packet_time_code / 160); + DEBUGC(DBCLASS_RTPBABL, "synthetic latency = %f, %f, %f, %i, %i", + calculatedtime-currenttime, calculatedtime, + currenttime, packet_time_code, + packet_time_code / 160); + } + + split_double_time(calculatedtime, &output_r_tv); + add_time_values(&output_r_tv,&(tc->starttime),ttv); +} + + + +/* + * Add timeval times + */ +static void add_time_values(const struct timeval *a, + const struct timeval *b, struct timeval *r) { + r->tv_sec = a->tv_sec + b->tv_sec; + r->tv_usec = a->tv_usec + b->tv_usec; + if (r->tv_usec >= 1000000) { + r->tv_usec -= 1000000; + r->tv_sec++; + } +} + +/* + * Subtract timeval values + */ +static void sub_time_values(const struct timeval *a, + const struct timeval *b, struct timeval *r) { + if ((a->tv_sec < b->tv_sec) || + ((a->tv_sec == b->tv_sec) && (a->tv_usec < b->tv_usec))) { + r->tv_usec = 0; + r->tv_sec = 0; + return; + } + if (a->tv_usec < b->tv_usec) { + r->tv_sec = a->tv_sec - b->tv_sec - 1; + r->tv_usec = a->tv_usec + 1000000 - b->tv_usec; + } else { + r->tv_sec = a->tv_sec - b->tv_sec; + r->tv_usec = a->tv_usec - b->tv_usec; + } +} +/* + * Compare timeval values + */ +static int cmp_time_values(const struct timeval *a, const struct timeval *b) { + if (a->tv_sec < b->tv_sec) return -1; + if (a->tv_sec > b->tv_sec) return 1; + if (a->tv_usec < b->tv_usec) return -1; + if (a->tv_usec > b->tv_usec) return 1; + return 0; +} + +/* + * Convert TIMEVAL to DOUBLE + */ +static double make_double_time(const struct timeval *tv) { + return 1000000.0 * tv->tv_sec + tv->tv_usec; +} + +/* + * Send Top of queue + * + * nolock: 1 - do not lock the mutex - lock already owned! + */ +static void send_top_of_que(int nolock) { + rtp_delayed_message *m; + int sts; + + if (msg_que) { + m = msg_que; + msg_que = m->next; + m->next = free_memory; + free_memory = m; + + if ((m->errret != NULL) && (m->errret->rtp_tx_sock)) { + sts = sendto(m->socked, &(m->rtp_buff), m->message_len, + m->flags, (const struct sockaddr *)&(m->dst_addr), + (socklen_t)sizeof(m->dst_addr)); + if ((sts == -1) && (m->errret != NULL) && (errno != ECONNREFUSED)) { + osip_call_id_t callid; + + ERROR("sendto() [%s:%i size=%i] delayed call failed: %s", + utils_inet_ntoa(m->errret->remote_ipaddr), + m->errret->remote_port, m->message_len, strerror(errno)); + + /* if sendto() fails with bad filedescriptor, + * this means that the opposite stream has been + * canceled or timed out. + * we should then cancel this stream as well.*/ + + WARN("stopping opposite stream"); + + callid.number=m->errret->callid_number; + callid.host=m->errret->callid_host; + + /* caller tells us if we must lock the fdset mutex */ + sts = rtp_relay_stop_fwd(&callid, + m->errret->direction, + m->errret->media_stream_no, + nolock); + if (sts != STS_SUCCESS) { + /* force the streams to timeout on next occasion */ + m->errret->timestamp=0; + } + } /* if sendto fails */ + } + } /* if (msg_que) */ +} + +/* + * Convert DOUBLE time into TIMEVAL + */ +static void split_double_time(double d, struct timeval *tv) { + tv->tv_sec = d / 1000000.0; + tv->tv_usec = d - 1000000.0 * tv->tv_sec; +} + +/* + * + */ +static int fetch_missalign_long_network_oder(char *where) { + int i = 0; + int k; + int j; + + for (j=0;j<4;j++) { + k = *where; + i = (i<<8) | (0xFF & k); + where ++; + } + return i; +} + +#endif diff --git a/src/dejitter.h b/src/dejitter.h new file mode 100644 index 0000000..b6972bf --- /dev/null +++ b/src/dejitter.h @@ -0,0 +1,52 @@ +/* + Copyright (C) 2006-2007 Hans Carlos Hofmann , + Thomas Ries + + This file is part of Siproxd. + + Siproxd is free software; you can redistribute it and/or modify + it under the terms of the GNU General Public License as published by + the Free Software Foundation; either version 2 of the License, or + (at your option) any later version. + + Siproxd is distributed in the hope that it will be useful, + but WITHOUT ANY WARRANTY; without even the implied warrantry of + MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + GNU General Public License for more details. + + You should have received a copy of the GNU General Public License + along with Siproxd; if not, write to the Free Software + Foundation, Inc., 59 Temple Place, Suite 330, Boston, MA 02111-1307 USA +*/ + +#ifdef GPL +#define USE_DEJITTER + +typedef struct { + void *next; /* next free or next que element */ + int socked; /* socket number */ + size_t message_len; /* length of message */ + int flags; /* flags */ + struct sockaddr_in dst_addr; /* where shall i send */ + struct timeval transm_time; /* when shall i send */ + rtp_proxytable_t *errret; /* deliver error status */ + rtp_buff_t rtp_buff; /* Data storage */ +} rtp_delayed_message; + + +/* dejitter */ +void dejitter_init(void); +void dejitter_delayedsendto(int s, const void *msg, size_t len, int flags, + const struct sockaddr_in *to, + const struct timeval *tv, + const struct timeval *current_tv, + rtp_proxytable_t *errret, int nolock); +void dejitter_cancel(rtp_proxytable_t *dropentry); +void dejitter_flush(struct timeval *current_tv, int nolock); +int dejitter_delay_of_next_tx(struct timeval *tv, struct timeval *current_tv); +void dejitter_init_time(timecontrol_t *tc, int dejitter); +void dejitter_calc_tx_time(rtp_buff_t *rtp_buff, timecontrol_t *tc, + struct timeval *input_tv, + struct timeval *ttv); + +#endif diff --git a/src/rtpproxy.c b/src/rtpproxy.c index cdb0009..7a853ee 100644 --- a/src/rtpproxy.c +++ b/src/rtpproxy.c @@ -1,5 +1,5 @@ /* - Copyright (C) 2003-2005 Thomas Ries + Copyright (C) 2003-2007 Thomas Ries This file is part of Siproxd. @@ -116,7 +116,7 @@ int rtp_stop_fwd (osip_call_id_t *callid, int direction) { if (configuration.rtp_proxy_enable == 0) { sts = STS_SUCCESS; } else if (configuration.rtp_proxy_enable == 1) { // Relay - sts = rtp_relay_stop_fwd(callid, direction, -1, 0); + sts = rtp_relay_stop_fwd(callid, direction, -1, LOCK_FDSET); } else { ERROR("CONFIG: rtp_proxy_enable has invalid value: %d", configuration.rtp_proxy_enable); diff --git a/src/rtpproxy.h b/src/rtpproxy.h index 17f7700..8f4d915 100644 --- a/src/rtpproxy.h +++ b/src/rtpproxy.h @@ -1,5 +1,5 @@ /* - Copyright (C) 2003-2005 Thomas Ries + Copyright (C) 2003-2007 Thomas Ries This file is part of Siproxd. @@ -23,19 +23,22 @@ #define CALLIDNUM_SIZE 256 #define CALLIDHOST_SIZE 128 -typedef struct { - struct timeval starttime ; - int calccount ; - int dejitter ; - struct timeval dejitter_tv ; - double dejitter_d ; - int time_code_a ; - double received_a ; /* time in µsec sience epoch */ - int time_code_b ; - double received_b ; /* time in µsec sience epoch */ - int time_code_c ; - double received_c ; /* time in µsec sience epoch */ +/* Buffer structure to store an RTP frame */ +typedef char rtp_buff_t[RTP_BUFFER_SIZE]; +/* Time control structure to be used with De-jitter feature */ +typedef struct { + struct timeval starttime ; + int calccount ; + int dejitter ; + struct timeval dejitter_tv ; + double dejitter_d ; + int time_code_a ; + double received_a ; /* time in µsec since epoch */ + int time_code_b ; + double received_b ; /* time in µsec since epoch */ + int time_code_c ; + double received_c ; /* time in µsec since epoch */ } timecontrol_t ; typedef struct { @@ -48,7 +51,7 @@ typedef struct { client_id_t client_id; int direction; /* Direction of RTP stream */ int media_stream_no; - timecontrol_t tc; + timecontrol_t tc; /* de-jitter feature */ struct in_addr local_ipaddr; /* local IP */ int local_port; /* local allocated port */ struct in_addr remote_ipaddr; /* remote IP */ @@ -68,3 +71,6 @@ int rtp_relay_start_fwd (osip_call_id_t *callid, client_id_t client_id, int dejitter); int rtp_relay_stop_fwd (osip_call_id_t *callid, int rtp_direction, int media_stream_no, int nolock); + +#define NOLOCK_FDSET 1 +#define LOCK_FDSET 0 diff --git a/src/rtpproxy_relay.c b/src/rtpproxy_relay.c index 8621b41..122e36c 100644 --- a/src/rtpproxy_relay.c +++ b/src/rtpproxy_relay.c @@ -1,5 +1,5 @@ /* - Copyright (C) 2003-2005 Thomas Ries + Copyright (C) 2003-2007 Thomas Ries This file is part of Siproxd. @@ -42,6 +42,10 @@ #include "rtpproxy.h" #include "log.h" +#ifdef GPL + #include "dejitter.h" +#endif + #if !defined(SOL_IP) #define SOL_IP IPPROTO_IP #endif @@ -71,65 +75,13 @@ static pthread_t rtpproxy_tid=0; static fd_set master_fdset; static int master_fd_max; -/* - * RTP buffers for dejitter - */ -typedef char rtp_buff_t[RTP_BUFFER_SIZE]; -typedef struct { - void *next; /* next free or next que element */ - int socked; /* socket number */ - size_t message_len; /* length of message */ - int flags; /* flags */ - struct sockaddr_in dst_addr; /* where shall i send */ - struct timeval transm_time; /* when shall i send */ - rtp_proxytable_t *errret; /* deliver error status */ - rtp_buff_t rtp_buff; /* Data storage */ -} rtp_delayed_message; - - -/* - * table to buffer date for dejitter function - */ -#define NUMBER_OF_BUFFER (10*RTPPROXY_SIZE) -static rtp_delayed_message rtp_buffer_area[NUMBER_OF_BUFFER]; - -static rtp_delayed_message *free_memory; -static rtp_delayed_message *msg_que; - -static struct timeval current_tv; -static struct timeval minstep; - /* * forward declarations of internal functions */ -static void *rtpproxy_main(void *i); -static int rtp_recreate_fdset(void); -void rtpproxy_kill( void ); static void sighdl_alm(int sig) {/* just wake up from select() */}; - -/* dejitter */ -static void rtp_buffer_init (); -static void add_time_values(const struct timeval *a, - const struct timeval *b, struct timeval *r); -static void sub_time_values(const struct timeval *a, const struct timeval *b, - struct timeval *r); -static int cmp_time_values(const struct timeval *a, const struct timeval *b); -static double make_double_time(const struct timeval *tv); -static void send_top_of_que(void); -static void delayedsendto(int s, const void *msg, size_t len, int flags, - const struct sockaddr_in *to, - const struct timeval *tv, rtp_proxytable_t *errret); -static void cancelmessages(rtp_proxytable_t *dropentry); -static void flushbuffers(void); -static int delay_of_next_transmission(struct timeval *tv); -static void split_double_time(double d, struct timeval *tv); -static void init_calculate_transmit_time(timecontrol_t *tc, int dejitter); -static int fetch_missalign_long_network_oder(char *where); -static void calculate_transmit_time(rtp_buff_t *rtp_buff, timecontrol_t *tc, - const struct timeval *input_tv, - struct timeval *ttv); - -/* */ +static void *rtpproxy_main(void *i); +static void rtpproxy_kill( void ); +static int rtp_recreate_fdset(void); static void match_socket (int rtp_proxytable_idx); static void error_handler (int rtp_proxytable_idx, int socket_type); @@ -145,7 +97,9 @@ int rtp_relay_init( void ) { int arg=0; struct sigaction sigact; - rtp_buffer_init(); +#ifdef USE_DEJITTER + dejitter_init(); +#endif atexit(rtpproxy_kill); /* cancel RTP thread at exit */ @@ -219,13 +173,13 @@ int rtp_relay_init( void ) { static void *rtpproxy_main(void *arg) { fd_set fdset; int fd_max; - int i; + int i, sts; int num_fd; static rtp_buff_t rtp_buff; int count; struct timeval last_tv ; struct timeval sleep_tv ; - struct timeval input_tv ; + struct timeval current_tv ; struct timezone tz ; memcpy(&fdset, &master_fdset, sizeof(fdset)); @@ -236,21 +190,24 @@ static void *rtpproxy_main(void *arg) { /* loop forever... */ for (;;) { - if (!delay_of_next_transmission(&sleep_tv)) - { -// DEBUGC(DBCLASS_RTP, "rtp que empty") ; - sleep_tv.tv_sec = 5; - sleep_tv.tv_usec = 0; +#ifdef USE_DEJITTER + /* calculate time until next packet to send from dejitter buffer */ + if (!dejitter_delay_of_next_tx(&sleep_tv, ¤t_tv)) { + sleep_tv.tv_sec = 5; + sleep_tv.tv_usec = 0; }; +#else + sleep_tv.tv_sec = 5; + sleep_tv.tv_usec = 0; +#endif num_fd=select(fd_max+1, &fdset, NULL, NULL, &sleep_tv); - gettimeofday(&input_tv,&tz); - current_tv = input_tv; + gettimeofday(¤t_tv, &tz); - /* - * Send delayed Packets - */ - flushbuffers(); +#ifdef USE_DEJITTER + /* Send delayed Packets that are timed to be send */ + dejitter_flush(¤t_tv, LOCK_FDSET); +#endif /* exit point for this thread in case of program terminaction */ pthread_testcancel(); @@ -305,10 +262,6 @@ static void *rtpproxy_main(void *arg) { if (rtp_proxytable[i].rtp_con_tx_sock != 0) { struct sockaddr_in dst_addr; - struct timeval ttv ; - - add_time_values(&(rtp_proxytable[i].tc.dejitter_tv), - &input_tv,&ttv) ; /* write to dest via socket rtp_con_tx_sock */ dst_addr.sin_family = AF_INET; @@ -316,18 +269,18 @@ static void *rtpproxy_main(void *arg) { &rtp_proxytable[i].remote_ipaddr, sizeof(struct in_addr)); dst_addr.sin_port= htons(rtp_proxytable[i].remote_port+1); - delayedsendto(rtp_proxytable[i].rtp_con_tx_sock, rtp_buff, - count, 0, &dst_addr, &ttv, &rtp_proxytable[i]) ; + + /* Don't dejitter RTCP packets */ + sts = sendto(rtp_proxytable[i].rtp_con_tx_sock, rtp_buff, + count, 0, (const struct sockaddr *)&dst_addr, + (socklen_t)sizeof(dst_addr)); + /* ignore errors here. We don't know if the remote + site does receive RTCP messages at all (or reject + them with ICMP-whatever). If it fails, it is lost. + Basta, end of story. */ } } /* count > 0 */ - /* update timestamp of last usage for both (RX and TX) entries. - * This allows silence (no data) on one stream without breaking - * the connection after the RTP timeout */ - rtp_proxytable[i].timestamp=current_tv.tv_sec; - if (rtp_proxytable[i].opposite_entry > 0) { - rtp_proxytable[rtp_proxytable[i].opposite_entry-1].timestamp= - current_tv.tv_sec; - } + /* RTCP does not wind up the keepalive timestamp. */ } /* if */ /* @@ -361,10 +314,9 @@ static void *rtpproxy_main(void *arg) { if (rtp_proxytable[i].rtp_tx_sock != 0) { struct sockaddr_in dst_addr; +#ifdef USE_DEJITTER struct timeval ttv; - - calculate_transmit_time(&rtp_buff,&(rtp_proxytable[i].tc), - &input_tv,&ttv) ; +#endif /* write to dest via socket rtp_tx_sock */ dst_addr.sin_family = AF_INET; @@ -372,8 +324,45 @@ static void *rtpproxy_main(void *arg) { &rtp_proxytable[i].remote_ipaddr, sizeof(struct in_addr)); dst_addr.sin_port= htons(rtp_proxytable[i].remote_port); - delayedsendto(rtp_proxytable[i].rtp_tx_sock, rtp_buff, - count, 0, &dst_addr, &ttv, &rtp_proxytable[i]); + +#ifdef USE_DEJITTER + dejitter_calc_tx_time(&rtp_buff, &(rtp_proxytable[i].tc), + ¤t_tv, &ttv); + dejitter_delayedsendto(rtp_proxytable[i].rtp_tx_sock, + rtp_buff, count, 0, &dst_addr, + &ttv, ¤t_tv, + &rtp_proxytable[i], NOLOCK_FDSET); +#else + sts = sendto(rtp_proxytable[i].rtp_tx_sock, rtp_buff, + count, 0, (const struct sockaddr *)&dst_addr, + (socklen_t)sizeof(dst_addr)); + if (sts == -1) { + if (errno != ECONNREFUSED) { + osip_call_id_t callid; + + ERROR("sendto() [%s:%i size=%i] call failed: %s", + utils_inet_ntoa(rtp_proxytable[i].remote_ipaddr), + rtp_proxytable[i].remote_port, count, strerror(errno)); + + /* if sendto() fails with bad filedescriptor, + * this means that the opposite stream has been + * canceled or timed out. + * we should then cancel this stream as well.*/ + + WARN("stopping opposite stream"); + callid.number=rtp_proxytable[i].callid_number; + callid.host=rtp_proxytable[i].callid_host; + /* don't lock the mutex, as we own the lock already */ + sts = rtp_relay_stop_fwd(&callid, + rtp_proxytable[i].direction, + -1, NOLOCK_FDSET); + if (sts != STS_SUCCESS) { + /* force the streams to timeout on next occasion */ + rtp_proxytable[i].timestamp=0; + } + } + } +#endif } } /* count > 0 */ /* update timestamp of last usage for both (RX and TX) entries. @@ -401,7 +390,9 @@ static void *rtpproxy_main(void *arg) { /* this one has expired, clean it up */ callid.number=rtp_proxytable[i].callid_number; callid.host=rtp_proxytable[i].callid_host; - cancelmessages(&rtp_proxytable[i]); +#ifdef USE_DEJITTER + dejitter_cancel(&rtp_proxytable[i]); +#endif INFO("RTP stream %s@%s (media=%i) has expired", callid.number, callid.host, rtp_proxytable[i].media_stream_no); @@ -415,7 +406,8 @@ static void *rtpproxy_main(void *arg) { * This may be a multiple stream conversation (audio/video) and * just one (unused?) has timed out. Seen with VoIPEX PBX! */ rtp_relay_stop_fwd(&callid, rtp_proxytable[i].direction, - rtp_proxytable[i].media_stream_no, 1); + rtp_proxytable[i].media_stream_no, + NOLOCK_FDSET); } /* if */ } /* for i */ } /* if (t>...) */ @@ -536,10 +528,10 @@ int rtp_relay_start_fwd (osip_call_id_t *callid, client_id_t client_id, sizeof(remote_ipaddr)); } - /* - * set up timecrontrol for dejitter function - */ - init_calculate_transmit_time(&rtp_proxytable[i].tc,dejitter); +#ifdef USE_DEJITTER + /* Initialize up timecrontrol for dejitter function */ + dejitter_init_time(&rtp_proxytable[i].tc, dejitter); +#endif /* return the already known local port number */ @@ -699,10 +691,10 @@ int rtp_relay_start_fwd (osip_call_id_t *callid, client_id_t client_id, rtp_proxytable[freeidx].remote_port=remote_port; time(&rtp_proxytable[freeidx].timestamp); - /* - * set up timecrontrol for dejitter function - */ - init_calculate_transmit_time(&rtp_proxytable[freeidx].tc,dejitter); +#ifdef USE_DEJITTER + /* Initialize up timecrontrol for dejitter function */ + dejitter_init_time(&rtp_proxytable[freeidx].tc, dejitter); +#endif *local_port=port; @@ -922,7 +914,7 @@ static int rtp_recreate_fdset(void) { * RETURNS * - */ -void rtpproxy_kill( void ) { +static void rtpproxy_kill( void ) { void *thread_status; osip_call_id_t cid; int i, sts; @@ -933,7 +925,8 @@ void rtpproxy_kill( void ) { cid.number = rtp_proxytable[i].callid_number; cid.host = rtp_proxytable[i].callid_host; sts = rtp_relay_stop_fwd(&cid, rtp_proxytable[i].direction, - rtp_proxytable[i].media_stream_no, 0); + rtp_proxytable[i].media_stream_no, + LOCK_FDSET); } } @@ -950,382 +943,6 @@ void rtpproxy_kill( void ) { } -/*********** - * De-Jitter - ***********/ - -/* - * Initialize RTP dejitter - */ -static void rtp_buffer_init () { - int i; - rtp_delayed_message *m; - - memset (&rtp_buffer_area, 0, sizeof(rtp_buffer_area)); - free_memory = NULL; - msg_que = NULL; - for (i=0,m=&rtp_buffer_area[0];inext = free_memory; - free_memory = m; - } -} - -/* - * Add timeval times - */ -static void add_time_values(const struct timeval *a, - const struct timeval *b, struct timeval *r) { - r->tv_sec = a->tv_sec + b->tv_sec; - r->tv_usec = a->tv_usec + b->tv_usec; - if (r->tv_usec >= 1000000) { - r->tv_usec -= 1000000; - r->tv_sec++; - } -} - -/* - * Subtract timeval values - */ -static void sub_time_values(const struct timeval *a, - const struct timeval *b, struct timeval *r) { - if ((a->tv_sec < b->tv_sec) || - ((a->tv_sec == b->tv_sec) && (a->tv_usec < b->tv_usec))) { - r->tv_usec = 0; - r->tv_sec = 0; - return; - } - if (a->tv_usec < b->tv_usec) { - r->tv_sec = a->tv_sec - b->tv_sec - 1; - r->tv_usec = a->tv_usec + 1000000 - b->tv_usec; - } else { - r->tv_sec = a->tv_sec - b->tv_sec; - r->tv_usec = a->tv_usec - b->tv_usec; - } -} - -/* - * Compare timeval values - */ -static int cmp_time_values(const struct timeval *a, const struct timeval *b) { - if (a->tv_sec < b->tv_sec) return -1; - if (a->tv_sec > b->tv_sec) return 1; - if (a->tv_usec < b->tv_usec) return -1; - if (a->tv_usec > b->tv_usec) return 1; - return 0; -} - -/* - * Convert TIMEVAL to DOUBLE - */ -static double make_double_time(const struct timeval *tv) { - return 1000000.0 * tv->tv_sec + tv->tv_usec; -} - -/* - * Send Top of queue - */ -static void send_top_of_que (void) { - rtp_delayed_message *m; - int sts; - - if (msg_que) { - m = msg_que; - msg_que = m->next; - m->next = free_memory; - free_memory = m; - - if ((m->errret != NULL) && (m->errret->rtp_tx_sock)) { - sts = sendto(m->socked, &(m->rtp_buff), m->message_len, - m->flags, (const struct sockaddr *)&(m->dst_addr), - (socklen_t)sizeof(m->dst_addr)); - if ((sts == -1) && (m->errret != NULL) && (errno != ECONNREFUSED)) { - osip_call_id_t callid; - - ERROR("sendto() [%s:%i size=%i] delayed call failed: %s", - utils_inet_ntoa(m->errret->remote_ipaddr), - m->errret->remote_port, m->message_len, strerror(errno)); - - /* if sendto() fails with bad filedescriptor, - * this means that the opposite stream has been - * canceled or timed out. - * we should then cancel this stream as well.*/ - - WARN("stopping opposite stream"); - - callid.number=m->errret->callid_number; - callid.host=m->errret->callid_host; - /* don't lock the mutex, as we own the lock */ - if (STS_SUCCESS != rtp_relay_stop_fwd(&callid, - m->errret->direction, - m->errret->media_stream_no, 1)) { - ERROR("fatal error in delayed error close! [%s:%i size=%i]", - utils_inet_ntoa(m->errret->remote_ipaddr), - m->errret->remote_port, m->message_len); - /* brute force protection agains looping errors */ - m->errret->rtp_rx_sock = 0; - rtp_recreate_fdset(); - } /* if stp_stop */ - } /* if sendto fails */ - } - } /* if (msg_que) */ -} - -/* - * Delayed send - */ -static void delayedsendto(int s, const void *msg, size_t len, int flags, - const struct sockaddr_in *to, - const struct timeval *tv, rtp_proxytable_t *errret) { - rtp_delayed_message *m; - rtp_delayed_message **linkin; - - if (!free_memory) send_top_of_que(); - - m = free_memory; - - m->socked = s; - memcpy(&(m->rtp_buff), msg, m->message_len = len); - m->flags = flags; - m->dst_addr = *to; - m->transm_time = *tv; - m->errret = errret; - - free_memory = m->next; - - if (cmp_time_values(¤t_tv,tv) >= 0) { - m->next = msg_que; - msg_que = m; - send_top_of_que(); - } else { - linkin = &msg_que; - while ((*linkin != NULL) && - (cmp_time_values(&((*linkin)->transm_time),tv) < 0)) { - linkin = (rtp_delayed_message **)&((*linkin)->next); - } - m->next = *linkin; - *linkin = m; - } -} - -/* - * Cancel a message - */ -static void cancelmessages(rtp_proxytable_t *dropentry) { - rtp_delayed_message **linkout; - rtp_delayed_message *m; - - linkout = &msg_que; - - while (*linkout != NULL) { - if ((*linkout)->errret == dropentry) { - m = *linkout; - *linkout = m->next; - m->next = free_memory; - free_memory = m; - } else { - linkout = (rtp_delayed_message **)&((*linkout)->next); - } - } -} - -/* - * Flush buffers - */ -static void flushbuffers(void) { - struct timezone tz; - - while (msg_que && - (cmp_time_values(&(msg_que->transm_time),¤t_tv)<=0)) { - send_top_of_que(); - gettimeofday(¤t_tv,&tz); - } -} - -/* - * Delay of next transmission - */ -static int delay_of_next_transmission(struct timeval *tv) { - struct timezone tz ; - - if (msg_que) { - gettimeofday(¤t_tv,&tz); - sub_time_values(&(msg_que->transm_time),¤t_tv,tv); - if (cmp_time_values(tv,&minstep)<=0) { - *tv = minstep ; - } - return -1; - } - return 0; -} - -/* - * Convert DOUBLE time into TIMEVAL - */ -static void split_double_time(double d, struct timeval *tv) { - tv->tv_sec = d / 1000000.0; - tv->tv_usec = d - 1000000.0 * tv->tv_sec; -} - -/* - * Initialize calculation of transmit the frame - */ -static void init_calculate_transmit_time(timecontrol_t *tc, int dejitter) { - struct timezone tz; - - minstep.tv_sec = 0; - minstep.tv_usec = 6000; - memset(tc, 0, sizeof(*tc)); - if (dejitter>0) { - gettimeofday(&(tc->starttime),&tz); - - tc->dejitter = dejitter; - tc->dejitter_d = dejitter; - split_double_time(tc->dejitter_d, &(tc->dejitter_tv)); - } -} - -/* - * - */ -static int fetch_missalign_long_network_oder(char *where) { - int i = 0; - int k; - int j; - - for (j=0;j<4;j++) { - k = *where; - i = (i<<8) | (0xFF & k); - where ++; - } - return i; -} - -/* - * Calculate transmit time - */ -static void calculate_transmit_time(rtp_buff_t *rtp_buff, timecontrol_t *tc, - const struct timeval *input_tv, - struct timeval *ttv) { - int packet_time_code; - double currenttime; - double calculatedtime = 0; - double calculatedtime2 = 0; - struct timeval input_r_tv; - struct timeval output_r_tv; - - if (!tc || !tc->dejitter) { - *ttv = current_tv; - return; - } - - - /* I hate this computer language ... :-/ quite confuse ! look modula */ - packet_time_code = fetch_missalign_long_network_oder(&((*rtp_buff)[4])); - -/*&&&& beware, it seems that when sending RTP events (payload type -telephone-event) the timestamp does not increment and stays the same. -The sequence number however DOES increment. This could lead to confusion when -transmitting RTP events (like DTMF). How can we handle this? Check for RTP event -and then do an "educated guess" for the to-be timestamp? -*/ - if (tc->calccount == 0) { - DEBUGC(DBCLASS_RTP, "initialise time calculatin"); - tc->starttime = *input_tv; - tc->time_code_a = packet_time_code; - } - - sub_time_values(input_tv,&(tc->starttime),&input_r_tv); - - calculatedtime = currenttime = make_double_time(&input_r_tv); - if (tc->calccount < 10) { - DEBUGC(DBCLASS_RTP, "initial data stage 1 %f usec", currenttime); - tc->received_a = currenttime / (packet_time_code - tc->time_code_a); - } else if (tc->calccount < 20) { - tc->received_a = 0.95 * tc->received_a + 0.05 * currenttime / - (packet_time_code - tc->time_code_a); - } else { - tc->received_a = 0.99 * tc->received_a + 0.01 * currenttime / - (packet_time_code - tc->time_code_a); - } - if (tc->calccount > 20) { - if (!tc->time_code_b) { - tc->time_code_b = packet_time_code; - tc->received_b = currenttime; - } else if (tc->time_code_b < packet_time_code) { - calculatedtime = tc->received_b = tc->received_b + - (packet_time_code - tc->time_code_b) * tc->received_a; - tc->time_code_b = packet_time_code; - if (tc->calccount < 28) { - tc->received_b = 0.90 * tc->received_b + 0.1 * currenttime; - } else if (tc->calccount < 300) { - tc->received_b = 0.95 * tc->received_b + 0.05 * currenttime; - } else { - tc->received_b = 0.99 * tc->received_b + 0.01 * currenttime; - } - } else { - calculatedtime = tc->received_b + - (packet_time_code - tc->time_code_b) * tc->received_a; - } - } - tc->received_c = currenttime; - tc->time_code_c = packet_time_code; - - if (tc->calccount < 30) { - /* - * But in the start phase, - * we asume every packet as not delayed. - */ - calculatedtime = currenttime; - } - - /* - ** theoretical value for F1000 Phone - */ - //calculatedtime = (tc->received_a = 125.) * packet_time_code; - - tc->calccount ++; - calculatedtime += tc->dejitter_d; - - if (calculatedtime < currenttime) { - calculatedtime = currenttime; - } else if (calculatedtime > currenttime + 2.* tc->dejitter_d) { - calculatedtime = currenttime + 2.* tc->dejitter_d; - } - - /* every 500 counts show statistics */ - if (tc->calccount % 500 == 0) { - DEBUGC(DBCLASS_RTPBABL, "currenttime = %f", currenttime); - DEBUGC(DBCLASS_RTPBABL, "packetcode = %i", packet_time_code); - DEBUGC(DBCLASS_RTPBABL, "timecodes %i, %i, %i", - tc->time_code_a, tc->time_code_b, tc->time_code_c); - DEBUGC(DBCLASS_RTPBABL, "measuredtimes %f usec, %f usec, %f usec", - tc->received_a, tc->received_b, tc->received_c); - DEBUGC(DBCLASS_RTPBABL, "p2 - p1 = (%i,%f usec)", - tc->time_code_b - tc->time_code_a, - tc->received_b - tc->received_a); - if (tc->time_code_c) { - DEBUGC(DBCLASS_RTPBABL, "p3 - p2 = (%i,%f usec)", - tc->time_code_c - tc->time_code_b, - tc->received_c - tc->received_b); - } - DEBUGC(DBCLASS_RTPBABL, "calculatedtime = %f", calculatedtime); - if (calculatedtime2) { - DEBUGC(DBCLASS_RTPBABL, "calculatedtime2 = %f", calculatedtime2); - } - DEBUGC(DBCLASS_RTPBABL, "transmtime = %f (%f)", calculatedtime / - (160. * tc->received_a) - packet_time_code / 160, - currenttime / (160. * tc->received_a) - - packet_time_code / 160); - DEBUGC(DBCLASS_RTPBABL, "synthetic latency = %f, %f, %f, %i, %i", - calculatedtime-currenttime, calculatedtime, - currenttime, packet_time_code, - packet_time_code / 160); - } - - split_double_time(calculatedtime, &output_r_tv); - add_time_values(&output_r_tv,&(tc->starttime),ttv); -} - /* * match_socket * matches and cross connects two rtp_proxytable entries @@ -1368,6 +985,7 @@ static void match_socket (int rtp_proxytable_idx) { } } + /* * error_handler * @@ -1435,6 +1053,3 @@ static void error_handler (int rtp_proxytable_idx, int socket_type) { } /* if errno != ECONNREFUSED */ } - - - diff --git a/src/siproxd.h b/src/siproxd.h index 76ed866..5f1c03f 100644 --- a/src/siproxd.h +++ b/src/siproxd.h @@ -1,5 +1,5 @@ /* - Copyright (C) 2002-2005 Thomas Ries + Copyright (C) 2002-2007 Thomas Ries This file is part of Siproxd. @@ -208,7 +208,6 @@ int rtp_start_fwd (osip_call_id_t *callid, client_id_t client_id, /*X*/ struct in_addr lcl_client_ipaddr, int lcl_clientport, int isrtp); int rtp_stop_fwd (osip_call_id_t *callid, int direction); /*X*/ -void rtpproxy_kill( void ); /*X*/ /* accessctl.c */ int accesslist_check(struct sockaddr_in from); @@ -321,4 +320,13 @@ int sip_message_set_body(osip_message_t * sip, const char *buf, size_t len); if ((last+(a)) <= now) {last=now; dolog=1;} \ if (dolog) - +/* + * if the following symbol 'GPL' is defined, building siproxd will + * include all features. If not defined, some features that will + * conflict with a non-GPL distribution license will be disabled. + * + * If you wish to distribute siproxd under another license than GPL + * (commercial License for example), contact the author to elaborate + * the details. + */ +#define GPL