shaibu01/Titan-Engine
0
1"""2TITAN QUANT-OS: SOVEREIGN WEB CONSOLE (v401.0 - FULL STACK EDITION)3TYPE: Institutional Command Console4ARCH: Streamlit + ZMQ Async Daemon + Watchdog + Asymmetric Grid5UPGRADE: ZMQ Multi-Port Triangulation (5556, 5557, 5558). Co-boots Alt-Data Oracle.6"""7import config8import os9import sys10import subprocess11import streamlit as st12try:13 from streamlit_autorefresh import st_autorefresh14except Exception:15 st_autorefresh = None16import threading17import zmq18from titan_transport import TitanTransportConfig, configure_socket19import json20import time21import pandas as pd22import numpy as np23import plotly.graph_objects as go24import plotly.express as px25import psutil26import base6427import warnings28import html29from pathlib import Path30 31try:32 import titan_storage as _titan_storage33except Exception:34 _titan_storage = None35 36try:37 from titan_cloud_beast import CloudTitanBrain, ZmqTradeTransport38except Exception as cloud_beast_import_error:39 CloudTitanBrain = None40 ZmqTradeTransport = None41 42try:43 from titan_provider_fabric import (44 PROVIDER_ALPACA,45 PROVIDER_PUPRIME,46 canonical_provider,47 merge_provider_account_maps,48 )49except Exception:50 PROVIDER_ALPACA = "ALPACA"51 PROVIDER_PUPRIME = "PU_PRIME_MT5"52 def canonical_provider(value, default=""):53 text = str(value or "").upper()54 if "ALPACA" in text or "APCA" in text:55 return PROVIDER_ALPACA56 if "PU" in text or "MT5" in text:57 return PROVIDER_PUPRIME58 return default59 def merge_provider_account_maps(*account_maps):60 merged = {}61 for account_map in account_maps:62 if not isinstance(account_map, dict):63 continue64 for provider_name, account in account_map.items():65 if not isinstance(account, dict):66 continue67 provider = canonical_provider(provider_name or account.get("provider"), default="")68 if provider:69 merged[provider] = {**account, "provider": provider}70 return merged71 72from titan_accounting import normalize_provider_account_money, positive_account_value73from titan_focus_control import clear_focus_target, load_focus_target, save_focus_target74 75# ๐จ SILENCE STREAMLIT DEPRECATION SPAM & STOP CPU DRAG ๐จ76warnings.filterwarnings("ignore")77os.environ["STREAMLIT_BROWSER_GATHER_USAGE_STATS"] = "false"78 79ENABLE_BACKGROUND_CORE = os.getenv("TITAN_ENABLE_BACKGROUND_CORE", "1").strip().lower() in {"1", "true", "yes", "on"}80CORE_SUPERVISED = os.getenv("TITAN_CORE_SUPERVISED", "0").strip().lower() in {"1", "true", "yes", "on"}81BRAIN_BUCKET_PREFIX = os.getenv("TITAN_BRAIN_BUCKET_PREFIX", "brain_models").strip("/") or "brain_models"82# Provider intake, execution supervision, and risk exits must never depend on a83# training artifact. Keep the legacy flag visible for diagnostics, but do not84# let it stop the cloud core; neural readiness is enforced at the entry gate.85STRICT_NEURAL_BOOT_LOCK_REQUESTED = os.getenv("TITAN_STRICT_NEURAL_BOOT_LOCK", "0").strip().lower() in {"1", "true", "yes", "on"}86STRICT_NEURAL_BOOT_LOCK = False87NEURAL_FORGE_RETRY_SEC = float(os.getenv("TITAN_NEURAL_FORGE_RETRY_SEC", "60"))88TITAN_RUNTIME_DIR = os.getenv("TITAN_RUNTIME_DIR", "runtime").strip() or "runtime"89CLOUD_RUNTIME_SNAPSHOT_PATH = os.path.join(TITAN_RUNTIME_DIR, "cloud_runtime_snapshot.json")90TRAINER_STATUS_FILE = os.getenv("TITAN_TRAINER_STATUS_FILE", "neural_forge_status.json").strip() or "neural_forge_status.json"91OPTUNA_EVENTS_FILE = os.getenv("TITAN_OPTUNA_EVENTS_FILE", "optuna_trials.jsonl").strip() or "optuna_trials.jsonl"92TRAINER_LOG_FILE = os.getenv("TITAN_TRAINER_LOG_FILE", "neural_trainer_boot.log").strip() or "neural_trainer_boot.log"93CORE_LOG_FILE = os.getenv("TITAN_CORE_LOG_FILE", "titan_core.log").strip() or "titan_core.log"94ALT_DATA_LOG_FILE = os.getenv("TITAN_ALT_DATA_LOG_FILE", "titan_alt_data.log").strip() or "titan_alt_data.log"95 96def _default_brain_models_dir():97 configured = os.getenv("TITAN_MODELS_DIR", "").strip()98 if configured:99 return configured100 if _titan_storage is not None and hasattr(_titan_storage, "local_path"):101 try:102 return str(_titan_storage.local_path(BRAIN_BUCKET_PREFIX))103 except Exception:104 pass105 return "brain_models"106 107BRAIN_MODELS_DIR = _default_brain_models_dir()108_BRAIN_RESTORE_LOCK = threading.Lock()109_BRAIN_LAST_RESTORE_TS = 0.0110_BRAIN_RESTORE_TTL_SEC = float(os.getenv("TITAN_BRAIN_RESTORE_TTL_SEC", "300"))111 112def _brain_file(name):113 return os.path.join(BRAIN_MODELS_DIR, str(name))114 115def _expected_brain_files():116 return ["optimal_params.json"] + [f"expert_{i}_gru.pth" for i in range(1, 51)]117 118def _missing_brain_files():119 return [name for name in _expected_brain_files() if not os.path.exists(_brain_file(name))]120 121def restore_brain_artifacts(force=False):122 global _BRAIN_LAST_RESTORE_TS123 os.makedirs(BRAIN_MODELS_DIR, exist_ok=True)124 if _titan_storage is None or not hasattr(_titan_storage, "restore_named_files"):125 return 0126 missing = _missing_brain_files()127 if not missing:128 return 0129 now = time.time()130 if not force and now - _BRAIN_LAST_RESTORE_TS < _BRAIN_RESTORE_TTL_SEC:131 return 0132 try:133 with _BRAIN_RESTORE_LOCK:134 missing = _missing_brain_files()135 if not missing:136 return 0137 now = time.time()138 if not force and now - _BRAIN_LAST_RESTORE_TS < _BRAIN_RESTORE_TTL_SEC:139 return 0140 restored = int(_titan_storage.restore_named_files(BRAIN_BUCKET_PREFIX, missing, BRAIN_MODELS_DIR))141 _BRAIN_LAST_RESTORE_TS = now142 if restored:143 print(f"๐ง TITAN-NN: Restored {restored}/{len(missing)} missing brain artifact(s) from durable storage.", flush=True)144 return restored145 except Exception as exc:146 print(f"โ ๏ธ TITAN-NN: Durable brain restore skipped: {exc}", flush=True)147 return 0148 149def brain_swarm_status():150 os.makedirs(BRAIN_MODELS_DIR, exist_ok=True)151 experts = sum(1 for i in range(1, 51) if os.path.exists(_brain_file(f"expert_{i}_gru.pth")))152 params_ready = os.path.exists(_brain_file("optimal_params.json"))153 return experts, params_ready, bool(experts >= 50 and params_ready)154 155def write_boot_log(message, mode="a"):156 try:157 log_path = CORE_LOG_FILE158 log_dir = os.path.dirname(os.path.abspath(log_path))159 if log_dir:160 os.makedirs(log_dir, exist_ok=True)161 with open(log_path, mode, encoding="utf-8") as f:162 f.write(str(message).rstrip() + "\n")163 except Exception:164 pass165 166def _boot_bool(name, default=False):167 raw = os.getenv(name, "1" if default else "0").strip().lower()168 return raw in {"1", "true", "yes", "on"}169 170def _boot_float(name, default):171 try:172 return float(os.getenv(name, str(default)))173 except Exception:174 return float(default)175 176def _run_boot_command(label, cmd, timeout_sec=45.0, env=None):177 write_boot_log(f"[{time.strftime('%H:%M:%S')}] {label}: starting {' '.join(map(str, cmd))}")178 try:179 proc = subprocess.run(180 list(cmd),181 stdout=subprocess.PIPE,182 stderr=subprocess.STDOUT,183 text=True,184 timeout=max(1.0, float(timeout_sec)),185 check=False,186 env=env,187 )188 tail = "\n".join((proc.stdout or "").splitlines()[-12:])189 if tail:190 write_boot_log(f"[{time.strftime('%H:%M:%S')}] {label}: output tail:\n{tail}")191 write_boot_log(f"[{time.strftime('%H:%M:%S')}] {label}: exit={proc.returncode}")192 return proc.returncode == 0193 except subprocess.TimeoutExpired as exc:194 partial = exc.stdout or ""195 if isinstance(partial, bytes):196 partial = partial.decode("utf-8", "replace")197 tail = "\n".join(str(partial).splitlines()[-8:])198 write_boot_log(f"[{time.strftime('%H:%M:%S')}] {label}: timed out after {timeout_sec:.1f}s; continuing without blocking core.")199 if tail:200 write_boot_log(f"[{time.strftime('%H:%M:%S')}] {label}: timeout output tail:\n{tail}")201 return False202 except Exception as exc:203 write_boot_log(f"[{time.strftime('%H:%M:%S')}] {label}: failed: {exc}; continuing.")204 return False205 206def _tail_for_boot_log(path, lines=10):207 try:208 if os.path.exists(path):209 with open(path, "r", encoding="utf-8", errors="ignore") as f:210 return "\n".join(f.readlines()[-int(lines):]).strip()211 except Exception:212 return ""213 return ""214 215def _read_json_file(path, default=None):216 try:217 if os.path.exists(path):218 with open(path, "r", encoding="utf-8") as f:219 return json.load(f)220 except Exception:221 return default if default is not None else {}222 return default if default is not None else {}223 224def _tail_text_file(path, lines=30):225 try:226 if os.path.exists(path):227 with open(path, "r", encoding="utf-8", errors="ignore") as f:228 return [line.rstrip("\n") for line in f.readlines()[-int(lines):]]229 except Exception as exc:230 return [f"Unable to read {path}: {exc}"]231 return []232 233def _file_age_sec(path):234 try:235 if os.path.exists(path):236 return max(0.0, time.time() - os.path.getmtime(path))237 except Exception:238 pass239 return None240 241def _tail_jsonl_file(path, lines=12):242 events = []243 for raw in _tail_text_file(path, lines=lines):244 try:245 events.append(json.loads(raw))246 except Exception:247 if raw:248 events.append({"raw": raw})249 return events250 251def _html_tail(path, lines=30, empty="No logs yet."):252 rows = _tail_text_file(path, lines=lines)253 if not rows:254 return html.escape(empty)255 return "<br>".join(html.escape(row) for row in rows)256 257def _short_json(value, max_len=220):258 try:259 text = json.dumps(value, ensure_ascii=False, sort_keys=True)260 except Exception:261 text = str(value)262 if len(text) > max_len:263 return text[:max_len - 3] + "..."264 return text265 266def _render_optuna_events(events):267 if not events:268 return "<span style='color:#8A8D9E;'>No Optuna trial events yet.</span>"269 rows = []270 for event in events:271 if "raw" in event:272 rows.append(html.escape(str(event["raw"])))273 continue274 ev = str(event.get("event", "event"))275 trial = event.get("trial", "?")276 total = event.get("total_trials", "?")277 value = event.get("value", event.get("best_value"))278 best = event.get("best_value")279 value_txt = "n/a" if value is None else f"{_safe_float(value, 0.0):.6f}"280 best_txt = "" if best is None else f" best={_safe_float(best, 0.0):.6f}"281 params = event.get("params") or event.get("best_params") or {}282 rows.append(283 html.escape(284 f"{event.get('ts', '')} | {ev} | trial={trial}/{total} "285 f"value={value_txt}{best_txt} params={_short_json(params, 160)}"286 )287 )288 return "<br>".join(rows)289 290def brain_env(extra=None):291 env = os.environ.copy()292 env["TITAN_BRAIN_BUCKET_PREFIX"] = BRAIN_BUCKET_PREFIX293 env["TITAN_MODELS_DIR"] = BRAIN_MODELS_DIR294 env["TITAN_RUNTIME_DIR"] = TITAN_RUNTIME_DIR295 env["TITAN_TRAINER_STATUS_FILE"] = TRAINER_STATUS_FILE296 env["TITAN_OPTUNA_EVENTS_FILE"] = OPTUNA_EVENTS_FILE297 env["TITAN_TRAINER_LOG_FILE"] = TRAINER_LOG_FILE298 env["TITAN_CORE_LOG_FILE"] = CORE_LOG_FILE299 env["TITAN_ALT_DATA_LOG_FILE"] = ALT_DATA_LOG_FILE300 if extra:301 env.update({str(k): str(v) for k, v in dict(extra).items()})302 return env303 304# ==============================================================================305# 1. PAGE SETUP & ASSET LOADING306# ==============================================================================307st.set_page_config(page_title="TITAN NEXUS | CORE", layout="wide", initial_sidebar_state="collapsed")308if st_autorefresh is not None and os.getenv("TITAN_AUTO_REFRESH", "1").strip().lower() in {"1", "true", "yes", "on"}:309 # A full Streamlit rerun faster than the page can render causes partial panels,310 # scroll jumps, and overlapping refreshes on the HF CPU tier.311 refresh_seconds = max(15.0, float(os.getenv("TITAN_AUTO_REFRESH_SEC", "20")))312 st_autorefresh(interval=int(refresh_seconds * 1000), limit=None, key="titan_live_console_refresh")313 314LOGO_FILENAME = "Gemini_Generated_Image_gucvqzgucvqzgucv.png"315 316def load_logo_base64():317 if Path(LOGO_FILENAME).exists():318 with open(LOGO_FILENAME, "rb") as f: return base64.b64encode(f.read()).decode()319 return None320 321logo_b64 = load_logo_base64()322if logo_b64:323 header_logo = f'<img src="data:image/png;base64,{logo_b64}" style="height: 65px; width: 65px; border-radius: 14px; margin-right: 18px; box-shadow: 0 4px 20px rgba(0, 229, 255, 0.4); border: 1px solid rgba(0,255,170,0.4);">'324else:325 header_logo = '<div style="height: 65px; width: 65px; border-radius: 14px; margin-right: 18px; background: #00FFAA; display: flex; align-items: center; justify-content: center; font-size: 30px; box-shadow: 0 4px 20px rgba(0, 255, 170, 0.4);">๐ </div>'326 327def _safe_float(value, default=0.0):328 try:329 out = float(value)330 return out if np.isfinite(out) else float(default)331 except Exception:332 return float(default)333 334def _packet_route_provider(packet):335 if not isinstance(packet, dict):336 return "--"337 extras = packet.get("extras", {}) if isinstance(packet.get("extras", {}), dict) else {}338 provider = (339 packet.get("route_provider")340 or extras.get("route_provider")341 or packet.get("provider")342 or packet.get("execution_provider")343 or packet.get("broker_provider")344 or extras.get("execution_provider")345 or extras.get("broker_provider")346 or extras.get("provider")347 or "--"348 )349 return canonical_provider(provider, default=str(provider or "--")).upper()350 351def _allocator_provider_rows(provider_allocators, providers_snapshot):352 rows = []353 if not isinstance(provider_allocators, dict):354 return rows355 providers_snapshot = providers_snapshot if isinstance(providers_snapshot, dict) else {}356 for provider_name, allocator in sorted(provider_allocators.items()):357 if not isinstance(allocator, dict):358 continue359 provider = canonical_provider(provider_name or allocator.get("provider"), default=str(provider_name or "--")).upper()360 provider_state = providers_snapshot.get(provider, {}) if isinstance(providers_snapshot.get(provider, {}), dict) else {}361 daily_pnl = _safe_float(allocator.get("daily_pnl_pct"), 0.0)362 reserve_pct = _safe_float(allocator.get("cash_reserve_pct"), 0.0)363 cooldown_left = max(0.0, _safe_float(allocator.get("cooldown_until"), 0.0) - time.time())364 rows.append({365 "Provider": provider,366 "Lane state": provider_state.get("state", "--"),367 "Equity": _safe_float(provider_state.get("equity"), 0.0),368 "Available": _safe_float(369 provider_state.get("available_funds", provider_state.get("buying_power", provider_state.get("margin_free", 0.0))),370 0.0,371 ),372 "Daily PnL": daily_pnl,373 "Reserve": reserve_pct,374 "Margin level": _safe_float(allocator.get("broker_margin_level_pct", allocator.get("margin_level_pct")), 0.0),375 "Margin heat": _safe_float(allocator.get("broker_margin_heat", allocator.get("margin_utilization")), 0.0),376 "Loss streak": int(_safe_float(allocator.get("loss_streak"), 0.0)),377 "Trades/hr": len(allocator.get("trade_times", []) or []),378 "Cooldown sec": cooldown_left,379 "Telemetry": bool(allocator.get("telemetry_ready")),380 "Reset reason": allocator.get("baseline_reset_reason", ""),381 })382 return rows383 384def _provider_feed_summary(provider_states):385 """Summarize provider-scoped lanes for UI truth without merging prices."""386 summary = {387 "provider_count": 0,388 "online_count": 0,389 "ready_count": 0,390 "quote_count": 0,391 "bar_symbols": 0,392 "real_bar_rows": 0,393 "fresh_quote_count": 0,394 "fresh_account_count": 0,395 "min_quote_age_sec": 9999.0,396 "max_account_age_sec": 0.0,397 "status": "AWAITING_PROVIDER_FEEDS",398 }399 if not isinstance(provider_states, dict):400 return summary401 now = time.time()402 for provider, state in provider_states.items():403 if not isinstance(state, dict):404 continue405 provider_name = str(provider or state.get("provider") or "").upper()406 if not provider_name:407 continue408 summary["provider_count"] += 1409 lane_state = str(state.get("state") or "").upper()410 connected = bool(state.get("connected") or state.get("present") or lane_state in {"READY", "ONLINE", "DEGRADED", "QUOTE_STALE", "ACCOUNT_STALE", "WARMING_BARS", "WAITING_FEED"})411 if connected:412 summary["online_count"] += 1413 quotes = int(max(0.0, _safe_float(state.get("quotes"), 0.0)))414 bars = int(max(0.0, _safe_float(state.get("bars"), 0.0)))415 real_rows = int(max(0.0, _safe_float(state.get("real_bar_rows"), state.get("max_bar_rows", 0.0))))416 quote_age = _safe_float(state.get("quote_age_sec"), 9999.0)417 account_age = _safe_float(state.get("account_age_sec"), 9999.0)418 summary["quote_count"] += quotes419 summary["bar_symbols"] += bars420 summary["real_bar_rows"] = max(summary["real_bar_rows"], real_rows)421 summary["min_quote_age_sec"] = min(summary["min_quote_age_sec"], quote_age)422 summary["max_account_age_sec"] = max(summary["max_account_age_sec"], account_age if account_age < 9999.0 else 0.0)423 if quote_age <= 6.0 and quotes > 0:424 summary["fresh_quote_count"] += 1425 if account_age <= 30.0:426 summary["fresh_account_count"] += 1427 if bool(state.get("candidate_ready")) or lane_state == "READY":428 summary["ready_count"] += 1429 if summary["ready_count"] > 1:430 summary["status"] = "ONLINE_DUAL_PROVIDER"431 elif summary["ready_count"] == 1:432 summary["status"] = "ONLINE_PROVIDER"433 elif summary["online_count"] and summary["bar_symbols"] > 0 and summary["quote_count"] > 0:434 summary["status"] = "PROVIDER_SCOPED_FEEDS_WARMING"435 elif summary["online_count"] and summary["bar_symbols"] > 0:436 summary["status"] = "PROVIDER_BARS_ONLINE_WAITING_QUOTES"437 elif summary["online_count"]:438 summary["status"] = "BROKER_TELEMETRY_ONLINE_WAITING_PROVIDER_FEEDS"439 return summary440 441def _provider_diagnostic_freshness(candidate_diagnostics):442 """Use provider-scoped candidate diagnostics as a freshness hint only."""443 freshness = {}444 for diagnostic in candidate_diagnostics or []:445 if not isinstance(diagnostic, dict):446 continue447 provider = str(diagnostic.get("provider", "") or "").strip().upper()448 if not provider:449 continue450 slot = freshness.setdefault(provider, {"quote_age_sec": 9999.0, "bar_count": 0})451 slot["quote_age_sec"] = min(452 _safe_float(slot.get("quote_age_sec"), 9999.0),453 _safe_float(diagnostic.get("quote_age_sec"), 9999.0),454 )455 slot["bar_count"] = max(456 int(_safe_float(slot.get("bar_count"), 0)),457 int(_safe_float(diagnostic.get("bar_count"), 0)),458 )459 return freshness460 461def _safe_prob_map(value, fallback=None):462 fallback = fallback or {"bull": 0.33, "bear": 0.33, "static": 0.34}463 if not isinstance(value, dict) or not value:464 return dict(fallback)465 cleaned = {str(k): max(0.0, _safe_float(v, 0.0)) for k, v in value.items()}466 total = sum(cleaned.values())467 if total <= 0.0:468 return dict(fallback)469 return {k: v / total for k, v in cleaned.items()}470 471def _display_cell(value):472 if value is None:473 return ""474 try:475 if isinstance(value, float) and not np.isfinite(value):476 return ""477 except Exception:478 pass479 if isinstance(value, (dict, list, tuple)):480 try:481 return json.dumps(value, ensure_ascii=False)482 except Exception:483 return str(value)484 if isinstance(value, (bool, np.bool_)):485 return "true" if bool(value) else "false"486 return str(value)487 488def _arrow_safe_display_df(df):489 if not isinstance(df, pd.DataFrame) or df.empty:490 return df491 out = df.copy()492 for col in out.columns:493 if out[col].dtype == "object" or str(out[col].dtype).startswith(("bool", "boolean")):494 out[col] = out[col].map(_display_cell).astype("string")495 return out496 497def _trace_signature(rec):498 if not isinstance(rec, dict):499 return str(rec)500 payload = rec.get("payload", {})501 try:502 payload_key = json.dumps(payload, sort_keys=True, default=str, separators=(",", ":"))503 except Exception:504 payload_key = str(payload)505 return "|".join([506 str(rec.get("kind", "")).upper(),507 str(rec.get("symbol", "")).upper(),508 payload_key,509 ])510 511def _dedupe_trace_items(trace_items, limit=50):512 if not isinstance(trace_items, list):513 return []514 seen = set()515 cleaned = []516 for rec in reversed(trace_items):517 key = _trace_signature(rec)518 if key in seen:519 continue520 seen.add(key)521 cleaned.append(rec)522 if len(cleaned) >= max(1, int(limit)):523 break524 return list(reversed(cleaned))525 526def _trade_widget_key(prefix, *parts):527 raw = "_".join(str(part or "") for part in parts)528 cleaned = "".join(ch if ch.isalnum() else "_" for ch in raw)529 return f"{prefix}_{cleaned[:90]}"530 531def _trade_progress_figure(trade, title="Trade Progress"):532 if not isinstance(trade, dict):533 return None534 side = str(trade.get("side", "") or "").upper()535 direction = -1.0 if side == "SELL" else 1.0536 entry = _safe_float(537 trade.get(538 "entry_price",539 trade.get("price_open", trade.get("open_price", trade.get("avg_price", 0.0))),540 ),541 0.0,542 )543 mark = _safe_float(544 trade.get(545 "last_mark_price",546 trade.get("current_price", trade.get("mark_price", trade.get("price", 0.0))),547 ),548 0.0,549 )550 if entry <= 0.0 and mark > 0.0:551 entry = mark552 if mark <= 0.0 and entry > 0.0:553 mark = entry554 if entry <= 0.0 or mark <= 0.0:555 return None556 stop = _safe_float(trade.get("sl_price", trade.get("stop_price", trade.get("stop_loss_price", 0.0))), 0.0)557 target = _safe_float(trade.get("tp_price", trade.get("target_price", trade.get("take_profit_price", 0.0))), 0.0)558 stop_pct = abs(_safe_float(trade.get("stop_loss_pct"), 0.0))559 target_pct = abs(_safe_float(trade.get("take_profit_pct", trade.get("min_net_pnl_pct", 0.0)), 0.0))560 if stop <= 0.0 and stop_pct > 0.0:561 stop = entry * (1.0 - direction * stop_pct)562 if target <= 0.0 and target_pct > 0.0:563 target = entry * (1.0 + direction * target_pct)564 benchmark = _safe_float(trade.get("benchmark_price"), 0.0)565 benchmark_pct = abs(_safe_float(trade.get("min_net_pnl_pct"), 0.0))566 if benchmark <= 0.0 and benchmark_pct > 0.0:567 benchmark = entry * (1.0 + direction * benchmark_pct)568 net_pct = _safe_float(trade.get("last_net_pnl_pct", trade.get("pnl_pct", 0.0)), 0.0)569 line_color = "#00D084" if net_pct >= 0.0 else "#FF5A52"570 fig = go.Figure()571 fig.add_trace(go.Scatter(572 x=["Entry", "Current"],573 y=[entry, mark],574 mode="lines+markers",575 name="Mark",576 line=dict(color=line_color, width=3),577 marker=dict(color=line_color, size=8),578 hovertemplate="%{x}: %{y:.6g}<extra></extra>",579 ))580 reference_lines = [581 ("Entry", entry, "#8A8D9E", "dot"),582 ("Benchmark", benchmark, "#55C2FF", "dash"),583 ("Take Profit", target, "#00D084", "dash"),584 ("Stop Loss", stop, "#FF5A52", "dash"),585 ]586 price_values = [entry, mark]587 for label, value, color, dash in reference_lines:588 if value <= 0.0:589 continue590 price_values.append(value)591 fig.add_trace(go.Scatter(592 x=["Entry", "Current"],593 y=[value, value],594 mode="lines",595 name=label,596 line=dict(color=color, width=1.6, dash=dash),597 hovertemplate=f"{label}: %{{y:.6g}}<extra></extra>",598 ))599 y_min, y_max = min(price_values), max(price_values)600 pad = max((y_max - y_min) * 0.18, abs(entry) * 0.0005, 0.000001)601 fig.update_layout(602 **plotly_layout,603 height=185,604 margin=dict(l=4, r=4, t=28, b=4),605 title=dict(text=title, font=dict(size=12), x=0.0),606 xaxis=dict(showgrid=False, color="#8A8D9E"),607 yaxis=dict(showgrid=True, gridcolor="rgba(255,255,255,0.06)", range=[y_min - pad, y_max + pad], color="#8A8D9E"),608 legend=dict(orientation="h", yanchor="bottom", y=1.02, xanchor="right", x=1.0, font=dict(size=9)),609 showlegend=True,610 )611 return fig612 613def _candidate_preview_trade(candidate, selected_diag=None, default_min_net_pct=0.01):614 if not isinstance(candidate, dict) or not candidate:615 return {}616 selected_diag = selected_diag if isinstance(selected_diag, dict) else {}617 provider = str(candidate.get("provider") or selected_diag.get("provider") or "--").strip().upper()618 symbol = str(candidate.get("symbol") or selected_diag.get("symbol") or "--").strip().upper()619 side = str(candidate.get("side") or selected_diag.get("side") or "BUY").strip().upper()620 diag_matches = (621 selected_diag622 and str(selected_diag.get("provider", "")).strip().upper() == provider623 and str(selected_diag.get("symbol", "")).strip().upper() == symbol624 )625 diag = selected_diag if diag_matches else {}626 entry = _safe_float(627 candidate.get(628 "entry_price",629 candidate.get("price", candidate.get("last_mark_price", diag.get("entry_price", diag.get("price", 0.0)))),630 ),631 0.0,632 )633 mark = _safe_float(634 candidate.get("last_mark_price", candidate.get("price", diag.get("last_mark_price", diag.get("price", entry)))),635 entry,636 )637 if entry <= 0.0 and mark > 0.0:638 entry = mark639 if entry <= 0.0:640 return {}641 min_net = max(0.0, _safe_float(candidate.get("min_net_pnl_pct"), default_min_net_pct))642 loss_limit = abs(_safe_float(candidate.get("loss_limit_pct"), 0.006))643 direction = -1.0 if side == "SELL" else 1.0644 return {645 "provider": provider,646 "symbol": symbol,647 "side": side,648 "qty": _safe_float(candidate.get("qty"), 0.0),649 "entry_price": entry,650 "last_mark_price": mark if mark > 0.0 else entry,651 "price": mark if mark > 0.0 else entry,652 "sl_price": _safe_float(candidate.get("sl_price", candidate.get("stop_price")), 0.0),653 "tp_price": _safe_float(candidate.get("tp_price", candidate.get("target_price")), 0.0),654 "stop_loss_pct": loss_limit,655 "take_profit_pct": max(min_net, _safe_float(candidate.get("take_profit_pct"), 0.0)),656 "min_net_pnl_pct": min_net,657 "benchmark_price": _safe_float(candidate.get("benchmark_price"), entry * (1.0 + direction * min_net)),658 "last_net_pnl_pct": _safe_float(candidate.get("last_net_pnl_pct"), 0.0),659 "last_gross_pnl_pct": _safe_float(candidate.get("last_gross_pnl_pct"), 0.0),660 "last_quote_age_sec": _safe_float(candidate.get("last_quote_age_sec", candidate.get("quote_age_sec")), 0.0),661 "selection_score": _safe_float(candidate.get("selection_score"), 0.0),662 "selection_reasons": candidate.get("selection_reasons", []) if isinstance(candidate.get("selection_reasons", []), list) else [],663 }664 665# ==============================================================================666# 2. DEEP CSS OVERRIDE667# ==============================================================================668st.markdown("""669<style>670 @import url('https://fonts.googleapis.com/css2?family=Inter:wght@300;400;600;800&family=JetBrains+Mono:wght@400;700&display=swap');671 672 .stApp { background-color: #020306; color: #F5F5F7; font-family: 'Inter', sans-serif; }673 header {visibility: hidden;} #MainMenu {visibility: hidden;} footer {visibility: hidden;}674 .block-container { padding: 1.5rem 2.5rem 3rem 2.5rem !important; max-width: 100%; }675 676 h1, h2, h3 { font-family: 'Inter', sans-serif; font-weight: 800; letter-spacing: -0.03em; color: #FFFFFF; margin-bottom: 0.1rem; }677 h4, h5, h6 { font-family: 'Inter', sans-serif; font-weight: 600; color: #8A8D9E; text-transform: uppercase; letter-spacing: 1.5px; font-size: 0.8rem; margin-top: 0; margin-bottom: 12px;}678 679 .section-title { font-size: 1.4rem; font-weight: 800; color: #FFFFFF; padding-top: 2.5rem; padding-bottom: 0.5rem; margin-top: 1rem; margin-bottom: 1.5rem; border-bottom: 1px solid rgba(255,255,255,0.08); text-shadow: 0 0 20px rgba(0, 229, 255, 0.1); }680 681 .nexus-top-bar { display: flex; justify-content: space-between; align-items: center; background: linear-gradient(180deg, rgba(12,13,18,0.9) 0%, rgba(4,5,8,0.95) 100%); border: 1px solid rgba(255,255,255,0.06); border-radius: 18px; padding: 20px 30px; margin-bottom: 25px; box-shadow: 0 10px 40px rgba(0,0,0,0.8); }682 .nexus-brand { display: flex; align-items: center; }683 .nexus-status-grid { display: flex; gap: 40px; }684 .status-item { display: flex; flex-direction: column; }685 .status-label { font-size: 10px; color: #8A8D9E; font-weight: 800; text-transform: uppercase; letter-spacing: 1px; }686 .status-val { font-size: 14px; font-weight: 700; font-family: 'JetBrains Mono', monospace; }687 688 div[data-testid="metric-container"] { background: linear-gradient(145deg, rgba(15,16,22,0.8) 0%, rgba(5,6,10,0.9) 100%); backdrop-filter: blur(25px); -webkit-backdrop-filter: blur(25px); border: 1px solid rgba(255, 255, 255, 0.05); border-radius: 14px; padding: 18px 25px; box-shadow: 0 8px 25px rgba(0,0,0,0.5), inset 0 1px 0 rgba(255,255,255,0.05); transition: transform 0.2s ease, border-color 0.2s ease; }689 div[data-testid="metric-container"]:hover { border-color: rgba(0, 229, 255, 0.3); transform: translateY(-2px); box-shadow: 0 12px 30px rgba(0, 229, 255, 0.1); }690 .stMetric-value { color: #00FFAA !important; font-weight: 700; font-size: 2.2rem !important; font-family: 'JetBrains Mono', monospace !important; text-shadow: 0 0 10px rgba(0, 255, 170, 0.2); }691 .stMetric-label { color: #8A8D9E !important; font-size: 0.8rem !important; font-weight: 700; text-transform: uppercase; letter-spacing: 1px; }692 693 .console-box { background: rgba(6, 7, 10, 0.85); backdrop-filter: blur(15px); border: 1px solid rgba(255, 255, 255, 0.05); border-radius: 12px; padding: 18px; font-family: 'JetBrains Mono', monospace; font-size: 11.5px; overflow-y: auto; box-shadow: inset 0 0 25px rgba(0,0,0,0.6); line-height: 1.5; margin-bottom: 20px; }694 .worm-box { border-left: 3px solid #00FFAA; height: 260px; color: #FFFFFF; }695 .cognition-box { border-left: 3px solid #FFB000; height: 260px; color: #00E5FF; }696 .uplink-box { border-left: 3px solid #B026FF; height: 260px; color: #E8E8ED; }697 .debug-box { border-left: 3px solid #FFB000; color: #FFB000; background: rgba(255, 176, 0, 0.05); font-size: 13px; margin-bottom: 20px;}698 699 ::-webkit-scrollbar { width: 5px; height: 5px; }700 ::-webkit-scrollbar-track { background: transparent; }701 ::-webkit-scrollbar-thumb { background: rgba(138, 141, 158, 0.3); border-radius: 10px; }702 ::-webkit-scrollbar-thumb:hover { background: rgba(0, 229, 255, 0.5); }703 704 [data-testid="stDataFrame"] { border-radius: 12px; overflow: hidden; border: 1px solid rgba(255,255,255,0.05); }705</style>706""", unsafe_allow_html=True)707 708# ==============================================================================709# 3. ZMQ ASYNC DAEMON710# ==============================================================================711class ZMQManager:712 def __init__(self):713 self.transport = TitanTransportConfig.from_env()714 self.hub_ip = self.transport.host715 self.core_port = int(os.getenv("TITAN_ZMQ_CORE_PORT", str(self.transport.bus_out_port)))716 self.data_port = int(os.getenv("TITAN_ZMQ_DATA_PORT", "5557"))717 self.alt_port = int(os.getenv("TITAN_ZMQ_ALT_PORT", "5558"))718 self.legacy_multiport = os.getenv("TITAN_ZMQ_LEGACY_MULTIPORT", "0").strip().lower() in {"1", "true", "yes", "on"}719 self.cache = {720 "cloud_telem": {}, "mt5_telem": {}, "uplink": {}, "alt_data": {}, "ticks": {}, "last_sync": 0,721 "start_time": time.time(), "thread_status": "Booting", 722 "last_error": "None", "msgs_received": 0, "boot_phase": "SYSTEM CHECKS",723 "provider_records_hydrated": 0,724 "last_topic": "NONE",725 "provider_accounts": {},726 "provider_accounts_last_good": {},727 "cloud_beast": CloudTitanBrain() if CloudTitanBrain else None,728 "cloud_guardian": {},729 "provider_snapshot_hydrated_ts": 0.0,730 }731 self.trade_transport = ZmqTradeTransport(host=self.hub_ip) if ZmqTradeTransport else None732 if self.trade_transport is not None:733 try:734 self.trade_transport.warmup()735 except Exception as exc:736 self.cache["last_error"] = f"Trade Transport Warmup Error: {exc}"737 self.last_failover_tick = 0.0738 self.last_entry_tick = 0.0739 self.lock = threading.Lock()740 self.start_listener()741 742 def start_listener(self):743 threading.Thread(target=self.listener_thread, daemon=True).start()744 threading.Thread(target=self.guardian_thread, daemon=True).start()745 threading.Thread(target=self.entry_thread, daemon=True).start()746 747 def guardian_thread(self):748 while True:749 try:750 self.maybe_run_cloud_guardian()751 except Exception as exc:752 with self.lock:753 self.cache["last_error"] = f"Cloud Guardian Timer Error: {exc}"754 time.sleep(3.0)755 756 def entry_thread(self):757 while True:758 try:759 self.maybe_run_cloud_entry()760 except Exception as exc:761 with self.lock:762 self.cache["last_error"] = f"Cloud Entry Timer Error: {exc}"763 time.sleep(3.0)764 765 def listener_thread(self):766 self.cache["thread_status"] = "Running"767 context = zmq.Context()768 poller = zmq.Poller()769 sub = context.socket(zmq.SUB)770 configure_socket(sub, zmq, "RECEIVE", self.transport)771 try:772 # The relay publishes every provider/control topic on one outbound773 # bus. Connecting a SUB socket to ingress ports creates reconnect774 # churn and never yields valid packets.775 sub.connect(f"tcp://{self.hub_ip}:{self.core_port}")776 if self.legacy_multiport:777 sub.connect(f"tcp://{self.hub_ip}:{self.data_port}")778 sub.connect(f"tcp://{self.hub_ip}:{self.alt_port}")779 poller.register(sub, zmq.POLLIN)780 except Exception as e:781 self.cache["last_error"] = f"Bind Error: {e}"782 783 sub.setsockopt_string(zmq.SUBSCRIBE, "")784 785 while True:786 try:787 events = dict(poller.poll(1000))788 except Exception as e:789 self.cache["last_error"] = f"Poll Error: {e}"790 time.sleep(0.1)791 continue792 793 for sock in (sub,):794 if sock not in events:795 continue796 try:797 msg_bytes = sock.recv()798 self.cache["msgs_received"] += 1799 except Exception as e:800 self.cache["last_error"] = f"Recv Error: {e}"801 continue802 803 try:804 msg = msg_bytes.decode('utf-8', errors='ignore')805 if " " not in msg:806 continue807 topic, payload_str = msg.split(" ", 1)808 data = json.loads(payload_str)809 810 with self.lock:811 self.cache["last_topic"] = topic812 beast = self.cache.get("cloud_beast")813 accept_external_alpaca = os.getenv("TITAN_CLOUD_ACCEPT_EXTERNAL_ALPACA_STATE", "0").strip().lower() in {"1", "true", "yes", "on"}814 if beast and (topic != "ALPACA_STATE" or accept_external_alpaca):815 beast.ingest_bus_message(topic, data)816 if topic == "CLOUD_TELEM":817 self.cache["cloud_telem"] = data818 self.cache["last_sync"] = time.time()819 self.cache["boot_phase"] = "LIVE"820 elif topic in {"TELEMETRY", "ACCOUNT", "MT5_STATE", "BROKER_ACCOUNT", "ACCOUNT_STATE"}:821 now = time.time()822 broker_payload = dict(data or {})823 claimed_provider = canonical_provider(824 broker_payload.get("provider") or broker_payload.get("broker_provider"),825 default="",826 )827 if claimed_provider and claimed_provider != PROVIDER_PUPRIME:828 self.cache["last_error"] = f"Provider identity quarantine: TELEMETRY claimed {claimed_provider}"829 continue830 provider = PROVIDER_PUPRIME831 broker_payload.update({832 "provider": provider,833 "broker_provider": provider,834 "data_provider": provider,835 "execution_provider": provider,836 })837 broker_payload = normalize_provider_account_money(broker_payload, provider=provider)838 if _safe_float(broker_payload.get("equity", 0.0), 0.0) <= 0.0 and not str(broker_payload.get("account_id", broker_payload.get("login", "")) or "").strip():839 continue840 broker_payload["telemetry_ready"] = True841 broker_payload["received_at"] = now842 self.cache["provider_accounts"][provider] = broker_payload843 self.cache["provider_accounts_last_good"][provider] = dict(broker_payload)844 self.cache["mt5_telem"] = broker_payload845 if topic == "MT5_STATE":846 self.cache["mt5_state"] = dict(data or {})847 self.cache["last_sync"] = now848 elif topic == "ALPACA_STATE":849 if not accept_external_alpaca:850 continue851 now = time.time()852 account_payload = dict(data or {})853 account_payload["provider"] = PROVIDER_ALPACA854 account_payload["received_at"] = now855 self.cache["provider_accounts"][PROVIDER_ALPACA] = account_payload856 self.cache["provider_accounts_last_good"][PROVIDER_ALPACA] = dict(account_payload)857 self.cache["last_sync"] = now858 elif topic == "CLOUD_ALPACA_ACCOUNT":859 now = time.time()860 account_payload = normalize_provider_account_money(dict(data or {}), provider=PROVIDER_ALPACA)861 account_payload.update({"provider": PROVIDER_ALPACA, "received_at": now})862 self.cache["provider_accounts"][PROVIDER_ALPACA] = account_payload863 self.cache["provider_accounts_last_good"][PROVIDER_ALPACA] = dict(account_payload)864 self.cache["last_sync"] = now865 elif topic in {"CLOUD_ALPACA_QUOTE", "CLOUD_ALPACA_BARS"}:866 self.cache["last_sync"] = time.time()867 elif topic == "MAC_UPLINK":868 self.cache["uplink"] = data869 elif topic == "MAC_ALT":870 self.cache["alt_data"] = data871 self.maybe_run_cloud_guardian()872 except Exception as e:873 self.cache["last_error"] = f"Parse Error: {e}"874 875 def maybe_run_cloud_guardian(self):876 if os.getenv("TITAN_CLOUD_FAILOVER_AUTORUN", "1") == "0":877 return878 beast = self.cache.get("cloud_beast")879 if not beast or ZmqTradeTransport is None:880 return881 now = time.time()882 try:883 interval = float(os.getenv("TITAN_CLOUD_FAILOVER_INTERVAL_SEC", "15"))884 except Exception:885 interval = 15.0886 if now - self.last_failover_tick < max(3.0, interval):887 return888 self.last_failover_tick = now889 try:890 sender = self.trade_transport or ZmqTradeTransport(host=self.hub_ip)891 result = beast.deploy_failover_guardian(sender)892 with self.lock:893 self.cache["cloud_guardian"] = result894 except Exception as e:895 with self.lock:896 self.cache["last_error"] = f"Cloud Guardian Error: {e}"897 898 def maybe_run_cloud_entry(self):899 autorun = os.getenv(900 "TITAN_CLOUD_BEAST_AUTORUN",901 os.getenv("TITAN_CLOUD_AUTONOMOUS_ENTRIES_ENABLED", "1"),902 ).strip().lower()903 if autorun in {"0", "false", "no", "off", "disabled"}:904 return905 beast = self.cache.get("cloud_beast")906 if not beast or ZmqTradeTransport is None:907 return908 now = time.time()909 try:910 interval = float(os.getenv("TITAN_CLOUD_BEAST_ENTRY_INTERVAL_SEC", "12"))911 except Exception:912 interval = 12.0913 if now - self.last_entry_tick < max(5.0, interval):914 return915 if int(self.cache.get("msgs_received", 0) or 0) <= 0 and os.getenv("TITAN_CLOUD_BEAST_ALLOW_NO_BUS", "0") != "1":916 return917 918 allow_concurrent = os.getenv("TITAN_CLOUD_BEAST_ALLOW_CONCURRENT", "0").strip().lower() in {"1", "true", "yes", "on"}919 if not allow_concurrent:920 if hasattr(beast, "active_order_counts"):921 order_counts = beast.active_order_counts()922 active_order_count = int(order_counts.get("broker_active_orders", order_counts.get("active_orders", 0)) or 0)923 pending_handoff_count = int(order_counts.get("pending_handoffs", 0) or 0)924 else:925 terminal_statuses = {926 "FILLED", "REJECTED", "CANCELLED", "CANCELED", "EXPIRED", "DONE",927 "CONFIRM_TIMEOUT", "EXPIRED_UNCONFIRMED", "BROKER_REJECTED", "TRANSPORT_FAILED",928 }929 active_order_count = sum(930 1 for order in getattr(beast, "orders", {}).values()931 if str(getattr(order, "status", "")).upper() not in terminal_statuses932 and "AWAITING" not in str(getattr(order, "status", "")).upper()933 )934 pending_handoff_count = sum(935 1 for order in getattr(beast, "orders", {}).values()936 if "AWAITING" in str(getattr(order, "status", "")).upper()937 )938 if active_order_count or pending_handoff_count:939 wait_reason = "ACTIVE_ORDER_WAIT" if active_order_count else "BROKER_CONFIRM_WAIT"940 order_states = beast.active_order_states(limit=8) if hasattr(beast, "active_order_states") else []941 with self.lock:942 self.cache["cloud_beast_entry"] = {943 "sent": False,944 "reason": wait_reason,945 "active_orders": active_order_count,946 "pending_handoffs": pending_handoff_count,947 "order_states": order_states,948 "policy": "do not send another trade until the broker confirms or rejects the previous handoff",949 "dry_run": True,950 }951 self.last_entry_tick = now952 return953 954 self.last_entry_tick = now955 try:956 sender = self.trade_transport or ZmqTradeTransport(host=self.hub_ip)957 result = beast.deploy_top(sender)958 with self.lock:959 self.cache["cloud_beast_entry"] = result960 except Exception as e:961 with self.lock:962 self.cache["last_error"] = f"Cloud Entry Error: {e}"963 964@st.cache_resource965def get_zmq_manager(): 966 return ZMQManager()967 968# ==============================================================================969# 4. IMMORTAL WATCHDOG (STRICT BOOT LOCK)970# ==============================================================================971@st.cache_resource972def get_watchdog():973 class CoreWatchdog:974 def __init__(self):975 self.core_proc = None976 self.trainer_proc = None977 self.alt_proc = None978 self.is_booting = False 979 self.neural_forge_required = False980 self.last_neural_forge_start = 0.0981 self.start_core()982 threading.Thread(target=self.monitor_loop, daemon=True).start()983 984 def start_neural_forge(self, reason="missing_swarm", force=True):985 restore_brain_artifacts(force=True)986 experts_found, params_ready, swarm_ready = brain_swarm_status()987 if swarm_ready:988 self.neural_forge_required = False989 return True990 now = time.time()991 if getattr(self, 'trainer_proc', None) is not None and self.trainer_proc.poll() is None:992 return False993 if now - float(self.last_neural_forge_start or 0.0) < NEURAL_FORGE_RETRY_SEC:994 return False995 self.neural_forge_required = True996 self.last_neural_forge_start = now997 zmq_manager.cache["boot_phase"] = "FORGING 50-EXPERT SWARM"998 env = brain_env({999 "TITAN_FORCE_RETRAIN": "1" if force else os.getenv("TITAN_FORCE_RETRAIN", "0"),1000 "TITAN_NEURAL_FORGE_REASON": reason,1001 })1002 write_boot_log(1003 f"[{time.strftime('%H:%M:%S')}] ๐ง Neural swarm missing "1004 f"(experts={experts_found}/50 params={params_ready}); launching forced forge. "1005 f"models_dir={BRAIN_MODELS_DIR} prefix={BRAIN_BUCKET_PREFIX} reason={reason}"1006 )1007 try:1008 status_dir = os.path.dirname(os.path.abspath(TRAINER_STATUS_FILE))1009 if status_dir:1010 os.makedirs(status_dir, exist_ok=True)1011 with open(TRAINER_STATUS_FILE, "w", encoding="utf-8") as f:1012 json.dump({1013 "stage": "LAUNCHING_TRAINER",1014 "updated_at": time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime()),1015 "reason": reason,1016 "experts_found": int(experts_found),1017 "params_ready": bool(params_ready),1018 "models_dir": BRAIN_MODELS_DIR,1019 "brain_bucket_prefix": BRAIN_BUCKET_PREFIX,1020 }, f, indent=2, sort_keys=True)1021 except Exception:1022 pass1023 print(1024 f"๐ง TITAN-NN: Forced neural forge started. "1025 f"experts={experts_found}/50 params={params_ready} models_dir={BRAIN_MODELS_DIR}",1026 flush=True,1027 )1028 trainer_log_dir = os.path.dirname(os.path.abspath(TRAINER_LOG_FILE))1029 if trainer_log_dir:1030 os.makedirs(trainer_log_dir, exist_ok=True)1031 with open(TRAINER_LOG_FILE, "a", encoding="utf-8") as trainer_log:1032 self.trainer_proc = subprocess.Popen(1033 [sys.executable, "-u", "neural_trainer.py"],1034 stdout=trainer_log,1035 stderr=subprocess.STDOUT,1036 env=env,1037 )1038 return False1039 1040 def start_core(self):1041 if self.is_booting: return 1042 self.is_booting = True1043 1044 def _async_boot():1045 try:1046 zmq_manager.cache["boot_phase"] = "BOOTSTRAPPING"1047 write_boot_log(f"[{time.strftime('%H:%M:%S')}] โก TITAN CORE BOOTSTRAP STARTED โก", mode="w")1048 if getattr(self, 'core_proc', None):1049 try:1050 self.core_proc.kill()1051 write_boot_log(f"[{time.strftime('%H:%M:%S')}] Previous titan_core.py process stopped.")1052 except Exception as exc:1053 write_boot_log(f"[{time.strftime('%H:%M:%S')}] Previous core stop skipped: {exc}")1054 if getattr(self, 'alt_proc', None):1055 try:1056 self.alt_proc.kill()1057 write_boot_log(f"[{time.strftime('%H:%M:%S')}] Previous titan_alt_data.py process stopped.")1058 except Exception as exc:1059 write_boot_log(f"[{time.strftime('%H:%M:%S')}] Previous alt-data stop skipped: {exc}")1060 1061 # PHASE 1: C++ COMPILATION1062 if not any(f.startswith('titan_engine') and f.endswith('.so') for f in os.listdir('.')):1063 zmq_manager.cache["boot_phase"] = "COMPILING C++ KERNEL"1064 if _boot_bool("TITAN_SKIP_CPP_BOOT_BUILD", False):1065 write_boot_log(f"[{time.strftime('%H:%M:%S')}] C++ kernel build skipped by TITAN_SKIP_CPP_BOOT_BUILD=1; Python fallback remains active.")1066 else:1067 ok = _run_boot_command(1068 "C++ kernel build",1069 [sys.executable, "setup.py", "build_ext", "--inplace"],1070 timeout_sec=_boot_float("TITAN_CPP_BOOT_BUILD_TIMEOUT_SEC", 45.0),1071 )1072 if not ok:1073 write_boot_log(f"[{time.strftime('%H:%M:%S')}] C++ kernel unavailable; continuing with Python/provider-feed core.")1074 else:1075 write_boot_log(f"[{time.strftime('%H:%M:%S')}] C++ kernel already present; skipping build.")1076 1077 # PHASE 2: RESTORE/FORGE NEURAL SWARM BEFORE CORE BOOT1078 write_boot_log(f"[{time.strftime('%H:%M:%S')}] Checking neural brain artifacts in {BRAIN_MODELS_DIR}...")1079 if _boot_bool("TITAN_BOOT_SYNC_BRAIN_RESTORE", False):1080 write_boot_log(f"[{time.strftime('%H:%M:%S')}] Synchronous brain restore requested.")1081 restore_brain_artifacts(force=True)1082 else:1083 write_boot_log(f"[{time.strftime('%H:%M:%S')}] Synchronous brain restore skipped; restore/forge will continue in the background so provider feeds can boot now.")1084 threading.Thread(target=lambda: restore_brain_artifacts(force=True), name="titan_async_brain_restore", daemon=True).start()1085 experts_found, params_ready, swarm_ready = brain_swarm_status()1086 trainer_running = getattr(self, 'trainer_proc', None) is not None and self.trainer_proc.poll() is None1087 1088 if not swarm_ready or (self.neural_forge_required and trainer_running):1089 if not swarm_ready:1090 self.start_neural_forge("boot_missing_swarm", force=True)1091 msg = (1092 f"๐ก TITAN-NN: Neural swarm incomplete "1093 f"(experts={experts_found}/50 params={params_ready}). "1094 )1095 if STRICT_NEURAL_BOOT_LOCK_REQUESTED:1096 msg += "Legacy strict boot lock was requested but suppressed so provider feeds and risk supervision remain online. "1097 msg += "Cloud Core will launch in provider-feed mode while neural forge continues in the background."1098 zmq_manager.cache["boot_phase"] = "NEURAL FORGE + CORE BOOT"1099 write_boot_log(f"[{time.strftime('%H:%M:%S')}] {msg}")1100 print(msg, flush=True)1101 else:1102 ready_msg = f"โ
TITAN-NN: 50-Expert Swarm found in {BRAIN_MODELS_DIR}. Core ignition cleared."1103 write_boot_log(f"[{time.strftime('%H:%M:%S')}] {ready_msg}")1104 print(ready_msg, flush=True)1105 if getattr(self, 'trainer_proc', None) is not None and self.trainer_proc.poll() is not None:1106 self.trainer_proc = None1107 1108 # PHASE 3: ALT-DATA ORACLE IGNITION1109 zmq_manager.cache["boot_phase"] = "IGNITING ALT-DATA ORACLE"1110 write_boot_log(f"[{time.strftime('%H:%M:%S')}] ๐ Booting Institutional Alt-Data Oracle...")1111 print("๐ TITAN-ALT: Booting Institutional Alt-Data Oracle...")1112 try:1113 alt_log_dir = os.path.dirname(os.path.abspath(ALT_DATA_LOG_FILE))1114 if alt_log_dir:1115 os.makedirs(alt_log_dir, exist_ok=True)1116 with open(ALT_DATA_LOG_FILE, "a", encoding="utf-8") as alt_log:1117 self.alt_proc = subprocess.Popen(1118 [sys.executable, "-u", "titan_alt_data.py"],1119 stdout=alt_log,1120 stderr=subprocess.STDOUT,1121 env=brain_env(),1122 )1123 write_boot_log(f"[{time.strftime('%H:%M:%S')}] Alt-data oracle launched pid={getattr(self.alt_proc, 'pid', 'n/a')}.")1124 except Exception as exc:1125 write_boot_log(f"[{time.strftime('%H:%M:%S')}] Alt-data oracle launch failed but core will continue: {exc}")1126 1127 # PHASE 4: CLOUD CORE IGNITION1128 zmq_manager.cache["boot_phase"] = "IGNITING CLOUD CORE"1129 write_boot_log(f"[{time.strftime('%H:%M:%S')}] โก TITAN CORE IGNITION SEQUENCE INITIATED โก")1130 core_env = brain_env()1131 core_log_dir = os.path.dirname(os.path.abspath(CORE_LOG_FILE))1132 if core_log_dir:1133 os.makedirs(core_log_dir, exist_ok=True)1134 with open(CORE_LOG_FILE, "a", encoding="utf-8") as log_file:1135 self.core_proc = subprocess.Popen(1136 [sys.executable, "-u", "titan_core.py"],1137 stdout=log_file,1138 stderr=subprocess.STDOUT,1139 env=core_env,1140 )1141 write_boot_log(f"[{time.strftime('%H:%M:%S')}] titan_core.py launched pid={getattr(self.core_proc, 'pid', 'n/a')}; checking early health...")1142 time.sleep(_boot_float("TITAN_CORE_EARLY_HEALTH_WAIT_SEC", 2.0))1143 if getattr(self, "core_proc", None) is not None and self.core_proc.poll() is not None:1144 rc = self.core_proc.returncode1145 tail = _tail_for_boot_log(CORE_LOG_FILE, lines=18)1146 zmq_manager.cache["boot_phase"] = "CORE EXITED"1147 zmq_manager.cache["last_error"] = f"titan_core.py exited early rc={rc}"1148 write_boot_log(f"[{time.strftime('%H:%M:%S')}] โ titan_core.py exited early rc={rc}. Recent log:\n{tail}")1149 self.is_booting = False1150 return1151 1152 # PHASE 5: FULLY OPERATIONAL1153 zmq_manager.cache["boot_phase"] = "LIVE"1154 write_boot_log(f"[{time.strftime('%H:%M:%S')}] โ
TITAN CORE LIVE: provider feeds, ZMQ intake, and cloud Beast are running.")1155 except Exception as boot_exc:1156 zmq_manager.cache["boot_phase"] = "BOOT ERROR"1157 zmq_manager.cache["last_error"] = f"Boot Error: {boot_exc}"1158 write_boot_log(f"[{time.strftime('%H:%M:%S')}] โ BOOT ERROR: {boot_exc}")1159 print(f"โ TITAN BOOT ERROR: {boot_exc}", flush=True)1160 finally:1161 self.is_booting = False1162 1163 # Fire the boot sequence into an independent thread so Streamlit doesn't timeout1164 threading.Thread(target=_async_boot, daemon=True).start()1165 1166 def monitor_loop(self):1167 while True:1168 time.sleep(20)1169 if self.is_booting: continue1170 1171 restore_brain_artifacts()1172 experts_found, params_ready, swarm_ready = brain_swarm_status()1173 if not swarm_ready:1174 trainer_running = getattr(self, 'trainer_proc', None) is not None and self.trainer_proc.poll() is None1175 zmq_manager.cache["boot_phase"] = "CORE LIVE / NEURAL FORGE WARMING"1176 if not trainer_running:1177 self.start_neural_forge("watchdog_missing_swarm", force=False)1178 elif self.neural_forge_required:1179 trainer_running = getattr(self, 'trainer_proc', None) is not None and self.trainer_proc.poll() is None1180 if trainer_running:1181 zmq_manager.cache["boot_phase"] = "NEURAL FORGE FINALIZING"1182 continue1183 self.neural_forge_required = False1184 write_boot_log(1185 f"[{time.strftime('%H:%M:%S')}] โ
Neural forge complete "1186 f"(experts={experts_found}/50 params={params_ready}); restarting core to load full swarm."1187 )1188 self.start_core()1189 continue1190 1191 core_crashed = getattr(self, 'core_proc', None) is not None and self.core_proc.poll() is not None1192 alt_crashed = getattr(self, 'alt_proc', None) is not None and self.alt_proc.poll() is not None1193 1194 if core_crashed or alt_crashed:1195 print(f"โ ๏ธ Watchdog: Subsystem crashed (Core: {core_crashed}, Alt: {alt_crashed}). Initiating clean reboot...")1196 self.start_core()1197 1198 return CoreWatchdog()1199 1200zmq_manager = get_zmq_manager()