Spaces:
Running
Running
Download core/operation_loop_v4.py from Qalam/Nuclear-Intelligence: direct link, hf CLI and curl.
- Browser
- Download file 23.4 kB
-
https://huggingface.co/spaces/Qalam/Nuclear-Intelligence/resolve/main/core/operation_loop_v4.py
- Command line
-
hf download hf://spaces/Qalam/Nuclear-Intelligence/core/operation_loop_v4.py
-
curl -L -o operation_loop_v4.py https://huggingface.co/spaces/Qalam/Nuclear-Intelligence/resolve/main/core/operation_loop_v4.py
23.4 kB
| """ | |
| 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 | |
| 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 | |
| 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'] |