MYTHOSLIVE / live_runner.py
3VVM's picture
Upload MYTHOSLIVE project from A2 ZIP
9831ced verified
Raw History Blame Contribute Delete
5.78 kB
"""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()