import os import time import base64 import json import requests try: from crewai.tools import tool except Exception: # Fallback no-op decorator when crewai is not installed (allows offline testing) def tool(name_or_fn=None): def _decorator(fn=None): return fn if callable(name_or_fn): return name_or_fn return _decorator try: from tavily import TavilyClient except Exception: TavilyClient = None try: from groq import Groq except Exception: Groq = None import io from PIL import Image import hashlib import shutil import subprocess from datetime import datetime # ========================================== # 1. FORENSICS TOOLS (Used by LangGraph Router) # ========================================== def run_hf_deepfake(image_path: str) -> dict: """Sends an image to Hugging Face to check for deepfake manipulation.""" hf_token = os.getenv("HUGGINGFACE_API_KEY") if not hf_token: return {"fake_prob": 0.0, "evidence": "Hugging Face API key missing."} headers = {"Authorization": f"Bearer {hf_token}"} models = [ "prithivMLmods/deepfake-detector-model-v1", "dima806/deepfake_vs_real_image_detection" ] try: with open(image_path, "rb") as f: image_data = f.read() except Exception as e: return {"fake_prob": 0.0, "evidence": f"Failed to read image: {e}"} for model_id in models: url = f"https://router.huggingface.co/hf-inference/models/{model_id}" print(f" ☁️ Querying HF Deepfake Model: {model_id}...") try: response = requests.post(url, headers=headers, data=image_data) # Handle cold-start loading if response.status_code == 503: print(" 💤 Model is loading. Waiting 3 seconds...") time.sleep(3) response = requests.post(url, headers=headers, data=image_data) if response.status_code == 200: scores = response.json() if isinstance(scores, list) and isinstance(scores[0], list): scores = scores[0] for item in scores: label = item.get('label', '').lower() if label in ['fake', 'artificial', 'deepfake', 'ai', 'generated']: prob = item.get('score', 0.0) return {"fake_prob": prob, "evidence": f"Flagged by {model_id} ({prob:.1%} confidence)"} elif label in ['real', 'human', 'realism'] and item.get('score', 0.0) > 0.9: return {"fake_prob": 0.0, "evidence": f"Classified as Real ({item.get('score'):.1%} confidence)"} return {"fake_prob": 0.0, "evidence": "Inconclusive results."} except Exception as e: print(f" ❌ HF Connection Error on {model_id}: {e}") return {"fake_prob": 0.0, "evidence": "All Deepfake models failed or timed out."} def run_groq_vision(image_path: str) -> str: """Uses Groq Vision (Llama-4-Scout) to generate a highly detailed caption of the media.""" groq_key = os.getenv("GROQ_API_KEY") if not groq_key: return "Groq API key missing. Cannot generate context." print(" 👁️ Asking Groq Vision to analyze scene context...") client = Groq(api_key=groq_key) try: with open(image_path, "rb") as image_file: b64_image = base64.b64encode(image_file.read()).decode('utf-8') prompt = "Describe this image in detail. What is happening? Who is in it? Return ONLY a JSON object: {\"caption\": \"detailed description\", \"subject\": \"name or 'Unknown'\"}" completion = client.chat.completions.create( model="meta-llama/llama-4-scout-17b-16e-instruct", # Updated to the active 2026 model messages=[ { "role": "user", "content": [ {"type": "text", "text": prompt}, {"type": "image_url", "image_url": {"url": f"data:image/jpeg;base64,{b64_image}"}}, ], } ], temperature=0.1, response_format={"type": "json_object"} ) data = json.loads(completion.choices[0].message.content) return f"{data.get('caption', 'Image content')}. Subject: {data.get('subject', 'Unknown')}." except Exception as e: print(f" ❌ Groq Vision Error: {e}") return "Failed to extract vision context." # ========================================== # 0. HUGGING FACE INFERENCE API HELPERS # ========================================== def hf_inference_request(model_id: str, data: bytes, content_type: str = None, timeout: int = 60): """Generic helper to call the Hugging Face Inference API router. Returns parsed JSON on success, or raises an Exception. """ hf_token = os.getenv("HUGGINGFACE_API_KEY") if not hf_token: raise RuntimeError("HUGGINGFACE_API_KEY not set") headers = {"Authorization": f"Bearer {hf_token}"} if content_type: headers["Content-Type"] = content_type urls = [ f"https://router.huggingface.co/hf-inference/models/{model_id}", f"https://api-inference.huggingface.co/models/{model_id}", ] last_error = None for url in urls: for attempt in range(2): resp = requests.post(url, headers=headers, data=data, timeout=timeout) if resp.status_code == 503 and attempt == 0: time.sleep(3) continue if resp.status_code == 200: try: return resp.json() except Exception: return resp.content msg = f"HF inference error {resp.status_code}: {resp.text}" last_error = msg unsupported = "not supported by provider" in (resp.text or "").lower() if unsupported and "router.huggingface.co" in url: break if attempt == 0 and resp.status_code >= 500: continue break raise RuntimeError(last_error or "HF inference failed") def hf_image_to_text(model_id: str, pil_image: Image.Image, max_new_tokens: int = 40) -> str: """Call HF inference API for image-to-text models and return generated caption/text. This avoids local model downloads by using the hosted inference endpoint. """ buf = io.BytesIO() pil_image.save(buf, format="JPEG") img_bytes = buf.getvalue() resp = hf_inference_request(model_id, img_bytes, content_type="application/octet-stream") # Common HF responses for vision captioners include 'generated_text' or list of dicts if isinstance(resp, dict) and "generated_text" in resp: return resp["generated_text"] if isinstance(resp, list) and len(resp) > 0: first = resp[0] if isinstance(first, dict) and "generated_text" in first: return first["generated_text"] # Some models return simple text if isinstance(first, str): return first # Fallback: if response is bytes or unknown, decode to string if isinstance(resp, (bytes, bytearray)): try: return resp.decode("utf-8") except Exception: return "" def hf_audio_asr(model_id: str, audio_path: str) -> str: """Send audio bytes to HF inference API for ASR and return the transcript text.""" with open(audio_path, "rb") as f: audio_bytes = f.read() resp = hf_inference_request(model_id, audio_bytes, content_type="application/octet-stream") if isinstance(resp, dict): # Many ASR endpoints return {'text': '...'} if "text" in resp: return resp.get("text", "") # Or return a list of chunks if "chunks" in resp: return " ".join([c.get("text", "") for c in resp.get("chunks", [])]) if isinstance(resp, str): return resp return "" def groq_audio_asr(audio_path: str, model_id: str = "whisper-large-v3-turbo") -> str: """Transcribe audio using Groq Speech-to-Text and return transcript text.""" groq_key = os.getenv("GROQ_API_KEY") if not groq_key: raise RuntimeError("GROQ_API_KEY not set") if Groq is None: raise RuntimeError("groq package not available") client = Groq(api_key=groq_key) with open(audio_path, "rb") as audio_file: transcript = client.audio.transcriptions.create( file=audio_file, model=model_id, response_format="json", ) text = getattr(transcript, "text", "") if isinstance(text, str): return text if isinstance(transcript, dict): return transcript.get("text", "") return "" def hf_audio_classify(model_id: str, audio_path: str): """Send audio to HF audio-classification endpoints and return parsed JSON result.""" with open(audio_path, "rb") as f: audio_bytes = f.read() resp = hf_inference_request(model_id, audio_bytes, content_type="application/octet-stream") return resp # ========================================== # 2. OSINT TOOLS (Used by CrewAI Agents) # ========================================== @tool("Tavily Web Search Tool") def tavily_search_tool(query: str) -> str: """ Searches the web for facts, news articles, and debunks regarding a specific claim. Always pass a detailed search query to this tool. """ tavily_key = os.getenv("TAVILY_API_KEY") if not tavily_key: return "Error: Tavily API key is missing." print(f" 🔎 CrewAI Agent Searching Web: '{query}'") try: client = TavilyClient(api_key=tavily_key) # We append 'Fact check:' to nudge the search engine toward verification articles response = client.search(query=f"Fact check: {query}", search_depth="basic", max_results=4) results_text = "" for item in response.get('results', []): results_text += f"Source ({item['url']}): {item['content']}\n\n" return results_text if results_text else "No relevant web search results found." except Exception as e: return f"Web search failed: {str(e)}" def serp_search_tool(query: str, num: int = 5) -> str: """Legacy helper kept for compatibility; SerpAPI is no longer used.""" return ( "SerpAPI has been disabled in this build. Use Tavily search instead: " f"{query}" ) def serp_reverse_image(image_path: str, num: int = 5) -> str: """Legacy helper kept for compatibility; SerpAPI is no longer used.""" return "SerpAPI reverse-image search has been disabled in this build." # ========================================== # Provenance & Forensics Helpers # ========================================== def compute_checksum(file_path: str, algo: str = "sha256") -> str: """Compute a hex checksum for a file using the selected algorithm.""" h = hashlib.new(algo) with open(file_path, "rb") as f: for chunk in iter(lambda: f.read(8192), b""): h.update(chunk) return h.hexdigest() def run_ffprobe(file_path: str) -> dict: """Run `ffprobe` to extract container and stream metadata. Returns parsed JSON or error info.""" cmd = [ "ffprobe", "-v", "error", "-show_format", "-show_streams", "-print_format", "json", file_path, ] try: proc = subprocess.run(cmd, capture_output=True, text=True, timeout=20) if proc.returncode != 0: return {"error": "ffprobe failed", "stderr": proc.stderr} try: return json.loads(proc.stdout) except Exception: return {"error": "failed to parse ffprobe output", "stdout": proc.stdout, "stderr": proc.stderr} except FileNotFoundError: return {"error": "ffprobe not installed"} except Exception as e: return {"error": str(e)} def preserve_original_copy(src_path: str, dest_dir: str) -> dict: """Copy the original media into `dest_dir` and record basic provenance metadata. Returns a dict with keys: copied_path, checksum, ffprobe, metadata_path """ os.makedirs(dest_dir, exist_ok=True) checksum = compute_checksum(src_path) base = os.path.basename(src_path) ts = datetime.utcnow().strftime("%Y%m%dT%H%M%SZ") name = f"{ts}_{checksum[:8]}_{base}" dest_path = os.path.join(dest_dir, name) shutil.copy2(src_path, dest_path) ff = run_ffprobe(dest_path) meta = { "original_path": os.path.abspath(src_path), "copied_path": os.path.abspath(dest_path), "checksum": checksum, "collected_at": datetime.utcnow().isoformat() + "Z", "ffprobe": ff, } meta_path = dest_path + ".metadata.json" try: with open(meta_path, "w") as mf: json.dump(meta, mf, indent=2) except Exception: pass return {"copied_path": dest_path, "checksum": checksum, "ffprobe": ff, "metadata_path": meta_path} def extract_audio_with_ffmpeg(src_path: str, out_audio_path: str, sample_rate: int = 16000) -> dict: """Extract audio to a WAV file using ffmpeg and return run diagnostics (stderr/stdout). This helper captures ffmpeg stderr which is useful for diagnosing silent/absent tracks. """ cmd = [ "ffmpeg", "-y", "-i", src_path, "-vn", "-acodec", "pcm_s16le", "-ar", str(sample_rate), "-ac", "1", out_audio_path, ] try: proc = subprocess.run(cmd, capture_output=True, text=True, timeout=60) return {"returncode": proc.returncode, "stderr": proc.stderr, "stdout": proc.stdout} except FileNotFoundError: return {"error": "ffmpeg not installed"} except Exception as e: return {"error": str(e)}