Coverage for /pythoncovmergedfiles/medio/medio/usr/local/lib/python3.11/site-packages/nfstream/streamer.py: 67%
Shortcuts on this page
r m x toggle line displays
j k next/prev highlighted chunk
0 (zero) top of page
1 (one) first highlighted chunk
Shortcuts on this page
r m x toggle line displays
j k next/prev highlighted chunk
0 (zero) top of page
1 (one) first highlighted chunk
1"""
2------------------------------------------------------------------------------------------------------------------------
3streamer.py
4Copyright (C) 2019-22 - NFStream Developers
5This file is part of NFStream, a Flexible Network Data Analysis Framework (https://www.nfstream.org/).
6NFStream is free software: you can redistribute it and/or modify it under the terms of the GNU Lesser General Public
7License as published by the Free Software Foundation, either version 3 of the License, or (at your option) any later
8version.
9NFStream is distributed in the hope that it will be useful, but WITHOUT ANY WARRANTY; without even the implied warranty
10of MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU Lesser General Public License for more details.
11You should have received a copy of the GNU Lesser General Public License along with NFStream.
12If not, see <http://www.gnu.org/licenses/>.
13------------------------------------------------------------------------------------------------------------------------
14"""
16from multiprocessing import get_context
17import threading
18import pandas as pd
19import time as tm
20import os
21import platform
22import psutil
23import csv
24from collections.abc import Iterable
25from os.path import isfile
26from .meter import meter_workflow
27from .anonymizer import NFAnonymizer
28from .engine import is_interface
29from .plugin import NFPlugin
30from .utils import (
31 csv_converter,
32 open_file,
33 RepeatedTimer,
34 update_performances,
35 set_affinity,
36 available_cpus_count,
37)
38from .utils import (
39 validate_flows_per_file,
40 NFMode,
41 create_csv_file_path,
42 NFEvent,
43 validate_rotate_files,
44)
45from .system import system_socket_worflow, match_flow_conn
48class NFStreamer(object):
49 streamer_id = 0 # class id generator
50 glock = threading.Lock()
51 is_windows = "windows" in platform.system().lower()
53 """ Network Flow Streamer
55 Examples:
57 >>> from nfstream import NFStreamer
58 >>> # Streamer object for reading traffic from a PCAP
59 >>> streamer = NFStreamer(source='path/to/file.pcap')
60 >>> # Converting data to pandas dataframe
61 >>> df = streamer.to_pandas()
63 """
65 def __init__(
66 self,
67 source=None,
68 decode_tunnels=True,
69 bpf_filter=None,
70 promiscuous_mode=True,
71 snapshot_length=1536,
72 socket_buffer_size=0,
73 idle_timeout=120, # https://www.kernel.org/doc/Documentation/networking/nf_conntrack-sysctl.txt
74 active_timeout=1800,
75 accounting_mode=0,
76 udps=None,
77 n_dissections=20,
78 statistical_analysis=False,
79 splt_analysis=0,
80 n_meters=0,
81 max_nflows=0,
82 performance_report=0,
83 system_visibility_mode=0,
84 system_visibility_poll_ms=100,
85 ):
86 with NFStreamer.glock:
87 NFStreamer.streamer_id += 1
88 self._idx = NFStreamer.streamer_id
89 self._mode = NFMode.SINGLE_FILE
90 self.source = source
91 self.decode_tunnels = decode_tunnels
92 self.bpf_filter = bpf_filter
93 self.promiscuous_mode = promiscuous_mode
94 self.snapshot_length = snapshot_length
95 self.idle_timeout = idle_timeout
96 self.active_timeout = active_timeout
97 self.accounting_mode = accounting_mode
98 self.udps = udps
99 self.n_dissections = n_dissections
100 self.statistical_analysis = statistical_analysis
101 self.splt_analysis = splt_analysis
102 self.n_meters = n_meters
103 self.max_nflows = max_nflows
104 self.performance_report = performance_report
105 self.system_visibility_mode = system_visibility_mode
106 self.system_visibility_poll_ms = system_visibility_poll_ms
108 # NIC socket buffer size. Default is 0, which means that the pcap default value is used.
109 # The default values may vary depending on the OS and CPU architecture.
110 # Range: 0 - 2^31-1
111 self.socket_buffer_size = socket_buffer_size
113 if NFStreamer.is_windows:
114 self._mp_context = get_context("spawn")
115 else:
116 self._mp_context = get_context("fork")
118 @property
119 def source(self):
120 return self._source
122 @source.setter
123 def source(self, value):
124 if type(value) == list: # List of pcap files to consider as a single one.
125 if len(value) == 0:
126 raise ValueError("Please provide a non-empty list of sources.")
127 else:
128 for i in range(len(value)):
129 try:
130 value[i] = str(os.fspath(value[i]))
131 if not isfile(value[i]):
132 raise TypeError
133 except TypeError:
134 raise ValueError(
135 "Invalid pcap file path at index: " + str(i) + "."
136 )
137 self._mode = NFMode.MULTIPLE_FILES
138 else:
139 try:
140 value = str(os.fspath(value))
141 except TypeError:
142 raise ValueError(
143 "Please specify a pcap file path or a valid network interface name as source."
144 )
145 if isfile(value):
146 self._mode = NFMode.SINGLE_FILE
147 else:
148 interface = is_interface(value)
149 if interface is not None:
150 self._mode = NFMode.INTERFACE
151 value = interface
152 else:
153 raise ValueError(
154 "Please specify a pcap file path or a valid network interface name as source."
155 )
156 self._source = value
158 @property
159 def decode_tunnels(self):
160 return self._decode_tunnels
162 @decode_tunnels.setter
163 def decode_tunnels(self, value):
164 if not isinstance(value, bool):
165 raise ValueError(
166 "Please specify a valid decode_tunnels parameter (possible values: True, False)."
167 )
168 self._decode_tunnels = value
170 @property
171 def bpf_filter(self):
172 return self._bpf_filter
174 @bpf_filter.setter
175 def bpf_filter(self, value):
176 if not isinstance(value, str) and value is not None:
177 raise ValueError("Please specify a valid bpf_filter format.")
178 self._bpf_filter = value
180 @property
181 def promiscuous_mode(self):
182 return self._promiscuous_mode
184 @promiscuous_mode.setter
185 def promiscuous_mode(self, value):
186 if not isinstance(value, bool):
187 raise ValueError(
188 "Please specify a valid promiscuous_mode parameter (possible values: True, False)."
189 )
190 self._promiscuous_mode = value
192 @property
193 def snapshot_length(self):
194 return self._snapshot_length
196 @snapshot_length.setter
197 def snapshot_length(self, value):
198 if not isinstance(value, int) or value <= 0:
199 raise ValueError(
200 "Please specify a valid snapshot_length parameter (positive integer)."
201 )
202 self._snapshot_length = value
204 @property
205 def socket_buffer_size(self):
206 return self._socket_buffer_size
208 @socket_buffer_size.setter
209 def socket_buffer_size(self, value):
210 if not isinstance(value, int) or (value < 0 or value > 2**31 - 1):
211 raise ValueError(
212 "Please specify a valid socket_buffer_size parameter (positive integer <= 2^31-1)."
213 )
214 self._socket_buffer_size = value
216 @property
217 def idle_timeout(self):
218 return self._idle_timeout
220 @idle_timeout.setter
221 def idle_timeout(self, value):
222 if not isinstance(value, int) or (
223 (value < 0) or (value * 1000) > 18446744073709551615
224 ): # max uint64_t
225 raise ValueError(
226 "Please specify a valid idle_timeout parameter (positive integer in seconds)."
227 )
228 self._idle_timeout = value
230 @property
231 def active_timeout(self):
232 return self._active_timeout
234 @active_timeout.setter
235 def active_timeout(self, value):
236 if not isinstance(value, int) or (
237 (value < 0) or (value * 1000) > 18446744073709551615
238 ): # max uint64_t
239 raise ValueError(
240 "Please specify a valid active_timeout parameter (positive integer in seconds)."
241 )
242 self._active_timeout = value
244 @property
245 def accounting_mode(self):
246 return self._accounting_mode
248 @accounting_mode.setter
249 def accounting_mode(self, value):
250 if not isinstance(value, int) or (value not in [0, 1, 2, 3]):
251 raise ValueError(
252 "Please specify a valid accounting_mode parameter (possible values: 0, 1, 2, 3)."
253 )
254 self._accounting_mode = value
256 @property
257 def udps(self):
258 return self._udps
260 @udps.setter
261 def udps(self, value):
262 multiple = isinstance(value, Iterable)
263 if multiple:
264 for plugin in value:
265 if isinstance(plugin, NFPlugin):
266 pass
267 else:
268 raise ValueError(
269 "User defined plugins must inherit from NFPlugin type."
270 )
271 self._udps = value
272 else:
273 if isinstance(value, NFPlugin):
274 self._udps = (value,)
275 else:
276 if value is None:
277 self._udps = ()
278 else:
279 raise ValueError(
280 "User defined plugins must inherit from NFPlugin type."
281 )
283 @property
284 def n_dissections(self):
285 return self._n_dissections
287 @n_dissections.setter
288 def n_dissections(self, value):
289 if not isinstance(value, int) or (value < 0 or value > 255):
290 raise ValueError(
291 "Please specify a valid n_dissections parameter (possible values in : [0,...,255])."
292 )
293 self._n_dissections = value
295 @property
296 def statistical_analysis(self):
297 return self._statistical_analysis
299 @statistical_analysis.setter
300 def statistical_analysis(self, value):
301 if not isinstance(value, bool):
302 raise ValueError(
303 "Please specify a valid statistical_analysis parameter (possible values: True, False)."
304 )
305 self._statistical_analysis = value
307 @property
308 def splt_analysis(self):
309 return self._splt_analysis
311 @splt_analysis.setter
312 def splt_analysis(self, value):
313 if not isinstance(value, int) or (value < 0 or value > 65535):
314 raise ValueError(
315 "Please specify a valid splt_analysis parameter (possible values in : [0,...,65535])"
316 )
317 if value > 255:
318 print(
319 "[WARNING]: The specified splt_analysis parameter is higher than 255. High values can impact the performance of the tool."
320 )
321 self._splt_analysis = value
323 @property
324 def n_meters(self):
325 return self._n_meters
327 @n_meters.setter
328 def n_meters(self, value):
329 if isinstance(value, int) and value >= 0:
330 pass
331 else:
332 raise ValueError(
333 "Please specify a valid n_meters parameter (>=1 or 0 for auto scaling)."
334 )
335 c_cpus, c_cores = available_cpus_count(), psutil.cpu_count(logical=False)
336 if (
337 c_cores is None
338 ): # Patch for platforms returning None (https://github.com/giampaolo/psutil/issues/1078)
339 c_cores = c_cpus
340 if value == 0:
341 if platform.system() == "Linux" and self._mode == NFMode.INTERFACE:
342 self._n_meters = (
343 c_cpus - 1
344 ) # We are in live capture mode and kernel fanout will be available
345 # only on Linux, we set the n_meters to detected logical CPUs -1
346 else: # Windows, MacOS, offline capture
347 if c_cpus >= c_cores:
348 if (
349 c_cpus == 2 * c_cores or c_cpus == c_cores
350 ): # multi-thread or single threaded
351 self._n_meters = c_cores - 1
352 else:
353 self._n_meters = int(divmod(c_cpus / 2, 1)[0]) - 1
354 else: # weird case, fallback on cpu count.
355 self._n_meters = c_cpus - 1
356 else:
357 if (value + 1) <= c_cpus:
358 self._n_meters = value
359 else: # avoid contention
360 print(
361 "WARNING: n_meters set to :{} in order to avoid contention.".format(
362 c_cpus - 1
363 )
364 )
365 self._n_meters = c_cpus - 1
366 if self._n_meters == 0: # one CPU case
367 self._n_meters = 1
369 @property
370 def max_nflows(self):
371 return self._max_nflows
373 @max_nflows.setter
374 def max_nflows(self, value):
375 if isinstance(value, int) and value >= 0:
376 self._max_nflows = value - 1
377 else:
378 raise ValueError("Please specify a valid max_nflows parameter (>=0).")
380 @property
381 def performance_report(self):
382 return self._performance_report
384 @performance_report.setter
385 def performance_report(self, value):
386 if isinstance(value, int) and value >= 0:
387 pass
388 else:
389 raise ValueError(
390 "Please specify a valid performance_report parameter (>=1 for reporting interval (seconds)"
391 " or 0 to disable). [Available only for Live capture]"
392 )
393 self._performance_report = value
395 @property
396 def system_visibility_mode(self):
397 return self._system_visibility_mode
399 @system_visibility_mode.setter
400 def system_visibility_mode(self, value):
401 if isinstance(value, int) and value in [0, 1]:
402 if self._mode == NFMode.SINGLE_FILE and value > 0:
403 print(
404 "WARNING: system_visibility_mode switched to 0 in offline capture "
405 "(available only for live capture)"
406 )
407 value = 0
408 else:
409 pass
410 else:
411 raise ValueError(
412 "Please specify a valid system_visibility_mode parameter\n"
413 "0: disable\n"
414 "1: process information\n"
415 "[Available only for live capture on the system generating the traffic]"
416 )
417 self._system_visibility_mode = value
419 @property
420 def system_visibility_poll_ms(self):
421 return self._system_visibility_poll_ms
423 @system_visibility_poll_ms.setter
424 def system_visibility_poll_ms(self, value):
425 if isinstance(value, int) and value >= 0:
426 pass
427 else:
428 raise ValueError(
429 "Please specify a valid system_visibility_poll_ms parameter "
430 "(positive integer in milliseconds)"
431 )
432 self._system_visibility_poll_ms = value
434 def __iter__(self):
435 lock = self._mp_context.Lock()
436 lock.acquire()
437 meters = []
438 performances = []
439 n_terminated = 0
440 child_error = None
441 rt = None
442 socket_listener = None
443 browser_listener = None
444 conn_cache = {}
446 # To avoid issues on PyPy on Windows (See https://foss.heptapod.net/pypy/pypy/-/issues/3488), All
447 # multiprocessing Value invocation must be performed before the call to Queue.
448 n_meters = self.n_meters
449 idx_generator = self._mp_context.Value("i", 0)
450 for i in range(n_meters):
451 performances.append(
452 [
453 self._mp_context.Value("I", 0),
454 self._mp_context.Value("I", 0),
455 self._mp_context.Value("I", 0),
456 ]
457 )
458 channel = self._mp_context.Queue(maxsize=32767) # Backpressure strategy.
459 # We set it to (2^15-1) to cope with OSX max semaphore value.
460 group_id = os.getpid() + self._idx # Used for fanout on Linux systems
461 try:
462 for i in range(n_meters):
463 meters.append(
464 self._mp_context.Process(
465 target=meter_workflow,
466 args=(
467 self.source,
468 self.snapshot_length,
469 self.decode_tunnels,
470 self.bpf_filter,
471 self.promiscuous_mode,
472 n_meters,
473 i,
474 self._mode,
475 self.idle_timeout * 1000,
476 self.active_timeout * 1000,
477 self.accounting_mode,
478 self.udps,
479 self.n_dissections,
480 self.statistical_analysis,
481 self.splt_analysis,
482 channel,
483 performances[i],
484 lock,
485 group_id,
486 self.system_visibility_mode,
487 self.socket_buffer_size,
488 ),
489 )
490 )
491 meters[i].daemon = True # demonize meter
492 meters[i].start()
493 if self._mode == NFMode.INTERFACE and self.performance_report > 0:
494 if platform.system() == "Linux":
495 rt = RepeatedTimer(
496 self.performance_report,
497 update_performances,
498 performances,
499 True,
500 idx_generator,
501 )
502 else:
503 rt = RepeatedTimer(
504 self.performance_report,
505 update_performances,
506 performances,
507 False,
508 idx_generator,
509 )
510 if self._mode == NFMode.INTERFACE and self.system_visibility_mode:
511 socket_listener = self._mp_context.Process(
512 target=system_socket_worflow,
513 args=(
514 channel,
515 self.idle_timeout * 1000,
516 self.system_visibility_poll_ms / 1000,
517 ),
518 )
519 socket_listener.daemon = True # demonize socket_listener
520 socket_listener.start()
522 while True:
523 try:
524 recv = channel.get()
525 if recv is None: # termination and stats
526 n_terminated += 1
527 if n_terminated == n_meters:
528 break # We finish up when all metering jobs are terminated
529 else:
530 if recv.id == NFEvent.ERROR: # Error message
531 for i in range(n_meters): # We break workflow loop
532 meters[i].terminate()
533 child_error = recv.message
534 break
535 elif recv.id == NFEvent.ALL_AFFINITY_SET:
536 set_affinity(
537 0
538 ) # we pin streamer to core 0 as it's the less intensive task and several services runs
539 # by default on this core.
540 elif recv.id == NFEvent.SOCKET_CREATE:
541 conn_cache[recv.key] = [recv.process_name, recv.process_pid]
542 elif recv.id == NFEvent.SOCKET_REMOVE:
543 del conn_cache[recv.key]
544 else: # NFEvent.FLOW
545 recv.id = idx_generator.value # Unify ID
546 idx_generator.value = idx_generator.value + 1
547 if (
548 self._mode == NFMode.INTERFACE
549 and self.system_visibility_mode
550 ):
551 recv = match_flow_conn(conn_cache, recv)
552 yield recv
553 if recv.id == self.max_nflows:
554 raise KeyboardInterrupt # We reached the maximum flows count defined by the user.
555 except KeyboardInterrupt:
556 for i in range(n_meters): # We break workflow loop
557 meters[i].terminate()
558 break
559 for i in range(n_meters):
560 if meters[i].is_alive():
561 meters[i].join() # Join metering jobs
562 if self._mode == NFMode.INTERFACE and self.performance_report > 0:
563 rt.stop()
564 if self._mode == NFMode.INTERFACE and self.system_visibility_mode:
565 socket_listener.terminate()
566 channel.close() # We close the queue
567 channel.join_thread() # and we join its thread
568 if child_error is not None:
569 raise ValueError(child_error)
570 except (
571 ValueError
572 ) as observer_error: # job initiation failed due to some bad observer parameters.
573 raise ValueError(observer_error)
575 def to_csv(
576 self, path=None, columns_to_anonymize=(), flows_per_file=0, rotate_files=0
577 ):
578 validate_flows_per_file(flows_per_file)
579 validate_rotate_files(rotate_files)
580 chunked, chunk_idx = True, -1
581 if flows_per_file == 0:
582 chunked = False
583 output_path = create_csv_file_path(path, self.source)
584 total_flows, chunk_flows = 0, 0
585 anon = NFAnonymizer(cols_names=columns_to_anonymize)
586 f = None
587 writer = None
589 for flow in self:
590 try:
591 if total_flows == 0 or (
592 chunked and (chunk_flows > flows_per_file)
593 ): # header creation
594 if f is not None:
595 f.close()
596 chunk_flows = 1
597 chunk_idx += 1
598 f = open_file(output_path, chunked, chunk_idx, rotate_files)
599 writer = csv.writer(f)
600 writer.writerow(flow.keys())
602 values = anon.process(flow)
603 writer.writerow(values)
604 total_flows += 1
605 chunk_flows += 1
606 except KeyboardInterrupt:
607 pass
608 if f is not None:
609 if not f.closed:
610 f.close()
612 return total_flows
614 def to_pandas(self, columns_to_anonymize=()):
615 """streamer to pandas function"""
616 temp_file_path = "nfstream-{pid}-{iid}-{ts}.csv".format(
617 pid=os.getpid(), iid=NFStreamer.streamer_id, ts=tm.time()
618 )
619 total_flows = self.to_csv(
620 path=temp_file_path,
621 columns_to_anonymize=columns_to_anonymize,
622 flows_per_file=0,
623 )
624 if total_flows > 0: # If there is flows, return Dataframe else return None.
625 df = pd.read_csv(
626 temp_file_path, engine="c"
627 ) # Use C engine for superior performance (non-experimental)
628 if total_flows != df.shape[0]:
629 print(
630 "WARNING: {} flows ignored by pandas type conversion. Consider using to_csv() "
631 "method if drops are critical.".format(
632 abs(df.shape[0] - total_flows)
633 )
634 )
635 else:
636 df = None
637 if os.path.exists(temp_file_path):
638 os.remove(temp_file_path)
639 return df