Subversion Repositories HomeAutomation

Rev

Rev 1122 | Rev 1203 | Go to most recent revision | Details | Compare with Previous | Last modification | View Log | SVN | RSS feed

Rev Author Line No. Line
969 runge 1
/***************************************************************************
984 runge 2
 *  Copyright (C) December 6, 2008 by Mattias Runge               *
3
 *  mattias@runge.se                           *
4
 *  asyncsocket.cpp                      *
5
 *                                     *
6
 *  This program is free software; you can redistribute it and/or modify *
7
 *  it under the terms of the GNU General Public License as published by *
8
 *  the Free Software Foundation; either version 2 of the License, or   *
9
 *  (at your option) any later version.                  *
10
 *                                     *
11
 *  This program is distributed in the hope that it will be useful,    *
12
 *  but WITHOUT ANY WARRANTY; without even the implied warranty of    *
13
 *  MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the     *
14
 *  GNU General Public License for more details.             *
15
 *                                     *
16
 *  You should have received a copy of the GNU General Public License   *
17
 *  along with this program; if not, write to the             *
18
 *  Free Software Foundation, Inc.,                    *
19
 *  59 Temple Place - Suite 330, Boston, MA 02111-1307, USA.       *
969 runge 20
 ***************************************************************************/
21
 
999 runge 22
#include "socketeventcallback.h"
1122 linlun 23
#include "../Logger/logger.h"
999 runge 24
 
969 runge 25
#include "asyncsocket.h"
976 runge 26
 
983 runge 27
int socketCount = 1;
969 runge 28
 
29
AsyncSocket::AsyncSocket()
30
{
983 runge 31
    myId = socketCount++;
969 runge 32
    mySocket = -1;
976 runge 33
    myReconnectTimeout = 0;
974 runge 34
    myForceReconnect = false;
999 runge 35
    myEventCallback = NULL;
969 runge 36
}
37
 
38
AsyncSocket::~AsyncSocket()
39
{
40
    stop();
984 runge 41
    silentClose();
969 runge 42
}
43
 
44
void AsyncSocket::run()
45
{
984 runge 46
    // If connection is already up then we should not connect again
981 runge 47
    if (!isConnected())
976 runge 48
    {
984 runge 49
        // If we have choosen to not use automatic reconnect do not start reconnect loop
981 runge 50
        if (myReconnectTimeout == 0)
51
        {
984 runge 52
            // Connect to somewhere
981 runge 53
            connect();
54
        }
55
        else
56
        {
984 runge 57
            // Start the reconnect loop
981 runge 58
            reconnectLoop();
59
        }
976 runge 60
    }
969 runge 61
 
983 runge 62
    bool loop = true;
969 runge 63
    try
64
    {
983 runge 65
        while (loop)
969 runge 66
        {
984 runge 67
            // If we have triggered a forced reconnect do it here
974 runge 68
            if (myForceReconnect)
69
            {
984 runge 70
                myForceReconnect = false;
974 runge 71
                reconnectLoop();
72
            }
73
 
984 runge 74
            // Receive data
75
            loop = receiveData();
76
        }
77
    }
78
    catch (SocketException *e)
79
    {
1124 runge 80
        cout << "DEBUG: socket got an exception: " << e->getDescription() << endl;
984 runge 81
        // Something bad happend and we can not continue
82
        eventAdd(SocketEvent::TYPE_CONNECTION_DIED, e->getDescription());
83
    }
969 runge 84
 
984 runge 85
    // Clean up socket if we would want to restart
86
    silentClose();
87
}
969 runge 88
 
984 runge 89
bool AsyncSocket::receiveData()
90
{
91
    char buffer[MAXBUFFER + 1];
92
    memset(buffer, 0, MAXBUFFER + 1);
969 runge 93
 
1124 runge 94
    //cout << "recv start" << endl;
984 runge 95
    int status = ::recv(mySocket, buffer, MAXBUFFER, 0);
1124 runge 96
    //cout << "recv end" << endl;
97
 
984 runge 98
    if (status == -1)
99
    {
100
        switch (errno)
101
        {
102
            case EAGAIN:
103
            throw new SocketException("The socket is marked non-blocking and the receive operation would block, or a receive timeout had been set and the timeout expired before data was received.");
969 runge 104
 
984 runge 105
            case EBADF:
106
            throw new SocketException("The argument s is an invalid descriptor.");
969 runge 107
 
984 runge 108
            case ECONNREFUSED:
109
            throw new SocketException("A remote host refused to allow the network connection (typically because it is not running the requested service).");
969 runge 110
 
984 runge 111
            case EFAULT:
112
            throw new SocketException("The receive buffer pointer(s) point outside the process's address space.");
969 runge 113
 
984 runge 114
            case EINTR:
115
            throw new SocketException("The receive was interrupted by delivery of a signal before any data were available; see signal(7).");
969 runge 116
 
984 runge 117
            case EINVAL:
118
            throw new SocketException("Invalid argument passed.");
969 runge 119
 
984 runge 120
            case ENOMEM:
121
            throw new SocketException("Could not allocate memory for recvmsg().");
969 runge 122
 
984 runge 123
            case ENOTCONN:
124
            throw new SocketException("The socket is associated with a connection-oriented protocol and has not been connected (see connect(2) and accept(2)).");
125
 
126
            case ENOTSOCK:
127
            throw new SocketException("The argument s does not refer to a socket.");
128
 
129
            case ECONNRESET:
130
            eventAdd(SocketEvent::TYPE_CONNECTION_RESET);
131
 
132
            // If we have automatic reconnect we want to start it now
133
            if (myReconnectTimeout == 0)
134
            {
135
                reconnectLoop();
983 runge 136
            }
984 runge 137
            else
983 runge 138
            {
984 runge 139
                // Otherwise we would like to end the main loop
140
                return false;
983 runge 141
            }
984 runge 142
            break;
976 runge 143
 
984 runge 144
            default:
145
            throw new SocketException("Unknow exception: " + itos(errno));
146
            break;
147
        }
148
    }
149
    else if (status == 0)
150
    {
151
        // Remote host have done a normal shutdown
152
        eventAdd(SocketEvent::TYPE_CONNECTION_CLOSED);
969 runge 153
 
984 runge 154
        if (myReconnectTimeout > 0)
155
        {
156
            reconnectLoop();
969 runge 157
        }
984 runge 158
        else
159
        {
160
            return false;
161
        }
969 runge 162
    }
984 runge 163
    else if (status > 0)
969 runge 164
    {
984 runge 165
        // We have received data
166
        eventAdd(SocketEvent::TYPE_DATA, buffer);
969 runge 167
    }
168
 
984 runge 169
    return true;
969 runge 170
}
171
 
984 runge 172
void AsyncSocket::silentClose()
173
{
174
    if (mySocket != -1)
175
    {
176
        ::close(mySocket);
177
        mySocket = -1;
178
    }
179
}
180
 
181
void AsyncSocket::close()
182
{
183
    silentClose();
184
    eventAdd(SocketEvent::TYPE_CONNECTION_CLOSED);
185
}
186
 
969 runge 187
void AsyncSocket::reconnectLoop()
188
{
984 runge 189
    while (true)
969 runge 190
    {
191
        try
192
        {
193
            connect();
984 runge 194
            return;
969 runge 195
        }
196
        catch (SocketException *e)
197
        {
984 runge 198
            eventAdd(SocketEvent::TYPE_CONNECTION_FAILED, e->getDescription());
199
            eventAdd(SocketEvent::TYPE_WAITING_RECONNECT);
969 runge 200
            sleep(myReconnectTimeout);
201
        }
202
    }
203
}
204
 
981 runge 205
void AsyncSocket::create()
969 runge 206
{
984 runge 207
    silentClose();
976 runge 208
 
969 runge 209
    mySocket = ::socket(AF_INET, SOCK_STREAM, 0);
210
 
211
    int on = 1;
212
    int status = setsockopt(mySocket, SOL_SOCKET, SO_REUSEADDR, (const char*)&on, sizeof(on));
213
    if (status == -1)
214
    {
984 runge 215
        silentClose();
981 runge 216
        throw new SocketException("Create:Reuseaddress: " + itos(errno));
969 runge 217
    }
981 runge 218
}
969 runge 219
 
981 runge 220
void AsyncSocket::startListen()
221
{
222
    create();
223
 
224
    myAddressStruct.sin_family = AF_INET;
225
    myAddressStruct.sin_addr.s_addr = INADDR_ANY;
226
    myAddressStruct.sin_port = htons(myPort);
227
 
228
    int status = ::bind(mySocket, (struct sockaddr *)&myAddressStruct, sizeof(myAddressStruct));
229
 
230
    if (status == -1)
231
    {
984 runge 232
        silentClose();
981 runge 233
        switch (errno)
234
        {
235
            case EACCES:
236
            throw new SocketException("The address is protected, and the user is not the superuser.");
237
 
238
            case EADDRINUSE:
239
            throw new SocketException("The given address is already in use.");
240
 
241
            case EBADF:
242
            throw new SocketException("sockfd is not a valid descriptor.");
243
 
244
            case EINVAL:
245
            throw new SocketException("The socket is already bound to an address.");
246
 
247
            case ENOTSOCK:
248
            throw new SocketException("sockfd is a descriptor for a file, not a socket.");
249
 
250
            //case EACCES:
251
            //throw new SocketException("Search permission is denied on a component of the path prefix. (See also path_resolution(7).)");
252
 
253
            case EADDRNOTAVAIL:
254
            throw new SocketException("A nonexistent interface was requested or the requested address was not local.");
255
 
256
            case EFAULT:
257
            throw new SocketException("addr points outside the user's accessible address space.");
258
 
259
            //case EINVAL:
260
            //throw new SocketException("The addrlen is wrong, or the socket was not in the AF_UNIX family.");
261
 
262
            case ELOOP:
263
            throw new SocketException("Too many symbolic links were encountered in resolving addr.");
264
 
265
            case ENAMETOOLONG:
266
            throw new SocketException("addr is too long.");
267
 
268
            case ENOENT:
269
            throw new SocketException("The file does not exist.");
270
 
271
            case ENOMEM:
272
            throw new SocketException("Insufficient kernel memory was available.");
273
 
274
            case ENOTDIR:
275
            throw new SocketException("A component of the path prefix is not a directory.");
276
 
277
            case EROFS:
278
            throw new SocketException("The socket inode would reside on a read-only file system.");
279
 
280
            default:
281
            throw new SocketException("Unknow exception: " + itos(errno));
282
            break;
283
        }
284
    }
285
 
286
    status = ::listen(mySocket, MAXCONNECTIONS);
287
 
288
    if (status == -1)
289
    {
984 runge 290
        silentClose();
981 runge 291
        switch (errno)
292
        {
293
            case EADDRINUSE:
294
            throw new SocketException("Another socket is already listening on the same port.");
295
 
296
            case EBADF:
297
            throw new SocketException("The argument sockfd is not a valid descriptor.");
298
 
299
            case ENOTSOCK:
300
            throw new SocketException("The argument sockfd is not a socket.");
301
 
302
            case EOPNOTSUPP:
303
            throw new SocketException("The socket is not of a type that supports the listen() operation.");
304
 
305
            default:
306
            throw new SocketException("Unknow exception: " + itos(errno));
307
            break;
308
        }
309
    }
310
}
311
 
312
bool AsyncSocket::accept(AsyncSocket* newSocket)
313
{
314
    int addr_length = sizeof(myAddressStruct);
315
    int socket = ::accept(mySocket, (sockaddr*)&myAddressStruct, (socklen_t*)&addr_length);
316
 
317
    if (socket > 0)
318
    {
319
        newSocket->setSocket(socket);
320
        return true;
321
    }
322
 
323
    return false;
324
}
325
 
326
void AsyncSocket::connect()
327
{
984 runge 328
    eventAdd(SocketEvent::TYPE_CONNECTING);
981 runge 329
 
330
    create();
331
 
969 runge 332
    memset(&myAddressStruct, 0, sizeof(myAddressStruct));
333
 
334
    myAddressStruct.sin_family = AF_INET;
335
    myAddressStruct.sin_port = htons(myPort);
336
 
976 runge 337
    struct hostent *hptr = gethostbyname(myAddress.c_str());
338
    if (hptr == NULL)
339
    {
984 runge 340
        silentClose();
976 runge 341
        throw new SocketException("Connect: Could not resolv ip address");
342
    }
343
 
344
    memcpy(&myAddressStruct.sin_addr, hptr->h_addr, hptr->h_length);
969 runge 345
 
983 runge 346
    int status = ::connect(mySocket, (sockaddr*)&myAddressStruct, sizeof(myAddressStruct));
976 runge 347
 
969 runge 348
    if (status == -1)
349
    {
984 runge 350
        silentClose();
969 runge 351
        switch (errno)
352
        {
353
            case EACCES:
984 runge 354
            throw new SocketException("For Unix domain sockets, which are identified by pathname: Write permission is denied on the socket file, or search per- mission is denied for one of the directories in the path prefix. (See also path_resolution(7).)");
969 runge 355
 
356
            case EPERM:
984 runge 357
            throw new SocketException("The user tried to connect to a broadcast address without having the socket broadcast flag enabled or the connection request failed because of a local firewall rule.");
969 runge 358
 
359
            case EADDRINUSE:
360
            throw new SocketException("Local address is already in use.");
361
 
362
            case EAFNOSUPPORT:
363
            throw new SocketException("The passed address didn't have the correct address family in its sa_family field.");
364
 
365
            case EAGAIN:
984 runge 366
            throw new SocketException("No more free local ports or insufficient entries in the routing cache. For AF_INET see the net.ipv4.ip_local_port_range sysctl in ip(7) on how to increase the number of local ports.");
969 runge 367
 
368
            case EALREADY:
369
            throw new SocketException("The socket is non-blocking and a previous connection attempt has not yet been completed.");
370
 
371
            case EBADF:
372
            throw new SocketException("The file descriptor is not a valid index in the descriptor table.");
373
 
374
            case ECONNREFUSED:
375
            throw new SocketException("No-one listening on the remote address.");
376
 
377
            case EFAULT:
378
            throw new SocketException("The socket structure address is outside the user's address space.");
379
 
380
            case EINPROGRESS:
984 runge 381
            throw new SocketException("The socket is non-blocking and the connection cannot be completed immediately. It is possible to select(2) or poll(2) for completion by selecting the socket for writing. After select(2) indicates writability, use getsockopt(2) to read the SO_ERROR option at level SOL_SOCKET to determine whether connect() completed successfully (SO_ERROR is zero) or unsuccessfully (SO_ERROR is one of the usual error codes listed here, explaining the reason for the failure).");
969 runge 382
 
383
            case EINTR:
384
            throw new SocketException("The system call was interrupted by a signal that was caught; see signal(7).");
385
 
386
            case EISCONN:
387
            throw new SocketException("The socket is already connected.");
388
 
389
            case ENETUNREACH:
390
            throw new SocketException("Network is unreachable.");
391
 
392
            case ENOTSOCK:
393
            throw new SocketException("The file descriptor is not associated with a socket.");
394
 
395
            case ETIMEDOUT:
984 runge 396
            throw new SocketException("Timeout while attempting connection. The server may be too busy to accept new connections. Note that for IP sockets the timeout may be very long when syncookies are enabled on the server.");
969 runge 397
 
398
            default:
399
            throw new SocketException("Other errors may be generated by the underlying protocol modules. : " + itos(errno));
400
        }
401
    }
402
 
984 runge 403
    eventAdd(SocketEvent::TYPE_CONNECTED);
969 runge 404
}
405
 
984 runge 406
void AsyncSocket::sendData(string data)
969 runge 407
{
1033 runge 408
    mySendMutex.lock();
409
 
1124 runge 410
    int status = ::send(mySocket, data.c_str(), data.size(), MSG_NOSIGNAL);
1122 linlun 411
//Logger::getInstance().add("Sent: \"" + data + "\" status was " + itos(status) + "\n");
1033 runge 412
    mySendMutex.unlock();
413
 
1036 runge 414
    //Logger &log = Logger::getInstance();
415
    //log.add("Sent: \"" + data + "\" status was " + itos(status) + "\n");
1033 runge 416
 
984 runge 417
    if (status == -1)
981 runge 418
    {
984 runge 419
        switch (errno)
981 runge 420
        {
984 runge 421
            case EACCES:
422
            throw new SocketException("(For Unix domain sockets, which are identified by pathname) Write permission is denied on the destination socket file, or search permission is denied for one of the directories the path prefix. (See path_resolution(7).)");
981 runge 423
 
984 runge 424
            case EAGAIN:
425
            throw new SocketException("The socket is marked non-blocking and the requested operation would block.");
981 runge 426
 
984 runge 427
            case EBADF:
428
            throw new SocketException("An invalid descriptor was specified.");
981 runge 429
 
984 runge 430
            case ECONNRESET:
431
            throw new SocketException("Connection reset by peer.");
981 runge 432
 
984 runge 433
            case EDESTADDRREQ:
434
            throw new SocketException("The socket is not connection-mode, and no peer address is set.");
981 runge 435
 
984 runge 436
            case EFAULT:
437
            throw new SocketException("An invalid user space address was specified for an argument.");
981 runge 438
 
984 runge 439
            case EINTR:
440
            throw new SocketException("A signal occurred before any data was transmitted; see signal(7).");
981 runge 441
 
984 runge 442
            case EINVAL:
443
            throw new SocketException("1Invalid argument passed.");
981 runge 444
 
984 runge 445
            case EISCONN:
446
            throw new SocketException("The connection-mode socket was connected already but a recipient was specified. (Now either this error is returned, or the recipient specification is ignored.)");
981 runge 447
 
984 runge 448
            case EMSGSIZE:
449
            throw new SocketException("The socket type requires that message be sent atomically, and the size of the message to be sent made this impossible.");
981 runge 450
 
984 runge 451
            case ENOBUFS:
452
            throw new SocketException("The output queue for a network interface was full. This generally indicates that the interface has stopped sending, but may be caused by transient congestion. (Normally, this does not occur in Linux. Packets are just silently dropped when a device queue overflows.)");
981 runge 453
 
984 runge 454
            case ENOMEM:
455
            throw new SocketException("No memory available.");
981 runge 456
 
984 runge 457
            case ENOTCONN:
458
            throw new SocketException("The socket is not connected, and no target has been given.");
981 runge 459
 
984 runge 460
            case ENOTSOCK:
461
            throw new SocketException("The argument s is not a socket.");
981 runge 462
 
984 runge 463
            case EOPNOTSUPP:
464
            throw new SocketException("Some bit in the flags argument is inappropriate for the socket type.");
981 runge 465
 
984 runge 466
            case EPIPE:
1124 runge 467
            close();
468
            stop();
469
            //throw new SocketException("The local end has been shut down on a connection oriented socket. In this case the process will also receive a SIGPIPE unless MSG_NOSIGNAL is set.");
470
            break;
981 runge 471
 
984 runge 472
            default:
473
            throw new SocketException("Unknow exception: " + itos(errno));
981 runge 474
        }
475
    }
476
}
477
 
984 runge 478
void AsyncSocket::eventAdd(unsigned int eventType)
969 runge 479
{
984 runge 480
    eventAdd(eventType, "");
969 runge 481
}
482
 
984 runge 483
void AsyncSocket::eventAdd(unsigned int eventType, string eventData)
974 runge 484
{
984 runge 485
    SocketEvent socketEvent(eventType, eventData);
999 runge 486
 
487
    if (myEventCallback == NULL)
488
    {
489
        myEventQueue.push(socketEvent);
490
        myEventSemaphore.broadcast();
491
    }
492
    else
493
    {
494
        myEventCallback->handleEvent(myId, socketEvent);
495
    }
974 runge 496
}
999 runge 497
 
498
void AsyncSocket::eventSetCallback(SocketEventCallback* eventCallback)
499
{
500
    myEventCallback = eventCallback;
501
}