Subversion Repositories HomeAutomation

Rev

Go to most recent revision | Details | 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
 
22
#include "asyncsocket.h"
976 runge 23
 
983 runge 24
int socketCount = 1;
969 runge 25
 
26
AsyncSocket::AsyncSocket()
27
{
983 runge 28
    myId = socketCount++;
969 runge 29
    mySocket = -1;
976 runge 30
    myReconnectTimeout = 0;
974 runge 31
    myForceReconnect = false;
969 runge 32
}
33
 
34
AsyncSocket::~AsyncSocket()
35
{
36
    stop();
984 runge 37
    silentClose();
969 runge 38
}
39
 
40
void AsyncSocket::run()
41
{
984 runge 42
    // If connection is already up then we should not connect again
981 runge 43
    if (!isConnected())
976 runge 44
    {
984 runge 45
        // If we have choosen to not use automatic reconnect do not start reconnect loop
981 runge 46
        if (myReconnectTimeout == 0)
47
        {
984 runge 48
            // Connect to somewhere
981 runge 49
            connect();
50
        }
51
        else
52
        {
984 runge 53
            // Start the reconnect loop
981 runge 54
            reconnectLoop();
55
        }
976 runge 56
    }
969 runge 57
 
983 runge 58
    bool loop = true;
969 runge 59
    try
60
    {
983 runge 61
        while (loop)
969 runge 62
        {
984 runge 63
            // If we have triggered a forced reconnect do it here
974 runge 64
            if (myForceReconnect)
65
            {
984 runge 66
                myForceReconnect = false;
974 runge 67
                reconnectLoop();
68
            }
69
 
984 runge 70
            // Receive data
71
            loop = receiveData();
72
        }
73
    }
74
    catch (SocketException *e)
75
    {
76
        // Something bad happend and we can not continue
77
        eventAdd(SocketEvent::TYPE_CONNECTION_DIED, e->getDescription());
78
    }
969 runge 79
 
984 runge 80
    // Clean up socket if we would want to restart
81
    silentClose();
82
}
969 runge 83
 
984 runge 84
bool AsyncSocket::receiveData()
85
{
86
    char buffer[MAXBUFFER + 1];
87
    memset(buffer, 0, MAXBUFFER + 1);
969 runge 88
 
984 runge 89
    int status = ::recv(mySocket, buffer, MAXBUFFER, 0);
969 runge 90
 
984 runge 91
    if (status == -1)
92
    {
93
        switch (errno)
94
        {
95
            case EAGAIN:
96
            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 97
 
984 runge 98
            case EBADF:
99
            throw new SocketException("The argument s is an invalid descriptor.");
969 runge 100
 
984 runge 101
            case ECONNREFUSED:
102
            throw new SocketException("A remote host refused to allow the network connection (typically because it is not running the requested service).");
969 runge 103
 
984 runge 104
            case EFAULT:
105
            throw new SocketException("The receive buffer pointer(s) point outside the process's address space.");
969 runge 106
 
984 runge 107
            case EINTR:
108
            throw new SocketException("The receive was interrupted by delivery of a signal before any data were available; see signal(7).");
969 runge 109
 
984 runge 110
            case EINVAL:
111
            throw new SocketException("Invalid argument passed.");
969 runge 112
 
984 runge 113
            case ENOMEM:
114
            throw new SocketException("Could not allocate memory for recvmsg().");
969 runge 115
 
984 runge 116
            case ENOTCONN:
117
            throw new SocketException("The socket is associated with a connection-oriented protocol and has not been connected (see connect(2) and accept(2)).");
118
 
119
            case ENOTSOCK:
120
            throw new SocketException("The argument s does not refer to a socket.");
121
 
122
            case ECONNRESET:
123
            eventAdd(SocketEvent::TYPE_CONNECTION_RESET);
124
 
125
            // If we have automatic reconnect we want to start it now
126
            if (myReconnectTimeout == 0)
127
            {
128
                reconnectLoop();
983 runge 129
            }
984 runge 130
            else
983 runge 131
            {
984 runge 132
                // Otherwise we would like to end the main loop
133
                return false;
983 runge 134
            }
984 runge 135
            break;
976 runge 136
 
984 runge 137
            default:
138
            throw new SocketException("Unknow exception: " + itos(errno));
139
            break;
140
        }
141
    }
142
    else if (status == 0)
143
    {
144
        // Remote host have done a normal shutdown
145
        eventAdd(SocketEvent::TYPE_CONNECTION_CLOSED);
969 runge 146
 
984 runge 147
        if (myReconnectTimeout > 0)
148
        {
149
            reconnectLoop();
969 runge 150
        }
984 runge 151
        else
152
        {
153
            return false;
154
        }
969 runge 155
    }
984 runge 156
    else if (status > 0)
969 runge 157
    {
984 runge 158
        // We have received data
159
        eventAdd(SocketEvent::TYPE_DATA, buffer);
969 runge 160
    }
161
 
984 runge 162
    return true;
969 runge 163
}
164
 
984 runge 165
void AsyncSocket::silentClose()
166
{
167
    if (mySocket != -1)
168
    {
169
        ::close(mySocket);
170
        mySocket = -1;
171
    }
172
}
173
 
174
void AsyncSocket::close()
175
{
176
    silentClose();
177
    eventAdd(SocketEvent::TYPE_CONNECTION_CLOSED);
178
}
179
 
969 runge 180
void AsyncSocket::reconnectLoop()
181
{
984 runge 182
    while (true)
969 runge 183
    {
184
        try
185
        {
186
            connect();
984 runge 187
            return;
969 runge 188
        }
189
        catch (SocketException *e)
190
        {
984 runge 191
            eventAdd(SocketEvent::TYPE_CONNECTION_FAILED, e->getDescription());
192
            eventAdd(SocketEvent::TYPE_WAITING_RECONNECT);
969 runge 193
            sleep(myReconnectTimeout);
194
        }
195
    }
196
}
197
 
981 runge 198
void AsyncSocket::create()
969 runge 199
{
984 runge 200
    silentClose();
976 runge 201
 
969 runge 202
    mySocket = ::socket(AF_INET, SOCK_STREAM, 0);
203
 
204
    int on = 1;
205
    int status = setsockopt(mySocket, SOL_SOCKET, SO_REUSEADDR, (const char*)&on, sizeof(on));
206
    if (status == -1)
207
    {
984 runge 208
        silentClose();
981 runge 209
        throw new SocketException("Create:Reuseaddress: " + itos(errno));
969 runge 210
    }
981 runge 211
}
969 runge 212
 
981 runge 213
void AsyncSocket::startListen()
214
{
215
    create();
216
 
217
    myAddressStruct.sin_family = AF_INET;
218
    myAddressStruct.sin_addr.s_addr = INADDR_ANY;
219
    myAddressStruct.sin_port = htons(myPort);
220
 
221
    int status = ::bind(mySocket, (struct sockaddr *)&myAddressStruct, sizeof(myAddressStruct));
222
 
223
    if (status == -1)
224
    {
984 runge 225
        silentClose();
981 runge 226
        switch (errno)
227
        {
228
            case EACCES:
229
            throw new SocketException("The address is protected, and the user is not the superuser.");
230
 
231
            case EADDRINUSE:
232
            throw new SocketException("The given address is already in use.");
233
 
234
            case EBADF:
235
            throw new SocketException("sockfd is not a valid descriptor.");
236
 
237
            case EINVAL:
238
            throw new SocketException("The socket is already bound to an address.");
239
 
240
            case ENOTSOCK:
241
            throw new SocketException("sockfd is a descriptor for a file, not a socket.");
242
 
243
            //case EACCES:
244
            //throw new SocketException("Search permission is denied on a component of the path prefix. (See also path_resolution(7).)");
245
 
246
            case EADDRNOTAVAIL:
247
            throw new SocketException("A nonexistent interface was requested or the requested address was not local.");
248
 
249
            case EFAULT:
250
            throw new SocketException("addr points outside the user's accessible address space.");
251
 
252
            //case EINVAL:
253
            //throw new SocketException("The addrlen is wrong, or the socket was not in the AF_UNIX family.");
254
 
255
            case ELOOP:
256
            throw new SocketException("Too many symbolic links were encountered in resolving addr.");
257
 
258
            case ENAMETOOLONG:
259
            throw new SocketException("addr is too long.");
260
 
261
            case ENOENT:
262
            throw new SocketException("The file does not exist.");
263
 
264
            case ENOMEM:
265
            throw new SocketException("Insufficient kernel memory was available.");
266
 
267
            case ENOTDIR:
268
            throw new SocketException("A component of the path prefix is not a directory.");
269
 
270
            case EROFS:
271
            throw new SocketException("The socket inode would reside on a read-only file system.");
272
 
273
            default:
274
            throw new SocketException("Unknow exception: " + itos(errno));
275
            break;
276
        }
277
    }
278
 
279
    status = ::listen(mySocket, MAXCONNECTIONS);
280
 
281
    if (status == -1)
282
    {
984 runge 283
        silentClose();
981 runge 284
        switch (errno)
285
        {
286
            case EADDRINUSE:
287
            throw new SocketException("Another socket is already listening on the same port.");
288
 
289
            case EBADF:
290
            throw new SocketException("The argument sockfd is not a valid descriptor.");
291
 
292
            case ENOTSOCK:
293
            throw new SocketException("The argument sockfd is not a socket.");
294
 
295
            case EOPNOTSUPP:
296
            throw new SocketException("The socket is not of a type that supports the listen() operation.");
297
 
298
            default:
299
            throw new SocketException("Unknow exception: " + itos(errno));
300
            break;
301
        }
302
    }
303
}
304
 
305
bool AsyncSocket::accept(AsyncSocket* newSocket)
306
{
307
    int addr_length = sizeof(myAddressStruct);
308
    int socket = ::accept(mySocket, (sockaddr*)&myAddressStruct, (socklen_t*)&addr_length);
309
 
310
    if (socket > 0)
311
    {
312
        newSocket->setSocket(socket);
313
        return true;
314
    }
315
 
316
    return false;
317
}
318
 
319
void AsyncSocket::connect()
320
{
984 runge 321
    eventAdd(SocketEvent::TYPE_CONNECTING);
981 runge 322
 
323
    create();
324
 
969 runge 325
    memset(&myAddressStruct, 0, sizeof(myAddressStruct));
326
 
327
    myAddressStruct.sin_family = AF_INET;
328
    myAddressStruct.sin_port = htons(myPort);
329
 
976 runge 330
    struct hostent *hptr = gethostbyname(myAddress.c_str());
331
    if (hptr == NULL)
332
    {
984 runge 333
        silentClose();
976 runge 334
        throw new SocketException("Connect: Could not resolv ip address");
335
    }
336
 
337
    memcpy(&myAddressStruct.sin_addr, hptr->h_addr, hptr->h_length);
969 runge 338
 
983 runge 339
    int status = ::connect(mySocket, (sockaddr*)&myAddressStruct, sizeof(myAddressStruct));
976 runge 340
 
969 runge 341
    if (status == -1)
342
    {
984 runge 343
        silentClose();
969 runge 344
        switch (errno)
345
        {
346
            case EACCES:
984 runge 347
            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 348
 
349
            case EPERM:
984 runge 350
            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 351
 
352
            case EADDRINUSE:
353
            throw new SocketException("Local address is already in use.");
354
 
355
            case EAFNOSUPPORT:
356
            throw new SocketException("The passed address didn't have the correct address family in its sa_family field.");
357
 
358
            case EAGAIN:
984 runge 359
            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 360
 
361
            case EALREADY:
362
            throw new SocketException("The socket is non-blocking and a previous connection attempt has not yet been completed.");
363
 
364
            case EBADF:
365
            throw new SocketException("The file descriptor is not a valid index in the descriptor table.");
366
 
367
            case ECONNREFUSED:
368
            throw new SocketException("No-one listening on the remote address.");
369
 
370
            case EFAULT:
371
            throw new SocketException("The socket structure address is outside the user's address space.");
372
 
373
            case EINPROGRESS:
984 runge 374
            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 375
 
376
            case EINTR:
377
            throw new SocketException("The system call was interrupted by a signal that was caught; see signal(7).");
378
 
379
            case EISCONN:
380
            throw new SocketException("The socket is already connected.");
381
 
382
            case ENETUNREACH:
383
            throw new SocketException("Network is unreachable.");
384
 
385
            case ENOTSOCK:
386
            throw new SocketException("The file descriptor is not associated with a socket.");
387
 
388
            case ETIMEDOUT:
984 runge 389
            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 390
 
391
            default:
392
            throw new SocketException("Other errors may be generated by the underlying protocol modules. : " + itos(errno));
393
        }
394
    }
395
 
984 runge 396
    eventAdd(SocketEvent::TYPE_CONNECTED);
969 runge 397
}
398
 
984 runge 399
void AsyncSocket::sendData(string data)
969 runge 400
{
981 runge 401
    int status = ::send(mySocket, data.c_str(), data.size(), 0);
402
 
984 runge 403
    if (status == -1)
981 runge 404
    {
984 runge 405
        switch (errno)
981 runge 406
        {
984 runge 407
            case EACCES:
408
            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 409
 
984 runge 410
            case EAGAIN:
411
            throw new SocketException("The socket is marked non-blocking and the requested operation would block.");
981 runge 412
 
984 runge 413
            case EBADF:
414
            throw new SocketException("An invalid descriptor was specified.");
981 runge 415
 
984 runge 416
            case ECONNRESET:
417
            throw new SocketException("Connection reset by peer.");
981 runge 418
 
984 runge 419
            case EDESTADDRREQ:
420
            throw new SocketException("The socket is not connection-mode, and no peer address is set.");
981 runge 421
 
984 runge 422
            case EFAULT:
423
            throw new SocketException("An invalid user space address was specified for an argument.");
981 runge 424
 
984 runge 425
            case EINTR:
426
            throw new SocketException("A signal occurred before any data was transmitted; see signal(7).");
981 runge 427
 
984 runge 428
            case EINVAL:
429
            throw new SocketException("1Invalid argument passed.");
981 runge 430
 
984 runge 431
            case EISCONN:
432
            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 433
 
984 runge 434
            case EMSGSIZE:
435
            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 436
 
984 runge 437
            case ENOBUFS:
438
            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 439
 
984 runge 440
            case ENOMEM:
441
            throw new SocketException("No memory available.");
981 runge 442
 
984 runge 443
            case ENOTCONN:
444
            throw new SocketException("The socket is not connected, and no target has been given.");
981 runge 445
 
984 runge 446
            case ENOTSOCK:
447
            throw new SocketException("The argument s is not a socket.");
981 runge 448
 
984 runge 449
            case EOPNOTSUPP:
450
            throw new SocketException("Some bit in the flags argument is inappropriate for the socket type.");
981 runge 451
 
984 runge 452
            case EPIPE:
453
            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.");
981 runge 454
 
984 runge 455
            default:
456
            throw new SocketException("Unknow exception: " + itos(errno));
981 runge 457
        }
458
    }
459
}
460
 
984 runge 461
void AsyncSocket::eventAdd(unsigned int eventType)
969 runge 462
{
984 runge 463
    eventAdd(eventType, "");
969 runge 464
}
465
 
984 runge 466
void AsyncSocket::eventAdd(unsigned int eventType, string eventData)
974 runge 467
{
984 runge 468
    SocketEvent socketEvent(eventType, eventData);
469
    myEventQueue.push(socketEvent);
470
    myEventSemaphore.broadcast();
974 runge 471
}