1358 lines
52 KiB
Python
1358 lines
52 KiB
Python
"""Round book — execution trail DB executor (psycopg).
|
|
|
|
The "round" is the execution projection of a scheduled/manual algo run:
|
|
|
|
scheduler_runs ──► ROUND ──► rd_experiments
|
|
|
|
Schema owned by tac-app's Drizzle migrations (`trading_round_book`, migration
|
|
0005) and mirrored here as ``CREATE TABLE IF NOT EXISTS`` so the agent can
|
|
write the same rows before the app has ever booted against a fresh database
|
|
(the ``tac-qlib-custom`` trace_db.py precedent). Four layers:
|
|
|
|
fact_events append-only evidence collected during the run (signals,
|
|
quotes, news, account/position state) — never mutated.
|
|
round_intents versioned target portfolios; a newer version supersedes
|
|
the previous one for the round.
|
|
round_decisions per-symbol actions decided by the gates (placed OR
|
|
deliberately skipped, each with a reason).
|
|
round_orders execution state of placed decisions (Alpaca order + fills).
|
|
|
|
Env: DATABASE_URL (Postgres). APCA_API_KEY_ID + APCA_API_SECRET_KEY +
|
|
APCA_API_BASE_URL are needed only for the Alpaca fill sync.
|
|
|
|
Usage::
|
|
|
|
.venv/bin/python -m tac_qlib.book_db init # ensure round tables exist
|
|
.venv/bin/python -m tac_qlib.book_db ls # recent rounds
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import argparse
|
|
import datetime as _dt
|
|
import decimal
|
|
import json
|
|
import os
|
|
import sys
|
|
from typing import Any, Dict, Iterable, List, Optional, Sequence
|
|
|
|
import psycopg
|
|
from psycopg.rows import dict_row
|
|
|
|
ROUND_STATUSES = ("open", "locked", "settled", "aborted")
|
|
DECISION_STATUSES = (
|
|
"intended",
|
|
"placed",
|
|
"accepted",
|
|
"filled",
|
|
"partially_filled",
|
|
"cancelled",
|
|
"superseded",
|
|
"skipped",
|
|
"rejected",
|
|
)
|
|
ORDER_STATUSES = ("accepted", "filled", "partially_filled", "cancelled", "superseded")
|
|
|
|
|
|
# --------------------------------------------------------------------------- env / connection
|
|
|
|
def _load_repo_env() -> None:
|
|
"""Load the repo-root `.env` (like rd_server.py) so DATABASE_URL + Alpaca
|
|
creds resolve whether the module runs from the repo, a wheel, or a container.
|
|
The repo `.env` is authoritative for any var that is currently *unset or
|
|
empty* — opencode spawns MCP servers with ``DATABASE_URL=""`` when the
|
|
``{env:...}`` substitution finds nothing in its own env, and a plain
|
|
``load_dotenv(override=False)`` would leave that empty value in place."""
|
|
from dotenv import dotenv_values
|
|
|
|
here = os.path.dirname(os.path.abspath(__file__))
|
|
for root in (here, os.getcwd()):
|
|
candidate = os.path.join(root, ".env")
|
|
if not os.path.isfile(candidate):
|
|
continue
|
|
for key, value in dotenv_values(candidate).items():
|
|
if value is None:
|
|
continue
|
|
if not (os.environ.get(key) or "").strip():
|
|
os.environ[key] = value
|
|
return
|
|
|
|
|
|
def _db_url() -> str:
|
|
url = (os.environ.get("DATABASE_URL") or "").strip()
|
|
if not url:
|
|
# Lazy env load: lets scripts use the library directly without first
|
|
# calling `_load_repo_env()` (the MCP server loads it at import).
|
|
_load_repo_env()
|
|
url = (os.environ.get("DATABASE_URL") or "").strip()
|
|
if not url:
|
|
raise RuntimeError("DATABASE_URL is not set")
|
|
return url
|
|
|
|
|
|
def connect() -> psycopg.Connection:
|
|
return psycopg.connect(_db_url(), row_factory=dict_row)
|
|
|
|
|
|
def _jsonable(value: Any) -> Any:
|
|
"""Deep-convert Decimals/dates/datetimes (psycopg dict_row output) to JSON-safe types."""
|
|
if isinstance(value, decimal.Decimal):
|
|
return float(value)
|
|
if isinstance(value, (_dt.date, _dt.datetime)):
|
|
return value.isoformat()
|
|
if isinstance(value, dict):
|
|
return {k: _jsonable(v) for k, v in value.items()}
|
|
if isinstance(value, (list, tuple)):
|
|
return [_jsonable(v) for v in value]
|
|
return value
|
|
|
|
|
|
# Evidence gate: locking / settling a round requires the evidence kinds that let
|
|
# the fact table answer *why the targets were chosen* and *what the market looked
|
|
# like*. Per-symbol evidence (symbol_features, decision_justification) is only
|
|
# required when the round actually has an intent targeting that symbol.
|
|
def _assert_evidence_gate(
|
|
round_id: int,
|
|
status: str,
|
|
*,
|
|
conn: Optional[psycopg.Connection] = None,
|
|
) -> None:
|
|
required_once = ("signal_score", "market_snapshot", "strategy_config", "account_state")
|
|
per_target = ("symbol_features", "decision_justification")
|
|
own = conn is None
|
|
c = conn or connect()
|
|
try:
|
|
with c.cursor() as cur:
|
|
cur.execute("SELECT kind, symbol FROM fact_events WHERE round_id = %s", (round_id,))
|
|
rows = cur.fetchall()
|
|
cur.execute(
|
|
"SELECT target_portfolio FROM round_intents WHERE round_id = %s ORDER BY version DESC LIMIT 1",
|
|
(round_id,),
|
|
)
|
|
active_intent = cur.fetchone()
|
|
finally:
|
|
if own:
|
|
c.close()
|
|
kinds = {r["kind"] for r in rows}
|
|
missing = [k for k in required_once if k not in kinds]
|
|
if active_intent:
|
|
portfolio = active_intent["target_portfolio"]
|
|
if isinstance(portfolio, str):
|
|
try:
|
|
portfolio = json.loads(portfolio)
|
|
except ValueError:
|
|
portfolio = []
|
|
target_symbols = {
|
|
t["symbol"] for t in portfolio or [] if isinstance(t, dict) and t.get("symbol")
|
|
}
|
|
per_sym = {s: set() for s in target_symbols}
|
|
for r in rows:
|
|
if r["kind"] in per_target and r["symbol"] in per_sym:
|
|
per_sym[r["symbol"]].add(r["kind"])
|
|
for s, got in sorted(per_sym.items()):
|
|
for k in per_target:
|
|
if k not in got:
|
|
missing.append(f"{k} for {s}")
|
|
if missing:
|
|
raise ValueError(
|
|
f"cannot {status} round {round_id}: evidence gate not met — missing fact kinds: {', '.join(missing)}. "
|
|
"Record them with fact_record (e.g. get_lake_ta/get_lake_sp persist:true + get_lake_features + rd_exp_model "
|
|
"for symbol_features, rd_exp_model tree path for decision_justification), then retry."
|
|
)
|
|
|
|
|
|
# --------------------------------------------------------------------------- DDL (mirrors migration 0005)
|
|
|
|
SCHEMA_DDL: List[str] = [
|
|
"""
|
|
CREATE TABLE IF NOT EXISTS trading_rounds (
|
|
id bigserial PRIMARY KEY NOT NULL,
|
|
source text DEFAULT 'scheduled' NOT NULL,
|
|
target_date date NOT NULL,
|
|
signal_date date,
|
|
scheduler_run_id bigint,
|
|
rd_experiment_id bigint,
|
|
experiment_name text,
|
|
run_id text,
|
|
model_path text,
|
|
strategy_snapshot jsonb,
|
|
account_equity_at_sizing numeric,
|
|
status text DEFAULT 'open' NOT NULL,
|
|
locked_intent_id bigint,
|
|
summary_metrics jsonb,
|
|
feedback_note text,
|
|
created_at timestamp with time zone DEFAULT now() NOT NULL,
|
|
updated_at timestamp with time zone DEFAULT now() NOT NULL
|
|
)
|
|
""",
|
|
"""
|
|
CREATE TABLE IF NOT EXISTS fact_events (
|
|
id bigserial PRIMARY KEY NOT NULL,
|
|
round_id bigint NOT NULL,
|
|
kind text NOT NULL,
|
|
symbol text,
|
|
payload jsonb,
|
|
source text,
|
|
at timestamp with time zone DEFAULT now() NOT NULL
|
|
)
|
|
""",
|
|
"""
|
|
CREATE TABLE IF NOT EXISTS round_intents (
|
|
id bigserial PRIMARY KEY NOT NULL,
|
|
round_id bigint NOT NULL,
|
|
version bigint NOT NULL,
|
|
supersedes_intent_id bigint,
|
|
target_portfolio jsonb,
|
|
raw_strategy_output jsonb,
|
|
reason text,
|
|
created_at timestamp with time zone DEFAULT now() NOT NULL
|
|
)
|
|
""",
|
|
"""
|
|
CREATE TABLE IF NOT EXISTS round_decisions (
|
|
id bigserial PRIMARY KEY NOT NULL,
|
|
round_id bigint NOT NULL,
|
|
intent_id bigint,
|
|
symbol text NOT NULL,
|
|
side text NOT NULL,
|
|
qty numeric,
|
|
order_type text,
|
|
expected_price numeric,
|
|
status text DEFAULT 'intended' NOT NULL,
|
|
reason text,
|
|
reason_detail text,
|
|
superseded_by_decision_id bigint,
|
|
created_at timestamp with time zone DEFAULT now() NOT NULL,
|
|
updated_at timestamp with time zone DEFAULT now() NOT NULL
|
|
)
|
|
""",
|
|
"""
|
|
CREATE TABLE IF NOT EXISTS round_orders (
|
|
id bigserial PRIMARY KEY NOT NULL,
|
|
decision_id bigint NOT NULL,
|
|
round_id bigint NOT NULL,
|
|
alpaca_order_id text,
|
|
client_order_id text,
|
|
qty_intended numeric,
|
|
qty_filled numeric DEFAULT '0' NOT NULL,
|
|
avg_fill_price numeric,
|
|
status text DEFAULT 'accepted' NOT NULL,
|
|
superseded_by_order_id bigint,
|
|
created_at timestamp with time zone DEFAULT now() NOT NULL,
|
|
updated_at timestamp with time zone DEFAULT now() NOT NULL
|
|
)
|
|
""",
|
|
]
|
|
|
|
SCHEMA_INDEXES: List[str] = [
|
|
"CREATE INDEX IF NOT EXISTS trading_rounds_target_date_idx ON trading_rounds (target_date)",
|
|
"CREATE INDEX IF NOT EXISTS trading_rounds_scheduler_run_id_idx ON trading_rounds (scheduler_run_id)",
|
|
"CREATE INDEX IF NOT EXISTS trading_rounds_rd_experiment_id_idx ON trading_rounds (rd_experiment_id)",
|
|
"CREATE INDEX IF NOT EXISTS trading_rounds_locked_intent_id_idx ON trading_rounds (locked_intent_id)",
|
|
"CREATE INDEX IF NOT EXISTS fact_events_round_id_idx ON fact_events (round_id)",
|
|
"CREATE INDEX IF NOT EXISTS fact_events_round_kind_idx ON fact_events (round_id, kind)",
|
|
"CREATE INDEX IF NOT EXISTS round_intents_round_id_idx ON round_intents (round_id)",
|
|
"CREATE INDEX IF NOT EXISTS round_intents_round_version_idx ON round_intents (round_id, version)",
|
|
"CREATE INDEX IF NOT EXISTS round_decisions_round_id_idx ON round_decisions (round_id)",
|
|
"CREATE INDEX IF NOT EXISTS round_decisions_round_symbol_idx ON round_decisions (round_id, symbol)",
|
|
"CREATE INDEX IF NOT EXISTS round_decisions_intent_id_idx ON round_decisions (intent_id)",
|
|
"CREATE INDEX IF NOT EXISTS round_orders_round_id_idx ON round_orders (round_id)",
|
|
"CREATE INDEX IF NOT EXISTS round_orders_decision_id_idx ON round_orders (decision_id)",
|
|
"CREATE INDEX IF NOT EXISTS round_orders_alpaca_order_id_idx ON round_orders (alpaca_order_id)",
|
|
]
|
|
|
|
|
|
def ensure_schema(conn: Optional[psycopg.Connection] = None) -> None:
|
|
"""Create the round-book tables + indexes (idempotent). Owns its own
|
|
connection when none is passed so the CLI / server can call it freely."""
|
|
own = conn is None
|
|
c = conn or connect()
|
|
try:
|
|
with c.cursor() as cur:
|
|
for ddl in SCHEMA_DDL:
|
|
cur.execute(ddl)
|
|
for ddl in SCHEMA_INDEXES:
|
|
cur.execute(ddl)
|
|
c.commit()
|
|
finally:
|
|
if own:
|
|
c.close()
|
|
|
|
|
|
# --------------------------------------------------------------------------- rounds
|
|
|
|
def create_round(
|
|
source: str = "scheduled",
|
|
target_date: Optional[str] = None,
|
|
signal_date: Optional[str] = None,
|
|
scheduler_run_id: Optional[int] = None,
|
|
rd_experiment_id: Optional[int] = None,
|
|
experiment_name: Optional[str] = None,
|
|
run_id: Optional[str] = None,
|
|
model_path: Optional[str] = None,
|
|
strategy_snapshot: Optional[Dict[str, Any]] = None,
|
|
account_equity_at_sizing: Optional[float] = None,
|
|
*,
|
|
reopen: bool = True,
|
|
conn: Optional[psycopg.Connection] = None,
|
|
) -> Dict[str, Any]:
|
|
"""Open a round window for a target trading date.
|
|
|
|
``reopen=True`` (default) makes this idempotent per (source, target_date):
|
|
if an open round already exists for the window it is returned unchanged
|
|
(``reused=True``) instead of creating a duplicate window. A settled window
|
|
is never reopened — a fresh round is created.
|
|
"""
|
|
if not target_date:
|
|
raise ValueError("target_date is required (YYYY-MM-DD, the trading day being executed)")
|
|
source = source or "scheduled"
|
|
own = conn is None
|
|
c = conn or connect()
|
|
try:
|
|
if reopen:
|
|
row = c.execute(
|
|
"SELECT * FROM trading_rounds WHERE source = %s AND target_date = %s AND status = 'open' "
|
|
"ORDER BY id DESC LIMIT 1",
|
|
(source, target_date),
|
|
).fetchone()
|
|
if row:
|
|
return {**_jsonable(row), "reused": True}
|
|
|
|
with c.cursor() as cur:
|
|
cur.execute(
|
|
"""
|
|
INSERT INTO trading_rounds
|
|
(source, target_date, signal_date, scheduler_run_id, rd_experiment_id,
|
|
experiment_name, run_id, model_path, strategy_snapshot, account_equity_at_sizing)
|
|
VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
|
|
RETURNING *
|
|
""",
|
|
(
|
|
source,
|
|
target_date,
|
|
signal_date,
|
|
scheduler_run_id,
|
|
rd_experiment_id,
|
|
experiment_name,
|
|
run_id,
|
|
model_path,
|
|
json.dumps(strategy_snapshot) if strategy_snapshot is not None else None,
|
|
account_equity_at_sizing,
|
|
),
|
|
)
|
|
row = cur.fetchone()
|
|
c.commit()
|
|
return {**_jsonable(row), "reused": False}
|
|
finally:
|
|
if own:
|
|
c.close()
|
|
|
|
|
|
def update_round_status(
|
|
round_id: int,
|
|
status: Optional[str] = None,
|
|
locked_intent_id: Optional[int] = None,
|
|
summary_metrics: Optional[Dict[str, Any]] = None,
|
|
feedback_note: Optional[str] = None,
|
|
*,
|
|
conn: Optional[psycopg.Connection] = None,
|
|
) -> Dict[str, Any]:
|
|
"""Advance a round (open → locked → settled | aborted). ``locked_intent_id``
|
|
pins the intent the reconciliation should reconcile against."""
|
|
if status is not None and status not in ROUND_STATUSES:
|
|
raise ValueError(f"invalid round status {status!r}; expected one of {ROUND_STATUSES}")
|
|
if status in ("locked", "settled"):
|
|
_assert_evidence_gate(round_id, status, conn=conn)
|
|
fields = []
|
|
params: List[Any] = []
|
|
if status is not None:
|
|
fields.append("status = %s")
|
|
params.append(status)
|
|
if locked_intent_id is not None:
|
|
fields.append("locked_intent_id = %s")
|
|
params.append(locked_intent_id)
|
|
if summary_metrics is not None:
|
|
fields.append("summary_metrics = %s")
|
|
params.append(json.dumps(summary_metrics))
|
|
if feedback_note is not None:
|
|
fields.append("feedback_note = %s")
|
|
params.append(feedback_note)
|
|
fields.append("updated_at = now()")
|
|
params.append(round_id)
|
|
own = conn is None
|
|
c = conn or connect()
|
|
try:
|
|
row = c.execute(
|
|
f"UPDATE trading_rounds SET {', '.join(fields)} WHERE id = %s RETURNING *",
|
|
params,
|
|
).fetchone()
|
|
c.commit()
|
|
if not row:
|
|
raise ValueError(f"round {round_id} not found")
|
|
return _jsonable(row)
|
|
finally:
|
|
if own:
|
|
c.close()
|
|
|
|
|
|
def update_round(
|
|
round_id: int,
|
|
*,
|
|
source: Optional[str] = None,
|
|
target_date: Optional[str] = None,
|
|
signal_date: Optional[str] = None,
|
|
scheduler_run_id: Optional[int] = None,
|
|
rd_experiment_id: Optional[int] = None,
|
|
experiment_name: Optional[str] = None,
|
|
run_id: Optional[str] = None,
|
|
model_path: Optional[str] = None,
|
|
strategy_snapshot: Optional[Dict[str, Any]] = None,
|
|
account_equity_at_sizing: Optional[float] = None,
|
|
conn: Optional[psycopg.Connection] = None,
|
|
) -> Dict[str, Any]:
|
|
"""Update metadata fields of a round window (e.g. fill in the new training
|
|
run / model path after the retrain completes)."""
|
|
fields, params = [], []
|
|
if source is not None:
|
|
fields.append("source = %s")
|
|
params.append(source)
|
|
if target_date is not None:
|
|
fields.append("target_date = %s")
|
|
params.append(target_date)
|
|
if signal_date is not None:
|
|
fields.append("signal_date = %s")
|
|
params.append(signal_date)
|
|
if scheduler_run_id is not None:
|
|
fields.append("scheduler_run_id = %s")
|
|
params.append(scheduler_run_id)
|
|
if rd_experiment_id is not None:
|
|
fields.append("rd_experiment_id = %s")
|
|
params.append(rd_experiment_id)
|
|
if experiment_name is not None:
|
|
fields.append("experiment_name = %s")
|
|
params.append(experiment_name)
|
|
if run_id is not None:
|
|
fields.append("run_id = %s")
|
|
params.append(run_id)
|
|
if model_path is not None:
|
|
fields.append("model_path = %s")
|
|
params.append(model_path)
|
|
if strategy_snapshot is not None:
|
|
fields.append("strategy_snapshot = %s")
|
|
params.append(json.dumps(strategy_snapshot))
|
|
if account_equity_at_sizing is not None:
|
|
fields.append("account_equity_at_sizing = %s")
|
|
params.append(account_equity_at_sizing)
|
|
if not fields:
|
|
return get_round(round_id, conn=conn) or {}
|
|
fields.append("updated_at = now()")
|
|
params.append(round_id)
|
|
own = conn is None
|
|
c = conn or connect()
|
|
try:
|
|
row = c.execute(
|
|
f"UPDATE trading_rounds SET {', '.join(fields)} WHERE id = %s RETURNING *",
|
|
params,
|
|
).fetchone()
|
|
c.commit()
|
|
if not row:
|
|
raise ValueError(f"round {round_id} not found")
|
|
return _jsonable(row)
|
|
finally:
|
|
if own:
|
|
c.close()
|
|
|
|
|
|
def get_round(round_id: int, *, conn: Optional[psycopg.Connection] = None) -> Optional[Dict[str, Any]]:
|
|
c = conn or connect()
|
|
try:
|
|
row = c.execute("SELECT * FROM trading_rounds WHERE id = %s", (round_id,)).fetchone()
|
|
return _jsonable(row) if row else None
|
|
finally:
|
|
if conn is None:
|
|
c.close()
|
|
|
|
|
|
def list_rounds(
|
|
source: Optional[str] = None,
|
|
target_date: Optional[str] = None,
|
|
status: Optional[str] = None,
|
|
limit: int = 20,
|
|
*,
|
|
with_detail: bool = False,
|
|
conn: Optional[psycopg.Connection] = None,
|
|
) -> List[Dict[str, Any]]:
|
|
clauses = []
|
|
params: List[Any] = []
|
|
if source:
|
|
clauses.append("source = %s")
|
|
params.append(source)
|
|
if target_date:
|
|
clauses.append("target_date = %s")
|
|
params.append(target_date)
|
|
if status:
|
|
clauses.append("status = %s")
|
|
params.append(status)
|
|
where = f"WHERE {' AND '.join(clauses)}" if clauses else ""
|
|
params.append(max(1, min(int(limit), 200)))
|
|
own = conn is None
|
|
c = conn or connect()
|
|
try:
|
|
rows = c.execute(
|
|
f"SELECT * FROM trading_rounds {where} ORDER BY target_date DESC, id DESC LIMIT %s",
|
|
params,
|
|
).fetchall()
|
|
rounds = [_jsonable(r) for r in rows]
|
|
if with_detail:
|
|
for r in rounds:
|
|
try:
|
|
r["funnel"] = funnel(r["id"], conn=c)
|
|
r["metrics"] = metrics(r["id"], conn=c)
|
|
except Exception: # noqa: BLE001 - keep the window even if a roll-up fails
|
|
pass
|
|
return rounds
|
|
finally:
|
|
if own:
|
|
c.close()
|
|
|
|
|
|
# --------------------------------------------------------------------------- fact_events
|
|
|
|
def record_fact(
|
|
round_id: int,
|
|
kind: str,
|
|
payload: Optional[Dict[str, Any]] = None,
|
|
symbol: Optional[str] = None,
|
|
source: Optional[str] = None,
|
|
*,
|
|
conn: Optional[psycopg.Connection] = None,
|
|
) -> Dict[str, Any]:
|
|
"""Append an evidence event to the round (never mutated, only appended)."""
|
|
if not kind:
|
|
raise ValueError("kind is required (signal_score | quote | news_sentiment | account_state | position_state | strategy_config | ...)")
|
|
own = conn is None
|
|
c = conn or connect()
|
|
try:
|
|
with c.cursor() as cur:
|
|
cur.execute(
|
|
"INSERT INTO fact_events (round_id, kind, symbol, payload, source) VALUES (%s, %s, %s, %s, %s) RETURNING *",
|
|
(round_id, kind, symbol, json.dumps(payload) if payload is not None else None, source),
|
|
)
|
|
row = cur.fetchone()
|
|
c.commit()
|
|
return _jsonable(row)
|
|
finally:
|
|
if own:
|
|
c.close()
|
|
|
|
|
|
def query_facts(
|
|
round_id: int,
|
|
kind: Optional[str] = None,
|
|
symbol: Optional[str] = None,
|
|
limit: int = 200,
|
|
*,
|
|
conn: Optional[psycopg.Connection] = None,
|
|
) -> List[Dict[str, Any]]:
|
|
clauses = ["round_id = %s"]
|
|
params: List[Any] = [round_id]
|
|
if kind:
|
|
clauses.append("kind = %s")
|
|
params.append(kind)
|
|
if symbol:
|
|
clauses.append("symbol = %s")
|
|
params.append(symbol)
|
|
params.append(max(1, min(int(limit), 1000)))
|
|
c = conn or connect()
|
|
try:
|
|
rows = c.execute(
|
|
f"SELECT * FROM fact_events WHERE {' AND '.join(clauses)} ORDER BY id DESC LIMIT %s",
|
|
params,
|
|
).fetchall()
|
|
return [_jsonable(r) for r in rows]
|
|
finally:
|
|
if conn is None:
|
|
c.close()
|
|
|
|
|
|
# --------------------------------------------------------------------------- round_intents
|
|
|
|
def set_intent(
|
|
round_id: int,
|
|
target_portfolio: Sequence[Dict[str, Any]],
|
|
raw_strategy_output: Optional[Dict[str, Any]] = None,
|
|
reason: Optional[str] = None,
|
|
*,
|
|
conn: Optional[psycopg.Connection] = None,
|
|
) -> Dict[str, Any]:
|
|
"""Write the next target-portfolio version for a round. Version numbers
|
|
increment per round; a new version supersedes the previous active one
|
|
(``supersedes_intent_id``), and the round's ``updated_at`` bumps so list
|
|
ordering reflects the latest intent."""
|
|
if not isinstance(target_portfolio, (list, tuple)) or not target_portfolio:
|
|
raise ValueError("target_portfolio must be a non-empty list of {symbol, side, ...} rows")
|
|
own = conn is None
|
|
c = conn or connect()
|
|
try:
|
|
prev = c.execute(
|
|
"SELECT id, version FROM round_intents WHERE round_id = %s ORDER BY version DESC LIMIT 1",
|
|
(round_id,),
|
|
).fetchone()
|
|
version = int((prev or {}).get("version", 0)) + 1
|
|
with c.cursor() as cur:
|
|
cur.execute(
|
|
"""
|
|
INSERT INTO round_intents
|
|
(round_id, version, supersedes_intent_id, target_portfolio, raw_strategy_output, reason)
|
|
VALUES (%s, %s, %s, %s, %s, %s) RETURNING *
|
|
""",
|
|
(
|
|
round_id,
|
|
version,
|
|
prev["id"] if prev else None,
|
|
json.dumps(list(target_portfolio)),
|
|
json.dumps(raw_strategy_output) if raw_strategy_output is not None else None,
|
|
reason,
|
|
),
|
|
)
|
|
row = cur.fetchone()
|
|
cur.execute("UPDATE trading_rounds SET updated_at = now() WHERE id = %s", (round_id,))
|
|
c.commit()
|
|
return _jsonable(row)
|
|
finally:
|
|
if own:
|
|
c.close()
|
|
|
|
|
|
def get_intent(round_id: int, version: Optional[int] = None, *, conn: Optional[psycopg.Connection] = None) -> Optional[Dict[str, Any]]:
|
|
c = conn or connect()
|
|
try:
|
|
if version is not None:
|
|
row = c.execute(
|
|
"SELECT * FROM round_intents WHERE round_id = %s AND version = %s", (round_id, version)
|
|
).fetchone()
|
|
else:
|
|
row = c.execute(
|
|
"SELECT * FROM round_intents WHERE round_id = %s ORDER BY version DESC LIMIT 1", (round_id,)
|
|
).fetchone()
|
|
return _jsonable(row) if row else None
|
|
finally:
|
|
if conn is None:
|
|
c.close()
|
|
|
|
|
|
def list_intents(round_id: int, *, conn: Optional[psycopg.Connection] = None) -> List[Dict[str, Any]]:
|
|
c = conn or connect()
|
|
try:
|
|
rows = c.execute(
|
|
"SELECT * FROM round_intents WHERE round_id = %s ORDER BY version ASC", (round_id,)
|
|
).fetchall()
|
|
return [_jsonable(r) for r in rows]
|
|
finally:
|
|
if conn is None:
|
|
c.close()
|
|
|
|
|
|
# --------------------------------------------------------------------------- round_decisions / round_orders
|
|
|
|
def record_decision(
|
|
round_id: int,
|
|
symbol: str,
|
|
side: str,
|
|
qty: Optional[float] = None,
|
|
order_type: Optional[str] = None,
|
|
expected_price: Optional[float] = None,
|
|
status: str = "intended",
|
|
reason: Optional[str] = None,
|
|
reason_detail: Optional[str] = None,
|
|
intent_id: Optional[int] = None,
|
|
supersedes_decision_id: Optional[int] = None,
|
|
alpaca_order_id: Optional[str] = None,
|
|
client_order_id: Optional[str] = None,
|
|
*,
|
|
conn: Optional[psycopg.Connection] = None,
|
|
) -> Dict[str, Any]:
|
|
"""Record one per-symbol decision by the gates — either a placed order
|
|
(status ``placed``/``accepted``/...) or a deliberate skip (status
|
|
``skipped``/``rejected`` with a ``reason``). When ``alpaca_order_id`` is
|
|
given and the decision is not a skip, a ``round_orders`` execution row is
|
|
created alongside."""
|
|
if status not in DECISION_STATUSES:
|
|
raise ValueError(f"invalid decision status {status!r}; expected one of {DECISION_STATUSES}")
|
|
if side not in ("buy", "sell"):
|
|
raise ValueError(f"side must be 'buy' or 'sell', got {side!r}")
|
|
own = conn is None
|
|
c = conn or connect()
|
|
try:
|
|
with c.cursor() as cur:
|
|
if supersedes_decision_id is not None:
|
|
cur.execute(
|
|
"UPDATE round_decisions SET status = 'superseded', superseded_by_decision_id = %s, "
|
|
"updated_at = now() WHERE id = %s",
|
|
(supersedes_decision_id, supersedes_decision_id),
|
|
)
|
|
cur.execute(
|
|
"""
|
|
INSERT INTO round_decisions
|
|
(round_id, intent_id, symbol, side, qty, order_type, expected_price, status,
|
|
reason, reason_detail, superseded_by_decision_id)
|
|
VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
|
|
RETURNING *
|
|
""",
|
|
(
|
|
round_id,
|
|
intent_id,
|
|
symbol,
|
|
side,
|
|
qty,
|
|
order_type,
|
|
expected_price,
|
|
status,
|
|
reason,
|
|
reason_detail,
|
|
supersedes_decision_id,
|
|
),
|
|
)
|
|
decision = cur.fetchone()
|
|
decision_id = decision["id"]
|
|
order = None
|
|
if alpaca_order_id is not None and status not in ("skipped", "rejected", "superseded"):
|
|
cur.execute(
|
|
"""
|
|
INSERT INTO round_orders
|
|
(decision_id, round_id, alpaca_order_id, client_order_id, qty_intended, status)
|
|
VALUES (%s, %s, %s, %s, %s, %s) RETURNING *
|
|
""",
|
|
(decision_id, round_id, alpaca_order_id, client_order_id, qty, "accepted"),
|
|
)
|
|
order = cur.fetchone()
|
|
cur.execute("UPDATE trading_rounds SET updated_at = now() WHERE id = %s", (round_id,))
|
|
c.commit()
|
|
result = _jsonable(decision)
|
|
result["order"] = _jsonable(order) if order else None
|
|
return result
|
|
finally:
|
|
if own:
|
|
c.close()
|
|
|
|
|
|
def link_order(
|
|
round_id: int,
|
|
decision_id: int,
|
|
alpaca_order_id: Optional[str] = None,
|
|
client_order_id: Optional[str] = None,
|
|
qty_filled: Optional[float] = None,
|
|
avg_fill_price: Optional[float] = None,
|
|
status: Optional[str] = None,
|
|
*,
|
|
conn: Optional[psycopg.Connection] = None,
|
|
) -> Dict[str, Any]:
|
|
"""Create or update the execution row for a placed decision (idempotent per
|
|
alpaca_order_id / client_order_id / decision_id)."""
|
|
own = conn is None
|
|
c = conn or connect()
|
|
try:
|
|
existing = c.execute(
|
|
"SELECT * FROM round_orders WHERE decision_id = %s AND round_id = %s ORDER BY id DESC LIMIT 1",
|
|
(decision_id, round_id),
|
|
).fetchone()
|
|
with c.cursor() as cur:
|
|
if existing is not None:
|
|
fields, params = [], []
|
|
if alpaca_order_id is not None:
|
|
fields.append("alpaca_order_id = %s")
|
|
params.append(alpaca_order_id)
|
|
if client_order_id is not None:
|
|
fields.append("client_order_id = %s")
|
|
params.append(client_order_id)
|
|
if qty_filled is not None:
|
|
fields.append("qty_filled = %s")
|
|
params.append(qty_filled)
|
|
if avg_fill_price is not None:
|
|
fields.append("avg_fill_price = %s")
|
|
params.append(avg_fill_price)
|
|
if status is not None:
|
|
if status not in ORDER_STATUSES:
|
|
raise ValueError(f"invalid order status {status!r}; expected one of {ORDER_STATUSES}")
|
|
fields.append("status = %s")
|
|
params.append(status)
|
|
fields.append("updated_at = now()")
|
|
params.extend([decision_id, round_id])
|
|
cur.execute(
|
|
f"UPDATE round_orders SET {', '.join(fields)} WHERE decision_id = %s AND round_id = %s RETURNING *",
|
|
params,
|
|
)
|
|
row = cur.fetchone()
|
|
else:
|
|
if status is not None and status not in ORDER_STATUSES:
|
|
raise ValueError(f"invalid order status {status!r}; expected one of {ORDER_STATUSES}")
|
|
cur.execute(
|
|
"""
|
|
INSERT INTO round_orders
|
|
(decision_id, round_id, alpaca_order_id, client_order_id, qty_intended,
|
|
qty_filled, avg_fill_price, status)
|
|
VALUES (%s, %s, %s, %s, %s, %s, %s, %s) RETURNING *
|
|
""",
|
|
(
|
|
decision_id,
|
|
round_id,
|
|
alpaca_order_id,
|
|
client_order_id,
|
|
None,
|
|
qty_filled if qty_filled is not None else 0,
|
|
avg_fill_price,
|
|
status or "accepted",
|
|
),
|
|
)
|
|
row = cur.fetchone()
|
|
cur.execute("UPDATE trading_rounds SET updated_at = now() WHERE id = %s", (round_id,))
|
|
c.commit()
|
|
return _jsonable(row)
|
|
finally:
|
|
if own:
|
|
c.close()
|
|
|
|
|
|
def query_decisions(
|
|
round_id: int,
|
|
symbol: Optional[str] = None,
|
|
include_superseded: bool = True,
|
|
*,
|
|
conn: Optional[psycopg.Connection] = None,
|
|
) -> List[Dict[str, Any]]:
|
|
clauses = ["round_id = %s"]
|
|
params: List[Any] = [round_id]
|
|
if symbol:
|
|
clauses.append("symbol = %s")
|
|
params.append(symbol)
|
|
if not include_superseded:
|
|
clauses.append("status <> 'superseded'")
|
|
c = conn or connect()
|
|
try:
|
|
rows = c.execute(
|
|
f"SELECT * FROM round_decisions WHERE {' AND '.join(clauses)} ORDER BY id ASC",
|
|
params,
|
|
).fetchall()
|
|
return [_jsonable(r) for r in rows]
|
|
finally:
|
|
if conn is None:
|
|
c.close()
|
|
|
|
|
|
def list_orders(
|
|
round_id: int,
|
|
decision_id: Optional[int] = None,
|
|
*,
|
|
conn: Optional[psycopg.Connection] = None,
|
|
) -> List[Dict[str, Any]]:
|
|
"""Execution rows for a round (optionally one decision)."""
|
|
clauses = ["round_id = %s"]
|
|
params: List[Any] = [round_id]
|
|
if decision_id is not None:
|
|
clauses.append("decision_id = %s")
|
|
params.append(decision_id)
|
|
c = conn or connect()
|
|
try:
|
|
rows = c.execute(
|
|
f"SELECT * FROM round_orders WHERE {' AND '.join(clauses)} ORDER BY id ASC",
|
|
params,
|
|
).fetchall()
|
|
return [_jsonable(r) for r in rows]
|
|
finally:
|
|
if conn is None:
|
|
c.close()
|
|
|
|
|
|
def _orders_for_round(round_id: int, conn: psycopg.Connection) -> List[Dict[str, Any]]:
|
|
rows = conn.execute(
|
|
"SELECT * FROM round_orders WHERE round_id = %s ORDER BY id ASC", (round_id,)
|
|
).fetchall()
|
|
return [_jsonable(r) for r in rows]
|
|
|
|
|
|
# --------------------------------------------------------------------------- Alpaca fill sync
|
|
|
|
_ALPACA_STATUS_MAP = {
|
|
"filled": "filled",
|
|
"partially_filled": "partially_filled",
|
|
"canceled": "cancelled",
|
|
"cancelled": "cancelled",
|
|
"expired": "cancelled",
|
|
"replaced": "superseded",
|
|
"superseded": "superseded",
|
|
}
|
|
|
|
|
|
def _alpaca_orders_from_api(*, feed: str = "iex") -> List[Dict[str, Any]]:
|
|
"""Fetch recent closed+open orders from the Alpaca REST API (paper by
|
|
default). Returns raw order dicts (already JSON-safe)."""
|
|
import httpx
|
|
|
|
key_id = (os.environ.get("APCA_API_KEY_ID") or "").strip()
|
|
secret = (os.environ.get("APCA_API_SECRET_KEY") or "").strip()
|
|
base = (os.environ.get("APCA_API_BASE_URL") or "https://paper-api.alpaca.markets").strip()
|
|
if not key_id or not secret:
|
|
raise RuntimeError(
|
|
"APCA_API_KEY_ID / APCA_API_SECRET_KEY not set — pass `orders` explicitly instead of fetching from Alpaca"
|
|
)
|
|
headers = {"APCA-API-KEY-ID": key_id, "APCA-API-SECRET-KEY": secret}
|
|
url = f"{base.rstrip('/')}/v2/orders"
|
|
resp = httpx.get(
|
|
url,
|
|
params={"status": "all", "limit": 500, "sort": "desc", "nested": "true", "feed": feed},
|
|
headers=headers,
|
|
timeout=30,
|
|
)
|
|
resp.raise_for_status()
|
|
return resp.json()
|
|
|
|
|
|
def _normalize_orders(orders: Iterable[Dict[str, Any]]) -> List[Dict[str, Any]]:
|
|
"""Normalize Alpaca order dicts (from the API or passed via the `orders`
|
|
param, e.g. the tac-engine list_orders output) to a single shape."""
|
|
out = []
|
|
for o in orders:
|
|
out.append(
|
|
{
|
|
"alpaca_order_id": o.get("id") or o.get("alpaca_order_id"),
|
|
"client_order_id": o.get("client_order_id"),
|
|
"symbol": o.get("symbol"),
|
|
"side": o.get("side"),
|
|
"qty": o.get("qty") or o.get("qty_intended"),
|
|
"filled_qty": o.get("filled_qty"),
|
|
"filled_avg_price": o.get("filled_avg_price"),
|
|
"status": o.get("status"),
|
|
}
|
|
)
|
|
return out
|
|
|
|
|
|
def sync_fills(
|
|
round_id: int,
|
|
orders: Optional[Sequence[Dict[str, Any]]] = None,
|
|
feed: str = "iex",
|
|
*,
|
|
conn: Optional[psycopg.Connection] = None,
|
|
) -> Dict[str, Any]:
|
|
"""Pull Alpaca order state into the round's execution rows and compute the
|
|
supersession delta against the effective (locked or latest) intent.
|
|
|
|
- ``orders``: list of Alpaca order dicts (tac-engine ``list_orders`` output
|
|
or the REST shape). When omitted, fetched from Alpaca directly using
|
|
APCA_API_* env vars.
|
|
- Orders match a ``round_orders`` row by alpaca_order_id, then
|
|
client_order_id. Orders that match nothing are returned as ``unmatched``
|
|
(they may belong to another round or predate the round book).
|
|
- A placed order whose symbol+side is no longer in the effective intent is
|
|
marked ``superseded`` (unless already filled); its decision is superseded
|
|
too.
|
|
"""
|
|
own = conn is None
|
|
c = conn or connect()
|
|
try:
|
|
if orders is None:
|
|
orders = _normalize_orders(_alpaca_orders_from_api(feed=feed))
|
|
else:
|
|
orders = _normalize_orders(orders)
|
|
|
|
existing = {r["alpaca_order_id"]: r for r in _orders_for_round(round_id, c) if r["alpaca_order_id"]}
|
|
existing_client = {
|
|
r["client_order_id"]: r for r in _orders_for_round(round_id, c) if r["client_order_id"] and r["client_order_id"] not in existing
|
|
}
|
|
|
|
intent = get_intent(round_id, conn=c)
|
|
target_set = set()
|
|
if intent and intent.get("target_portfolio"):
|
|
target_set = {(t.get("symbol"), t.get("side")) for t in intent["target_portfolio"] if t.get("symbol")}
|
|
|
|
updated, unmatched, superseded = [], [], []
|
|
with c.cursor() as cur:
|
|
for o in orders:
|
|
row = existing.get(o["alpaca_order_id"]) or existing_client.get(o["client_order_id"])
|
|
if row is None:
|
|
unmatched.append(o)
|
|
continue
|
|
status = _ALPACA_STATUS_MAP.get((o.get("status") or "").lower(), "accepted")
|
|
filled_qty = o.get("filled_qty")
|
|
avg_price = o.get("filled_avg_price")
|
|
cur.execute(
|
|
"""
|
|
UPDATE round_orders SET status = %s, qty_filled = %s, avg_fill_price = %s, updated_at = now()
|
|
WHERE id = %s RETURNING *
|
|
""",
|
|
(status, filled_qty if filled_qty is not None else row.get("qty_filled", 0), avg_price, row["id"]),
|
|
)
|
|
ord_row = cur.fetchone()
|
|
updated.append(_jsonable(ord_row))
|
|
|
|
# Supersession: the effective intent no longer wants this symbol+side.
|
|
key = (o.get("symbol"), o.get("side"))
|
|
if target_set and key not in target_set and status in ("accepted", "partially_filled"):
|
|
cur.execute(
|
|
"UPDATE round_orders SET status = 'superseded', updated_at = now() WHERE id = %s RETURNING *",
|
|
(row["id"],),
|
|
)
|
|
cur.execute(
|
|
"UPDATE round_decisions SET status = 'superseded', updated_at = now() WHERE id = %s",
|
|
(row["decision_id"],),
|
|
)
|
|
superseded.append(_jsonable(cur.fetchone()))
|
|
|
|
cur.execute("UPDATE trading_rounds SET updated_at = now() WHERE id = %s", (round_id,))
|
|
c.commit()
|
|
return {
|
|
"round_id": round_id,
|
|
"matched": len(updated),
|
|
"updated": updated,
|
|
"superseded": superseded,
|
|
"unmatched": unmatched,
|
|
}
|
|
finally:
|
|
if own:
|
|
c.close()
|
|
|
|
|
|
# --------------------------------------------------------------------------- reconciliation / trail / metrics
|
|
|
|
def _effective_intent(round_id: int, conn: psycopg.Connection) -> Optional[Dict[str, Any]]:
|
|
round_row = get_round(round_id, conn=conn)
|
|
if not round_row:
|
|
raise ValueError(f"round {round_id} not found")
|
|
locked = round_row.get("locked_intent_id")
|
|
if locked is not None:
|
|
row = conn.execute("SELECT * FROM round_intents WHERE id = %s", (locked,)).fetchone()
|
|
if row:
|
|
return _jsonable(row)
|
|
return get_intent(round_id, conn=conn)
|
|
|
|
|
|
def funnel(round_id: int, *, conn: Optional[psycopg.Connection] = None) -> Dict[str, Any]:
|
|
"""Count the decision funnel: intended targets → decided → placed → filled.
|
|
|
|
Order-derived counts (placed/filled/partial/live/cancelled) come from the
|
|
``round_orders`` execution rows; decision-derived counts (decided/skipped/
|
|
superseded) come from ``round_decisions``."""
|
|
own = conn is None
|
|
c = conn or connect()
|
|
try:
|
|
intent = _effective_intent(round_id, c)
|
|
target_symbols = set()
|
|
if intent and intent.get("target_portfolio"):
|
|
target_symbols = {(t.get("symbol"), t.get("side")) for t in intent["target_portfolio"] if t.get("symbol")}
|
|
decisions = query_decisions(round_id, conn=c)
|
|
orders = _orders_for_round(round_id, c)
|
|
|
|
decided = [d for d in decisions if d["status"] != "superseded"]
|
|
skipped = [d for d in decisions if d["status"] in ("skipped", "rejected")]
|
|
superseded = [d for d in decisions if d["status"] == "superseded"]
|
|
|
|
order_status = {"filled": 0, "partially_filled": 0, "accepted": 0, "cancelled": 0, "superseded": 0}
|
|
for o in orders:
|
|
order_status[o["status"]] = order_status.get(o["status"], 0) + 1
|
|
|
|
skipped_reasons: Dict[str, int] = {}
|
|
for d in skipped:
|
|
skipped_reasons[d.get("reason") or "no_reason"] = skipped_reasons.get(d.get("reason") or "no_reason", 0) + 1
|
|
|
|
return _jsonable(
|
|
{
|
|
"round_id": round_id,
|
|
"intent_version": intent.get("version") if intent else None,
|
|
"target_symbols": len(target_symbols),
|
|
"decided": len(decided),
|
|
"placed": len(orders),
|
|
"filled": order_status["filled"],
|
|
"partially_filled": order_status["partially_filled"],
|
|
"live": order_status["accepted"] + order_status["partially_filled"],
|
|
"cancelled": order_status["cancelled"],
|
|
"skipped": len(skipped),
|
|
"skipped_reasons": skipped_reasons,
|
|
"superseded": len(superseded),
|
|
"superseded_orders": order_status["superseded"],
|
|
"no_decision_targets": len(target_symbols - {(d.get("symbol"), d.get("side")) for d in decided}),
|
|
}
|
|
)
|
|
finally:
|
|
if conn is None:
|
|
c.close()
|
|
|
|
|
|
def reconcile(round_id: int, *, conn: Optional[psycopg.Connection] = None) -> Dict[str, Any]:
|
|
"""Reconcile the effective intent against decisions + fills.
|
|
|
|
Returns a per-symbol residual table (target qty minus filled qty, with the
|
|
reason the target did not fill), plus roll-ups for cash/BP impact, slippage
|
|
and estimated cost.
|
|
"""
|
|
own = conn is None
|
|
c = conn or connect()
|
|
try:
|
|
round_row = get_round(round_id, conn=c)
|
|
if not round_row:
|
|
raise ValueError(f"round {round_id} not found")
|
|
intent = _effective_intent(round_id, c)
|
|
decisions = query_decisions(round_id, conn=c)
|
|
orders = _orders_for_round(round_id, c)
|
|
decision_by = {}
|
|
for d in decisions:
|
|
decision_by.setdefault((d.get("symbol"), d.get("side")), []).append(d)
|
|
order_by_decision = {}
|
|
for o in orders:
|
|
order_by_decision.setdefault(o.get("decision_id"), []).append(o)
|
|
|
|
equity = float(round_row.get("account_equity_at_sizing") or 0) or None
|
|
rows = []
|
|
gross_traded = 0.0
|
|
cash_impact = 0.0
|
|
slippage_bps = 0.0
|
|
slippage_w = 0.0
|
|
cost_est = 0.0
|
|
open_cost = float((round_row.get("strategy_snapshot") or {}).get("open_cost") or 0.0005)
|
|
close_cost = float((round_row.get("strategy_snapshot") or {}).get("close_cost") or 0.0015)
|
|
min_cost = float((round_row.get("strategy_snapshot") or {}).get("min_cost") or 5)
|
|
|
|
for target in intent.get("target_portfolio") or []:
|
|
symbol = target.get("symbol")
|
|
side = target.get("side")
|
|
if not symbol:
|
|
continue
|
|
qty_target = float(target.get("qty") or 0) or None
|
|
cands = decision_by.get((symbol, side), [])
|
|
cand = next((d for d in cands if d["status"] != "superseded"), cands[-1] if cands else None)
|
|
if cand is None:
|
|
rows.append(
|
|
{
|
|
"symbol": symbol,
|
|
"side": side,
|
|
"target_qty": qty_target,
|
|
"filled_qty": 0,
|
|
"residual_qty": qty_target,
|
|
"status": "no_decision",
|
|
"reason": None,
|
|
}
|
|
)
|
|
continue
|
|
order = None
|
|
for o in order_by_decision.get(cand["id"], []):
|
|
if o["status"] != "superseded":
|
|
order = o
|
|
break
|
|
if order is None:
|
|
order = (order_by_decision.get(cand["id"]) or [None])[0]
|
|
filled_qty = float(order.get("qty_filled") or 0) if order else 0.0
|
|
residual = (qty_target - filled_qty) if qty_target is not None else None
|
|
rows.append(
|
|
{
|
|
"symbol": symbol,
|
|
"side": side,
|
|
"target_qty": qty_target,
|
|
"filled_qty": filled_qty,
|
|
"residual_qty": residual,
|
|
"status": (order.get("status") if order else cand.get("status")),
|
|
"reason": cand.get("reason"),
|
|
"expected_price": target.get("expected_price") or cand.get("expected_price"),
|
|
"avg_fill_price": order.get("avg_fill_price") if order else None,
|
|
}
|
|
)
|
|
notional = filled_qty * float(order.get("avg_fill_price") or 0) if order and order.get("avg_fill_price") else 0.0
|
|
gross_traded += notional
|
|
cash_impact += notional if side == "sell" else -notional
|
|
expected = target.get("expected_price") or cand.get("expected_price")
|
|
if order and order.get("avg_fill_price") and expected:
|
|
px = float(order["avg_fill_price"])
|
|
ref = float(expected)
|
|
if ref:
|
|
bp = (px - ref) / ref * 1e4
|
|
slippage_bps += abs(bp) * notional
|
|
slippage_w += notional
|
|
if order:
|
|
q = float(order.get("qty_filled") or 0)
|
|
px = float(order.get("avg_fill_price") or 0)
|
|
val = q * px
|
|
rate = close_cost if side == "sell" else open_cost
|
|
cost_est += max(val * rate, min_cost if q else 0)
|
|
|
|
inv = round_row.get("summary_metrics") or {}
|
|
return _jsonable(
|
|
{
|
|
"round_id": round_id,
|
|
"intent_version": intent.get("version") if intent else None,
|
|
"account_equity_at_sizing": equity,
|
|
"target_symbols": len(intent.get("target_portfolio") or []),
|
|
"rows": rows,
|
|
"rollup": {
|
|
"gross_traded_notional": round(gross_traded, 2),
|
|
"net_cash_impact": round(cash_impact, 2),
|
|
"slippage_bps": round(slippage_bps / slippage_w, 2) if slippage_w else None,
|
|
"estimated_cost": round(cost_est, 2),
|
|
"residual_symbols": sum(1 for r in rows if (r.get("residual_qty") or 0) > 0),
|
|
},
|
|
"summary_metrics": inv,
|
|
}
|
|
)
|
|
finally:
|
|
if conn is None:
|
|
c.close()
|
|
|
|
|
|
def trail(round_id: int, symbol: Optional[str] = None, *, conn: Optional[psycopg.Connection] = None) -> List[Dict[str, Any]]:
|
|
"""Waterfall per symbol: intent target → decision → order → fill.
|
|
|
|
The target row is resolved from the decision's own intent (via
|
|
``intent_id``) so superseded symbols keep their original target context
|
|
even after a newer intent dropped them; the effective (locked/latest)
|
|
intent is the fallback."""
|
|
own = conn is None
|
|
c = conn or connect()
|
|
try:
|
|
intents = list_intents(round_id, conn=c)
|
|
active_intent = _effective_intent(round_id, c)
|
|
by_intent = {}
|
|
intent_version = {}
|
|
for iv in intents:
|
|
index = {}
|
|
for t in iv.get("target_portfolio") or []:
|
|
if t.get("symbol"):
|
|
index.setdefault(t["symbol"], t)
|
|
by_intent[iv["id"]] = index
|
|
intent_version[iv["id"]] = iv.get("version")
|
|
active_index = {}
|
|
if active_intent:
|
|
for t in active_intent.get("target_portfolio") or []:
|
|
if t.get("symbol"):
|
|
active_index[t["symbol"]] = t
|
|
decisions = query_decisions(round_id, symbol=symbol, conn=c)
|
|
orders = _orders_for_round(round_id, c)
|
|
order_by_decision = {}
|
|
for o in orders:
|
|
order_by_decision.setdefault(o.get("decision_id"), []).append(o)
|
|
|
|
rows = []
|
|
for d in decisions:
|
|
sym = d.get("symbol")
|
|
if symbol and sym != symbol:
|
|
continue
|
|
target = (by_intent.get(d.get("intent_id")) or {}).get(sym) or active_index.get(sym) or {}
|
|
rows.append(
|
|
{
|
|
"symbol": sym,
|
|
"side": d.get("side"),
|
|
"decision_id": d.get("id"),
|
|
"decision_status": d.get("status"),
|
|
"reason": d.get("reason"),
|
|
"reason_detail": d.get("reason_detail"),
|
|
"target_qty": target.get("qty"),
|
|
"target_score": target.get("score"),
|
|
"target_rank": target.get("rank"),
|
|
"intent_version": intent_version.get(d.get("intent_id")),
|
|
"decided_qty": d.get("qty"),
|
|
"expected_price": d.get("expected_price"),
|
|
"order_status": None,
|
|
"qty_filled": None,
|
|
"avg_fill_price": None,
|
|
}
|
|
)
|
|
for o in order_by_decision.get(d.get("id"), []):
|
|
rows[-1].update(
|
|
{
|
|
"order_status": o.get("status"),
|
|
"qty_filled": o.get("qty_filled"),
|
|
"avg_fill_price": o.get("avg_fill_price"),
|
|
}
|
|
)
|
|
return _jsonable(rows)
|
|
finally:
|
|
if conn is None:
|
|
c.close()
|
|
|
|
|
|
def metrics(round_id: int, *, conn: Optional[psycopg.Connection] = None) -> Dict[str, Any]:
|
|
"""Roll-up metrics for the round: invested notional, turnover, slippage bps,
|
|
estimated cost and filled count. Safe to call before any fills.
|
|
|
|
Cost model mirrors ``reconcile``: per filled share, ``open_cost`` on buys
|
|
and ``close_cost`` on sells (floored at ``min_cost`` per order), slippage
|
|
measured vs the expected price recorded on the linked decision. The
|
|
``cost_pct_of_gross`` figure is the "is the edge eaten?" number — the ratio
|
|
of total estimated cost to gross traded notional.
|
|
"""
|
|
own = conn is None
|
|
c = conn or connect()
|
|
try:
|
|
round_row = get_round(round_id, conn=c)
|
|
if not round_row:
|
|
raise ValueError(f"round {round_id} not found")
|
|
orders = _orders_for_round(round_id, c)
|
|
filled = [o for o in orders if o["status"] in ("filled", "partially_filled")]
|
|
invested = sum(float(o.get("qty_filled") or 0) * float(o.get("avg_fill_price") or 0) for o in filled)
|
|
equity = float(round_row.get("account_equity_at_sizing") or 0) or 0
|
|
snap = round_row.get("strategy_snapshot") or {}
|
|
open_cost = float(snap.get("open_cost") or 0.0005)
|
|
close_cost = float(snap.get("close_cost") or 0.0015)
|
|
min_cost = float(snap.get("min_cost") or 5)
|
|
|
|
decision_by_id = {d["id"]: d for d in query_decisions(round_id, conn=c)}
|
|
gross_traded = 0.0
|
|
slippage_bps = 0.0
|
|
slippage_w = 0.0
|
|
cost_est = 0.0
|
|
for o in orders:
|
|
q = float(o.get("qty_filled") or 0)
|
|
px = float(o.get("avg_fill_price") or 0)
|
|
if not q or not px:
|
|
continue
|
|
val = q * px
|
|
gross_traded += val
|
|
dec = decision_by_id.get(o.get("decision_id")) or {}
|
|
expected = dec.get("expected_price")
|
|
if expected:
|
|
ref = float(expected)
|
|
if ref:
|
|
bp = (px - ref) / ref * 1e4
|
|
slippage_bps += abs(bp) * val
|
|
slippage_w += val
|
|
side = (dec.get("side") or "buy").lower()
|
|
rate = close_cost if side == "sell" else open_cost
|
|
cost_est += max(val * rate, min_cost if q else 0)
|
|
|
|
return _jsonable(
|
|
{
|
|
"round_id": round_id,
|
|
"placed_orders": len(orders),
|
|
"filled_orders": len([o for o in orders if o["status"] == "filled"]),
|
|
"partially_filled_orders": len([o for o in orders if o["status"] == "partially_filled"]),
|
|
"invested_notional": round(invested, 2),
|
|
"turnover": round(invested / equity, 4) if equity else None,
|
|
"gross_traded_notional": round(gross_traded, 2),
|
|
"slippage_bps": round(slippage_bps / slippage_w, 2) if slippage_w else None,
|
|
"estimated_cost": round(cost_est, 2),
|
|
"cost_pct_of_gross": round(cost_est / gross_traded * 1e4, 2) if gross_traded else None,
|
|
"status": round_row.get("status"),
|
|
}
|
|
)
|
|
finally:
|
|
if conn is None:
|
|
c.close()
|
|
|
|
|
|
# --------------------------------------------------------------------------- CLI
|
|
|
|
def _cli() -> None:
|
|
parser = argparse.ArgumentParser(description="Round book DB utilities")
|
|
parser.add_argument("cmd", choices=("init", "ls"))
|
|
parser.add_argument("--source", default=None)
|
|
parser.add_argument("--limit", type=int, default=10)
|
|
args = parser.parse_args()
|
|
|
|
if args.cmd == "init":
|
|
ensure_schema()
|
|
print("round-book schema ready")
|
|
return
|
|
for r in list_rounds(source=args.source, limit=args.limit):
|
|
print(f"#{r['id']} {r['target_date']} {r['source']} {r['status']} {r.get('experiment_name')} run {str(r.get('run_id'))[:8]}")
|
|
|
|
|
|
if __name__ == "__main__":
|
|
_load_repo_env()
|
|
_cli()
|