Compare commits
3 commits
d1056ed486
...
bec445b599
| Author | SHA1 | Date | |
|---|---|---|---|
| bec445b599 | |||
| 16932f85ec | |||
| bb84b6e73f |
7 changed files with 188 additions and 10 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():
|
||||||
|
|
|
||||||
|
|
@ -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.
|
||||||
|
|
|
||||||
|
|
@ -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 };
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -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