SniperGold_ML/ml/p3/p3_s251_external_ingest/s251_aggregate.py

121 lines
No EOL
4 KiB
Python

# -*- coding: utf-8 -*-
"""P3-S25.1 STREAMING AGGREGATION (PRODUCTION PATH).
Deterministic bucket membership:
M15 bucket = floor(epoch_utc / 900 ) * 900 -> boundary: tick AT a 15-min
boundary belongs to the bucket beginning at that boundary.
M30 bucket = floor(M15_open / 1800) * 1800 -> two M15 sub-buckets.
OHLC uses the reference price = bid (consistent with P3-S25: bid<=last<=ask).
open=first bid in bucket, high=max, low=min, close=last, tick_count.
A bucket is 'completed' when the first tick of a strictly later bucket is
observed; a bucket still open at end-of-input is 'incomplete' (reported
separately, never forward-filled). Only current open bucket state + a bounded
buffer of this chunk's completed rows is retained -> bounded memory.
"""
from __future__ import annotations
import s251_config as CFG
M15 = CFG.M15
M30 = CFG.M30
M15_FIELDS = ("ts", "open", "high", "low", "close", "ticks",
"first_ts", "last_ts")
M30_FIELDS = ("ts", "open", "high", "low", "close", "n_m15")
def _new_m15(bk, price, epo):
return {"ts": bk, "open": price, "high": price, "low": price,
"close": price, "ticks": 1, "first_ts": epo, "last_ts": epo}
def _m15_carry_ts(state):
return state["ts"] if state else None
def process_chunk(ticks, m15_carry, m30_carry):
"""Process ticks (list of (epo, price), ascending) with carry states.
ticks should already be sorted ascending by epoch (see ingest).
m15_carry: dict open M15 bucket from prior chunk or None.
m30_carry: dict open M30 bucket from prior chunk or None.
Returns dict: m15_rows, m30_rows (completed rows emitted in this chunk),
m15_carry, m30_carry (open buckets to persist across chunks).
"""
m15_rows = []
m30_rows = []
cur = m15_carry
cur30 = m30_carry
if cur is not None:
cur["ts"] = (int(cur["ts"]) // M15) * M15 # normalise carry ts
for epo, price in ticks:
bk = (int(epo) // M15) * M15
if cur is None:
cur = _new_m15(bk, price, epo)
continue
if bk == cur["ts"]:
cur["ticks"] += 1
cur["last_ts"] = epo
if price > cur["high"]:
cur["high"] = price
if price < cur["low"]:
cur["low"] = price
cur["close"] = price
continue
# close cur bucket (completed), fold to M30, open new
done = cur
cur = _new_m15(bk, price, epo)
m15_rows.append(_bundle_m15(done))
cur30 = _fold_m15_to_m30(done, cur30, m30_rows)
# if this new M15 crosses into a later M30 window, the window that
# just received `done` is complete -> emit it.
if (int(done["ts"]) // M30) != (bk // M30):
if cur30 is not None and cur30["ts"] == (int(done["ts"]) // M30) * M30:
m30_rows.append(_bundle_m30(cur30))
cur30 = None
return {"m15_rows": m15_rows, "m30_rows": m30_rows,
"m15_carry": cur, "m30_carry": cur30}
def _bundle_m15(c):
return (c["ts"], c["open"], c["high"], c["low"], c["close"],
c["ticks"], c["first_ts"], c["last_ts"])
def _bundle_m30(c):
return (c["ts"], c["open"], c["high"], c["low"], c["close"], c["n_m15"])
def _fold_m15_to_m30(m15b, cur30, out30):
ts, o, h, l, c, cnt, f, lst = _bundle_m15(m15b)
bk = (int(ts) // M30) * M30
if cur30 is None:
return {"ts": bk, "open": o, "high": h, "low": l, "close": c,
"n_m15": 1}
if bk == cur30["ts"]:
cur30["n_m15"] += 1
if h > cur30["high"]:
cur30["high"] = h
if l < cur30["low"]:
cur30["low"] = l
cur30["close"] = c
return cur30
done = cur30
out30.append((done["ts"], done["open"], done["high"], done["low"],
done["close"], done["n_m15"]))
return {"ts": bk, "open": o, "high": h, "low": l, "close": c, "n_m15": 1}
def bundle_m15(c):
return _bundle_m15(c)
def bundle_m30(c):
if c is None:
return None
return _bundle_m30(c)