5 Commits

Author SHA1 Message Date
leetcrypt 51ba69a541 flock_tap: audit-integrity + data-minimization for cloud-vs-local verdicts
Makes the passive tap a defensible auditing instrument for the
"does this camera phone home to Flock?" question, while shrinking what it
retains about bystanders. Nothing here widens third-party capture.

- Unified verdict vocabulary. A single private _classify() folds flow-stat
  and device-level signals into exactly one of CLOUD_CONNECTED /
  LOCAL_STATION / INDETERMINATE. classify_device/classify_camera become thin
  wrappers. Previously the two classifiers returned divergent sets
  (OFFLINE_OR_UNMONITORED, NO_DATA, UNKNOWN) that fell out of every summary
  bucket, so totals never reconciled.
- Report counts now reconcile. Summary buckets partition every tracked IP;
  counts_reconcile asserts cloud+local+indeterminate == total_ips_tracked.
- Explainable cloud verdicts. Record known-Flock IPs a device contacted
  (cloud_ip_contacts) and surface CLOUD_IP(...) evidence, plus a
  TLS_CATEGORY(...) fallback, so a cloud verdict is never unexplained.
- Bounded captures. Per-device raw sample lists cap at max_samples
  (default 500, --tap-max-samples). Exact integer counters increment
  regardless, so verdict inputs and report totals stay precise past the cap.
- Bystander protection. FRP auth detection stores a non-content descriptor
  (matched keyword + payload length) instead of 64 raw payload bytes.
  Optional --tap-redact-ips salt-hashes non-Flock destination IPs in stored
  samples; known-Flock IPs stay legible so cloud evidence is readable.
- Metadata-only destination_summary per cloud device (DNS names, TLS SNIs,
  known-Flock IPs) — who it talks to, never payload contents.

Tests extend the merged suite with audit-integrity coverage (vocabulary
agreement, bucket reconciliation, indeterminate counting, CLOUD_IP evidence,
sample capping vs exact counters, FRP descriptor, IP redaction,
metadata-only summary). Stale-vocabulary assertions updated to INDETERMINATE.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-07-22 09:04:06 -07:00
ek0ms savi0r 4c605b3093 Merge pull request 'Add flock_tap test suite and .gitignore' (#2) from Trilltechnician/Flock_SCAN-fix:tests-hygiene/flock-tap into main
Reviewed-on: ek0mssavi0r/Flock_SCAN#2
2026-07-21 21:08:20 +00:00
leetcrypt 4081e55ba8 Add flock_tap test suite and .gitignore
Adds pytest coverage for flock_tap.py's parsing, classification, and
byte/FRP/cloud-IP accounting, plus a .gitignore for Python caches and
capture/scan output artifacts.

The packet-level tests double as regression guards for the recently fixed
detection bugs (FRP auth-payload firing on data segments not just SYN,
no bytes_up inflation, and known-Flock-IP correlation); they fail against
the pre-fix code and pass after it. scapy-dependent tests skip cleanly when
scapy is unavailable.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-07-20 08:57:44 -07:00
ek0ms savi0r 0817ab9713 Merge pull request 'Fix FRP payload detection, byte accounting, and cloud-IP correlation in flock_tap' (#1) from Trilltechnician/Flock_SCAN-fix:fix/flock-tap-detection-bugs into main
Reviewed-on: ek0mssavi0r/Flock_SCAN#1
2026-07-20 06:05:02 +00:00
leetcrypt 20d30815a3 Fix FRP payload detection, byte accounting, and cloud-IP correlation in flock_tap
- Move the FRP auth-payload scan out of the SYN-only branch. SYN packets carry
  no payload, so payload-based FRP detection never fired; it now runs on every
  TCP segment where the frp/auth/proxy_type bytes actually appear.
- Remove the per-packet `bytes_up += 64` approximation in on_tcp_connect, which
  double-counted against real ip.len accounting and corrupted bandwidth stats.
- Correlate connections against known Flock cloud IPs (static seed list + IPs
  learned from DNS answers) so cameras that reach cloud IPs without their own
  DNS/SNI are classified CLOUD_CONNECTED.
- Match the Flock auth0 tenant in DNS detection, consistent with the SNI path.
- Install a SIGINT handler so the report is always produced on Ctrl+C.
- Drop unused scapy TLS imports (SNI is parsed manually).

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-07-19 22:25:40 -07:00
4 changed files with 583 additions and 112 deletions
+13
View File
@@ -0,0 +1,13 @@
# Python
__pycache__/
*.py[cod]
*.egg-info/
.pytest_cache/
.venv/
venv/
# Capture / scan outputs
*.pcap
report.json
creds_*.txt
s3_urls_*.txt
+232 -112
View File
@@ -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
@@ -41,11 +42,15 @@ FLOCK_CLOUD_IPS = [
"104.18.16.189", "104.18.17.189", "104.18.16.189", "104.18.17.189",
] ]
# ── Canonical verdict vocabulary (single source of truth) ──
# Every device maps to exactly one of these so audit summary counts reconcile.
VERDICT_CLOUD = "CLOUD_CONNECTED"
VERDICT_LOCAL = "LOCAL_STATION"
VERDICT_INDETERMINATE = "INDETERMINATE"
# 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
@@ -73,7 +78,8 @@ class FlockTrafficTap:
""" """
def __init__(self, interface=None, pcap=None, pipe=False, verbose=False, def __init__(self, interface=None, pcap=None, pipe=False, verbose=False,
output_file=None, timeout=None, filter_expr=None): output_file=None, timeout=None, filter_expr=None,
max_samples=500, redact_non_flock_ips=False):
self.interface = interface self.interface = interface
self.pcap = pcap self.pcap = pcap
self.pipe = pipe self.pipe = pipe
@@ -82,6 +88,16 @@ class FlockTrafficTap:
self.timeout = timeout self.timeout = timeout
self.filter_expr = filter_expr self.filter_expr = filter_expr
# Data-minimization knobs (see AUDIT_INTEGRITY_SPEC.md).
# max_samples caps per-device raw sample lists so a long audit does not
# accrete an unbounded pile of raw records; integer counters below stay
# exact regardless. redact_non_flock_ips salt-hashes non-Flock dst IPs
# in stored samples to protect bystanders (off by default: hashing
# removes reproducibility a Flock-side auditor may need).
self.max_samples = max_samples
self.redact_non_flock_ips = redact_non_flock_ips
self._redact_salt = os.urandom(16)
# ── flow_stats: matches pseudocode structure ── # ── flow_stats: matches pseudocode structure ──
self.flow_stats = defaultdict(lambda: { self.flow_stats = defaultdict(lambda: {
"cloud_dns": 0, "cloud_dns": 0,
@@ -103,6 +119,10 @@ class FlockTrafficTap:
"frp_tunnels": [], "frp_tunnels": [],
"tls_snis": [], "tls_snis": [],
"http_requests": [], "http_requests": [],
"cloud_ip_contacts": [],
# Exact running totals — always incremented even after the sample
# lists above stop growing at max_samples, so counts/verdicts hold.
"counts": {"dns": 0, "conn": 0, "frp": 0, "sni": 0, "flock_dns": 0},
"bytes_up": 0, "bytes_up": 0,
"bytes_down": 0, "bytes_down": 0,
"packets_seen": 0, "packets_seen": 0,
@@ -112,6 +132,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 +145,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
@@ -159,36 +191,57 @@ class FlockTrafficTap:
# CLASSIFICATION — matches pseudocode interface # CLASSIFICATION — matches pseudocode interface
# ═══════════════════════════════════════════════════ # ═══════════════════════════════════════════════════
def _classify(self, ip):
"""Single source of truth for a device verdict.
Folds flow-stat signals and device-level raw fallback into exactly one
canonical verdict so audit summary counts always reconcile.
"""
stats = self.flow_stats[ip]
# Cloud signals from the per-flow counters.
if (stats.get("cloud_dns", 0) > 0 or stats.get("frp_tunnel")
or stats.get("s3_uploads", 0) > 0 or stats.get("auth_tls", 0) > 0
or stats.get("cloud_api", 0) > 0):
return VERDICT_CLOUD
# Cloud signals from the raw device store (fallback when flow_stats is
# sparse — e.g. SNI/FRP seen but no counter set).
dev = self.devices.get(ip, {})
if dev:
has_flock_sni = any(
"flock" in s.get("sni", "").lower() or "auth0" in s.get("sni", "").lower()
for s in dev.get("tls_snis", [])
)
if has_flock_sni or dev.get("frp_tunnels") or dev.get("cloud_ip_contacts"):
return VERDICT_CLOUD
# Traffic seen but no cloud signal → stays on-prem.
if stats.get("connections", 0) > 0 or (dev and dev.get("connections")):
return VERDICT_LOCAL
return VERDICT_INDETERMINATE
# Public names kept as thin wrappers so existing callers/tests stay valid.
def classify_device(self, camera_ip): def classify_device(self, camera_ip):
"""Generate a full traffic profile for a camera (from pseudocode).""" return self._classify(camera_ip)
stats = self.flow_stats[camera_ip]
if stats.get("cloud_dns", 0) > 0 or stats.get("frp_tunnel"):
return "CLOUD_CONNECTED"
if stats.get("s3_uploads", 0) > 0 or stats.get("auth_tls", 0) > 0:
return "CLOUD_CONNECTED"
if stats.get("connections", 0) > 0:
return "LOCAL_STATION"
return "OFFLINE_OR_UNMONITORED"
def classify_camera(self, ip): def classify_camera(self, ip):
"""Return verdict for a single camera IP (raw data fallback).""" return self._classify(ip)
verdict = self.classify_device(ip)
if verdict != "OFFLINE_OR_UNMONITORED":
return verdict
dev = self.devices.get(ip, {}) # ── data-minimization helpers ──
if not dev: def _append_capped(self, lst, item):
return "NO_DATA" """Append a raw sample only while under the cap. Counters (kept
has_flock_sni = any( separately) stay exact, so bounding samples never skews verdicts."""
"flock" in s.get("sni", "").lower() or "auth0" in s.get("sni", "").lower() if len(lst) < self.max_samples:
for s in dev.get("tls_snis", []) lst.append(item)
)
has_frp = len(dev.get("frp_tunnels", [])) > 0 def _redact_ip(self, ip):
if has_flock_sni or has_frp: """Salt-hash a non-Flock IP for bystander protection. Known-Flock IPs
return "CLOUD_CONNECTED" pass through unchanged so cloud evidence stays legible."""
if dev.get("connections"): if not self.redact_non_flock_ips or ip in self.known_flock_ips:
return "LOCAL_STATION" return ip
return "UNKNOWN" import hashlib
return "redacted:" + hashlib.sha256(self._redact_salt + ip.encode()).hexdigest()[:12]
# ═══════════════════════════════════════════════════ # ═══════════════════════════════════════════════════
# PACKET HANDLER — scapy # PACKET HANDLER — scapy
@@ -233,7 +286,10 @@ class FlockTrafficTap:
pass pass
self.on_dns_query(qname, src, resolved_ips, ts) self.on_dns_query(qname, src, resolved_ips, ts)
dev["dns_queries"].append({ dev["counts"]["dns"] += 1
if "flock" in qname.lower() or "auth0" in qname.lower():
dev["counts"]["flock_dns"] += 1
self._append_capped(dev["dns_queries"], {
"timestamp": ts, "src": src, "timestamp": ts, "src": src,
"query": qname, "resolved_ips": resolved_ips, "query": qname, "resolved_ips": resolved_ips,
}) })
@@ -255,46 +311,64 @@ 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 dst not in dev["cloud_ip_contacts"]:
dev["cloud_ip_contacts"].append(dst)
if tcp.flags & 0x02: # SYN if tcp.flags & 0x02: # SYN
dev["connections"].append({ dev["counts"]["conn"] += 1
self._append_capped(dev["connections"], {
"timestamp": ts, "src": src, "sport": sport, "timestamp": ts, "src": src, "sport": sport,
"dst": dst, "dport": dport, "bytes": ip.len, "dst": self._redact_ip(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["counts"]["frp"] += 1
self._append_capped(dev["frp_tunnels"], {
"timestamp": ts, "src": src, "dst": self._redact_ip(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
low = raw.lower()
matched = next((k for k in (b"frp", b"auth", b"proxy_type") if k in low), None)
if matched:
self.flow_stats[src]["frp_tunnel"] = True self.flow_stats[src]["frp_tunnel"] = True
dev["frp_tunnels"].append({ dev["counts"]["frp"] += 1
"timestamp": ts, "src": src, "dst": dst, # Store a non-content descriptor, not the raw bytes: which
"port": dport, "type": "frp_tunnel", # keyword matched + payload length. Proves FRP detection
# without retaining any of the payload (bystander protection).
self._append_capped(dev["frp_tunnels"], {
"timestamp": ts, "src": src, "dst": self._redact_ip(dst),
"port": dport, "type": "frp_auth_payload",
"matched_keyword": matched.decode(),
"payload_len": len(raw),
}) })
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):
@@ -305,7 +379,11 @@ class FlockTrafficTap:
sni = self._extract_sni(raw) sni = self._extract_sni(raw)
if sni: if sni:
self.on_tls_sni(sni, src, dst, ts) self.on_tls_sni(sni, src, dst, ts)
dev["tls_snis"].append({"timestamp": ts, "src": src, "dst": dst, "sni": sni}) dev["counts"]["sni"] += 1
self._append_capped(dev["tls_snis"], {
"timestamp": ts, "src": src,
"dst": self._redact_ip(dst), "sni": sni,
})
if self.verbose: if self.verbose:
cat = self._categorize_sni(sni) cat = self._categorize_sni(sni)
c = C.R if cat else C.CY c = C.R if cat else C.CY
@@ -394,10 +472,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:
@@ -459,14 +549,18 @@ class FlockTrafficTap:
if parsed["type"] == "dns": if parsed["type"] == "dns":
query = parsed.get("query", "") query = parsed.get("query", "")
if query: if query:
dev["dns_queries"].append({ dev["counts"]["dns"] += 1
is_flock = "flock" in query.lower() or "auth0" in query.lower()
if is_flock:
dev["counts"]["flock_dns"] += 1
self._append_capped(dev["dns_queries"], {
"timestamp": time.time(), "timestamp": time.time(),
"src": src, "src": src,
"query": query, "query": query,
"resolved_ips": [], "resolved_ips": [],
}) })
self.seen_domains.add(query) self.seen_domains.add(query)
if "flock" in query.lower(): if is_flock:
self.flow_stats[src]["cloud_dns"] += 1 self.flow_stats[src]["cloud_dns"] += 1
if self.verbose: if self.verbose:
self.log(f"{C.R}Flock DNS {src} -> {query}{C.END}") self.log(f"{C.R}Flock DNS {src} -> {query}{C.END}")
@@ -474,9 +568,10 @@ class FlockTrafficTap:
self.log(f"{C.CY}DNS {src} -> {query}{C.END}") self.log(f"{C.CY}DNS {src} -> {query}{C.END}")
if parsed["type"] == "syn" and dport in self.frp_ports: if parsed["type"] == "syn" and dport in self.frp_ports:
self.flow_stats[src]["frp_tunnel"] = True self.flow_stats[src]["frp_tunnel"] = True
dev["frp_tunnels"].append({ dev["counts"]["frp"] += 1
self._append_capped(dev["frp_tunnels"], {
"timestamp": time.time(), "src": src, "timestamp": time.time(), "src": src,
"dst": dst, "port": dport, "type": "frp_tunnel", "dst": self._redact_ip(dst), "port": dport, "type": "frp_tunnel",
}) })
if self.verbose: if self.verbose:
self.log(f"{C.R}FRP {src} -> {dst}:{dport}{C.END}") self.log(f"{C.R}FRP {src} -> {dst}:{dport}{C.END}")
@@ -504,18 +599,12 @@ class FlockTrafficTap:
"summary": {}, "summary": {},
} }
cloud_count = 0 counts = {VERDICT_CLOUD: 0, VERDICT_LOCAL: 0, VERDICT_INDETERMINATE: 0}
local_count = 0
unknown_count = 0
for ip in sorted(set(list(self.devices.keys()) + list(self.flow_stats.keys()))): all_ips = sorted(set(list(self.devices.keys()) + list(self.flow_stats.keys())))
verdict = self.classify_camera(ip) for ip in all_ips:
if verdict == "CLOUD_CONNECTED": verdict = self._classify(ip)
cloud_count += 1 counts[verdict] += 1
elif verdict == "LOCAL_STATION":
local_count += 1
else:
unknown_count += 1
# Always include flow_stats # Always include flow_stats
report["flow_stats"][ip] = dict(self.flow_stats[ip]) report["flow_stats"][ip] = dict(self.flow_stats[ip])
@@ -524,8 +613,10 @@ class FlockTrafficTap:
if dev.get("packets_seen", 0) < 5: if dev.get("packets_seen", 0) < 5:
continue continue
snis = list(set(s["sni"] for s in dev.get("tls_snis", []))) c = dev.get("counts", {})
dns_list = list(set(q["query"] for q in dev.get("dns_queries", []))) snis = sorted(set(s["sni"] for s in dev.get("tls_snis", [])))
dns_list = sorted(set(q["query"] for q in dev.get("dns_queries", [])))
cloud_ips = list(dev.get("cloud_ip_contacts", []))
report["devices"][ip] = { report["devices"][ip] = {
"verdict": verdict, "verdict": verdict,
@@ -534,31 +625,44 @@ class FlockTrafficTap:
"bytes_down": dev["bytes_down"], "bytes_down": dev["bytes_down"],
"first_seen": dev["first_seen"], "first_seen": dev["first_seen"],
"last_seen": dev["last_seen"], "last_seen": dev["last_seen"],
"dns_queries_count": len(dev["dns_queries"]), # Exact counters (not len() of the capped sample lists).
"tls_snis_count": len(dev.get("tls_snis", [])), "dns_queries_count": c.get("dns", len(dev["dns_queries"])),
"frp_tunnels_count": len(dev.get("frp_tunnels", [])), "tls_snis_count": c.get("sni", len(dev.get("tls_snis", []))),
"connections_count": len(dev.get("connections", [])), "frp_tunnels_count": c.get("frp", len(dev.get("frp_tunnels", []))),
"connections_count": c.get("conn", len(dev.get("connections", []))),
"samples_capped_at": self.max_samples,
"flow_stats": dict(self.flow_stats[ip]), "flow_stats": dict(self.flow_stats[ip]),
"tls_snis": snis, # Metadata-only destination evidence — who the device talks to,
# never payload contents.
"destination_summary": {
"dns_domains": dns_list,
"tls_snis": snis,
"flock_cloud_ips": cloud_ips,
},
"frp_tunnels": dev.get("frp_tunnels", []), "frp_tunnels": dev.get("frp_tunnels", []),
"dns_domains": dns_list,
} }
report["summary"] = { report["summary"] = {
"total_ips_tracked": len(self.devices), "total_ips_tracked": len(all_ips),
"cloud_connected": cloud_count, "cloud_connected": counts[VERDICT_CLOUD],
"local_station": local_count, "local_station": counts[VERDICT_LOCAL],
"unknown": unknown_count, "indeterminate": counts[VERDICT_INDETERMINATE],
# Invariant: the three buckets partition every tracked IP.
"counts_reconcile": (
counts[VERDICT_CLOUD] + counts[VERDICT_LOCAL]
+ counts[VERDICT_INDETERMINATE] == len(all_ips)
),
"flock_domains_resolved": sorted( "flock_domains_resolved": sorted(
d for d in self.seen_domains if "flock" in d.lower() d for d in self.seen_domains
if "flock" in d.lower() or "auth0" in d.lower()
), ),
"total_flock_queries": sum( "total_flock_queries": sum(
1 for dev in self.devices.values() dev.get("counts", {}).get("flock_dns", 0)
for q in dev.get("dns_queries", []) for dev in self.devices.values()
if "flock" in q.get("query", "").lower()
), ),
"total_frp_tunnels": sum( "total_frp_tunnels": sum(
len(dev.get("frp_tunnels", [])) for dev in self.devices.values() dev.get("counts", {}).get("frp", 0)
for dev in self.devices.values()
), ),
} }
return report return report
@@ -601,15 +705,16 @@ class FlockTrafficTap:
print(f" Duration: {elapsed:.1f}s") print(f" Duration: {elapsed:.1f}s")
print(f" Packets: {self.packet_count:,}") print(f" Packets: {self.packet_count:,}")
cloud_cams = [] # Partition every tracked IP into exactly one bucket (counts reconcile).
for ip in sorted(self.devices): verdicts = {ip: self._classify(ip) for ip in self.devices}
if self.classify_camera(ip) == "CLOUD_CONNECTED": cloud_cams = [ip for ip, v in sorted(verdicts.items()) if v == VERDICT_CLOUD]
cloud_cams.append(ip) local_n = sum(1 for v in verdicts.values() if v == VERDICT_LOCAL)
indet_n = sum(1 for v in verdicts.values() if v == VERDICT_INDETERMINATE)
print(f"\n{C.CY}Classification:{C.END}") print(f"\n{C.CY}Classification:{C.END}")
print(f" Cloud-connected: {len(cloud_cams)}") print(f" Cloud-connected: {len(cloud_cams)}")
print(f" Local station: {sum(1 for ip in self.devices if self.classify_camera(ip) == 'LOCAL_STATION')}") print(f" Local station: {local_n}")
print(f" Unknown: {sum(1 for ip in self.devices if self.classify_camera(ip) == 'UNKNOWN')}") print(f" Indeterminate: {indet_n}")
if cloud_cams: if cloud_cams:
print(f"\n{C.R}Cloud-Connected Cameras:{C.END}") print(f"\n{C.R}Cloud-Connected Cameras:{C.END}")
@@ -617,17 +722,27 @@ class FlockTrafficTap:
dev = self.devices[ip] dev = self.devices[ip]
evidence = [] evidence = []
dns_list = list(set(q["query"] for q in dev.get("dns_queries", []) dns_list = list(set(q["query"] for q in dev.get("dns_queries", [])
if "flock" in q["query"].lower())) if "flock" in q["query"].lower() or "auth0" in q["query"].lower()))
snis = list(set(s["sni"] for s in dev.get("tls_snis", []) snis = list(set(s["sni"] for s in dev.get("tls_snis", [])
if "flock" in s["sni"].lower() or "auth0" in s["sni"].lower())) if "flock" in s["sni"].lower() or "auth0" in s["sni"].lower()))
frps = dev.get("frp_tunnels", []) frps = dev.get("frp_tunnels", [])
cloud_ips = dev.get("cloud_ip_contacts", [])
if dns_list: if dns_list:
evidence.append(f"DNS({', '.join(dns_list[:3])})") evidence.append(f"DNS({', '.join(dns_list[:3])})")
if snis: if snis:
evidence.append(f"SNI({', '.join(snis[:3])})") evidence.append(f"SNI({', '.join(snis[:3])})")
if frps: if frps:
evidence.append(f"FRP({len(frps)} tunnels)") evidence.append(f"FRP({dev.get('counts', {}).get('frp', len(frps))} tunnels)")
print(f" {C.R}[!] {ip:<16}{C.END} {' | '.join(evidence[:3])}") if cloud_ips:
evidence.append(f"CLOUD_IP({', '.join(cloud_ips[:3])})")
if not evidence:
# Cloud via an S3/auth/api SNI category counter with no
# retained hostname sample — name the signal so the verdict
# is never unexplained.
fs = self.flow_stats[ip]
sig = [k for k in ("s3_uploads", "auth_tls", "cloud_api") if fs.get(k)]
evidence.append(f"TLS_CATEGORY({', '.join(sig)})" if sig else "cloud signal")
print(f" {C.R}[!] {ip:<16}{C.END} {' | '.join(evidence[:4])}")
# All Flock domains seen # All Flock domains seen
flock_domains = sorted( flock_domains = sorted(
@@ -639,8 +754,8 @@ class FlockTrafficTap:
for d in flock_domains: for d in flock_domains:
print(f"{d}") print(f"{d}")
# FRP tunnels # FRP tunnels (exact counter; the per-tunnel list below is a capped sample)
total_frp = sum(len(dev.get("frp_tunnels", [])) for dev in self.devices.values()) total_frp = sum(dev.get("counts", {}).get("frp", 0) for dev in self.devices.values())
if total_frp > 0: if total_frp > 0:
print(f"\n{C.R}FRP Tunnels: {total_frp}{C.END}") print(f"\n{C.R}FRP Tunnels: {total_frp}{C.END}")
for ip, dev in sorted(self.devices.items()): for ip, dev in sorted(self.devices.items()):
@@ -697,12 +812,17 @@ def main():
parser.add_argument("--tap-verbose", "-v", action="store_true", help="Verbose output") parser.add_argument("--tap-verbose", "-v", action="store_true", help="Verbose output")
parser.add_argument("--tap-timeout", type=int, default=None, help="Capture duration (seconds)") parser.add_argument("--tap-timeout", type=int, default=None, help="Capture duration (seconds)")
parser.add_argument("--tap-filter", default=None, help="BPF filter expression") parser.add_argument("--tap-filter", default=None, help="BPF filter expression")
parser.add_argument("--tap-max-samples", type=int, default=500,
help="Cap on per-device raw samples kept (data minimization). Default 500.")
parser.add_argument("--tap-redact-ips", action="store_true",
help="Salt-hash non-Flock destination IPs in stored samples (bystander protection).")
args = parser.parse_args() args = parser.parse_args()
tap = FlockTrafficTap( tap = FlockTrafficTap(
interface=args.tap_interface, pcap=args.tap_pcap, pipe=args.tap_pipe, interface=args.tap_interface, pcap=args.tap_pcap, pipe=args.tap_pipe,
verbose=args.tap_verbose, output_file=args.tap_output, timeout=args.tap_timeout, verbose=args.tap_verbose, output_file=args.tap_output, timeout=args.tap_timeout,
filter_expr=args.tap_filter, filter_expr=args.tap_filter, max_samples=args.tap_max_samples,
redact_non_flock_ips=args.tap_redact_ips,
) )
if not args.tap_interface and not args.tap_pcap and not args.tap_pipe: if not args.tap_interface and not args.tap_pcap and not args.tap_pipe:
parser.print_help() parser.print_help()
+6
View File
@@ -0,0 +1,6 @@
import os
import sys
# Make the repo root importable so `import flock_tap` works regardless of the
# directory pytest is invoked from.
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
+332
View File
@@ -0,0 +1,332 @@
"""
Unit tests for flock_tap.py — the passive Flock camera traffic monitor.
These lock in the detection/accounting behavior, and in particular guard the
three correctness bugs that were fixed:
* FRP auth-payload detection must fire on data segments, not only on SYN.
* on_tcp_connect must not inflate bytes_up (real byte counts come from ip.len).
* connections to a known Flock cloud IP must be flagged even with no SNI/DNS.
Pure-logic tests run everywhere; packet-level tests are skipped when scapy is
not installed.
"""
import pytest
import flock_tap
from flock_tap import FlockTrafficTap
HAVE_SCAPY = flock_tap.HAVE_SCAPY
requires_scapy = pytest.mark.skipif(not HAVE_SCAPY, reason="scapy not installed")
if HAVE_SCAPY:
from scapy.all import IP, TCP, Raw
def _rebuild(pkt):
"""Serialize and re-parse so IP.len and offsets are populated like a real capture."""
return IP(bytes(pkt))
@pytest.fixture
def tap():
return FlockTrafficTap()
# ── _categorize_sni ─────────────────────────────────────────────────────────
@pytest.mark.parametrize("sni,expected", [
("login.flocksafety.com", "auth"),
("prod-flock-cd-xxx.edge.tenants.auth0.com", "auth"),
("flock-hibiki-inbox.s3.us-east-1.amazonaws.com", "s3_upload"),
("api.flocksafety.com", "cloud_api"),
("websockets.flocksafety.com", "cloud_api"),
("example.com", None),
("google.com", None),
])
def test_categorize_sni(tap, sni, expected):
assert tap._categorize_sni(sni) == expected
# ── on_dns_query ────────────────────────────────────────────────────────────
def test_on_dns_query_flocksafety_counts_cloud_dns(tap):
tap.on_dns_query("api.flocksafety.com", "10.0.0.5")
assert tap.flow_stats["10.0.0.5"]["cloud_dns"] == 1
def test_on_dns_query_auth0_counts_cloud_dns(tap):
# Regression: DNS detection used to match only "flocksafety", ignoring the
# Flock auth0 tenant that the SNI path already recognized.
tap.on_dns_query("prod-flock-cd-xxx.edge.tenants.auth0.com", "10.0.0.5")
assert tap.flow_stats["10.0.0.5"]["cloud_dns"] == 1
def test_on_dns_query_learns_resolved_ips(tap):
# Regression: resolved IPs for Flock domains are unioned into known_flock_ips
# so later direct-to-IP connections (no SNI/DNS) can still be correlated.
tap.on_dns_query("api.flocksafety.com", "10.0.0.5", resolved_ips=["203.0.113.9"])
assert "203.0.113.9" in tap.known_flock_ips
def test_on_dns_query_ignores_non_flock(tap):
tap.on_dns_query("example.com", "10.0.0.5", resolved_ips=["203.0.113.9"])
assert tap.flow_stats["10.0.0.5"]["cloud_dns"] == 0
assert "203.0.113.9" not in tap.known_flock_ips
# ── on_tcp_connect ──────────────────────────────────────────────────────────
def test_on_tcp_connect_does_not_inflate_bytes(tap):
# Regression: on_tcp_connect used to add a fixed 64 bytes per packet on top
# of real ip.len accounting, corrupting bandwidth stats.
tap.on_tcp_connect("10.0.0.5", "8.8.8.8", 51000, 443)
assert tap.flow_stats["10.0.0.5"]["bytes_up"] == 0
assert tap.flow_stats["10.0.0.5"]["connections"] == 1
def test_on_tcp_connect_frp_port_sets_tunnel(tap):
tap.on_tcp_connect("10.0.0.5", "10.0.0.9", 51000, 7000)
assert tap.flow_stats["10.0.0.5"]["frp_tunnel"] is True
def test_on_tcp_connect_non_frp_port_no_tunnel(tap):
tap.on_tcp_connect("10.0.0.5", "10.0.0.9", 51000, 8080)
assert tap.flow_stats["10.0.0.5"]["frp_tunnel"] is False
# ── on_tls_sni ──────────────────────────────────────────────────────────────
def test_on_tls_sni_categories(tap):
tap.on_tls_sni("login.flocksafety.com", "10.0.0.5", "1.1.1.1")
tap.on_tls_sni("flock-hibiki-inbox.s3.us-east-1.amazonaws.com", "10.0.0.5", "1.1.1.1")
tap.on_tls_sni("api.flocksafety.com", "10.0.0.5", "1.1.1.1")
fs = tap.flow_stats["10.0.0.5"]
assert fs["auth_tls"] == 1
assert fs["s3_uploads"] == 1
assert fs["cloud_api"] == 1
# ── classify_device ─────────────────────────────────────────────────────────
def test_classify_device_cloud_dns(tap):
tap.flow_stats["ip"]["cloud_dns"] = 1
assert tap.classify_device("ip") == "CLOUD_CONNECTED"
def test_classify_device_frp(tap):
tap.flow_stats["ip"]["frp_tunnel"] = True
assert tap.classify_device("ip") == "CLOUD_CONNECTED"
def test_classify_device_cloud_api_only(tap):
# The cloud_api-only branch is what surfaces known-Flock-IP correlation.
tap.flow_stats["ip"]["cloud_api"] = 1
assert tap.classify_device("ip") == "CLOUD_CONNECTED"
def test_classify_device_local_station(tap):
tap.flow_stats["ip"]["connections"] = 3
assert tap.classify_device("ip") == "LOCAL_STATION"
def test_classify_device_offline(tap):
# Unified vocabulary: no traffic → INDETERMINATE (was OFFLINE_OR_UNMONITORED).
assert tap.classify_device("ip") == "INDETERMINATE"
# ── classify_camera (fallback via device data) ──────────────────────────────
def test_classify_camera_no_data(tap):
# Unified vocabulary: no data → INDETERMINATE (was NO_DATA).
assert tap.classify_camera("10.0.0.99") == "INDETERMINATE"
def test_classify_camera_flock_sni_fallback(tap):
tap.devices["ip"]["tls_snis"] = [{"sni": "login.flocksafety.com"}]
assert tap.classify_camera("ip") == "CLOUD_CONNECTED"
def test_classify_camera_local_fallback(tap):
tap.devices["ip"]["connections"] = [{"dst": "10.0.0.9"}]
assert tap.classify_camera("ip") == "LOCAL_STATION"
# ── _parse_tcpdump_line ─────────────────────────────────────────────────────
def test_parse_tcpdump_dns(tap):
line = "13:37:00.000000 IP 10.0.0.5.54321 > 8.8.8.8.53: 12345+ A? api.flocksafety.com. (36)"
parsed = tap._parse_tcpdump_line(line)
assert parsed["type"] == "dns"
assert parsed["src"] == "10.0.0.5"
assert parsed["query"] == "api.flocksafety.com"
def test_parse_tcpdump_syn(tap):
line = "13:37:00.000000 IP 10.0.0.5.44444 > 10.0.0.9.7000: Flags [S], seq 1, win 64240, length 0"
parsed = tap._parse_tcpdump_line(line)
assert parsed["type"] == "syn"
assert parsed["dport"] == 7000
def test_parse_tcpdump_other(tap):
line = "13:37:00.000000 IP 10.0.0.5.44444 > 10.0.0.9.8080: Flags [P.], seq 1:10, length 9"
parsed = tap._parse_tcpdump_line(line)
assert parsed["type"] == "other"
def test_parse_tcpdump_garbage(tap):
assert tap._parse_tcpdump_line("not a tcpdump line") is None
# ── packet-level regression guards (scapy) ──────────────────────────────────
@requires_scapy
def test_byte_accounting_matches_ip_len(tap):
# bytes_up on the source device must equal the sum of real ip.len values,
# with no per-packet inflation.
pkts = [
_rebuild(IP(src="10.0.0.5", dst="8.8.8.8") / TCP(dport=443, flags="S")),
_rebuild(IP(src="10.0.0.5", dst="8.8.8.8") / TCP(dport=443, flags="PA") / Raw(b"x" * 100)),
]
expected = sum(int(p[IP].len) for p in pkts)
for p in pkts:
tap._handle_packet_scapy(p)
assert tap.devices["10.0.0.5"]["bytes_up"] == expected
@requires_scapy
def test_frp_payload_detected_on_non_syn(tap):
# Regression: the FRP auth-payload scan used to be nested in the SYN-only
# branch, so it never fired (SYN packets carry no payload). A data segment
# (PSH/ACK) carrying the auth JSON must now be detected.
pkt = _rebuild(
IP(src="10.0.0.5", dst="10.0.0.9")
/ TCP(sport=51000, dport=12345, flags="PA")
/ Raw(b'{"proxy_type":"tcp","auth":"token"}')
)
tap._handle_packet_scapy(pkt)
assert tap.flow_stats["10.0.0.5"]["frp_tunnel"] is True
types = [t["type"] for t in tap.devices["10.0.0.5"]["frp_tunnels"]]
assert "frp_auth_payload" in types
@requires_scapy
def test_frp_port_detected_on_syn(tap):
pkt = _rebuild(IP(src="10.0.0.5", dst="10.0.0.9") / TCP(sport=51000, dport=7000, flags="S"))
tap._handle_packet_scapy(pkt)
assert tap.flow_stats["10.0.0.5"]["frp_tunnel"] is True
types = [t["type"] for t in tap.devices["10.0.0.5"]["frp_tunnels"]]
assert "frp_tunnel" in types
@requires_scapy
def test_known_flock_ip_correlation(tap):
# A camera talking straight to a known Flock cloud IP (no SNI, no DNS of its
# own) must still be counted as cloud_api.
flock_ip = flock_tap.FLOCK_CLOUD_IPS[0]
pkt = _rebuild(IP(src="10.0.0.5", dst=flock_ip) / TCP(sport=51000, dport=443, flags="S"))
tap._handle_packet_scapy(pkt)
assert tap.flow_stats["10.0.0.5"]["cloud_api"] >= 1
assert tap.classify_device("10.0.0.5") == "CLOUD_CONNECTED"
# ── audit-integrity: unified verdict vocabulary ─────────────────────────────
def test_device_and_camera_classifiers_agree(tap):
"""Both public names fold into one canonical verdict."""
tap.flow_stats["10.0.0.5"]["cloud_dns"] = 1
assert tap.classify_device("10.0.0.5") == tap.classify_camera("10.0.0.5") == "CLOUD_CONNECTED"
# ── audit-integrity: report buckets partition every tracked IP ──────────────
def test_summary_buckets_reconcile(tap):
tap.start_time = 0
tap.flow_stats["10.0.0.1"]["cloud_dns"] = 1 # cloud
tap.flow_stats["10.0.0.2"]["connections"] = 2 # local
_ = tap.flow_stats["10.0.0.3"] # indeterminate (touched only)
s = tap.generate_report()["summary"]
assert s["counts_reconcile"] is True
assert s["cloud_connected"] + s["local_station"] + s["indeterminate"] == s["total_ips_tracked"]
def test_indeterminate_device_still_counted(tap):
"""Regression: the old report dropped NO_DATA/OFFLINE from every bucket."""
tap.start_time = 0
_ = tap.flow_stats["10.0.0.9"]
assert tap.generate_report()["summary"]["indeterminate"] == 1
# ── audit-integrity: CLOUD_IP evidence is recorded ──────────────────────────
@requires_scapy
def test_cloud_ip_contact_recorded_for_evidence(tap):
dst = flock_tap.FLOCK_CLOUD_IPS[0]
pkt = _rebuild(IP(src="192.168.1.120", dst=dst) / TCP(sport=47000, dport=443, flags="A"))
tap._handle_packet_scapy(pkt)
assert dst in tap.devices["192.168.1.120"]["cloud_ip_contacts"]
# ── audit-integrity: memory capping keeps counters exact ────────────────────
@requires_scapy
def test_samples_capped_but_counter_exact():
tap = FlockTrafficTap(max_samples=3)
for i in range(10):
pkt = _rebuild(IP(src="192.168.1.130", dst="10.0.0.20") /
TCP(sport=48000 + i, dport=8080, flags="S"))
tap._handle_packet_scapy(pkt)
dev = tap.devices["192.168.1.130"]
assert len(dev["connections"]) == 3 # sample list bounded
assert dev["counts"]["conn"] == 10 # counter exact
rep = tap.generate_report()
assert rep["devices"]["192.168.1.130"]["connections_count"] == 10
# ── audit-integrity: FRP auth stores a descriptor, not raw payload ──────────
@requires_scapy
def test_frp_auth_stores_descriptor_not_payload(tap):
pkt = _rebuild(IP(src="192.168.1.140", dst="10.0.0.30") /
TCP(sport=49000, dport=8080, flags="PA") /
Raw(load=b'{"proxy_type":"tcp","secret":"topsecret"}'))
tap._handle_packet_scapy(pkt)
rec = [f for f in tap.devices["192.168.1.140"]["frp_tunnels"]
if f["type"] == "frp_auth_payload"][0]
assert "payload_preview" not in rec
assert rec["matched_keyword"] == "proxy_type"
assert rec["payload_len"] == 41
assert "topsecret" not in str(rec) # secret bytes must not be retained
# ── audit-integrity: optional non-Flock IP redaction ────────────────────────
@requires_scapy
def test_redaction_hashes_non_flock_ip_keeps_flock_clear():
tap = FlockTrafficTap(redact_non_flock_ips=True)
flock = flock_tap.FLOCK_CLOUD_IPS[0]
p1 = _rebuild(IP(src="192.168.1.150", dst="8.8.8.8") / TCP(sport=50000, dport=8080, flags="S"))
tap._handle_packet_scapy(p1)
assert tap.devices["192.168.1.150"]["connections"][0]["dst"].startswith("redacted:")
p2 = _rebuild(IP(src="192.168.1.151", dst=flock) / TCP(sport=50001, dport=7000, flags="S"))
tap._handle_packet_scapy(p2)
assert tap.devices["192.168.1.151"]["frp_tunnels"][0]["dst"] == flock
# ── audit-integrity: destination summary is metadata-only ───────────────────
@requires_scapy
def test_destination_summary_is_metadata_only(tap):
tap.start_time = 0
for i in range(6): # clear the <5-packet filter
tap._handle_packet_scapy(
_rebuild(IP(src="192.168.1.160", dst="10.0.0.40") /
TCP(sport=51000 + i, dport=8080, flags="S")))
tap.on_tls_sni("api.flocksafety.com", "192.168.1.160", "10.0.0.40")
tap.devices["192.168.1.160"]["tls_snis"].append(
{"timestamp": 0, "src": "192.168.1.160", "dst": "10.0.0.40", "sni": "api.flocksafety.com"})
tap.devices["192.168.1.160"]["counts"]["sni"] += 1
dsum = tap.generate_report()["devices"]["192.168.1.160"]["destination_summary"]
assert set(dsum.keys()) == {"dns_domains", "tls_snis", "flock_cloud_ips"}
assert "api.flocksafety.com" in dsum["tls_snis"]