]> git.stg.codes - stg.git/blob - projects/rscriptd/listener.cpp
Усунуто проблему race conditions у rscriptd
[stg.git] / projects / rscriptd / listener.cpp
1 /*
2  *    This program is free software; you can redistribute it and/or modify
3  *    it under the terms of the GNU General Public License as published by
4  *    the Free Software Foundation; either version 2 of the License, or
5  *    (at your option) any later version.
6  *
7  *    This program is distributed in the hope that it will be useful,
8  *    but WITHOUT ANY WARRANTY; without even the implied warranty of
9  *    MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE.  See the
10  *    GNU General Public License for more details.
11  *
12  *    You should have received a copy of the GNU General Public License
13  *    along with this program; if not, write to the Free Software
14  *    Foundation, Inc., 59 Temple Place, Suite 330, Boston, MA  02111-1307  USA
15  */
16
17 /*
18  *    Author : Boris Mikhailenko <stg34@stargazer.dp.ua>
19  *    Author : Maxim Mamontov <faust@stargazer.dp.ua>
20  */
21
22 #include <arpa/inet.h>
23 #include <sys/uio.h> // readv
24 #include <sys/types.h> // for historical versions of BSD
25 #include <sys/socket.h>
26 #include <netinet/in.h>
27 #include <unistd.h>
28
29 #include <csignal>
30 #include <cerrno>
31 #include <ctime>
32 #include <cstring>
33 #include <sstream>
34 #include <algorithm>
35
36 #include "listener.h"
37 #include "script_executer.h"
38 #include "stg_locker.h"
39 #include "common.h"
40
41 void InitEncrypt(BLOWFISH_CTX * ctx, const std::string & password);
42 void Decrypt(BLOWFISH_CTX * ctx, char * dst, const char * src, int len8);
43
44 //-----------------------------------------------------------------------------
45 LISTENER::LISTENER()
46     : WriteServLog(GetStgLogger()),
47       port(0),
48       running(false),
49       receiverStopped(true),
50       processorStopped(true),
51       userTimeout(0),
52       listenSocket(0)
53 {
54 version = "rscriptd listener v.1.2";
55
56 pthread_mutex_init(&mutex, NULL);
57 }
58 //-----------------------------------------------------------------------------
59 void LISTENER::SetPassword(const std::string & p)
60 {
61 password = p;
62 printfd(__FILE__, "Encryption initiated with password \'%s\'\n", password.c_str());
63 InitEncrypt(&ctxS, password);
64 }
65 //-----------------------------------------------------------------------------
66 bool LISTENER::Start()
67 {
68 printfd(__FILE__, "LISTENER::Start()\n");
69 running = true;
70
71 if (PrepareNet())
72     {
73     return true;
74     }
75
76 if (receiverStopped)
77     {
78     if (pthread_create(&receiverThread, NULL, Run, this))
79         {
80         errorStr = "Cannot create thread.";
81         return true;
82         }
83     }
84
85 if (processorStopped)
86     {
87     if (pthread_create(&processorThread, NULL, RunProcessor, this))
88         {
89         errorStr = "Cannot create thread.";
90         return true;
91         }
92     }
93
94 errorStr = "";
95
96 return false;
97 }
98 //-----------------------------------------------------------------------------
99 bool LISTENER::Stop()
100 {
101 running = false;
102
103 printfd(__FILE__, "LISTENER::Stop()\n");
104
105 usleep(500000);
106
107 if (!processorStopped)
108     {
109     //5 seconds to thread stops itself
110     for (int i = 0; i < 25 && !processorStopped; i++)
111         {
112         usleep(200000);
113         }
114
115     //after 5 seconds waiting thread still running. now killing it
116     if (!processorStopped)
117         {
118         //TODO pthread_cancel()
119         if (pthread_kill(processorThread, SIGINT))
120             {
121             errorStr = "Cannot kill thread.";
122             return true;
123             }
124         printfd(__FILE__, "LISTENER killed Timeouter\n");
125         }
126     }
127
128 if (!receiverStopped)
129     {
130     //5 seconds to thread stops itself
131     for (int i = 0; i < 25 && !receiverStopped; i++)
132         {
133         usleep(200000);
134         }
135
136     //after 5 seconds waiting thread still running. now killing it
137     if (!receiverStopped)
138         {
139         //TODO pthread_cancel()
140         if (pthread_kill(receiverThread, SIGINT))
141             {
142             errorStr = "Cannot kill thread.";
143             return true;
144             }
145         printfd(__FILE__, "LISTENER killed Run\n");
146         }
147     }
148
149 pthread_join(receiverThread, NULL);
150 pthread_join(processorThread, NULL);
151
152 pthread_mutex_destroy(&mutex);
153
154 FinalizeNet();
155
156 std::for_each(users.begin(), users.end(), DisconnectUser(*this));
157
158 printfd(__FILE__, "LISTENER::Stoped successfully.\n");
159
160 return false;
161 }
162 //-----------------------------------------------------------------------------
163 void * LISTENER::Run(void * d)
164 {
165 LISTENER * ia = static_cast<LISTENER *>(d);
166
167 ia->Runner();
168
169 return NULL;
170 }
171 //-----------------------------------------------------------------------------
172 void LISTENER::Runner()
173 {
174 receiverStopped = false;
175
176 while (running)
177     {
178     RecvPacket();
179     }
180
181 receiverStopped = true;
182 }
183 //-----------------------------------------------------------------------------
184 void * LISTENER::RunProcessor(void * d)
185 {
186 LISTENER * ia = static_cast<LISTENER *>(d);
187
188 ia->ProcessorRunner();
189
190 return NULL;
191 }
192 //-----------------------------------------------------------------------------
193 void LISTENER::ProcessorRunner()
194 {
195 processorStopped = false;
196
197 while (running)
198     {
199     usleep(500000);
200     if (!pending.empty())
201         ProcessPending();
202     ProcessTimeouts();
203     }
204
205 processorStopped = true;
206 }
207 //-----------------------------------------------------------------------------
208 bool LISTENER::PrepareNet()
209 {
210 listenSocket = socket(AF_INET, SOCK_DGRAM, 0);
211
212 if (listenSocket < 0)
213     {
214     errorStr = "Cannot create socket.";
215     return true;
216     }
217
218 printfd(__FILE__, "Port: %d\n", port);
219
220 struct sockaddr_in listenAddr;
221 listenAddr.sin_family = AF_INET;
222 listenAddr.sin_port = htons(port);
223 listenAddr.sin_addr.s_addr = inet_addr("0.0.0.0");
224
225 if (bind(listenSocket, (struct sockaddr*)&listenAddr, sizeof(listenAddr)) < 0)
226     {
227     errorStr = "LISTENER: Bind failed.";
228     return true;
229     }
230
231 printfd(__FILE__, "LISTENER::PrepareNet() >>>> Start successfull.\n");
232
233 return false;
234 }
235 //-----------------------------------------------------------------------------
236 bool LISTENER::FinalizeNet()
237 {
238 close(listenSocket);
239
240 return false;
241 }
242 //-----------------------------------------------------------------------------
243 bool LISTENER::RecvPacket()
244 {
245 struct iovec iov[2];
246
247 char buffer[RS_MAX_PACKET_LEN];
248 RS_PACKET_HEADER packetHead;
249
250 iov[0].iov_base = reinterpret_cast<char *>(&packetHead);
251 iov[0].iov_len = sizeof(packetHead);
252 iov[1].iov_base = buffer;
253 iov[1].iov_len = sizeof(buffer);
254
255 size_t dataLen = 0;
256 while (dataLen < sizeof(buffer))
257     {
258     if (!WaitPackets(listenSocket))
259         {
260         if (!running)
261             return false;
262         continue;
263         }
264     int portion = readv(listenSocket, iov, 2);
265     if (portion < 0)
266         {
267         return true;
268         }
269     dataLen += portion;
270     }
271
272 if (CheckHeader(packetHead))
273     {
274     printfd(__FILE__, "Invalid packet or incorrect protocol version!\n");
275     return true;
276     }
277
278 std::string userLogin((char *)packetHead.login);
279 PendingData data;
280 data.login = userLogin;
281 data.ip = ntohl(packetHead.ip);
282 data.id = ntohl(packetHead.id);
283
284 if (packetHead.packetType == RS_ALIVE_PACKET)
285     {
286     data.type = PendingData::ALIVE;
287     }
288 else if (packetHead.packetType == RS_CONNECT_PACKET)
289     {
290     data.type = PendingData::CONNECT;
291     if (GetParams(buffer, data))
292         {
293         return true;
294         }
295     }
296 else if (packetHead.packetType == RS_DISCONNECT_PACKET)
297     {
298     data.type = PendingData::DISCONNECT;
299     if (GetParams(buffer, data))
300         {
301         return true;
302         }
303     }
304
305 STG_LOCKER lock(&mutex, __FILE__, __LINE__);
306 pending.push_back(data);
307
308 return false;
309 }
310 //-----------------------------------------------------------------------------
311 bool LISTENER::GetParams(char * buffer, UserData & data)
312 {
313 RS_PACKET_TAIL packetTail;
314
315 Decrypt(&ctxS, (char *)&packetTail, buffer, sizeof(packetTail) / 8);
316
317 if (strncmp((char *)packetTail.magic, RS_ID, RS_MAGIC_LEN))
318     {
319     printfd(__FILE__, "Invalid crypto magic\n");
320     return true;
321     }
322
323 std::stringstream params;
324 params << data.login << " "
325        << inet_ntostring(data.ip) << " "
326        << data.id << " "
327        << (char *)packetTail.params;
328
329 data.params = params.str();
330
331 return false;
332 }
333 //-----------------------------------------------------------------------------
334 void LISTENER::ProcessPending()
335 {
336 std::list<PendingData> localPending;
337
338     {
339     STG_LOCKER lock(&mutex, __FILE__, __LINE__);
340     printfd(__FILE__, "Pending data size: %d\n", pending.size());
341     localPending.swap(pending);
342     }
343
344 std::list<PendingData>::iterator it(localPending.begin());
345 while (it != localPending.end())
346     {
347     std::vector<AliveData>::iterator uit(
348             std::lower_bound(
349                 users.begin(),
350                 users.end(),
351                 it->login)
352             );
353     if (it->type == PendingData::CONNECT)
354         {
355         if (uit == users.end() || uit->login != it->login)
356             {
357             // Add new user
358             Connect(*it);
359             users.insert(uit, AliveData(static_cast<UserData>(*it)));
360             }
361         else if (uit->login == it->login)
362             {
363             // Update already existing user
364             time(&uit->lastAlive);
365             uit->params = it->params;
366             }
367         }
368     else if (it->type == PendingData::ALIVE)
369         {
370         if (uit != users.end() && uit->login == it->login)
371             {
372             // Update existing user
373             time(&uit->lastAlive);
374             }
375         }
376     else if (it->type == PendingData::DISCONNECT)
377         {
378         if (uit != users.end() && uit->login == it->login.c_str())
379             {
380             // Disconnect existing user
381             Disconnect(*uit);
382             users.erase(uit);
383             }
384         }
385     ++it;
386     }
387 }
388 //-----------------------------------------------------------------------------
389 void LISTENER::ProcessTimeouts()
390 {
391 const std::vector<AliveData>::iterator it(
392         std::stable_partition(
393             users.begin(),
394             users.end(),
395             IsNotTimedOut(userTimeout)
396         )
397     );
398
399 if (it != users.end())
400     {
401     printfd(__FILE__, "Total users: %d, users to disconnect: %d\n", users.size(), std::distance(it, users.end()));
402
403     std::for_each(
404             it,
405             users.end(),
406             DisconnectUser(*this)
407         );
408
409     users.erase(it, users.end());
410     }
411 }
412 //-----------------------------------------------------------------------------
413 bool LISTENER::Connect(const UserData & data) const
414 {
415 printfd(__FILE__, "Connect %s\n", data.login.c_str());
416 if (access(scriptOnConnect.c_str(), X_OK) == 0)
417     {
418     if (ScriptExec(scriptOnConnect + " " + data.params))
419         {
420         WriteServLog("Script %s cannot be executed for an unknown reason.", scriptOnConnect.c_str());
421         return true;
422         }
423     }
424 else
425     {
426     WriteServLog("Script %s cannot be executed. File not found.", scriptOnConnect.c_str());
427     return true;
428     }
429 return false;
430 }
431 //-----------------------------------------------------------------------------
432 bool LISTENER::Disconnect(const UserData & data) const
433 {
434 printfd(__FILE__, "Disconnect %s\n", data.login.c_str());
435 if (access(scriptOnDisconnect.c_str(), X_OK) == 0)
436     {
437     if (ScriptExec(scriptOnDisconnect + " " + data.params))
438         {
439         WriteServLog("Script %s cannot be executed for an unknown reson.", scriptOnDisconnect.c_str());
440         return true;
441         }
442     }
443 else
444     {
445     WriteServLog("Script %s cannot be executed. File not found.", scriptOnDisconnect.c_str());
446     return true;
447     }
448 return false;
449 }
450 //-----------------------------------------------------------------------------
451 bool LISTENER::CheckHeader(const RS_PACKET_HEADER & header) const
452 {
453 if (strncmp((char *)header.magic, RS_ID, RS_MAGIC_LEN))
454     {
455     return true;
456     }
457 if (strncmp((char *)header.protoVer, "02", RS_PROTO_VER_LEN))
458     {
459     return true;
460     }
461 return false;
462 }
463 //-----------------------------------------------------------------------------
464 bool LISTENER::WaitPackets(int sd) const
465 {
466 fd_set rfds;
467 FD_ZERO(&rfds);
468 FD_SET(sd, &rfds);
469
470 struct timeval tv;
471 tv.tv_sec = 0;
472 tv.tv_usec = 500000;
473
474 int res = select(sd + 1, &rfds, NULL, NULL, &tv);
475 if (res == -1) // Error
476     {
477     if (errno != EINTR)
478         {
479         printfd(__FILE__, "Error on select: '%s'\n", strerror(errno));
480         }
481     return false;
482     }
483
484 if (res == 0) // Timeout
485     {
486     return false;
487     }
488
489 return true;
490 }
491 //-----------------------------------------------------------------------------
492 inline
493 void InitEncrypt(BLOWFISH_CTX * ctx, const std::string & password)
494 {
495 unsigned char keyL[PASSWD_LEN];
496 memset(keyL, 0, PASSWD_LEN);
497 strncpy((char *)keyL, password.c_str(), PASSWD_LEN);
498 Blowfish_Init(ctx, keyL, PASSWD_LEN);
499 }
500 //-----------------------------------------------------------------------------
501 inline
502 void Decrypt(BLOWFISH_CTX * ctx, char * dst, const char * src, int len8)
503 {
504 if (dst != src)
505     memcpy(dst, src, len8 * 8);
506
507 for (int i = 0; i < len8; i++)
508     Blowfish_Decrypt(ctx, (uint32_t *)(dst + i * 8), (uint32_t *)(dst + i * 8 + 4));
509 }
510 //-----------------------------------------------------------------------------