diff --git a/app/config.py b/app/config.py index c96df35..a407a3e 100644 --- a/app/config.py +++ b/app/config.py @@ -44,8 +44,10 @@ class Settings(BaseSettings): schwab_symbol: str = "/ES" # Seconds between forming-bar emissions built from LEVEL_ONE_FUTURES ticks. # 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. - schwab_tick_seconds: float = 1.0 + # only. Measured in regular hours, /ES supplies a price-changing trade far + # 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 # 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 diff --git a/app/market/schwab.py b/app/market/schwab.py index 27ba965..c7e7c19 100644 --- a/app/market/schwab.py +++ b/app/market/schwab.py @@ -37,6 +37,7 @@ FIELD_VOLUME = "VOLUME" FIELD_LAST_PRICE = "LAST_PRICE" FIELD_LAST_SIZE = "LAST_SIZE" FIELD_TRADE_TIME = "TRADE_TIME_MILLIS" +FIELD_TOTAL_VOLUME = "TOTAL_VOLUME" 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) size = content.get(FIELD_LAST_SIZE) 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 millis = traded_at if millis is None: diff --git a/app/runtime.py b/app/runtime.py index 4dbae0b..9e5d397 100644 --- a/app/runtime.py +++ b/app/runtime.py @@ -1,10 +1,11 @@ import asyncio import logging -from dataclasses import dataclass, field +from dataclasses import dataclass, field, replace from app.analysis.alerts import Alert, AlertEngine from app.bars.models import Bar, Timeframe 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.horizontals import build_prior_day_levels from app.analysis.levels import Level @@ -66,6 +67,9 @@ class Runtime: self.store.put(bar) self.price = bar.c self.broadcast({"type": "bar", "bar": bar}) + for provisional in self.provisional_higher(bar): + self.store.put(provisional) + self.broadcast({"type": "bar", "bar": provisional}) return evaluate_alerts = False @@ -85,6 +89,41 @@ class Runtime: self.rebuild_levels() 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: for queue in self.subscribers.copy(): if queue.full(): diff --git a/tests/test_runtime_alerts.py b/tests/test_runtime_alerts.py index 6e3f2c1..2ca27cc 100644 --- a/tests/test_runtime_alerts.py +++ b/tests/test_runtime_alerts.py @@ -72,3 +72,32 @@ def test_blank_topic_sends_nothing(tmp_path, monkeypatch): monkeypatch.setattr("app.runtime.send_ntfy", record) asyncio.run(instance.notify("anything")) 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 diff --git a/tests/test_schwab_source.py b/tests/test_schwab_source.py index 9dbce2c..06cabd6 100644 --- a/tests/test_schwab_source.py +++ b/tests/test_schwab_source.py @@ -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(): 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)]