malkhut(assets)+uv: system-wide asset directory + UV universe init

MALKHUT asset directory (malkhut/assets/): normalized canonical symbols,
KNOWN_EXCHANGES aux table (BINANCE/BINGX/BINGX_VST), per-exchange listing
status (TRADING/OFFLINE/UNKNOWN), JSON-backed, built for full-Binance-500
scale. Seeded with BLUE's NG7 feed universe (50 symbols) + live VST
contracts probe: 35 TRADING / 15 OFFLINE on VST (BAND, CELR, COS, CVC,
DENT, FUN, HOT, ICX, TFUEL, TUSD, USDC, WAN, WIN, XTZ, ZIL).
uv_asset_universe.init_asset_universe() = runtime tradable set for the
execution exchange; --onboard runs MALKHUT AssetCompiler for full taxonomy.
12 new tests, mutation-RED verified (status-filter + tradability guard).
This commit is contained in:
Codex
2026-07-11 21:48:05 +02:00
parent aaaf326abf
commit a062507656
7 changed files with 1687 additions and 0 deletions

View File

@@ -0,0 +1,203 @@
"""
UV asset-universe init — MALKHUT asset directory as the single source of truth.
Operator directive 2026-07-11: the MALKHUT asset store (malkhut/assets/) is
the system-wide known-asset universe. UV/VIOLET initialize their tradable
set from it against the CURRENT execution exchange, so a dead venue symbol
(e.g. BAND-USDT offline on VST) can never reach the kernel.
Flow:
1. import_blue_universe() — seed canonical symbols from the DOLPHIN NG7
scan feed (BINANCE listings; that feed IS the
Binance universe BLUE trades from).
2. refresh_vst_listings() — probe the execution venue's contracts endpoint
and stamp BINGX_VST TRADING/OFFLINE per symbol.
3. init_asset_universe() — the runtime call: returns the frozen tradable
set for the execution exchange. The promotion
bridge / seam refuses ENTER intents outside it.
Built for the full-Binance (~500+) universe north-star; today's feed carries
50 symbols — the store scales, the feed is the bottleneck.
CLI:
uv_asset_universe.py --seed-from-scans --refresh-vst --show
"""
from __future__ import annotations
import argparse
import json
import os
import sys
import urllib.request
from datetime import datetime, timezone
from typing import Any, Iterable
MALKHUT_ROOT = os.environ.get("MALKHUT_ROOT", "/mnt/dolphinng5_predict/MALKHUT")
if MALKHUT_ROOT not in sys.path:
sys.path.insert(0, MALKHUT_ROOT)
from malkhut.assets.directory import ( # noqa: E402
AssetDirectory,
ListingStatus,
normalize_symbol,
)
from uv_anomaly_sentinel import CHQuerier, CH_DB_UV, VST_BASE_URL # noqa: E402
FEED_SOURCE = "dolphin_ng7_scan_feed"
PROBE_SOURCE = "vst_contracts_probe"
def _now_iso() -> str:
return datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")
# ── step 1: seed from the scan feed (the BLUE/Binance universe) ────────────
def blue_universe_from_scans(ch: CHQuerier | None = None) -> list[str]:
ch = ch or CHQuerier()
rows = ch.query(
f"SELECT DISTINCT arrayJoin(JSONExtract(assets,'Array(String)')) AS a"
f" FROM {CH_DB_UV}.prime_scans ORDER BY a FORMAT JSONEachRow"
)
return [r["a"] for r in rows]
def import_blue_universe(directory: AssetDirectory, symbols: Iterable[str]) -> int:
return directory.import_symbols(
[normalize_symbol(s) for s in symbols],
"BINANCE",
status=ListingStatus.TRADING,
checked_at=_now_iso(),
source=FEED_SOURCE,
)
# ── step 2: stamp execution-venue listings from the contracts probe ───────
def fetch_vst_contracts(base_url: str = VST_BASE_URL, timeout_s: int = 15) -> list[dict[str, Any]]:
req = urllib.request.Request(f"{base_url}/openApi/swap/v2/quote/contracts")
body = urllib.request.urlopen(req, timeout=timeout_s).read().decode("utf-8")
parsed = json.loads(body)
return parsed.get("data") or []
def contract_is_tradable(contract: dict[str, Any]) -> bool:
"""BingX contracts row -> tradable? status 1 = online; apiStateOpen gates API entry."""
return (
int(contract.get("status") or 0) == 1
and str(contract.get("apiStateOpen", "true")).lower() == "true"
)
def refresh_vst_listings(
directory: AssetDirectory, contracts: list[dict[str, Any]]
) -> tuple[int, int]:
"""Stamp BINGX_VST status for every symbol already in the directory.
Present + tradable -> TRADING; absent or non-tradable -> OFFLINE.
Returns (trading, offline) counts over the directory.
"""
now = _now_iso()
tradable: dict[str, str] = {} # canonical -> venue symbol
for c in contracts:
vsym = str(c.get("symbol") or "")
if vsym and contract_is_tradable(c):
tradable[normalize_symbol(vsym)] = vsym
n_trading = n_offline = 0
for sym in list(directory.records):
if sym in tradable:
directory.set_listing(sym, "BINGX_VST", venue_symbol=tradable[sym],
status=ListingStatus.TRADING,
checked_at=now, source=PROBE_SOURCE)
n_trading += 1
else:
directory.set_listing(sym, "BINGX_VST",
status=ListingStatus.OFFLINE,
checked_at=now, source=PROBE_SOURCE)
n_offline += 1
return n_trading, n_offline
# ── step 2b: full-taxonomy onboarding via MALKHUT's AssetCompiler ──────────
def onboard_profiles(directory: AssetDirectory, *, compiler: Any = None,
out_path: str = "") -> tuple[int, int]:
"""Run MALKHUT's asset onboarder over every directory symbol lacking a
profile_ref; persist compiled taxonomy JSON beside the directory.
Returns (compiled, failed). Failures are noted on the record, not fatal —
the directory stays authoritative for existence/listing even when the
taxonomy fetch fails.
"""
from dataclasses import asdict
if compiler is None:
from malkhut.training.asset_compiler import AssetCompiler
compiler = AssetCompiler()
path = out_path or str(directory.path.parent / "compiled_profiles.json")
try:
with open(path, encoding="utf-8") as fh:
store = json.load(fh)
except (OSError, json.JSONDecodeError):
store = {}
done = failed = 0
for sym in sorted(directory.records):
rec = directory.records[sym]
if rec.profile_ref and sym in store:
continue
result = compiler.compile_and_register(sym)
if result.asset_profile is None:
failed += 1
rec.notes = (rec.notes + " | " if rec.notes else "") + \
f"onboard failed {_now_iso()}: {'; '.join(result.warnings)[:120]}"
continue
entry = {"profile": asdict(result.asset_profile), "_compiled_at": _now_iso()}
if result.asset_behavior is not None:
entry["behavior"] = asdict(result.asset_behavior)
store[sym] = json.loads(json.dumps(entry, default=str))
rec.profile_ref = sym
done += 1
with open(path, "w", encoding="utf-8") as fh:
json.dump(store, fh, indent=1, sort_keys=True, default=str)
return done, failed
# ── step 3: the runtime call ───────────────────────────────────────────────
def init_asset_universe(execution_exchange: str = "BINGX_VST",
directory: AssetDirectory | None = None) -> frozenset[str]:
"""Tradable canonical symbols for the execution exchange in use."""
directory = directory or AssetDirectory()
return frozenset(directory.symbols_for_exchange(execution_exchange))
def main(argv: list[str] | None = None) -> int:
p = argparse.ArgumentParser(description="UV asset-universe init (MALKHUT store)")
p.add_argument("--seed-from-scans", action="store_true")
p.add_argument("--refresh-vst", action="store_true")
p.add_argument("--onboard", action="store_true",
help="full-taxonomy compile via MALKHUT AssetCompiler (rate-limited)")
p.add_argument("--show", action="store_true")
a = p.parse_args(argv)
d = AssetDirectory()
if a.seed_from_scans:
n = import_blue_universe(d, blue_universe_from_scans())
d.save()
print(f"seeded {n} BINANCE listings from {FEED_SOURCE} -> {d.path}")
if a.refresh_vst:
t, o = refresh_vst_listings(d, fetch_vst_contracts())
d.save()
print(f"BINGX_VST refresh: {t} TRADING, {o} OFFLINE")
if o:
print("OFFLINE:", ", ".join(d.symbols_for_exchange(
"BINGX_VST", status=ListingStatus.OFFLINE)))
if a.onboard:
done, failed = onboard_profiles(d)
d.save()
print(f"onboarded taxonomy: {done} compiled, {failed} failed")
if a.show:
uni = init_asset_universe("BINGX_VST", d)
print(f"directory: {len(d)} assets @ {d.path}")
print(f"BINGX_VST tradable universe ({len(uni)}):", ", ".join(sorted(uni)))
return 0
if __name__ == "__main__":
sys.exit(main())