SniperGold_ML/engine/verify/vaggregate.py

127 lines
4.9 KiB
Python

"""Independent aggregator (Layer 12, spec 18.1).
Re-implementation of the bar aggregation rules (spec 11) with distinct code:
aggregation state is kept per timeframe in flat dicts, finalization uses an
explicit bar-emission event, and canonical bar serialization uses f-strings
(spec 12.4). Must produce byte-identical rows to engine/aggregate.py.
"""
from .. import versions
def v_tick_serialize(ts, bid, ask, vol):
return ("%d|%d|%d|%d\n" % (ts, bid, ask, vol)).encode("utf-8")
_BAR_FIELDS14 = (0, 1, 4, 5, 6, 7, 8, 9, 10, 11, 12, 14, 15, 16)
def v_bar_serialize(row):
"""Canonical bar serialization (spec 12.4) over a full 17-field CBS row:
bar_idx|period_id|open_u2|high_u2|low_u2|close_u2|tick_count|vol_sum|
spread_min_u|spread_max_u|spread_sum_u|is_final|first_src_line|
last_src_line, is_final as 1/0."""
sel = [row[i] for i in _BAR_FIELDS14]
(bar_idx, period_id, open_u2, high_u2, low_u2, close_u2, tick_count,
vol_sum, spread_min_u, spread_max_u, spread_sum_u, is_final,
first_src_line, last_src_line) = sel
return ("%d|%d|%d|%d|%d|%d|%d|%d|%d|%d|%d|%d|%d|%d\n"
% (bar_idx, period_id, open_u2, high_u2, low_u2, close_u2,
tick_count, vol_sum, spread_min_u, spread_max_u, spread_sum_u,
1 if is_final else 0, first_src_line, last_src_line)).encode("utf-8")
def v_aggregate(ticks, carries, cfg):
"""Independent sequential aggregation.
``ticks``: iterable of (ts, bid, ask, spread, vol, src_line, rec_ord).
``carries``: {tf: carry_dict or None}.
Returns (bars_by_tf, carries_out, counts_by_tf) where bars_by_tf maps tf
to the list of finalized bar rows; the EOF partial bars are appended with
is_final False (caller calls v_aggregate_eof for the tail).
"""
state = {}
for tf in cfg["timeframes"]:
c = carries.get(tf) if carries else None
state[tf] = {
"carry": dict(c) if c else None,
"count": 0,
"rows": [],
"last_period": None,
}
for tick in ticks:
ts, bid, ask, spread, vol, src_line, _ = tick
mid = bid + ask
for tf in cfg["timeframes"]:
period_ms = versions.TIMEFRAMES[tf]
s = state[tf]
pid = ts // period_ms
if s["carry"] is None:
s["carry"] = {
"pid": pid, "start_ms": pid * period_ms,
"o": mid, "h": mid, "l": mid, "c": mid,
"n": 1, "v": vol, "smn": spread, "smx": spread,
"ss": spread, "last_ts": ts, "last_line": src_line,
"first_line": src_line,
}
continue
car = s["carry"]
if pid == car["pid"]:
car["h"] = mid if mid > car["h"] else car["h"]
car["l"] = mid if mid < car["l"] else car["l"]
car["c"] = mid
car["n"] += 1
car["v"] += vol
car["smn"] = spread if spread < car["smn"] else car["smn"]
car["smx"] = spread if spread > car["smx"] else car["smx"]
car["ss"] += spread
car["last_ts"] = ts
car["last_line"] = src_line
continue
assert pid > car["pid"], "v_aggregate: order violation"
assert s["last_period"] is None or car["pid"] > s["last_period"]
s["last_period"] = car["pid"]
s["rows"].append(_v_row(car, s["count"], period_ms, True))
s["count"] += 1
s["carry"] = {
"pid": pid, "start_ms": pid * period_ms,
"o": mid, "h": mid, "l": mid, "c": mid,
"n": 1, "v": vol, "smn": spread, "smx": spread,
"ss": spread, "last_ts": ts, "last_line": src_line,
"first_line": src_line,
}
return ({tf: state[tf]["rows"] for tf in cfg["timeframes"]},
{tf: state[tf]["carry"] for tf in cfg["timeframes"]},
{tf: state[tf]["count"] for tf in cfg["timeframes"]})
def v_aggregate_eof(cfg, carries, counts):
"""Flush EOF partial bars with is_final=False (spec 11.5)."""
rows = {}
counts_out = dict(counts)
carries_out = {}
for tf in cfg["timeframes"]:
period_ms = versions.TIMEFRAMES[tf]
carry = carries.get(tf)
r = []
if carry is not None:
r.append(_v_row(carry, counts_out[tf], period_ms, False))
counts_out[tf] += 1
carry = None
rows[tf] = r
carries_out[tf] = carry
return rows, carries_out, counts_out
def _v_row(car, bar_idx, period_ms, is_final):
return (bar_idx,
car["pid"],
car["start_ms"],
car["start_ms"] + period_ms,
car["o"], car["h"], car["l"], car["c"],
car["n"], car["v"],
car["smn"], car["smx"], car["ss"],
car["ss"] // car["n"],
is_final,
car["first_line"],
car["last_line"])