Rev 1314 | Details | Compare with Previous | Last modification | View Log | SVN | RSS feed
| Rev | Author | Line No. | Line |
|---|---|---|---|
| 1309 | runge | 1 | /* |
| 2 | * UdpServer.cpp |
||
| 3 | * |
||
| 4 | * Created on: Jul 20, 2009 |
||
| 5 | * Author: Mattias Runge |
||
| 6 | */ |
||
| 7 | |||
| 8 | #include "UdpServer.h" |
||
| 9 | |||
| 10 | namespace atom { |
||
| 1317 | runge | 11 | namespace net { |
| 1309 | runge | 12 | |
| 13 | UdpServer::UdpServer(unsigned int port) : |
||
| 14 | mySocket(this->myIoService, udp::endpoint(udp::v4(), port)) |
||
| 15 | { |
||
| 16 | LOG.setName("UdpServer"); |
||
| 17 | |||
| 18 | this->receiveFrom(); |
||
| 19 | |||
| 20 | this->start(); |
||
| 21 | } |
||
| 22 | |||
| 23 | UdpServer::~UdpServer() |
||
| 24 | { |
||
| 25 | this->stop(); |
||
| 26 | this->mySocket.cancel(); |
||
| 27 | this->mySocket.close(); |
||
| 28 | } |
||
| 29 | |||
| 30 | UdpServer::pointer UdpServer::getInstance(unsigned int port) |
||
| 31 | { |
||
| 32 | static map<unsigned int, UdpServer::pointer> servers; // FIXME Make thread-safe |
||
| 33 | |||
| 34 | map<unsigned int, UdpServer::pointer>::iterator iter = servers.find(port); |
||
| 35 | |||
| 36 | if (iter != servers.end()) |
||
| 37 | { |
||
| 38 | return iter->second; |
||
| 39 | } |
||
| 40 | |||
| 41 | UdpServer::pointer udpServer(new UdpServer(port)); |
||
| 42 | |||
| 1310 | runge | 43 | servers[port] = udpServer; |
| 1309 | runge | 44 | |
| 45 | return udpServer; |
||
| 46 | } |
||
| 47 | |||
| 48 | udp::endpoint UdpServer::getEndpoint(string address) |
||
| 49 | { |
||
| 1314 | runge | 50 | LOG.info("Resolving " + address + ":" + convert::int2string(this->getPort()) + "..."); |
| 1309 | runge | 51 | udp::resolver resolver(this->myIoService); |
| 1314 | runge | 52 | udp::resolver::query query(udp::v4(), address, convert::int2string(this->getPort())); |
| 1309 | runge | 53 | return *resolver.resolve(query); |
| 54 | } |
||
| 55 | |||
| 56 | unsigned int UdpServer::getPort() |
||
| 57 | { |
||
| 58 | return this->mySocket.local_endpoint().port(); |
||
| 59 | } |
||
| 60 | |||
| 61 | void UdpServer::sendTo(udp::endpoint receiver, byte_list bytes) |
||
| 62 | { |
||
| 63 | try |
||
| 64 | { |
||
| 1314 | runge | 65 | LOG.info("Sending packet, length " + convert::int2string(bytes.size()) + ", to " + receiver.address().to_string() + ":" + convert::int2string(this->getPort()) + "..."); |
| 1309 | runge | 66 | this->mySocket.send_to(boost::asio::buffer(bytes), receiver); |
| 67 | } |
||
| 68 | catch (std::exception& e) |
||
| 69 | { |
||
| 70 | LOG.error(e.what()); |
||
| 71 | } |
||
| 72 | } |
||
| 73 | |||
| 74 | void UdpServer::receiveFrom() |
||
| 75 | { |
||
| 1314 | runge | 76 | LOG.info("Listening for incoming data on port " + convert::int2string(this->getPort()) + "..."); |
| 1309 | runge | 77 | this->mySocket.async_receive_from(boost::asio::buffer(this->myDataBuffer, MAX_LENGTH), this->mySender, boost::bind(&UdpServer::handleReceiveFrom, this, boost::asio::placeholders::error, boost::asio::placeholders::bytes_transferred)); |
| 78 | } |
||
| 79 | |||
| 80 | void UdpServer::handleReceiveFrom(const boost::system::error_code& error, size_t bytesReceived) |
||
| 81 | { |
||
| 82 | if (error) |
||
| 83 | { |
||
| 1314 | runge | 84 | LOG.error("handleReceiveFrom got the error code " + convert::int2string(error.value()) + " and the message " + error.message()); |
| 1309 | runge | 85 | } |
| 86 | else if (bytesReceived == 0) |
||
| 87 | { |
||
| 88 | LOG.info("handleReceiveFrom got 0 bytes."); |
||
| 89 | } |
||
| 90 | else |
||
| 91 | { |
||
| 1314 | runge | 92 | LOG.info("Got " + convert::int2string(bytesReceived) + " bytes from " + this->mySender.address().to_string()); |
| 1309 | runge | 93 | |
| 94 | byte_list bytes(this->myDataBuffer, this->myDataBuffer + bytesReceived); |
||
| 95 | |||
| 96 | this->onNewData(this->mySender, bytes); |
||
| 97 | } |
||
| 98 | |||
| 99 | this->receiveFrom(); |
||
| 100 | } |
||
| 101 | |||
| 102 | void UdpServer::connect(const OnDataSignal::slot_type & slot) |
||
| 103 | { |
||
| 104 | LOG.info("New listener connected."); |
||
| 105 | this->onNewData.connect(slot); |
||
| 106 | } |
||
| 107 | |||
| 108 | void UdpServer::run() |
||
| 109 | { |
||
| 110 | LOG.info("Starting thread..."); |
||
| 111 | this->myIoService.run(); |
||
| 112 | LOG.info("Thread completed."); |
||
| 113 | } |
||
| 114 | |||
| 115 | } |
||
| 116 | } |