Compare commits

...

3 commits

Author SHA1 Message Date
bec445b599 Draw volume, which travelled the whole pipeline unseen
Volume is parsed from both Schwab services, aggregated into every timeframe,
stored and broadcast in every bar payload — and nothing ever drew it.

An overlay histogram on its own hidden scale, confined to the bottom fifth and
tinted by each bar's direction. An overlay rather than a second pane, and
deliberately not on the price scale: volumes are five figures against
four-figure prices, so sharing a scale would flatten the candles into a line.
Verified in a browser that the price scales are untouched — the same price still
maps to the identical coordinate through both.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-10 12:08:59 -05:00
16932f85ec Move daily context to the left price scale
The right scale was crowding: prior-day levels, session VWAP and five daily
moving averages competing with the live price and hand-drawn intraday levels.
Daily and session context moves left, and the right is left for intraday.

A price scale takes its range from the series on it, so moving levels across
would have drawn them against a different range and put them at the wrong
height — the failure this codebase has already paid for once. A transparent
candlestick mirror on the left scale gives it exactly the same input as the
right. Verified in a browser: the same price maps to the identical y coordinate
through both scales, a delta of zero pixels.

Prior-day levels move; hand-drawn price levels stay on the right, since those
are the intraday markers the space is being cleared for. priceScaleId is fixed
when a series is created, so it is passed at construction and left out of the
options reapplied afterwards.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-10 12:07:52 -05:00
bb84b6e73f 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>
2026-08-10 12:07:52 -05:00
7 changed files with 188 additions and 10 deletions

View file

@ -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

View file

@ -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:

View file

@ -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():

View file

@ -1380,3 +1380,39 @@ 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,12 +82,38 @@ 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());
@ -171,7 +197,10 @@ class ConfluenceChart {
setBars(bars) { setBars(bars) {
this.bars = bars; this.bars = bars;
this.candles.setData(bars.map(this.toCandle)); const candleData = 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
@ -199,6 +228,8 @@ 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();
@ -215,9 +246,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, line] of this.priceLines) { for (const [id, entry] of this.priceLines) {
if (!wanted.has(id)) { if (!wanted.has(id)) {
this.candles.removePriceLine(line); entry.host.removePriceLine(entry.line);
this.priceLines.delete(id); this.priceLines.delete(id);
} }
} }
@ -231,8 +262,14 @@ 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.applyOptions(options); if (existing) existing.line.applyOptions(options);
else this.priceLines.set(level.id, this.candles.createPriceLine(options)); else {
// 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) });
}
} }
} }
@ -302,7 +339,14 @@ class ConfluenceChart {
autoscaleInfoProvider: () => null, autoscaleInfoProvider: () => null,
}; };
if (!entry) { if (!entry) {
entry = { series: this.chart.addSeries(LightweightCharts.LineSeries, options) }; // Daily averages and session VWAP are daily/session context, so their
// 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);
@ -696,6 +740,15 @@ 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,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

View file

@ -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)]