From 34665040da316e8c44e9b527da334f3ccd601fbd Mon Sep 17 00:00:00 2001 From: Raphael Coeffic Date: Wed, 2 Feb 2011 17:11:09 +0100 Subject: [PATCH] Wip: start multiple SIP/UDP servers on multiple interfaces. This only starts the UDP servers: proper interface selection for sent messages / RTP is still missing. --- core/SipCtrlInterface.cpp | 66 +++++++++++++++++++++++++++++---------- core/SipCtrlInterface.h | 9 ++++-- core/sems.cpp | 4 +-- core/sip/trans_layer.cpp | 46 ++++++++++++++------------- core/sip/trans_layer.h | 16 ++++++++-- 5 files changed, 96 insertions(+), 45 deletions(-) diff --git a/core/SipCtrlInterface.cpp b/core/SipCtrlInterface.cpp index f49023b6..6aa70977 100644 --- a/core/SipCtrlInterface.cpp +++ b/core/SipCtrlInterface.cpp @@ -119,7 +119,8 @@ int SipCtrlInterface::load() } SipCtrlInterface::SipCtrlInterface() - : stopped(false), udp_servers(NULL), udp_socket(NULL) + : stopped(false), udp_servers(NULL), udp_sockets(NULL), + nr_udp_sockets(0), nr_udp_servers(0) { trans_layer::instance()->register_ua(this); } @@ -229,28 +230,50 @@ int SipCtrlInterface::send(AmSipRequest &req, return res; } -void SipCtrlInterface::run(const string& bind_addr, unsigned short bind_port) +int SipCtrlInterface::run() { - DBG("Starting SIP control interface\n"); + //AmConfig::LocalSIPIP(), AmConfig::LocalSIPPort() - udp_socket = new udp_trsp_socket; - udp_socket->bind(bind_addr,bind_port); + DBG("Starting SIP control interface\n"); - trans_layer::instance()->register_transport(udp_socket); - udp_servers = new udp_trsp*[AmConfig::SIPServerThreads]; + udp_sockets = new udp_trsp_socket*[AmConfig::Ifs.size()]; + udp_servers = new udp_trsp*[AmConfig::SIPServerThreads * AmConfig::Ifs.size()]; wheeltimer::instance()->start(); - for(int i=0; istart(); + // Init transport instances + for(unsigned int i=0; ibind(AmConfig::Ifs[i].LocalSIPIP, + AmConfig::Ifs[i].LocalSIPPort) < 0){ + + ERROR("Could not bind SIP/UDP socket to %s:%i", + AmConfig::Ifs[i].LocalSIPIP.c_str(), + AmConfig::Ifs[i].LocalSIPPort); + + delete udp_socket; + return -1; + } + + trans_layer::instance()->register_transport(udp_socket); + udp_sockets[i] = udp_socket; + nr_udp_sockets++; + + for(int j=0; jstart(); + nr_udp_servers++; + } } - + while (!stopped.get()) { stopped.wait_for(); } DBG("SIP control interface ending\n"); + return 0; } void SipCtrlInterface::stop() @@ -263,17 +286,28 @@ void SipCtrlInterface::cleanup() DBG("Stopping SIP control interface threads\n"); if (NULL != udp_servers) { - for(int i=0; istop(); udp_servers[i]->join(); delete udp_servers[i]; } + + delete [] udp_servers; + udp_servers = NULL; + nr_udp_servers = 0; } - trans_layer::instance()->register_transport(NULL); - - delete [] udp_servers; - delete udp_socket; + trans_layer::instance()->clear_transports(); + + if (NULL != udp_sockets) { + for(int i=0; i stopped; - udp_trsp_socket* udp_socket; - udp_trsp** udp_servers; + unsigned short nr_udp_sockets; + udp_trsp_socket** udp_sockets; + + unsigned short nr_udp_servers; + udp_trsp** udp_servers; public: @@ -77,7 +80,7 @@ public: int load(); - void run(const string& bind_addr, unsigned short bind_port); + int run(); void stop(); void cleanup(); diff --git a/core/sems.cpp b/core/sems.cpp index 3e3e05c2..9cdc9bbe 100644 --- a/core/sems.cpp +++ b/core/sems.cpp @@ -628,9 +628,9 @@ int main(int argc, char* argv[]) INFO("Starting SIP stack (control interface)\n"); sip_ctrl.load(); - sip_ctrl.run(AmConfig::LocalSIPIP(), AmConfig::LocalSIPPort()); - success = true; + if(sip_ctrl.run() != -1) + success = true; // session container stops active sessions INFO("Disposing session container\n"); diff --git a/core/sip/trans_layer.cpp b/core/sip/trans_layer.cpp index a27e76a1..a3f8094d 100644 --- a/core/sip/trans_layer.cpp +++ b/core/sip/trans_layer.cpp @@ -61,7 +61,7 @@ bool _trans_layer::accept_fr_without_totag = false; _trans_layer::_trans_layer() : ua(NULL), - transport(NULL) + transports() { } @@ -76,9 +76,13 @@ void _trans_layer::register_ua(sip_ua* ua) void _trans_layer::register_transport(trsp_socket* trsp) { - transport = trsp; + transports.push_back(trsp); } +void _trans_layer::clear_transports() +{ + transports.clear(); +} int _trans_layer::send_reply(trans_ticket* tt, @@ -327,7 +331,7 @@ int _trans_layer::send_reply(trans_ticket* tt, memcpy(c,body.s,body.len); } - assert(transport); + assert(transport()); int err = -1; @@ -372,7 +376,7 @@ int _trans_layer::send_reply(trans_ticket* tt, ntohs(((sockaddr_in*)&remote_ip)->sin_port), 50 /* preview - instead of p_msg->len */,reply_buf); - err = transport->send(&remote_ip,reply_buf,reply_len); + err = transport()->send(&remote_ip,reply_buf,reply_len); if(err < 0){ delete [] reply_buf; goto end; @@ -541,9 +545,9 @@ int _trans_layer::send_sl_reply(sip_msg* req, int reply_code, memcpy(c,body.s,body.len); } - assert(transport); + assert(transport()); - int err = transport->send(&req->remote_ip,reply_buf,reply_len); + int err = transport()->send(&req->remote_ip,reply_buf,reply_len); delete [] reply_buf; return err; @@ -808,7 +812,7 @@ int _trans_layer::send_request(sip_msg* msg, trans_ticket* tt, // Supported / Require // Content-Length / Content-Type - assert(transport); + assert(transport()); assert(msg); assert(tt); @@ -835,9 +839,9 @@ int _trans_layer::send_request(sip_msg* msg, trans_ticket* tt, compute_branch(branch_buf,msg->callid->value,msg->cseq->value); cstring branch(branch_buf,BRANCH_BUF_LEN); - string via(transport->get_ip()); - if(transport->get_port() != 5060) - via += ":" + int2str(transport->get_port()); + string via(transport()->get_ip()); + if(transport()->get_port() != 5060) + via += ":" + int2str(transport()->get_port()); // add 'rport' parameter defaultwise? yes, for now request_len += via_len(stl2cstr(via),branch,true); @@ -913,7 +917,7 @@ int _trans_layer::send_request(sip_msg* msg, trans_ticket* tt, get_cseq(p_msg)->num_str); tt->_bucket->lock(); - int send_err = transport->send(&p_msg->remote_ip,p_msg->buf,p_msg->len); + int send_err = transport()->send(&p_msg->remote_ip,p_msg->buf,p_msg->len); if(send_err < 0){ ERROR("Error from transport layer\n"); delete p_msg; @@ -985,9 +989,9 @@ int _trans_layer::cancel(trans_ticket* tt) compute_branch(branch_buf,req->callid->value,get_cseq(req)->num_str); cstring branch(branch_buf,BRANCH_BUF_LEN); - string via(transport->get_ip()); - if(transport->get_port() != 5060) - via += ":" + int2str(transport->get_port()); + string via(transport()->get_ip()); + if(transport()->get_port() != 5060) + via += ":" + int2str(transport()->get_port()); //TODO: add 'rport' parameter by default? @@ -1045,7 +1049,7 @@ int _trans_layer::cancel(trans_ticket* tt) if(bucket != n_bucket) n_bucket->lock(); - int send_err = transport->send(&p_msg->remote_ip,p_msg->buf,p_msg->len); + int send_err = transport()->send(&p_msg->remote_ip,p_msg->buf,p_msg->len); if(send_err < 0){ ERROR("Error from transport layer\n"); delete p_msg; @@ -1619,8 +1623,8 @@ void _trans_layer::send_non_200_ack(sip_msg* reply, sip_trans* t) DBG("About to send ACK\n"); - assert(transport); - int send_err = transport->send(&inv->remote_ip,ack_buf,ack_len); + assert(transport()); + int send_err = transport()->send(&inv->remote_ip,ack_buf,ack_len); if(send_err < 0){ ERROR("Error from transport layer\n"); delete ack_buf; @@ -1635,13 +1639,13 @@ void _trans_layer::send_non_200_ack(sip_msg* reply, sip_trans* t) void _trans_layer::retransmit(sip_trans* t) { - assert(transport); + assert(transport()); if(!t->retr_buf || !t->retr_len){ // there is nothing to re-transmit yet!!! return; } - int send_err = transport->send(&t->retr_addr,t->retr_buf,t->retr_len); + int send_err = transport()->send(&t->retr_addr,t->retr_buf,t->retr_len); if(send_err < 0){ ERROR("Error from transport layer\n"); } @@ -1649,8 +1653,8 @@ void _trans_layer::retransmit(sip_trans* t) void _trans_layer::retransmit(sip_msg* msg) { - assert(transport); - int send_err = transport->send(&msg->remote_ip,msg->buf,msg->len); + assert(transport()); + int send_err = transport()->send(&msg->remote_ip,msg->buf,msg->len); if(send_err < 0){ ERROR("Error from transport layer\n"); } diff --git a/core/sip/trans_layer.h b/core/sip/trans_layer.h index 9308767b..2ebc2a1c 100644 --- a/core/sip/trans_layer.h +++ b/core/sip/trans_layer.h @@ -35,6 +35,9 @@ #include using std::list; +#include +using std::vector; + struct sip_msg; struct sip_uri; struct sip_trans; @@ -70,10 +73,15 @@ public: /** * Register a transport instance. - * This method MUST be called ONCE. + * This method MUST be called at least once. */ void register_transport(trsp_socket* trsp); + /** + * Clears all registered transport instances. + */ + void clear_transports(); + /** * Sends a UAS reply. * If a body is included, the hdrs parameter should @@ -116,9 +124,11 @@ public: */ void timer_expired(timer* t, trans_bucket* bucket, sip_trans* tr); - sip_ua* ua; - trsp_socket* transport; + sip_ua* ua; + vector transports; + trsp_socket* transport() { return transports[0]; } + /** * Implements the state changes for the UAC state machine * @return -1 if errors