Fix FRP payload detection, byte accounting, and cloud-IP correlation in flock_tap #1
+67
-35
@@ -20,6 +20,7 @@ import os
|
|||||||
import sys
|
import sys
|
||||||
import json
|
import json
|
||||||
import time
|
import time
|
||||||
|
import signal
|
||||||
import threading
|
import threading
|
||||||
import subprocess
|
import subprocess
|
||||||
from datetime import datetime
|
from datetime import datetime
|
||||||
@@ -44,8 +45,6 @@ FLOCK_CLOUD_IPS = [
|
|||||||
# Scapy availability
|
# Scapy availability
|
||||||
try:
|
try:
|
||||||
from scapy.all import sniff, IP, TCP, UDP, DNS, DNSQR, Raw, conf
|
from scapy.all import sniff, IP, TCP, UDP, DNS, DNSQR, Raw, conf
|
||||||
from scapy.layers.tls.all import TLS
|
|
||||||
from scapy.layers.tls.handshake import TLSClientHello
|
|
||||||
HAVE_SCAPY = True
|
HAVE_SCAPY = True
|
||||||
except ImportError:
|
except ImportError:
|
||||||
HAVE_SCAPY = False
|
HAVE_SCAPY = False
|
||||||
@@ -112,6 +111,10 @@ class FlockTrafficTap:
|
|||||||
|
|
||||||
self.seen_domains = set()
|
self.seen_domains = set()
|
||||||
self.frp_ports = {7000, 7500, 7001, 7002}
|
self.frp_ports = {7000, 7500, 7001, 7002}
|
||||||
|
# Known Flock cloud IPs: static seed list + IPs learned from DNS answers
|
||||||
|
# for Flock domains. Used to catch cameras that talk to cloud IPs without
|
||||||
|
# exposing an SNI or issuing their own DNS query.
|
||||||
|
self.known_flock_ips = set(FLOCK_CLOUD_IPS)
|
||||||
self.running = False
|
self.running = False
|
||||||
self.packet_count = 0
|
self.packet_count = 0
|
||||||
self.start_time = None
|
self.start_time = None
|
||||||
@@ -121,15 +124,23 @@ class FlockTrafficTap:
|
|||||||
# ═══════════════════════════════════════════════════
|
# ═══════════════════════════════════════════════════
|
||||||
|
|
||||||
def on_dns_query(self, hostname, src_ip, resolved_ips=None, timestamp=None):
|
def on_dns_query(self, hostname, src_ip, resolved_ips=None, timestamp=None):
|
||||||
"""Called for every DNS query. Updates flow_stats."""
|
"""Called for every DNS query. Updates flow_stats and learns Flock IPs."""
|
||||||
if "flocksafety" in hostname.lower():
|
h = hostname.lower()
|
||||||
|
if "flocksafety" in h or "auth0.com" in h:
|
||||||
self.flow_stats[src_ip]["cloud_dns"] += 1
|
self.flow_stats[src_ip]["cloud_dns"] += 1
|
||||||
|
# Learn the resolved addresses so we can later flag cameras that
|
||||||
|
# connect straight to these IPs without their own DNS/SNI.
|
||||||
|
for rip in (resolved_ips or []):
|
||||||
|
self.known_flock_ips.add(rip)
|
||||||
|
|
||||||
def on_tcp_connect(self, src, dst, sport, dport, timestamp=None):
|
def on_tcp_connect(self, src, dst, sport, dport, timestamp=None):
|
||||||
"""Called for every TCP SYN. Updates flow_stats."""
|
"""Called for every TCP packet. Updates flow_stats.
|
||||||
|
|
||||||
|
Byte accounting is handled from real ``ip.len`` in the packet handler;
|
||||||
|
this only tracks connection signals so it does not inflate bandwidth.
|
||||||
|
"""
|
||||||
fs = self.flow_stats[src]
|
fs = self.flow_stats[src]
|
||||||
fs["connections"] += 1
|
fs["connections"] += 1
|
||||||
fs["bytes_up"] += 64 # approximate
|
|
||||||
if dport in self.frp_ports:
|
if dport in self.frp_ports:
|
||||||
fs["frp_tunnel"] = True
|
fs["frp_tunnel"] = True
|
||||||
|
|
||||||
@@ -166,6 +177,8 @@ class FlockTrafficTap:
|
|||||||
return "CLOUD_CONNECTED"
|
return "CLOUD_CONNECTED"
|
||||||
if stats.get("s3_uploads", 0) > 0 or stats.get("auth_tls", 0) > 0:
|
if stats.get("s3_uploads", 0) > 0 or stats.get("auth_tls", 0) > 0:
|
||||||
return "CLOUD_CONNECTED"
|
return "CLOUD_CONNECTED"
|
||||||
|
if stats.get("cloud_api", 0) > 0:
|
||||||
|
return "CLOUD_CONNECTED"
|
||||||
if stats.get("connections", 0) > 0:
|
if stats.get("connections", 0) > 0:
|
||||||
return "LOCAL_STATION"
|
return "LOCAL_STATION"
|
||||||
return "OFFLINE_OR_UNMONITORED"
|
return "OFFLINE_OR_UNMONITORED"
|
||||||
@@ -255,46 +268,53 @@ class FlockTrafficTap:
|
|||||||
|
|
||||||
self.on_tcp_connect(src, dst, sport, dport, ts)
|
self.on_tcp_connect(src, dst, sport, dport, ts)
|
||||||
|
|
||||||
|
# ── Connection to a known Flock cloud IP (no SNI/DNS needed) ──
|
||||||
|
if dst in self.known_flock_ips:
|
||||||
|
self.flow_stats[src]["cloud_api"] += 1
|
||||||
|
|
||||||
if tcp.flags & 0x02: # SYN
|
if tcp.flags & 0x02: # SYN
|
||||||
dev["connections"].append({
|
dev["connections"].append({
|
||||||
"timestamp": ts, "src": src, "sport": sport,
|
"timestamp": ts, "src": src, "sport": sport,
|
||||||
"dst": dst, "dport": dport, "bytes": ip.len,
|
"dst": dst, "dport": dport, "bytes": ip.len,
|
||||||
})
|
})
|
||||||
|
|
||||||
# ── FRP Tunnel Detection ──
|
# ── FRP Tunnel Detection ──
|
||||||
# FRP handshake pattern:
|
# FRP handshake pattern:
|
||||||
# 1. Camera connects to port 7000-7500 on a remote server
|
# 1. Camera connects to port 7000-7500 on a remote server
|
||||||
# 2. Sends auth + proxy configuration JSON
|
# 2. Sends auth + proxy configuration JSON
|
||||||
# 3. Server opens reverse tunnel
|
# 3. Server opens reverse tunnel
|
||||||
#
|
#
|
||||||
# Detection signature (Zeek-equivalent):
|
# Detection signature (Zeek-equivalent):
|
||||||
# signature frp-tunnel {
|
# signature frp-tunnel {
|
||||||
# ip-proto == tcp
|
# ip-proto == tcp
|
||||||
# dst-port in [7000, 7500, 7001, 7002]
|
# dst-port in [7000, 7500, 7001, 7002]
|
||||||
# payload /frp|auth|proxy_type/
|
# payload /frp|auth|proxy_type/
|
||||||
# event "FRP TUNNEL DETECTED"
|
# event "FRP TUNNEL DETECTED"
|
||||||
# }
|
# }
|
||||||
if dport in self.frp_ports:
|
# Port-based signal is evaluated per-packet (fires on the SYN too).
|
||||||
|
if dport in self.frp_ports:
|
||||||
|
self.flow_stats[src]["frp_tunnel"] = True
|
||||||
|
dev["frp_tunnels"].append({
|
||||||
|
"timestamp": ts, "src": src, "dst": dst,
|
||||||
|
"port": dport, "type": "frp_tunnel",
|
||||||
|
})
|
||||||
|
if self.verbose:
|
||||||
|
print(f" {C.R}FRP {src:<16} -> {dst}:{dport} [FRP TUNNEL]{C.END}")
|
||||||
|
|
||||||
|
# Raw payload scan for FRP auth. The auth/proxy JSON is sent in a
|
||||||
|
# data segment *after* the handshake, so this must run on every TCP
|
||||||
|
# packet — not only on the SYN (SYN packets carry no payload).
|
||||||
|
if tcp.haslayer(Raw):
|
||||||
|
raw = tcp[Raw].load
|
||||||
|
if b"frp" in raw.lower() or b"auth" in raw.lower() or b"proxy_type" in raw.lower():
|
||||||
self.flow_stats[src]["frp_tunnel"] = True
|
self.flow_stats[src]["frp_tunnel"] = True
|
||||||
dev["frp_tunnels"].append({
|
dev["frp_tunnels"].append({
|
||||||
"timestamp": ts, "src": src, "dst": dst,
|
"timestamp": ts, "src": src, "dst": dst,
|
||||||
"port": dport, "type": "frp_tunnel",
|
"port": dport, "type": "frp_auth_payload",
|
||||||
|
"payload_preview": raw[:64].decode(errors="replace"),
|
||||||
})
|
})
|
||||||
if self.verbose:
|
if self.verbose:
|
||||||
print(f" {C.R}FRP {src:<16} -> {dst}:{dport} [FRP TUNNEL]{C.END}")
|
self.log(f"{C.R}FRP_AUTH {src} -> {dst}:{dport} - auth payload{C.END}")
|
||||||
|
|
||||||
# Raw payload scan for FRP auth
|
|
||||||
if tcp.haslayer(Raw):
|
|
||||||
raw = tcp[Raw].load
|
|
||||||
if b"frp" in raw.lower() or b"auth" in raw.lower() or b"proxy_type" in raw.lower():
|
|
||||||
self.flow_stats[src]["frp_tunnel"] = True
|
|
||||||
dev["frp_tunnels"].append({
|
|
||||||
"timestamp": ts, "src": src, "dst": dst,
|
|
||||||
"port": dport, "type": "frp_auth_payload",
|
|
||||||
"payload_preview": raw[:64].decode(errors="replace"),
|
|
||||||
})
|
|
||||||
if self.verbose:
|
|
||||||
self.log(f"{C.R}FRP_AUTH {src} -> {dst}:{dport} - auth payload{C.END}")
|
|
||||||
|
|
||||||
# ── TLS SNI ──
|
# ── TLS SNI ──
|
||||||
if pkt.haslayer(TCP) and pkt.haslayer(Raw):
|
if pkt.haslayer(TCP) and pkt.haslayer(Raw):
|
||||||
@@ -394,10 +414,22 @@ class FlockTrafficTap:
|
|||||||
ts = datetime.now().strftime("%H:%M:%S")
|
ts = datetime.now().strftime("%H:%M:%S")
|
||||||
print(f"{C.CY}[{ts}]{C.END} {msg}")
|
print(f"{C.CY}[{ts}]{C.END} {msg}")
|
||||||
|
|
||||||
|
def _install_sigint(self):
|
||||||
|
"""Stop capture cleanly on Ctrl+C so the report is always produced."""
|
||||||
|
def _handler(signum, frame):
|
||||||
|
self.running = False
|
||||||
|
raise KeyboardInterrupt
|
||||||
|
try:
|
||||||
|
signal.signal(signal.SIGINT, _handler)
|
||||||
|
except (ValueError, RuntimeError):
|
||||||
|
# Not on the main thread — fall back to scapy's own handling.
|
||||||
|
pass
|
||||||
|
|
||||||
def start(self):
|
def start(self):
|
||||||
"""Start capture in the foreground."""
|
"""Start capture in the foreground."""
|
||||||
self.start_time = time.time()
|
self.start_time = time.time()
|
||||||
self.running = True
|
self.running = True
|
||||||
|
self._install_sigint()
|
||||||
if self.pcap:
|
if self.pcap:
|
||||||
self._run_pcap()
|
self._run_pcap()
|
||||||
elif self.pipe:
|
elif self.pipe:
|
||||||
|
|||||||
Reference in New Issue
Block a user