The candle still paused for ten to twenty seconds at a time. Instrumenting the raw Level 1 stream settled why: 87 messages in 90 seconds, only 33 carrying LAST_PRICE. Most of the remainder is bid and ask movement, correctly ignored, but a seventh carry LAST_SIZE, TRADE_TIME_MILLIS and TOTAL_VOLUME with no LAST_PRICE — trades that printed at the price of the one before, so the field did not change and Level 1 did not resend it. Requiring LAST_PRICE discarded those trades and their volume with them. parse_level_one now recognises size-plus-trade-time as a trade and returns a null price, which stream() fills from the forming bar. A quote carrying neither a price nor any trade field is still skipped: a bid is not a trade and must not extend a candle's high or low. Measured on the live feed: median gap between updates 3.1s to 2.0s, worst gap 21.5s to 8.1s, roughly 9 updates a minute to 22, and bar volume climbs within the minute instead of standing still. The pauses that remain are the market rather than the pipe. Thin pre-open tape goes seconds without a price-changing trade and then moves several ticks at once, which is what a gap up after a quiet spell is. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
269 lines
8.6 KiB
Python
269 lines
8.6 KiB
Python
import asyncio
|
|
|
|
import pytest
|
|
|
|
from app.bars.models import Timeframe
|
|
from app.config import Settings
|
|
from app.market.factory import live_source, seed_source
|
|
from app.market.schwab import SchwabSource, parse_chart_futures, parse_level_one
|
|
|
|
# Shape taken from a live CHART_FUTURES message, not invented.
|
|
LIVE_MESSAGE = {
|
|
"service": "CHART_FUTURES",
|
|
"timestamp": 1786356976793,
|
|
"command": "SUBS",
|
|
"content": [
|
|
{
|
|
"seq": 50,
|
|
"key": "/ES",
|
|
"CHART_TIME_MILLIS": 1786356900000,
|
|
"OPEN_PRICE": 7786.75,
|
|
"HIGH_PRICE": 7787,
|
|
"LOW_PRICE": 7786.5,
|
|
"CLOSE_PRICE": 7787,
|
|
"VOLUME": 107,
|
|
}
|
|
],
|
|
}
|
|
|
|
|
|
def test_parses_a_live_chart_futures_message():
|
|
bar = parse_chart_futures(LIVE_MESSAGE, "/ES")[0]
|
|
|
|
assert bar.tf is Timeframe.M1
|
|
assert bar.t == 1786356900 # milliseconds down to seconds
|
|
assert (bar.o, bar.h, bar.l, bar.c) == (7786.75, 7787.0, 7786.5, 7787.0)
|
|
assert bar.v == 107
|
|
assert bar.symbol == "/ES"
|
|
assert bar.source == "schwab"
|
|
# The minute has elapsed by the time the message arrives.
|
|
assert bar.closed is True
|
|
|
|
|
|
def test_incomplete_content_is_skipped_not_defaulted():
|
|
# A bar invented from partial data is indistinguishable downstream from a
|
|
# real one, which is worse than having no bar.
|
|
for missing in ("CHART_TIME_MILLIS", "OPEN_PRICE", "CLOSE_PRICE"):
|
|
content = dict(LIVE_MESSAGE["content"][0])
|
|
del content[missing]
|
|
assert parse_chart_futures({"content": [content]}, "/ES") == []
|
|
|
|
|
|
def test_schwab_offers_no_history():
|
|
source = SchwabSource(Settings())
|
|
assert source.supports_history() is False
|
|
assert source.supports_stream() is True
|
|
assert asyncio.run(source.history("/ES", Timeframe.M1, None, None)) == []
|
|
|
|
|
|
def test_seeding_stays_on_yahoo_even_when_live_is_schwab(tmp_path):
|
|
# Schwab has no history, so the pairing is the intended configuration
|
|
# rather than a fallback.
|
|
settings = Settings(
|
|
live_source="schwab",
|
|
seed_source="schwab",
|
|
manual_lines_path=tmp_path / "lines.json",
|
|
)
|
|
assert seed_source(settings).name == "yahoo"
|
|
assert live_source(settings).name == "schwab"
|
|
|
|
|
|
def test_live_symbol_follows_the_live_source():
|
|
# Yahoo says ES=F where Schwab says /ES; seeding always uses the Yahoo one.
|
|
assert Settings(live_source="yahoo").live_symbol == "ES=F"
|
|
assert Settings(live_source="schwab").live_symbol == "/ES"
|
|
|
|
|
|
def test_stream_yields_bars_from_the_socket():
|
|
class FakeStreamClient:
|
|
def __init__(self):
|
|
self.handler = None
|
|
self.subscribed = []
|
|
self.quote_handler = None
|
|
self.quote_subscribed = []
|
|
|
|
async def login(self):
|
|
return None
|
|
|
|
def add_chart_futures_handler(self, handler):
|
|
self.handler = handler
|
|
|
|
async def chart_futures_subs(self, symbols):
|
|
self.subscribed = list(symbols)
|
|
|
|
def add_level_one_futures_handler(self, handler):
|
|
self.quote_handler = handler
|
|
|
|
async def level_one_futures_subs(self, symbols):
|
|
self.quote_subscribed = list(symbols)
|
|
|
|
async def handle_message(self):
|
|
# One message, then idle rather than returning — a real socket
|
|
# never stops on its own.
|
|
if self.handler:
|
|
self.handler(LIVE_MESSAGE)
|
|
self.handler = None
|
|
await asyncio.sleep(3600)
|
|
|
|
fake = FakeStreamClient()
|
|
source = SchwabSource(Settings(), stream_client_factory=lambda: fake)
|
|
|
|
async def first_bar():
|
|
async for bar in source.stream("/ES"):
|
|
return bar
|
|
|
|
bar = asyncio.run(asyncio.wait_for(first_bar(), timeout=10))
|
|
assert bar.c == 7787.0
|
|
assert fake.subscribed == ["/ES"]
|
|
|
|
|
|
def test_socket_failure_surfaces_rather_than_hanging():
|
|
class ExplodingStreamClient:
|
|
async def login(self):
|
|
return None
|
|
|
|
def add_chart_futures_handler(self, handler):
|
|
pass
|
|
|
|
def add_level_one_futures_handler(self, handler):
|
|
pass
|
|
|
|
async def level_one_futures_subs(self, symbols):
|
|
return None
|
|
|
|
async def chart_futures_subs(self, symbols):
|
|
pass
|
|
|
|
async def handle_message(self):
|
|
raise RuntimeError("socket closed")
|
|
|
|
source = SchwabSource(Settings(), stream_client_factory=ExplodingStreamClient)
|
|
|
|
async def drain():
|
|
async for _ in source.stream("/ES"):
|
|
pass
|
|
|
|
with pytest.raises(RuntimeError, match="socket closed"):
|
|
asyncio.run(asyncio.wait_for(drain(), timeout=15))
|
|
|
|
|
|
# Shape taken from a live LEVEL_ONE_FUTURES message.
|
|
QUOTE_MESSAGE = {
|
|
"service": "LEVEL_ONE_FUTURES",
|
|
"command": "SUBS",
|
|
"content": [
|
|
{"key": "/ES", "LAST_PRICE": 7786.25, "LAST_SIZE": 3, "TRADE_TIME_MILLIS": 1786356930000}
|
|
],
|
|
}
|
|
|
|
|
|
def test_parses_a_level_one_trade():
|
|
assert parse_level_one(QUOTE_MESSAGE) == [(1786356930000, 7786.25, 3)]
|
|
|
|
|
|
def test_quotes_without_a_trade_are_skipped():
|
|
# A bid-only update is not a trade and must not extend a candle's range.
|
|
bid_only = {"content": [{"key": "/ES", "BID_PRICE": 7786.0, "ASK_PRICE": 7786.5}]}
|
|
assert parse_level_one(bid_only) == []
|
|
|
|
|
|
def test_ticks_build_an_unclosed_bar_for_the_current_minute():
|
|
class QuotingClient:
|
|
def __init__(self):
|
|
self.quote_handler = None
|
|
|
|
async def login(self):
|
|
return None
|
|
|
|
def add_chart_futures_handler(self, handler):
|
|
pass
|
|
|
|
async def chart_futures_subs(self, symbols):
|
|
return None
|
|
|
|
def add_level_one_futures_handler(self, handler):
|
|
self.quote_handler = handler
|
|
|
|
async def level_one_futures_subs(self, symbols):
|
|
return None
|
|
|
|
async def handle_message(self):
|
|
if self.quote_handler:
|
|
self.quote_handler(QUOTE_MESSAGE)
|
|
self.quote_handler = None
|
|
await asyncio.sleep(3600)
|
|
|
|
source = SchwabSource(Settings(schwab_tick_seconds=0), stream_client_factory=QuotingClient)
|
|
|
|
async def first_bar():
|
|
async for bar in source.stream("/ES"):
|
|
return bar
|
|
|
|
bar = asyncio.run(asyncio.wait_for(first_bar(), timeout=10))
|
|
# Bucketed to its minute, and explicitly not closed — the minute is still
|
|
# running, and a closed flag would let it into the aggregator.
|
|
assert bar.t == 1786356900
|
|
assert bar.closed is False
|
|
assert (bar.o, bar.h, bar.l, bar.c) == (7786.25, 7786.25, 7786.25, 7786.25)
|
|
|
|
|
|
def test_a_tick_for_an_already_closed_minute_is_ignored():
|
|
# CHART_FUTURES is authoritative. A late tick for a minute it has already
|
|
# settled would otherwise overwrite a real bar with a partial one.
|
|
class LateTickClient:
|
|
def __init__(self):
|
|
self.chart_handler = None
|
|
self.quote_handler = None
|
|
|
|
async def login(self):
|
|
return None
|
|
|
|
def add_chart_futures_handler(self, handler):
|
|
self.chart_handler = handler
|
|
|
|
async def chart_futures_subs(self, symbols):
|
|
return None
|
|
|
|
def add_level_one_futures_handler(self, handler):
|
|
self.quote_handler = handler
|
|
|
|
async def level_one_futures_subs(self, symbols):
|
|
return None
|
|
|
|
async def handle_message(self):
|
|
if self.chart_handler:
|
|
self.chart_handler(LIVE_MESSAGE) # closes 1786356900
|
|
self.quote_handler(QUOTE_MESSAGE) # tick inside it
|
|
self.chart_handler = None
|
|
await asyncio.sleep(3600)
|
|
|
|
source = SchwabSource(Settings(schwab_tick_seconds=0), stream_client_factory=LateTickClient)
|
|
|
|
async def two_bars():
|
|
seen = []
|
|
async for bar in source.stream("/ES"):
|
|
seen.append(bar)
|
|
if len(seen) == 1:
|
|
# Give the late tick a chance to be wrongly emitted.
|
|
await asyncio.sleep(0.2)
|
|
break
|
|
return seen
|
|
|
|
seen = asyncio.run(asyncio.wait_for(two_bars(), timeout=10))
|
|
assert [bar.closed for bar in seen] == [True]
|
|
|
|
|
|
def test_a_same_price_trade_keeps_its_volume():
|
|
# Level 1 sends only changed fields, so a trade at the price of the one
|
|
# before carries size and trade time but no LAST_PRICE. Measured live at a
|
|
# fifth of all trades; dropping them lost that volume from the bar.
|
|
same_price = {
|
|
"content": [
|
|
{"key": "/ES", "LAST_SIZE": 4, "TRADE_TIME_MILLIS": 1786356931000, "TOTAL_VOLUME": 9}
|
|
]
|
|
}
|
|
assert parse_level_one(same_price) == [(1786356931000, None, 4)]
|
|
|
|
|
|
def test_a_quote_with_neither_price_nor_trade_is_still_skipped():
|
|
assert parse_level_one({"content": [{"key": "/ES", "BID_SIZE": 12}]}) == []
|