from __future__ import annotations import datetime as dt import json import time import uuid from dataclasses import dataclass from typing import Dict, List, Optional, Tuple import yaml from kalshi_client import KalshiClient, KalshiConfig from storage import Storage, DecisionRow def _parse_iso_z(s: str) -> dt.datetime: # example: "2026-02-10T15:29:16Z" return dt.datetime.fromisoformat(s.replace("Z", "+00:00")) def _market_prob_from_asks(m: dict) -> Tuple[Optional[float], Optional[float]]: """ Return (yes_ask_prob, no_ask_prob) in 0..1 using whatever fields are present. Kalshi market objects often include yes_ask/yes_bid (cents) and/or yes_ask_dollars (float). """ def read_prob(prefix: str) -> Optional[float]: cents_key = f"{prefix}_ask" dollars_key = f"{prefix}_ask_dollars" if m.get(dollars_key) is not None: return float(m[dollars_key]) if m.get(cents_key) is not None: return float(m[cents_key]) / 100.0 return None yes_p = read_prob("yes") no_p = read_prob("no") return yes_p, no_p def _yes_bid_ask_spread(m: dict) -> Tuple[float, float, float]: def read(prefix: str, side: str) -> Optional[float]: cents_key = f"{prefix}_{side}" dollars_key = f"{prefix}_{side}_dollars" if m.get(dollars_key) is not None: return float(m[dollars_key]) if m.get(cents_key) is not None: return float(m[cents_key]) / 100.0 return None yes_bid = read("yes", "bid") or 0.0 yes_ask = read("yes", "ask") or 1.0 return yes_bid, yes_ask, max(0.0, yes_ask - yes_bid) def _stake_for_prob(prob: float, tiers: List[dict]) -> Optional[float]: for t in tiers: if float(t["min"]) <= prob < float(t["max"]): return float(t["stake_dollars"]) # allow exact upper bound match (e.g., prob==0.95) for t in tiers: if abs(prob - float(t["max"])) < 1e-12: return float(t["stake_dollars"]) return None def _dollars_to_cents_price(p: float) -> int: c = int(round(p * 100)) return max(1, min(99, c)) def discover_open_crypto_15m_markets(client: KalshiClient, symbols: list[str]) -> dict[str, dict]: """ Scan open markets and pick the nearest-to-close market for each symbol. Uses title heuristics so you don't need series discovery (more robust). """ data = client.list_open_markets(limit=1000) markets = data.get("markets") or [] now = dt.datetime.now(dt.timezone.utc) per_symbol: dict[str, dict] = {} for m in markets: title = (m.get("title") or "").upper() # Heuristics: 15-min, Up/Down, and the symbol present if "UP OR DOWN" not in title: continue if "15" not in title: continue close_time_s = m.get("close_time") if not close_time_s: continue close_time = _parse_iso_z(close_time_s) if close_time <= now: continue for sym in symbols: if sym.upper() not in title: continue prev = per_symbol.get(sym) if prev is None or _parse_iso_z(prev["close_time"]) > close_time: per_symbol[sym] = m return per_symbol def main() -> None: cfg = yaml.safe_load(open("config.yaml", "r")) mode = cfg["mode"].lower() symbols = cfg["symbols"] lead_seconds = int(cfg["lead_seconds"]) trade_window_seconds = int(cfg["trade_window_seconds"]) min_prob = float(cfg["min_prob"]) max_prob = float(cfg["max_prob"]) tiers = list(cfg["tiers"]) max_spread = float(cfg["max_spread_dollars"]) improve_cents = int(cfg.get("improve_cents", 0)) tif = str(cfg.get("time_in_force", "fill_or_kill")) poll_seconds = float(cfg.get("poll_seconds", 2)) storage = Storage(cfg.get("sqlite_path", "storage.sqlite")) client = KalshiClient(KalshiConfig(env=str((__import__("os").getenv("KALSHI_ENV") or "demo")))) print(f"[init] mode={mode} symbols={symbols} lead_seconds={lead_seconds} window={trade_window_seconds}s") last_traded_market: Dict[str, str] = {} # sym -> market_ticker while True: now = dt.datetime.now(dt.timezone.utc) print(f"[heartbeat] {now.isoformat().replace('+00:00','Z')}") try: markets_by_symbol = discover_open_crypto_15m_markets(client, symbols) except Exception as e: print(f"[kalshi] market discovery failed: {e}") time.sleep(poll_seconds) continue for sym, m in markets_by_symbol.items(): try: market_ticker = m["ticker"] close_time = _parse_iso_z(m["close_time"]) # Fire only in a narrow window at T - lead_seconds start = close_time - dt.timedelta(seconds=lead_seconds) end = start + dt.timedelta(seconds=trade_window_seconds) if not (start <= now <= end): continue # Dedup per symbol per market if last_traded_market.get(sym) == market_ticker: continue # Pull full market object (more reliable fields) market_full = client.get_market(market_ticker).get("market") or m yes_p, no_p = _market_prob_from_asks(market_full) if yes_p is None or no_p is None: storage.log_decision( DecisionRow( ts=time.time(), symbol=sym, market_ticker=market_ticker, close_time=market_full.get("close_time", ""), side="", prob=0.0, yes_bid=0.0, yes_ask=0.0, spread=0.0, stake_dollars=0.0, count=0, limit_cents=0, reason="skip: missing yes/no ask fields", ) ) continue yes_bid, yes_ask, spread = _yes_bid_ask_spread(market_full) if spread > max_spread: storage.log_decision( DecisionRow( ts=time.time(), symbol=sym, market_ticker=market_ticker, close_time=market_full.get("close_time",""), side="", prob=0.0, yes_bid=yes_bid, yes_ask=yes_ask, spread=spread, stake_dollars=0.0, count=0, limit_cents=0, reason=f"skip: spread {spread:.4f} > {max_spread:.4f}", ) ) continue # Choose side(s) that qualify: any of the 6 outcomes (YES/NO for each symbol) candidates: List[Tuple[str, float]] = [] if min_prob <= yes_p <= max_prob: candidates.append(("yes", yes_p)) if min_prob <= no_p <= max_prob: candidates.append(("no", no_p)) if not candidates: storage.log_decision( DecisionRow( ts=time.time(), symbol=sym, market_ticker=market_ticker, close_time=market_full.get("close_time",""), side="", prob=max(yes_p, no_p), yes_bid=yes_bid, yes_ask=yes_ask, spread=spread, stake_dollars=0.0, count=0, limit_cents=0, reason=f"skip: prob not in [{min_prob:.2f},{max_prob:.2f}] (yes={yes_p:.3f} no={no_p:.3f})", ) ) continue # If both somehow qualify (rare/unexpected), prefer the higher probability side side, prob = sorted(candidates, key=lambda x: x[1], reverse=True)[0] stake = _stake_for_prob(prob, tiers) if stake is None: storage.log_decision( DecisionRow( ts=time.time(), symbol=sym, market_ticker=market_ticker, close_time=market_full.get("close_time",""), side=side, prob=prob, yes_bid=yes_bid, yes_ask=yes_ask, spread=spread, stake_dollars=0.0, count=0, limit_cents=0, reason="skip: no tier matched prob", ) ) continue # Determine limit price (in cents) using ask for that side if side == "yes": ask = float(market_full.get("yes_ask_dollars") or (market_full.get("yes_ask", 99) / 100.0)) else: ask = float(market_full.get("no_ask_dollars") or (market_full.get("no_ask", 99) / 100.0)) ask_cents = _dollars_to_cents_price(ask) limit_cents = max(1, ask_cents - improve_cents) # Size contracts to spend up to stake_dollars max_cost_cents = int(round(stake * 100)) count = max_cost_cents // limit_cents if count <= 0: storage.log_decision( DecisionRow( ts=time.time(), symbol=sym, market_ticker=market_ticker, close_time=market_full.get("close_time",""), side=side, prob=prob, yes_bid=yes_bid, yes_ask=yes_ask, spread=spread, stake_dollars=stake, count=0, limit_cents=limit_cents, reason="skip: count computed as 0", ) ) continue storage.log_decision( DecisionRow( ts=time.time(), symbol=sym, market_ticker=market_ticker, close_time=market_full.get("close_time",""), side=side, prob=prob, yes_bid=yes_bid, yes_ask=yes_ask, spread=spread, stake_dollars=stake, count=count, limit_cents=limit_cents, reason="trade", ) ) client_order_id = f"{sym}-{side}-{uuid.uuid4().hex[:12]}" if mode == "paper": print( f"[PAPER] {sym} {side.upper()} prob={prob:.3f} " f"{market_ticker} count={count} limit={limit_cents}c stake=${stake:.2f}" ) storage.log_order( market_ticker=market_ticker, mode="paper", client_order_id=client_order_id, order_id=None, status="simulated", details=json.dumps( {"symbol": sym, "side": side, "prob": prob, "count": count, "limit_cents": limit_cents, "stake": stake} ), ) else: order = { "ticker": market_ticker, "action": "buy", "side": side, "count": int(count), "type": "limit", "client_order_id": client_order_id, "time_in_force": tif, } if side == "yes": order["yes_price"] = int(limit_cents) else: order["no_price"] = int(limit_cents) print( f"[LIVE] {sym} {side.upper()} prob={prob:.3f} " f"{market_ticker} count={count} limit={limit_cents}c stake=${stake:.2f}" ) resp = client.create_order(order) o = resp.get("order", {}) storage.log_order( market_ticker=market_ticker, mode="live", client_order_id=client_order_id, order_id=o.get("order_id"), status=o.get("status", "unknown"), details=json.dumps(o)[:2000], ) last_traded_market[sym] = market_ticker except Exception as e: print(f"[loop] error sym={sym}: {e}") time.sleep(poll_seconds) if __name__ == "__main__": main()