Files
sentiment-engine/MALKHUT/malkhut/launch_smoke_test.py
Codex 8af7e3bce8 malkhut(T9): smoke test launchers
launch_smoke_test.py: 10-min quick smoke.
smoke_test_60min.py: 60-min full smoke with checkpoints.
2026-07-11 10:41:31 +02:00

350 lines
13 KiB
Python

#!/usr/bin/env python3
"""
MALKHUT Training Pipeline Launcher — 10-minute smoke test.
Launches the full training pipeline with bounded resources:
- CMA-ES training (generations, evals, time budget)
- Strategy generation (genetic operators)
- Pipeline logging (JSONL)
- CPU/RAM monitoring
- Results summary
Usage:
python -m malkhut.launch_smoke_test [--duration 600] [--evals 50]
Naming convention for discovered strategies:
{strategy_type}_{generation}_{timestamp}
e.g., SM_MCTS_gen3_20260707_034500
"""
from __future__ import annotations
import argparse
import json
import os
import resource
import sys
import time
from dataclasses import dataclass, field
from typing import Any, List, Mapping
# ── Setup paths ──────────────────────────────────────────────────────────────
_HERE = os.path.dirname(os.path.abspath(__file__))
if _HERE not in sys.path:
sys.path.insert(0, _HERE)
from malkhut.state import FulfilmentPolicyParams
from malkhut.training.pipeline import TrainingPipeline, PipelineConfig
from malkhut.training.generator import StrategyGenerator, GeneratorConfig
from malkhut.training.registry import PolicyRegistry
from malkhut.training.cma_trainer import ScenarioFactory
from malkhut.storage.ch_store import MalkhutCHStore
# ── Configuration ────────────────────────────────────────────────────────────
def _baseline() -> FulfilmentPolicyParams:
return FulfilmentPolicyParams(
version="baseline", ucb_c=1.414, max_sims=256, max_depth=3,
rollout_depth=3, root_temperature=0.5, min_root_entropy=0.25,
quote_offsets_ticks=(0, 1, 2), quote_size_fractions=(0.1, 0.25, 0.5),
passive_ttl_ms=200, aggressive_ttl_ms=50,
maker_edge_min_bps=0.5, cross_spread_edge_min_bps=5.0,
adverse_toxicity_cancel_threshold=0.5, queue_churn_cancel_threshold=0.5,
mae_tail_cut_bps=50.0, mfe_giveback_cut_fraction=0.5,
max_time_in_loss_s=300.0, failed_recovery_cut_count=3,
recovery_velocity_min_bps_per_s=0.0,
max_symbol_notional_fraction=0.20, max_single_order_notional_fraction=0.05,
reduce_when_global_up_fraction=0.30, session_profit_lock_fraction=0.02,
w_expected_pnl=1.0, w_fill_probability=0.5, w_adverse_selection=2.0,
w_queue_priority=0.5, w_inventory_risk=1.5, w_tail_loss=5.0,
w_fee_quality=0.5, w_time_decay=0.3, w_policy_entropy=0.5,
robust_tail_weight=2.0, toxic_counterparty_weight=3.0,
low_liquidity_weight=2.0, latency_stress_weight=1.0,
)
@dataclass
class SmokeTestResult:
"""Results from a smoke test run."""
duration_s: float
generations_run: int
total_evals: int
best_score: float
strategies_developed: int
builtin_strategies_parsed: int
genetic_strategies_evolved: int
peak_cpu_pct: float
peak_ram_mb: float
avg_cpu_pct: float
avg_ram_mb: float
events_logged: int
registry_records: int
strategy_names: List[str] = field(default_factory=list)
# ── CPU/RAM Monitor ──────────────────────────────────────────────────────────
class ResourceMonitor:
"""Track CPU and RAM usage during the run."""
def __init__(self) -> None:
self._samples: list[tuple[float, float]] = [] # (cpu%, ram_mb)
self._start_time = time.time()
def sample(self) -> tuple[float, float]:
"""Sample current CPU% and RAM MB."""
usage = resource.getrusage(resource.RUSAGE_SELF)
ram_mb = usage.ru_maxrss / 1024 # KB → MB (Linux)
# CPU% from /proc/self/stat (user + system time)
try:
with open("/proc/self/stat") as f:
fields = f.read().split()
utime = int(fields[13]) # user time (ticks)
stime = int(fields[14]) # system time (ticks)
total_ticks = utime + stime
elapsed = time.time() - self._start_time
# Approximate CPU% (ticks are ~10ms on Linux)
cpu_pct = min(100.0, (total_ticks * 10.0) / max(elapsed * 1000.0, 1.0) * 100.0)
except Exception:
cpu_pct = 0.0
self._samples.append((cpu_pct, ram_mb))
return cpu_pct, ram_mb
@property
def peak_cpu(self) -> float:
return max((s[0] for s in self._samples), default=0.0)
@property
def peak_ram(self) -> float:
return max((s[1] for s in self._samples), default=0.0)
@property
def avg_cpu(self) -> float:
if not self._samples:
return 0.0
return sum(s[0] for s in self._samples) / len(self._samples)
@property
def avg_ram(self) -> float:
if not self._samples:
return 0.0
return sum(s[1] for s in self._samples) / len(self._samples)
# ── Strategy Naming ──────────────────────────────────────────────────────────
def name_strategy(
strategy_type: str,
generation: int,
fitness: float,
parent_ids: tuple[str, ...] = (),
) -> str:
"""
Name a discovered strategy.
Convention:
{type}_gen{N}_{timestamp}
Examples:
SM_MCTS_gen3_20260707_034500
UCB1_gen1_20260707_034515
"""
ts = time.strftime("%Y%m%d_%H%M%S")
return f"{strategy_type}_gen{generation}_{ts}"
# ── Main Smoke Test ──────────────────────────────────────────────────────────
def run_smoke_test(duration_s: int = 600, max_evals: int = 50) -> SmokeTestResult:
"""
Run a 10-minute smoke test of the full training pipeline.
Returns SmokeTestResult with metrics.
"""
print("=" * 70)
print("MALKHUT SMOKE TEST — Training Pipeline")
print(f"Duration: {duration_s}s | Max evals: {max_evals}")
print("=" * 70)
monitor = ResourceMonitor()
t0 = time.time()
# ── Setup ────────────────────────────────────────────────────────────────
print("\n[1/5] Setting up infrastructure...")
monitor.sample()
store = MalkhutCHStore()
store.ensure_tables()
registry = PolicyRegistry(store=store)
# ── Training Pipeline ────────────────────────────────────────────────────
print("[2/5] Running training pipeline...")
pipeline_config = PipelineConfig(
max_generations=5,
max_evals_per_generation=max_evals // 5,
max_time_s=duration_s * 0.6, # 60% of time for training
auto_promote=True,
)
pipeline = TrainingPipeline(
config=pipeline_config, registry=registry,
log_path=os.path.join(_HERE, "training.log"),
)
monitor.sample()
pipeline_result = pipeline.run(
incumbent=_baseline(),
symbols=("BTCUSDT",),
)
monitor.sample()
print(f" Generations: {pipeline_result.generations_run}")
print(f" Evals: {pipeline_result.total_evals}")
print(f" Best score: {pipeline_result.best_score:.4f}")
print(f" Duration: {pipeline_result.duration_s:.1f}s")
# ── Strategy Generation ──────────────────────────────────────────────────
print("[3/5] Running strategy generator...")
remaining_time = duration_s * 0.3 - pipeline_result.duration_s
if remaining_time > 10:
gen_config = GeneratorConfig(
population_size=10,
generations=2,
tournament_size=3,
elitism_count=2,
)
generator = StrategyGenerator(config=gen_config, registry=registry)
scenarios = ScenarioFactory().build_suite(symbols=("BTCUSDT",), steps_per_scenario=5)
gen_population = generator.evolve(_baseline(), scenarios)
# Name discovered strategies
strategy_names = []
for genome in gen_population:
name = name_strategy(
genome.strategy_type.value,
genome.generation,
genome.fitness,
)
strategy_names.append(name)
generator.add_to_pool(genome)
genetic_count = len([g for g in gen_population if g.generation > 0])
print(f" Population: {len(gen_population)} strategies")
print(f" Genetic strategies evolved: {genetic_count}")
else:
gen_population = []
strategy_names = []
genetic_count = 0
print(" Skipped (time budget exhausted)")
monitor.sample()
# ── Builtin Strategy Parsing ─────────────────────────────────────────────
print("[4/5] Parsing builtin strategies...")
from malkhut.training.dsl import StrategyDSLCompiler, list_builtin_strategies
compiler = StrategyDSLCompiler()
builtin_count = 0
for name in list_builtin_strategies():
from malkhut.training.dsl import get_builtin_strategy
text = get_builtin_strategy(name)
if text:
template = compiler.compile(text)
builtin_count += 1
print(f" Builtin strategies parsed: {builtin_count}")
# ── Summary ──────────────────────────────────────────────────────────────
duration = time.time() - t0
monitor.sample()
print("[5/5] Summary...")
print()
print("=" * 70)
print("SMOKE TEST RESULTS")
print("=" * 70)
print(f"Duration: {duration:.1f}s")
print(f"Generations run: {pipeline_result.generations_run}")
print(f"Total evals: {pipeline_result.total_evals}")
print(f"Best score: {pipeline_result.best_score:.4f}")
print(f"Strategies developed: {len(gen_population)} (genetic: {genetic_count})")
print(f"Builtin strategies: {builtin_count}")
print(f"Registry records: {registry.record_count}")
print(f"Events logged: {len(pipeline_result.events)}")
print()
print("RESOURCE USAGE")
print(f"Peak CPU: {monitor.peak_cpu:.1f}%")
print(f"Avg CPU: {monitor.avg_cpu:.1f}%")
print(f"Peak RAM: {monitor.peak_ram:.1f} MB")
print(f"Avg RAM: {monitor.avg_ram:.1f} MB")
print()
print("STRATEGY NAMING CONVENTION")
print(" {strategy_type}_gen{generation}_{timestamp}")
print(" Examples:")
for name in strategy_names[:5]:
print(f" {name}")
if len(strategy_names) > 5:
print(f" ... and {len(strategy_names) - 5} more")
print()
print("STRATEGY TYPES DISCOVERED")
if gen_population:
types = set(g.strategy_type.value for g in gen_population)
for t in types:
count = sum(1 for g in gen_population if g.strategy_type.value == t)
print(f" {t}: {count}")
print("=" * 70)
return SmokeTestResult(
duration_s=duration,
generations_run=pipeline_result.generations_run,
total_evals=pipeline_result.total_evals,
best_score=pipeline_result.best_score,
strategies_developed=len(gen_population),
builtin_strategies_parsed=builtin_count,
genetic_strategies_evolved=genetic_count,
peak_cpu_pct=monitor.peak_cpu,
peak_ram_mb=monitor.peak_ram,
avg_cpu_pct=monitor.avg_cpu,
avg_ram_mb=monitor.avg_ram,
events_logged=len(pipeline_result.events),
registry_records=registry.record_count,
strategy_names=strategy_names,
)
# ── Entry Point ──────────────────────────────────────────────────────────────
def main():
parser = argparse.ArgumentParser(description="MALKHUT Smoke Test")
parser.add_argument("--duration", type=int, default=600, help="Duration in seconds")
parser.add_argument("--evals", type=int, default=50, help="Max evaluations")
args = parser.parse_args()
result = run_smoke_test(duration_s=args.duration, max_evals=args.evals)
# Write results to JSON
output = {
"duration_s": result.duration_s,
"generations_run": result.generations_run,
"total_evals": result.total_evals,
"best_score": result.best_score,
"strategies_developed": result.strategies_developed,
"builtin_strategies_parsed": result.builtin_strategies_parsed,
"genetic_strategies_evolved": result.genetic_strategies_evolved,
"peak_cpu_pct": result.peak_cpu_pct,
"peak_ram_mb": result.peak_ram_mb,
"avg_cpu_pct": result.avg_cpu_pct,
"avg_ram_mb": result.avg_ram_mb,
"events_logged": result.events_logged,
"registry_records": result.registry_records,
"strategy_names": result.strategy_names,
}
with open(os.path.join(_HERE, "smoke_test_results.json"), "w") as f:
json.dump(output, f, indent=2)
print(f"\nResults saved to smoke_test_results.json")
if __name__ == "__main__":
main()