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