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.
sayer/1.4-spce2.6
Raphael Coeffic 16 years ago
parent aafc7ba080
commit 34665040da

@ -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; i<AmConfig::SIPServerThreads;i++){
udp_servers[i] = new udp_trsp(udp_socket);
udp_servers[i]->start();
// Init transport instances
for(unsigned int i=0; i<AmConfig::Ifs.size();i++) {
udp_trsp_socket* udp_socket = new udp_trsp_socket;
if(udp_socket->bind(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; j<AmConfig::SIPServerThreads;j++){
udp_servers[i*AmConfig::SIPServerThreads + j] = new udp_trsp(udp_socket);
udp_servers[i*AmConfig::SIPServerThreads + j]->start();
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; i<AmConfig::SIPServerThreads;i++){
for(int i=0; i<nr_udp_servers;i++){
udp_servers[i]->stop();
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<nr_udp_sockets;i++){
delete udp_sockets[i];
}
delete [] udp_sockets;
udp_sockets = NULL;
nr_udp_sockets = 0;
}
}
int SipCtrlInterface::send(const AmSipReply &rep,

@ -61,8 +61,11 @@ class SipCtrlInterface:
AmCondition<bool> 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();

@ -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");

@ -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");
}

@ -35,6 +35,9 @@
#include <list>
using std::list;
#include <vector>
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<trsp_socket*> transports;
trsp_socket* transport() { return transports[0]; }
/**
* Implements the state changes for the UAC state machine
* @return -1 if errors

Loading…
Cancel
Save