sentinel / gatekeeper.py
kswffs's picture
Upload folder using huggingface_hub
b96103d verified
Raw
History Blame Contribute Delete
21 kB
import cv2
import csv
import math
import time
import numpy as np
from dataclasses import dataclass
# Structured logger setup
try:
import structlog
logger = structlog.get_logger()
except ImportError:
import logging
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger("sentinel-gatekeeper")
# Wrap heavy ML framework imports to ensure robustness in different environments
try:
import tensorflow as tf
import tensorflow_hub as hub
TF_AVAILABLE = True
except ImportError:
TF_AVAILABLE = False
logger.warn("TensorFlow or TensorFlow Hub is not installed. AudioMonitor will run in stub mode.")
try:
import mediapipe as mp
MP_AVAILABLE = True
except ImportError:
MP_AVAILABLE = False
logger.warn("MediaPipe is not installed. PoseAnalyzer will run in stub mode.")
# --- DATACLASSES ---
@dataclass
class FrameChangeResult:
"""
Data structure containing motion telemetry metrics.
"""
change_pct: float
is_significant: bool # True if change_pct > threshold_used
threshold_used: float
motion_regions: int # Number of distinct moving blobs of significant size
@dataclass
class Detection:
"""
Object detection representation mapping to target VLM triggers.
"""
class_name: str
confidence: float
bbox: tuple[int, int, int, int] # (x1, y1, x2, y2)
area_pct: float # Bounding box area as percentage of total frame area
is_trigger: bool = True # Always True since we only yield trigger detections
@dataclass
class AudioClass:
"""
Audio classification output mapping to critical trigger classes.
"""
class_name: str
confidence: float
alert_level: str # "critical" | "warning" | "info"
@dataclass
class PoseData:
"""
Human pose estimation landmarks and fall indicators.
"""
head_y: float # Normalized Y coordinate (0.0 to 1.0)
hip_y: float
shoulder_y: float
is_fall: bool # True if fall conditions are met
is_sitting: bool # True if person appears sitting
confidence: float
@dataclass
class SensorSnapshot:
"""
Mobile companion telemetry data snapshot.
"""
accelerometer: tuple[float, float, float] | None = None
gyroscope: tuple[float, float, float] | None = None
gps: tuple[float, float] | None = None
light_level: float | None = None
battery_pct: float | None = None
heading: float | None = None
@dataclass
class TriggerDecision:
"""
Decision wrapper containing VLM activation outcome and telemetry.
"""
should_trigger: bool
reason: str
urgency: str # "critical" | "warning" | "info" | "none"
confidence: float # Combined confidence rating from active signals
signals: dict # Dictionary of raw input signals for logger
cooldown_remaining: float # Seconds until next non-critical VLM call is permitted
# --- MODULE CLASSES ---
class FrameChangeDetector:
"""
CPU-bound motion detection using traditional computer vision.
Resizes frames to 320x240 and computes absolute differences to filter idle scenes.
"""
def __init__(self, threshold_pct: float = 5.0, adaptive: bool = True):
self.previous_frame: np.ndarray | None = None
self.threshold_pct = threshold_pct
self.adaptive = adaptive
self.change_history: list[float] = []
def reset(self) -> None:
self.previous_frame = None
self.change_history = []
def detect(self, frame: np.ndarray) -> FrameChangeResult:
if frame is None:
return FrameChangeResult(0.0, False, self.threshold_pct, 0)
gray = cv2.cvtColor(frame, cv2.COLOR_BGR2GRAY)
gray_resized = cv2.resize(gray, (320, 240))
if self.previous_frame is None:
self.previous_frame = gray_resized
return FrameChangeResult(0.0, False, self.threshold_pct, 0)
diff = cv2.absdiff(self.previous_frame, gray_resized)
_, thresh = cv2.threshold(diff, 25, 255, cv2.THRESH_BINARY)
non_zero = cv2.countNonZero(thresh)
total_pixels = 320 * 240
change_pct = (non_zero / total_pixels) * 100.0
threshold_used = self.threshold_pct
if self.adaptive:
self.change_history.append(change_pct)
if len(self.change_history) > 60:
self.change_history.pop(0)
if len(self.change_history) >= 10:
mean_change = np.mean(self.change_history)
std_change = np.std(self.change_history)
adaptive_thresh = mean_change + (2.0 * std_change)
threshold_used = max(adaptive_thresh, self.threshold_pct)
is_significant = change_pct > threshold_used
contours, _ = cv2.findContours(thresh, cv2.RETR_EXTERNAL, cv2.CHAIN_APPROX_SIMPLE)
motion_regions = 0
for cnt in contours:
if cv2.contourArea(cnt) > 500.0:
motion_regions += 1
self.previous_frame = gray_resized
return FrameChangeResult(
change_pct=round(change_pct, 2),
is_significant=is_significant,
threshold_used=round(threshold_used, 2),
motion_regions=motion_regions
)
class ObjectDetector:
"""
CPU-optimized object detection using YOLO11n.
Downscales frames to 416x416 and outputs critical safety trigger classes.
"""
def __init__(self, model_size: str = "yolo11n"):
from ultralytics import YOLO
self.model = YOLO(f"{model_size}.pt")
self.trigger_classes = {
"person": 0.7,
"fire": 0.5,
"knife": 0.8,
"dog": 0.6,
"car": 0.7,
"bicycle": 0.7,
"suitcase": 0.5,
"chair": 0.6,
}
def detect(self, frame: np.ndarray) -> list[Detection]:
if frame is None:
return []
h, w = frame.shape[:2]
resized = cv2.resize(frame, (416, 416))
results = self.model.predict(resized, imgsz=416, conf=0.3, verbose=False)
detections = []
if len(results) == 0:
return []
result = results[0]
boxes = result.boxes
for box in boxes:
cls_id = int(box.cls[0])
class_name = self.model.names[cls_id]
conf = float(box.conf[0])
if class_name in self.trigger_classes and conf >= self.trigger_classes[class_name]:
xyxy = box.xyxy[0].tolist()
x1, y1, x2, y2 = int(xyxy[0]), int(xyxy[1]), int(xyxy[2]), int(xyxy[3])
bbox_area = (x2 - x1) * (y2 - y1)
area_pct = (bbox_area / (416.0 * 416.0)) * 100.0
if area_pct >= 5.0:
detections.append(Detection(
class_name=class_name,
confidence=conf,
bbox=(x1, y1, x2, y2),
area_pct=round(area_pct, 2),
is_trigger=True
))
return detections
def get_trigger_summary(self, detections: list[Detection]) -> str:
if not detections:
return "No critical objects detected."
summaries = []
for det in detections:
cx = (det.bbox[0] + det.bbox[2]) / 2.0
if cx < 0.33 * 416.0:
pos = "left side"
elif cx > 0.66 * 416.0:
pos = "right side"
else:
pos = "center"
summaries.append(
f"{det.class_name.capitalize()} detected ({int(det.confidence * 100)}% conf) at the {pos} of the frame."
)
return " ".join(summaries)
class AudioMonitor:
"""
Audio classification gatekeeper loading YAMNet.
Monitors audio streams on CPU and triggers on critical safety noises.
"""
def __init__(self):
self.model = None
self.class_names = []
self.trigger_sounds = {
"Scream": {"threshold": 0.6, "level": "critical"},
"Glass": {"threshold": 0.5, "level": "warning"},
"Alarm": {"threshold": 0.5, "level": "warning"},
"Crash": {"threshold": 0.6, "level": "critical"},
"Explosion": {"threshold": 0.5, "level": "critical"},
"Baby cry": {"threshold": 0.7, "level": "info"},
"Dog bark": {"threshold": 0.7, "level": "info"},
"Siren": {"threshold": 0.6, "level": "warning"},
}
if TF_AVAILABLE:
try:
logger.info("Initializing YAMNet from TensorFlow Hub...")
self.model = hub.load("https://tfhub.dev/google/yamnet/1")
class_map_path = self.model.class_map_path().numpy().decode("utf-8")
with open(class_map_path) as f:
reader = csv.reader(f)
next(reader)
for row in reader:
self.class_names.append(row[2])
except Exception as e:
logger.error("YAMNet load failed. AudioMonitor will default to stub mode.", error=str(e))
self.model = None
def classify(self, audio_samples: np.ndarray, sample_rate: int = 16000) -> list[AudioClass]:
if self.model is None or not TF_AVAILABLE:
return []
try:
if len(audio_samples) < 15600:
audio_samples = np.pad(audio_samples, (0, 15600 - len(audio_samples)), mode='constant')
scores, _, _ = self.model(audio_samples)
scores_np = scores.numpy()
max_scores = np.max(scores_np, axis=0)
results = []
for i, class_name in enumerate(self.class_names):
conf = float(max_scores[i])
for trigger_name, trigger_info in self.trigger_sounds.items():
if trigger_name.lower() in class_name.lower() and conf >= trigger_info["threshold"]:
results.append(AudioClass(
class_name=trigger_name,
confidence=conf,
alert_level=trigger_info["level"]
))
return results
except Exception as e:
logger.error("Error during audio classification in YAMNet", error=str(e))
return []
class PoseAnalyzer:
"""
Postural analysis and fall detection module leveraging MediaPipe Pose.
Processes single frames on CPU to identify falls and sit-down stances.
"""
def __init__(self):
self.pose = None
self.mp_pose = None
if MP_AVAILABLE:
try:
self.mp_pose = mp.solutions.pose
self.pose = self.mp_pose.Pose(
static_image_mode=True,
model_complexity=0,
min_detection_confidence=0.5
)
except Exception as e:
logger.error("MediaPipe initialization failed. PoseAnalyzer will run in stub mode.", error=str(e))
self.pose = None
def analyze(self, frame: np.ndarray) -> PoseData | None:
if self.pose is None or not MP_AVAILABLE or frame is None:
return None
try:
rgb_frame = cv2.cvtColor(frame, cv2.COLOR_BGR2RGB)
results = self.pose.process(rgb_frame)
if not results.pose_landmarks:
return None
landmarks = results.pose_landmarks.landmark
nose = landmarks[self.mp_pose.PoseLandmark.NOSE]
left_hip = landmarks[self.mp_pose.PoseLandmark.LEFT_HIP]
right_hip = landmarks[self.mp_pose.PoseLandmark.RIGHT_HIP]
left_shoulder = landmarks[self.mp_pose.PoseLandmark.LEFT_SHOULDER]
right_shoulder = landmarks[self.mp_pose.PoseLandmark.RIGHT_SHOULDER]
head_y = nose.y
hip_y = (left_hip.y + right_hip.y) / 2.0
shoulder_y = (left_shoulder.y + right_shoulder.y) / 2.0
is_fall = head_y > (hip_y + 0.1) and head_y > 0.6
is_sitting = hip_y > 0.7
confidence = (nose.visibility + left_hip.visibility + right_hip.visibility) / 3.0
return PoseData(
head_y=round(head_y, 4),
hip_y=round(hip_y, 4),
shoulder_y=round(shoulder_y, 4),
is_fall=is_fall,
is_sitting=is_sitting,
confidence=round(confidence, 4)
)
except Exception as e:
logger.error("Error analyzing human pose in MediaPipe", error=str(e))
return None
def close(self):
if self.pose is not None:
self.pose.close()
class GatekeeperDecision:
"""
Advanced Signal Aggregation Decision Engine for Sentinel.
Combines CPU-bound frame delta checks, object detections, auditory events,
and posture models to determine whether to trigger Tier 2 VLM analysis.
Implements a strict, configurable cooldown protocol.
"""
def __init__(self, cooldown_seconds: float = 30.0):
"""
Initializes the GatekeeperDecision engine.
"""
self.cooldown_seconds = cooldown_seconds
self.last_trigger_time: float = 0.0
self.trigger_history: list[TriggerDecision] = []
self.last_gps: tuple[float, float] | None = None
def decide(
self,
frame_change: FrameChangeResult,
detections: list[Detection],
audio_classes: list[AudioClass],
pose_data: PoseData | None,
user_query: str | None = None,
sensor_data: SensorSnapshot | None = None
) -> TriggerDecision:
"""
Evaluates current sensor feeds. Returns a TriggerDecision.
"""
current_time = time.time()
should_trigger = False
reason = "Monitoring active: environment stable."
urgency = "none"
confidences = []
# 1. MANUAL QUERY (Highest priority)
if user_query is not None and user_query.strip() != "":
should_trigger = True
reason = f"Manual user query received: '{user_query.strip()}'"
urgency = "info"
confidences.append(1.0)
# 2. CRITICAL AUDIO
if not should_trigger:
critical_sounds = [a for a in audio_classes if a.alert_level == "critical"]
if critical_sounds:
should_trigger = True
max_audio = max(critical_sounds, key=lambda x: x.confidence)
reason = f"Critical sound detected: {max_audio.class_name} ({int(max_audio.confidence * 100)}% conf)"
urgency = "critical"
confidences.extend([a.confidence for a in critical_sounds])
# 3. FALL DETECTION
if not should_trigger:
if pose_data and pose_data.is_fall and pose_data.confidence > 0.6:
should_trigger = True
reason = f"Fall detected (head_y: {pose_data.head_y:.2f} > hip_y: {pose_data.hip_y:.2f} + 0.1)"
urgency = "critical"
confidences.append(pose_data.confidence)
# 4. DANGEROUS COMBINATION (Significant Motion + Person Detected)
if not should_trigger:
person_detections = [d for d in detections if d.class_name == "person"]
if frame_change.is_significant and person_detections:
should_trigger = True
max_person = max(person_detections, key=lambda x: x.confidence)
reason = f"Motion + Person presence: frame delta {frame_change.change_pct:.1f}% and person conf {max_person.confidence:.2f}"
urgency = "warning"
confidences.append(frame_change.change_pct / 100.0)
confidences.extend([d.confidence for d in person_detections])
# 5. SENSOR ANOMALY (Drop or Impact)
if not should_trigger and sensor_data is not None:
if sensor_data.accelerometer is not None and sensor_data.gyroscope is not None:
ax, ay, az = sensor_data.accelerometer
gx, gy, gz = sensor_data.gyroscope
accel_mag = math.sqrt(ax**2 + ay**2 + az**2)
gyro_mag = math.sqrt(gx**2 + gy**2 + gz**2)
# Check for accelerometer spike (>3g = ~29.4 m/s^2 or 3.0 g) and high gyro spin (>5 rad/s)
is_accel_spike = accel_mag > 29.4 or (3.0 < accel_mag < 9.0)
is_gyro_spin = gyro_mag > 5.0
if is_accel_spike and is_gyro_spin:
should_trigger = True
reason = f"Sensor anomaly: Acceleration magnitude {accel_mag:.2f} and gyroscope spin {gyro_mag:.2f} rad/s"
urgency = "critical"
confidences.append(0.9)
# 6. ENVIRONMENTAL CHANGE (Walking into sudden darkness)
if not should_trigger and sensor_data is not None:
if sensor_data.light_level is not None and sensor_data.light_level < 5.0 and sensor_data.gps is not None:
gps_moved = False
lat, lng = sensor_data.gps
if self.last_gps is not None:
plat, plng = self.last_gps
# Estimate horizontal distance change (0.00002 degrees ~ 2.2m)
dist = math.sqrt((lat - plat)**2 + (lng - plng)**2)
if dist > 0.00002:
gps_moved = True
self.last_gps = (lat, lng)
if gps_moved:
should_trigger = True
reason = f"Environmental change: Sudden darkness ({sensor_data.light_level:.1f} lux) and user is walking/moving"
urgency = "info"
confidences.append(0.8)
elif sensor_data.gps is not None:
self.last_gps = sensor_data.gps
# Calculate combined confidence
combined_confidence = 0.0
if confidences:
num_signals = len(confidences)
combined_confidence = max(confidences) + 0.1 * (num_signals - 1)
combined_confidence = min(combined_confidence, 1.0)
# Calculate cooldown remaining
time_since_trigger = current_time - self.last_trigger_time
cooldown_remaining = max(0.0, self.cooldown_seconds - time_since_trigger)
# Cooldown enforcement logic
if should_trigger:
# Critical alerts bypass cooldown; warnings and info checks respect it
if urgency != "critical" and cooldown_remaining > 0.0:
should_trigger = False
reason = f"Trigger suppressed due to cooldown ({cooldown_remaining:.1f}s remaining). Event: {reason}"
else:
self.last_trigger_time = current_time
# Raw signals dictionary for logging and RAG context injection
signals_dict = {
"frame_change_pct": frame_change.change_pct,
"motion_significant": frame_change.is_significant,
"motion_regions": frame_change.motion_regions,
"detections": [{"class": d.class_name, "conf": d.confidence} for d in detections],
"audio_events": [{"class": a.class_name, "conf": a.confidence, "level": a.alert_level} for a in audio_classes],
"pose": {
"head_y": pose_data.head_y,
"hip_y": pose_data.hip_y,
"is_fall": pose_data.is_fall,
"confidence": pose_data.confidence
} if pose_data else None,
"sensor": {
"accel": sensor_data.accelerometer,
"gyro": sensor_data.gyroscope,
"gps": sensor_data.gps,
"light": sensor_data.light_level
} if sensor_data else None,
"user_query": user_query
}
# Structlog execution logging
logger.info(
"Gatekeeper decision evaluated",
should_trigger=should_trigger,
urgency=urgency,
confidence=round(combined_confidence, 2),
cooldown_remaining=round(cooldown_remaining, 1),
reason=reason
)
decision = TriggerDecision(
should_trigger=should_trigger,
reason=reason,
urgency=urgency,
confidence=round(combined_confidence, 2),
signals=signals_dict,
cooldown_remaining=round(cooldown_remaining, 1)
)
self.trigger_history.append(decision)
if len(self.trigger_history) > 500:
self.trigger_history = self.trigger_history[-250:]
return decision