"""OGuard trading API client for Python 3.9+. One file. REST needs only the standard library. Optional extras: pip install websockets # MarketStream and HistoryStream pip install cryptography # Ed25519 keys (HMAC keys need nothing) Three products, each with its own methods: from oguard_client import OGuardClient og = OGuardClient(api_key="ogk_...", secret="ogs_...") print(og.balances()) ack = og.perps.place_order("BTCUSDT.P", "buy", "limit", "0.001", price="50000") print(og.perps.order(ack["orderId"])) og.stocks.place_order("NVDAXUSDT", "buy", "market", "10", market_unit="quoteCoin") print(og.stocks.holdings()) og.cfd.place_order("EURUSD", "buy", "market", "0.10") print(og.cfd.positions()) Check a key from a terminal (reads OGUARD_API_KEY and OGUARD_API_SECRET): python oguard_client.py doctor Guide: https://oguard.io/developers """ from __future__ import annotations import asyncio import base64 import hashlib import hmac import http.client import json import os import random import re import secrets import ssl import sys import threading import time from decimal import Decimal from typing import Any, AsyncIterator, Dict, Iterable, List, Optional, Tuple, Union from urllib.parse import quote, urlencode, urlsplit __version__ = "2.1.0" __all__ = [ "OGuardClient", "Perpetuals", "TokenizedStocks", "CFDs", "KEEP", "CancelAllAfterHeartbeat", "OGuardError", "OGuardAPIError", "OGuardNetworkError", "MarketStream", "HistoryStream", "CfdStream", "generate_ed25519_keypair", "DEFAULT_BASE_URL", ] DEFAULT_BASE_URL = "https://api.oguard.io" Number = Union[str, int, Decimal] _CANONICAL = re.compile(r"^(0|[1-9]\d*)(\.\d+)?$") _CLIENT_ORDER_ID = re.compile(r"^[A-Za-z0-9_-]{1,30}$") _RETRYABLE_STATUS = {429, 500, 502, 503, 504} _IDLE_REUSE_S = 25.0 # reopen a kept-alive connection idle longer than this _CLOCK_RESYNC_S = 300.0 # re-read the server clock this often _MAX_RETRY_AFTER_S = 30.0 # a longer Retry-After is raised to the caller, not slept class OGuardError(Exception): """Base class for every error this module raises.""" class OGuardAPIError(OGuardError): """The API answered with a non-2xx status. `code` is the API's error code (snake_case, stable: branch on it); `message` is the readable explanation, when there is one.""" def __init__(self, status: int, code: Optional[str], body: Any, retry_after: Optional[float] = None, client_order_id: Optional[str] = None): self.status, self.code, self.body = status, code, body self.retry_after, self.client_order_id = retry_after, client_order_id self.message: Optional[str] = body.get("message") if isinstance(body, dict) else None super().__init__(f"HTTP {status}: {code or body}" + (f" ({self.message})" if self.message else "")) class OGuardNetworkError(OGuardError): """No usable response. The server may still have acted: for a placement, look it up with og..order_by_client_id(err.client_order_id).""" def __init__(self, message: str, timed_out: bool = False, client_order_id: Optional[str] = None): self.timed_out, self.client_order_id = timed_out, client_order_id super().__init__(message) # ── Values ──────────────────────────────────────────────────────────────────── def decimal_str(value: Number, field: str = "value") -> str: """A price or quantity as the API's canonical decimal string. Floats are refused: most decimal prices have no exact float, and the API never rounds.""" if isinstance(value, bool) or isinstance(value, float): raise TypeError(f"{field} must be str, int or Decimal, not {type(value).__name__}: " f"use Decimal('0.001') or '0.001'") if isinstance(value, int): text = str(value) elif isinstance(value, Decimal): if not value.is_finite(): raise ValueError(f"{field} must be finite") text = format(value, "f") elif isinstance(value, str): text = value else: raise TypeError(f"{field} must be str, int or Decimal") if not _CANONICAL.match(text): raise ValueError(f"{field}={value!r} is not a canonical decimal like '0.01' or '65000.5'") return text def new_client_order_id() -> str: """A fresh id for an order: 'py-' and 22 random url-safe characters.""" return "py-" + base64.urlsafe_b64encode(secrets.token_bytes(16)).decode().rstrip("=")[:22] def _check_client_order_id(cid: Optional[str]) -> str: if cid is None: return new_client_order_id() if not isinstance(cid, str) or not _CLIENT_ORDER_ID.match(cid): raise ValueError("client_order_id must be 1-30 characters of [A-Za-z0-9_-]") return cid # ── Signing ─────────────────────────────────────────────────────────────────── def sign_hmac(secret: str, pre_sign: str) -> str: return hmac.new(secret.encode(), pre_sign.encode(), hashlib.sha256).hexdigest() def _load_ed25519(private_key: Union[str, bytes]) -> Any: try: from cryptography.hazmat.primitives import serialization from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey except ImportError as exc: # pragma: no cover - depends on the environment raise OGuardError("Ed25519 keys need the 'cryptography' package: pip install cryptography") from exc data = private_key.encode() if isinstance(private_key, str) else private_key if b"-----BEGIN" in data: return serialization.load_pem_private_key(data, password=None) raw = base64.b64decode(data) if len(raw) != 32: raise ValueError("private_key must be a PEM or the base64 of a raw 32-byte Ed25519 seed") return Ed25519PrivateKey.from_private_bytes(raw) def generate_ed25519_keypair() -> Tuple[str, str]: """(private key PEM, public key base64). Keep the PEM; paste the base64 on the API keys page.""" from cryptography.hazmat.primitives import serialization from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey key = Ed25519PrivateKey.generate() pem = key.private_bytes(serialization.Encoding.PEM, serialization.PrivateFormat.PKCS8, serialization.NoEncryption()).decode() pub = key.public_key().public_bytes(serialization.Encoding.Raw, serialization.PublicFormat.Raw) return pem, base64.b64encode(pub).decode() # ── Transport: kept-alive HTTP(S) connections, safe to share across threads ─── class _Transport: def __init__(self, base_url: str, timeout: float): u = urlsplit(base_url) if u.scheme not in ("http", "https") or not u.hostname: raise ValueError(f"base_url must be http(s)://host, got {base_url!r}") self.https, self.host = u.scheme == "https", u.hostname self.port = u.port or (443 if self.https else 80) self.prefix = u.path.rstrip("/") self.timeout = timeout self._ssl = ssl.create_default_context() if self.https else None self._idle: List[Tuple[http.client.HTTPConnection, float]] = [] self._lock = threading.Lock() def _connect(self) -> http.client.HTTPConnection: if self.https: return http.client.HTTPSConnection(self.host, self.port, timeout=self.timeout, context=self._ssl) return http.client.HTTPConnection(self.host, self.port, timeout=self.timeout) def _acquire(self) -> Tuple[http.client.HTTPConnection, bool]: with self._lock: while self._idle: conn, since = self._idle.pop() if time.monotonic() - since < _IDLE_REUSE_S: return conn, True conn.close() return self._connect(), False def send(self, method: str, target: str, body: Optional[bytes], headers: Dict[str, str], resend_on_stale: bool) -> Tuple[int, Dict[str, str], bytes]: """(status, headers, body). Raises OSError / HTTPException on failure. A kept-alive connection the server has already closed fails at once with a reset. When the call is safe to repeat it is resent once on a fresh connection; otherwise the error is raised, because the server may have read the request before closing.""" conn, reused = self._acquire() try: conn.request(method, self.prefix + target, body=body, headers=headers) resp = conn.getresponse() data = resp.read() except (http.client.RemoteDisconnected, ConnectionResetError, BrokenPipeError, ConnectionAbortedError): conn.close() if not (reused and resend_on_stale): raise conn = self._connect() try: conn.request(method, self.prefix + target, body=body, headers=headers) resp = conn.getresponse() data = resp.read() except BaseException: conn.close() raise except BaseException: conn.close() raise hdrs = {k.lower(): v for k, v in resp.getheaders()} if resp.will_close: conn.close() else: with self._lock: self._idle.append((conn, time.monotonic())) return resp.status, hdrs, data def close(self) -> None: with self._lock: for conn, _ in self._idle: conn.close() self._idle.clear() # ── REST client ─────────────────────────────────────────────────────────────── class OGuardClient: """Signed REST client. Thread-safe; share one per key. Provide `secret` for an HMAC key or `private_key` (PEM or base64 seed) for an Ed25519 key. Each product has its own methods: og.perps, og.stocks and og.cfd. Account methods (balances, transactions, wallet reads) sit on the client. Perpetual and stock market data work without a key. """ def __init__(self, api_key: Optional[str] = None, secret: Optional[str] = None, *, private_key: Union[str, bytes, None] = None, base_url: str = DEFAULT_BASE_URL, recv_window: int = 5000, timeout: float = 10.0, max_retries: int = 2, sync_clock: bool = True, decimals: bool = True): if secret and private_key: raise ValueError("provide secret (HMAC) or private_key (Ed25519), not both") if (secret or private_key) and not api_key: raise ValueError("api_key is required with a secret or private_key") if not 1 <= recv_window <= 60000: raise ValueError("recv_window must be 1..60000 ms") self.api_key, self.base_url = api_key, base_url.rstrip("/") self._secret, self._ed = secret, _load_ed25519(private_key) if private_key else None self.recv_window, self.max_retries, self.sync_clock = recv_window, max(0, max_retries), sync_clock self._parse_float = Decimal if decimals else float self._transport = _Transport(self.base_url, timeout) self._clock_lock = threading.Lock() self._clock: Optional[Tuple[float, float]] = None # (server ms, monotonic s) at the last sync self.clock_offset_ms: Optional[float] = None # server minus local, for information self.last_rate_limit: Dict[str, int] = {} self.perps = Perpetuals(self) self.stocks = TokenizedStocks(self) self.cfd = CFDs(self) # ── plumbing ────────────────────────────────────────────────────────────── def close(self) -> None: self._transport.close() def __enter__(self) -> "OGuardClient": return self def __exit__(self, *exc: Any) -> None: self.close() def sync_time(self) -> float: """Read the server clock and sign with it from now on. Returns the offset (ms).""" t0 = time.monotonic() server = self.request("GET", "/api/v1/time", signed=False)["serverTime"] t1 = time.monotonic() with self._clock_lock: self._clock = (float(server), (t0 + t1) / 2) self.clock_offset_ms = float(server) - time.time() * 1000 + (time.monotonic() - (t0 + t1) / 2) * 1000 return self.clock_offset_ms def _now_ms(self) -> int: if not self.sync_clock: return int(time.time() * 1000) with self._clock_lock: clock = self._clock if clock is None or time.monotonic() - clock[1] > _CLOCK_RESYNC_S: self.sync_time() with self._clock_lock: clock = self._clock assert clock is not None return int(clock[0] + (time.monotonic() - clock[1]) * 1000) def _signed_headers(self, payload: str) -> Dict[str, str]: if not self.api_key or not (self._secret or self._ed): raise OGuardError("this call needs an API key: pass api_key with secret or private_key") ts, window = str(self._now_ms()), str(self.recv_window) pre = ts + self.api_key + window + payload if self._secret: sig, kind = sign_hmac(self._secret, pre), "hmac" else: sig, kind = base64.b64encode(self._ed.sign(pre.encode())).decode(), "ed25519" return {"X-OG-API-KEY": self.api_key, "X-OG-TIMESTAMP": ts, "X-OG-RECV-WINDOW": window, "X-OG-SIGN-TYPE": kind, "X-OG-SIGN": sig} def request(self, method: str, path: str, params: Optional[Dict[str, Any]] = None, body: Optional[Dict[str, Any]] = None, *, signed: bool = True, retry: Optional[bool] = None, client_order_id: Optional[str] = None) -> Any: """Send one API call and return the parsed JSON. Reads (GET) and placements retry on network errors, 5xx and 429; other writes never retry.""" method = method.upper() retry = (method == "GET") if retry is None else retry attempts, clock_retried = (self.max_retries + 1 if retry else 1), False attempt = 0 while True: try: return self._once(method, path, params, body, signed, client_order_id, retry) except OGuardAPIError as err: if err.code == "timestamp_invalid" and self.sync_clock and not clock_retried: clock_retried = True # refused before anything ran: safe to resend self.sync_time() continue attempt += 1 if not retry or attempt >= attempts or err.status not in _RETRYABLE_STATUS: raise if err.retry_after is not None and err.retry_after > _MAX_RETRY_AFTER_S: raise time.sleep(err.retry_after if err.retry_after else self._backoff(attempt)) except OGuardNetworkError: attempt += 1 if not retry or attempt >= attempts: raise time.sleep(self._backoff(attempt)) @staticmethod def _backoff(attempt: int) -> float: base = 0.25 * 2 ** (attempt - 1) return base / 2 + random.random() * base / 2 def _once(self, method: str, path: str, params: Optional[Dict[str, Any]], body: Optional[Dict[str, Any]], signed: bool, cid: Optional[str], may_resend: bool) -> Any: clean = {k: _param(v) for k, v in (params or {}).items() if v is not None} qs = urlencode(clean, quote_via=quote, safe=",") is_write = method in ("POST", "PUT", "PATCH") raw = json.dumps(body or {}, separators=(",", ":")) if is_write else "" headers = {"User-Agent": f"oguard-python/{__version__} (Python {sys.version.split()[0]})", "Accept": "application/json"} if signed: headers.update(self._signed_headers(raw if is_write else qs)) if is_write: headers["Content-Type"] = "application/json" target = path + (f"?{qs}" if qs else "") try: status, hdrs, data = self._transport.send(method, target, raw.encode() if is_write else None, headers, may_resend) except (OSError, http.client.HTTPException) as exc: timed_out = isinstance(exc, TimeoutError) or "timed out" in str(exc) raise OGuardNetworkError(f"{method} {path}: {exc}", timed_out, cid) from exc self._record_limits(hdrs) try: parsed = json.loads(data.decode("utf-8"), parse_float=self._parse_float) if data else None except ValueError: parsed = data.decode("utf-8", "replace")[:500] if status >= 400: code = parsed.get("error") if isinstance(parsed, dict) and isinstance(parsed.get("error"), str) else None ra = hdrs.get("retry-after") raise OGuardAPIError(status, code, parsed, float(ra) if ra else None, cid) return parsed def _record_limits(self, hdrs: Dict[str, str]) -> None: names = {"x-og-used-weight-1m": "used_weight_1m", "x-og-limit-weight-1m": "limit_weight_1m", "x-og-order-count-10s": "order_count_10s", "x-og-limit-order-10s": "limit_order_10s", "x-og-order-count-1m": "order_count_1m", "x-og-limit-order-1m": "limit_order_1m"} got = {v: int(hdrs[k]) for k, v in names.items() if k in hdrs and hdrs[k].lstrip("-").isdigit()} if got: self.last_rate_limit = got def get(self, path: str, params: Optional[Dict[str, Any]] = None, *, signed: bool = True) -> Any: return self.request("GET", path, params, signed=signed) def post(self, path: str, body: Optional[Dict[str, Any]] = None) -> Any: return self.request("POST", path, body=body) @staticmethod def _retry_404(call: Any, wait: float) -> Any: """Run `call`; while it answers 404 and `wait` seconds have not passed, run it again.""" deadline = time.monotonic() + max(0.0, wait) delay = 0.05 while True: try: return call() except OGuardAPIError as err: if err.status != 404 or time.monotonic() + delay > deadline: raise time.sleep(delay) delay = min(delay * 2, 0.4) # ── account (every product) ─────────────────────────────────────────────── def server_time(self) -> int: return self.request("GET", "/api/v1/time", signed=False)["serverTime"] def balances(self) -> Any: """Balance, equity and margin of the trading account perpetuals and tokenized stocks share. CFDs have their own wallet: see og.cfd.account().""" return self.get("/api/v1/account/balances") def transactions(self, days: Optional[int] = None) -> Any: return self.get("/api/v1/account/transactions", {"days": days}) def deposit_info(self) -> Any: return self.get("/api/v1/account/wallet/deposit-info") def withdraw_info(self) -> Any: return self.get("/api/v1/account/wallet/withdraw-info") def withdrawal_addresses(self) -> Any: return self.get("/api/v1/account/wallet/withdrawal-addresses") def cancel_all_after(self, timeout_ms: int) -> Any: """Arm the dead-man's switch: unless this is called again within `timeout_ms` (5000 to 600000), every open order on every product is cancelled. Positions are left alone. 0 disarms. Call it as a heartbeat, or let CancelAllAfterHeartbeat do it for you.""" return self.request("POST", "/api/v1/account/cancel-all-after", body={"timeoutMs": int(timeout_ms)}) def cancel_all_after_state(self) -> Any: return self.get("/api/v1/account/cancel-all-after") def ws_ticket(self) -> str: """A 60-second ticket for /ws/history or /ws/cfd. The streams fetch their own.""" return self.get("/api/v1/ws/ticket")["token"] class _Product: prefix = "" def __init__(self, client: OGuardClient): self._c = client def _get(self, path: str, params: Optional[Dict[str, Any]] = None, *, signed: bool = True) -> Any: return self._c.request("GET", self.prefix + path, params, signed=signed) def _history(self, path: str, days: Optional[int], window: Optional[str], symbol: Optional[str], limit: Optional[int]) -> Any: return self._get(path, {"days": days, "window": window, "symbol": symbol, "limit": limit}) class _OrderBook(_Product): """Reads shared by perpetuals and tokenized stocks: OGuard's record of the product's orders and fills, which follows the market within about a second.""" def markets(self) -> Any: return self._get("/markets", signed=False) def tickers(self, symbols: Optional[Iterable[str]] = None) -> Any: return self._get("/tickers", {"symbols": ",".join(symbols) if symbols else None}, signed=False) def kline(self, symbol: str, interval: str, start: Optional[int] = None, end: Optional[int] = None, limit: Optional[int] = None) -> Any: """interval: 1 3 5 15 30 60 120 240 360 720 D W M. start/end: Unix ms.""" return self._get(f"/markets/{_seg(symbol)}/kline", {"interval": interval, "start": start, "end": end, "limit": limit}, signed=False) def orderbook(self, symbol: str, limit: int = 25) -> Any: """Order book snapshot: {"bids": [{"price", "size"}], "asks": [...], "ts"}. limit 1..200.""" return self._get(f"/markets/{_seg(symbol)}/orderbook", {"limit": limit}, signed=False) def recent_trades(self, symbol: str, limit: int = 50) -> Any: """Recent public trades, newest first.""" return self._get(f"/markets/{_seg(symbol)}/trades", {"limit": limit}, signed=False) def open_orders(self) -> Any: return self._get("/orders") def order(self, order_id: str, wait: float = 2.0) -> Any: """One order by id. A new order is recorded a moment after it is acknowledged, so a 404 is retried for up to `wait` seconds; wait=0 asks once.""" return self._c._retry_404(lambda: self._get(f"/orders/{_seg(order_id)}"), wait) def order_by_client_id(self, client_order_id: str, wait: float = 2.0) -> Any: """One order by your clientOrderId: how to resolve a placement whose response you never got. Still 404 after `wait` seconds means it was not placed.""" return self._c._retry_404(lambda: self._get(f"/orders/by-client-id/{_seg(client_order_id)}"), wait) def order_history(self, days: Optional[int] = None, window: Optional[str] = None, symbol: Optional[str] = None, limit: Optional[int] = None) -> Any: """days: 1 7 15 30 45 60, or window='24h'.""" return self._history("/orders/history", days, window, symbol, limit) def trades(self, days: Optional[int] = None, window: Optional[str] = None, symbol: Optional[str] = None, limit: Optional[int] = None) -> Any: return self._history("/trades", days, window, symbol, limit) class Perpetuals(_OrderBook): """USDT perpetuals ('BTCUSDT.P'), under /api/v1/perps. Accounts trade in hedge mode: position_idx 1 is the long side, 2 the short side.""" prefix = "/api/v1/perps" def contract_details(self, symbol: str) -> Any: return self._get(f"/markets/{_seg(symbol)}/contract-details", signed=False) def funding_history(self, symbol: str, limit: Optional[int] = None) -> Any: return self._get(f"/markets/{_seg(symbol)}/funding-history", {"limit": limit}, signed=False) def open_interest(self, symbol: str, interval: str = "1h", limit: int = 50) -> Any: """interval: 5min 15min 30min 1h 4h 1d.""" return self._get(f"/markets/{_seg(symbol)}/open-interest", {"interval": interval, "limit": limit}, signed=False) def place_order(self, symbol: str, side: str, order_type: str, qty: Number, *, price: Optional[Number] = None, tif: Optional[str] = None, reduce_only: Optional[bool] = None, position_idx: Optional[int] = None, take_profit: Optional[Number] = None, stop_loss: Optional[Number] = None, trigger_price: Optional[Number] = None, trigger_by: Optional[str] = None, trigger_direction: Optional[int] = None, leverage: Optional[int] = None, client_order_id: Optional[str] = None) -> Any: """Place a perpetual order. Idempotent: every attempt carries the same client_order_id (yours, or one generated here), so a retry returns the order already placed (deduplicated: True) instead of placing a second one.""" cid = _check_client_order_id(client_order_id) body = _drop_none({ "symbol": symbol, "side": side, "type": order_type, "qty": decimal_str(qty, "qty"), "price": _opt_dec(price, "price"), "tif": tif, "reduceOnly": reduce_only, "positionIdx": position_idx, "takeProfit": _opt_dec(take_profit, "take_profit"), "stopLoss": _opt_dec(stop_loss, "stop_loss"), "triggerPrice": _opt_dec(trigger_price, "trigger_price"), "triggerBy": trigger_by, "triggerDirection": trigger_direction, "leverage": leverage, "clientOrderId": cid, }) return self._c.request("POST", f"{self.prefix}/orders", body=body, retry=True, client_order_id=cid) # Amends and cancels act on OGuard's record of the order, made a moment after # the placement is acknowledged. A 404 means nothing was done, so these retry # it for up to `wait` seconds: placing and then at once cancelling works. def amend_order(self, order_id: str, *, qty: Optional[Number] = None, price: Optional[Number] = None, trigger_price: Optional[Number] = None, take_profit: Optional[Number] = None, stop_loss: Optional[Number] = None, wait: float = 2.0) -> Any: """Change what you pass. take_profit / stop_loss '0' clears that leg.""" body = _drop_none({"qty": _opt_dec(qty, "qty"), "price": _opt_dec(price, "price"), "triggerPrice": _opt_dec(trigger_price, "trigger_price"), "takeProfit": _opt_dec(take_profit, "take_profit"), "stopLoss": _opt_dec(stop_loss, "stop_loss")}) return self._c._retry_404( lambda: self._c.request("POST", f"{self.prefix}/orders/{_seg(order_id)}/amend", body=body), wait) def cancel_order(self, order_id: str, wait: float = 2.0) -> Any: return self._c._retry_404( lambda: self._c.request("POST", f"{self.prefix}/orders/{_seg(order_id)}/cancel", body={}), wait) def cancel_order_by_client_id(self, client_order_id: str) -> Any: return self._c.request("POST", f"{self.prefix}/orders/by-client-id/{_seg(client_order_id)}/cancel", body={}) def amend_order_by_client_id(self, client_order_id: str, *, qty: Optional[Number] = None, price: Optional[Number] = None, trigger_price: Optional[Number] = None, take_profit: Optional[Number] = None, stop_loss: Optional[Number] = None) -> Any: body = _drop_none({"qty": _opt_dec(qty, "qty"), "price": _opt_dec(price, "price"), "triggerPrice": _opt_dec(trigger_price, "trigger_price"), "takeProfit": _opt_dec(take_profit, "take_profit"), "stopLoss": _opt_dec(stop_loss, "stop_loss")}) return self._c.request("POST", f"{self.prefix}/orders/by-client-id/{_seg(client_order_id)}/amend", body=body) def cancel_all_orders(self) -> Any: return self._c.request("POST", f"{self.prefix}/orders/cancel-all", body={}) def set_leverage(self, symbol: str, leverage: int) -> Any: return self._c.request("POST", f"{self.prefix}/leverage", body={"symbol": symbol, "leverage": int(leverage)}) def leverage_info(self, symbol: str) -> Any: return self._get("/leverage-info", {"symbol": symbol}) def leverage_impact(self, symbol: str, leverage: int) -> Any: return self._get("/leverage-impact", {"symbol": symbol, "leverage": int(leverage)}) def positions(self) -> Any: return self._get("/positions") def position(self, symbol: str, side: Optional[str] = None, wait: float = 0.0) -> Optional[Dict[str, Any]]: """The open position on `symbol` ('buy' is the long side, 'sell' the short), or None. After an order, pass wait=5 to give the fill time to be recorded.""" deadline = time.monotonic() + max(0.0, wait) while True: for p in self.positions().get("positions", []): if p.get("symbol") == symbol and (side is None or p.get("side") == side) \ and Decimal(str(p.get("quantity", "0"))) > 0: return p if time.monotonic() >= deadline: return None time.sleep(0.25) def reverse_position(self, position_id: str) -> Any: return self._c.request("POST", f"{self.prefix}/positions/{_seg(position_id)}/reverse", body={}) def set_tpsl(self, position_id: str, mode: str = "entire", *, take_profit_pct: Optional[Number] = None, stop_loss_pct: Optional[Number] = None, partial_qty: Optional[Number] = None) -> Any: """mode 'entire' or 'partial'. Percentages are of the last price, e.g. '2.5'.""" body = _drop_none({"mode": mode, "takeProfitPct": _opt_dec(take_profit_pct, "take_profit_pct"), "stopLossPct": _opt_dec(stop_loss_pct, "stop_loss_pct"), "partialQty": _opt_dec(partial_qty, "partial_qty")}) return self._c.request("POST", f"{self.prefix}/positions/{_seg(position_id)}/tpsl", body=body) def clear_tpsl(self, position_id: str) -> Any: return self._c.request("DELETE", f"{self.prefix}/positions/{_seg(position_id)}/tpsl") def closed_pnl(self, days: Optional[int] = None, window: Optional[str] = None, symbol: Optional[str] = None, limit: Optional[int] = None) -> Any: return self._history("/closed-pnl", days, window, symbol, limit) class TokenizedStocks(_OrderBook): """Tokenized stocks ('NVDAXUSDT'), under /api/v1/stocks. Bought and sold outright from the trading balance; what you own is in holdings().""" prefix = "/api/v1/stocks" def place_order(self, symbol: str, side: str, order_type: str, qty: Number, market_unit: str = "baseCoin", *, price: Optional[Number] = None, trigger_price: Optional[Number] = None, client_order_id: Optional[str] = None) -> Any: """market_unit: 'baseCoin' (qty in shares' tokens) or 'quoteCoin' (qty in USDT, market buys only); a limit order's qty is always in the base coin. trigger_price makes it conditional: an order of `order_type` is placed when the last price reaches it. Idempotent, like Perpetuals.place_order.""" cid = _check_client_order_id(client_order_id) conditional = None if trigger_price is None else { "triggerPrice": decimal_str(trigger_price, "trigger_price"), "mode": order_type} body = _drop_none({"symbol": symbol, "side": side, "orderType": order_type, "qty": decimal_str(qty, "qty"), "marketUnit": market_unit, "price": _opt_dec(price, "price"), "conditional": conditional, "clientOrderId": cid}) return self._c.request("POST", f"{self.prefix}/orders", body=body, retry=True, client_order_id=cid) def cancel_order(self, order_id: str, wait: float = 2.0) -> Any: return self._c._retry_404( lambda: self._c.request("POST", f"{self.prefix}/orders/{_seg(order_id)}/cancel", body={}), wait) def cancel_order_by_client_id(self, client_order_id: str) -> Any: return self._c.request("POST", f"{self.prefix}/orders/by-client-id/{_seg(client_order_id)}/cancel", body={}) def cancel_all_orders(self) -> Any: """Cancel every open tokenized-stock order.""" return self._c.request("POST", f"{self.prefix}/orders/cancel-all", body={}) def holdings(self) -> Any: return self._get("/holdings") # Sentinel for CFDs.set_protection: leave this level as it is. KEEP: Any = object() class CFDs(_Product): """CFDs on FX, metals, indices and energy ('EURUSD', 'XAUUSD'), under /api/v1/cfd. They trade on their own wallet: move money into it from the OGuard website (a key never moves funds). Sizes are lots. Order ids are integers. Reads come straight from the ledger the order was written to, so there is no delay between placing an order and reading it back.""" prefix = "/api/v1/cfd" def instruments(self) -> Any: return self._get("/instruments") def quotes(self, symbols: Optional[Iterable[str]] = None) -> Any: return self._get("/quotes", {"symbols": ",".join(symbols) if symbols else None}) def timeframes(self) -> Any: return self._get("/timeframes") def bars(self, symbol: str, timeframe: str = "1h") -> Any: """timeframe: an id from timeframes(): 1m 5m 15m 30m 1h 4h 6h 12h 1d 1w 1M.""" return self._get("/bars", {"symbol": symbol, "timeframe": timeframe}) def account(self) -> Any: """Balance, equity, margin, free margin and margin level of the CFD wallet.""" return self._get("/account") def positions(self) -> Any: """Open positions, priced against the live market.""" return self._get("/positions") def pending_orders(self) -> Any: return self._get("/orders", {"status": "pending"}) def closed_orders(self, days: Optional[int] = None, limit: Optional[int] = None) -> Any: """The latest `limit` (default 100, at most 1000) closed orders, or those that closed in the last `days` (1 to 366).""" return self._get("/orders", {"status": "closed", "days": days, "limit": limit}) def order_by_client_id(self, client_order_id: str) -> Any: """The order placed under your clientOrderId, whatever its state now. 404 means nothing was placed under it.""" return self._get(f"/orders/by-client-id/{_seg(client_order_id)}") def preview(self, symbol: str, side: str, volume: Number, order_type: str = "market", entry_price: Optional[Number] = None) -> Any: """Margin and value an order would take, without placing it.""" return self._get("/order-preview", {"symbol": symbol, "side": side, "volume": decimal_str(volume, "volume"), "orderType": order_type, "entryPrice": _opt_dec(entry_price, "entry_price")}) def place_order(self, symbol: str, side: str, order_type: str = "market", volume: Optional[Number] = None, *, value: Optional[Number] = None, entry_price: Optional[Number] = None, stop_limit_price: Optional[Number] = None, take_profit: Optional[Number] = None, stop_loss: Optional[Number] = None, trailing_stop_points: Optional[int] = None, trailing_step_points: Optional[int] = None, expires_at: Optional[int] = None, client_order_id: Optional[str] = None) -> Any: """Open a position (market) or rest an order (limit, stop, stop_limit at entry_price; stop_limit also takes stop_limit_price). Size is `volume` in lots, or `value` in account currency. expires_at: Unix ms, pending orders only. Idempotent: a retry with the same client_order_id returns the order already placed (deduplicated: True).""" if (volume is None) == (value is None): raise ValueError("pass volume (lots) or value, not both") cid = _check_client_order_id(client_order_id) body = _drop_none({ "symbol": symbol, "side": side, "orderType": order_type, "volume": _opt_dec(volume, "volume"), "value": _opt_dec(value, "value"), "entryPrice": _opt_dec(entry_price, "entry_price"), "stopLimitPrice": _opt_dec(stop_limit_price, "stop_limit_price"), "takeProfit": _opt_dec(take_profit, "take_profit"), "stopLoss": _opt_dec(stop_loss, "stop_loss"), "trailingStopPoints": _opt_points(trailing_stop_points, "trailing_stop_points"), "trailingStepPoints": _opt_points(trailing_step_points, "trailing_step_points"), "expiresAt": None if expires_at is None else int(expires_at), "clientOrderId": cid, }) return self._c.request("POST", f"{self.prefix}/orders", body=body, retry=True, client_order_id=cid) def close(self, order_id: int, volume: Optional[Number] = None) -> Any: """Close a position, or `volume` lots of it. Not retried: a repeated partial close would close more.""" return self._c.request("POST", f"{self.prefix}/orders/{_seg(str(order_id))}/close", body=_drop_none({"volume": _opt_dec(volume, "volume")})) def cancel_order(self, order_id: int) -> Any: """Cancel a pending order.""" return self._c.request("DELETE", f"{self.prefix}/orders/{_seg(str(order_id))}") def close_by_client_id(self, client_order_id: str, volume: Optional[Number] = None) -> Any: return self._c.request("POST", f"{self.prefix}/orders/by-client-id/{_seg(client_order_id)}/close", body=_drop_none({"volume": _opt_dec(volume, "volume")})) def cancel_order_by_client_id(self, client_order_id: str) -> Any: return self._c.request("DELETE", f"{self.prefix}/orders/by-client-id/{_seg(client_order_id)}") def cancel_all_orders(self) -> Any: """Cancel every pending CFD order.""" return self._c.request("POST", f"{self.prefix}/orders/cancel-all", body={}) def modify_order(self, order_id: int, *, entry_price: Optional[Number] = None, stop_limit_price: Optional[Number] = None) -> Any: """Move a pending order's entry (and a stop-limit's limit) price.""" body = _drop_none({"entryPrice": _opt_dec(entry_price, "entry_price"), "stopLimitPrice": _opt_dec(stop_limit_price, "stop_limit_price")}) return self._c.request("PATCH", f"{self.prefix}/orders/{_seg(str(order_id))}/pending", body=body) def set_protection(self, order_id: int, *, take_profit: Any = KEEP, stop_loss: Any = KEEP, trailing_stop_points: Any = KEEP, trailing_step_points: Any = KEEP) -> Any: """Take-profit, stop-loss and trailing stop on a position or pending order. A level left as KEEP is unchanged; None clears it.""" body: Dict[str, Any] = {} for name, field, value in (("takeProfit", "take_profit", take_profit), ("stopLoss", "stop_loss", stop_loss)): if value is not KEEP: body[name] = None if value is None else decimal_str(value, field) for name, field, value in (("trailingStopPoints", "trailing_stop_points", trailing_stop_points), ("trailingStepPoints", "trailing_step_points", trailing_step_points)): if value is not KEEP: body[name] = _opt_points(value, field) return self._c.request("PATCH", f"{self.prefix}/orders/{_seg(str(order_id))}/protection", body=body) def _seg(value: str) -> str: return quote(str(value), safe="") def _param(v: Any) -> str: if isinstance(v, bool): return "true" if v else "false" return v if isinstance(v, str) else str(v) def _opt_dec(v: Optional[Number], field: str) -> Optional[str]: return None if v is None else decimal_str(v, field) def _opt_points(v: Any, field: str) -> Optional[str]: """A trailing distance in whole points, as the string the API reads.""" if v is None: return None if isinstance(v, bool) or not (isinstance(v, int) or (isinstance(v, str) and v.isdigit())) or int(v) <= 0: raise ValueError(f"{field} must be a positive whole number of points") return str(int(v)) def _drop_none(d: Dict[str, Any]) -> Dict[str, Any]: return {k: v for k, v in d.items() if v is not None} class CancelAllAfterHeartbeat: """Keep cancel-all-after armed from a background thread while your program runs. with CancelAllAfterHeartbeat(og, timeout_ms=15000, every_s=5): run_strategy() # if this process dies or hangs, orders are cancelled # within 15 s; on a clean exit the switch is disarmed A heartbeat that fails is retried at the next interval; `last_error` holds the latest failure, and `last_state` the latest armed state. """ def __init__(self, client: OGuardClient, timeout_ms: int = 15_000, every_s: float = 5.0): if every_s * 1000 >= timeout_ms: raise ValueError("every_s must be well under timeout_ms, or the switch fires between beats") self.client, self.timeout_ms, self.every_s = client, timeout_ms, every_s self.last_state: Any = None self.last_error: Optional[BaseException] = None self._stop = threading.Event() self._thread: Optional[threading.Thread] = None def _beat(self) -> None: try: self.last_state = self.client.cancel_all_after(self.timeout_ms) self.last_error = None except OGuardError as exc: self.last_error = exc def _run(self) -> None: while not self._stop.wait(self.every_s): self._beat() def start(self) -> "CancelAllAfterHeartbeat": self.last_state = self.client.cancel_all_after(self.timeout_ms) # fail loudly if it cannot arm self._thread = threading.Thread(target=self._run, name="oguard-cancel-all-after", daemon=True) self._thread.start() return self def stop(self, disarm: bool = True) -> None: self._stop.set() if self._thread is not None: self._thread.join(timeout=self.every_s + 1) if disarm: self.client.cancel_all_after(0) def __enter__(self) -> "CancelAllAfterHeartbeat": return self.start() def __exit__(self, *exc: Any) -> None: self.stop() # ── Streams (pip install websockets) ────────────────────────────────────────── def _ws_origin(base_url: str) -> str: return re.sub(r"^http", "ws", base_url.rstrip("/")) class _Stream: """Reconnecting WebSocket with an app-level ping. Iterate with `async for`.""" path = "" app_ping = False def __init__(self, base_url: str, ping_interval: float = 20.0, max_backoff: float = 30.0): try: import websockets # noqa: F401 except ImportError as exc: # pragma: no cover - depends on the environment raise OGuardError("streams need the 'websockets' package: pip install websockets") from exc self.origin, self.ping_interval, self.max_backoff = _ws_origin(base_url), ping_interval, max_backoff self._ws: Any = None self._closed = False self._queue: "Optional[asyncio.Queue[Optional[Dict[str, Any]]]]" = None self._runner: Optional["asyncio.Task[None]"] = None async def __aenter__(self) -> "_Stream": await self.start() return self async def __aexit__(self, *exc: Any) -> None: await self.close() def __aiter__(self) -> AsyncIterator[Dict[str, Any]]: return self._iterate() async def _iterate(self) -> AsyncIterator[Dict[str, Any]]: if self._queue is None: await self.start() assert self._queue is not None while True: msg = await self._queue.get() if msg is None: return yield msg async def start(self) -> None: if self._runner is None: self._queue = asyncio.Queue() self._runner = asyncio.ensure_future(self._run()) async def close(self) -> None: """Stop for good; an `async for` over the stream ends.""" if self._closed and self._runner is None: return self._closed = True runner, self._runner = self._runner, None if runner is not None: runner.cancel() try: await runner except (asyncio.CancelledError, Exception): pass if self._queue is not None: self._queue.put_nowait(None) def _emit(self, msg: Optional[Dict[str, Any]]) -> None: if self._queue is not None: self._queue.put_nowait(msg) def _lifecycle(self, name: str, **extra: Any) -> Dict[str, Any]: """A local event: {"kind": name} on the history stream, {"ch": name} on market.""" return {"kind": name, **extra} async def _url(self) -> str: return self.origin + self.path async def _on_open(self) -> None: pass def _on_frame(self, msg: Dict[str, Any]) -> None: self._emit(msg) @property def connected(self) -> bool: return self._ws is not None async def send(self, obj: Dict[str, Any]) -> None: """Send a raw frame now. Raises OGuardError while disconnected; subscribe() is the call that survives reconnects.""" if self._ws is None: raise OGuardError("stream is not connected") await self._ws.send(json.dumps(obj)) async def _run(self) -> None: import websockets attempt = 0 while not self._closed: try: url = await self._url() async with websockets.connect(url, max_size=2 ** 23, close_timeout=5) as ws: self._ws, attempt = ws, 0 self._emit(self._lifecycle("open")) await self._on_open() pinger = asyncio.ensure_future(self._ping()) if self.app_ping else None try: async for raw in ws: try: msg = json.loads(raw, parse_float=Decimal) except ValueError: continue if isinstance(msg, dict): self._on_frame(msg) finally: if pinger is not None: pinger.cancel() self._ws = None except asyncio.CancelledError: raise except OGuardAPIError as exc: self._emit(self._lifecycle("error", error=str(exc), status=exc.status, code=exc.code)) if exc.status in (401, 403): # key revoked, expired or refused: retrying cannot help self._closed = True self._emit(None) return except Exception as exc: # network or handshake failure: reconnect self._emit(self._lifecycle("error", error=str(exc))) if self._closed: break self._emit(self._lifecycle("close")) attempt += 1 delay = min(self.max_backoff, 0.5 * 2 ** min(attempt, 6)) await asyncio.sleep(delay / 2 + random.random() * delay / 2) async def _ping(self) -> None: while True: await asyncio.sleep(self.ping_interval) if self.connected: await self.send({"op": "ping"}) class MarketStream(_Stream): """Public market data: tickers, klines, order book depth, trades. No key needed. async with MarketStream() as market: await market.subscribe("BTCUSDT.P") # ticker await market.subscribe("BTCUSDT.P", "kline", "1") # 1-minute candles async for msg in market: print(msg["ch"], msg.get("data")) """ path = "/ws/market" app_ping = True def __init__(self, base_url: str = DEFAULT_BASE_URL, **kw: Any): super().__init__(base_url, **kw) self._subs: List[Dict[str, Any]] = [] def _lifecycle(self, name: str, **extra: Any) -> Dict[str, Any]: return {"ch": name, **extra} @staticmethod def _op(op: str, symbol: str, channel: str, interval: Optional[str]) -> Dict[str, Any]: msg: Dict[str, Any] = {"op": op, "symbol": symbol} if channel != "ticker": msg["ch"] = channel if interval is not None: msg["interval"] = interval return msg async def subscribe(self, symbol: str, channel: str = "ticker", interval: Optional[str] = None) -> None: """channel: ticker, kline (needs interval), depth or trade. Kept across reconnects.""" msg = self._op("sub", symbol, channel, interval) if msg not in self._subs: self._subs.append(msg) if self.connected: await self.send(msg) async def unsubscribe(self, symbol: str, channel: str = "ticker", interval: Optional[str] = None) -> None: sub = self._op("sub", symbol, channel, interval) if sub in self._subs: self._subs.remove(sub) if self.connected: await self.send(self._op("unsub", symbol, channel, interval)) async def _on_open(self) -> None: for msg in self._subs: await self.send(msg) class HistoryStream(_Stream): """Your perpetual and tokenized-stock events: order, trade, pnl, position, balance, funding. CFD events come on CfdStream. Frames carry {epoch, seq}. When frames may have been missed (a gap, a new epoch, a counter reset, or a reconnect whose hello is ahead), a synthetic {"kind": "resync", "reason": ...} is yielded: re-read positions, open orders, balances and recent history over REST. async with HistoryStream(client) as stream: async for event in stream: if event["kind"] == "resync": reload_from_rest() elif event["kind"] == "order": on_order(event["data"]) """ path = "/ws/history" def __init__(self, client: OGuardClient, **kw: Any): super().__init__(client.base_url, **kw) self.client = client self.cursor: Optional[Tuple[str, int]] = None async def _url(self) -> str: ticket = await asyncio.get_running_loop().run_in_executor(None, self.client.ws_ticket) return f"{self.origin}{self.path}?ticket={quote(ticket, safe='')}" def _on_frame(self, msg: Dict[str, Any]) -> None: kind, epoch, seq = msg.get("kind"), msg.get("epoch"), msg.get("seq") if isinstance(epoch, str) and epoch and isinstance(seq, int) and not isinstance(seq, bool) and kind != "error": reason = self._advance(epoch, seq, hello=kind == "hello") if reason: self._emit({"kind": "resync", "reason": reason}) self._emit(msg) def _advance(self, epoch: str, seq: int, hello: bool) -> Optional[str]: prev = self.cursor if prev is None: self.cursor = (epoch, seq) return None if prev[0] != epoch: self.cursor = (epoch, seq) return "epoch-change" if hello: if seq > prev[1]: self.cursor = (epoch, seq) return "missed" return None if seq == prev[1] + 1: self.cursor = (epoch, seq) return None if seq > prev[1] + 1: self.cursor = (epoch, seq) return "gap" if seq == 1 and prev[1] > 1: self.cursor = (epoch, seq) return "reset" return None # late or duplicate frame: harmless class CfdStream(_Stream): """Live CFD quotes on request, and what the book does to your orders. Frames are {"kind": ..., "data": ...}: a "hello" on every connection, "quote" batches for the symbols you subscribe to, and the book's own events: "order" (a pending order filled or was retired), "protection" (a take-profit, stop-loss or trailing stop fired), "stop_out", each followed by "account" with the new balance and margin. Your own placements, closes and changes answer in their REST response and are not echoed here. Frames are not numbered: after a reconnect (an {"kind": "open"} event that is not the first), re-read og.cfd.positions() and og.cfd.pending_orders(). async with CfdStream(client) as stream: await stream.subscribe_quotes(["EURUSD", "XAUUSD"]) async for event in stream: print(event["kind"], event.get("data")) """ path = "/ws/cfd" def __init__(self, client: OGuardClient, **kw: Any): super().__init__(client.base_url, **kw) self.client = client self._quote_symbols: List[str] = [] async def _url(self) -> str: ticket = await asyncio.get_running_loop().run_in_executor(None, self.client.ws_ticket) return f"{self.origin}{self.path}?ticket={quote(ticket, safe='')}" async def subscribe_quotes(self, symbols: Iterable[str]) -> None: """Receive quotes for exactly these symbols (replaces the previous set). Kept across reconnects.""" self._quote_symbols = list(symbols) if self.connected: await self.send({"op": "subscribe", "symbols": self._quote_symbols}) async def _on_open(self) -> None: if self._quote_symbols: await self.send({"op": "subscribe", "symbols": self._quote_symbols}) # ── Command line: python oguard_client.py doctor | time | keygen ────────────── def _doctor(base_url: str) -> int: key, secret = os.environ.get("OGUARD_API_KEY"), os.environ.get("OGUARD_API_SECRET") pem = os.environ.get("OGUARD_PRIVATE_KEY_FILE") print(f"API: {base_url}") og = OGuardClient(key, secret, private_key=open(pem).read() if pem else None, base_url=base_url) offset = og.sync_time() print(f"clock: yours is {-offset / 1000:+.3f} s from the server (signing uses the server's)") print(f"markets: {len(og.perps.markets().get('symbols', []))} perpetuals, " f"{len(og.stocks.markets().get('symbols', []))} tokenized stocks listed") if not key: print("key: none (set OGUARD_API_KEY and OGUARD_API_SECRET to check one)") return 0 try: bal = og.balances() except OGuardAPIError as err: print(f"key: refused, {err.status} {err.code}") return 1 print(f"key: accepted; {len(bal.get('balances', []))} balance row(s)") print(f"budget: {og.last_rate_limit}") return 0 def main(argv: List[str]) -> int: base = os.environ.get("OGUARD_API_URL", DEFAULT_BASE_URL) cmd = argv[1] if len(argv) > 1 else "doctor" if cmd == "time": print(OGuardClient(base_url=base).server_time()) return 0 if cmd == "keygen": pem, pub = generate_ed25519_keypair() print(pem + "\nPublic key (paste this on the API keys page):\n" + pub) return 0 if cmd == "doctor": return _doctor(base) print("usage: python oguard_client.py [doctor | time | keygen]") return 2 if __name__ == "__main__": sys.exit(main(sys.argv))