Compare commits
No commits in common. "bec445b5990c20f5ca9b7f45cbc3593e2579ad9e" and "d1056ed4869a0320e5e662d48fbd7257fe86b260" have entirely different histories.
bec445b599
...
d1056ed486
7 changed files with 10 additions and 188 deletions
|
|
@ -44,10 +44,8 @@ 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. 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
|
||||
# only. 1.0 is a candle that visibly moves without a broadcast per trade.
|
||||
schwab_tick_seconds: float = 1.0
|
||||
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
|
||||
|
|
|
|||
|
|
@ -37,7 +37,6 @@ 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]]:
|
||||
|
|
@ -58,13 +57,7 @@ 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)
|
||||
# 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):
|
||||
if price is None and (size is None or traded_at is None):
|
||||
continue
|
||||
millis = traded_at
|
||||
if millis is None:
|
||||
|
|
|
|||
|
|
@ -1,11 +1,10 @@
|
|||
import asyncio
|
||||
import logging
|
||||
from dataclasses import dataclass, field, replace
|
||||
from dataclasses import dataclass, field
|
||||
|
||||
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
|
||||
|
|
@ -67,9 +66,6 @@ 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
|
||||
|
|
@ -89,41 +85,6 @@ 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():
|
||||
|
|
|
|||
|
|
@ -1380,39 +1380,3 @@ trendline 147 points away.
|
|||
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
|
||||
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,38 +82,12 @@ class ConfluenceChart {
|
|||
}),
|
||||
},
|
||||
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, {
|
||||
upColor: '#27825c', downColor: '#bd4545', borderVisible: true,
|
||||
borderUpColor: '#1d6849', borderDownColor: '#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.chart.applyOptions({ width: el.clientWidth, height: el.clientHeight });
|
||||
requestAnimationFrame(() => this.renderAnchorHandles());
|
||||
|
|
@ -197,10 +171,7 @@ class ConfluenceChart {
|
|||
|
||||
setBars(bars) {
|
||||
this.bars = bars;
|
||||
const candleData = bars.map(this.toCandle);
|
||||
this.candles.setData(candleData);
|
||||
this.leftMirror.setData(candleData);
|
||||
this.volume.setData(bars.map(ConfluenceChart.toVolume));
|
||||
this.candles.setData(bars.map(this.toCandle));
|
||||
// Anchored by time, not by logical index. A logical index addresses the
|
||||
// chart's *shared* scale — the union of every series' time points — not
|
||||
// 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];
|
||||
if (last && bar.t < last.t) return;
|
||||
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;
|
||||
else this.bars.push(bar);
|
||||
this.renderAnchorHandles();
|
||||
|
|
@ -246,9 +215,9 @@ class ConfluenceChart {
|
|||
syncPriceLines(levels) {
|
||||
const flat = levels.filter(level => ConfluenceChart.isFlat(level) && !level.hidden);
|
||||
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)) {
|
||||
entry.host.removePriceLine(entry.line);
|
||||
this.candles.removePriceLine(line);
|
||||
this.priceLines.delete(id);
|
||||
}
|
||||
}
|
||||
|
|
@ -262,14 +231,8 @@ class ConfluenceChart {
|
|||
title: level.label,
|
||||
};
|
||||
const existing = this.priceLines.get(level.id);
|
||||
if (existing) existing.line.applyOptions(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) });
|
||||
}
|
||||
if (existing) existing.applyOptions(options);
|
||||
else this.priceLines.set(level.id, this.candles.createPriceLine(options));
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -339,14 +302,7 @@ class ConfluenceChart {
|
|||
autoscaleInfoProvider: () => null,
|
||||
};
|
||||
if (!entry) {
|
||||
// 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 }),
|
||||
};
|
||||
entry = { series: this.chart.addSeries(LightweightCharts.LineSeries, options) };
|
||||
this.levelSeries.set(level.id, entry);
|
||||
} else {
|
||||
entry.series.applyOptions(options);
|
||||
|
|
@ -740,15 +696,6 @@ class ConfluenceChart {
|
|||
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) {
|
||||
return { time: bar.t, open: bar.o, high: bar.h, low: bar.l, close: bar.c };
|
||||
}
|
||||
|
|
|
|||
|
|
@ -72,32 +72,3 @@ 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
|
||||
|
|
|
|||
|
|
@ -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():
|
||||
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