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
}