Compare commits

..

No commits in common. "bec445b5990c20f5ca9b7f45cbc3593e2579ad9e" and "d1056ed4869a0320e5e662d48fbd7257fe86b260" have entirely different histories.

7 changed files with 10 additions and 188 deletions

View file

@ -44,10 +44,8 @@ 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. Measured in regular hours, /ES supplies a price-changing trade far # only. 1.0 is a candle that visibly moves without a broadcast per trade.
# faster than this, so the value is the update rate: at 1.0 the throttle was schwab_tick_seconds: float = 1.0
# 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

View file

@ -37,7 +37,6 @@ 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]]:
@ -58,13 +57,7 @@ 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)
# A trade stamp alongside a moved cumulative volume is a trade even when if price is None and (size is None or traded_at is None):
# 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:

View file

@ -1,11 +1,10 @@
import asyncio import asyncio
import logging import logging
from dataclasses import dataclass, field, replace from dataclasses import dataclass, field
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
@ -67,9 +66,6 @@ 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
@ -89,41 +85,6 @@ 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():

View file

@ -1380,39 +1380,3 @@ trendline 147 points away.
One limit worth knowing: tick *placement* is still computed on UTC days, so the One limit worth knowing: tick *placement* is still computed on UTC days, so the
day-change divider sits at 00:00 UTC rather than local midnight, labelled with day-change divider sits at 00:00 UTC rather than local midnight, labelled with
the local date. The labels are right; the divider is in the UTC place. the local date. The labels are right; the divider is in the UTC place.
### 2026-08-10 (afternoon) — update rate, the left scale, and volume
**Schwab conflates Level 1 to one update per second.** Chasing "still slow in
market hours" ended at a hard ceiling rather than a bug. In regular hours the
gaps between updates are whole multiples of 1.005s — 2.01, 3.02, 4.03 — which
only happens if the source emits on a one-second cadence and some seconds carry
no trade. `SCHWAB_TICK_SECONDS` was the limiter at 1.0 and is now 0.25, where it
no longer binds. **One update per second is the source's ceiling.** Anything
faster would mean inventing prices between trades, which a chart must not do.
Two real losses were found on the way and fixed:
- Higher timeframes only moved once a minute, because tick bars are 1m and the
socket filters by subscriber timeframe. `Runtime.provisional_higher` now
combines the aggregator's committed state with the live minute — without
mutating it, since the aggregator accumulates volume and would double count.
- Trades carrying only a trade stamp and a moved `TOTAL_VOLUME` — no
`LAST_PRICE`, no `LAST_SIZE` — were skipped. 66 → 74 updates per 90s.
**Daily context moved to the left price scale.** The right had prior-day levels,
session VWAP and five daily MAs competing with the live price and hand-drawn
intraday levels. The trap: a price scale takes its range from the series on it,
so moving levels across draws them against a different range and puts them at
the wrong height. A transparent candlestick mirror on the left scale feeds it
exactly the right scale's input; verified as a zero-pixel delta between the two.
Hand-drawn levels stay right, which is the space being cleared. `priceScaleId`
is fixed at series creation, so it is passed at construction and kept out of the
options reapplied afterwards.
**Volume is finally drawn.** It travelled the entire pipeline — parsed from both
Schwab services, aggregated, stored, broadcast in every bar — and nothing
rendered it. Now an overlay histogram on its own hidden scale in the bottom
fifth. An overlay rather than a pane, and emphatically not the price scale:
volumes are five figures against four-figure prices, and sharing a scale would
flatten the candles to a line.

View file

@ -82,38 +82,12 @@ class ConfluenceChart {
}), }),
}, },
rightPriceScale: { borderVisible: false }, rightPriceScale: { borderVisible: false },
// Daily context lives on the left, intraday on the right.
leftPriceScale: { visible: true, borderVisible: false },
}); });
this.candles = this.chart.addSeries(LightweightCharts.CandlestickSeries, { this.candles = this.chart.addSeries(LightweightCharts.CandlestickSeries, {
upColor: '#27825c', downColor: '#bd4545', borderVisible: true, upColor: '#27825c', downColor: '#bd4545', borderVisible: true,
borderUpColor: '#1d6849', borderDownColor: '#963737', borderUpColor: '#1d6849', borderDownColor: '#963737',
wickUpColor: '#1d6849', wickDownColor: '#963737', wickUpColor: '#1d6849', wickDownColor: '#963737',
}); });
// A price scale derives its range from the series on it, so levels moved to
// the left would be drawn against a different range and sit at the wrong
// height. This transparent copy of the candles gives the left scale exactly
// the same input as the right, which keeps one price at one y.
this.leftMirror = this.chart.addSeries(LightweightCharts.CandlestickSeries, {
priceScaleId: 'left',
upColor: 'transparent', downColor: 'transparent', borderVisible: false,
wickUpColor: 'transparent', wickDownColor: 'transparent',
lastValueVisible: false, priceLineVisible: false,
});
// Volume as an overlay on its own hidden scale, confined to the bottom
// fifth. An overlay rather than a pane so it cannot alter the price scale:
// volumes are five figures and prices four, and sharing a scale would
// flatten the candles into a line.
this.volume = this.chart.addSeries(LightweightCharts.HistogramSeries, {
priceScaleId: 'volume',
priceFormat: { type: 'volume' },
lastValueVisible: false,
priceLineVisible: false,
});
this.chart.priceScale('volume').applyOptions({
scaleMargins: { top: 0.82, bottom: 0 },
visible: false,
});
this.resizeObserver = new ResizeObserver(() => { this.resizeObserver = new ResizeObserver(() => {
this.chart.applyOptions({ width: el.clientWidth, height: el.clientHeight }); this.chart.applyOptions({ width: el.clientWidth, height: el.clientHeight });
requestAnimationFrame(() => this.renderAnchorHandles()); requestAnimationFrame(() => this.renderAnchorHandles());
@ -197,10 +171,7 @@ class ConfluenceChart {
setBars(bars) { setBars(bars) {
this.bars = bars; this.bars = bars;
const candleData = bars.map(this.toCandle); this.candles.setData(bars.map(this.toCandle));
this.candles.setData(candleData);
this.leftMirror.setData(candleData);
this.volume.setData(bars.map(ConfluenceChart.toVolume));
// Anchored by time, not by logical index. A logical index addresses the // Anchored by time, not by logical index. A logical index addresses the
// chart's *shared* scale — the union of every series' time points — not // chart's *shared* scale — the union of every series' time points — not
// this array. The daily MAs land straight after with hundreds of points // this array. The daily MAs land straight after with hundreds of points
@ -228,8 +199,6 @@ class ConfluenceChart {
const last = this.bars[this.bars.length - 1]; const last = this.bars[this.bars.length - 1];
if (last && bar.t < last.t) return; if (last && bar.t < last.t) return;
this.candles.update(this.toCandle(bar)); this.candles.update(this.toCandle(bar));
this.leftMirror.update(this.toCandle(bar));
this.volume.update(ConfluenceChart.toVolume(bar));
if (this.bars.length && this.bars[this.bars.length - 1].t === bar.t) this.bars[this.bars.length - 1] = bar; if (this.bars.length && this.bars[this.bars.length - 1].t === bar.t) this.bars[this.bars.length - 1] = bar;
else this.bars.push(bar); else this.bars.push(bar);
this.renderAnchorHandles(); this.renderAnchorHandles();
@ -246,9 +215,9 @@ class ConfluenceChart {
syncPriceLines(levels) { syncPriceLines(levels) {
const flat = levels.filter(level => ConfluenceChart.isFlat(level) && !level.hidden); const flat = levels.filter(level => ConfluenceChart.isFlat(level) && !level.hidden);
const wanted = new Set(flat.map(level => level.id)); const wanted = new Set(flat.map(level => level.id));
for (const [id, entry] of this.priceLines) { for (const [id, line] of this.priceLines) {
if (!wanted.has(id)) { if (!wanted.has(id)) {
entry.host.removePriceLine(entry.line); this.candles.removePriceLine(line);
this.priceLines.delete(id); this.priceLines.delete(id);
} }
} }
@ -262,14 +231,8 @@ class ConfluenceChart {
title: level.label, title: level.label,
}; };
const existing = this.priceLines.get(level.id); const existing = this.priceLines.get(level.id);
if (existing) existing.line.applyOptions(options); if (existing) existing.applyOptions(options);
else { else this.priceLines.set(level.id, this.candles.createPriceLine(options));
// Prior-day levels are daily context and move to the left scale; a
// hand-drawn price level is intraday and keeps the right, which is the
// side being kept clear for it.
const host = level.kind === 'horizontal' ? this.leftMirror : this.candles;
this.priceLines.set(level.id, { host, line: host.createPriceLine(options) });
}
} }
} }
@ -339,14 +302,7 @@ class ConfluenceChart {
autoscaleInfoProvider: () => null, autoscaleInfoProvider: () => null,
}; };
if (!entry) { if (!entry) {
// Daily averages and session VWAP are daily/session context, so their entry = { series: this.chart.addSeries(LightweightCharts.LineSeries, options) };
// last-value labels belong on the left. priceScaleId is fixed at
// creation, which is why it is not in the options reapplied below.
const scale = (isMa || level.kind === 'vwap') ? 'left' : 'right';
entry = {
series: this.chart.addSeries(LightweightCharts.LineSeries,
{ ...options, priceScaleId: scale }),
};
this.levelSeries.set(level.id, entry); this.levelSeries.set(level.id, entry);
} else { } else {
entry.series.applyOptions(options); entry.series.applyOptions(options);
@ -740,15 +696,6 @@ class ConfluenceChart {
if (this.onLineEnd) this.onLineEnd({ ...level }); if (this.onLineEnd) this.onLineEnd({ ...level });
} }
// Tinted by the bar's own direction, muted so the candles stay the subject.
static toVolume(bar) {
return {
time: bar.t,
value: bar.v,
color: bar.c >= bar.o ? 'rgba(39,130,92,.38)' : 'rgba(189,69,69,.38)',
};
}
toCandle(bar) { toCandle(bar) {
return { time: bar.t, open: bar.o, high: bar.h, low: bar.l, close: bar.c }; return { time: bar.t, open: bar.o, high: bar.h, low: bar.l, close: bar.c };
} }

View file

@ -72,32 +72,3 @@ 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

View file

@ -267,15 +267,3 @@ 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)]