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

361 statements  

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""" 

15 

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 

46 

47 

48class NFStreamer(object): 

49 streamer_id = 0 # class id generator 

50 glock = threading.Lock() 

51 is_windows = "windows" in platform.system().lower() 

52 

53 """ Network Flow Streamer 

54 

55 Examples: 

56 

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() 

62 

63 """ 

64 

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 

107 

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 

112 

113 if NFStreamer.is_windows: 

114 self._mp_context = get_context("spawn") 

115 else: 

116 self._mp_context = get_context("fork") 

117 

118 @property 

119 def source(self): 

120 return self._source 

121 

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 

157 

158 @property 

159 def decode_tunnels(self): 

160 return self._decode_tunnels 

161 

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 

169 

170 @property 

171 def bpf_filter(self): 

172 return self._bpf_filter 

173 

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 

179 

180 @property 

181 def promiscuous_mode(self): 

182 return self._promiscuous_mode 

183 

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 

191 

192 @property 

193 def snapshot_length(self): 

194 return self._snapshot_length 

195 

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 

203 

204 @property 

205 def socket_buffer_size(self): 

206 return self._socket_buffer_size 

207 

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 

215 

216 @property 

217 def idle_timeout(self): 

218 return self._idle_timeout 

219 

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 

229 

230 @property 

231 def active_timeout(self): 

232 return self._active_timeout 

233 

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 

243 

244 @property 

245 def accounting_mode(self): 

246 return self._accounting_mode 

247 

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 

255 

256 @property 

257 def udps(self): 

258 return self._udps 

259 

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 ) 

282 

283 @property 

284 def n_dissections(self): 

285 return self._n_dissections 

286 

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 

294 

295 @property 

296 def statistical_analysis(self): 

297 return self._statistical_analysis 

298 

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 

306 

307 @property 

308 def splt_analysis(self): 

309 return self._splt_analysis 

310 

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 

322 

323 @property 

324 def n_meters(self): 

325 return self._n_meters 

326 

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 

368 

369 @property 

370 def max_nflows(self): 

371 return self._max_nflows 

372 

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).") 

379 

380 @property 

381 def performance_report(self): 

382 return self._performance_report 

383 

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 

394 

395 @property 

396 def system_visibility_mode(self): 

397 return self._system_visibility_mode 

398 

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 

418 

419 @property 

420 def system_visibility_poll_ms(self): 

421 return self._system_visibility_poll_ms 

422 

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 

433 

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 = {} 

445 

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() 

521 

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) 

574 

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 

588 

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()) 

601 

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() 

611 

612 return total_flows 

613 

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