· 10 years ago · Sep 01, 2016, 09:10 AM
1#define _DEFAULT_SOURCE 600
2#include <stdio.h>
3#include <stdlib.h>
4#include <signal.h>
5#include <string.h>
6#include <syslog.h>
7#include <glib.h>
8#include <glib/gprintf.h>
9#include <libpq-fe.h>
10#include <arpa/inet.h>
11#include <sys/stat.h>
12
13#define FLOWSDIR "/backups/flows/users/" // Directory for flows output
14#define UNRELFLOWS "/backups/flows/unrelated_flows.txt" // Unrelated flows temp filename
15#define UNRELFDIR "/backups/flows/unrelated" // Unrelated flows directory
16#define PGCONNSTR "user=radius dbname=radius" // PostgreSQL connection string
17#define TZOFFSET 25200 // GMT+7
18#define ONLINEQUERY "SELECT radacct.username, usergroup.id, radacct.framedipaddress \
19FROM radacct LEFT JOIN usergroup ON radacct.username = usergroup.username \
20WHERE radacct.acctstoptime IS NULL"
21
22// PostgreSQL variables
23PGconn *conn;
24PGresult *res;
25
26// hash tables
27GHashTable *online_ht, *traffic_ht;
28
29// online structure
30struct Online {
31 gchar username[32];
32 gint uid;
33 FILE *file;
34};
35
36// traffic structure
37struct Traffic {
38 // Glib type guint64 differs on x86/x86_64
39 long long unsigned int octetsin;
40 long long unsigned int octetsout;
41};
42
43// Unrelated flows file
44FILE *unrel_file;
45
46// Match if is ip in subnet
47gint ip_in_subnet(gchar *ip, gchar *subnet, guint netmask)
48{
49
50 struct in_addr s_ip, s_subnet, s_netmask;
51 guint octets;
52 inet_aton(ip, &s_ip);
53 inet_aton(subnet, &s_subnet);
54 if ( netmask < 0 || netmask > 32 )
55 {
56 return -1;
57 }
58 octets = (netmask + 7) / 8;
59 s_netmask.s_addr = 0;
60 if (octets > 0)
61 {
62 memset(&s_netmask.s_addr, 255, (gsize)octets - 1);
63 memset((guchar *)&s_netmask.s_addr + (octets - 1), (256 - (1 << (32 - netmask) % 8)), 1);
64 }
65 return ((s_ip.s_addr & s_netmask.s_addr) == (s_subnet.s_addr & s_netmask.s_addr));
66
67}
68
69// Our nets is here
70gint is_client_ip(gchar *ip)
71{
72
73 return ( ip_in_subnet(ip, "93.95.156.0", 24) ||
74 ip_in_subnet(ip, "31.13.178.0", 24) ||
75 ip_in_subnet(ip, "93.171.236.0", 22) );
76
77}
78
79// Iterate through hashtable
80void iterator(gpointer key, gpointer value, gpointer user_data)
81{
82 g_printf("The UID %d traffic is %llu in %llu out\n", *(gint*)key, ((struct Traffic*)value)->octetsin, ((struct Traffic*)value)->octetsout);
83}
84
85// Correctly close all file descriptors in online struct
86void free_online(gpointer value)
87{
88
89 fclose(((struct Online*)value)->file);
90 g_free(value);
91
92}
93
94// For PostgreSQL unexpected termination
95void pg_exit()
96{
97
98 fprintf(stderr, "PostgreSQL error: %s\n", PQerrorMessage(conn));
99 PQclear(res);
100 PQfinish(conn);
101 if ( NULL != unrel_file ) fclose(unrel_file);
102 g_hash_table_remove_all(traffic_ht);
103 g_hash_table_remove_all(online_ht);
104 g_hash_table_destroy(online_ht);
105 g_hash_table_destroy(traffic_ht);
106 closelog();
107 exit(1);
108
109}
110
111// Insert collected traffic data
112void traffic_insert(gpointer key, gpointer value, struct tm *date )
113{
114
115 gchar query[512];
116 gchar date_now[11];
117 gint hours;
118
119 strftime(date_now, 20, "%F", date);
120 hours = (0 < date->tm_hour) ? (date->tm_hour - 1) : 23;
121
122 g_sprintf(query, "INSERT INTO dayflowtemp (username, \"day\", hours, \"in\", \"out\") VALUES (%d, '%s', %d, %llu, %llu)", *(gint*)key, date_now, hours, ((struct Traffic*)value)->octetsin, ((struct Traffic*)value)->octetsout);
123 res = PQexec(conn, query);
124 if ( PQresultStatus(res) != PGRES_COMMAND_OK )
125 {
126 pg_exit();
127 }
128 PQclear(res);
129
130 g_sprintf(query, "INSERT INTO flow (username, \"day\", hours, \"in\", \"out\") VALUES (%d, '%s', %d, %llu, %llu)", *(gint*)key, date_now, hours, ((struct Traffic*)value)->octetsin, ((struct Traffic*)value)->octetsout);
131 res = PQexec(conn, query);
132 if ( PQresultStatus(res) != PGRES_COMMAND_OK )
133 {
134 pg_exit();
135 }
136 PQclear(res);
137
138}
139
140// Interrupt signal handler
141void sigintHandler(int sig_num)
142{
143 g_printf("\n Termination attempt using Ctrl+C \n Freeing resources \n");
144 g_hash_table_foreach(traffic_ht, (GHFunc)iterator, NULL);
145 fflush(stdout);
146 PQfinish(conn);
147 fclose(unrel_file);
148 g_hash_table_remove_all(traffic_ht);
149 g_hash_table_remove_all(online_ht);
150 g_hash_table_destroy(online_ht);
151 g_hash_table_destroy(traffic_ht);
152 closelog();
153 exit(1);
154}
155
156
157// Main loop
158int main()
159{
160
161 // Our custom SIGINT handler
162 signal(SIGINT, sigintHandler);
163 // For syslog logging
164 openlog("billing", 0, LOG_USER);
165
166 // stdin line buffer
167 gchar line[256];
168 // for incoming data
169 gchar *inbuff;
170 time_t unix_time;
171 struct tm * timeinfo;
172 guint octets, octetsin, octetsout, srcport, dstport, prot;
173 gchar date[20], srcaddr[16], dstaddr[16], *userip, *host;
174 // HashTables
175 online_ht = g_hash_table_new_full(g_str_hash, g_str_equal, g_free, free_online);
176 traffic_ht = g_hash_table_new_full(g_int_hash, g_int_equal, g_free, g_free);
177 // misc
178 guint not_used = 0,
179 unrelated = 0,
180 lines = 0,
181 i, rows, hours;
182 guint total = 0;
183 gchar *framedipaddr, *filename, *dirname;
184 struct Online *online, *online_cmp;
185 struct Traffic *traffic;
186 gint *uid;
187 gchar date_now[11];
188 struct tm *tm_now, *tm_date;
189
190 // current date
191 unix_time = time(NULL) + TZOFFSET;
192 tm_date = g_memdup(gmtime(&unix_time), sizeof(struct tm));
193 tm_now = g_memdup(gmtime(&unix_time), sizeof(struct tm));
194 strftime(date_now, 20, "%F", tm_date);
195
196 // Connect to PostgreSQL
197 if ( PQstatus(conn = PQconnectdb(PGCONNSTR)) == CONNECTION_BAD )
198 {
199 pg_exit();
200 }
201
202 // Filling hash tables
203 // Query
204 res = PQexec(conn, ONLINEQUERY);
205 if ( PQresultStatus(res) != PGRES_TUPLES_OK )
206 {
207 pg_exit();
208 }
209
210 rows = PQntuples(res);
211 filename = g_malloc(128 * sizeof(gchar));
212 dirname = g_malloc(128 * sizeof(gchar));
213
214 for ( i=0; i<rows; i++ )
215 {
216
217 // Filling online hash table
218 online = g_malloc(sizeof(struct Online));
219 strncpy(online->username, PQgetvalue(res, i, 0), 32);
220 online->uid = atoi(PQgetvalue(res, i, 1));
221 framedipaddr = g_strdup(PQgetvalue(res, i, 2));
222 // Constructing directory name
223 g_sprintf(dirname, "%s%s", FLOWSDIR, online->username);
224 // If directory does not exists, create it
225 if ( !g_file_test(dirname, G_FILE_TEST_IS_DIR) )
226 {
227 mkdir(dirname, 0755);
228 }
229 // Constructing full path to filename
230 g_sprintf(filename, "%s/%s.txt", dirname, date_now);
231
232 // debug
233 //printf("%s\n", filename);
234
235 // Open file descriptor
236 online->file = fopen(filename, "a");
237 // Insert data in hash table
238 g_hash_table_insert(online_ht, framedipaddr, online);
239
240 // Filling traffic hash table
241 uid = g_memdup(&(online->uid), sizeof(gint)); // table keys must be allocated
242 traffic = g_malloc(sizeof(struct Traffic));
243 traffic->octetsin = 0;
244 traffic->octetsout = 0;
245 g_hash_table_insert(traffic_ht, uid, traffic);
246
247 }
248
249 PQclear(res);
250 g_free(filename);
251 g_free(dirname);
252
253 // Open temp file for unrelated flows
254 unrel_file = fopen(UNRELFLOWS, "a");
255
256 // Skip first line
257 fgets(line, 256, stdin);
258
259 while ( fgets(line, 256, stdin) != NULL )
260 {
261 // date
262 unix_time = g_ascii_strtoll(strtok(line, ","), NULL, 10) + TZOFFSET; // GMT+TZOFFSET hack
263 // localtime doing too many allocs/frees so using gmtime with offset
264 strftime(date, 20, "%F %T", gmtime(&unix_time));
265 // octets
266 if ( NULL == (inbuff = strtok(NULL, ",")) ) continue;
267 octets = strtoul(inbuff, NULL, 10);
268 // srcip
269 if ( NULL == (inbuff = strtok(NULL, ",")) ) continue;
270 strncpy(srcaddr, inbuff, 16);
271 // dstip
272 if ( NULL == (inbuff = strtok(NULL, ",")) ) continue;
273 strncpy(dstaddr, inbuff, 16);
274 // srcport
275 if ( NULL == (inbuff = strtok(NULL, ",")) ) continue;
276 srcport = strtoul(inbuff, NULL, 10);
277 // dstport
278 if ( NULL == (inbuff = strtok(NULL, ",")) ) continue;
279 dstport = strtoul(inbuff, NULL, 10);
280 // proto
281 if ( NULL == (inbuff = strtok(NULL, "\n")) ) continue;
282 prot = strtoul(inbuff, NULL, 10);
283
284 // Unneeded traffic
285 if ( 0 == strcmp(srcaddr, "93.95.156.250") ||
286 0 == strcmp(dstaddr, "93.95.156.250") ||
287 0 == strcmp(srcaddr, "93.171.239.250") ||
288 0 == strcmp(dstaddr, "93.171.239.250") )
289 {
290 not_used++;
291 continue;
292 }
293
294 // Check if srcaddr or dstaddr is client ip
295 if ( is_client_ip(srcaddr) )
296 {
297 userip = srcaddr;
298 host = dstaddr;
299 octetsin = 0;
300 octetsout = octets;
301 }
302 else if ( is_client_ip(dstaddr) )
303 {
304 userip = dstaddr;
305 host = srcaddr;
306 octetsin = octets;
307 octetsout = 0;
308 }
309 else
310 {
311 not_used++;
312 // debug
313 //printf("%s\t%s\t%s\t%d\t%d\t%d\t%d\n", date, srcaddr, dstaddr, srcport, dstport, octets, prot);
314 continue;
315 }
316
317 // If user with that ip is found
318 if ( online = g_hash_table_lookup(online_ht, userip) )
319 {
320 // Increment traffic
321 traffic = g_hash_table_lookup(traffic_ht, &(online->uid));
322 traffic->octetsin += octetsin;
323 traffic->octetsout += octetsout;
324 // Write flow to file
325 fprintf(online->file, "%s\t%s\t%s\t%u\t%u\t%u\t%u\t%u\n", date, userip, host, srcport, dstport, octetsin, octetsout, prot);
326 // debug
327 //printf("%llu\t%llu\t%s\t%s\t%s\t%s\t%u\t%u\t%u\t%u\t%u\n", traffic->octetsin, traffic->octetsout, online->username, date, userip, host, srcport, dstport, octetsin, octetsout, prot);
328 }
329 else
330 {
331 fprintf(unrel_file, "%s\t%s\t%s\t%u\t%u\t%u\t%u\t%u\n", date, userip, host, srcport, dstport, octetsin, octetsout, prot);
332 unrelated++;
333 }
334
335 lines++;
336 total++;
337
338 if ( 50000 < lines )
339 {
340
341 lines = 0;
342 unix_time = time(NULL) + TZOFFSET;
343 g_free(tm_date);
344 tm_date = g_memdup(gmtime(&unix_time), sizeof(struct tm));
345
346 printf("%d hours, 50000 lines, DB work:\n", tm_date->tm_hour);
347
348 // If next day
349 if ( tm_date->tm_yday != tm_now->tm_yday )
350 {
351 g_printf("End of a day (%d -> %d), rotating unrelated flows file\n", tm_now->tm_yday, tm_date->tm_yday);
352 // Move temp unrelated flows file
353 filename = g_malloc(128 * sizeof(gchar));
354 // Constructing full path to new filename
355 strftime(date_now, 20, "%F", tm_now);
356 g_sprintf(filename, "%s/%s.txt", UNRELFDIR, date_now);
357 fclose(unrel_file);
358 rename(UNRELFLOWS, filename);
359 unrel_file = fopen(UNRELFLOWS, "a");
360 g_free(filename);
361
362 }
363
364 // If next hour
365 if ( tm_date->tm_hour != tm_now->tm_hour )
366 {
367 g_printf(" Next hour (%d -> %d). Flushing data...\n", tm_now->tm_hour, tm_date->tm_hour);
368 // Insert traffic data
369 g_hash_table_foreach(traffic_ht, (GHFunc)traffic_insert, tm_date);
370 // Clear traffic data
371 g_hash_table_remove_all(traffic_ht);
372 // Clear online data
373 g_hash_table_remove_all(online_ht);
374 // Renew date_now just in case
375 strftime(date_now, 20, "%F", tm_date);
376 // TODO: Message to SYSLOG must be here
377 syslog(LOG_INFO, "flowtosql: обработано %u NetFlow Ñтрок, неÑвÑзанных %u, неиÑпользованных %u", total, unrelated, not_used);
378 not_used = unrelated = total = 0;
379 g_free(tm_now);
380 tm_now = g_memdup(gmtime(&unix_time), sizeof(struct tm));
381 }
382
383 res = PQexec(conn, ONLINEQUERY);
384 if ( PQresultStatus(res) != PGRES_TUPLES_OK )
385 {
386 pg_exit();
387 }
388
389 rows = PQntuples(res);
390 filename = g_malloc(256 * sizeof(gchar));
391 dirname = g_malloc(256 * sizeof(gchar));
392 for ( i=0; i<rows; i++ )
393 {
394
395 // Filling online hash table
396 online = g_malloc(sizeof(struct Online));
397 strncpy(online->username, PQgetvalue(res, i, 0), 32);
398 online->uid = atoi(PQgetvalue(res, i, 1));
399 framedipaddr = PQgetvalue(res, i, 2);
400
401 // if item with that ip already exists
402 if ( online_cmp = g_hash_table_lookup(online_ht, framedipaddr) )
403 {
404 // but if it belongs to another user
405 if ( 0 != strcmp(online_cmp->username, online->username) )
406 {
407 g_printf(" IP address %s now belongs to user %s\n", framedipaddr, online->username);
408 syslog(LOG_NOTICE, "flowtosql: ip %s отношение изменено к %s", framedipaddr, online->username);
409 // Constructing directory name
410 g_sprintf(dirname, "%s%s", FLOWSDIR, online->username);
411 // If directory does not exists, create it
412 if ( !g_file_test(dirname, G_FILE_TEST_IS_DIR) )
413 {
414 mkdir(dirname, 0755);
415 }
416 // Constructing full path to filename
417 g_sprintf(filename, "%s/%s.txt", dirname, date_now);
418 // Open new file descriptor
419 online->file = fopen(filename, "a");
420 // Update online
421 g_hash_table_replace(online_ht, g_strdup(framedipaddr), online);
422 // If there is not traffic entry for new user
423 if ( NULL == (traffic = g_hash_table_lookup(traffic_ht, &(online->uid))) )
424 {
425 // ...then create it
426 uid = g_memdup(&(online->uid), sizeof(gint)); // key must be allocated
427 traffic = g_malloc(sizeof(struct Traffic));
428 traffic->octetsin = 0;
429 traffic->octetsout = 0;
430 g_hash_table_insert(traffic_ht, uid, traffic);
431 }
432 }
433 else
434 {
435 g_free(online);
436 }
437 }
438 // else if ip address is new
439 else
440 {
441
442 if ( tm_date->tm_hour == tm_now->tm_hour ) printf(" New user connected: %s %s\n", framedipaddr, online->username);
443 // Constructing directory name
444 g_sprintf(dirname, "%s%s", FLOWSDIR, online->username);
445 // If directory does not exists, create it
446 if ( !g_file_test(dirname, G_FILE_TEST_IS_DIR) )
447 {
448 mkdir(dirname, 0755);
449 }
450 // Constructing full path to filename
451 g_sprintf(filename, "%s/%s.txt", dirname, date_now);
452 // Open file descriptor
453 online->file = fopen(filename, "a");
454 // Insert data in online hash table
455 g_hash_table_insert(online_ht, g_strdup(framedipaddr), online);
456
457 // Insert data in traffic hash table
458 uid = g_memdup(&(online->uid), sizeof(gint)); // table key must be allocated
459 traffic = g_malloc(sizeof(struct Traffic));
460 traffic->octetsin = 0;
461 traffic->octetsout = 0;
462 g_hash_table_insert(traffic_ht, uid, traffic);
463
464 }
465
466 }
467 PQclear(res);
468 g_free(filename);
469 g_free(dirname);
470
471 g_printf(" %d items in the hash table, continuing\n", g_hash_table_size(online_ht));
472
473 }
474
475 }
476
477 PQfinish(conn);
478
479 // Correctly free and destroy hash tables
480 g_hash_table_remove_all(traffic_ht);
481 g_hash_table_remove_all(online_ht);
482 g_hash_table_destroy(traffic_ht);
483 g_hash_table_destroy(online_ht);
484
485 g_free(tm_now);
486 g_free(tm_date);
487
488 fclose(unrel_file);
489 closelog();
490
491 return 0;
492
493}