Compare commits
No commits in common. "fced18528b8f32c00ccc370b4094c6065163fbeb" and "d0a0f9f6d6b5bb38158d20c4c67d15afce3c29d0" have entirely different histories.
fced18528b
...
d0a0f9f6d6
25 changed files with 36 additions and 774 deletions
|
|
@ -14,7 +14,6 @@ deferred fixes. Mobile interaction work also has its own detailed plan in
|
||||||
[`docs/vite_build.md`](docs/vite_build.md). Light/dark theme constraints are
|
[`docs/vite_build.md`](docs/vite_build.md). Light/dark theme constraints are
|
||||||
[`docs/plan_light_dark_themes.md`](docs/plan_light_dark_themes.md).
|
[`docs/plan_light_dark_themes.md`](docs/plan_light_dark_themes.md).
|
||||||
Daily MA alert toggles are [`docs/plan_dma_alerts.md`](docs/plan_dma_alerts.md).
|
Daily MA alert toggles are [`docs/plan_dma_alerts.md`](docs/plan_dma_alerts.md).
|
||||||
Adding `/NQ` `/GC` `/CL` is [`docs/investigate_added_symbols.md`](docs/investigate_added_symbols.md).
|
|
||||||
|
|
||||||
## Tests earn their place by catching a real bug
|
## Tests earn their place by catching a real bug
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -7,7 +7,6 @@ from pathlib import Path
|
||||||
|
|
||||||
from app.analysis.confluence import Cluster
|
from app.analysis.confluence import Cluster
|
||||||
from app.analysis.levels import Level, LevelKind, Side
|
from app.analysis.levels import Level, LevelKind, Side
|
||||||
from app.instrument import DEFAULT_SYMBOL, instrument_for_symbol
|
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
@ -44,7 +43,6 @@ class Alert:
|
||||||
class _Fired:
|
class _Fired:
|
||||||
center: float
|
center: float
|
||||||
at: int
|
at: int
|
||||||
symbol: str = DEFAULT_SYMBOL
|
|
||||||
|
|
||||||
|
|
||||||
class AlertEngine:
|
class AlertEngine:
|
||||||
|
|
@ -85,14 +83,7 @@ class AlertEngine:
|
||||||
if isinstance(payload, dict):
|
if isinstance(payload, dict):
|
||||||
self._next_number = int(payload.get("next_number", 1))
|
self._next_number = int(payload.get("next_number", 1))
|
||||||
payload = payload.get("fired", [])
|
payload = payload.get("fired", [])
|
||||||
return [
|
return [_Fired(float(item["center"]), int(item["at"])) for item in payload]
|
||||||
_Fired(
|
|
||||||
float(item["center"]),
|
|
||||||
int(item["at"]),
|
|
||||||
str(item.get("symbol") or DEFAULT_SYMBOL),
|
|
||||||
)
|
|
||||||
for item in payload
|
|
||||||
]
|
|
||||||
except Exception:
|
except Exception:
|
||||||
# Corrupt state costs one burst of duplicate alerts, which is a far
|
# Corrupt state costs one burst of duplicate alerts, which is a far
|
||||||
# better failure than refusing to start the stream.
|
# better failure than refusing to start the stream.
|
||||||
|
|
@ -124,12 +115,7 @@ class AlertEngine:
|
||||||
{
|
{
|
||||||
"next_number": self._next_number,
|
"next_number": self._next_number,
|
||||||
"fired": [
|
"fired": [
|
||||||
{
|
{"center": entry.center, "at": entry.at} for entry in self._fired
|
||||||
"center": entry.center,
|
|
||||||
"at": entry.at,
|
|
||||||
"symbol": entry.symbol,
|
|
||||||
}
|
|
||||||
for entry in self._fired
|
|
||||||
],
|
],
|
||||||
},
|
},
|
||||||
indent=2,
|
indent=2,
|
||||||
|
|
@ -156,7 +142,6 @@ class AlertEngine:
|
||||||
tolerance = 0.5 * atr15
|
tolerance = 0.5 * atr15
|
||||||
if tolerance <= 0:
|
if tolerance <= 0:
|
||||||
return []
|
return []
|
||||||
root = instrument_for_symbol(symbol).schwab_symbol
|
|
||||||
# Two zones within an ATR of each other are the same zone as far as
|
# Two zones within an ATR of each other are the same zone as far as
|
||||||
# being told about them goes.
|
# being told about them goes.
|
||||||
merge_distance = 2 * tolerance
|
merge_distance = 2 * tolerance
|
||||||
|
|
@ -210,12 +195,10 @@ class AlertEngine:
|
||||||
# was oscillating on re-alerted on every crossing — which is exactly
|
# was oscillating on re-alerted on every crossing — which is exactly
|
||||||
# when a level is least newsworthy, not most.
|
# when a level is least newsworthy, not most.
|
||||||
if any(
|
if any(
|
||||||
entry.symbol == root
|
abs(entry.center - cluster.center) <= merge_distance for entry in self._fired
|
||||||
and abs(entry.center - cluster.center) <= merge_distance
|
|
||||||
for entry in self._fired
|
|
||||||
):
|
):
|
||||||
continue
|
continue
|
||||||
self._fired.append(_Fired(cluster.center, now, root))
|
self._fired.append(_Fired(cluster.center, now))
|
||||||
changed = True
|
changed = True
|
||||||
direction = "BEARISH" if cluster.side.value == "resistance" else "BULLISH"
|
direction = "BEARISH" if cluster.side.value == "resistance" else "BULLISH"
|
||||||
timeframes = ", ".join(dict.fromkeys(member.tf.value for member in cluster.members))
|
timeframes = ", ".join(dict.fromkeys(member.tf.value for member in cluster.members))
|
||||||
|
|
@ -256,13 +239,9 @@ class AlertEngine:
|
||||||
price = level.current_p if level.current_p is not None else level.anchor_p
|
price = level.current_p if level.current_p is not None else level.anchor_p
|
||||||
if abs(price - current_price) > tolerance:
|
if abs(price - current_price) > tolerance:
|
||||||
continue
|
continue
|
||||||
if any(
|
if any(abs(entry.center - price) <= merge_distance for entry in self._fired):
|
||||||
entry.symbol == root
|
|
||||||
and abs(entry.center - price) <= merge_distance
|
|
||||||
for entry in self._fired
|
|
||||||
):
|
|
||||||
continue
|
continue
|
||||||
self._fired.append(_Fired(price, now, root))
|
self._fired.append(_Fired(price, now))
|
||||||
changed = True
|
changed = True
|
||||||
direction = "BEARISH" if level.side is Side.RESISTANCE else "BULLISH"
|
direction = "BEARISH" if level.side is Side.RESISTANCE else "BULLISH"
|
||||||
name = f"{level.period} DMA" if level.period else level.label
|
name = f"{level.period} DMA" if level.period else level.label
|
||||||
|
|
|
||||||
|
|
@ -3,8 +3,6 @@ import logging
|
||||||
import time
|
import time
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
|
|
||||||
from app.instrument import DEFAULT_SYMBOL
|
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
DAY_SECONDS = 86400
|
DAY_SECONDS = 86400
|
||||||
|
|
@ -41,16 +39,12 @@ class EventLog:
|
||||||
except Exception:
|
except Exception:
|
||||||
logger.warning("Could not persist event log", exc_info=True)
|
logger.warning("Could not persist event log", exc_info=True)
|
||||||
|
|
||||||
def add(
|
def add(self, kind: str, message: str, *, number: int | None = None, at: int | None = None) -> dict:
|
||||||
self, kind: str, message: str, *, number: int | None = None,
|
|
||||||
at: int | None = None, symbol: str | None = None,
|
|
||||||
) -> dict:
|
|
||||||
entry = {
|
entry = {
|
||||||
"kind": kind,
|
"kind": kind,
|
||||||
"message": message,
|
"message": message,
|
||||||
"number": number,
|
"number": number,
|
||||||
"at": int(at if at is not None else time.time()),
|
"at": int(at if at is not None else time.time()),
|
||||||
"symbol": symbol or DEFAULT_SYMBOL,
|
|
||||||
}
|
}
|
||||||
self._entries.append(entry)
|
self._entries.append(entry)
|
||||||
self._save()
|
self._save()
|
||||||
|
|
|
||||||
|
|
@ -3,7 +3,6 @@ from enum import Enum
|
||||||
from typing import Any
|
from typing import Any
|
||||||
|
|
||||||
from app.bars.models import Timeframe
|
from app.bars.models import Timeframe
|
||||||
from app.instrument import DEFAULT_SYMBOL
|
|
||||||
|
|
||||||
|
|
||||||
class LevelKind(str, Enum):
|
class LevelKind(str, Enum):
|
||||||
|
|
@ -56,7 +55,6 @@ class Level:
|
||||||
# price this line safely. Such a line remains visible but cannot cluster or
|
# price this line safely. Such a line remains visible but cannot cluster or
|
||||||
# alert using the absolute-time fallback.
|
# alert using the absolute-time fallback.
|
||||||
geometry_resolved: bool = True
|
geometry_resolved: bool = True
|
||||||
symbol: str = DEFAULT_SYMBOL
|
|
||||||
|
|
||||||
def price_at(self, t: int) -> float:
|
def price_at(self, t: int) -> float:
|
||||||
return self.anchor_p + self.slope * (t - self.anchor_t)
|
return self.anchor_p + self.slope * (t - self.anchor_t)
|
||||||
|
|
|
||||||
|
|
@ -6,7 +6,6 @@ from threading import RLock
|
||||||
from app.analysis.levels import Level, LevelKind, Side
|
from app.analysis.levels import Level, LevelKind, Side
|
||||||
from app.bars.models import Timeframe
|
from app.bars.models import Timeframe
|
||||||
from app.config import TIMEFRAME_WEIGHT
|
from app.config import TIMEFRAME_WEIGHT
|
||||||
from app.instrument import DEFAULT_SYMBOL
|
|
||||||
|
|
||||||
|
|
||||||
@dataclass(slots=True)
|
@dataclass(slots=True)
|
||||||
|
|
@ -38,7 +37,6 @@ class ManualLine:
|
||||||
# Visual multiplier for marks. 1 is the original 30px glyph; the pin stays
|
# Visual multiplier for marks. 1 is the original 30px glyph; the pin stays
|
||||||
# on (anchor_t, anchor_p) regardless of this value.
|
# on (anchor_t, anchor_p) regardless of this value.
|
||||||
scale: float = 1.0
|
scale: float = 1.0
|
||||||
symbol: str = DEFAULT_SYMBOL
|
|
||||||
|
|
||||||
@property
|
@property
|
||||||
def drawing_kind(self) -> str:
|
def drawing_kind(self) -> str:
|
||||||
|
|
@ -95,7 +93,6 @@ class ManualLine:
|
||||||
cutoff_t=self.cutoff_t,
|
cutoff_t=self.cutoff_t,
|
||||||
armed=self.armed,
|
armed=self.armed,
|
||||||
alert_early_points=self.alert_early_points,
|
alert_early_points=self.alert_early_points,
|
||||||
symbol=self.symbol,
|
|
||||||
)
|
)
|
||||||
|
|
||||||
def to_dict(self) -> dict:
|
def to_dict(self) -> dict:
|
||||||
|
|
@ -133,7 +130,6 @@ class ManualLine:
|
||||||
if value.get("alert_early_points") is not None else None
|
if value.get("alert_early_points") is not None else None
|
||||||
),
|
),
|
||||||
scale=float(value.get("scale", 1.0) or 1.0),
|
scale=float(value.get("scale", 1.0) or 1.0),
|
||||||
symbol=str(value.get("symbol") or DEFAULT_SYMBOL),
|
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -132,7 +132,6 @@ class LineRestore(BaseModel):
|
||||||
icon: str = ""
|
icon: str = ""
|
||||||
alert_early_points: float | None = None
|
alert_early_points: float | None = None
|
||||||
scale: float = Field(1.0, ge=0.5, le=3)
|
scale: float = Field(1.0, ge=0.5, le=3)
|
||||||
symbol: str = ""
|
|
||||||
|
|
||||||
|
|
||||||
class LinePatch(BaseModel):
|
class LinePatch(BaseModel):
|
||||||
|
|
@ -176,7 +175,6 @@ def status(request: Request):
|
||||||
"worst": round(runtime.loop_lag_worst * 1000, 1),
|
"worst": round(runtime.loop_lag_worst * 1000, 1),
|
||||||
},
|
},
|
||||||
"needs_login": runtime.needs_login(),
|
"needs_login": runtime.needs_login(),
|
||||||
"instrument": runtime.settings.profile.payload(),
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
|
@ -281,7 +279,6 @@ def create_line(request: Request, payload: LineCreate):
|
||||||
cutoff_t=payload.cutoff_t,
|
cutoff_t=payload.cutoff_t,
|
||||||
armed=payload.armed,
|
armed=payload.armed,
|
||||||
kind=payload.kind,
|
kind=payload.kind,
|
||||||
symbol=request.app.state.runtime.settings.profile.schwab_symbol,
|
|
||||||
)
|
)
|
||||||
runtime = request.app.state.runtime
|
runtime = request.app.state.runtime
|
||||||
line = runtime.manual_lines.add(line)
|
line = runtime.manual_lines.add(line)
|
||||||
|
|
@ -322,7 +319,6 @@ def restore_line(request: Request, payload: LineRestore):
|
||||||
icon=payload.icon,
|
icon=payload.icon,
|
||||||
alert_early_points=payload.alert_early_points,
|
alert_early_points=payload.alert_early_points,
|
||||||
scale=payload.scale,
|
scale=payload.scale,
|
||||||
symbol=payload.symbol or runtime.settings.profile.schwab_symbol,
|
|
||||||
)
|
)
|
||||||
line = runtime.manual_lines.add(line)
|
line = runtime.manual_lines.add(line)
|
||||||
runtime.rebuild_levels()
|
runtime.rebuild_levels()
|
||||||
|
|
@ -352,7 +348,6 @@ def create_price_alert(request: Request, payload: PriceAlertCreate):
|
||||||
color=payload.color,
|
color=payload.color,
|
||||||
line_width=payload.line_width,
|
line_width=payload.line_width,
|
||||||
alert_early_points=payload.alert_early_points,
|
alert_early_points=payload.alert_early_points,
|
||||||
symbol=runtime.settings.profile.schwab_symbol,
|
|
||||||
)
|
)
|
||||||
line = runtime.manual_lines.add(line)
|
line = runtime.manual_lines.add(line)
|
||||||
runtime.rebuild_levels()
|
runtime.rebuild_levels()
|
||||||
|
|
@ -405,7 +400,6 @@ def create_comment(request: Request, payload: CommentCreate):
|
||||||
y=payload.y,
|
y=payload.y,
|
||||||
# A comment must never alert, whatever else changes around it.
|
# A comment must never alert, whatever else changes around it.
|
||||||
armed=False,
|
armed=False,
|
||||||
symbol=runtime.settings.profile.schwab_symbol,
|
|
||||||
)
|
)
|
||||||
line = runtime.manual_lines.add(line)
|
line = runtime.manual_lines.add(line)
|
||||||
return line.to_dict()
|
return line.to_dict()
|
||||||
|
|
|
||||||
|
|
@ -193,7 +193,6 @@ def snapshot(runtime, tf: Timeframe, prefs: dict | None = None) -> dict:
|
||||||
"future_times": displayed_future_times(runtime, tf),
|
"future_times": displayed_future_times(runtime, tf),
|
||||||
"events": events,
|
"events": events,
|
||||||
"events_more": events_more,
|
"events_more": events_more,
|
||||||
"instrument": runtime.settings.profile.payload(),
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
|
@ -251,12 +250,7 @@ async def websocket_endpoint(websocket: WebSocket):
|
||||||
event = await queue.get()
|
event = await queue.get()
|
||||||
if event["type"] == "disconnect":
|
if event["type"] == "disconnect":
|
||||||
break
|
break
|
||||||
if event["type"] == "resync":
|
if event["type"] == "bar":
|
||||||
# History changed behind the live edge (a backfilled outage).
|
|
||||||
# Bar deltas only move the tail, so the whole series is resent.
|
|
||||||
last_bar_t.clear()
|
|
||||||
await websocket.send_json(snapshot(runtime, tf, prefs))
|
|
||||||
elif event["type"] == "bar":
|
|
||||||
bar = event["bar"]
|
bar = event["bar"]
|
||||||
full = last_bar_t.get(bar.tf) != bar.t
|
full = last_bar_t.get(bar.tf) != bar.t
|
||||||
if full:
|
if full:
|
||||||
|
|
|
||||||
|
|
@ -46,35 +46,6 @@ class InMemoryBarStore:
|
||||||
# Buckets are ordered, so nothing further back can match.
|
# Buckets are ordered, so nothing further back can match.
|
||||||
return
|
return
|
||||||
|
|
||||||
def fill(self, bars: list[Bar]) -> int:
|
|
||||||
"""Insert history into buckets the store has no bar for.
|
|
||||||
|
|
||||||
``put`` only lands a bar at the tail or a few buckets behind it, so a
|
|
||||||
stretch missed while the stream was down cannot reach it — the live
|
|
||||||
bars that arrived on reconnect are already newer. Existing bars always
|
|
||||||
win: they are the live source's own figures, and the bucket either side
|
|
||||||
of the hole is the live aggregator's to finish. Returns how many bars
|
|
||||||
were inserted.
|
|
||||||
"""
|
|
||||||
added = 0
|
|
||||||
by_tf: dict[Timeframe, list[Bar]] = defaultdict(list)
|
|
||||||
for bar in bars:
|
|
||||||
by_tf[bar.tf].append(bar)
|
|
||||||
for tf, incoming in by_tf.items():
|
|
||||||
held = self._bars[tf]
|
|
||||||
merged = {bar.t: bar for bar in held}
|
|
||||||
for bar in incoming:
|
|
||||||
if bar.t not in merged:
|
|
||||||
merged[bar.t] = bar
|
|
||||||
added += 1
|
|
||||||
if len(merged) == len(held):
|
|
||||||
continue
|
|
||||||
ordered = [merged[t] for t in sorted(merged)]
|
|
||||||
held.clear()
|
|
||||||
# maxlen keeps the newest, which is the history the chart shows.
|
|
||||||
held.extend(ordered)
|
|
||||||
return added
|
|
||||||
|
|
||||||
def get(self, tf: Timeframe, limit: int | None = None) -> list[Bar]:
|
def get(self, tf: Timeframe, limit: int | None = None) -> list[Bar]:
|
||||||
bars = list(self._bars[tf])
|
bars = list(self._bars[tf])
|
||||||
return bars[-limit:] if limit is not None else bars
|
return bars[-limit:] if limit is not None else bars
|
||||||
|
|
|
||||||
|
|
@ -3,7 +3,6 @@ from pathlib import Path
|
||||||
from pydantic_settings import BaseSettings, SettingsConfigDict
|
from pydantic_settings import BaseSettings, SettingsConfigDict
|
||||||
|
|
||||||
from app.bars.models import Timeframe
|
from app.bars.models import Timeframe
|
||||||
from app.instrument import get_instrument
|
|
||||||
|
|
||||||
|
|
||||||
DEFAULT_MAX_BARS_PER_TF = 5000
|
DEFAULT_MAX_BARS_PER_TF = 5000
|
||||||
|
|
@ -25,7 +24,6 @@ class Settings(BaseSettings):
|
||||||
|
|
||||||
live_source: str = "yahoo"
|
live_source: str = "yahoo"
|
||||||
seed_source: str = "yahoo"
|
seed_source: str = "yahoo"
|
||||||
instrument: str = "es"
|
|
||||||
yahoo_symbol: str = "ES=F"
|
yahoo_symbol: str = "ES=F"
|
||||||
yahoo_poll_seconds: float = 20
|
yahoo_poll_seconds: float = 20
|
||||||
seed_1h_range: str = "730d"
|
seed_1h_range: str = "730d"
|
||||||
|
|
@ -75,10 +73,6 @@ class Settings(BaseSettings):
|
||||||
alert_timezone: str = "America/Chicago"
|
alert_timezone: str = "America/Chicago"
|
||||||
replay_file: Path | None = None
|
replay_file: Path | None = None
|
||||||
|
|
||||||
@property
|
|
||||||
def profile(self):
|
|
||||||
return get_instrument(self.instrument)
|
|
||||||
|
|
||||||
@property
|
@property
|
||||||
def live_symbol(self) -> str:
|
def live_symbol(self) -> str:
|
||||||
"""What the live source calls the instrument.
|
"""What the live source calls the instrument.
|
||||||
|
|
|
||||||
|
|
@ -1,51 +0,0 @@
|
||||||
from dataclasses import dataclass
|
|
||||||
|
|
||||||
|
|
||||||
DEFAULT_SYMBOL = "/ES"
|
|
||||||
|
|
||||||
|
|
||||||
@dataclass(frozen=True)
|
|
||||||
class Instrument:
|
|
||||||
id: str
|
|
||||||
yahoo_symbol: str
|
|
||||||
schwab_symbol: str
|
|
||||||
tick: float
|
|
||||||
decimals: int = 2
|
|
||||||
session: str = "globex_18_17"
|
|
||||||
rth: str = "spy_rth"
|
|
||||||
|
|
||||||
def snap(self, price: float) -> float:
|
|
||||||
return round(round(price / self.tick) * self.tick, self.decimals)
|
|
||||||
|
|
||||||
def payload(self) -> dict:
|
|
||||||
return {
|
|
||||||
"id": self.id,
|
|
||||||
"yahoo_symbol": self.yahoo_symbol,
|
|
||||||
"schwab_symbol": self.schwab_symbol,
|
|
||||||
"tick": self.tick,
|
|
||||||
"decimals": self.decimals,
|
|
||||||
"session": self.session,
|
|
||||||
"rth": self.rth,
|
|
||||||
}
|
|
||||||
|
|
||||||
|
|
||||||
INSTRUMENTS = {
|
|
||||||
"es": Instrument("es", "ES=F", "/ES", 0.25, rth="spy_rth"),
|
|
||||||
"nq": Instrument("nq", "NQ=F", "/NQ", 0.25, rth="spy_rth"),
|
|
||||||
"gc": Instrument("gc", "GC=F", "/GC", 0.10, rth="none"),
|
|
||||||
"cl": Instrument("cl", "CL=F", "/CL", 0.01, rth="nymex_day"),
|
|
||||||
}
|
|
||||||
|
|
||||||
|
|
||||||
def get_instrument(instrument_id: str) -> Instrument:
|
|
||||||
try:
|
|
||||||
return INSTRUMENTS[instrument_id]
|
|
||||||
except KeyError:
|
|
||||||
raise ValueError(f"Unknown instrument: {instrument_id}") from None
|
|
||||||
|
|
||||||
|
|
||||||
def instrument_for_symbol(symbol: str) -> Instrument:
|
|
||||||
for instrument in INSTRUMENTS.values():
|
|
||||||
if symbol in (instrument.id, instrument.schwab_symbol, instrument.yahoo_symbol):
|
|
||||||
return instrument
|
|
||||||
return INSTRUMENTS["es"]
|
|
||||||
|
|
@ -19,12 +19,6 @@ class StreamService:
|
||||||
self._handlers: list[BarHandler] = []
|
self._handlers: list[BarHandler] = []
|
||||||
self._stop = asyncio.Event()
|
self._stop = asyncio.Event()
|
||||||
self.on_drop = None
|
self.on_drop = None
|
||||||
# Called with (last bar before the outage, first bar after it) when a
|
|
||||||
# new connection opens further past the last bar than this. Reconnect
|
|
||||||
# alone only resumes the present; nothing else fetches what was missed.
|
|
||||||
self.on_resume: Callable[[int, int], None] | None = None
|
|
||||||
self.resume_gap_seconds = 120
|
|
||||||
self.reconnect_seconds = 5.0
|
|
||||||
|
|
||||||
def add_handler(self, handler: BarHandler) -> None:
|
def add_handler(self, handler: BarHandler) -> None:
|
||||||
self._handlers.append(handler)
|
self._handlers.append(handler)
|
||||||
|
|
@ -51,21 +45,11 @@ class StreamService:
|
||||||
|
|
||||||
async def run(self) -> None:
|
async def run(self) -> None:
|
||||||
while not self._stop.is_set():
|
while not self._stop.is_set():
|
||||||
first = True
|
|
||||||
try:
|
try:
|
||||||
async for bar in self.source.stream(self.symbol):
|
async for bar in self.source.stream(self.symbol):
|
||||||
self.status = "replay" if self.source.name == "replay" else "connected"
|
self.status = "replay" if self.source.name == "replay" else "connected"
|
||||||
self.last_error = None
|
self.last_error = None
|
||||||
before = self.last_bar_t
|
|
||||||
await self._emit(bar)
|
await self._emit(bar)
|
||||||
if first:
|
|
||||||
first = False
|
|
||||||
if (
|
|
||||||
self.on_resume is not None
|
|
||||||
and before is not None
|
|
||||||
and bar.t - before > self.resume_gap_seconds
|
|
||||||
):
|
|
||||||
self.on_resume(before, bar.t)
|
|
||||||
if self._stop.is_set():
|
if self._stop.is_set():
|
||||||
break
|
break
|
||||||
if self.source.name == "replay":
|
if self.source.name == "replay":
|
||||||
|
|
@ -80,7 +64,7 @@ class StreamService:
|
||||||
self.on_drop(str(exc))
|
self.on_drop(str(exc))
|
||||||
self.status = "disconnected"
|
self.status = "disconnected"
|
||||||
try:
|
try:
|
||||||
await asyncio.wait_for(self._stop.wait(), timeout=self.reconnect_seconds)
|
await asyncio.wait_for(self._stop.wait(), timeout=5)
|
||||||
except TimeoutError:
|
except TimeoutError:
|
||||||
pass
|
pass
|
||||||
|
|
||||||
|
|
|
||||||
111
app/runtime.py
111
app/runtime.py
|
|
@ -27,15 +27,6 @@ from app.notify.ntfy import send_ntfy
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
def _range_seconds(range_: str) -> int:
|
|
||||||
"""How far back a Yahoo range reaches: "8d" is eight days."""
|
|
||||||
units = {"d": 86400, "wk": 7 * 86400, "mo": 31 * 86400, "y": 366 * 86400}
|
|
||||||
for suffix, seconds in units.items():
|
|
||||||
if range_.endswith(suffix) and range_[: -len(suffix)].isdigit():
|
|
||||||
return int(range_[: -len(suffix)]) * seconds
|
|
||||||
return 8 * 86400
|
|
||||||
|
|
||||||
|
|
||||||
@dataclass
|
@dataclass
|
||||||
class Runtime:
|
class Runtime:
|
||||||
settings: Settings
|
settings: Settings
|
||||||
|
|
@ -53,7 +44,6 @@ class Runtime:
|
||||||
alert_engine: AlertEngine = field(init=False)
|
alert_engine: AlertEngine = field(init=False)
|
||||||
_sent_levels: dict[str, dict] = field(default_factory=dict)
|
_sent_levels: dict[str, dict] = field(default_factory=dict)
|
||||||
_notify_tasks: set[asyncio.Task] = field(default_factory=set)
|
_notify_tasks: set[asyncio.Task] = field(default_factory=set)
|
||||||
_backfill_tasks: set[asyncio.Task] = field(default_factory=set)
|
|
||||||
# The loop that owns the subscriber queues. Set once the app is running;
|
# The loop that owns the subscriber queues. Set once the app is running;
|
||||||
# None while a test drives the runtime directly.
|
# None while a test drives the runtime directly.
|
||||||
_loop: asyncio.AbstractEventLoop | None = None
|
_loop: asyncio.AbstractEventLoop | None = None
|
||||||
|
|
@ -91,7 +81,6 @@ class Runtime:
|
||||||
self.stream = StreamService(live_source(self.settings), self.settings.live_symbol)
|
self.stream = StreamService(live_source(self.settings), self.settings.live_symbol)
|
||||||
self.stream.add_handler(self.on_bar)
|
self.stream.add_handler(self.on_bar)
|
||||||
self.stream.on_drop = self._on_stream_drop
|
self.stream.on_drop = self._on_stream_drop
|
||||||
self.stream.on_resume = self._on_stream_resume
|
|
||||||
|
|
||||||
async def on_bar(self, bar: Bar) -> None:
|
async def on_bar(self, bar: Bar) -> None:
|
||||||
# A tick-built bar is provisional and arrives many times a minute. It
|
# A tick-built bar is provisional and arrives many times a minute. It
|
||||||
|
|
@ -328,108 +317,12 @@ class Runtime:
|
||||||
|
|
||||||
def _on_stream_drop(self, error: str) -> None:
|
def _on_stream_drop(self, error: str) -> None:
|
||||||
kind = "auth" if "invalid_grant" in error or "Refresh token" in error else "stream"
|
kind = "auth" if "invalid_grant" in error or "Refresh token" in error else "stream"
|
||||||
self.events.add(
|
self.events.add(kind, error.split("\n", 1)[0][:200])
|
||||||
kind, error.split("\n", 1)[0][:200],
|
|
||||||
symbol=self.settings.profile.schwab_symbol,
|
|
||||||
)
|
|
||||||
|
|
||||||
def _on_stream_resume(self, after: int, before: int) -> None:
|
|
||||||
task = asyncio.create_task(self.backfill_gap(after, before), name="gap-backfill")
|
|
||||||
self._backfill_tasks.add(task)
|
|
||||||
task.add_done_callback(self._backfill_tasks.discard)
|
|
||||||
|
|
||||||
async def backfill_gap(self, after: int, before: int) -> None:
|
|
||||||
"""Fill the history missed while the stream was down.
|
|
||||||
|
|
||||||
The seed runs once, at startup, so an outage between deploys used to
|
|
||||||
stay a hole until the next one — fifteen days of it after a Schwab
|
|
||||||
refresh token expired. Yahoo serves futures about ten minutes late, so
|
|
||||||
the minutes just before reconnect are fetched again once they exist.
|
|
||||||
"""
|
|
||||||
try:
|
|
||||||
await self.fill_gap(after, before)
|
|
||||||
delay = getattr(seed_source(self.settings), "delay_minutes", 0) or 0
|
|
||||||
if delay:
|
|
||||||
await asyncio.sleep(delay * 60 + 60)
|
|
||||||
await self.fill_gap(max(after, before - (delay + 5) * 60), before)
|
|
||||||
except asyncio.CancelledError:
|
|
||||||
raise
|
|
||||||
except Exception:
|
|
||||||
# A failed backfill leaves the hole it found; the live stream is fine.
|
|
||||||
logger.exception("Gap backfill failed")
|
|
||||||
|
|
||||||
async def fill_gap(self, after: int, before: int) -> int:
|
|
||||||
source = seed_source(self.settings)
|
|
||||||
if source is None or not source.supports_history():
|
|
||||||
return 0
|
|
||||||
symbol = self.settings.yahoo_symbol
|
|
||||||
now = int(time.time())
|
|
||||||
# Native coarse history first: those buckets are complete, where one
|
|
||||||
# rebuilt from 1m is only as old as Yahoo's 1m reach.
|
|
||||||
passes = (
|
|
||||||
(Timeframe.H1, self.settings.seed_1h_range),
|
|
||||||
(Timeframe.M30, self.settings.seed_30m_range),
|
|
||||||
(Timeframe.M1, self.settings.seed_1m_range),
|
|
||||||
)
|
|
||||||
added = 0
|
|
||||||
for tf, range_ in passes:
|
|
||||||
start = max(bucket_start(after, tf), now - _range_seconds(range_))
|
|
||||||
if start >= before:
|
|
||||||
continue
|
|
||||||
try:
|
|
||||||
bars = await source.history(symbol, tf, start, before)
|
|
||||||
except Exception:
|
|
||||||
logger.exception("Gap backfill: %s history failed", tf.value)
|
|
||||||
continue
|
|
||||||
# A fresh aggregator: the live one is already past the hole, and
|
|
||||||
# feeding it history would reopen buckets it has closed.
|
|
||||||
aggregator = Aggregator(self.settings.enabled_timeframes)
|
|
||||||
derived: dict[tuple[Timeframe, int], Bar] = {}
|
|
||||||
for bar in bars:
|
|
||||||
if not start <= bar.t < before:
|
|
||||||
continue
|
|
||||||
for aggregated in aggregator.update(replace(bar, closed=True)):
|
|
||||||
derived[(aggregated.tf, aggregated.t)] = aggregated
|
|
||||||
added += self.store.fill(list(derived.values()))
|
|
||||||
if added:
|
|
||||||
logger.info("Gap backfill: %d bars between %d and %d", added, after, before)
|
|
||||||
self.refold_forming(before)
|
|
||||||
values = atr(self.store.get(Timeframe.M15), 14)
|
|
||||||
self.atr15 = next((value for value in reversed(values) if value is not None), 0.0)
|
|
||||||
# Levels only, never alerts: a touch during the outage is not news.
|
|
||||||
self.rebuild_levels()
|
|
||||||
self.broadcast({"type": "resync"})
|
|
||||||
return added
|
|
||||||
|
|
||||||
def refold_forming(self, before: int) -> None:
|
|
||||||
"""Rebuild the live buckets that opened before the stream came back.
|
|
||||||
|
|
||||||
The live aggregator started today's daily bar — and the current hour —
|
|
||||||
from the first minute after reconnect, so its open, high and low
|
|
||||||
ignore everything the backfill just recovered. Refolding from the stored
|
|
||||||
minutes is idempotent, so the delayed second pass can run it again.
|
|
||||||
"""
|
|
||||||
minutes = [bar for bar in self.store.get(Timeframe.M1) if bar.closed]
|
|
||||||
for tf, forming in self.aggregator.forming.items():
|
|
||||||
if forming.t >= before:
|
|
||||||
continue
|
|
||||||
inside = [bar for bar in minutes if bucket_start(bar.t, tf) == forming.t]
|
|
||||||
if not inside or inside[0].t >= before:
|
|
||||||
continue
|
|
||||||
forming.o = inside[0].o
|
|
||||||
forming.h = max(bar.h for bar in inside)
|
|
||||||
forming.l = min(bar.l for bar in inside)
|
|
||||||
forming.c = inside[-1].c
|
|
||||||
forming.v = sum(bar.v for bar in inside)
|
|
||||||
self.store.put(replace(forming))
|
|
||||||
|
|
||||||
def dispatch_alerts(self, alerts: list[Alert]) -> None:
|
def dispatch_alerts(self, alerts: list[Alert]) -> None:
|
||||||
tripped: set[str] = set()
|
tripped: set[str] = set()
|
||||||
for alert in alerts:
|
for alert in alerts:
|
||||||
self.events.add(
|
self.events.add("alert", alert.message, number=alert.number, at=alert.at)
|
||||||
"alert", alert.message, number=alert.number, at=alert.at,
|
|
||||||
symbol=self.settings.profile.schwab_symbol,
|
|
||||||
)
|
|
||||||
self.broadcast({
|
self.broadcast({
|
||||||
"type": "alert",
|
"type": "alert",
|
||||||
"cluster": alert.cluster,
|
"cluster": alert.cluster,
|
||||||
|
|
|
||||||
|
|
@ -1418,41 +1418,3 @@ which clips drawing. Slope and both handles stay. Extend clears the cutoff.
|
||||||
The menu now offers End here for any click after the earlier of the two
|
The menu now offers End here for any click after the earlier of the two
|
||||||
anchors, including on the body between them and at `last_t` itself (that
|
anchors, including on the body between them and at `last_t` itself (that
|
||||||
stops the projection without moving the end handle).
|
stops the projection without moving the end handle).
|
||||||
|
|
||||||
### 2026-08-31 — instrument profile and symbol stamps
|
|
||||||
|
|
||||||
Drawings, alert-state rows, and events now carry `symbol`, default `/ES`.
|
|
||||||
A missing field loads as ES so production JSON does not need a rewrite
|
|
||||||
before the first save. Fired zones for `/ES` and `ES=F` are the same
|
|
||||||
root; `/GC` at the same price is not. The browser snap grid reads `tick`
|
|
||||||
from the snapshot/status profile instead of a chart constant. No
|
|
||||||
switcher and no second stream — see `docs/investigate_added_symbols.md`.
|
|
||||||
|
|
||||||
### 2026-09-24 — a Schwab outage stayed a hole until the next deploy
|
|
||||||
|
|
||||||
The chart was missing Sept 9 → 24. The Schwab refresh token was dead
|
|
||||||
(`invalid_grant` every 5 s). After reauth, live bars resumed but the gap did
|
|
||||||
not fill. History was only ever fetched by the startup seed; the stream's
|
|
||||||
reconnect loop just resumes the present. The container had been up since the
|
|
||||||
Sept 4 deploy.
|
|
||||||
|
|
||||||
`put()` could not have taken the recovered bars anyway: it only lands a bar at
|
|
||||||
the tail or within 8 buckets of it, and the live bars were already newer. Hence
|
|
||||||
`InMemoryBarStore.fill`, which inserts into empty buckets only and never
|
|
||||||
replaces a live bar. A `deque` with `maxlen` also raises on `insert` when full,
|
|
||||||
so it rebuilds the sorted deque and lets `maxlen` drop the oldest.
|
|
||||||
|
|
||||||
A reconnect more than 120 s past the last bar now triggers a backfill (`docs/plan.md`, after
|
|
||||||
"In production both run at once"). Measured against real Yahoo for this outage: 11,173 bars in about
|
|
||||||
1 s including network, 3 ms of that in `fill`. 1h/30m/1d cover all 15 days. 1m
|
|
||||||
only reaches back 8 days, so 1m/5m/15m start Sept 17, and the 5,000-bar 1m cap
|
|
||||||
keeps about 3.5 days of that.
|
|
||||||
|
|
||||||
Two things that looked fine and were not:
|
|
||||||
- Yahoo `ES=F` is ~10 minutes late. The first pass always leaves the last
|
|
||||||
~10 minutes before reconnect empty, so a second pass runs after
|
|
||||||
`delay_minutes`.
|
|
||||||
- The live aggregator opened today's daily bar (and the current hour) at
|
|
||||||
the reconnect minute, so filling holes alone left the day's open/high/low
|
|
||||||
wrong, and tomorrow's prior-day levels would inherit that. `refold_forming`
|
|
||||||
rebuilds those buckets from the stored minutes, idempotently.
|
|
||||||
|
|
|
||||||
|
|
@ -1,160 +0,0 @@
|
||||||
# Added symbols — `/ES`, `/NQ`, `/GC`, `/CL`
|
|
||||||
|
|
||||||
**Status: placeholders in, switcher not.** Agreed 2026-08-31. Profile
|
|
||||||
and `symbol` stamps are live; still one stream and no UI switcher.
|
|
||||||
|
|
||||||
The process is one chart of one contract. Yahoo and Schwab already know
|
|
||||||
the other roots. Almost everything after the stream is *this* instrument.
|
|
||||||
The likely set is four: **ES, NQ, gold, oil**. Design for N, not a
|
|
||||||
boolean ES/gold switch.
|
|
||||||
|
|
||||||
---
|
|
||||||
|
|
||||||
## Decision
|
|
||||||
|
|
||||||
**Persist `symbol` now. Do not split `Runtime` now.**
|
|
||||||
|
|
||||||
Same move as `user_id: "shared"` in `docs/multi_user.md`: cheap while
|
|
||||||
there is one value, expensive after two files exist.
|
|
||||||
|
|
||||||
### Do first (placeholders, still one live series)
|
|
||||||
|
|
||||||
1. **Instrument profile** the settings and the browser both see:
|
|
||||||
|
|
||||||
```
|
|
||||||
id: "es" | "nq" | "gc" | "cl"
|
|
||||||
yahoo_symbol: "ES=F" | "NQ=F" | "GC=F" | "CL=F"
|
|
||||||
schwab_symbol: "/ES" | "/NQ" | "/GC" | "/CL"
|
|
||||||
tick: 0.25 | 0.25 | 0.10 | 0.01
|
|
||||||
decimals: 2
|
|
||||||
session: globex_18_17
|
|
||||||
rth: spy_rth | spy_rth | none | nymex_day
|
|
||||||
```
|
|
||||||
|
|
||||||
Snap, nudge, alert inputs, and the status label read `tick` / names
|
|
||||||
from here. Kill `ConfluenceChart.TICK = 0.25`. Session code stays
|
|
||||||
shared. RTH follows the profile (`none` hides SPY marks on gold).
|
|
||||||
|
|
||||||
2. **Stamp persistence.** Drawings, alert-state rows, and events get
|
|
||||||
`symbol` (Schwab root, e.g. `/ES`). Missing field means `/ES`. New
|
|
||||||
writes always stamp the current profile. Do not wait for a second
|
|
||||||
chart.
|
|
||||||
|
|
||||||
3. **Keep one store, one stream, one seed.** `Bar.symbol` already exists.
|
|
||||||
`InMemoryBarStore` stays `tf → bars` until something actually switches.
|
|
||||||
|
|
||||||
Today’s env still selects the one live profile (`YAHOO_SYMBOL` /
|
|
||||||
`SCHWAB_SYMBOL` or an `INSTRUMENT=es` key). Default remains ES.
|
|
||||||
|
|
||||||
### Then (after placeholders have been live)
|
|
||||||
|
|
||||||
4. Prove a second root as a **replace**: point env at NQ or GC, restart,
|
|
||||||
confirm Yahoo history + Schwab stream + tick snaps. NQ is the cheap
|
|
||||||
proof (same tick and RTH as ES). GC or CL is the proof that tick/RTH
|
|
||||||
actually split.
|
|
||||||
|
|
||||||
5. **Switcher** last: UI picks the profile; snapshot replace (`setBars`),
|
|
||||||
not a tick; filter drawings/alerts by `symbol`; lazy-seed the other
|
|
||||||
series; do not Schwab-sub the hidden root.
|
|
||||||
|
|
||||||
### Not in this plan
|
|
||||||
|
|
||||||
Two live streams (ES and gold on screen together). That is a second tick
|
|
||||||
path and a second 5k-bar store inside Stay cheap, plus an unverified
|
|
||||||
double `CHART_FUTURES` sub on one token. Not until the switcher has been
|
|
||||||
used.
|
|
||||||
|
|
||||||
---
|
|
||||||
|
|
||||||
## Instrument table
|
|
||||||
|
|
||||||
| | `/ES` | `/NQ` | `/GC` | `/CL` |
|
|
||||||
|---|---|---|---|---|
|
|
||||||
| Tick | 0.25 | 0.25 | 0.10 | 0.01 |
|
|
||||||
| Display | 2 dp | 2 dp | 2 dp | 2 dp |
|
|
||||||
| Yahoo | `ES=F` | `NQ=F` | `GC=F` | `CL=F` |
|
|
||||||
| Schwab | `/ES` | `/NQ` | `/GC` | `/CL` |
|
|
||||||
| Globex 18:00–17:00 | yes | yes | yes | yes |
|
|
||||||
| SPY RTH overlay | yes | yes | no | no (NYMEX day 9:00–14:30 ET) |
|
|
||||||
| Options UI | keep | hide or later | hide | hide |
|
|
||||||
|
|
||||||
NQ is the cheapest second chart. Gold and oil are why tick cannot stay a
|
|
||||||
chart constant. `toFixed(2)` covers all four.
|
|
||||||
|
|
||||||
---
|
|
||||||
|
|
||||||
## Why the placeholders
|
|
||||||
|
|
||||||
The current process is one `Runtime`, one Schwab socket, one bar store,
|
|
||||||
one `manual_lines.json`, one alert-state file, one ntfy topic.
|
|
||||||
|
|
||||||
| File | Today | After step 2 |
|
|
||||||
|---|---|---|
|
|
||||||
| `data/manual_lines.json` | no symbol | each row `symbol: "/ES"` |
|
|
||||||
| `data/alert_state.json` | zones by price | zone + symbol |
|
|
||||||
| `data/events.json` | one log | tagged or filtered by symbol |
|
|
||||||
| `data/user_prefs.json` | global | leave global until the switcher |
|
|
||||||
| `localStorage` | layers, theme | leave until the switcher |
|
|
||||||
|
|
||||||
Without `symbol` on drawings, a switcher would mix ES lines onto NQ.
|
|
||||||
Without it on alert state, gold 2650 would be silenced by an old ES
|
|
||||||
2650. Adding the field later is a migration of production JSON.
|
|
||||||
|
|
||||||
`ConfluenceChart.TICK = 0.25` feeds every snap, keyboard nudge, and the
|
|
||||||
price-alert `step`. Gold cannot ship with that literal. Pulling it into
|
|
||||||
the profile is not scaffolding — it is deleting a lie.
|
|
||||||
|
|
||||||
Schwab continuous roots (`/ES`, `/NQ`, `/GC`, `/CL`) should auto-resolve
|
|
||||||
the front month the same way `/ES` → `/ESU26`. Verify each with
|
|
||||||
`scripts/check_stream.py` before trusting it. Singular `get_quote("/GC")`
|
|
||||||
is still the equity slash trap; always `get_quotes`.
|
|
||||||
|
|
||||||
---
|
|
||||||
|
|
||||||
## What stays generic
|
|
||||||
|
|
||||||
Aggregator, VWAP, daily MAs, prior-day H/L/C (session is shared Globex),
|
|
||||||
WebSocket snapshot shape, drawing tools, Fibonacci, comments, confluence
|
|
||||||
*math*, ntfy, auth. They work if the bars and the profile are the
|
|
||||||
instrument’s.
|
|
||||||
|
|
||||||
What does not: RTH marks, tick grid, options panel, confluence *score*
|
|
||||||
(28 and the 4h cooldown were calibrated on ES — re-run
|
|
||||||
`scripts/calibrate_alerts.py` per tape before turning confluence on).
|
|
||||||
|
|
||||||
Yahoo daily bars stay unused. 1h → session 1d, same as ES.
|
|
||||||
|
|
||||||
---
|
|
||||||
|
|
||||||
## Stay cheap
|
|
||||||
|
|
||||||
The tick budget is one series: forming tick = candle + price label.
|
|
||||||
A symbol switch is `setBars`. Do not subscribe Schwab to a hidden root.
|
|
||||||
Do not 2× the ~82s Yahoo seed; lazy-load the next instrument on first
|
|
||||||
view.
|
|
||||||
|
|
||||||
Two symbols on one Schwab socket is unverified. Two *processes* both
|
|
||||||
opening a stream still kick each other off.
|
|
||||||
|
|
||||||
---
|
|
||||||
|
|
||||||
## What not to do
|
|
||||||
|
|
||||||
- Do not split `Runtime` or the bar store until the switcher exists.
|
|
||||||
- Do not add a disabled switcher, a second seed, or a second Schwab sub
|
|
||||||
in the placeholder change.
|
|
||||||
- Do not add a second `package.json`, Vue app, or process “for gold.”
|
|
||||||
- Do not leave `TICK = 0.25` and “just chart gold.”
|
|
||||||
- Do not reuse ES confluence calibration.
|
|
||||||
- Do not show SPY RTH on metals or oil.
|
|
||||||
- Do not build gold/oil/NQ options in the same change as the chart.
|
|
||||||
- Do not put two symbols in one `manual_lines.json` without a symbol key.
|
|
||||||
- Do not start two live streams.
|
|
||||||
|
|
||||||
---
|
|
||||||
|
|
||||||
## Open, before the switcher (not before placeholders)
|
|
||||||
|
|
||||||
1. First extra root to prove as a replace: NQ (easy) or GC (forces tick)?
|
|
||||||
2. Does `/NQ` `/GC` `/CL` on this account stream `delayed: false`?
|
|
||||||
3. Does each Yahoo `*F` 1h series go back far enough for a daily 200 SMA?
|
|
||||||
14
docs/plan.md
14
docs/plan.md
|
|
@ -176,16 +176,6 @@ Nothing downstream of these may know which source it is using. Selection is one
|
||||||
var. **In production both run at once:** Yahoo seeds history at startup, Schwab
|
var. **In production both run at once:** Yahoo seeds history at startup, Schwab
|
||||||
provides the live tail.
|
provides the live tail.
|
||||||
|
|
||||||
When the live stream reconnects more than two minutes past its last bar,
|
|
||||||
`Runtime.backfill_gap` fetches the missed stretch from the same seed source
|
|
||||||
(1h, 30m, 1m, each limited to its seed range), rebuilds higher timeframes with a
|
|
||||||
fresh aggregator, and inserts only into empty buckets (`InMemoryBarStore.fill`).
|
|
||||||
Live bars always win. The live aggregator's still-forming buckets are then
|
|
||||||
refolded from the stored minutes, levels are rebuilt without evaluating alerts,
|
|
||||||
and every socket gets a `resync` → full `snapshot`. Yahoo is ~10 minutes late,
|
|
||||||
so a second pass runs after that delay for the minutes just before reconnect.
|
|
||||||
The bucket that was forming when the stream *died* keeps only what it had.
|
|
||||||
|
|
||||||
### Do not use Yahoo's daily bars
|
### Do not use Yahoo's daily bars
|
||||||
|
|
||||||
Yahoo anchors `ES=F` daily bars to **midnight ET**, but the CME futures session runs
|
Yahoo anchors `ES=F` daily bars to **midnight ET**, but the CME futures session runs
|
||||||
|
|
@ -826,8 +816,8 @@ Server → client:
|
||||||
Rules:
|
Rules:
|
||||||
- Send `bar` on **every** update of the forming bar (that is the live chart) but batch
|
- Send `bar` on **every** update of the forming bar (that is the live chart) but batch
|
||||||
`levels` — they only change on higher-TF closes.
|
`levels` — they only change on higher-TF closes.
|
||||||
- Always send a full `snapshot` on connect and after any reconnect, and after the
|
- Always send a full `snapshot` on connect and after any reconnect. The client must
|
||||||
server backfills a gap (`resync`). The client must never try to reconcile a gap.
|
never try to reconcile a gap.
|
||||||
- Levels are sent for **all** timeframes regardless of the displayed timeframe. That is
|
- Levels are sent for **all** timeframes regardless of the displayed timeframe. That is
|
||||||
the entire point: a 4h line drawn through a 1m chart.
|
the entire point: a 4h line drawn through a 1m chart.
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -154,7 +154,6 @@ createApp({
|
||||||
const animateCurrentPrice = ref(localStorage.getItem('chart-animate-current-price') !== 'false');
|
const animateCurrentPrice = ref(localStorage.getItem('chart-animate-current-price') !== 'false');
|
||||||
const autoScrollLivePrice = ref(localStorage.getItem('chart-auto-scroll-live-price') !== 'false');
|
const autoScrollLivePrice = ref(localStorage.getItem('chart-auto-scroll-live-price') !== 'false');
|
||||||
const extraDetail = ref(localStorage.getItem('chart-extra-detail') === 'true');
|
const extraDetail = ref(localStorage.getItem('chart-extra-detail') === 'true');
|
||||||
const tick = ref(0.25);
|
|
||||||
const sessionRange = ref(localStorage.getItem('chart-session-range') !== 'false');
|
const sessionRange = ref(localStorage.getItem('chart-session-range') !== 'false');
|
||||||
const hideLowerTfDrawings = ref(localStorage.getItem('chart-hide-lower-tf-drawings') !== 'false');
|
const hideLowerTfDrawings = ref(localStorage.getItem('chart-hide-lower-tf-drawings') !== 'false');
|
||||||
const optionPrefs = (() => {
|
const optionPrefs = (() => {
|
||||||
|
|
@ -366,13 +365,7 @@ createApp({
|
||||||
|
|
||||||
async function refreshStatus() {
|
async function refreshStatus() {
|
||||||
const response = await apiFetch('/api/status');
|
const response = await apiFetch('/api/status');
|
||||||
if (response.ok) {
|
if (response.ok) status.value = await response.json();
|
||||||
status.value = await response.json();
|
|
||||||
if (status.value.instrument && chartApi) {
|
|
||||||
chartApi.setInstrument(status.value.instrument);
|
|
||||||
tick.value = status.value.instrument.tick || tick.value;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
function selectedExpiration() {
|
function selectedExpiration() {
|
||||||
|
|
@ -532,10 +525,6 @@ createApp({
|
||||||
if (message.type === 'snapshot') {
|
if (message.type === 'snapshot') {
|
||||||
dataReceivedAt.value = Date.now();
|
dataReceivedAt.value = Date.now();
|
||||||
chartApi.setTrendlineGeometry(message.trendline_geometry);
|
chartApi.setTrendlineGeometry(message.trendline_geometry);
|
||||||
if (message.instrument) {
|
|
||||||
chartApi.setInstrument(message.instrument);
|
|
||||||
tick.value = message.instrument.tick || tick.value;
|
|
||||||
}
|
|
||||||
chartApi.setDisplayFutureTimes(message.future_times);
|
chartApi.setDisplayFutureTimes(message.future_times);
|
||||||
chartApi.setBars(message.bars);
|
chartApi.setBars(message.bars);
|
||||||
levels.value = message.levels || [];
|
levels.value = message.levels || [];
|
||||||
|
|
@ -1365,14 +1354,14 @@ createApp({
|
||||||
if (item.line) {
|
if (item.line) {
|
||||||
const line = { ...item.line };
|
const line = { ...item.line };
|
||||||
const previous = { ...item.line };
|
const previous = { ...item.line };
|
||||||
if (vertical) line.anchor_p = chartApi.snapPrice(line.anchor_p + vertical * chartApi.tick);
|
if (vertical) line.anchor_p = ConfluenceChart.snapToTick(line.anchor_p + vertical * ConfluenceChart.TICK);
|
||||||
if (line.slope === 0) {
|
if (line.slope === 0) {
|
||||||
if (vertical) await updateLineGeometry(line, previous);
|
if (vertical) await updateLineGeometry(line, previous);
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
const endPrice = item.line.anchor_p
|
const endPrice = item.line.anchor_p
|
||||||
+ item.line.slope * (item.line.last_t - item.line.anchor_t)
|
+ item.line.slope * (item.line.last_t - item.line.anchor_t)
|
||||||
+ vertical * chartApi.tick;
|
+ vertical * ConfluenceChart.TICK;
|
||||||
if (horizontal) {
|
if (horizontal) {
|
||||||
if (line.geometry_resolved === false) continue;
|
if (line.geometry_resolved === false) continue;
|
||||||
const anchorT = chartApi.shiftLineTime(line, line.anchor_t, horizontal);
|
const anchorT = chartApi.shiftLineTime(line, line.anchor_t, horizontal);
|
||||||
|
|
@ -1394,8 +1383,8 @@ createApp({
|
||||||
if (comment.pinned) {
|
if (comment.pinned) {
|
||||||
const changes = {};
|
const changes = {};
|
||||||
if (vertical) {
|
if (vertical) {
|
||||||
changes.anchor_p = chartApi.snapPrice(
|
changes.anchor_p = ConfluenceChart.snapToTick(
|
||||||
comment.anchor_p + vertical * chartApi.tick,
|
comment.anchor_p + vertical * ConfluenceChart.TICK,
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
if (horizontal) {
|
if (horizontal) {
|
||||||
|
|
@ -1708,6 +1697,6 @@ createApp({
|
||||||
window.removeEventListener('keydown', handleKeydown);
|
window.removeEventListener('keydown', handleKeydown);
|
||||||
});
|
});
|
||||||
|
|
||||||
return { status, price, sessionOpen, quoteChange, animateCurrentPrice, autoScrollLivePrice, extraDetail, sessionRange, hideLowerTfDrawings, confluenceAlerts, tick, barAge, dataUpdatedAt, buildStamp, timeframe, timeframes, drawingColors, drawingColorRows, drawingColorName, colorRowLabels, symbolChoices, selectedSymbol, symbolColor, symbolScale, symbolScales, symbolPanelOpen, prefs, clusters, clustersByPrice, events, eventsMore, loadOlderEvents, diagnosticMode, captureBusy, captureCountdown, captureDiagnostic, armedTool, drawName, drawColor, drawWidth, drawSide, snap, selectedDrawing, selectedDrawings, drawingList, startDrawingListResize, manualLines, hasDrawingSelection, allShownSelected, selectedAreHidden, selectedTrendline, duplicateSelected, alertPrice, alertNote, alertEarlyPoints, levelColor, levelWidth, addPriceAlert, armTool, selectTimeframe, allEnabled, toggleGroup, maValue, maAlertOn, toggleMaAlert, deleteSelected, deleteLine, toggleDrawingSelection, activateDrawing, toggleSelectAll, toggleSelectedVisibility, renameLine, updateLineStyle, updateLevelNumber, setArmed, commentText, commentFloat, comments, drawings, filteredDrawings, drawingFilter, drawingKind, drawingTf, deleteDrawing, toggleComment, togglePinned, chooseSymbol, toggleSymbolPanel, startSymbolDrag, dropSymbol, optionExpirations, optionExpiryId, optionSide, optionMode, optionMin, optionMax, optionContracts, optionUnderlying, optionBusy, optionError, optionSearched, optionCopied, onOptionsToggle, searchOptions, copyOption, canUndo, undoTitle, undo };
|
return { status, price, sessionOpen, quoteChange, animateCurrentPrice, autoScrollLivePrice, extraDetail, sessionRange, hideLowerTfDrawings, confluenceAlerts, barAge, dataUpdatedAt, buildStamp, timeframe, timeframes, drawingColors, drawingColorRows, drawingColorName, colorRowLabels, symbolChoices, selectedSymbol, symbolColor, symbolScale, symbolScales, symbolPanelOpen, prefs, clusters, clustersByPrice, events, eventsMore, loadOlderEvents, diagnosticMode, captureBusy, captureCountdown, captureDiagnostic, armedTool, drawName, drawColor, drawWidth, drawSide, snap, selectedDrawing, selectedDrawings, drawingList, startDrawingListResize, manualLines, hasDrawingSelection, allShownSelected, selectedAreHidden, selectedTrendline, duplicateSelected, alertPrice, alertNote, alertEarlyPoints, levelColor, levelWidth, addPriceAlert, armTool, selectTimeframe, allEnabled, toggleGroup, maValue, maAlertOn, toggleMaAlert, deleteSelected, deleteLine, toggleDrawingSelection, activateDrawing, toggleSelectAll, toggleSelectedVisibility, renameLine, updateLineStyle, updateLevelNumber, setArmed, commentText, commentFloat, comments, drawings, filteredDrawings, drawingFilter, drawingKind, drawingTf, deleteDrawing, toggleComment, togglePinned, chooseSymbol, toggleSymbolPanel, startSymbolDrag, dropSymbol, optionExpirations, optionExpiryId, optionSide, optionMode, optionMin, optionMax, optionContracts, optionUnderlying, optionBusy, optionError, optionSearched, optionCopied, onOptionsToggle, searchOptions, copyOption, canUndo, undoTitle, undo };
|
||||||
},
|
},
|
||||||
}).mount('#app');
|
}).mount('#app');
|
||||||
|
|
|
||||||
|
|
@ -294,8 +294,6 @@ class ConfluenceChart {
|
||||||
this.lastCurrentPrice = null;
|
this.lastCurrentPrice = null;
|
||||||
this.pendingBar = null;
|
this.pendingBar = null;
|
||||||
this.pendingBarFrame = null;
|
this.pendingBarFrame = null;
|
||||||
this.tick = 0.25;
|
|
||||||
this.rthMode = 'spy_rth';
|
|
||||||
this.previewLine = null;
|
this.previewLine = null;
|
||||||
this.lineBridgeLayer = null;
|
this.lineBridgeLayer = null;
|
||||||
this.bars = [];
|
this.bars = [];
|
||||||
|
|
@ -386,6 +384,8 @@ class ConfluenceChart {
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
|
static TICK = 0.25;
|
||||||
|
|
||||||
// Pixels of travel between press and release that make a gesture a drag
|
// Pixels of travel between press and release that make a gesture a drag
|
||||||
// rather than a click. Wide enough to survive a twitch on a deliberate click.
|
// rather than a click. Wide enough to survive a twitch on a deliberate click.
|
||||||
static DRAG_THRESHOLD = 12;
|
static DRAG_THRESHOLD = 12;
|
||||||
|
|
@ -432,20 +432,8 @@ class ConfluenceChart {
|
||||||
return localStorage.getItem('chart-diag') === '1';
|
return localStorage.getItem('chart-diag') === '1';
|
||||||
}
|
}
|
||||||
|
|
||||||
static snapToTick(price, tick = 0.25) {
|
static snapToTick(price) {
|
||||||
const step = tick > 0 ? tick : 0.25;
|
return Math.round(price / ConfluenceChart.TICK) * ConfluenceChart.TICK;
|
||||||
return Math.round(price / step) * step;
|
|
||||||
}
|
|
||||||
|
|
||||||
snapPrice(price) {
|
|
||||||
return ConfluenceChart.snapToTick(price, this.tick);
|
|
||||||
}
|
|
||||||
|
|
||||||
setInstrument(instrument) {
|
|
||||||
if (!instrument) return;
|
|
||||||
if (instrument.tick > 0) this.tick = Number(instrument.tick);
|
|
||||||
if (instrument.rth) this.rthMode = instrument.rth;
|
|
||||||
this.syncRthLines();
|
|
||||||
}
|
}
|
||||||
|
|
||||||
create(el) {
|
create(el) {
|
||||||
|
|
@ -1425,7 +1413,7 @@ class ConfluenceChart {
|
||||||
if (point && this.withinPlot(point) && this.onCommentMove) {
|
if (point && this.withinPlot(point) && this.onCommentMove) {
|
||||||
this.onCommentMove(node.comment, {
|
this.onCommentMove(node.comment, {
|
||||||
anchor_t: Math.round(point.t),
|
anchor_t: Math.round(point.t),
|
||||||
anchor_p: this.snapPrice(point.p),
|
anchor_p: ConfluenceChart.snapToTick(point.p),
|
||||||
});
|
});
|
||||||
} else {
|
} else {
|
||||||
this.renderComments();
|
this.renderComments();
|
||||||
|
|
@ -1840,7 +1828,7 @@ class ConfluenceChart {
|
||||||
if (!point || !this.withinPlot(point) || point.t == null) return null;
|
if (!point || !this.withinPlot(point) || point.t == null) return null;
|
||||||
return {
|
return {
|
||||||
time: point.t,
|
time: point.t,
|
||||||
price: this.snapPrice(point.p),
|
price: ConfluenceChart.snapToTick(point.p),
|
||||||
x: Math.min(Math.max(point.x / this.overlayLayer.clientWidth, 0), 1),
|
x: Math.min(Math.max(point.x / this.overlayLayer.clientWidth, 0), 1),
|
||||||
y: Math.min(Math.max(point.y / this.overlayLayer.clientHeight, 0), 1),
|
y: Math.min(Math.max(point.y / this.overlayLayer.clientHeight, 0), 1),
|
||||||
};
|
};
|
||||||
|
|
@ -2168,7 +2156,7 @@ class ConfluenceChart {
|
||||||
this.onToolComplete?.({
|
this.onToolComplete?.({
|
||||||
tool,
|
tool,
|
||||||
time: end.t,
|
time: end.t,
|
||||||
price: this.snapPrice(end.p),
|
price: ConfluenceChart.snapToTick(end.p),
|
||||||
x: Math.min(Math.max(end.x / this.chartEl.clientWidth, 0), 1),
|
x: Math.min(Math.max(end.x / this.chartEl.clientWidth, 0), 1),
|
||||||
y: Math.min(Math.max(end.y / this.chartEl.clientHeight, 0), 1),
|
y: Math.min(Math.max(end.y / this.chartEl.clientHeight, 0), 1),
|
||||||
});
|
});
|
||||||
|
|
@ -2178,7 +2166,7 @@ class ConfluenceChart {
|
||||||
if (tool === 'level') {
|
if (tool === 'level') {
|
||||||
// A click with no drag is a valid placement; the drag is only there to
|
// A click with no drag is a valid placement; the drag is only there to
|
||||||
// let you fine-tune the price before committing.
|
// let you fine-tune the price before committing.
|
||||||
this.onToolComplete?.({ tool, price: this.snapPrice(end.p) });
|
this.onToolComplete?.({ tool, price: ConfluenceChart.snapToTick(end.p) });
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -2261,7 +2249,7 @@ class ConfluenceChart {
|
||||||
const last = this.bars[this.bars.length - 1];
|
const last = this.bars[this.bars.length - 1];
|
||||||
if (fallbackT > last.t) {
|
if (fallbackT > last.t) {
|
||||||
const index = Math.round(this.indexAt(fallbackT));
|
const index = Math.round(this.indexAt(fallbackT));
|
||||||
return { t: this.timeAtIndex(index), p: this.snapPrice(point.p), snappedSide: null };
|
return { t: this.timeAtIndex(index), p: ConfluenceChart.snapToTick(point.p), snappedSide: null };
|
||||||
}
|
}
|
||||||
// An already-snapped point carries no cursor position; return it untouched
|
// An already-snapped point carries no cursor position; return it untouched
|
||||||
// rather than measuring against undefined.
|
// rather than measuring against undefined.
|
||||||
|
|
@ -2290,7 +2278,7 @@ class ConfluenceChart {
|
||||||
if (!this.gesture) return;
|
if (!this.gesture) return;
|
||||||
const { start, end } = this.gesture;
|
const { start, end } = this.gesture;
|
||||||
if (this.armedTool === 'level') {
|
if (this.armedTool === 'level') {
|
||||||
const price = this.snapPrice(end.p);
|
const price = ConfluenceChart.snapToTick(end.p);
|
||||||
const y = this.candles.priceToCoordinate(price);
|
const y = this.candles.priceToCoordinate(price);
|
||||||
if (y == null) return;
|
if (y == null) return;
|
||||||
this.previewLine.removeAttribute('hidden');
|
this.previewLine.removeAttribute('hidden');
|
||||||
|
|
@ -3189,7 +3177,7 @@ class ConfluenceChart {
|
||||||
? sourcePoint - sourceStart
|
? sourcePoint - sourceStart
|
||||||
: this.indexAt(point.t) - this.indexAt(start.t),
|
: this.indexAt(point.t) - this.indexAt(start.t),
|
||||||
);
|
);
|
||||||
const priceShift = this.snapPrice(point.p - start.p);
|
const priceShift = ConfluenceChart.snapToTick(point.p - start.p);
|
||||||
const anchorT = source
|
const anchorT = source
|
||||||
? this.shiftLineTime(original, original.anchor_t, indexShift)
|
? this.shiftLineTime(original, original.anchor_t, indexShift)
|
||||||
: this.timeAtIndex(this.indexAt(original.anchor_t) + indexShift);
|
: this.timeAtIndex(this.indexAt(original.anchor_t) + indexShift);
|
||||||
|
|
@ -3255,7 +3243,7 @@ class ConfluenceChart {
|
||||||
}
|
}
|
||||||
const snapped = this.snapToBars
|
const snapped = this.snapToBars
|
||||||
? this.snapPoint(point)
|
? this.snapPoint(point)
|
||||||
: { ...point, p: this.snapPrice(point.p) };
|
: { ...point, p: ConfluenceChart.snapToTick(point.p) };
|
||||||
if (snapped.t == null || snapped.p == null) return;
|
if (snapped.t == null || snapped.p == null) return;
|
||||||
const time = snapped.t;
|
const time = snapped.t;
|
||||||
const price = snapped.p;
|
const price = snapped.p;
|
||||||
|
|
@ -3383,8 +3371,7 @@ class ConfluenceChart {
|
||||||
|
|
||||||
syncRthLines() {
|
syncRthLines() {
|
||||||
if (!this.rthPrimitive) return;
|
if (!this.rthPrimitive) return;
|
||||||
if (!this.rthEnabled || this.rthMode !== 'spy_rth'
|
if (!this.rthEnabled || !this.bars.length || this.bars[0]?.tf === '1d') {
|
||||||
|| !this.bars.length || this.bars[0]?.tf === '1d') {
|
|
||||||
this.rthPrimitive.setMarks([]);
|
this.rthPrimitive.setMarks([]);
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -138,9 +138,9 @@
|
||||||
<label>Colour<input type="color" v-model="levelColor" aria-label="Level colour"></label>
|
<label>Colour<input type="color" v-model="levelColor" aria-label="Level colour"></label>
|
||||||
<label>Width<select v-model.number="levelWidth" aria-label="Level width"><option v-for="width in 9" :value="width">{{ width }}px</option></select></label>
|
<label>Width<select v-model.number="levelWidth" aria-label="Level width"><option v-for="width in 9" :value="width">{{ width }}px</option></select></label>
|
||||||
</div>
|
</div>
|
||||||
<label>Alert early (pts)<input type="number" min="0" :step="tick" v-model.number="alertEarlyPoints" placeholder="ATR default" aria-label="Alert early points"></label>
|
<label>Alert early (pts)<input type="number" min="0" step="0.25" v-model.number="alertEarlyPoints" placeholder="ATR default" aria-label="Alert early points"></label>
|
||||||
<form class="row price-row" @submit.prevent="addPriceAlert">
|
<form class="row price-row" @submit.prevent="addPriceAlert">
|
||||||
<label>Price<input type="number" :step="tick" v-model.number="alertPrice" :placeholder="price == null ? '0.00' : price.toFixed(2)" aria-label="Level price"></label>
|
<label>Price<input type="number" step="0.25" v-model.number="alertPrice" :placeholder="price == null ? '0.00' : price.toFixed(2)" aria-label="Level price"></label>
|
||||||
<button type="submit" :disabled="!alertPrice">Add</button>
|
<button type="submit" :disabled="!alertPrice">Add</button>
|
||||||
</form>
|
</form>
|
||||||
<p class="hint">Drag on the chart, or type an exact price. Alerts whenever price reaches it, whatever the confluence score.</p>
|
<p class="hint">Drag on the chart, or type an exact price. Alerts whenever price reaches it, whatever the confluence score.</p>
|
||||||
|
|
@ -325,11 +325,11 @@
|
||||||
<span v-else-if="item.kind === 'symbol'">#{{ item.number }} · SYMBOL · {{ item.comment.note }}</span>
|
<span v-else-if="item.kind === 'symbol'">#{{ item.number }} · SYMBOL · {{ item.comment.note }}</span>
|
||||||
<span v-else-if="item.kind === 'fibonacci'">#{{ item.number }} · FIB · {{ item.label }}</span>
|
<span v-else-if="item.kind === 'fibonacci'">#{{ item.number }} · FIB · {{ item.label }}</span>
|
||||||
<span v-else-if="item.kind === 'level'" class="level-editors">
|
<span v-else-if="item.kind === 'level'" class="level-editors">
|
||||||
<label>Price<input type="number" :min="tick" :step="tick" :value="item.line.anchor_p"
|
<label>Price<input type="number" min="0.25" step="0.25" :value="item.line.anchor_p"
|
||||||
aria-label="Level price in drawing list"
|
aria-label="Level price in drawing list"
|
||||||
@keydown.enter="$event.target.blur()"
|
@keydown.enter="$event.target.blur()"
|
||||||
@change="updateLevelNumber(item.line, 'anchor_p', $event.target.value)"></label>
|
@change="updateLevelNumber(item.line, 'anchor_p', $event.target.value)"></label>
|
||||||
<label>Early<input type="number" min="0" :step="tick" :value="item.line.alert_early_points ?? ''"
|
<label>Early<input type="number" min="0" step="0.25" :value="item.line.alert_early_points ?? ''"
|
||||||
placeholder="ATR" aria-label="Level alert early points"
|
placeholder="ATR" aria-label="Level alert early points"
|
||||||
@keydown.enter="$event.target.blur()"
|
@keydown.enter="$event.target.blur()"
|
||||||
@change="updateLevelNumber(item.line, 'alert_early_points', $event.target.value)"></label>
|
@change="updateLevelNumber(item.line, 'alert_early_points', $event.target.value)"></label>
|
||||||
|
|
|
||||||
|
|
@ -58,24 +58,6 @@ def test_corrupt_state_does_not_prevent_alerting(tmp_path):
|
||||||
assert len(engine(tmp_path).evaluate(zone(), 100, 1, 0, "/ES")) == 1
|
assert len(engine(tmp_path).evaluate(zone(), 100, 1, 0, "/ES")) == 1
|
||||||
|
|
||||||
|
|
||||||
def test_an_es_zone_does_not_suppress_the_same_price_on_gold(tmp_path):
|
|
||||||
one = engine(tmp_path)
|
|
||||||
assert len(one.evaluate(zone(), 100, 1, 0, "/ES")) == 1
|
|
||||||
assert len(one.evaluate(zone(), 100, 1, 60, "/GC")) == 1
|
|
||||||
assert one.evaluate(zone(), 100, 1, 90, "ES=F") == []
|
|
||||||
|
|
||||||
|
|
||||||
def test_legacy_fired_rows_without_symbol_are_es(tmp_path):
|
|
||||||
path = tmp_path / "alert_state.json"
|
|
||||||
path.write_text(
|
|
||||||
'{"next_number": 2, "fired": [{"center": 100.0, "at": 0}]}\n',
|
|
||||||
encoding="utf-8",
|
|
||||||
)
|
|
||||||
two = engine(tmp_path)
|
|
||||||
assert two.evaluate(zone(), 100, 1, 60, "/ES") == []
|
|
||||||
assert len(two.evaluate(zone(), 100, 1, 60, "/GC")) == 1
|
|
||||||
|
|
||||||
|
|
||||||
def test_state_file_records_centre_and_time(tmp_path):
|
def test_state_file_records_centre_and_time(tmp_path):
|
||||||
engine(tmp_path).evaluate(zone(), 100, 1, 42, "/ES")
|
engine(tmp_path).evaluate(zone(), 100, 1, 42, "/ES")
|
||||||
payload = json.loads((tmp_path / "alert_state.json").read_text(encoding="utf-8"))
|
payload = json.loads((tmp_path / "alert_state.json").read_text(encoding="utf-8"))
|
||||||
|
|
@ -83,5 +65,4 @@ def test_state_file_records_centre_and_time(tmp_path):
|
||||||
# do not restart from 1 after a deploy and collide with a phone's history.
|
# do not restart from 1 after a deploy and collide with a phone's history.
|
||||||
assert len(payload["fired"]) == 1
|
assert len(payload["fired"]) == 1
|
||||||
assert payload["fired"][0]["at"] == 42
|
assert payload["fired"][0]["at"] == 42
|
||||||
assert payload["fired"][0]["symbol"] == "/ES"
|
|
||||||
assert payload["next_number"] == 2
|
assert payload["next_number"] == 2
|
||||||
|
|
|
||||||
|
|
@ -21,7 +21,6 @@ def test_log_keeps_old_entries_and_recent_is_the_last_day(tmp_path):
|
||||||
assert [entry["message"] for entry in day] == ["dropped", "yesterday"]
|
assert [entry["message"] for entry in day] == ["dropped", "yesterday"]
|
||||||
assert more is True
|
assert more is True
|
||||||
assert EventLog(tmp_path / "events.json")._entries[0]["message"] == "old"
|
assert EventLog(tmp_path / "events.json")._entries[0]["message"] == "old"
|
||||||
assert EventLog(tmp_path / "events.json")._entries[0]["symbol"] == "/ES"
|
|
||||||
|
|
||||||
|
|
||||||
def test_more_pages_older_than_the_cutoff(tmp_path):
|
def test_more_pages_older_than_the_cutoff(tmp_path):
|
||||||
|
|
@ -57,6 +56,3 @@ def test_dispatched_alerts_land_in_the_event_log(tmp_path):
|
||||||
assert "#27" in events[0]["message"]
|
assert "#27" in events[0]["message"]
|
||||||
assert "confluence" not in events[0]["message"]
|
assert "confluence" not in events[0]["message"]
|
||||||
assert snapshot(runtime, Timeframe.M1)["events"][0]["number"] == alerts[0].number
|
assert snapshot(runtime, Timeframe.M1)["events"][0]["number"] == alerts[0].number
|
||||||
assert snapshot(runtime, Timeframe.M1)["events"][0]["symbol"] == "/ES"
|
|
||||||
assert snapshot(runtime, Timeframe.M1)["instrument"]["id"] == "es"
|
|
||||||
assert snapshot(runtime, Timeframe.M1)["instrument"]["tick"] == 0.25
|
|
||||||
|
|
|
||||||
|
|
@ -1,99 +0,0 @@
|
||||||
import asyncio
|
|
||||||
import time
|
|
||||||
|
|
||||||
from app.bars.models import Bar, Timeframe
|
|
||||||
from app.bars.session import bucket_start
|
|
||||||
from app.config import Settings
|
|
||||||
from app.market.stream import StreamService
|
|
||||||
from app.runtime import Runtime
|
|
||||||
|
|
||||||
|
|
||||||
def minute(t, price, closed=True, source="schwab"):
|
|
||||||
return Bar(Timeframe.M1, t, price, price + 1, price - 1, price, 10, closed, "/ES", source)
|
|
||||||
|
|
||||||
|
|
||||||
class History:
|
|
||||||
"""A seed source holding only 1m history, like Yahoo's recent reach."""
|
|
||||||
|
|
||||||
name = "yahoo"
|
|
||||||
delay_minutes = 0
|
|
||||||
|
|
||||||
def __init__(self, bars):
|
|
||||||
self.bars = bars
|
|
||||||
|
|
||||||
def supports_history(self):
|
|
||||||
return True
|
|
||||||
|
|
||||||
async def history(self, symbol, tf, start, end, *, range_=None):
|
|
||||||
if tf is not Timeframe.M1:
|
|
||||||
return []
|
|
||||||
return [bar for bar in self.bars if start <= bar.t < end]
|
|
||||||
|
|
||||||
|
|
||||||
def runtime(tmp_path) -> Runtime:
|
|
||||||
return Runtime(
|
|
||||||
Settings(
|
|
||||||
manual_lines_path=tmp_path / "manual_lines.json",
|
|
||||||
alert_state_path=tmp_path / "alert_state.json",
|
|
||||||
user_prefs_path=tmp_path / "user_prefs.json",
|
|
||||||
events_path=tmp_path / "events.json",
|
|
||||||
)
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
def test_an_outage_is_backfilled_behind_the_live_bars(tmp_path, monkeypatch):
|
|
||||||
# The seed ran only at startup, so a stream that was down for days came
|
|
||||||
# back to live bars with the whole outage still missing.
|
|
||||||
day = bucket_start(int(time.time()) - 86400, Timeframe.D1)
|
|
||||||
instance = runtime(tmp_path)
|
|
||||||
missed = [minute(day + 60 * i, 5000 + i, source="yahoo") for i in range(1, 60)]
|
|
||||||
monkeypatch.setattr("app.runtime.seed_source", lambda settings: History(missed))
|
|
||||||
queue: asyncio.Queue = asyncio.Queue(maxsize=100)
|
|
||||||
instance.subscribers.add(queue)
|
|
||||||
|
|
||||||
async def scenario():
|
|
||||||
await instance.on_bar(minute(day, 4990))
|
|
||||||
# The stream comes back an hour later, into the same day.
|
|
||||||
await instance.on_bar(minute(day + 3600, 6000))
|
|
||||||
await instance.on_bar(minute(day + 3660, 6001))
|
|
||||||
return await instance.fill_gap(day, day + 3600)
|
|
||||||
|
|
||||||
added = asyncio.run(scenario())
|
|
||||||
|
|
||||||
held = [bar.t for bar in instance.store.get(Timeframe.M1)]
|
|
||||||
assert held == [day] + [bar.t for bar in missed] + [day + 3600, day + 3660]
|
|
||||||
assert added > len(missed), "higher timeframes are rebuilt from the recovered minutes"
|
|
||||||
assert day + 300 in [bar.t for bar in instance.store.get(Timeframe.M5)]
|
|
||||||
daily = instance.store.get(Timeframe.D1)[-1]
|
|
||||||
assert daily.t == day
|
|
||||||
assert daily.l == 4989, "the live daily bar must include the pre-reconnect low"
|
|
||||||
assert daily.o == 4990
|
|
||||||
events = []
|
|
||||||
while not queue.empty():
|
|
||||||
events.append(queue.get_nowait()["type"])
|
|
||||||
assert "resync" in events
|
|
||||||
assert "alert" not in events
|
|
||||||
|
|
||||||
|
|
||||||
def test_a_reconnect_past_a_gap_asks_for_a_backfill():
|
|
||||||
class Flaky:
|
|
||||||
name = "schwab"
|
|
||||||
|
|
||||||
def __init__(self):
|
|
||||||
self.connections = [[minute(60, 1)], [minute(60 + 86400, 2)]]
|
|
||||||
|
|
||||||
async def stream(self, symbol):
|
|
||||||
for bar in self.connections.pop(0):
|
|
||||||
yield bar
|
|
||||||
if not self.connections:
|
|
||||||
service.stop()
|
|
||||||
raise RuntimeError("socket closed")
|
|
||||||
|
|
||||||
service = StreamService(Flaky(), "/ES")
|
|
||||||
service.reconnect_seconds = 0
|
|
||||||
resumed = []
|
|
||||||
service.on_resume = lambda after, before: resumed.append((after, before))
|
|
||||||
|
|
||||||
asyncio.run(service.run())
|
|
||||||
|
|
||||||
assert resumed == [(60, 60 + 86400)]
|
|
||||||
|
|
@ -1,60 +0,0 @@
|
||||||
from app.config import Settings
|
|
||||||
from app.instrument import (
|
|
||||||
DEFAULT_SYMBOL, get_instrument, instrument_for_symbol, INSTRUMENTS,
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
def test_settings_default_to_the_es_profile():
|
|
||||||
settings = Settings()
|
|
||||||
assert settings.instrument == "es"
|
|
||||||
assert settings.profile.tick == 0.25
|
|
||||||
assert Settings(instrument="gc").profile.schwab_symbol == "/GC"
|
|
||||||
|
|
||||||
|
|
||||||
def test_es_is_the_default_profile():
|
|
||||||
es = get_instrument("es")
|
|
||||||
assert es.schwab_symbol == DEFAULT_SYMBOL
|
|
||||||
assert es.tick == 0.25
|
|
||||||
assert es.rth == "spy_rth"
|
|
||||||
|
|
||||||
|
|
||||||
def test_nq_shares_es_tick_and_rth():
|
|
||||||
nq = get_instrument("nq")
|
|
||||||
es = get_instrument("es")
|
|
||||||
assert nq.tick == es.tick
|
|
||||||
assert nq.rth == es.rth
|
|
||||||
assert nq.schwab_symbol == "/NQ"
|
|
||||||
|
|
||||||
|
|
||||||
def test_gold_and_oil_are_not_quarter_ticks():
|
|
||||||
assert get_instrument("gc").tick == 0.10
|
|
||||||
assert get_instrument("gc").rth == "none"
|
|
||||||
assert get_instrument("cl").tick == 0.01
|
|
||||||
assert get_instrument("cl").rth == "nymex_day"
|
|
||||||
|
|
||||||
|
|
||||||
def test_yahoo_and_schwab_names_map_to_the_same_root():
|
|
||||||
assert instrument_for_symbol("ES=F").schwab_symbol == "/ES"
|
|
||||||
assert instrument_for_symbol("/ES").id == "es"
|
|
||||||
assert instrument_for_symbol("GC=F").id == "gc"
|
|
||||||
assert instrument_for_symbol("unknown").id == "es"
|
|
||||||
|
|
||||||
|
|
||||||
def test_unknown_instrument_id_is_rejected():
|
|
||||||
try:
|
|
||||||
get_instrument("btc")
|
|
||||||
except ValueError as error:
|
|
||||||
assert "btc" in str(error)
|
|
||||||
else:
|
|
||||||
raise AssertionError("unknown id was accepted")
|
|
||||||
|
|
||||||
|
|
||||||
def test_gold_snap_is_not_a_quarter_point():
|
|
||||||
gc = get_instrument("gc")
|
|
||||||
assert gc.snap(3450.07) == 3450.10
|
|
||||||
assert get_instrument("es").snap(6400.10) == 6400.00
|
|
||||||
assert get_instrument("cl").snap(70.014) == 70.01
|
|
||||||
|
|
||||||
|
|
||||||
def test_every_planned_root_is_in_the_table():
|
|
||||||
assert set(INSTRUMENTS) == {"es", "nq", "gc", "cl"}
|
|
||||||
|
|
@ -9,25 +9,6 @@ def sample_line():
|
||||||
return ManualLine("ml_test", Timeframe.H1, Side.RESISTANCE, 100, 5000, -0.01, 200, 300, number=1)
|
return ManualLine("ml_test", Timeframe.H1, Side.RESISTANCE, 100, 5000, -0.01, 200, 300, number=1)
|
||||||
|
|
||||||
|
|
||||||
def test_a_drawing_without_symbol_loads_as_es(tmp_path):
|
|
||||||
path = tmp_path / "manual_lines.json"
|
|
||||||
path.write_text(
|
|
||||||
'[{"id":"ml_old","tf":"1h","side":"resistance","anchor_t":100,'
|
|
||||||
'"anchor_p":5000,"slope":-0.01,"last_t":200,"created_at":300,"number":1}]\n',
|
|
||||||
encoding="utf-8",
|
|
||||||
)
|
|
||||||
loaded = ManualLineStore(path).lines["ml_old"]
|
|
||||||
assert loaded.symbol == "/ES"
|
|
||||||
|
|
||||||
|
|
||||||
def test_a_gold_drawing_keeps_its_symbol_through_json(tmp_path):
|
|
||||||
path = tmp_path / "manual_lines.json"
|
|
||||||
line = sample_line()
|
|
||||||
line.symbol = "/GC"
|
|
||||||
ManualLineStore(path).add(line)
|
|
||||||
assert ManualLineStore(path).lines["ml_test"].symbol == "/GC"
|
|
||||||
|
|
||||||
|
|
||||||
def test_json_persistence_round_trip(tmp_path):
|
def test_json_persistence_round_trip(tmp_path):
|
||||||
path = tmp_path / "manual_lines.json"
|
path = tmp_path / "manual_lines.json"
|
||||||
store = ManualLineStore(path)
|
store = ManualLineStore(path)
|
||||||
|
|
@ -178,7 +159,6 @@ def test_a_null_cutoff_clears_an_ended_line(tmp_path):
|
||||||
"cutoff_t": 150,
|
"cutoff_t": 150,
|
||||||
}).json()
|
}).json()
|
||||||
assert created["cutoff_t"] == 150
|
assert created["cutoff_t"] == 150
|
||||||
assert created["symbol"] == "/ES"
|
|
||||||
|
|
||||||
cleared = client.patch(f"/api/lines/{created['id']}", json={"cutoff_t": None}).json()
|
cleared = client.patch(f"/api/lines/{created['id']}", json={"cutoff_t": None}).json()
|
||||||
assert cleared["cutoff_t"] is None
|
assert cleared["cutoff_t"] is None
|
||||||
|
|
@ -218,7 +198,6 @@ def test_restoring_a_deleted_line_keeps_its_id_and_number(tmp_path):
|
||||||
}).json()
|
}).json()
|
||||||
assert restored["id"] == line_id
|
assert restored["id"] == line_id
|
||||||
assert restored["number"] == number
|
assert restored["number"] == number
|
||||||
assert restored["symbol"] == "/ES"
|
|
||||||
assert client.post("/api/lines/restore", json={
|
assert client.post("/api/lines/restore", json={
|
||||||
"id": line_id,
|
"id": line_id,
|
||||||
"tf": "1m",
|
"tf": "1m",
|
||||||
|
|
@ -229,30 +208,3 @@ def test_restoring_a_deleted_line_keeps_its_id_and_number(tmp_path):
|
||||||
"last_t": 200,
|
"last_t": 200,
|
||||||
"number": number,
|
"number": number,
|
||||||
}).status_code == 409
|
}).status_code == 409
|
||||||
|
|
||||||
|
|
||||||
def test_restoring_keeps_a_non_es_symbol(tmp_path):
|
|
||||||
from fastapi import FastAPI
|
|
||||||
from fastapi.testclient import TestClient
|
|
||||||
|
|
||||||
from app.api.routes import router
|
|
||||||
from app.config import Settings
|
|
||||||
from app.runtime import Runtime
|
|
||||||
|
|
||||||
app = FastAPI()
|
|
||||||
app.include_router(router)
|
|
||||||
app.state.runtime = Runtime(Settings(manual_lines_path=tmp_path / "manual_lines.json"))
|
|
||||||
client = TestClient(app)
|
|
||||||
restored = client.post("/api/lines/restore", json={
|
|
||||||
"id": "ml_gold",
|
|
||||||
"tf": "1m",
|
|
||||||
"side": "support",
|
|
||||||
"anchor_t": 100,
|
|
||||||
"anchor_p": 1.0,
|
|
||||||
"slope": 0.01,
|
|
||||||
"last_t": 200,
|
|
||||||
"number": 9,
|
|
||||||
"symbol": "/GC",
|
|
||||||
}).json()
|
|
||||||
assert restored["symbol"] == "/GC"
|
|
||||||
assert ManualLineStore(tmp_path / "manual_lines.json").lines["ml_gold"].symbol == "/GC"
|
|
||||||
|
|
|
||||||
|
|
@ -48,22 +48,3 @@ def test_a_tick_cannot_overwrite_a_settled_bar():
|
||||||
|
|
||||||
held = store.get(Timeframe.M1)[0]
|
held = store.get(Timeframe.M1)[0]
|
||||||
assert held.closed is True and held.v == 400
|
assert held.closed is True and held.v == 400
|
||||||
|
|
||||||
|
|
||||||
def test_a_backfilled_hole_lands_behind_live_bars_without_replacing_them():
|
|
||||||
# After an outage the live stream is already newer than the hole, so put()
|
|
||||||
# dropped every recovered bar and fifteen missed days stayed missing.
|
|
||||||
store = InMemoryBarStore(4)
|
|
||||||
store.put(bar(60))
|
|
||||||
store.put(bar(600, close=7))
|
|
||||||
|
|
||||||
added = store.fill([bar(120), bar(180), bar(600, close=99)])
|
|
||||||
|
|
||||||
assert added == 2
|
|
||||||
assert [value.t for value in store.get(Timeframe.M1)] == [60, 120, 180, 600]
|
|
||||||
assert store.get(Timeframe.M1)[-1].c == 7, "the live source's bar must win"
|
|
||||||
|
|
||||||
store.fill([bar(240)])
|
|
||||||
assert [value.t for value in store.get(Timeframe.M1)] == [120, 180, 240, 600], (
|
|
||||||
"the cap keeps the newest history"
|
|
||||||
)
|
|
||||||
|
|
|
||||||
|
|
@ -82,8 +82,6 @@ def test_quote_change_always_uses_the_daily_session_open(tmp_path):
|
||||||
message = snapshot(runtime, Timeframe.H1)
|
message = snapshot(runtime, Timeframe.H1)
|
||||||
|
|
||||||
assert message["session_open"] == 6123.25
|
assert message["session_open"] == 6123.25
|
||||||
assert message["instrument"]["schwab_symbol"] == "/ES"
|
|
||||||
assert message["instrument"]["tick"] == 0.25
|
|
||||||
|
|
||||||
|
|
||||||
def test_snapshot_session_range_uses_the_forming_daily_bar(tmp_path):
|
def test_snapshot_session_range_uses_the_forming_daily_bar(tmp_path):
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue