malkhut(T9): smoke test launchers
launch_smoke_test.py: 10-min quick smoke. smoke_test_60min.py: 60-min full smoke with checkpoints.
This commit is contained in:
349
MALKHUT/malkhut/launch_smoke_test.py
Normal file
349
MALKHUT/malkhut/launch_smoke_test.py
Normal file
@@ -0,0 +1,349 @@
|
|||||||
|
#!/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()
|
||||||
327
MALKHUT/malkhut/smoke_test_60min.py
Normal file
327
MALKHUT/malkhut/smoke_test_60min.py
Normal file
@@ -0,0 +1,327 @@
|
|||||||
|
#!/usr/bin/env python3
|
||||||
|
"""
|
||||||
|
MALKHUT 60-Minute Smoke Test — comprehensive system validation.
|
||||||
|
|
||||||
|
Runs the full pipeline for 60 minutes with:
|
||||||
|
- Training pipeline (CMA-ES + genetic programming)
|
||||||
|
- Strategy generator (evolving strategies)
|
||||||
|
- All 9 planner types cycling
|
||||||
|
- Performance metrics tracking
|
||||||
|
- Resource usage monitoring
|
||||||
|
- Strategy development tracking
|
||||||
|
- Improvement metrics
|
||||||
|
|
||||||
|
Usage:
|
||||||
|
python -m malkhut.smoke_test_60min
|
||||||
|
"""
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import json
|
||||||
|
import os
|
||||||
|
import resource
|
||||||
|
import sys
|
||||||
|
import threading
|
||||||
|
import time
|
||||||
|
from dataclasses import dataclass, field
|
||||||
|
from typing import Any, Dict, List, Optional
|
||||||
|
|
||||||
|
# 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.planner.alternatives import PLANNER_REGISTRY, create_planner
|
||||||
|
from malkhut.cwm.core import MinimalCryptoLOBCWM
|
||||||
|
from malkhut.counterparties import default_counterparty_ecology
|
||||||
|
from malkhut.storage.ch_store import MalkhutCHStore
|
||||||
|
|
||||||
|
|
||||||
|
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,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
# ── Resource Monitor ─────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
class ResourceMonitor:
|
||||||
|
"""Track CPU and RAM usage."""
|
||||||
|
|
||||||
|
def __init__(self):
|
||||||
|
self._samples: list = []
|
||||||
|
self._start = time.time()
|
||||||
|
self._peak_ram = 0.0
|
||||||
|
self._peak_cpu = 0.0
|
||||||
|
|
||||||
|
def sample(self):
|
||||||
|
usage = resource.getrusage(resource.RUSAGE_SELF)
|
||||||
|
ram_mb = usage.ru_maxrss / 1024
|
||||||
|
try:
|
||||||
|
with open("/proc/self/stat") as f:
|
||||||
|
fields = f.read().split()
|
||||||
|
utime = int(fields[13])
|
||||||
|
stime = int(fields[14])
|
||||||
|
elapsed = time.time() - self._start
|
||||||
|
cpu = min(100.0, ((utime + stime) * 10.0) / max(elapsed * 1000.0, 1.0) * 100.0)
|
||||||
|
except Exception:
|
||||||
|
cpu = 0.0
|
||||||
|
|
||||||
|
self._samples.append({"time": time.time() - self._start, "cpu": cpu, "ram_mb": ram_mb})
|
||||||
|
self._peak_ram = max(self._peak_ram, ram_mb)
|
||||||
|
self._peak_cpu = max(self._peak_cpu, cpu)
|
||||||
|
|
||||||
|
@property
|
||||||
|
def avg_cpu(self) -> float:
|
||||||
|
if not self._samples: return 0
|
||||||
|
return sum(s["cpu"] for s in self._samples) / len(self._samples)
|
||||||
|
|
||||||
|
@property
|
||||||
|
def avg_ram(self) -> float:
|
||||||
|
if not self._samples: return 0
|
||||||
|
return sum(s["ram_mb"] for s in self._samples) / len(self._samples)
|
||||||
|
|
||||||
|
|
||||||
|
# ── Strategy Tracker ────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
class StrategyTracker:
|
||||||
|
"""Track strategies developed and their improvement."""
|
||||||
|
|
||||||
|
def __init__(self):
|
||||||
|
self._strategies: list = []
|
||||||
|
self._scores: list = []
|
||||||
|
self._planner_types_used: dict = {}
|
||||||
|
self._best_score_history: list = []
|
||||||
|
|
||||||
|
def record(self, score: float, planner_type: str, generation: int):
|
||||||
|
self._strategies.append({"score": score, "planner": planner_type, "gen": generation})
|
||||||
|
self._scores.append(score)
|
||||||
|
self._planner_types_used[planner_type] = self._planner_types_used.get(planner_type, 0) + 1
|
||||||
|
self._best_score_history.append(max(self._scores) if self._scores else 0)
|
||||||
|
|
||||||
|
@property
|
||||||
|
def total_strategies(self) -> int:
|
||||||
|
return len(self._strategies)
|
||||||
|
|
||||||
|
@property
|
||||||
|
def best_score(self) -> float:
|
||||||
|
return max(self._scores) if self._scores else 0
|
||||||
|
|
||||||
|
@property
|
||||||
|
def improvement(self) -> float:
|
||||||
|
if len(self._scores) < 2: return 0
|
||||||
|
return self._best_score_history[-1] - self._best_score_history[0]
|
||||||
|
|
||||||
|
@property
|
||||||
|
def planner_usage(self) -> dict:
|
||||||
|
return dict(self._planner_types_used)
|
||||||
|
|
||||||
|
def summary(self) -> dict:
|
||||||
|
return {
|
||||||
|
"total_strategies": self.total_strategies,
|
||||||
|
"best_score": self.best_score,
|
||||||
|
"improvement": self.improvement,
|
||||||
|
"planner_usage": self.planner_usage,
|
||||||
|
"score_history_len": len(self._best_score_history),
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
# ── Main Smoke Test ──────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
def run_60min_smoke():
|
||||||
|
DURATION_S = 3600 # 60 minutes
|
||||||
|
|
||||||
|
print("=" * 70)
|
||||||
|
print("MALKHUT 60-MINUTE SMOKE TEST")
|
||||||
|
print(f"Duration: {DURATION_S}s ({DURATION_S // 60} minutes)")
|
||||||
|
print("=" * 70)
|
||||||
|
|
||||||
|
monitor = ResourceMonitor()
|
||||||
|
tracker = StrategyTracker()
|
||||||
|
t0 = time.time()
|
||||||
|
|
||||||
|
# Setup
|
||||||
|
print("\n[1/4] Setting up infrastructure...")
|
||||||
|
monitor.sample()
|
||||||
|
store = MalkhutCHStore()
|
||||||
|
store.ensure_tables()
|
||||||
|
registry = PolicyRegistry(store=store)
|
||||||
|
|
||||||
|
# Training pipeline
|
||||||
|
print("[2/4] Running training pipeline (cycles through ALL 9 planner types)...")
|
||||||
|
pipeline_config = PipelineConfig(
|
||||||
|
max_generations=20,
|
||||||
|
max_evals_per_generation=10,
|
||||||
|
max_time_s=DURATION_S * 0.6,
|
||||||
|
auto_promote=True,
|
||||||
|
)
|
||||||
|
pipeline = TrainingPipeline(
|
||||||
|
config=pipeline_config, registry=registry,
|
||||||
|
log_path=os.path.join(_HERE, "smoke_60min.log"),
|
||||||
|
)
|
||||||
|
monitor.sample()
|
||||||
|
|
||||||
|
# Run training
|
||||||
|
pipeline_result = pipeline.run(
|
||||||
|
incumbent=_baseline(),
|
||||||
|
symbols=("BTCUSDT",),
|
||||||
|
)
|
||||||
|
monitor.sample()
|
||||||
|
|
||||||
|
# Track strategies from training
|
||||||
|
for event in pipeline_result.events:
|
||||||
|
if event.event_type == "generation":
|
||||||
|
tracker.record(event.score, "cma_es", event.generation)
|
||||||
|
|
||||||
|
print(f" Generations: {pipeline_result.generations_run}")
|
||||||
|
print(f" Evals: {pipeline_result.total_evals}")
|
||||||
|
print(f" Best score: {pipeline_result.best_score:.2f}")
|
||||||
|
|
||||||
|
# Strategy generator
|
||||||
|
print("[3/4] Running strategy generator (genetic programming)...")
|
||||||
|
remaining_time = DURATION_S * 0.3 - pipeline_result.duration_s
|
||||||
|
if remaining_time > 30:
|
||||||
|
gen_config = GeneratorConfig(
|
||||||
|
population_size=15, generations=3, 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)
|
||||||
|
genetic_count = len([g for g in gen_population if g.generation > 0])
|
||||||
|
|
||||||
|
for genome in gen_population:
|
||||||
|
if genome.generation > 0:
|
||||||
|
tracker.record(genome.fitness, genome.strategy_type.value, genome.generation)
|
||||||
|
generator.add_to_pool(genome)
|
||||||
|
|
||||||
|
print(f" Population: {len(gen_population)} strategies")
|
||||||
|
print(f" Genetic strategies: {genetic_count}")
|
||||||
|
|
||||||
|
monitor.sample()
|
||||||
|
|
||||||
|
# Planner diversity test
|
||||||
|
print("[4/4] Testing all 9 planner types...")
|
||||||
|
planner_scores = {}
|
||||||
|
for name in PLANNER_REGISTRY.keys():
|
||||||
|
try:
|
||||||
|
cwm = MinimalCryptoLOBCWM()
|
||||||
|
planner = create_planner(name, cwm=cwm, counterparties=default_counterparty_ecology())
|
||||||
|
from malkhut.state import ExecutionIntent, IntentKind, MarketWorldState, Mode, OrderBookState, AccountState, PriceLevel
|
||||||
|
s = MarketWorldState(
|
||||||
|
ts_ns=1, mode=Mode.REPLAY_NO_IMPACT, venue=_venue(),
|
||||||
|
book=_book(), account=_account(),
|
||||||
|
intent=_intent(),
|
||||||
|
)
|
||||||
|
result = planner.plan(s, _baseline(), budget_ms=10)
|
||||||
|
planner_scores[name] = len(result.actions)
|
||||||
|
except Exception as e:
|
||||||
|
planner_scores[name] = f"error: {e}"
|
||||||
|
|
||||||
|
# Final metrics
|
||||||
|
duration = time.time() - t0
|
||||||
|
monitor.sample()
|
||||||
|
|
||||||
|
print()
|
||||||
|
print("=" * 70)
|
||||||
|
print("60-MINUTE SMOKE TEST RESULTS")
|
||||||
|
print("=" * 70)
|
||||||
|
print(f"Duration: {duration:.1f}s ({duration/60:.1f} min)")
|
||||||
|
print(f"Generations: {pipeline_result.generations_run}")
|
||||||
|
print(f"Total evals: {pipeline_result.total_evals}")
|
||||||
|
print(f"Best score: {pipeline_result.best_score:.2f}")
|
||||||
|
print(f"Strategies dev: {tracker.total_strategies}")
|
||||||
|
print(f"Improvement: {tracker.improvement:.2f}")
|
||||||
|
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("PLANNER USAGE")
|
||||||
|
for ptype, count in tracker.planner_usage.items():
|
||||||
|
print(f" {ptype:<20} {count} evaluations")
|
||||||
|
print()
|
||||||
|
print("PLANNER DIVERSITY")
|
||||||
|
for name, score in planner_scores.items():
|
||||||
|
print(f" {name:<20} {score} actions")
|
||||||
|
print()
|
||||||
|
print("EVENTS LOGGED")
|
||||||
|
print(f" Pipeline events: {len(pipeline_result.events)}")
|
||||||
|
print(f" Registry records: {registry.record_count}")
|
||||||
|
print("=" * 70)
|
||||||
|
|
||||||
|
# Save results
|
||||||
|
results = {
|
||||||
|
"duration_s": duration,
|
||||||
|
"generations": pipeline_result.generations_run,
|
||||||
|
"total_evals": pipeline_result.total_evals,
|
||||||
|
"best_score": pipeline_result.best_score,
|
||||||
|
"strategies_developed": tracker.total_strategies,
|
||||||
|
"improvement": tracker.improvement,
|
||||||
|
"peak_cpu_pct": monitor._peak_cpu,
|
||||||
|
"avg_cpu_pct": monitor.avg_cpu,
|
||||||
|
"peak_ram_mb": monitor._peak_ram,
|
||||||
|
"avg_ram_mb": monitor.avg_ram,
|
||||||
|
"planner_usage": tracker.planner_usage,
|
||||||
|
"planner_diversity": planner_scores,
|
||||||
|
"events_logged": len(pipeline_result.events),
|
||||||
|
"registry_records": registry.record_count,
|
||||||
|
}
|
||||||
|
with open(os.path.join(_HERE, "smoke_60min_results.json"), "w") as f:
|
||||||
|
json.dump(results, f, indent=2)
|
||||||
|
print(f"\nResults saved to smoke_60min_results.json")
|
||||||
|
|
||||||
|
|
||||||
|
def _venue():
|
||||||
|
from malkhut.state import VenueRules
|
||||||
|
return VenueRules(exchange="bingx", symbol="BTCUSDT", tick_size=0.1, lot_size=0.001,
|
||||||
|
min_qty=0.001, min_notional=5.0, maker_fee_bps=-0.2, taker_fee_bps=0.5,
|
||||||
|
post_only_supported=True, reduce_only_supported=True,
|
||||||
|
max_orders_per_second=100, max_cancels_per_minute=120)
|
||||||
|
|
||||||
|
|
||||||
|
def _book():
|
||||||
|
from malkhut.state import OrderBookState, PriceLevel
|
||||||
|
return OrderBookState(ts_ns=1, symbol="BTCUSDT",
|
||||||
|
bids=(PriceLevel(50000.0, 1.0),), asks=(PriceLevel(50001.0, 1.0),))
|
||||||
|
|
||||||
|
|
||||||
|
def _account():
|
||||||
|
from malkhut.state import AccountState
|
||||||
|
return AccountState(ts_ns=1, equity=10000.0, wallet_balance=10000.0,
|
||||||
|
available_balance=10000.0, margin_used=0.0, total_notional=0.0)
|
||||||
|
|
||||||
|
|
||||||
|
def _intent():
|
||||||
|
from malkhut.state import ExecutionIntent, IntentKind
|
||||||
|
return ExecutionIntent(
|
||||||
|
intent_id="smoke", ts_ns=1, symbol="BTCUSDT",
|
||||||
|
kind=IntentKind.ENTER_LONG, target_qty=0.01, max_notional=500.0,
|
||||||
|
urgency=0.5, alpha_horizon_s=60.0, alpha_bps=2.0,
|
||||||
|
max_slippage_bps=5.0, prefer_maker=True, reduce_only=False,
|
||||||
|
ttl_s=300.0, reason="smoke_test",
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
run_60min_smoke()
|
||||||
Reference in New Issue
Block a user