Subversion Repositories HomeAutomation

Rev

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