autool-dispatcher/calculate/get_traffic_summary.py
2026-06-17 19:50:39 +08:00

256 lines
8.7 KiB
Python

#!/usr/bin/env python3
"""
Merge all traffic_count*.txt files from the traffic_data directory into a single
traffic_summary.csv file.
Directory structure:
\\192.168.2.75\dpi-sync\autool_config\data\traffic_data\
└── <PC_MAC_IP>/
└── <package_name>/
└── traffic_count*.txt
Input line format (comma-separated):
App_ID, App_Name, Domain, Flow(IP:Port-IP:Port-Protocol), Transport_Protocol, App_Protocol, Traffic_Size, [Organization]
Output CSV format:
App ID, App Name, Domain, Traffic Size(Bytes), Flow Count, Traffic Ratio, Domain Traffic Ratio, Organization
"""
import os
import re
import csv
import glob
import argparse
from collections import defaultdict
from datetime import datetime
def parse_traffic_size(size_str):
"""Parse traffic size string like '128 B', '1.5 KB', '2.3 MB' into bytes."""
size_str = size_str.strip()
if not size_str:
return 0
# Match number and optional unit
match = re.match(r'([\d.]+)\s*([KMGT]?B?)', size_str, re.IGNORECASE)
if not match:
return 0
value = float(match.group(1))
unit = match.group(2).upper().strip()
multipliers = {
'': 1, 'B': 1,
'KB': 1024, 'K': 1024,
'MB': 1024 ** 2, 'M': 1024 ** 2,
'GB': 1024 ** 3, 'G': 1024 ** 3,
'TB': 1024 ** 4, 'T': 1024 ** 4,
}
return int(value * multipliers.get(unit, 1))
def determine_domain_or_model(fields):
"""
Determine the domain field value.
If App_Protocol is recognized (e.g., DNS, TLS, HTTP, QUIC...), use the domain field directly.
Otherwise, generate a model_data identifier from the flow info.
"""
domain = fields[2].strip() if len(fields) > 2 else ''
# app_protocol = fields[5].strip() if len(fields) > 5 else ''
flow_info = fields[3].strip() if len(fields) > 3 else ''
# # Known application protocols that have domain info
# known_protocols = {
# 'DNS', 'TLS', 'HTTP', 'HTTPS', 'QUIC', 'HTTP/S', 'SSL',
# 'NTP', 'STUN', 'DTLS', 'MQTT', 'MDNS', 'SSDP', 'LLMNR',
# }
# if domain and app_protocol.upper() in known_protocols:
# return domain
# elif domain:
# return domain
if domain != 'model_data:':
return domain
else:
# Use flow info as model_data identifier
return f'model_data:{flow_info}'
def parse_traffic_file(filepath):
"""
Parse a single traffic_count file.
Returns a list of tuples: (app_id, app_name, domain, traffic_bytes, organization)
"""
records = []
try:
with open(filepath, 'r', encoding='utf-8', errors='ignore') as f:
for line_num, line in enumerate(f, 1):
line = line.strip()
if not line:
continue
# Split by comma, but be careful with trailing comma
fields = line.split(',')
# Need at least 7 fields: app_id, app_name, domain, flow, transport, app_protocol, size
if len(fields) < 7:
continue
app_id = fields[0].strip()
app_name = fields[1].strip()
domain = determine_domain_or_model(fields)
traffic_size_str = fields[6].strip()
organization = fields[7].strip() if len(fields) > 7 else ''
traffic_bytes = parse_traffic_size(traffic_size_str)
records.append((app_id, app_name, domain, traffic_bytes, organization))
except Exception as e:
print(f" [WARNING] Error reading {filepath}: {e}")
return records
def merge_traffic_data(base_dir, output_file):
"""
Scan all traffic_count*.txt files under base_dir and merge them into output_file.
"""
print(f"Scanning directory: {base_dir}")
print(f"Output file: {output_file}")
print()
# Find all traffic_count*.txt files
pattern = os.path.join(base_dir, '**', 'traffic_count*.txt')
traffic_files = glob.glob(pattern, recursive=True)
if not traffic_files:
print(f"No traffic_count files found in {base_dir}")
return
print(f"Found {len(traffic_files)} traffic_count files")
print()
# Aggregate data: key = (app_id, app_name, domain), value = {bytes, flow_count, organization}
aggregated = defaultdict(lambda: {'bytes': 0, 'flow_count': 0, 'organization': ''})
total_files_processed = 0
total_records = 0
for filepath in sorted(traffic_files):
records = parse_traffic_file(filepath)
if records:
total_files_processed += 1
total_records += len(records)
rel_path = os.path.relpath(filepath, base_dir)
# print(f" Processed: {rel_path} ({len(records)} records)")
for app_id, app_name, domain, traffic_bytes, organization in records:
key = (app_id, app_name, domain)
aggregated[key]['bytes'] += traffic_bytes
aggregated[key]['flow_count'] += 1
# Keep the first non-empty organization
if organization and not aggregated[key]['organization']:
aggregated[key]['organization'] = organization
print()
print(f"Total files processed: {total_files_processed}")
print(f"Total raw records: {total_records}")
print(f"Total aggregated entries: {len(aggregated)}")
# Pre-compute traffic totals per app (all entries and domain-only entries)
app_total_bytes = defaultdict(int)
app_domain_total_bytes = defaultdict(int)
print(f"start calculate traffic ratios per app")
for (app_id, app_name, domain), data in aggregated.items():
app_total_bytes[(app_id, app_name)] += data['bytes']
if not domain.startswith('model_data:'):
app_domain_total_bytes[(app_id, app_name)] += data['bytes']
# Sort by app_id, then by traffic bytes descending within each app
sorted_entries = sorted(
aggregated.items(),
key=lambda x: (x[0][0], -x[1]['bytes'])
)
print(f"end calculate traffic ratios per app")
# Write output CSV
os.makedirs(os.path.dirname(output_file), exist_ok=True)
print(f"start write output CSV")
with open(output_file, 'w', newline='', encoding='utf-8') as f:
writer = csv.writer(f)
writer.writerow([
'App ID', 'App Name', 'Domain',
'Traffic Size(Bytes)', 'Flow Count',
'Traffic Ratio', 'Domain Traffic Ratio', 'Organization'
])
for (app_id, app_name, domain), data in sorted_entries:
total_bytes = data['bytes']
flow_count = data['flow_count']
app_total = app_total_bytes[(app_id, app_name)]
# Traffic Ratio = this entry's bytes / total bytes for this app
if app_total > 0:
traffic_ratio = total_bytes / app_total * 100
traffic_ratio_str = f"{traffic_ratio:.2f}%"
else:
traffic_ratio_str = "0.00%"
# Domain Traffic Ratio: same as traffic ratio for domain entries,
# empty for model_data entries
if domain.startswith('model_data:'):
domain_traffic_ratio_str = ''
else:
app_domain_total = app_domain_total_bytes[(app_id, app_name)]
if app_domain_total > 0:
domain_ratio = total_bytes / app_domain_total * 100
domain_traffic_ratio_str = f"{domain_ratio:.2f}%"
else:
domain_traffic_ratio_str = "0.00%"
writer.writerow([
app_id, app_name, domain,
total_bytes, flow_count,
traffic_ratio_str, domain_traffic_ratio_str,
data['organization']
])
print(f"\nOutput written to: {output_file}")
def main():
parser = argparse.ArgumentParser(
description='Merge traffic_count files into traffic_summary.csv'
)
parser.add_argument(
'--input-dir',
default=r'\\192.168.2.75\dpi-sync\autool_config\data\traffic_data',
help='Root directory containing PC subdirectories with traffic data'
)
parser.add_argument(
'--output',
default=None,
help='Output CSV file path (default: TrafficData/traffic_summary/<date>/traffic_summary.csv)'
)
args = parser.parse_args()
# Default output path with date
if args.output is None:
script_dir = os.path.dirname(os.path.abspath(__file__))
project_dir = os.path.dirname(script_dir)
date_str = datetime.now().strftime('%Y%m%d')
# date_str = '20260224'
output_file = os.path.join(
project_dir, 'TrafficData', 'traffic_summary', date_str, 'traffic_summary.csv'
)
else:
output_file = args.output
merge_traffic_data(args.input_dir, output_file)
if __name__ == '__main__':
main()