""" Nuclear Intelligence v4.0 - Enhanced Operation Loop ═══════════════════════════════════════════════════════════════════ Autonomous research-to-tokenization with: - Multi-stage pipeline - Intelligent retry logic - Error recovery - Comprehensive reporting - Developer mode with deep analysis - Always-online optimized ═══════════════════════════════════════════════════════════════════ """ import time import os import json import threading import hashlib from datetime import datetime from typing import Dict, Any, List, Optional from loguru import logger from dataclasses import dataclass, field, asdict from core.nuclear_intelligence_v4 import NuclearIntelligenceCore, ResearchQuestion, ResearchAnswer, EvaluationScore from core.research_controller import ResearchController from blockchain.virtual_ledger import VirtualLedger @dataclass class OperationLoopConfig: """Configuration for the operation loop""" interval_minutes: int = 30 min_accuracy: float = 70.0 min_novelty: float = 50.0 min_usefulness: float = 50.0 min_overall: float = 55.0 min_completeness: float = 40.0 auto_start: bool = True questions_per_cycle: int = 1 developer_mode: bool = True web_search_enabled: bool = True save_reports: bool = True max_retries: int = 3 retry_delay: int = 10 # HF Sync sync_to_hf: bool = True sync_to_gh: bool = True evaluation_samples: int = 2 evaluation_agreement_threshold: float = 0.80 @dataclass class OperationCycleResult: """Result of a single operation cycle""" cycle_id: str timestamp: str question: Dict answer: Dict evaluation: Dict minted: bool tx_hash: Optional[str] = None developer_analysis: Optional[Dict] = None governance: Optional[Dict] = None execution_time_seconds: float = 0.0 retry_count: int = 0 error: Optional[str] = None stage_timings_seconds: Dict[str, float] = field(default_factory=dict) def to_dict(self) -> Dict: return asdict(self) class OperationLoop: """Advanced autonomous research loop v4.0""" def __init__( self, core: NuclearIntelligenceCore, ledger: VirtualLedger, config: Optional[OperationLoopConfig] = None, ): self.core = core self.ledger = ledger self.config = config or OperationLoopConfig() self.history: List[OperationCycleResult] = [] self.is_running = False self._thread: Optional[threading.Thread] = None self._total_cycles = 0 self._successful_cycles = 0 self.controller = ResearchController( evaluation_samples=self.config.evaluation_samples, agreement_threshold=self.config.evaluation_agreement_threshold, ) self._load_history() logger.info(f"⚙️ Operation Loop v4.0 initialized: interval={self.config.interval_minutes}min, threshold={self.config.min_accuracy}%") def _load_history(self): """Load cycle history from reports directory""" reports_dir = "reports" if os.path.exists(reports_dir): try: files = sorted( [f for f in os.listdir(reports_dir) if f.startswith("cycle_") and f.endswith(".json")], reverse=True )[:200] for filename in files: try: with open(os.path.join(reports_dir, filename), 'r', encoding='utf-8') as f: d = json.load(f) self.history.append(OperationCycleResult( cycle_id=d.get("cycle_id", filename), timestamp=d.get("timestamp", ""), question=d.get("question", {}), answer=d.get("answer", {}), evaluation=d.get("evaluation", {}), minted=d.get("minted", False), tx_hash=d.get("tx_hash"), developer_analysis=d.get("developer_analysis"), governance=d.get("governance"), execution_time_seconds=d.get("execution_time_seconds", 0), retry_count=d.get("retry_count", 0), error=d.get("error"), stage_timings_seconds=d.get("stage_timings_seconds", {}), )) except Exception as e: logger.warning(f"Failed to load {filename}: {e}") logger.info(f"📜 Loaded {len(self.history)} history records") except Exception as e: logger.warning(f"History loading failed: {e}") def _should_mint( self, evaluation: EvaluationScore, answer: Optional[ResearchAnswer] = None, controller_approved: bool = True, ) -> Dict[str, Any]: """Determine if a genuinely researched, independently evaluated answer may be minted.""" overall = evaluation.overall_score() provider = (getattr(answer, "provider", "") or "").lower() evaluator_unavailable = "evaluation unavailable" in (evaluation.justification or "").lower() research_is_fallback = provider in {"fallback", "demo", "template_fallback", "unknown"} checks = { "accuracy": evaluation.scientific_accuracy >= self.config.min_accuracy, "novelty": evaluation.novelty_score >= self.config.min_novelty, "usefulness": evaluation.usefulness_score >= self.config.min_usefulness, "completeness": evaluation.completeness >= self.config.min_completeness, "overall": overall >= self.config.min_overall, "consistency": evaluation.self_consistency_check, } passed = sum(checks.values()) total = len(checks) threshold_pct = (passed / total) * 100 should_mint = ( checks["overall"] and checks["consistency"] and not evaluator_unavailable and not research_is_fallback and controller_approved ) if evaluator_unavailable or research_is_fallback or not controller_approved: logger.warning( "🚫 Mint blocked: real evaluator/provider required " f"(provider={provider or 'missing'}, evaluator_unavailable={evaluator_unavailable}, " f"controller_approved={controller_approved})" ) logger.info( f"📊 Minting Check: {passed}/{total} ({threshold_pct:.0f}%) | " f"Acc={evaluation.scientific_accuracy:.1f}% Novel={evaluation.novelty_score:.1f}% " f"Use={evaluation.usefulness_score:.1f}% Overall={overall:.1f}% → " f"{'✅ MINT' if should_mint else '❌ REJECT'}" ) return {"should_mint": should_mint, "checks": checks, "passed": passed, "total": total, "overall": overall} def _sync_to_huggingface(self, report_path: str): """Sync report to HuggingFace dataset""" try: from huggingface_hub import HfApi, create_repo hf_token = os.getenv("HF_TOKEN") if not hf_token or hf_token == "hf_placeholder": return False api = HfApi(token=hf_token) dataset_repo = os.getenv("HF_DATASET_REPO", "Qalam/nuclear-intelligence-dataset") try: create_repo(repo_id=dataset_repo, repo_type="dataset", token=hf_token, exist_ok=True) except: pass filename = os.path.basename(report_path) api.upload_file(path_or_fileobj=report_path, path_in_repo=f"reports/{filename}", repo_id=dataset_repo, repo_type="dataset") logger.info(f"📤 Synced to HF: {filename}") return True except Exception as e: logger.warning(f"HF sync failed: {e}") return False def _sync_to_github(self, report_path: str): """Sync report to GitHub""" try: from github import Github token = os.getenv("GITHUB_TOKEN") if not token or token == "ghp_placeholder": return False g = Github(token) repo_name = os.getenv("GH_REPO", "QalamHipHop/nuclear-intelligence") repo = g.get_repo(repo_name) filename = os.path.basename(report_path) content = open(report_path, 'r').read() path = f"reports/{filename}" try: existing = repo.get_contents(path) repo.update_file(path, f"Auto: NES cycle report", content, existing.sha) except: repo.create_file(path, f"Auto: NES cycle report", content) logger.info(f"📤 Synced to GH: {path}") return True except Exception as e: logger.warning(f"GitHub sync failed: {e}") return False def run_cycle(self, developer_mode: bool = False, force_category: str = "") -> OperationCycleResult: """Execute a single research cycle with retry logic""" cycle_id = hashlib.sha256(datetime.now().isoformat().encode()).hexdigest()[:16] start_time = time.time() logger.info(f"══════════════════════════════════════") logger.info(f"🔄 CYCLE {cycle_id} STARTING") logger.info(f"══════════════════════════════════════") retry_count = 0 last_error = None stage_timings: Dict[str, float] = {} while retry_count <= self.config.max_retries: try: # Step 1: Select a transparent research agenda and generate a question. stage_start = time.perf_counter() agenda = self.controller.select_next_category(self.history, force_category) logger.info( f"🧭 Agenda selected {agenda.selected_category} " f"(priority={agenda.priority:.3f})" ) logger.info(f"📝 Step 1: Generating question...") question = self.core.generate_question(category_hint=agenda.selected_category) if not question: raise RuntimeError("Question generation failed") stage_timings["agenda_and_question"] = round(time.perf_counter() - stage_start, 4) # Step 2: Conduct Research logger.info(f"🔬 Step 2: Conducting research...") stage_start = time.perf_counter() answer = self.core.conduct_research( question, use_web_search=self.config.web_search_enabled ) if not answer: raise RuntimeError("Research generation failed") stage_timings["research"] = round(time.perf_counter() - stage_start, 4) # Step 3: Evaluate independently, then apply the enhanced evidence gate. logger.info(f"📊 Step 3: Evaluating answer...") stage_start = time.perf_counter() evaluations = [self.core.evaluate_answer(question, answer)] for _ in range(1, self.config.evaluation_samples): evaluations.append(self.core.evaluate_answer(question, answer)) admission = self.controller.enhanced_gate( question, answer, evaluations, getattr(self.core, "kg", None), ) evaluation = admission["evaluation"] stage_timings["evaluation_and_gate"] = round(time.perf_counter() - stage_start, 4) # Step 4: Developer Mode Analysis dev_analysis = None if developer_mode or self.config.developer_mode: logger.info(f"🔬 Step 4: Developer mode analysis...") stage_start = time.perf_counter() dev_analysis = self.core.developer_mode_analysis(question, answer) stage_timings["developer_analysis"] = round(time.perf_counter() - stage_start, 4) # Step 5: Mint or reject. The controller may only tighten the gate. logger.info(f"💰 Step 5: Minting decision...") mint_check = self._should_mint( evaluation, answer, controller_approved=bool(admission["approved"]), ) minted = False tx_hash = None if mint_check["should_mint"]: logger.info(f"🎉 Minting NES token...") self.core.integrate_knowledge(question, answer, evaluation) tx_hash = self.ledger.mint_nes_token({ "cycle_id": cycle_id, "question": question.to_dict(), "answer": answer.to_dict(), "evaluation": evaluation.to_dict(), "overall_score": mint_check["overall"], "checks_passed": mint_check["passed"], "provider": answer.provider, "governance": { "agenda": agenda.to_dict(), "admission": {key: value for key, value in admission.items() if key != "evaluation"}, }, }) minted = True else: self.core.reject_answer(evaluation) # Calculate execution time elapsed = round(time.time() - start_time, 2) governance = { "agenda": agenda.to_dict(), "admission": {key: value for key, value in admission.items() if key != "evaluation"}, "development_proposals": self.controller.development_proposals(dev_analysis, cycle_id), } # Create result result = OperationCycleResult( cycle_id=cycle_id, timestamp=datetime.now().isoformat(), question=question.to_dict(), answer=answer.to_dict(), evaluation=evaluation.to_dict(), minted=minted, tx_hash=tx_hash, developer_analysis=dev_analysis, governance=governance, execution_time_seconds=elapsed, stage_timings_seconds=stage_timings, retry_count=retry_count, error=None, ) self.history.append(result) self._total_cycles += 1 if minted: self._successful_cycles += 1 if self.config.save_reports: self._save_report(result) logger.info(f"══════════════════════════════════════") logger.info(f"✅ CYCLE {cycle_id} COMPLETE | {'MINTED' if minted else 'REJECTED'} | {elapsed}s") logger.info(f"══════════════════════════════════════") return result except Exception as e: retry_count += 1 last_error = str(e) logger.error(f"⚠️ Cycle {cycle_id} failed (attempt {retry_count}): {e}") if retry_count <= self.config.max_retries: logger.info(f"🔄 Retrying in {self.config.retry_delay}s...") time.sleep(self.config.retry_delay) else: logger.error(f"❌ Cycle {cycle_id} failed after {retry_count} attempts") # All retries failed elapsed = round(time.time() - start_time, 2) result = OperationCycleResult( cycle_id=cycle_id, timestamp=datetime.now().isoformat(), question={"error": last_error}, answer={}, evaluation={}, minted=False, tx_hash=None, developer_analysis=None, governance={"error": "cycle failed before controller decision"}, execution_time_seconds=elapsed, retry_count=retry_count, error=last_error, stage_timings_seconds=stage_timings, ) self.history.append(result) self._total_cycles += 1 if self.config.save_reports: self._save_report(result, is_error=True) return result def _save_report(self, result: OperationCycleResult, is_error: bool = False): """Save cycle report to disk""" try: os.makedirs("reports", exist_ok=True) prefix = "cycle_error" if is_error else "cycle_minted" if result.minted else "cycle_rejected" filename = f"reports/{prefix}_{result.cycle_id}_{datetime.now().strftime('%Y%m%d_%H%M%S')}.json" with open(filename, 'w', encoding='utf-8') as f: json.dump(result.to_dict(), f, indent=4, ensure_ascii=False) logger.debug(f"💾 Report saved: {filename}") self._save_governance_summary() # Sync to HF and GH if result.minted and self.config.sync_to_hf: self._sync_to_huggingface(filename) if result.minted and self.config.sync_to_gh: self._sync_to_github(filename) except Exception as e: logger.error(f"Failed to save report: {e}") def _save_governance_summary(self) -> None: """Persist a compact controller status without leaking runtime secrets.""" reports_dir = "reports" os.makedirs(reports_dir, exist_ok=True) summary_path = os.path.join(reports_dir, "governance_summary.json") try: with open(summary_path, "w", encoding="utf-8") as summary_file: json.dump( self.controller.governance_snapshot(self.history), summary_file, indent=2, ensure_ascii=False, ) except Exception as exc: logger.warning(f"Failed to write governance summary: {exc}") def start(self): """Start the autonomous loop""" if self.is_running: logger.warning("Loop already running") return self.is_running = True logger.info(f"▶️ Loop started: interval={self.config.interval_minutes}min, threshold={self.config.min_accuracy}%") def loop(): while self.is_running: try: self.run_cycle(developer_mode=self.config.developer_mode) except Exception as e: logger.error(f"Cycle error: {e}") if self.is_running: sleep_time = self.config.interval_minutes * 60 logger.info(f"😴 Sleeping for {sleep_time}s until next cycle...") time.sleep(sleep_time) self._thread = threading.Thread(target=loop, daemon=True) self._thread.start() def stop(self): """Stop the autonomous loop""" self.is_running = False if self._thread: self._thread.join(timeout=10) logger.info("⏹️ Loop stopped") def pause(self): """Pause the loop (alias for stop)""" self.stop() def resume(self): """Resume the loop""" if not self.is_running: self.start() def get_stats(self) -> Dict[str, Any]: """Get comprehensive loop statistics""" total = len(self.history) minted = sum(1 for r in self.history if r.minted) rejected = total - minted total_time = sum(r.execution_time_seconds for r in self.history) avg_time = total_time / max(total, 1) recent_cycles = self.history[-10:] if len(self.history) > 10 else self.history recent_minted = sum(1 for r in recent_cycles if r.minted) recent_rate = (recent_minted / max(len(recent_cycles), 1)) * 100 return { "total_cycles": total, "tokens_minted": minted, "tokens_rejected": rejected, "approval_rate": f"{(minted / max(total, 1) * 100):.1f}%", "recent_approval_rate": f"{recent_rate:.1f}%", "average_cycle_time": f"{avg_time:.1f}s", "is_running": self.is_running, "config": { "interval_minutes": self.config.interval_minutes, "min_accuracy": self.config.min_accuracy, "min_novelty": self.config.min_novelty, "min_usefulness": self.config.min_usefulness, "min_overall": self.config.min_overall, "developer_mode": self.config.developer_mode, "web_search_enabled": self.config.web_search_enabled, "max_retries": self.config.max_retries, "sync_to_hf": self.config.sync_to_hf, "sync_to_gh": self.config.sync_to_gh, "evaluation_samples": self.config.evaluation_samples, "evaluation_agreement_threshold": self.config.evaluation_agreement_threshold, }, "governance": self.controller.governance_snapshot(self.history), "last_cycle": self.history[-1].to_dict() if self.history else None, } def get_recent_cycles(self, limit: int = 20) -> List[Dict]: """Get recent cycle results""" return [r.to_dict() for r in self.history[-limit:]] def get_cycle_by_id(self, cycle_id: str) -> Optional[OperationCycleResult]: """Get a specific cycle by ID""" for r in self.history: if r.cycle_id == cycle_id: return r return None def get_best_cycles(self, limit: int = 10) -> List[Dict]: """Get best performing cycles by overall score""" cycles_with_scores = [] for r in self.history: eval_data = r.evaluation if eval_data: score = ( eval_data.get('scientific_accuracy', 0) * 0.45 + eval_data.get('novelty_score', 0) * 0.25 + eval_data.get('usefulness_score', 0) * 0.20 ) cycles_with_scores.append((score, r.to_dict())) cycles_with_scores.sort(key=lambda x: x[0], reverse=True) return [c[1] for _, c in cycles_with_scores[:limit]] __all__ = ['OperationLoop', 'OperationLoopConfig', 'OperationCycleResult']