CoolFace
Apppublic

shaibu01/Titan-Engine

sourceHugging Faceupdated 2d agoView on Hugging Face
0likes
app.py4033 linesDownload Raw Back to root
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()

Showing the first 1,200 of 4033 lines. Download the file for the rest.