Update higher timeframes from ticks, and count volume-only trades
Two things kept the chart quieter than the feed. Higher timeframes only moved once a minute. Tick bars are 1m and the socket filters bar events by the subscriber's timeframe, so on the hourly chart every tick was discarded and only a closed minute passing through the aggregator showed up. They cannot simply be fed to the aggregator — it accumulates with current.v += incoming.v, so the same forming minute re-sent on each tick would add its volume to every higher timeframe again and again. provisional_higher combines the aggregator's committed state with the live minute instead, without mutating it; the next closed minute goes through normally and replaces the result, because the store keys on the bucket timestamp. A test pins the behaviour: five ticks in one minute leave the hour's volume at closed plus live, counted exactly once. Trades known only by their volume were skipped. Level 1 resends only changed fields, so some trades carry a trade stamp and a moved TOTAL_VOLUME with neither LAST_PRICE nor LAST_SIZE. Those now count, with size left at zero rather than guessed from the volume delta — CHART_FUTURES replaces the minute's volume with the exchange's own figure moments later, and two ways of counting the same trades is how double counting starts. Measured: 66 to 74 updates per 90s. The tick throttle drops to 0.25s, which no longer binds. Measured in regular hours the gaps between updates are whole multiples of 1.005s — 2.01, 3.02, 4.03 — which is Schwab conflating LEVEL_ONE_FUTURES to one update per second per symbol. One per second is the source's ceiling, not ours; the longer gaps are seconds in which their feed carried no trade. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
parent
d1056ed486
commit
bb84b6e73f
5 changed files with 93 additions and 4 deletions
|
|
@ -44,8 +44,10 @@ class Settings(BaseSettings):
|
||||||
schwab_symbol: str = "/ES"
|
schwab_symbol: str = "/ES"
|
||||||
# Seconds between forming-bar emissions built from LEVEL_ONE_FUTURES ticks.
|
# Seconds between forming-bar emissions built from LEVEL_ONE_FUTURES ticks.
|
||||||
# Set negative to drop the Level 1 subscription and take closed minute bars
|
# Set negative to drop the Level 1 subscription and take closed minute bars
|
||||||
# only. 1.0 is a candle that visibly moves without a broadcast per trade.
|
# only. Measured in regular hours, /ES supplies a price-changing trade far
|
||||||
schwab_tick_seconds: float = 1.0
|
# faster than this, so the value is the update rate: at 1.0 the throttle was
|
||||||
|
# the limiter and the chart felt sluggish.
|
||||||
|
schwab_tick_seconds: float = 0.25
|
||||||
confluence_min_score: float = 28
|
confluence_min_score: float = 28
|
||||||
# Four hours, chosen from the sweep in scripts/calibrate_alerts.py. Suppression is
|
# Four hours, chosen from the sweep in scripts/calibrate_alerts.py. Suppression is
|
||||||
# per price zone, so an unrelated zone still alerts immediately; this only
|
# per price zone, so an unrelated zone still alerts immediately; this only
|
||||||
|
|
|
||||||
|
|
@ -37,6 +37,7 @@ FIELD_VOLUME = "VOLUME"
|
||||||
FIELD_LAST_PRICE = "LAST_PRICE"
|
FIELD_LAST_PRICE = "LAST_PRICE"
|
||||||
FIELD_LAST_SIZE = "LAST_SIZE"
|
FIELD_LAST_SIZE = "LAST_SIZE"
|
||||||
FIELD_TRADE_TIME = "TRADE_TIME_MILLIS"
|
FIELD_TRADE_TIME = "TRADE_TIME_MILLIS"
|
||||||
|
FIELD_TOTAL_VOLUME = "TOTAL_VOLUME"
|
||||||
|
|
||||||
|
|
||||||
def parse_level_one(message: dict) -> list[tuple[int, float | None, int]]:
|
def parse_level_one(message: dict) -> list[tuple[int, float | None, int]]:
|
||||||
|
|
@ -57,7 +58,13 @@ def parse_level_one(message: dict) -> list[tuple[int, float | None, int]]:
|
||||||
price = content.get(FIELD_LAST_PRICE)
|
price = content.get(FIELD_LAST_PRICE)
|
||||||
size = content.get(FIELD_LAST_SIZE)
|
size = content.get(FIELD_LAST_SIZE)
|
||||||
traded_at = content.get(FIELD_TRADE_TIME)
|
traded_at = content.get(FIELD_TRADE_TIME)
|
||||||
if price is None and (size is None or traded_at is None):
|
# A trade stamp alongside a moved cumulative volume is a trade even when
|
||||||
|
# neither the price nor the size field was resent. Size is left at zero
|
||||||
|
# rather than guessed from the volume delta: CHART_FUTURES replaces the
|
||||||
|
# minute's volume with the exchange's own figure a moment later, and two
|
||||||
|
# ways of counting the same trades is how double counting starts.
|
||||||
|
traded = size is not None or content.get(FIELD_TOTAL_VOLUME) is not None
|
||||||
|
if price is None and not (traded and traded_at is not None):
|
||||||
continue
|
continue
|
||||||
millis = traded_at
|
millis = traded_at
|
||||||
if millis is None:
|
if millis is None:
|
||||||
|
|
|
||||||
|
|
@ -1,10 +1,11 @@
|
||||||
import asyncio
|
import asyncio
|
||||||
import logging
|
import logging
|
||||||
from dataclasses import dataclass, field
|
from dataclasses import dataclass, field, replace
|
||||||
|
|
||||||
from app.analysis.alerts import Alert, AlertEngine
|
from app.analysis.alerts import Alert, AlertEngine
|
||||||
from app.bars.models import Bar, Timeframe
|
from app.bars.models import Bar, Timeframe
|
||||||
from app.bars.aggregator import Aggregator
|
from app.bars.aggregator import Aggregator
|
||||||
|
from app.bars.session import bucket_start
|
||||||
from app.analysis.bar_space import price_in_bar_space
|
from app.analysis.bar_space import price_in_bar_space
|
||||||
from app.analysis.horizontals import build_prior_day_levels
|
from app.analysis.horizontals import build_prior_day_levels
|
||||||
from app.analysis.levels import Level
|
from app.analysis.levels import Level
|
||||||
|
|
@ -66,6 +67,9 @@ class Runtime:
|
||||||
self.store.put(bar)
|
self.store.put(bar)
|
||||||
self.price = bar.c
|
self.price = bar.c
|
||||||
self.broadcast({"type": "bar", "bar": bar})
|
self.broadcast({"type": "bar", "bar": bar})
|
||||||
|
for provisional in self.provisional_higher(bar):
|
||||||
|
self.store.put(provisional)
|
||||||
|
self.broadcast({"type": "bar", "bar": provisional})
|
||||||
return
|
return
|
||||||
|
|
||||||
evaluate_alerts = False
|
evaluate_alerts = False
|
||||||
|
|
@ -85,6 +89,41 @@ class Runtime:
|
||||||
self.rebuild_levels()
|
self.rebuild_levels()
|
||||||
self.rebuild_clusters(evaluate_alerts=True)
|
self.rebuild_clusters(evaluate_alerts=True)
|
||||||
|
|
||||||
|
def provisional_higher(self, bar: Bar) -> list[Bar]:
|
||||||
|
"""Higher-timeframe bars including the minute still being traded.
|
||||||
|
|
||||||
|
The aggregator cannot be asked for these: it accumulates with
|
||||||
|
``current.v += incoming.v``, so re-feeding the same forming minute on
|
||||||
|
every tick would add its volume to each higher timeframe again and
|
||||||
|
again. Its committed state — every minute that has actually closed — is
|
||||||
|
combined with the live minute here instead, without mutating it. The
|
||||||
|
next closed minute goes through the aggregator normally and replaces
|
||||||
|
what this produced, because the store keys on the bucket's timestamp.
|
||||||
|
"""
|
||||||
|
out: list[Bar] = []
|
||||||
|
for tf in self.settings.enabled_timeframes:
|
||||||
|
if tf is Timeframe.M1:
|
||||||
|
continue
|
||||||
|
start = bucket_start(bar.t, tf)
|
||||||
|
base = self.aggregator.forming.get(tf)
|
||||||
|
if base is None or start > base.t:
|
||||||
|
# The live minute opens a bucket the aggregator has not started.
|
||||||
|
out.append(replace(bar, tf=tf, t=start, closed=False))
|
||||||
|
continue
|
||||||
|
if start < base.t:
|
||||||
|
continue
|
||||||
|
out.append(
|
||||||
|
replace(
|
||||||
|
base,
|
||||||
|
h=max(base.h, bar.h),
|
||||||
|
l=min(base.l, bar.l),
|
||||||
|
c=bar.c,
|
||||||
|
v=base.v + bar.v,
|
||||||
|
closed=False,
|
||||||
|
)
|
||||||
|
)
|
||||||
|
return out
|
||||||
|
|
||||||
def broadcast(self, event: dict) -> None:
|
def broadcast(self, event: dict) -> None:
|
||||||
for queue in self.subscribers.copy():
|
for queue in self.subscribers.copy():
|
||||||
if queue.full():
|
if queue.full():
|
||||||
|
|
|
||||||
|
|
@ -72,3 +72,32 @@ def test_blank_topic_sends_nothing(tmp_path, monkeypatch):
|
||||||
monkeypatch.setattr("app.runtime.send_ntfy", record)
|
monkeypatch.setattr("app.runtime.send_ntfy", record)
|
||||||
asyncio.run(instance.notify("anything"))
|
asyncio.run(instance.notify("anything"))
|
||||||
assert calls == [""] # send_ntfy itself is the one that short-circuits
|
assert calls == [""] # send_ntfy itself is the one that short-circuits
|
||||||
|
|
||||||
|
|
||||||
|
def test_tick_bars_update_higher_timeframes_without_doubling_volume(tmp_path):
|
||||||
|
# The aggregator accumulates with `current.v += incoming.v`, so a forming
|
||||||
|
# minute re-sent on every tick would add its volume to each higher
|
||||||
|
# timeframe again and again. The provisional path must combine, not
|
||||||
|
# accumulate: five ticks in one minute leave the hour's volume equal to the
|
||||||
|
# closed minutes plus the live one, exactly once.
|
||||||
|
from app.bars.models import Bar
|
||||||
|
|
||||||
|
instance = runtime(tmp_path)
|
||||||
|
|
||||||
|
def minute(t, close, volume, closed=True):
|
||||||
|
return Bar(tf=Timeframe.M1, t=t, o=close, h=close, l=close, c=close,
|
||||||
|
v=volume, closed=closed, symbol="/ES", source="test")
|
||||||
|
|
||||||
|
base = 1786356000 # top of an hour
|
||||||
|
asyncio.run(instance.on_bar(minute(base, 100.0, 10)))
|
||||||
|
asyncio.run(instance.on_bar(minute(base + 60, 101.0, 20)))
|
||||||
|
settled = [b for b in instance.store.get(Timeframe.H1) if b.t == base][-1].v
|
||||||
|
assert settled == 30
|
||||||
|
|
||||||
|
for _ in range(5):
|
||||||
|
asyncio.run(instance.on_bar(minute(base + 120, 102.0, 7, closed=False)))
|
||||||
|
|
||||||
|
hour = [b for b in instance.store.get(Timeframe.H1) if b.t == base][-1]
|
||||||
|
assert hour.v == 37, "the live minute's volume must be added once, not per tick"
|
||||||
|
assert hour.c == 102.0
|
||||||
|
assert hour.closed is False
|
||||||
|
|
|
||||||
|
|
@ -267,3 +267,15 @@ def test_a_same_price_trade_keeps_its_volume():
|
||||||
|
|
||||||
def test_a_quote_with_neither_price_nor_trade_is_still_skipped():
|
def test_a_quote_with_neither_price_nor_trade_is_still_skipped():
|
||||||
assert parse_level_one({"content": [{"key": "/ES", "BID_SIZE": 12}]}) == []
|
assert parse_level_one({"content": [{"key": "/ES", "BID_SIZE": 12}]}) == []
|
||||||
|
|
||||||
|
|
||||||
|
def test_a_trade_known_only_by_its_volume_still_counts():
|
||||||
|
# Neither price nor size resent, but the trade stamp moved and cumulative
|
||||||
|
# volume rose: a trade happened, and the candle should learn about it.
|
||||||
|
volume_only = {
|
||||||
|
"content": [
|
||||||
|
{"key": "/ES", "TRADE_TIME_MILLIS": 1786356932000, "TOTAL_VOLUME": 41,
|
||||||
|
"BID_SIZE": 8}
|
||||||
|
]
|
||||||
|
}
|
||||||
|
assert parse_level_one(volume_only) == [(1786356932000, None, 0)]
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue