"""Poll real closed candles with the same selectable dashboard pipeline.""" import argparse import csv from dataclasses import asdict import json from pathlib import Path import time import pandas as pd from config import MODEL_NAMES, OUTPUT_DIR, PROVIDERS, TIMEFRAMES from pipeline import RunConfig, live_result, prepare_history from utils.helpers import price_direction def main(): parser = argparse.ArgumentParser(description='MYTHOS selectable closed-candle forecast logger.') parser.add_argument('--provider', choices=PROVIDERS, default='Yahoo Finance') parser.add_argument('--market', choices=['Forex', 'Crypto'], default='Crypto') parser.add_argument('--symbol', default='BTC-USD') parser.add_argument('--timeframe', choices=TIMEFRAMES, default='5m') parser.add_argument('--models', nargs='+', choices=MODEL_NAMES, default=['ARIMA-GARCH']) parser.add_argument('--lookback-days', type=int, default=30) parser.add_argument('--window', type=int, default=200) parser.add_argument('--horizon', type=int, default=1) parser.add_argument('--features', action='store_true') parser.add_argument('--poll-seconds', type=int, default=60) args = parser.parse_args() config = RunConfig(provider=args.provider, market=args.market, symbol=args.symbol, timeframe=args.timeframe, models=tuple(args.models), mode='Parallel models' if len(args.models) > 1 else 'Single model', lookback_days=args.lookback_days, window=args.window, horizon=args.horizon, use_features=args.features) run_id = pd.Timestamp.now(tz='UTC').strftime('%Y%m%dT%H%M%SZ') Path(OUTPUT_DIR).mkdir(parents=True, exist_ok=True) path = Path(OUTPUT_DIR) / f'live_logger_{run_id}.csv' path.with_suffix('.json').write_text(json.dumps(asdict(config), indent=2), encoding='utf-8') fields = ['model', 'predicted_at_utc', 'origin_open_utc', 'target_open_utc', 'current_close', 'predicted_close', 'predicted_direction', 'entry', 'stop_loss', 'take_profit', 'risk_reward', 'actual_close', 'actual_direction', 'correct', 'status'] def log(row): new_file = not path.exists() with path.open('a', newline='', encoding='utf-8') as file: writer = csv.DictWriter(file, fieldnames=fields) if new_file: writer.writeheader() writer.writerow(row) pending = [] last_origin = None print(f'MYTHOS {config.provider} · {config.symbol} · {config.timeframe} · {config.models}. Logging: {path}') try: while True: try: frame, features = prepare_history(config) still_pending = [] for row in pending: target = pd.Timestamp(row['target_open_utc']) if target in frame.index: actual = float(frame.loc[target, 'Close']) direction = price_direction(actual, row['current_close']) interval = frame.loc[pd.Timestamp(row['origin_open_utc']):target].index.to_series().diff().dropna() continuous = len(interval) == config.horizon and (interval == config.candle_step).all() verified = {**row, 'actual_close': actual, 'actual_direction': direction, 'correct': direction == row['predicted_direction'] if continuous else '', 'status': 'VERIFIED' if continuous else 'GAP — not scored'} log(verified) print(verified) elif frame.index[-1] > target: log({**row, 'status': 'MISSING TARGET — not scored'}) else: still_pending.append(row) pending = still_pending if last_origin != frame.index[-1]: result = live_result(config, frame, features) last_origin = frame.index[-1] for forecast in result['forecasts']: if not forecast['forecast_path'] or result['metadata']['stale']: print(f"{forecast['Model']}: {forecast['Status']} {forecast['Error']}") continue pending.append({ 'model': forecast['Model'], 'predicted_at_utc': result['metadata']['fetched_at_utc'], 'origin_open_utc': last_origin.isoformat(), 'target_open_utc': (last_origin + config.candle_step * config.horizon).isoformat(), 'current_close': forecast['Current close'], 'predicted_close': forecast['Predicted close'], 'predicted_direction': forecast['Direction'].lower(), 'entry': result['signal'].get('trade_levels', {}).get('entry'), 'stop_loss': result['signal'].get('trade_levels', {}).get('stop_loss'), 'take_profit': result['signal'].get('trade_levels', {}).get('take_profit'), 'risk_reward': result['signal'].get('trade_levels', {}).get('risk_reward'), }) print(f"{forecast['Model']}: {forecast['Direction']} → {forecast['Predicted close']:.8f}") except Exception as error: print(f'Pipeline error: {error}') time.sleep(max(1, args.poll_seconds)) except KeyboardInterrupt: for row in pending: log({**row, 'status': 'STOPPED BEFORE TARGET — not scored'}) print('Stopped.') if __name__ == '__main__': main()