Files

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()