File size: 5,688 Bytes
68025ee
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
"""
Daily live trading job — run once per day via cron.
Idempotent: skips a day already processed.
Usage: python live/daily_job.py
"""
import logging
import sys
import os

sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))

from datetime import datetime, date
from data.prices import fetch_ohlcv
from data.indicators import compute_indicators, get_latest_indicators
from data.news import fetch_news
from data.onchain import fetch_onchain_data
from backtest.portfolio import Portfolio
from agents.pipeline import build_pipeline
from db.store import init_db, _get_conn
from config import FREE_MODELS, ASSETS, BENCHMARKS

logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)

TODAY = date.today().isoformat()


def _already_processed(benchmark: str, model: str, asset: str, day: str) -> bool:
    conn = _get_conn()
    row = conn.execute(
        "SELECT id FROM decisions WHERE run_id IN (SELECT id FROM runs WHERE benchmark=? AND model=? AND asset=? AND status='live') AND date=?",
        (benchmark, model, asset, day),
    ).fetchone()
    conn.close()
    return row is not None


def _get_or_create_live_run(benchmark: str, model: str, asset: str) -> str:
    conn = _get_conn()
    row = conn.execute(
        "SELECT id FROM runs WHERE benchmark=? AND model=? AND asset=? AND status='live'",
        (benchmark, model, asset),
    ).fetchone()
    if row:
        run_id = row["id"]
    else:
        import uuid
        run_id = str(uuid.uuid4())
        from datetime import datetime
        conn.execute(
            "INSERT INTO runs (id, benchmark, model, asset, status, created_at) VALUES (?,?,?,?,?,?)",
            (run_id, benchmark, model, asset, "live", datetime.utcnow().isoformat()),
        )
        conn.commit()
    conn.close()
    return run_id


def run_daily():
    init_db()
    logger.info(f"Daily live job — {TODAY}")

    # Look back 60 days for indicators
    from datetime import datetime, timedelta
    lookback_start = (datetime.strptime(TODAY, "%Y-%m-%d") - timedelta(days=60)).strftime("%Y-%m-%d")

    for model in FREE_MODELS:
        for benchmark in BENCHMARKS:
            for asset in ASSETS:
                logger.info(f"Processing {benchmark}/{model}/{asset}")

                if _already_processed(benchmark, model, asset, TODAY):
                    logger.info(f"Already processed {benchmark}/{model}/{asset}/{TODAY}, skipping")
                    continue

                try:
                    df_raw = fetch_ohlcv(asset, lookback_start, TODAY)
                    df = compute_indicators(df_raw)
                    if df.empty:
                        logger.warning(f"No data for {asset}")
                        continue

                    row = df.iloc[-1]
                    price = float(row["close"])
                    indicators = get_latest_indicators(df)

                    run_id = _get_or_create_live_run(benchmark, model, asset)

                    # Load portfolio state from DB
                    conn = _get_conn()
                    last_snap = conn.execute(
                        "SELECT portfolio_value FROM decisions WHERE run_id=? ORDER BY date DESC LIMIT 1",
                        (run_id,),
                    ).fetchone()
                    conn.close()

                    from config import INITIAL_CAPITAL
                    portfolio_value = last_snap["portfolio_value"] if last_snap else INITIAL_CAPITAL
                    # Simplified snapshot (no persistent position tracking for now)
                    portfolio_snapshot = {
                        "cash": portfolio_value,
                        "position": 0.0,
                        "total_value": portfolio_value,
                        "drawdown": 0.0,
                    }

                    from data.prices import ohlcv_to_records
                    market_data = {
                        "asset": asset,
                        "current_price": price,
                        "date": TODAY,
                        "recent_ohlcv": ohlcv_to_records(df)[-30:],
                        "indicators": indicators,
                        "portfolio": portfolio_snapshot,
                    }

                    if benchmark in ("B", "C"):
                        market_data["news"] = fetch_news(asset, limit=5)
                    if benchmark == "C":
                        market_data["onchain"] = fetch_onchain_data(asset)

                    pipeline = build_pipeline(benchmark, model)
                    result = pipeline.decide(market_data)
                    decision = result["decision"]
                    agent_outputs = result.get("agent_outputs", {})

                    import json
                    conn = _get_conn()
                    conn.execute(
                        "INSERT INTO decisions (run_id, date, price, action, size, confidence, reason, agent_outputs, portfolio_value) VALUES (?,?,?,?,?,?,?,?,?)",
                        (
                            run_id, TODAY, price,
                            decision.get("action"), decision.get("size"),
                            decision.get("confidence"), decision.get("reason"),
                            json.dumps(agent_outputs), portfolio_value,
                        ),
                    )
                    conn.commit()
                    conn.close()
                    logger.info(f"Decision saved: {benchmark}/{model}/{asset}/{TODAY} -> {decision.get('action')}")

                except Exception as e:
                    logger.error(f"Error for {benchmark}/{model}/{asset}: {e}", exc_info=True)


if __name__ == "__main__":
    run_daily()