Skip to content
Merged

Dev #25

Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
160 changes: 158 additions & 2 deletions backends/native_event.py
Original file line number Diff line number Diff line change
Expand Up @@ -486,6 +486,7 @@ def run_stat_arb_pair_arbitrage(
close_dict = align_series(closes, symbols, idx)
contract_sizes = self._contract_size_for_spec(spec, contract_size)
fee_rates = self._fee_rate_for_spec(spec)
stat_funding = self._funding_for_spec(spec, funding_rate)
rebalance_threshold = spec.hedge_policy.rebalance_threshold
if not spec.hedge_policy.freeze_on_entry and rebalance_threshold is None:
rebalance_threshold = 0.0
Expand Down Expand Up @@ -523,18 +524,34 @@ def run_stat_arb_pair_arbitrage(
closes=close_dict,
highs=highs,
lows=lows,
funding_rate=funding_rate,
funding_rate=stat_funding,
contract_size=contract_sizes,
leverage=leverage,
fee_rate=fee_rates,
symbols=symbols,
)
funding_dict = prepare_funding(stat_funding if self.config.use_funding else 0.0, symbols, idx)
roles = self._stat_arb_roles(spec)
leg_pnl_report = self._leg_pnl_report(
idx=idx,
symbols=symbols,
roles=roles,
result=result,
closes=close_dict,
funding=funding_dict,
contract_sizes=contract_sizes,
)
package_report = self._package_pnl_report(idx, result, leg_pnl_report)
beta_drift_report = self._stat_arb_beta_drift_report(
idx=idx,
spec=spec,
plan=arb_plan,
rebalance_threshold=rebalance_threshold,
)
diagnostics = result.diagnostics.copy()
diagnostics["package_pnl"] = package_report["package_pnl"]
diagnostics["package_pnl_residual"] = package_report["pnl_residual"]
result.diagnostics = diagnostics
result.metadata.update(
{
"backend": "native_event",
Expand All @@ -547,6 +564,9 @@ def run_stat_arb_pair_arbitrage(
"basket_plan": plan,
"basket_target_units": arb_plan.target_units,
"beta_drift_report": beta_drift_report,
"spread_report": self._stat_arb_spread_report(idx, spec, close_dict, arb_plan),
"leg_pnl_report": leg_pnl_report,
"package_pnl_report": package_report,
"rebalance_threshold": rebalance_threshold,
"fee_rate_oneway": fee_rates,
"contract_size": contract_sizes,
Expand Down Expand Up @@ -671,7 +691,11 @@ def run_package_arbitrage(
"""
unsupported = (CrossExchangeArbSpec, TriangularArbSpec, OptionsVolArbSpec)
if isinstance(spec, unsupported):
raise NotImplementedError(f"{type(spec).__name__} requires a specialized Phase G+ engine")
raise NotImplementedError(
f"{type(spec).__name__} is schema-validated but requires a specialized arbitrage engine; "
"do not route it through generic package execution. "
"Use QuantBTEndpoint.arbitrage_support_matrix() to inspect supported routes."
)
supported = (CalendarSpreadSpec, FundingArbitrageSpec, SpotPerpCashCarrySpec, IndexBasketArbSpec)
if not isinstance(spec, supported):
raise TypeError("run_package_arbitrage requires a Phase G package-style arbitrage spec")
Expand Down Expand Up @@ -1048,6 +1072,15 @@ def _stat_arb_basket_from_spec(spec: StatArbPairSpec) -> BasketSpec:
},
)

@staticmethod
def _stat_arb_roles(spec: StatArbPairSpec) -> Dict[str, str]:
symbols = [leg.symbol for leg in spec.legs]
roles = {leg.symbol: str(leg.role or "leg") for leg in spec.legs}
if len(symbols) >= 2 and len(set(roles.values())) == 1:
roles[symbols[0]] = "leg"
roles[symbols[1]] = "hedge"
return roles

@staticmethod
def _stat_arb_beta_drift_report(
idx: pd.DatetimeIndex,
Expand Down Expand Up @@ -1095,6 +1128,129 @@ def _stat_arb_beta_drift_report(
)
return pd.DataFrame(rows)

@staticmethod
def _stat_arb_spread_report(
idx: pd.DatetimeIndex,
spec: StatArbPairSpec,
closes: Dict[str, pd.Series],
plan,
) -> pd.DataFrame:
symbols = [leg.symbol for leg in spec.legs]
leg_symbol = symbols[0]
hedge_symbol = symbols[1] if len(symbols) > 1 else symbols[0]
leg_close = closes[leg_symbol].astype(float)
hedge_close = closes[hedge_symbol].astype(float)
ref_ratio = plan.entry_ratios[leg_symbol].replace(0.0, np.nan).astype(float)
hedge_ratio = (plan.entry_ratios[hedge_symbol].astype(float) / ref_ratio).fillna(0.0)
spread = leg_close + hedge_ratio * hedge_close
return pd.DataFrame(
{
"leg_symbol": leg_symbol,
"hedge_symbol": hedge_symbol,
"leg_close": leg_close,
"hedge_close": hedge_close,
"hedge_ratio_to_leg": hedge_ratio,
"spread": spread,
"abs_spread": spread.abs(),
},
index=idx,
)

def _leg_pnl_report(
self,
idx: pd.DatetimeIndex,
symbols: List[str],
roles: Dict[str, str],
result: BacktestResultV2,
closes: Dict[str, pd.Series],
funding: Dict[str, pd.Series],
contract_sizes: Dict[str, float],
) -> pd.DataFrame:
fill_rows = {}
for fill in result.fills:
ts = pd.Timestamp(fill.timestamp)
if ts.tz is None:
ts = ts.tz_localize("UTC")
else:
ts = ts.tz_convert("UTC")
key = (ts, fill.symbol)
fee, fill_pnl = fill_rows.get(key, (0.0, 0.0))
close_price = float(closes[fill.symbol].loc[ts])
cs = float(contract_sizes[fill.symbol])
fill_pnl += fill.signed_qty * (close_price - float(fill.price)) * cs
fee += float(fill.fee)
fill_rows[key] = (fee, fill_pnl)

funding_mask = make_funding_mask(idx)
cumulative = {symbol: 0.0 for symbol in symbols}
rows = []
for i, ts in enumerate(idx):
for symbol in symbols:
cs = float(contract_sizes[symbol])
close_price = float(closes[symbol].iloc[i])
prev_units = 0.0 if i == 0 else float(result.positions[f"Position_{symbol}"].iloc[i - 1])
units = float(result.positions[f"Position_{symbol}"].iloc[i])
price_pnl = 0.0
if i > 0:
price_pnl = prev_units * (close_price - float(closes[symbol].iloc[i - 1])) * cs
funding_cost = 0.0
if self.config.use_funding and funding_mask[i]:
funding_cost = prev_units * close_price * cs * float(funding[symbol].iloc[i])
fee, fill_pnl = fill_rows.get((ts, symbol), (0.0, 0.0))
total_pnl = price_pnl + fill_pnl - fee - funding_cost
cumulative[symbol] += total_pnl
rows.append(
{
"timestamp": ts,
"symbol": symbol,
"role": roles.get(symbol, "leg"),
"units": units,
"close": close_price,
"notional": abs(units) * close_price * cs,
"price_pnl": price_pnl,
"fill_pnl": fill_pnl,
"fee": fee,
"funding_pnl": -funding_cost,
"total_pnl": total_pnl,
"cumulative_pnl": cumulative[symbol],
}
)
return pd.DataFrame(rows)

@staticmethod
def _package_pnl_report(idx: pd.DatetimeIndex, result: BacktestResultV2, leg_pnl_report: pd.DataFrame) -> pd.DataFrame:
grouped = leg_pnl_report.groupby("timestamp", sort=False)
package_pnl = grouped["total_pnl"].sum().reindex(idx, fill_value=0.0)
price_pnl = grouped["price_pnl"].sum().reindex(idx, fill_value=0.0)
fill_pnl = grouped["fill_pnl"].sum().reindex(idx, fill_value=0.0)
fees = grouped["fee"].sum().reindex(idx, fill_value=0.0)
funding_pnl = grouped["funding_pnl"].sum().reindex(idx, fill_value=0.0)
role_pnl = leg_pnl_report.pivot_table(
index="timestamp",
columns="role",
values="total_pnl",
aggfunc="sum",
fill_value=0.0,
).reindex(idx, fill_value=0.0)
leg_pnl = role_pnl["leg"] if "leg" in role_pnl else pd.Series(0.0, index=idx)
hedge_pnl = role_pnl["hedge"] if "hedge" in role_pnl else pd.Series(0.0, index=idx)
report = pd.DataFrame(
{
"price_pnl": price_pnl,
"fill_pnl": fill_pnl,
"fees": fees,
"funding_pnl": funding_pnl,
"leg_pnl": leg_pnl,
"hedge_pnl": hedge_pnl,
"spread_pnl": leg_pnl + hedge_pnl,
"package_pnl": package_pnl,
"equity_delta": result.equity.diff().fillna(0.0),
},
index=idx,
)
report["pnl_residual"] = report["equity_delta"] - report["package_pnl"]
return report

def _basis_leg_pnl_report(
self,
idx: pd.DatetimeIndex,
Expand Down
82 changes: 49 additions & 33 deletions backends/native_portfolio.py
Original file line number Diff line number Diff line change
Expand Up @@ -464,7 +464,8 @@ def _build_result(
exposure_report = self._build_exposure_report(
accepted_notional_arr=accepted_notional_arr,
target_notional_arr=target_notional_arr,
equity=equity,
equity_arr=equity_arr,
idx=idx,
leverages=leverages,
maintenance_ratio=maintenance_ratio,
betas=betas,
Expand All @@ -482,29 +483,41 @@ def _build_result(
fee_arr=fee_arr,
)
rebalance_report = self._build_rebalance_report(
target_units=target_units_report,
accepted_units=accepted_units_report,
closes=close_report,
contract_sizes=cs,
idx=idx,
symbols=symbol_list,
target_units_arr=target_m,
accepted_units_arr=pos_arr,
closes_arr=closes_m,
contract_sizes=contract_sizes,
)

positions = pd.DataFrame(pos_arr, index=idx, columns=[f"Position_{s}" for s in symbol_list], copy=False)
closes = pd.DataFrame(closes_m, index=idx, columns=[f"Close_{s}" for s in symbol_list], copy=False)
fees = pd.Series(fee_arr, index=idx, name="fees")
turnover = pd.Series(turnover_arr, index=idx, name="turnover")
funding_cost = symbol_pnl_report.groupby("timestamp", sort=False)["funding_cost"].sum().reindex(idx, fill_value=0.0)
prev_units = np.vstack([np.zeros((1, len(symbol_list)), dtype=np.float64), pos_arr[:-1]])
funding_cost_arr = prev_units * closes_m * cs_row * funding_m
funding_cost_arr = np.where(is_funding_bar.reshape(-1, 1).astype(bool), funding_cost_arr, 0.0).sum(axis=1)
margin = exposure_report[["initial_margin", "maintenance_margin"]].copy()
diagnostics = pd.DataFrame(
{
"turnover": turnover,
"rejected_rebalances": (target_units_report - accepted_units_report).abs().sum(axis=1) > 1e-10,
"turnover": turnover_arr,
"rejected_rebalances": np.abs(target_m - pos_arr).sum(axis=1) > 1e-10,
},
index=idx,
)
returns_arr = np.zeros_like(equity_arr, dtype=np.float64)
if len(equity_arr) > 1:
returns_arr[1:] = np.divide(
equity_arr[1:] - equity_arr[:-1],
equity_arr[:-1],
out=np.zeros(len(equity_arr) - 1, dtype=np.float64),
where=equity_arr[:-1] != 0.0,
)

return BacktestResultV2(
equity=equity,
returns=equity.pct_change().fillna(0.0),
returns=pd.Series(returns_arr, index=idx, name="returns"),
positions=positions,
closes=closes,
symbols=symbol_list,
Expand All @@ -513,7 +526,7 @@ def _build_result(
liquidated=liquidated,
liquidation_bar=liquidation_bar,
fees=fees,
funding=pd.Series(funding_cost.to_numpy(dtype=float), index=idx, name="funding"),
funding=pd.Series(funding_cost_arr, index=idx, name="funding"),
margin=margin,
diagnostics=diagnostics,
metadata={
Expand Down Expand Up @@ -598,13 +611,13 @@ def _build_exposure_report(
*,
accepted_notional_arr: np.ndarray,
target_notional_arr: np.ndarray,
equity: pd.Series,
equity_arr: np.ndarray,
idx: pd.DatetimeIndex,
leverages: np.ndarray,
maintenance_ratio: float,
betas: np.ndarray,
) -> pd.DataFrame:
abs_accepted = np.abs(accepted_notional_arr)
equity_arr = equity.to_numpy(dtype=np.float64)
gross = abs_accepted.sum(axis=1)
net = accepted_notional_arr.sum(axis=1)
initial_margin = (abs_accepted / leverages.reshape(1, -1)).sum(axis=1)
Expand All @@ -613,7 +626,9 @@ def _build_exposure_report(
target_gross = np.abs(target_notional_arr).sum(axis=1)
target_beta_exposure = (target_notional_arr * betas.reshape(1, -1)).sum(axis=1)
mean_leverage = float(np.mean(leverages))
out = pd.DataFrame(
gross_leverage = np.divide(gross, equity_arr, out=np.zeros_like(gross), where=equity_arr != 0.0)
net_exposure_pct = np.divide(net, equity_arr, out=np.zeros_like(net), where=equity_arr != 0.0)
return pd.DataFrame(
{
"long_notional": np.where(accepted_notional_arr > 0.0, accepted_notional_arr, 0.0).sum(axis=1),
"short_notional": np.where(accepted_notional_arr < 0.0, -accepted_notional_arr, 0.0).sum(axis=1),
Expand All @@ -627,38 +642,39 @@ def _build_exposure_report(
"equity": equity_arr,
"available_equity_after_im": equity_arr - initial_margin,
"buying_power": equity_arr * mean_leverage,
"gross_leverage": gross_leverage,
"net_exposure_pct": net_exposure_pct,
},
index=equity.index,
index=idx,
)
out["gross_leverage"] = out["gross_notional"] / out["equity"].replace(0.0, np.nan)
out["net_exposure_pct"] = out["net_notional"] / out["equity"].replace(0.0, np.nan)
return out.fillna(0.0)

@staticmethod
def _build_rebalance_report(
*,
target_units: pd.DataFrame,
accepted_units: pd.DataFrame,
closes: pd.DataFrame,
contract_sizes: pd.Series,
idx: pd.DatetimeIndex,
symbols: List[str],
target_units_arr: np.ndarray,
accepted_units_arr: np.ndarray,
closes_arr: np.ndarray,
contract_sizes: np.ndarray,
) -> pd.DataFrame:
diff = target_units - accepted_units
mask = diff.abs() > 1e-10
if not mask.to_numpy().any():
diff = target_units_arr - accepted_units_arr
row_idx, col_idx = np.nonzero(np.abs(diff) > 1e-10)
if len(row_idx) == 0:
return pd.DataFrame(
columns=["timestamp", "symbol", "target_units", "accepted_units", "unit_diff", "notional_diff", "reason"]
)
notional_diff = diff.mul(closes, axis=0).mul(contract_sizes, axis=1)
stacked = diff.where(mask).stack(future_stack=True).dropna()
index = stacked.index
unit_diff = diff[row_idx, col_idx]
notional_diff = unit_diff * closes_arr[row_idx, col_idx] * contract_sizes[col_idx]
symbol_arr = np.asarray(symbols, dtype=object)
return pd.DataFrame(
{
"timestamp": index.get_level_values(0),
"symbol": index.get_level_values(1),
"target_units": target_units.stack(future_stack=True).reindex(index).to_numpy(dtype=float),
"accepted_units": accepted_units.stack(future_stack=True).reindex(index).to_numpy(dtype=float),
"unit_diff": stacked.to_numpy(dtype=float),
"notional_diff": notional_diff.stack(future_stack=True).reindex(index).to_numpy(dtype=float),
"timestamp": idx.take(row_idx),
"symbol": symbol_arr[col_idx],
"target_units": target_units_arr[row_idx, col_idx],
"accepted_units": accepted_units_arr[row_idx, col_idx],
"unit_diff": unit_diff,
"notional_diff": notional_diff,
"reason": "margin_or_portfolio_gate",
}
)
Expand Down
Loading
Loading