Compare commits

..

No commits in common. "fced18528b8f32c00ccc370b4094c6065163fbeb" and "d0a0f9f6d6b5bb38158d20c4c67d15afce3c29d0" have entirely different histories.

25 changed files with 36 additions and 774 deletions

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

@ -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');

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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