import asyncio import logging import colorlog import re import json import random import time import sys from pathlib import Path from sqlalchemy import select, func, delete from sqlalchemy.exc import IntegrityError # Add project root to path sys.path.append(str(Path(__file__).resolve().parent)) import config from database import init_db, AsyncSessionLocal, Alpha from core.wq_client import AsyncBrainClient from core.generator import AlphaGenerator from core.evaluator import AlphaEvaluator, serialize_raw_payload from core.submitter import AlphaSubmitter # Configure main system logging via colorlog logger = logging.getLogger("Orchestrator") # Compatibility variables for app.py dashboard MINER_RUNNING = False MINER_TASK = None ACTIVE_LLM_MODEL = "Initializing..." # Updated by generator during runtime # Shared event: set by /login-cookie HTTP endpoint to wake up the Telegram recovery loop COOKIE_UPDATED_EVENT: asyncio.Event = asyncio.Event() VALID_OPERATORS = { # Mathematical "abs", "log", "sign", "sqrt", "exp", "power", "min", "max", # Cross-sectional "rank", "zscore", "group_rank", "group_neutralize", "group_zscore", # Time-series "ts_delta", "ts_mean", "ts_std_dev", "ts_decay_linear", "signed_power", "ts_corr", "ts_covariance", "ts_max", "ts_min", "ts_rank", "ts_sum", "ts_product", "ts_argmax", "ts_argmin", "ts_scale", "ts_av_diff", "ts_zscore", # Vector "vec_avg", "vec_sum", "vec_choose", # Constants/Groups "sector", "subindustry", "industry", "market", "cap", # Price-volume/basic fields "close", "open", "high", "low", "volume", "vwap", "returns", "adv20", "adv60", "sharesout", "adjfactor", # Other common operators "trade_when", "if_else", "nan_to_zero", "purify", "floor", "ceil", "round", "trunc", "sector_rank", "industry_rank", "subindustry_rank" } # Fast Expression arity contract used before a platform simulation is spent. # A range is used only where the platform accepts optional parameters. OPERATOR_ARITY = { "abs": (1, 1), "log": (1, 1), "sign": (1, 1), "sqrt": (1, 1), "exp": (1, 1), "power": (2, 2), "signed_power": (2, 2), "min": (2, 2), "max": (2, 2), "rank": (1, 1), "zscore": (1, 1), "group_rank": (1, 1), "group_zscore": (1, 1), "group_neutralize": (2, 2), "ts_delta": (2, 2), "ts_mean": (2, 2), "ts_std_dev": (2, 2), "ts_decay_linear": (2, 2), "ts_corr": (3, 3), "ts_covariance": (3, 3), "ts_max": (2, 2), "ts_min": (2, 2), "ts_rank": (2, 2), "ts_sum": (2, 2), "ts_product": (2, 2), "ts_argmax": (2, 2), "ts_argmin": (2, 2), "ts_scale": (2, 2), "ts_av_diff": (1, 1), "ts_zscore": (2, 2), "vec_avg": (1, 1), "vec_sum": (1, 1), "vec_choose": (2, 2), "if_else": (3, 3), "trade_when": (3, 3), "nan_to_zero": (1, 1), "purify": (1, 1), "floor": (1, 1), "ceil": (1, 1), "round": (1, 1), "trunc": (1, 1), "sector_rank": (1, 1), "industry_rank": (1, 1), "subindustry_rank": (1, 1), } def _split_top_level_args(inner: str) -> list[str]: args, start, depth = [], 0, 0 for index, char in enumerate(inner): if char == "(": depth += 1 elif char == ")": depth -= 1 elif char == "," and depth == 0: args.append(inner[start:index].strip()) start = index + 1 tail = inner[start:].strip() if tail or inner.strip(): args.append(tail) return args def validate_operator_arity(expr: str) -> tuple[bool, str]: """Validate function argument counts without attempting full evaluation.""" for match in re.finditer(r"\b([A-Za-z_][A-Za-z0-9_]*)\s*\(", expr): name = match.group(1).lower() if name not in OPERATOR_ARITY: continue open_index = expr.find("(", match.start(), match.end()) depth = 1 close_index = open_index + 1 while close_index < len(expr) and depth: if expr[close_index] == "(": depth += 1 elif expr[close_index] == ")": depth -= 1 close_index += 1 if depth: return False, f"Unbalanced call for operator {name}" args = _split_top_level_args(expr[open_index + 1:close_index - 1]) minimum, maximum = OPERATOR_ARITY[name] if not minimum <= len(args) <= maximum: return False, f"Operator {name} expects {minimum} input(s), received {len(args)}" if name == "group_neutralize" and args[1].lower() not in {"sector", "subindustry", "industry", "market"}: return False, "group_neutralize requires a literal group field" return True, "" def is_retryable_simulation_failure(result: dict) -> bool: """Return whether a failed simulation is an infrastructure, not research, failure.""" if result.get("status") != "FAILED": return False message = str(result.get("error_message", "")).lower() retryable_markers = ( "401", "authentication", "session", "timeout", "timed out", "network", "connection", "rate limit", "status 429", "status 5", "polling error", ) return any(marker in message for marker in retryable_markers) def sanitize_expression(expr: str) -> str: """ Sanitizes WorldQuant expressions to ensure they are compatible with the platform. Converts scientific notation like '1e-9' or '1e-5' to full decimal notation. """ if not expr: return expr def replace_sci(match): sci_str = match.group(0) try: val = float(sci_str) dec_str = f"{val:.15f}".rstrip('0') if dec_str.endswith('.'): dec_str += '0' return dec_str except ValueError: return sci_str pattern = r'\b\d+(?:\.\d+)?[eE][+-]?\d+\b' return re.sub(pattern, replace_sci, expr) def validate_expression(expr: str) -> bool: """ Validates basic syntax rules and variable names of a generated alpha formula before submitting for simulation. Enforces a strict allowlist. """ if not expr or len(expr.strip()) < 5: return False # Check balanced parentheses and brackets if expr.count('(') != expr.count(')'): return False if expr.count('[') != expr.count(']'): return False # Check for empty wrappers if "()" in expr or "[]" in expr: return False arity_ok, arity_error = validate_operator_arity(expr) if not arity_ok: logger.warning("Validation failed: %s in formula '%s'", arity_error, expr) return False # Check for redundant nested wrappers if "rank(rank(" in expr or "zscore(zscore(" in expr: return False # Block double-nested group_neutralize — reliably causes simulation polling timeout if expr.count("group_neutralize(") >= 2: return False valid_fields = { "close", "open", "high", "low", "volume", "vwap", "returns", "cap", "adv20", "adv60", "sharesout", "adjfactor" } # Try to load custom caching if available cache_file = config.DATA_DIR / "data_fields_cache.json" if cache_file.exists(): try: with open(cache_file, "r") as f: cached_data = json.load(f) if isinstance(cached_data, dict) and cached_data: for cat_fields in cached_data.values(): for f in cat_fields: fid = f.get("id") if fid: valid_fields.add(fid) except Exception as e: logger.warning(f"Error loading cache for variable validation: {e}") # Extract all word tokens using regex tokens = re.findall(r'\b[a-zA-Z_][a-zA-Z0-9_]*\b', expr) for t in tokens: t_lower = t.lower() if t_lower in VALID_OPERATORS or t in VALID_OPERATORS: continue if t_lower in ("sector", "subindustry", "industry", "market", "sector_rank", "group"): continue # Allow variables with valid dataset prefixes (e.g. mdl17_target_diff, ern4_eps_surprise, sir2_short_interest) if any(t_lower.startswith(p) for p in ["mdl17_", "mdl77_", "ern4_", "snt1_", "sir2_"]): continue if t not in valid_fields: logger.warning(f"Validation failed: Unknown variable '{t}' in formula '{expr}'") return False return True async def run_simulation_and_process( alpha_id: int, wq_client: AsyncBrainClient, evaluator: AlphaEvaluator, submitter: AlphaSubmitter, notifier = None ) -> None: """ Simulates a single Alpha record, passes it to the evaluator, and immediately submits it if it passes evaluation. """ # 1. Fetch Alpha record from database async with AsyncSessionLocal() as session: alpha = await session.get(Alpha, alpha_id) if not alpha: logger.warning(f"Alpha {alpha_id} not found in database.") return # Guard check if already simulating or processed if alpha.status != "PENDING_GEN": return # Verify expression syntax first — validate the sanitized version but # do NOT assign back to alpha.expression to avoid dirty-tracking the # expression column (which has a UNIQUE constraint and would crash on commit). sanitized = sanitize_expression(alpha.expression) if not validate_expression(sanitized): try: alpha.status = "TRASHED" alpha.error_log = "Expression syntax/allowlist validation failed." # Explicitly expire expression so SQLAlchemy does NOT include it in UPDATE session.expire(alpha, ["expression"]) await session.commit() except IntegrityError: await session.rollback() # Fall back to a targeted UPDATE that only touches status + error_log from sqlalchemy import text await session.execute( text("UPDATE alphas SET status='TRASHED', error_log='Validation failed (duplicate)' WHERE id=:id"), {"id": alpha_id} ) await session.commit() logger.warning(f"Alpha {alpha_id} failed expression validation. Trashed.") return # Transition status to SIMULATING to prevent duplicate pick-ups alpha.status = "SIMULATING" try: await session.commit() except IntegrityError: # Duplicate expression already exists in DB — trash this alpha and move on await session.rollback() alpha.status = "TRASHED" alpha.error_log = "Duplicate expression: identical formula already exists in database." await session.commit() logger.warning(f"Alpha {alpha_id} TRASHED: duplicate expression collision. Skipping.") return # Keep variable values for client payload # NOTE: use `sanitized` (not alpha.expression) so scientific notation # like 1e-6 is already converted to 0.000001 before reaching BRAIN. alpha_dict = { "expression": sanitized, "language": alpha.language, "region": alpha.region, "universe": alpha.universe, "delay": alpha.delay, "decay": alpha.decay, "truncation": alpha.truncation, "neutralization": alpha.neutralization, "pasteurization": alpha.pasteurization, "nan_handling": alpha.nan_handling, "unit_handling": alpha.unit_handling, "lookback": alpha.lookback, "test_period": alpha.test_period, "max_trade": alpha.max_trade, "max_position": alpha.max_position, } delay = alpha.delay # 2. Dispatch to Simulation Engine (Semaphored) try: logger.info(f"Dispatching simulation for Alpha {alpha_id}...") sim_result = await wq_client.simulate_alpha(alpha_dict) except Exception as e: logger.exception(f"Unexpected error during simulation of Alpha {alpha_id}: {e}") async with AsyncSessionLocal() as session: alpha = await session.get(Alpha, alpha_id) if alpha: alpha.status = "PENDING_GEN" # Reset status so we can retry alpha.retry_count += 1 alpha.error_log = f"Simulation crash: {str(e)}" await session.commit() return # An expired session, rate limit, timeout, or network problem says nothing # about the expression. Preserve the candidate for retry instead of # feeding it to the evaluator, which would permanently trash it. if is_retryable_simulation_failure(sim_result): failure = str(sim_result.get("error_message", "retryable simulation failure"))[:500] async with AsyncSessionLocal() as session: alpha = await session.get(Alpha, alpha_id) if alpha: alpha.retry_count += 1 alpha.error_log = f"Retryable simulation failure: {failure}" if "timed out" in failure.lower() or "timeout" in failure.lower(): alpha.status = "RESULT_TIMEOUT_REVIEW" logger.error("Alpha %s simulation timed out; marked for review.", alpha_id) else: retry_limit = 5 if alpha.retry_count >= retry_limit: alpha.status = "RETRY_EXHAUSTED" logger.error("Alpha %s exhausted retry budget; retaining it for inspection.", alpha_id) else: alpha.status = "PENDING_GEN" await session.commit() if wq_client.auth_dead: logger.warning("Authentication recovery required; candidate %s returned to queue.", alpha_id) else: logger.warning("Retryable simulation failure for Alpha %s; returned to queue.", alpha_id) return if sim_result.get("status") == "INCOMPLETE": async with AsyncSessionLocal() as session: alpha = await session.get(Alpha, alpha_id) if alpha: alpha.status = "RESULT_INCOMPLETE" alpha.error_log = str(sim_result.get("error_message", "Incomplete simulation result"))[:500] alpha.raw_result_json = serialize_raw_payload(sim_result.get("raw_data")) await session.commit() logger.warning("Alpha %s completed without a recognized metrics payload; retained for parser review.", alpha_id) return # 3. Evaluate Results try: new_status = await evaluator.evaluate_and_update(alpha_id, sim_result, delay) except Exception as e: logger.exception(f"Unexpected error during evaluation of Alpha {alpha_id}: {e}") return # 4. Immediate Submission if passed if new_status == "PENDING_SUBMIT": platform_id = sim_result.get("alpha_id") if platform_id: try: success = await submitter.submit_winning_alpha(alpha_id, platform_id, wq_client) if success and notifier: sharpe = sim_result.get("sharpe") or 0.0 fitness = sim_result.get("fitness") or 0.0 region = alpha_dict.get("region", "USA") asyncio.create_task( notifier.send_submission_alert( alpha_id=str(alpha_id), sharpe=sharpe, fitness=fitness, region=region ) ) except Exception as e: logger.exception(f"Unexpected error during auto-submission of Alpha {alpha_id}: {e}") else: logger.error(f"Alpha {alpha_id} passed evaluation but has no platform alpha_id returned. Cannot submit.") async def recover_incomplete_results( wq_client: AsyncBrainClient, evaluator: AlphaEvaluator, submitter: AlphaSubmitter, notifier=None, limit: int = 10, ) -> None: """Resolve Alpha metrics for completed simulations saved before the parser fix.""" async with AsyncSessionLocal() as session: rows = await session.execute( select(Alpha.id, Alpha.delay, Alpha.raw_result_json) .where(Alpha.status == "RESULT_INCOMPLETE") .order_by(Alpha.id) .limit(limit) ) candidates = rows.all() for alpha_id, delay, raw_result_json in candidates: try: raw_data = json.loads(raw_result_json or "{}") simulation_data = raw_data.get("simulation", raw_data) if not isinstance(simulation_data, dict): continue result = await wq_client._resolve_completed_simulation(simulation_data) if result.get("status") == "INCOMPLETE": continue new_status = await evaluator.evaluate_and_update(alpha_id, result, delay) if new_status == "PENDING_SUBMIT" and result.get("alpha_id"): submitted = await submitter.submit_winning_alpha(alpha_id, result["alpha_id"], wq_client) if submitted and notifier: asyncio.create_task( notifier.send_submission_alert( alpha_id=str(alpha_id), sharpe=result.get("sharpe") or 0.0, fitness=result.get("fitness") or 0.0, ) ) except Exception as exc: logger.warning("Could not recover incomplete Alpha %s: %s", alpha_id, exc) async def log_dashboard_metrics() -> None: """Queries current SQLite DB metrics and prints a clean, formatted terminal dashboard.""" async with AsyncSessionLocal() as session: # Get count for each status stmt = select(Alpha.status, func.count(Alpha.id)).group_by(Alpha.status) results = await session.execute(stmt) counts = {status: count for status, count in results.all()} pending_gen = counts.get("PENDING_GEN", 0) simulating = counts.get("SIMULATING", 0) pending_submit = counts.get("PENDING_SUBMIT", 0) submitted = counts.get("SUBMITTED", 0) trashed = counts.get("TRASHED", 0) failed_submit = counts.get("SUBMISSION_FAILED", 0) print("\n" + "=" * 70) print(" QUANTFORGE v2.0 - MINER STATUS DASHBOARD") print("=" * 70) print(f" [+] PENDING GENERATION : {pending_gen}") print(f" [+] SIMULATING : {simulating}") print(f" [+] PENDING SUBMISSION : {pending_submit}") print(f" [+] SUBMITTED ALPHAS : {submitted}") print(f" [+] SUBMISSION FAILURES : {failed_submit}") print(f" [+] TRASHED ALPHAS : {trashed}") print("=" * 70 + "\n") async def main_loop() -> None: """Primary infinite loop carrying out Phase 3 lifecycle and handling retries.""" logger.info("Initializing QuantForge v2.0 Database and Core Clients...") await init_db() # Instantiate modules generator = AlphaGenerator() evaluator = AlphaEvaluator() submitter = AlphaSubmitter() # Setup Telegram bot alert system notifier = None if config.TELEGRAM_BOT_TOKEN and config.TELEGRAM_CHAT_ID: try: from core.telegram_bot import TelegramNotifier notifier = TelegramNotifier( config.TELEGRAM_BOT_TOKEN, config.TELEGRAM_CHAT_ID, getattr(config, "TELEGRAM_API_PROXY", None) ) asyncio.create_task(notifier.send_startup()) logger.info("Telegram notifier system successfully initialized.") except Exception as te: logger.error(f"Failed to initialize Telegram notifier: {te}") async def handle_auth_recovery(wq_client) -> bool: if not notifier: logger.error("Authentication expired, but Telegram notifier is not configured.") return False logger.warning("Auth is dead/expired. Initializing Telegram authentication recovery loop...") await notifier.send_auth_expired_alert() # Wait up to 30 minutes (1800 seconds) for cookie reply new_cookie = await notifier.poll_for_cookie(timeout_seconds=1800, wake_event=COOKIE_UPDATED_EVENT) if new_cookie: logger.info("New cookie received from Telegram. Updating config and environment...") # Update memory and env config.WQ_SESSION_COOKIE = new_cookie import os os.environ["WQ_SESSION_COOKIE"] = new_cookie # Persist cookie to miner_settings.json and sync immediately to Hugging Face try: from config import save_wq_credentials save_wq_credentials(config.WQ_USERNAME, config.WQ_PASSWORD, new_cookie) logger.info("New cookie saved and synced to Hugging Face backup.") except Exception as ex: logger.error(f"Failed to save and sync new cookie: {ex}") # Re-verify and update current session headers session = await wq_client.get_session() if "session=" in new_cookie or "t=" in new_cookie or ";" in new_cookie: cookie_header = new_cookie if cookie_header.lower().startswith("cookie:"): cookie_header = cookie_header[7:].strip() else: cookie_header = f"t={new_cookie}; session={new_cookie}" session.headers["Cookie"] = cookie_header # Attempt verification/authentication authenticated = await wq_client.authenticate() if authenticated: wq_client.auth_dead = False await notifier.send_session_restored() return True return False logger.info("Connecting and Authenticating with WorldQuant BRAIN platform...") async with AsyncBrainClient() as wq_client: # Start keep-alive loop to prevent cookie expiration keep_alive_task = asyncio.create_task(wq_client.keep_alive()) authenticated = await wq_client.authenticate() if not authenticated: logger.warning("Initial authentication failed. Attempting Telegram recovery...") if notifier: recovered = await handle_auth_recovery(wq_client) if not recovered: logger.error("Initial Telegram recovery failed. Exiting to retry...") await asyncio.sleep(60) return else: logger.error("Platform authentication failed and Telegram notifier is not configured. Re-trying in 60s...") await asyncio.sleep(60) return logger.info("Successfully connected to BRAIN. Starting infinite execution loop...") # The platform record, not a 2xx response from POST /submit, is the # source of truth. Correct stale local flags before the dashboard or # Telegram can claim a submission that BRAIN has not accepted. await submitter.reconcile_local_submission_records(wq_client) # On every boot: rescue alphas stuck in SIMULATING from previous crashed sessions. async with AsyncSessionLocal() as session: stuck_stmt = select(func.count(Alpha.id)).where(Alpha.status == "SIMULATING") res = await session.execute(stuck_stmt) stuck_count = res.scalar() or 0 if stuck_count > 0: from sqlalchemy import update await session.execute( update(Alpha) .where(Alpha.status == "SIMULATING") .values(status="PENDING_GEN", error_log="Rescued: stuck SIMULATING from previous session") ) await session.commit() logger.info(f"Startup rescue: reset {stuck_count} stuck SIMULATING alphas back to PENDING_GEN.") last_incomplete_recovery = 0.0 generation_backoff_until = 0.0 while True: try: # Print Status Dashboard await log_dashboard_metrics() # If auth is flagged as dead mid-run, trigger recovery if wq_client.auth_dead: recovered = await handle_auth_recovery(wq_client) if not recovered: logger.error("Telegram auth recovery failed or timed out. Sleeping 60s...") await asyncio.sleep(60) continue # Revisit completed simulation envelopes at a bounded rate. # This recovers the backlog created before metric resolution # followed the Alpha ID returned by the simulations endpoint. now = time.monotonic() if now - last_incomplete_recovery >= 300: await recover_incomplete_results(wq_client, evaluator, submitter, notifier) last_incomplete_recovery = now # Step A: Generate new Alphas if queue is dry (< 50 candidates) async with AsyncSessionLocal() as session: pending_cnt_stmt = select(func.count(Alpha.id)).where(Alpha.status == "PENDING_GEN") res = await session.execute(pending_cnt_stmt) pending_cnt = res.scalar() or 0 if pending_cnt < 50 and time.monotonic() >= generation_backoff_until: needed = 50 - pending_cnt batch_size = min(needed, 5) logger.info(f"Queue count ({pending_cnt}) below threshold (50). Generating batch of {batch_size} candidates...") generated = await generator.generate_alpha_candidates(batch_size=batch_size) if not generated: generation_backoff_until = time.monotonic() + config.GENERATION_EMPTY_BACKOFF_SECONDS logger.warning( "No novel research candidates generated; pausing generation for %ss.", config.GENERATION_EMPTY_BACKOFF_SECONDS, ) # Do not dispatch newly generated work against a dead BRAIN # session. It creates avoidable timeout/retry noise and can # incorrectly bias family scores with infrastructure errors. if wq_client.auth_dead: logger.warning("Authentication expired during generation; pausing simulation dispatch until recovery.") await asyncio.sleep(5) continue # Step B: Get simulating alphas candidates async with AsyncSessionLocal() as session: alphas_stmt = select(Alpha.id).where(Alpha.status == "PENDING_GEN").limit(20) res = await session.execute(alphas_stmt) alphas_to_process = res.scalars().all() if alphas_to_process: logger.info(f"Dispatched simulation tasks for {len(alphas_to_process)} alphas...") tasks = [ run_simulation_and_process(alpha_id, wq_client, evaluator, submitter, notifier) for alpha_id in alphas_to_process ] await asyncio.gather(*tasks) # Check if submission quota freeze was set if submitter.quota_frozen: logger.warning("Auto-submission quota is currently FROZEN. WINNING alphas are backlogged.") # Politeness sleep between loop cycles await asyncio.sleep(10) except asyncio.CancelledError: logger.info("Orchestrator received cancel signal. Shutting down gracefully...") keep_alive_task.cancel() break except Exception as e: logger.error( f"CRITICAL ERROR encountered in main loop: {e}. " "Sleeping 60 seconds and retrying to ensure 24/7 continuous operation...", exc_info=True ) await asyncio.sleep(60) async def miner_loop() -> None: """Wrapper that manages MINER_RUNNING status state for compatibility.""" global MINER_RUNNING MINER_RUNNING = True try: await main_loop() finally: MINER_RUNNING = False def main() -> None: try: asyncio.run(miner_loop()) except KeyboardInterrupt: print("\n[!] Shutdown requested. Exiting QuantForge v2.0...") if __name__ == "__main__": main()