328 lines
13 KiB
Python
328 lines
13 KiB
Python
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()
|