forked from chiki2bum2/SniperGold_ML
121 lines
No EOL
4 KiB
Python
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) |