diff --git a/speed_bot.py b/speed_bot.py new file mode 100644 index 0000000..5da1fa6 --- /dev/null +++ b/speed_bot.py @@ -0,0 +1,273 @@ +""" +Kalshi Speed Bot — WebSocket-driven, single-model (qwen3.5:4b-mlx), split-second decisions. +""" +import os, sys, json, time, sqlite3, base64, threading, re +from datetime import datetime, timezone +from collections import deque +import requests +import websocket + +from cryptography.hazmat.primitives import serialization, hashes +from cryptography.hazmat.primitives.asymmetric import padding + +# ── Config ── +KALSHI_REST = "https://api.elections.kalshi.com/trade-api/v2" +KALSHI_WS = "wss://external-api-ws.kalshi.com/trade-api/ws/v2" +OLLAMA_URL = "http://localhost:11434" +MODEL = "qwen3.5:4b-mlx" + +# ── Auth (RSA-PSS-SHA256, same as existing bot) ── +KEY_PATH = os.path.expanduser("~/kalshi-bot/kalshi_private_key.pem") +KEY_ID = "28d5876b" + +with open(KEY_PATH, "rb") as f: + PK = serialization.load_pem_private_key(f.read(), password=None) + +def sign(method, path): + ts = str(int(time.time() * 1000)) + msg = (ts + method.upper() + path.split("?")[0]).encode() + sig = PK.sign(msg, padding.PSS(mgf=padding.MGF1(hashes.SHA256()), + salt_length=padding.PSS.DIGEST_LENGTH), hashes.SHA256()) + return { + "KALSHI-ACCESS-KEY": KEY_ID, + "KALSHI-ACCESS-SIGNATURE": base64.b64encode(sig).decode(), + "KALSHI-ACCESS-TIMESTAMP": ts, + } + +def kx(method, path, body=None): + h = sign(method, path) + h["Content-Type"] = "application/json" + r = requests.request(method, KALSHI_REST + path, headers=h, json=body, timeout=10) + return r.json() if r.status_code in (200,201) else {"error": r.status_code, "text": r.text[:200]} + +# ── Market Utils ── +def get_open_markets(): + """Get currently open 15-min crypto markets.""" + r = requests.get(f"{KALSHI_REST}/markets?status=open&limit=100", timeout=10) + if r.status_code != 200: + return [] + markets = r.json().get("markets", []) + return [m for m in markets if "15M" in m.get("ticker","") and m.get("ticker","").startswith("KX")] + +def get_ticker(market_ticker): + r = requests.get(f"{KALSHI_REST}/markets/{market_ticker}", timeout=5) + m = r.json().get("market", r.json()) + return m + +# ── Fast LLM Vote ── +def fast_vote(context: str) -> tuple[str, float, str]: + """Single call to qwen3.5:4b-mlx. Returns (UP|DOWN|SKIP, confidence, why).""" + prompt = f"""Crypto 15-min binary. {context} +Vote UP DOWN or SKIP. Respond ONLY with JSON: {{"vote":"UP|DOWN|SKIP","conf":0.0-1.0,"why":"<10 words>"}}""" + try: + r = requests.post(f"{OLLAMA_URL}/api/generate", + json={"model": MODEL, "prompt": prompt, "stream": False, + "options": {"temperature": 0.15, "num_predict": 40}}, + timeout=8, proxies={"http": None, "https": None}) + txt = r.json().get("response", "") + s, e = txt.find("{"), txt.rfind("}")+1 + d = json.loads(txt[s:e]) if s >= 0 else {} + v = str(d.get("vote", "SKIP")).upper() + if v not in ("UP", "DOWN", "SKIP"): + v = "SKIP" + return v, float(d.get("conf", 0.5)), str(d.get("why", ""))[:80] + except Exception as e: + return "SKIP", 0.0, f"llm_err:{e}"[:40] + +# ── Indicators ── +def compute_rsi(prices, period=14): + if len(prices) < period + 1: + return 50.0 + gains = [max(prices[i] - prices[i-1], 0) for i in range(1, len(prices))] + losses = [max(prices[i-1] - prices[i], 0) for i in range(1, len(prices))] + avg_gain = sum(gains[-period:]) / period + avg_loss = sum(losses[-period:]) / period + if avg_loss == 0: + return 100.0 + return 100.0 - (100.0 / (1.0 + avg_gain / avg_loss)) + +# ── WebSocket Client ── +class SpeedBot: + def __init__(self, coin="DOGE", dry_run=True, max_spend=50): + self.coin = coin.upper() + self.dry_run = dry_run + self.max_spend = max_spend + self.prices = deque(maxlen=60) # 5 min at 1 tick/sec + self.last_decision = 0 + self.decision_cooldown = 120 # seconds between LLM calls + self.ws = None + self.running = False + self.db = sqlite3.connect(os.path.expanduser(f"~/kalshi-bot/speed_{coin.lower()}.db")) + self._init_db() + + def _init_db(self): + self.db.execute("CREATE TABLE IF NOT EXISTS ticks (ts REAL, price REAL, volume INTEGER)") + self.db.execute("CREATE TABLE IF NOT EXISTS decisions (ts REAL, ticker TEXT, vote TEXT, conf REAL, why TEXT, action TEXT)") + self.db.commit() + + def on_open(self, ws): + print(f"[speed-{self.coin}] WebSocket connected") + # Subscribe to ticker for this coin's markets + sub = { + "type": "subscribe", + "channels": ["ticker"], + "params": {"market_tickers": []} # All markets initially + } + ws.send(json.dumps(sub)) + + def on_message(self, ws, raw): + try: + msg = json.loads(raw) + except json.JSONDecodeError: + return + + msg_type = msg.get("type", "") + + if msg_type == "subscribed": + sids = msg.get("sids", []) + print(f"[speed-{self.coin}] Subscribed: {sids}") + # Update subscription to only our coin's markets + markets = get_open_markets() + our_markets = [m["ticker"] for m in markets if self.coin in m.get("ticker", "")] + if our_markets and sids: + update = { + "type": "update_subscription", + "sids": sids, + "params": {"action": "add_markets", "market_tickers": our_markets} + } + ws.send(json.dumps(update)) + print(f"[speed-{self.coin}] Targeting: {our_markets}") + + elif msg_type == "ticker": + self._process_ticker(msg) + + def _process_ticker(self, msg): + ticker = msg.get("market_ticker", "") + if self.coin not in ticker: + return + + price = float(msg.get("last_price", 0) or 0) + if price <= 0: + return + + self.prices.append(price) + self.db.execute("INSERT INTO ticks VALUES (?,?,?)", + (time.time(), price, msg.get("volume", 0) or 0)) + self.db.commit() + + # Check decision cooldown + now = time.time() + if now - self.last_decision < self.decision_cooldown: + return + + # Need enough data + if len(self.prices) < 20: + return + + self.last_decision = now + self._decide(ticker, price) + + def _decide(self, ticker, price): + prices = list(self.prices) + rsi14 = compute_rsi(prices) + + # Compute momentum over last 5 prices + if len(prices) >= 6: + short_ma = sum(prices[-3:]) / 3 + long_ma = sum(prices[-6:]) / 6 + mom = (short_ma / long_ma - 1) * 100 + else: + mom = 0 + + # Gate: RSI extremes with momentum confirmation + if rsi14 > 92 and mom < 0: + vote, conf, why = "DOWN", 0.85, "RSI extreme+fading" + elif rsi14 < 8 and mom > 0: + vote, conf, why = "UP", 0.85, "RSI oversold+bouncing" + elif rsi14 > 85 or rsi14 < 15: + # Borderline — skip, let LLM handle next cycle + return + else: + # Build compact context + ctx = f"coin={self.coin} price={price:.2f} RSI14={rsi14:.1f} 5tick_mom={mom:+.2f}%" + vote, conf, why = fast_vote(ctx) + + # Minimum confidence + if conf < 0.55: + self.db.execute("INSERT INTO decisions VALUES (?,?,?,?,?,?)", + (time.time(), ticker, vote, conf, why, "SKIP_low_conf")) + self.db.commit() + return + + # Execute + side = "yes" if vote == "UP" else "no" + order = { + "ticker": ticker, + "client_order_id": f"speed_{int(time.time())}", + "side": side, + "type": "market", + "count": 1, + "buy_max_cost": self.max_spend, + } + + if self.dry_run: + print(f"[speed-{self.coin}] DRY {vote} {ticker} @ {price:.2f} | RSI={rsi14:.1f} mom={mom:+.2f}% | {why}") + else: + result = kx("POST", "/portfolio/orders", order) + action = f"{side}_order_{result.get('order_id','err')}" + print(f"[speed-{self.coin}] LIVE {vote} {ticker} {result}") + + self.db.execute("INSERT INTO decisions VALUES (?,?,?,?,?,?)", + (time.time(), ticker, vote, conf, why, + "DRY" if self.dry_run else f"LIVE_{side}")) + self.db.commit() + + def on_error(self, ws, error): + print(f"[speed-{self.coin}] WS error: {error}") + + def on_close(self, ws, code, msg): + print(f"[speed-{self.coin}] WS closed: {code} {msg}") + self.running = False + + def run(self): + self.running = True + print(f"[speed-{self.coin}] Starting speed bot (dry_run={self.dry_run}, model={MODEL})") + + while self.running: + try: + self.ws = websocket.WebSocketApp( + KALSHI_WS, + header=self._ws_headers(), + on_open=self.on_open, + on_message=self.on_message, + on_error=self.on_error, + on_close=self.on_close, + ) + self.ws.run_forever(ping_interval=30, ping_timeout=10) + except Exception as e: + print(f"[speed-{self.coin}] Reconnecting in 5s: {e}") + time.sleep(5) + + def _ws_headers(self): + ts = str(int(time.time() * 1000)) + path = "/trade-api/ws/v2" + msg = f"{ts}GET{path}".encode() + sig = base64.b64encode(hmac.HMAC(KEY, msg, hashlib.sha256).digest()).decode() + return { + "KALSHI-ACCESS-KEY": KEY_ID, + "KALSHI-ACCESS-SIGNATURE": sig, + "KALSHI-ACCESS-TIMESTAMP": ts, + } + + +if __name__ == "__main__": + import argparse + ap = argparse.ArgumentParser() + ap.add_argument("--coin", default="DOGE") + ap.add_argument("--dry", action="store_true", default=True) + ap.add_argument("--live", dest="dry", action="store_false") + ap.add_argument("--max", type=int, default=50, help="Max spend cents per trade") + args = ap.parse_args() + + bot = SpeedBot(coin=args.coin, dry_run=args.dry, max_spend=args.max) + bot.run()