///|
/// 连接统计
pub struct ConnStats {
packets : Int
bytes_in : Int64
bytes_out : Int64
first_seen : Int64
last_seen : Int64
} derive(Default)
///|
/// 连接跟踪器:五元组 → 统计(双向流量归同一连接)
pub struct ConnTracker {
conns : @hashmap.HashMap[FiveTuple, ConnStats]
idle_timeout : Int64 // 空闲清理阈值(秒)
}
///|
/// 新建跟踪器;idle_timeout 秒无流量后 cleanup 会移除连接
pub fn ConnTracker::new(idle_timeout : Int64) -> ConnTracker {
{ conns: @hashmap.HashMap([]), idle_timeout }
}
///|
/// 五元组归一化:按 (src_ip, src_port) < (dst_ip, dst_port) 排序作 key
fn normalize(t : FiveTuple) -> FiveTuple {
if t.src_ip < t.dst_ip || (t.src_ip == t.dst_ip && t.src_port <= t.dst_port) {
t
} else {
{
src_ip: t.dst_ip,
dst_ip: t.src_ip,
src_port: t.dst_port,
dst_port: t.src_port,
proto: t.proto,
}
}
}
///|
/// 更新连接统计;Local(回环)流量计入入方向
pub fn ConnTracker::update(self : ConnTracker, pkt : ParsedPacket) -> Unit {
let key = normalize(pkt.tuple)
let stats = self.conns.get_or_default(key, {
packets: 0,
bytes_in: 0L,
bytes_out: 0L,
first_seen: 0L,
last_seen: 0L,
})
let is_new = stats.packets == 0
let add = pkt.payload_len.to_int64()
let (bytes_in, bytes_out) = match pkt.direction {
In | Local => (stats.bytes_in + add, stats.bytes_out)
Out => (stats.bytes_in, stats.bytes_out + add)
}
self.conns[key] = {
packets: stats.packets + 1,
bytes_in,
bytes_out,
first_seen: if is_new {
pkt.ts_sec
} else {
stats.first_seen
},
last_seen: pkt.ts_sec,
}
}
///|
/// 连接快照(含归一化后的五元组)
pub fn ConnTracker::snapshot(
self : ConnTracker,
) -> Array[(FiveTuple, ConnStats)] {
self.conns.iter().to_array()
}
///|
/// 清理空闲连接(last_seen 距 now 超过 idle_timeout 秒)
pub fn ConnTracker::cleanup(self : ConnTracker, now : Int64) -> Unit {
let stale = self.conns
.iter()
.filter_map(fn(kv) {
if now - kv.1.last_seen > self.idle_timeout {
Some(kv.0)
} else {
None
}
})
stale.each(fn(k) { self.conns.remove(k) })
}
///|
pub fn ConnStats::to_string(self : ConnStats) -> String {
"\{self.packets}p in=\{self.bytes_in} out=\{self.bytes_out} age=\{(self.last_seen - self.first_seen)}s"
}
///|
/// 每秒一行 JSON 快照:{ts, conns:[{src_ip,src_port,dst_ip,dst_port,proto,packets,bytes_in,bytes_out,first_seen,last_seen}]}
/// 按流量(入+出)降序排列
pub fn ConnTracker::snapshot_json(self : ConnTracker, ts : Int64) -> String {
let snap = self.snapshot()
snap.sort_by(fn(a, b) {
let ta = a.1.bytes_in + a.1.bytes_out
let tb = b.1.bytes_in + b.1.bytes_out
if ta > tb {
-1
} else if ta < tb {
1
} else {
0
}
})
let conns = snap.map(fn(kv) { conn_json(kv.0, kv.1) })
Json::object(Map([("ts", json_num64(ts)), ("conns", Json::array(conns))])).stringify()
}
///|
fn conn_json(t : FiveTuple, s : ConnStats) -> Json {
Json::object(
Map([
("src_ip", Json::string(t.src_ip.to_string())),
("src_port", json_num(t.src_port)),
("dst_ip", Json::string(t.dst_ip.to_string())),
("dst_port", json_num(t.dst_port)),
("proto", Json::string(t.proto.to_string())),
("packets", json_num(s.packets)),
("bytes_in", json_num64(s.bytes_in)),
("bytes_out", json_num64(s.bytes_out)),
("first_seen", json_num64(s.first_seen)),
("last_seen", json_num64(s.last_seen)),
]),
)
}
///|
/// Int 转 JSON 数字(repr 保留整数格式,避免科学计数法)
fn json_num(v : Int) -> Json {
Json::number(v.to_double(), repr=v.to_string())
}
///|
fn json_num64(v : Int64) -> Json {
Json::number(v.to_double(), repr=v.to_string())
}