1143 lines
41 KiB
Python
1143 lines
41 KiB
Python
'''Run the focused four-arm daily portfolio-capacity research matrix.
|
|
|
|
The expensive point-in-time daily replay and full-universe ranks are cached
|
|
once. Empty-book monthly paths and warm-book weekly seeds are then evaluated
|
|
under cap 10, cap 15, cash-only unbounded, and weekly current-rank top 10.
|
|
'''
|
|
|
|
from __future__ import annotations
|
|
|
|
import argparse
|
|
import asyncio
|
|
import hashlib
|
|
import json
|
|
import multiprocessing
|
|
import os
|
|
import pickle
|
|
import platform
|
|
import subprocess
|
|
import sys
|
|
from collections import Counter
|
|
from concurrent.futures import ProcessPoolExecutor, as_completed
|
|
from datetime import date, datetime, timezone
|
|
from pathlib import Path
|
|
from typing import Any
|
|
|
|
from sqlalchemy import select
|
|
from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker, create_async_engine
|
|
|
|
ROOT = Path(__file__).resolve().parents[1]
|
|
if str(ROOT) not in sys.path:
|
|
sys.path.insert(0, str(ROOT))
|
|
|
|
from scripts.portfolio_capacity_research import ( # noqa: E402
|
|
ANCHOR_YEARS,
|
|
ARM_BY_ID,
|
|
ARMS,
|
|
BOOTSTRAP_REPLICATES,
|
|
BOOTSTRAP_SEED,
|
|
COSTS_PER_SIDE_PCT,
|
|
aggregate_results,
|
|
build_cells,
|
|
build_cohort_manifest,
|
|
median,
|
|
summarize_simulation,
|
|
validate_cohort_manifest,
|
|
)
|
|
from scripts.research_rankings import _live_universe_rank_map # noqa: E402
|
|
|
|
|
|
CACHE_VERSION = 'portfolio-capacity-candidates-v1-zero-horizon'
|
|
RUNNER_VERSION = 'portfolio-capacity-bracket-v2'
|
|
CONSTRUCTION_VIEW_VERSION = 'production-book-filter-v1'
|
|
MIN_PRODUCTION_UNIVERSE = 450
|
|
MAX_PRODUCTION_UNIVERSE = 600
|
|
SPEC_PATH = ROOT / 'docs' / 'research' / 'portfolio-capacity-bracket.md'
|
|
DEFAULT_RUN_ID = 'prod505-capacity-bracket-daily-v1'
|
|
_WORKER_CONTEXT: dict[str, Any] | None = None
|
|
|
|
|
|
def _parse_args() -> argparse.Namespace:
|
|
parser = argparse.ArgumentParser(description=__doc__)
|
|
parser.add_argument('snapshot', help='SQLite backtest snapshot')
|
|
parser.add_argument('--run-id', default=DEFAULT_RUN_ID)
|
|
parser.add_argument(
|
|
'--workers',
|
|
default='auto',
|
|
help='Worker count or auto',
|
|
)
|
|
parser.add_argument('--resume', action='store_true')
|
|
parser.add_argument('--validate-only', action='store_true')
|
|
parser.add_argument('--candidate-cache', default=None)
|
|
parser.add_argument('--checkpoint', default=None)
|
|
parser.add_argument('--out', default=None)
|
|
parser.add_argument('--quiet', action='store_true')
|
|
return parser.parse_args()
|
|
|
|
|
|
def _worker_count(raw: str) -> int:
|
|
if str(raw).lower() == 'auto':
|
|
return max(1, min(6, (multiprocessing.cpu_count() or 2) - 1))
|
|
value = int(raw)
|
|
if value <= 0:
|
|
raise ValueError('--workers must be positive or auto')
|
|
return value
|
|
|
|
|
|
def _sqlite_url(path: Path) -> str:
|
|
return f'sqlite+aiosqlite:///{path.resolve().as_posix()}'
|
|
|
|
|
|
def _sha256_file(path: Path) -> str:
|
|
digest = hashlib.sha256()
|
|
with path.open('rb') as handle:
|
|
for chunk in iter(lambda: handle.read(1024 * 1024), b''):
|
|
digest.update(chunk)
|
|
return digest.hexdigest()
|
|
|
|
|
|
def _json_hash(value: Any) -> str:
|
|
payload = json.dumps(
|
|
value,
|
|
sort_keys=True,
|
|
separators=(',', ':'),
|
|
default=str,
|
|
).encode('utf-8')
|
|
return hashlib.sha256(payload).hexdigest()
|
|
|
|
|
|
def _construction_universe_errors(manifest: dict[str, Any]) -> list[str]:
|
|
errors: list[str] = []
|
|
construction_rows = int(manifest['construction_ticker_rows'])
|
|
construction_priced = int(manifest['construction_symbols_with_prices'])
|
|
ranking_rows = int(manifest['ranking_ticker_rows'])
|
|
rank_only_rows = int(manifest['rank_only_ticker_rows'])
|
|
unknown_rank_only = int(manifest['rank_only_unknown_symbols'])
|
|
if not MIN_PRODUCTION_UNIVERSE <= construction_rows <= MAX_PRODUCTION_UNIVERSE:
|
|
errors.append(
|
|
'construction ticker universe must contain '
|
|
f'{MIN_PRODUCTION_UNIVERSE}-{MAX_PRODUCTION_UNIVERSE} symbols; '
|
|
f'found {construction_rows}'
|
|
)
|
|
if not MIN_PRODUCTION_UNIVERSE <= construction_priced <= MAX_PRODUCTION_UNIVERSE:
|
|
errors.append(
|
|
'priced construction universe must contain '
|
|
f'{MIN_PRODUCTION_UNIVERSE}-{MAX_PRODUCTION_UNIVERSE} symbols; '
|
|
f'found {construction_priced}'
|
|
)
|
|
if construction_rows + rank_only_rows != ranking_rows:
|
|
errors.append(
|
|
'construction and rank-only ticker partitions do not cover the '
|
|
'ranking universe'
|
|
)
|
|
if unknown_rank_only:
|
|
errors.append(
|
|
f'research_rank_only contains {unknown_rank_only} symbols absent '
|
|
'from tickers'
|
|
)
|
|
return errors
|
|
|
|
|
|
def _construction_candidate_view(
|
|
cached: dict[str, Any],
|
|
snapshot_data: dict[str, Any],
|
|
) -> dict[str, Any]:
|
|
'''Filter a reusable full-universe rank cache to production setup symbols.'''
|
|
allowed = set(snapshot_data['construction_symbols'])
|
|
raw_candidates = list(cached['qualified_candidates'])
|
|
qualified = [
|
|
row for row in raw_candidates if str(row['symbol']) in allowed
|
|
]
|
|
view_key = {
|
|
'version': CONSTRUCTION_VIEW_VERSION,
|
|
'base_cache_key_hash': _json_hash(cached['key']),
|
|
'construction_universe_manifest': (
|
|
snapshot_data['construction_universe_manifest']
|
|
),
|
|
}
|
|
return {
|
|
**cached,
|
|
'qualified_candidates': qualified,
|
|
'raw_full_universe_qualified_long_count': len(raw_candidates),
|
|
'qualified_long_count': len(qualified),
|
|
'filtered_rank_only_qualified_long_count': (
|
|
len(raw_candidates) - len(qualified)
|
|
),
|
|
'construction_view_key': view_key,
|
|
}
|
|
|
|
|
|
def _atomic_json(path: Path, value: Any) -> None:
|
|
path.parent.mkdir(parents=True, exist_ok=True)
|
|
temporary = path.with_suffix(path.suffix + '.tmp')
|
|
with temporary.open('w', encoding='utf-8', newline='\n') as handle:
|
|
json.dump(value, handle, indent=2, sort_keys=True, allow_nan=False)
|
|
handle.write('\n')
|
|
temporary.replace(path)
|
|
|
|
|
|
def _atomic_text(path: Path, value: str) -> None:
|
|
path.parent.mkdir(parents=True, exist_ok=True)
|
|
temporary = path.with_suffix(path.suffix + '.tmp')
|
|
with temporary.open('w', encoding='utf-8', newline='\n') as handle:
|
|
handle.write(value.rstrip() + '\n')
|
|
temporary.replace(path)
|
|
|
|
|
|
def _atomic_pickle(path: Path, value: Any) -> None:
|
|
path.parent.mkdir(parents=True, exist_ok=True)
|
|
temporary = path.with_suffix(path.suffix + '.tmp')
|
|
with temporary.open('wb') as handle:
|
|
pickle.dump(value, handle, protocol=pickle.HIGHEST_PROTOCOL)
|
|
temporary.replace(path)
|
|
|
|
|
|
def _git_output(*args: str) -> str:
|
|
completed = subprocess.run(
|
|
['git', *args],
|
|
cwd=ROOT,
|
|
check=True,
|
|
capture_output=True,
|
|
text=True,
|
|
)
|
|
return completed.stdout.strip()
|
|
|
|
|
|
def _assert_clean_worktree() -> None:
|
|
dirty = _git_output('status', '--porcelain')
|
|
if dirty:
|
|
raise SystemExit(
|
|
'Authoritative research refuses a dirty worktree; commit first:\n'
|
|
+ dirty
|
|
)
|
|
|
|
|
|
def _worker_init(context: dict[str, Any]) -> None:
|
|
global _WORKER_CONTEXT
|
|
os.environ['BACKTEST_SNAPSHOT_OFFLINE'] = '1'
|
|
os.environ['BACKTEST_ALLOW_SPAWN'] = '1'
|
|
from app.services import backtest_service as bt
|
|
|
|
context = dict(context)
|
|
context['post_stop_reentry_fn'] = bt._make_gate_reset_reentry_fn(
|
|
context['qualified_candidates'],
|
|
context['prices'],
|
|
cadence='daily',
|
|
ranking_key=context['ranking_key'],
|
|
evaluation_horizon_sessions=0,
|
|
)
|
|
_WORKER_CONTEXT = context
|
|
|
|
|
|
def _worker_run_cell(cell: dict[str, Any]) -> dict[str, Any]:
|
|
if _WORKER_CONTEXT is None:
|
|
raise RuntimeError('Portfolio-capacity worker was not initialized')
|
|
from app.services import backtest_service as bt
|
|
|
|
context = _WORKER_CONTEXT
|
|
arm = ARM_BY_ID[str(cell['arm_id'])]
|
|
measurement_start = date.fromisoformat(str(cell['measurement_start']))
|
|
hard_end = date.fromisoformat(str(cell['hard_end_exclusive']))
|
|
sim = bt._simulate_portfolio(
|
|
context['qualified_candidates'],
|
|
context['prices'],
|
|
context['benchmark_closes'],
|
|
context['exit_policy'],
|
|
context['hold_days'],
|
|
ranking_key=context['ranking_key'],
|
|
max_positions=arm['max_positions'],
|
|
risk_per_trade=context['risk_per_trade'],
|
|
atr_trail_multiplier=context['atr_trail_multiplier'],
|
|
cost_per_side=float(cell['cost_per_side_pct']) / 100.0,
|
|
post_stop_reentry_fn=context['post_stop_reentry_fn'],
|
|
start_date=date.fromisoformat(str(cell['simulation_start'])),
|
|
end_date=hard_end,
|
|
measurement_start_date=measurement_start,
|
|
hard_end_date=hard_end,
|
|
fill_mode=bt.FILL_MODE_CLOSE,
|
|
min_initial_risk_fraction=arm['min_initial_risk_fraction'],
|
|
weekly_top_n_rebalance=bool(arm['weekly_top_n_rebalance']),
|
|
daily_rank_map=(
|
|
context['daily_rank_map']
|
|
if arm['weekly_top_n_rebalance']
|
|
else None
|
|
),
|
|
include_curve=True,
|
|
include_trades=True,
|
|
include_capacity_diagnostics=True,
|
|
)
|
|
if sim is None:
|
|
raise RuntimeError(f'Cell produced no trades: {cell["cell_id"]}')
|
|
return {
|
|
**cell,
|
|
'metrics': summarize_simulation(sim),
|
|
}
|
|
|
|
|
|
async def _fetch_snapshot_columns(
|
|
db: AsyncSession,
|
|
ticker_id: int,
|
|
) -> tuple | None:
|
|
'''Load only the stable OHLCV columns required by the research replay.
|
|
|
|
Research snapshots can predate unrelated additions to the Ticker ORM model
|
|
(for example SEC CIK/SIC metadata). Keeping this query column-scoped avoids
|
|
requiring or mutating those newer application-schema fields.
|
|
'''
|
|
from app.models.ohlcv import OHLCVRecord
|
|
|
|
result = await db.execute(
|
|
select(
|
|
OHLCVRecord.date,
|
|
OHLCVRecord.open,
|
|
OHLCVRecord.high,
|
|
OHLCVRecord.low,
|
|
OHLCVRecord.close,
|
|
OHLCVRecord.volume,
|
|
)
|
|
.where(OHLCVRecord.ticker_id == ticker_id)
|
|
.order_by(OHLCVRecord.date)
|
|
)
|
|
rows = result.all()
|
|
if not rows:
|
|
return None
|
|
return (
|
|
[row[0].toordinal() for row in rows],
|
|
[float(row[1]) for row in rows],
|
|
[float(row[2]) for row in rows],
|
|
[float(row[3]) for row in rows],
|
|
[float(row[4]) for row in rows],
|
|
[int(row[5]) for row in rows],
|
|
)
|
|
|
|
|
|
async def _load_snapshot(
|
|
snapshot: Path,
|
|
*,
|
|
quiet: bool,
|
|
) -> dict[str, Any]:
|
|
from app.models.ticker import Ticker
|
|
from app.services import backtest_service as bt
|
|
from app.services.admin_service import get_activation_config
|
|
from app.services.paper_trade_service import get_exit_policy
|
|
from app.services.recommendation_service import get_recommendation_config
|
|
|
|
engine = create_async_engine(_sqlite_url(snapshot), pool_pre_ping=True)
|
|
Session = async_sessionmaker(
|
|
engine,
|
|
class_=AsyncSession,
|
|
expire_on_commit=False,
|
|
)
|
|
try:
|
|
async with Session() as db:
|
|
recommendation_config = await get_recommendation_config(db)
|
|
activation = await get_activation_config(db)
|
|
exit_config = await get_exit_policy(db)
|
|
benchmark_closes = await bt._load_benchmark_closes_for_backtest(
|
|
db,
|
|
days=None,
|
|
refresh=False,
|
|
)
|
|
ticker_result = await db.execute(
|
|
select(Ticker.id, Ticker.symbol).order_by(Ticker.symbol)
|
|
)
|
|
ticker_rows = [
|
|
(int(ticker_id), str(symbol))
|
|
for ticker_id, symbol in ticker_result.all()
|
|
]
|
|
symbols = [symbol for _ticker_id, symbol in ticker_rows]
|
|
prices: dict[str, tuple] = {}
|
|
for index, (ticker_id, symbol) in enumerate(ticker_rows, 1):
|
|
columns = await _fetch_snapshot_columns(db, ticker_id)
|
|
if columns is not None:
|
|
prices[symbol] = columns
|
|
if not quiet and index % 50 == 0:
|
|
print(
|
|
f'loaded prices: {index}/{len(symbols)}',
|
|
flush=True,
|
|
)
|
|
rank_only_symbols = await bt._load_research_rank_only_symbols(db)
|
|
finally:
|
|
await engine.dispose()
|
|
|
|
if not prices or not benchmark_closes:
|
|
raise SystemExit('Snapshot has no usable prices or benchmark history')
|
|
|
|
ticker_symbols = set(symbols)
|
|
price_symbols = set(prices)
|
|
known_rank_only = ticker_symbols & rank_only_symbols
|
|
unknown_rank_only = rank_only_symbols - ticker_symbols
|
|
construction_symbols = ticker_symbols - known_rank_only
|
|
priced_construction_symbols = price_symbols & construction_symbols
|
|
priced_rank_only_symbols = price_symbols & known_rank_only
|
|
construction_universe_manifest = {
|
|
'ranking_ticker_rows': len(ticker_symbols),
|
|
'ranking_symbols_with_prices': len(price_symbols),
|
|
'ranking_symbols_sha256': _json_hash(sorted(ticker_symbols)),
|
|
'construction_ticker_rows': len(construction_symbols),
|
|
'construction_symbols_with_prices': len(priced_construction_symbols),
|
|
'construction_symbols_sha256': _json_hash(sorted(construction_symbols)),
|
|
'rank_only_ticker_rows': len(known_rank_only),
|
|
'rank_only_symbols_with_prices': len(priced_rank_only_symbols),
|
|
'rank_only_symbols_sha256': _json_hash(sorted(known_rank_only)),
|
|
'rank_only_unknown_symbols': len(unknown_rank_only),
|
|
'rank_only_unknown_symbols_sha256': _json_hash(sorted(unknown_rank_only)),
|
|
}
|
|
|
|
strategy = next(
|
|
row
|
|
for row in bt.PORTFOLIO_MONITOR_STRATEGIES
|
|
if row.get('is_production')
|
|
)
|
|
entry_config = bt._entry_variant_config(str(strategy['entry_variant']))
|
|
if entry_config is None:
|
|
raise RuntimeError('Production entry configuration is missing')
|
|
ranking_key = str(
|
|
entry_config.get('ranking_key') or entry_config['percentile_key']
|
|
)
|
|
exit_policy = bt.LIVE_EXIT_MODE_TO_SIM.get(
|
|
str(exit_config.get('mode', 'atr_trailing')),
|
|
'atr_trail3',
|
|
)
|
|
hold_days = int(exit_config.get('hold_days', 30))
|
|
risk_per_trade = float(entry_config['risk_per_trade'])
|
|
atr_trail_multiplier = float(
|
|
exit_config.get('atr_multiplier', bt.ATR_TRAIL_MULTIPLIER)
|
|
)
|
|
threshold = float(
|
|
activation.get('min_momentum_percentile', 80.0)
|
|
)
|
|
expected = {
|
|
'ranking_key': bt.RESIDUAL_HIGH_VOL_BLEND_80_20_KEY,
|
|
'exit_policy': 'atr_trail3',
|
|
'hold_days': 30,
|
|
'risk_per_trade': 0.01,
|
|
'max_positions': 10,
|
|
'threshold': 80.0,
|
|
}
|
|
realized = {
|
|
'ranking_key': ranking_key,
|
|
'exit_policy': exit_policy,
|
|
'hold_days': hold_days,
|
|
'risk_per_trade': risk_per_trade,
|
|
'max_positions': int(entry_config['max_positions']),
|
|
'threshold': threshold,
|
|
}
|
|
if realized != expected:
|
|
raise SystemExit(
|
|
'Snapshot runtime configuration is not the frozen Phase A control: '
|
|
+ json.dumps({'expected': expected, 'realized': realized}, sort_keys=True)
|
|
)
|
|
|
|
return {
|
|
'recommendation_config': recommendation_config,
|
|
'activation': activation,
|
|
'exit_config': exit_config,
|
|
'benchmark_closes': benchmark_closes,
|
|
'prices': prices,
|
|
'symbols': symbols,
|
|
'construction_symbols': construction_symbols,
|
|
'universe_manifest': {
|
|
'ticker_rows': len(symbols),
|
|
'symbols_with_prices': len(prices),
|
|
'symbols_sha256': _json_hash(sorted(symbols)),
|
|
},
|
|
'construction_universe_manifest': construction_universe_manifest,
|
|
'ranking_key': ranking_key,
|
|
'exit_policy': exit_policy,
|
|
'hold_days': hold_days,
|
|
'risk_per_trade': risk_per_trade,
|
|
'atr_trail_multiplier': atr_trail_multiplier,
|
|
'threshold': threshold,
|
|
'runtime_config': realized,
|
|
}
|
|
|
|
|
|
def _build_candidate_cache(
|
|
snapshot_data: dict[str, Any],
|
|
*,
|
|
snapshot: Path,
|
|
snapshot_sha256: str,
|
|
cache_path: Path,
|
|
workers: int,
|
|
quiet: bool,
|
|
) -> dict[str, Any]:
|
|
from app.services import backtest_service as bt
|
|
|
|
cache_key = {
|
|
'version': CACHE_VERSION,
|
|
'snapshot': str(snapshot.resolve()),
|
|
'snapshot_sha256': snapshot_sha256,
|
|
'cadence': 'daily',
|
|
'outcome_horizon_sessions': 0,
|
|
'recommendation_config_hash': _json_hash(
|
|
snapshot_data['recommendation_config']
|
|
),
|
|
'activation_hash': _json_hash(snapshot_data['activation']),
|
|
'runtime_config': snapshot_data['runtime_config'],
|
|
'universe_manifest': snapshot_data['universe_manifest'],
|
|
}
|
|
if cache_path.exists():
|
|
with cache_path.open('rb') as handle:
|
|
cached = pickle.load(handle) # noqa: S301 - trusted local cache
|
|
if cached.get('key') == cache_key:
|
|
if not quiet:
|
|
print(f'loaded candidate/rank cache: {cache_path}', flush=True)
|
|
return cached
|
|
if not quiet:
|
|
print(f'candidate cache mismatch; rebuilding: {cache_path}', flush=True)
|
|
|
|
replay_rows: list[dict[str, Any]] = []
|
|
prices = snapshot_data['prices']
|
|
replay_args = [
|
|
(
|
|
symbol,
|
|
columns,
|
|
snapshot_data['recommendation_config'],
|
|
snapshot_data['activation'],
|
|
snapshot_data['benchmark_closes'],
|
|
date(1900, 1, 1),
|
|
'daily',
|
|
True,
|
|
True,
|
|
0,
|
|
)
|
|
for symbol, columns in prices.items()
|
|
]
|
|
if workers == 1:
|
|
for index, args in enumerate(replay_args, 1):
|
|
replay_rows.extend(bt._replay_candidates_for_period(*args))
|
|
if not quiet and index % 25 == 0:
|
|
print(
|
|
f'daily replay: {index}/{len(replay_args)} tickers',
|
|
flush=True,
|
|
)
|
|
else:
|
|
context = bt._mp_context() or multiprocessing.get_context('spawn')
|
|
with ProcessPoolExecutor(
|
|
max_workers=workers,
|
|
mp_context=context,
|
|
) as pool:
|
|
futures = [
|
|
pool.submit(bt._replay_candidates_for_period, *args)
|
|
for args in replay_args
|
|
]
|
|
for index, future in enumerate(as_completed(futures), 1):
|
|
replay_rows.extend(future.result())
|
|
if not quiet and index % 25 == 0:
|
|
print(
|
|
f'daily replay: {index}/{len(futures)} tickers',
|
|
flush=True,
|
|
)
|
|
|
|
setup_candidates = [
|
|
row for row in replay_rows if not row.get('_rank_only')
|
|
]
|
|
rank_observations = [
|
|
row for row in replay_rows if row.get('_universe_rank_observation')
|
|
]
|
|
daily_rank_map = _live_universe_rank_map(
|
|
rank_observations,
|
|
snapshot_data['benchmark_closes'],
|
|
bt.STRATEGY_RANK_MOMENTUM_WEIGHT,
|
|
)
|
|
qualified: list[dict[str, Any]] = []
|
|
for setup in setup_candidates:
|
|
if setup.get('direction') != 'long':
|
|
continue
|
|
identity = (str(setup['symbol']), str(setup['date']))
|
|
rank = daily_rank_map.get(identity)
|
|
if rank is None:
|
|
continue
|
|
candidate = {
|
|
key: value
|
|
for key, value in setup.items()
|
|
if not key.startswith('_universe_')
|
|
}
|
|
candidate[bt.PRODUCTION_PERCENTILE_KEY] = rank['momentum_percentile']
|
|
candidate[bt.VOL_PERCENTILE_KEY] = rank['volatility_percentile']
|
|
candidate[bt.RESIDUAL_HIGH_VOL_BLEND_80_20_KEY] = rank['strategy_rank']
|
|
candidate['qualified'] = bt._momentum_qualifies(
|
|
candidate,
|
|
snapshot_data['threshold'],
|
|
)
|
|
if candidate['qualified']:
|
|
qualified.append(candidate)
|
|
|
|
qualified.sort(
|
|
key=lambda row: (
|
|
str(row['date']),
|
|
-float(row.get(snapshot_data['ranking_key']) or 0.0),
|
|
str(row['symbol']),
|
|
float(row.get('entry') or 0.0),
|
|
float(row.get('stop') or 0.0),
|
|
float(row.get('target') or 0.0),
|
|
str(row.get('action') or ''),
|
|
)
|
|
)
|
|
rank_dates = sorted({identity[1] for identity in daily_rank_map})
|
|
cached = {
|
|
'key': cache_key,
|
|
'qualified_candidates': qualified,
|
|
'daily_rank_map': daily_rank_map,
|
|
'setup_candidate_count': len(setup_candidates),
|
|
'qualified_long_count': len(qualified),
|
|
'entry_candidates_by_direction': dict(
|
|
Counter(str(row['direction']) for row in setup_candidates)
|
|
),
|
|
'rank_observation_count': len(rank_observations),
|
|
'rank_first_date': rank_dates[0] if rank_dates else None,
|
|
'rank_last_date': rank_dates[-1] if rank_dates else None,
|
|
}
|
|
_atomic_pickle(cache_path, cached)
|
|
if not quiet:
|
|
print(f'wrote candidate/rank cache: {cache_path}', flush=True)
|
|
return cached
|
|
|
|
|
|
def _checkpoint_state(
|
|
checkpoint_dir: Path,
|
|
fingerprint: str,
|
|
*,
|
|
resume: bool,
|
|
) -> dict[str, dict[str, Any]]:
|
|
manifest_path = checkpoint_dir / 'manifest.json'
|
|
if checkpoint_dir.exists() and not resume:
|
|
existing_cells = list(checkpoint_dir.glob('cell-*.json'))
|
|
if existing_cells:
|
|
raise SystemExit(
|
|
f'Checkpoint cells already exist at {checkpoint_dir}; use --resume '
|
|
'or a new --run-id'
|
|
)
|
|
checkpoint_dir.mkdir(parents=True, exist_ok=True)
|
|
if manifest_path.exists():
|
|
existing = json.loads(manifest_path.read_text(encoding='utf-8'))
|
|
if existing.get('fingerprint') != fingerprint:
|
|
raise SystemExit(
|
|
f'Checkpoint fingerprint mismatch at {checkpoint_dir}'
|
|
)
|
|
else:
|
|
_atomic_json(
|
|
manifest_path,
|
|
{
|
|
'runner_version': RUNNER_VERSION,
|
|
'fingerprint': fingerprint,
|
|
'created_at': datetime.now(timezone.utc).isoformat(),
|
|
},
|
|
)
|
|
|
|
completed: dict[str, dict[str, Any]] = {}
|
|
if resume:
|
|
for path in sorted(checkpoint_dir.glob('cell-*.json')):
|
|
row = json.loads(path.read_text(encoding='utf-8'))
|
|
completed[str(row['cell_id'])] = row
|
|
return completed
|
|
|
|
|
|
def _write_cell_checkpoint(
|
|
checkpoint_dir: Path,
|
|
row: dict[str, Any],
|
|
) -> None:
|
|
name = hashlib.sha256(str(row['cell_id']).encode('utf-8')).hexdigest()
|
|
_atomic_json(checkpoint_dir / f'cell-{name}.json', row)
|
|
|
|
|
|
def _fmt(value: Any, digits: int = 3) -> str:
|
|
if value is None:
|
|
return 'n/a'
|
|
if isinstance(value, float):
|
|
return f'{value:.{digits}f}'
|
|
return str(value)
|
|
|
|
|
|
def _operational_summary(cells: list[dict[str, Any]]) -> list[dict[str, Any]]:
|
|
rows: list[dict[str, Any]] = []
|
|
for arm in ARMS:
|
|
arm_cells = [
|
|
row
|
|
for row in cells
|
|
if row['arm_id'] == arm['id']
|
|
and float(row['cost_per_side_pct']) == 0.1
|
|
]
|
|
weekly = [
|
|
row['metrics']['weekly_rebalance']
|
|
for row in arm_cells
|
|
if row['metrics'].get('weekly_rebalance')
|
|
]
|
|
rows.append({
|
|
'arm_id': arm['id'],
|
|
'paths': len(arm_cells),
|
|
'median_trades': median(
|
|
row['metrics'].get('trades') for row in arm_cells
|
|
),
|
|
'median_blocked_fraction': median(
|
|
row['metrics'].get('blocked_fraction') for row in arm_cells
|
|
),
|
|
'median_avg_positions': median(
|
|
row['metrics'].get('avg_positions') for row in arm_cells
|
|
),
|
|
'peak_positions': max(
|
|
(
|
|
int(row['metrics'].get('peak_positions') or 0)
|
|
for row in arm_cells
|
|
),
|
|
default=0,
|
|
),
|
|
'median_turnover_multiple': median(
|
|
row['metrics'].get('turnover_multiple') for row in arm_cells
|
|
),
|
|
'min_risk_rejections': sum(
|
|
int(row['metrics'].get('skipped_min_initial_risk') or 0)
|
|
for row in arm_cells
|
|
),
|
|
'weekly_zero_entrant_fraction': median(
|
|
item.get('zero_entrant_fraction') for item in weekly
|
|
),
|
|
'weekly_entrant_pool_median': median(
|
|
item.get('entrant_pool_median') for item in weekly
|
|
),
|
|
'weekly_replacements': sum(
|
|
int(item.get('replacements') or 0) for item in weekly
|
|
),
|
|
'weekly_reentries_within_10_sessions': sum(
|
|
int(item.get('reentries_within_10_sessions') or 0)
|
|
for item in weekly
|
|
),
|
|
})
|
|
return rows
|
|
|
|
|
|
def _markdown(report: dict[str, Any]) -> str:
|
|
lines = [
|
|
'# Focused daily portfolio-capacity matrix',
|
|
'',
|
|
f'Generated: {report["generated_at"]}',
|
|
'',
|
|
'## Question',
|
|
'',
|
|
'The current daily Phase A control admitted 472 trades and rejected 519 '
|
|
'qualified opportunities because the ten-slot book was full. This run '
|
|
'brackets the economic cost of that binding constraint; it has no formal '
|
|
'promotion gate.',
|
|
'',
|
|
'> Universe caveat: today\'s production membership is projected backward. '
|
|
'Use paired arm-versus-control differences, not absolute profitability, '
|
|
'for construction conclusions.',
|
|
'',
|
|
'## Paired annual medians',
|
|
'',
|
|
]
|
|
validation = report.get('validation') or {}
|
|
construction = validation.get('construction_universe_manifest') or {}
|
|
coverage = validation.get('candidate_rank_coverage') or {}
|
|
if construction:
|
|
construction_count = construction['construction_symbols_with_prices']
|
|
rank_only_count = construction['rank_only_symbols_with_prices']
|
|
ranking_count = construction['ranking_symbols_with_prices']
|
|
qualified_count = coverage.get('construction_qualified_longs')
|
|
filtered_count = coverage.get('filtered_rank_only_qualified_longs')
|
|
insertion = lines.index('## Paired annual medians')
|
|
lines[insertion:insertion] = [
|
|
'## Validated universes',
|
|
'',
|
|
f'- Tradable setup symbols with prices: '
|
|
f'{construction_count}.',
|
|
f'- Rank-only symbols with prices: '
|
|
f'{rank_only_count}.',
|
|
f'- Full ranking symbols with prices: '
|
|
f'{ranking_count}.',
|
|
f'- Tradable qualified longs: '
|
|
f'{qualified_count}.',
|
|
f'- Rank-only qualified rows removed: '
|
|
f'{filtered_count}.',
|
|
'',
|
|
]
|
|
paired = report['analysis']['paired_per_year']
|
|
for cost, protocol in (
|
|
(cost, protocol)
|
|
for cost in COSTS_PER_SIDE_PCT
|
|
for protocol in ('empty_book', 'warm_book')
|
|
):
|
|
lines.extend([
|
|
f'### {protocol.replace("_", " ").title()} — 0.10% per fill',
|
|
'',
|
|
'| Arm | ΔEV net R | 90% context | ΔCalmar | 90% context |',
|
|
'|---|---:|---:|---:|---:|',
|
|
])
|
|
lines[-4] = lines[-4].replace('0.10%', f'{cost:.2f}%')
|
|
for arm in ARMS:
|
|
row = next(
|
|
item
|
|
for item in paired
|
|
if item['arm_id'] == arm['id']
|
|
and item['protocol'] == protocol
|
|
and float(item['cost_per_side_pct']) == float(cost)
|
|
)
|
|
ev = row['headline']['ev_net_r']
|
|
calmar = row['headline']['calmar']
|
|
ev_ci = ev['bootstrap_90']
|
|
calmar_ci = calmar['bootstrap_90']
|
|
lines.append(
|
|
f'| {arm["id"]} | {_fmt(ev["paired_delta_median"])} | '
|
|
f'[{_fmt(ev_ci["p05"])}, {_fmt(ev_ci["p95"])}] | '
|
|
f'{_fmt(calmar["paired_delta_median"])} | '
|
|
f'[{_fmt(calmar_ci["p05"])}, {_fmt(calmar_ci["p95"])}] |'
|
|
)
|
|
lines.append('')
|
|
lines.extend([
|
|
'| Arm | ΔPF | ΔGain-to-Pain | ΔSortino | ΔCAGR pp | ΔMaxDD pp |',
|
|
'|---|---:|---:|---:|---:|---:|',
|
|
])
|
|
for arm in ARMS:
|
|
row = next(
|
|
item
|
|
for item in paired
|
|
if item['arm_id'] == arm['id']
|
|
and item['protocol'] == protocol
|
|
and float(item['cost_per_side_pct']) == float(cost)
|
|
)
|
|
headline = row['headline']
|
|
lines.append(
|
|
f'| {arm["id"]} | '
|
|
f'{_fmt(headline["profit_factor"]["paired_delta_median"])} | '
|
|
f'{_fmt(headline["gain_to_pain"]["paired_delta_median"])} | '
|
|
f'{_fmt(headline["sortino"]["paired_delta_median"])} | '
|
|
f'{_fmt(headline["cagr_pct"]["paired_delta_median"])} | '
|
|
f'{_fmt(headline["max_drawdown_pct"]["paired_delta_median"])} |'
|
|
)
|
|
lines.append('')
|
|
|
|
lines.extend([
|
|
'## Warm-seed initialization dispersion',
|
|
'',
|
|
'| Arm | Cost/fill | Median EV IQR ratio | Median Calmar IQR ratio |',
|
|
'|---|---:|---:|---:|',
|
|
])
|
|
for row in report['analysis']['warm_seed_dispersion']:
|
|
lines.append(
|
|
f'| {row["arm_id"]} | {row["cost_per_side_pct"]:.2f}% | '
|
|
f'{_fmt(row["headline"]["ev_net_r"]["median_iqr_ratio"])} | '
|
|
f'{_fmt(row["headline"]["calmar"]["median_iqr_ratio"])} |'
|
|
)
|
|
|
|
lines.extend([
|
|
'',
|
|
'## Capacity and operations — 0.10% per fill',
|
|
'',
|
|
'| Arm | Median trades | Median blocked | Median positions | Peak | '
|
|
'Turnover | Min-risk rejects |',
|
|
'|---|---:|---:|---:|---:|---:|---:|',
|
|
])
|
|
for row in report['operational_summary']:
|
|
blocked = row['median_blocked_fraction']
|
|
blocked_text = (
|
|
f'{float(blocked) * 100.0:.1f}%' if blocked is not None else 'n/a'
|
|
)
|
|
lines.append(
|
|
f'| {row["arm_id"]} | {_fmt(row["median_trades"], 1)} | '
|
|
f'{blocked_text} | {_fmt(row["median_avg_positions"], 2)} | '
|
|
f'{row["peak_positions"]} | '
|
|
f'{_fmt(row["median_turnover_multiple"], 2)} | '
|
|
f'{row["min_risk_rejections"]} |'
|
|
)
|
|
|
|
weekly = next(
|
|
row
|
|
for row in report['operational_summary']
|
|
if row['arm_id'] == 'cap10_weekly_top10'
|
|
)
|
|
lines.extend([
|
|
'',
|
|
'## Weekly-ranking opportunity set',
|
|
'',
|
|
f'- Median fresh entrant pool: '
|
|
f'{_fmt(weekly["weekly_entrant_pool_median"], 1)}.',
|
|
f'- Median zero-entrant fraction: '
|
|
f'{_fmt(weekly["weekly_zero_entrant_fraction"], 3)}.',
|
|
f'- Replacements across reported paths: {weekly["weekly_replacements"]}.',
|
|
f'- Same-symbol re-entries within 10 sessions: '
|
|
f'{weekly["weekly_reentries_within_10_sessions"]}.',
|
|
'',
|
|
'Bootstrap intervals above resample seven annual summaries and are '
|
|
'descriptive context only. They are not gates or independent-population '
|
|
'confidence claims.',
|
|
])
|
|
return '\n'.join(lines)
|
|
|
|
|
|
async def _main() -> None:
|
|
args = _parse_args()
|
|
snapshot = Path(args.snapshot)
|
|
if not snapshot.exists():
|
|
raise SystemExit(f'Snapshot does not exist: {snapshot}')
|
|
if not SPEC_PATH.exists():
|
|
raise SystemExit(f'Frozen specification is missing: {SPEC_PATH}')
|
|
|
|
os.environ['BACKTEST_SNAPSHOT_OFFLINE'] = '1'
|
|
os.environ['BACKTEST_ALLOW_SPAWN'] = '1'
|
|
workers = _worker_count(str(args.workers))
|
|
snapshot_sha256 = _sha256_file(snapshot)
|
|
specification_sha256 = _sha256_file(SPEC_PATH)
|
|
|
|
snapshot_data = await _load_snapshot(snapshot, quiet=bool(args.quiet))
|
|
cache_path = (
|
|
Path(args.candidate_cache)
|
|
if args.candidate_cache
|
|
else ROOT / 'reports' / '.cache' / f'{args.run_id}-candidates.pkl'
|
|
)
|
|
base_candidate_cache = _build_candidate_cache(
|
|
snapshot_data,
|
|
snapshot=snapshot,
|
|
snapshot_sha256=snapshot_sha256,
|
|
cache_path=cache_path,
|
|
workers=workers,
|
|
quiet=bool(args.quiet),
|
|
)
|
|
candidate_cache = _construction_candidate_view(
|
|
base_candidate_cache,
|
|
snapshot_data,
|
|
)
|
|
|
|
cohort_manifest = build_cohort_manifest(
|
|
snapshot_data['benchmark_closes'].keys()
|
|
)
|
|
cohort_errors = validate_cohort_manifest(cohort_manifest)
|
|
universe_errors = _construction_universe_errors(
|
|
snapshot_data['construction_universe_manifest']
|
|
)
|
|
validation_errors = [*universe_errors, *cohort_errors]
|
|
cells = build_cells(cohort_manifest)
|
|
validation_payload = {
|
|
'runner_version': RUNNER_VERSION,
|
|
'snapshot': str(snapshot.resolve()),
|
|
'snapshot_sha256': snapshot_sha256,
|
|
'snapshot_sessions': {
|
|
'count': cohort_manifest['session_count'],
|
|
'first': cohort_manifest['snapshot_first_session'],
|
|
'last': cohort_manifest['snapshot_last_session'],
|
|
},
|
|
'candidate_rank_coverage': {
|
|
'first': candidate_cache['rank_first_date'],
|
|
'last': candidate_cache['rank_last_date'],
|
|
'observations': candidate_cache['rank_observation_count'],
|
|
'raw_full_universe_qualified_longs': candidate_cache[
|
|
'raw_full_universe_qualified_long_count'
|
|
],
|
|
'filtered_rank_only_qualified_longs': candidate_cache[
|
|
'filtered_rank_only_qualified_long_count'
|
|
],
|
|
'construction_qualified_longs': candidate_cache[
|
|
'qualified_long_count'
|
|
],
|
|
},
|
|
'universe_manifest': snapshot_data['universe_manifest'],
|
|
'construction_universe_manifest': (
|
|
snapshot_data['construction_universe_manifest']
|
|
),
|
|
'empty_cluster_counts': cohort_manifest['empty_cluster_counts'],
|
|
'warm_seed_counts': cohort_manifest['warm_seed_counts'],
|
|
'empty_cluster_count': cohort_manifest['empty_cluster_count'],
|
|
'warm_cluster_count': cohort_manifest['warm_cluster_count'],
|
|
'expected_clusters': list(ANCHOR_YEARS),
|
|
'matrix_cells': len(cells),
|
|
'cache_path': str(cache_path.resolve()),
|
|
'cache_key_hash': _json_hash(candidate_cache['key']),
|
|
'construction_view_key_hash': _json_hash(
|
|
candidate_cache['construction_view_key']
|
|
),
|
|
'errors': validation_errors,
|
|
}
|
|
print(json.dumps(validation_payload, indent=2, sort_keys=True), flush=True)
|
|
if validation_errors:
|
|
raise SystemExit(
|
|
'Research validation failed; do not start the authoritative run'
|
|
)
|
|
if args.validate_only:
|
|
return
|
|
|
|
_assert_clean_worktree()
|
|
git_commit = _git_output('rev-parse', 'HEAD')
|
|
fingerprint_payload = {
|
|
'runner_version': RUNNER_VERSION,
|
|
'git_commit': git_commit,
|
|
'snapshot_sha256': snapshot_sha256,
|
|
'specification_sha256': specification_sha256,
|
|
'candidate_cache_key': candidate_cache['key'],
|
|
'construction_view_key': candidate_cache['construction_view_key'],
|
|
'universe_manifest': snapshot_data['universe_manifest'],
|
|
'construction_universe_manifest': (
|
|
snapshot_data['construction_universe_manifest']
|
|
),
|
|
'cohort_manifest': cohort_manifest,
|
|
'arms': list(ARMS),
|
|
'costs_per_side_pct': list(COSTS_PER_SIDE_PCT),
|
|
'bootstrap': {
|
|
'replicates': BOOTSTRAP_REPLICATES,
|
|
'seed': BOOTSTRAP_SEED,
|
|
},
|
|
}
|
|
fingerprint = _json_hash(fingerprint_payload)
|
|
checkpoint_dir = (
|
|
Path(args.checkpoint)
|
|
if args.checkpoint
|
|
else ROOT
|
|
/ 'reports'
|
|
/ '.cache'
|
|
/ f'{args.run_id}-{RUNNER_VERSION}-checkpoint'
|
|
)
|
|
completed = _checkpoint_state(
|
|
checkpoint_dir,
|
|
fingerprint,
|
|
resume=bool(args.resume),
|
|
)
|
|
expected_ids = {str(cell['cell_id']) for cell in cells}
|
|
unknown = set(completed) - expected_ids
|
|
if unknown:
|
|
raise SystemExit(
|
|
f'Checkpoint contains {len(unknown)} unknown matrix cells'
|
|
)
|
|
remaining = [
|
|
cell for cell in cells if str(cell['cell_id']) not in completed
|
|
]
|
|
if not args.quiet:
|
|
print(
|
|
f'matrix cells: {len(completed)} resumed, {len(remaining)} remaining',
|
|
flush=True,
|
|
)
|
|
|
|
worker_context = {
|
|
'qualified_candidates': candidate_cache['qualified_candidates'],
|
|
'daily_rank_map': candidate_cache['daily_rank_map'],
|
|
'prices': snapshot_data['prices'],
|
|
'benchmark_closes': snapshot_data['benchmark_closes'],
|
|
'ranking_key': snapshot_data['ranking_key'],
|
|
'exit_policy': snapshot_data['exit_policy'],
|
|
'hold_days': snapshot_data['hold_days'],
|
|
'risk_per_trade': snapshot_data['risk_per_trade'],
|
|
'atr_trail_multiplier': snapshot_data['atr_trail_multiplier'],
|
|
}
|
|
if workers == 1:
|
|
_worker_init(worker_context)
|
|
for index, cell in enumerate(remaining, 1):
|
|
row = _worker_run_cell(cell)
|
|
completed[str(row['cell_id'])] = row
|
|
_write_cell_checkpoint(checkpoint_dir, row)
|
|
if not args.quiet:
|
|
print(
|
|
f'portfolio cells: {len(completed)}/{len(cells)} '
|
|
f'({row["cell_id"]})',
|
|
flush=True,
|
|
)
|
|
elif remaining:
|
|
context = multiprocessing.get_context('spawn')
|
|
with ProcessPoolExecutor(
|
|
max_workers=workers,
|
|
mp_context=context,
|
|
initializer=_worker_init,
|
|
initargs=(worker_context,),
|
|
) as pool:
|
|
futures = {
|
|
pool.submit(_worker_run_cell, cell): str(cell['cell_id'])
|
|
for cell in remaining
|
|
}
|
|
for future in as_completed(futures):
|
|
row = future.result()
|
|
completed[str(row['cell_id'])] = row
|
|
_write_cell_checkpoint(checkpoint_dir, row)
|
|
if not args.quiet:
|
|
print(
|
|
f'portfolio cells: {len(completed)}/{len(cells)} '
|
|
f'({row["cell_id"]})',
|
|
flush=True,
|
|
)
|
|
|
|
missing = expected_ids - set(completed)
|
|
if missing:
|
|
raise RuntimeError(
|
|
f'Incomplete matrix: {len(missing)} cells are missing'
|
|
)
|
|
result_cells = sorted(
|
|
completed.values(),
|
|
key=lambda row: str(row['cell_id']),
|
|
)
|
|
analysis = aggregate_results(result_cells)
|
|
dependency_manifest = ROOT / 'pyproject.toml'
|
|
report: dict[str, Any] = {
|
|
'run_id': args.run_id,
|
|
'status': 'complete',
|
|
'generated_at': datetime.now(timezone.utc).isoformat(),
|
|
'research_question': (
|
|
'Bracket the economic cost of the binding ten-position cap and '
|
|
'test whether weekly current-rank selection beats arrival order.'
|
|
),
|
|
'decision_rule': (
|
|
'No formal promotion gate. Report paired annual medians, warm-seed '
|
|
'EV/Calmar IQR ratios, and simple bootstrap intervals as context.'
|
|
),
|
|
'survivorship_bias_caveat': (
|
|
'The current production universe is projected backward; construction '
|
|
'conclusions rely on paired relative comparisons, not absolute levels.'
|
|
),
|
|
'motivation': {
|
|
'source': 'reports/research-matrix-phase-a.json a0_control full window',
|
|
'trades': 472,
|
|
'skipped_book_full': 519,
|
|
'blocked_fraction': 519 / (519 + 472),
|
|
'stale_claim_corrected': (
|
|
'The older weekly pre-gate-reset claim that cap 10 never bound '
|
|
'does not apply to the current daily configuration.'
|
|
),
|
|
},
|
|
'fingerprint': fingerprint,
|
|
'fingerprint_payload': fingerprint_payload,
|
|
'environment': {
|
|
'python': sys.version,
|
|
'platform': platform.platform(),
|
|
'dependency_manifest': str(dependency_manifest.relative_to(ROOT)),
|
|
'dependency_manifest_sha256': (
|
|
_sha256_file(dependency_manifest)
|
|
if dependency_manifest.exists()
|
|
else None
|
|
),
|
|
'command': [sys.executable, *sys.argv],
|
|
},
|
|
'validation': validation_payload,
|
|
'runtime_config': snapshot_data['runtime_config'],
|
|
'universe_manifest': snapshot_data['universe_manifest'],
|
|
'construction_universe_manifest': (
|
|
snapshot_data['construction_universe_manifest']
|
|
),
|
|
'candidate_cache': {
|
|
key: value
|
|
for key, value in candidate_cache.items()
|
|
if key
|
|
not in {
|
|
'qualified_candidates',
|
|
'daily_rank_map',
|
|
}
|
|
},
|
|
'cohort_manifest': cohort_manifest,
|
|
'arms': list(ARMS),
|
|
'costs_per_side_pct': list(COSTS_PER_SIDE_PCT),
|
|
'cell_count': len(result_cells),
|
|
'cells': result_cells,
|
|
'analysis': analysis,
|
|
'operational_summary': _operational_summary(result_cells),
|
|
}
|
|
out_path = (
|
|
Path(args.out)
|
|
if args.out
|
|
else ROOT
|
|
/ 'reports'
|
|
/ f'portfolio-construction-{args.run_id}.json'
|
|
)
|
|
_atomic_json(out_path, report)
|
|
_atomic_text(out_path.with_suffix('.md'), _markdown(report))
|
|
if not args.quiet:
|
|
print(f'wrote {out_path}', flush=True)
|
|
print(f'wrote {out_path.with_suffix(".md")}', flush=True)
|
|
|
|
|
|
if __name__ == '__main__':
|
|
asyncio.run(_main())
|