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/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).
|
||||
Adding `/NQ` `/GC` `/CL` is [`docs/investigate_added_symbols.md`](docs/investigate_added_symbols.md).
|
||||
|
||||
## 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.levels import Level, LevelKind, Side
|
||||
from app.instrument import DEFAULT_SYMBOL, instrument_for_symbol
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
|
@ -44,7 +43,6 @@ class Alert:
|
|||
class _Fired:
|
||||
center: float
|
||||
at: int
|
||||
symbol: str = DEFAULT_SYMBOL
|
||||
|
||||
|
||||
class AlertEngine:
|
||||
|
|
@ -85,14 +83,7 @@ class AlertEngine:
|
|||
if isinstance(payload, dict):
|
||||
self._next_number = int(payload.get("next_number", 1))
|
||||
payload = payload.get("fired", [])
|
||||
return [
|
||||
_Fired(
|
||||
float(item["center"]),
|
||||
int(item["at"]),
|
||||
str(item.get("symbol") or DEFAULT_SYMBOL),
|
||||
)
|
||||
for item in payload
|
||||
]
|
||||
return [_Fired(float(item["center"]), int(item["at"])) for item in payload]
|
||||
except Exception:
|
||||
# Corrupt state costs one burst of duplicate alerts, which is a far
|
||||
# better failure than refusing to start the stream.
|
||||
|
|
@ -124,12 +115,7 @@ class AlertEngine:
|
|||
{
|
||||
"next_number": self._next_number,
|
||||
"fired": [
|
||||
{
|
||||
"center": entry.center,
|
||||
"at": entry.at,
|
||||
"symbol": entry.symbol,
|
||||
}
|
||||
for entry in self._fired
|
||||
{"center": entry.center, "at": entry.at} for entry in self._fired
|
||||
],
|
||||
},
|
||||
indent=2,
|
||||
|
|
@ -156,7 +142,6 @@ class AlertEngine:
|
|||
tolerance = 0.5 * atr15
|
||||
if tolerance <= 0:
|
||||
return []
|
||||
root = instrument_for_symbol(symbol).schwab_symbol
|
||||
# Two zones within an ATR of each other are the same zone as far as
|
||||
# being told about them goes.
|
||||
merge_distance = 2 * tolerance
|
||||
|
|
@ -210,12 +195,10 @@ class AlertEngine:
|
|||
# was oscillating on re-alerted on every crossing — which is exactly
|
||||
# when a level is least newsworthy, not most.
|
||||
if any(
|
||||
entry.symbol == root
|
||||
and abs(entry.center - cluster.center) <= merge_distance
|
||||
for entry in self._fired
|
||||
abs(entry.center - cluster.center) <= merge_distance for entry in self._fired
|
||||
):
|
||||
continue
|
||||
self._fired.append(_Fired(cluster.center, now, root))
|
||||
self._fired.append(_Fired(cluster.center, now))
|
||||
changed = True
|
||||
direction = "BEARISH" if cluster.side.value == "resistance" else "BULLISH"
|
||||
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
|
||||
if abs(price - current_price) > tolerance:
|
||||
continue
|
||||
if any(
|
||||
entry.symbol == root
|
||||
and abs(entry.center - price) <= merge_distance
|
||||
for entry in self._fired
|
||||
):
|
||||
if any(abs(entry.center - price) <= merge_distance for entry in self._fired):
|
||||
continue
|
||||
self._fired.append(_Fired(price, now, root))
|
||||
self._fired.append(_Fired(price, now))
|
||||
changed = True
|
||||
direction = "BEARISH" if level.side is Side.RESISTANCE else "BULLISH"
|
||||
name = f"{level.period} DMA" if level.period else level.label
|
||||
|
|
|
|||
|
|
@ -3,8 +3,6 @@ import logging
|
|||
import time
|
||||
from pathlib import Path
|
||||
|
||||
from app.instrument import DEFAULT_SYMBOL
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
DAY_SECONDS = 86400
|
||||
|
|
@ -41,16 +39,12 @@ class EventLog:
|
|||
except Exception:
|
||||
logger.warning("Could not persist event log", exc_info=True)
|
||||
|
||||
def add(
|
||||
self, kind: str, message: str, *, number: int | None = None,
|
||||
at: int | None = None, symbol: str | None = None,
|
||||
) -> dict:
|
||||
def add(self, kind: str, message: str, *, number: int | None = None, at: int | None = None) -> dict:
|
||||
entry = {
|
||||
"kind": kind,
|
||||
"message": message,
|
||||
"number": number,
|
||||
"at": int(at if at is not None else time.time()),
|
||||
"symbol": symbol or DEFAULT_SYMBOL,
|
||||
}
|
||||
self._entries.append(entry)
|
||||
self._save()
|
||||
|
|
|
|||
|
|
@ -3,7 +3,6 @@ from enum import Enum
|
|||
from typing import Any
|
||||
|
||||
from app.bars.models import Timeframe
|
||||
from app.instrument import DEFAULT_SYMBOL
|
||||
|
||||
|
||||
class LevelKind(str, Enum):
|
||||
|
|
@ -56,7 +55,6 @@ class Level:
|
|||
# price this line safely. Such a line remains visible but cannot cluster or
|
||||
# alert using the absolute-time fallback.
|
||||
geometry_resolved: bool = True
|
||||
symbol: str = DEFAULT_SYMBOL
|
||||
|
||||
def price_at(self, t: int) -> float:
|
||||
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.bars.models import Timeframe
|
||||
from app.config import TIMEFRAME_WEIGHT
|
||||
from app.instrument import DEFAULT_SYMBOL
|
||||
|
||||
|
||||
@dataclass(slots=True)
|
||||
|
|
@ -38,7 +37,6 @@ class ManualLine:
|
|||
# Visual multiplier for marks. 1 is the original 30px glyph; the pin stays
|
||||
# on (anchor_t, anchor_p) regardless of this value.
|
||||
scale: float = 1.0
|
||||
symbol: str = DEFAULT_SYMBOL
|
||||
|
||||
@property
|
||||
def drawing_kind(self) -> str:
|
||||
|
|
@ -95,7 +93,6 @@ class ManualLine:
|
|||
cutoff_t=self.cutoff_t,
|
||||
armed=self.armed,
|
||||
alert_early_points=self.alert_early_points,
|
||||
symbol=self.symbol,
|
||||
)
|
||||
|
||||
def to_dict(self) -> dict:
|
||||
|
|
@ -133,7 +130,6 @@ class ManualLine:
|
|||
if value.get("alert_early_points") is not None else None
|
||||
),
|
||||
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 = ""
|
||||
alert_early_points: float | None = None
|
||||
scale: float = Field(1.0, ge=0.5, le=3)
|
||||
symbol: str = ""
|
||||
|
||||
|
||||
class LinePatch(BaseModel):
|
||||
|
|
@ -176,7 +175,6 @@ def status(request: Request):
|
|||
"worst": round(runtime.loop_lag_worst * 1000, 1),
|
||||
},
|
||||
"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,
|
||||
armed=payload.armed,
|
||||
kind=payload.kind,
|
||||
symbol=request.app.state.runtime.settings.profile.schwab_symbol,
|
||||
)
|
||||
runtime = request.app.state.runtime
|
||||
line = runtime.manual_lines.add(line)
|
||||
|
|
@ -322,7 +319,6 @@ def restore_line(request: Request, payload: LineRestore):
|
|||
icon=payload.icon,
|
||||
alert_early_points=payload.alert_early_points,
|
||||
scale=payload.scale,
|
||||
symbol=payload.symbol or runtime.settings.profile.schwab_symbol,
|
||||
)
|
||||
line = runtime.manual_lines.add(line)
|
||||
runtime.rebuild_levels()
|
||||
|
|
@ -352,7 +348,6 @@ def create_price_alert(request: Request, payload: PriceAlertCreate):
|
|||
color=payload.color,
|
||||
line_width=payload.line_width,
|
||||
alert_early_points=payload.alert_early_points,
|
||||
symbol=runtime.settings.profile.schwab_symbol,
|
||||
)
|
||||
line = runtime.manual_lines.add(line)
|
||||
runtime.rebuild_levels()
|
||||
|
|
@ -405,7 +400,6 @@ def create_comment(request: Request, payload: CommentCreate):
|
|||
y=payload.y,
|
||||
# A comment must never alert, whatever else changes around it.
|
||||
armed=False,
|
||||
symbol=runtime.settings.profile.schwab_symbol,
|
||||
)
|
||||
line = runtime.manual_lines.add(line)
|
||||
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),
|
||||
"events": events,
|
||||
"events_more": events_more,
|
||||
"instrument": runtime.settings.profile.payload(),
|
||||
}
|
||||
|
||||
|
||||
|
|
@ -251,12 +250,7 @@ async def websocket_endpoint(websocket: WebSocket):
|
|||
event = await queue.get()
|
||||
if event["type"] == "disconnect":
|
||||
break
|
||||
if event["type"] == "resync":
|
||||
# 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":
|
||||
if event["type"] == "bar":
|
||||
bar = event["bar"]
|
||||
full = last_bar_t.get(bar.tf) != bar.t
|
||||
if full:
|
||||
|
|
|
|||
|
|
@ -46,35 +46,6 @@ class InMemoryBarStore:
|
|||
# Buckets are ordered, so nothing further back can match.
|
||||
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]:
|
||||
bars = list(self._bars[tf])
|
||||
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 app.bars.models import Timeframe
|
||||
from app.instrument import get_instrument
|
||||
|
||||
|
||||
DEFAULT_MAX_BARS_PER_TF = 5000
|
||||
|
|
@ -25,7 +24,6 @@ class Settings(BaseSettings):
|
|||
|
||||
live_source: str = "yahoo"
|
||||
seed_source: str = "yahoo"
|
||||
instrument: str = "es"
|
||||
yahoo_symbol: str = "ES=F"
|
||||
yahoo_poll_seconds: float = 20
|
||||
seed_1h_range: str = "730d"
|
||||
|
|
@ -75,10 +73,6 @@ class Settings(BaseSettings):
|
|||
alert_timezone: str = "America/Chicago"
|
||||
replay_file: Path | None = None
|
||||
|
||||
@property
|
||||
def profile(self):
|
||||
return get_instrument(self.instrument)
|
||||
|
||||
@property
|
||||
def live_symbol(self) -> str:
|
||||
"""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._stop = asyncio.Event()
|
||||
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:
|
||||
self._handlers.append(handler)
|
||||
|
|
@ -51,21 +45,11 @@ class StreamService:
|
|||
|
||||
async def run(self) -> None:
|
||||
while not self._stop.is_set():
|
||||
first = True
|
||||
try:
|
||||
async for bar in self.source.stream(self.symbol):
|
||||
self.status = "replay" if self.source.name == "replay" else "connected"
|
||||
self.last_error = None
|
||||
before = self.last_bar_t
|
||||
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():
|
||||
break
|
||||
if self.source.name == "replay":
|
||||
|
|
@ -80,7 +64,7 @@ class StreamService:
|
|||
self.on_drop(str(exc))
|
||||
self.status = "disconnected"
|
||||
try:
|
||||
await asyncio.wait_for(self._stop.wait(), timeout=self.reconnect_seconds)
|
||||
await asyncio.wait_for(self._stop.wait(), timeout=5)
|
||||
except TimeoutError:
|
||||
pass
|
||||
|
||||
|
|
|
|||
111
app/runtime.py
111
app/runtime.py
|
|
@ -27,15 +27,6 @@ from app.notify.ntfy import send_ntfy
|
|||
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
|
||||
class Runtime:
|
||||
settings: Settings
|
||||
|
|
@ -53,7 +44,6 @@ class Runtime:
|
|||
alert_engine: AlertEngine = field(init=False)
|
||||
_sent_levels: dict[str, dict] = field(default_factory=dict)
|
||||
_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;
|
||||
# None while a test drives the runtime directly.
|
||||
_loop: asyncio.AbstractEventLoop | None = None
|
||||
|
|
@ -91,7 +81,6 @@ class Runtime:
|
|||
self.stream = StreamService(live_source(self.settings), self.settings.live_symbol)
|
||||
self.stream.add_handler(self.on_bar)
|
||||
self.stream.on_drop = self._on_stream_drop
|
||||
self.stream.on_resume = self._on_stream_resume
|
||||
|
||||
async def on_bar(self, bar: Bar) -> None:
|
||||
# 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:
|
||||
kind = "auth" if "invalid_grant" in error or "Refresh token" in error else "stream"
|
||||
self.events.add(
|
||||
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))
|
||||
self.events.add(kind, error.split("\n", 1)[0][:200])
|
||||
|
||||
def dispatch_alerts(self, alerts: list[Alert]) -> None:
|
||||
tripped: set[str] = set()
|
||||
for alert in alerts:
|
||||
self.events.add(
|
||||
"alert", alert.message, number=alert.number, at=alert.at,
|
||||
symbol=self.settings.profile.schwab_symbol,
|
||||
)
|
||||
self.events.add("alert", alert.message, number=alert.number, at=alert.at)
|
||||
self.broadcast({
|
||||
"type": "alert",
|
||||
"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
|
||||
anchors, including on the body between them and at `last_t` itself (that
|
||||
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
|
||||
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
|
||||
|
||||
Yahoo anchors `ES=F` daily bars to **midnight ET**, but the CME futures session runs
|
||||
|
|
@ -826,8 +816,8 @@ Server → client:
|
|||
Rules:
|
||||
- Send `bar` on **every** update of the forming bar (that is the live chart) but batch
|
||||
`levels` — they only change on higher-TF closes.
|
||||
- Always send a full `snapshot` on connect and after any reconnect, and after the
|
||||
server backfills a gap (`resync`). The client must never try to reconcile a gap.
|
||||
- Always send a full `snapshot` on connect and after any reconnect. The client must
|
||||
never try to reconcile a gap.
|
||||
- Levels are sent for **all** timeframes regardless of the displayed timeframe. That is
|
||||
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 autoScrollLivePrice = ref(localStorage.getItem('chart-auto-scroll-live-price') !== 'false');
|
||||
const extraDetail = ref(localStorage.getItem('chart-extra-detail') === 'true');
|
||||
const tick = ref(0.25);
|
||||
const sessionRange = ref(localStorage.getItem('chart-session-range') !== 'false');
|
||||
const hideLowerTfDrawings = ref(localStorage.getItem('chart-hide-lower-tf-drawings') !== 'false');
|
||||
const optionPrefs = (() => {
|
||||
|
|
@ -366,13 +365,7 @@ createApp({
|
|||
|
||||
async function refreshStatus() {
|
||||
const response = await apiFetch('/api/status');
|
||||
if (response.ok) {
|
||||
status.value = await response.json();
|
||||
if (status.value.instrument && chartApi) {
|
||||
chartApi.setInstrument(status.value.instrument);
|
||||
tick.value = status.value.instrument.tick || tick.value;
|
||||
}
|
||||
}
|
||||
if (response.ok) status.value = await response.json();
|
||||
}
|
||||
|
||||
function selectedExpiration() {
|
||||
|
|
@ -532,10 +525,6 @@ createApp({
|
|||
if (message.type === 'snapshot') {
|
||||
dataReceivedAt.value = Date.now();
|
||||
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.setBars(message.bars);
|
||||
levels.value = message.levels || [];
|
||||
|
|
@ -1365,14 +1354,14 @@ createApp({
|
|||
if (item.line) {
|
||||
const line = { ...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 (vertical) await updateLineGeometry(line, previous);
|
||||
continue;
|
||||
}
|
||||
const endPrice = item.line.anchor_p
|
||||
+ item.line.slope * (item.line.last_t - item.line.anchor_t)
|
||||
+ vertical * chartApi.tick;
|
||||
+ vertical * ConfluenceChart.TICK;
|
||||
if (horizontal) {
|
||||
if (line.geometry_resolved === false) continue;
|
||||
const anchorT = chartApi.shiftLineTime(line, line.anchor_t, horizontal);
|
||||
|
|
@ -1394,8 +1383,8 @@ createApp({
|
|||
if (comment.pinned) {
|
||||
const changes = {};
|
||||
if (vertical) {
|
||||
changes.anchor_p = chartApi.snapPrice(
|
||||
comment.anchor_p + vertical * chartApi.tick,
|
||||
changes.anchor_p = ConfluenceChart.snapToTick(
|
||||
comment.anchor_p + vertical * ConfluenceChart.TICK,
|
||||
);
|
||||
}
|
||||
if (horizontal) {
|
||||
|
|
@ -1708,6 +1697,6 @@ createApp({
|
|||
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');
|
||||
|
|
|
|||
|
|
@ -294,8 +294,6 @@ class ConfluenceChart {
|
|||
this.lastCurrentPrice = null;
|
||||
this.pendingBar = null;
|
||||
this.pendingBarFrame = null;
|
||||
this.tick = 0.25;
|
||||
this.rthMode = 'spy_rth';
|
||||
this.previewLine = null;
|
||||
this.lineBridgeLayer = null;
|
||||
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
|
||||
// rather than a click. Wide enough to survive a twitch on a deliberate click.
|
||||
static DRAG_THRESHOLD = 12;
|
||||
|
|
@ -432,20 +432,8 @@ class ConfluenceChart {
|
|||
return localStorage.getItem('chart-diag') === '1';
|
||||
}
|
||||
|
||||
static snapToTick(price, tick = 0.25) {
|
||||
const step = tick > 0 ? tick : 0.25;
|
||||
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();
|
||||
static snapToTick(price) {
|
||||
return Math.round(price / ConfluenceChart.TICK) * ConfluenceChart.TICK;
|
||||
}
|
||||
|
||||
create(el) {
|
||||
|
|
@ -1425,7 +1413,7 @@ class ConfluenceChart {
|
|||
if (point && this.withinPlot(point) && this.onCommentMove) {
|
||||
this.onCommentMove(node.comment, {
|
||||
anchor_t: Math.round(point.t),
|
||||
anchor_p: this.snapPrice(point.p),
|
||||
anchor_p: ConfluenceChart.snapToTick(point.p),
|
||||
});
|
||||
} else {
|
||||
this.renderComments();
|
||||
|
|
@ -1840,7 +1828,7 @@ class ConfluenceChart {
|
|||
if (!point || !this.withinPlot(point) || point.t == null) return null;
|
||||
return {
|
||||
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),
|
||||
y: Math.min(Math.max(point.y / this.overlayLayer.clientHeight, 0), 1),
|
||||
};
|
||||
|
|
@ -2168,7 +2156,7 @@ class ConfluenceChart {
|
|||
this.onToolComplete?.({
|
||||
tool,
|
||||
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),
|
||||
y: Math.min(Math.max(end.y / this.chartEl.clientHeight, 0), 1),
|
||||
});
|
||||
|
|
@ -2178,7 +2166,7 @@ class ConfluenceChart {
|
|||
if (tool === 'level') {
|
||||
// A click with no drag is a valid placement; the drag is only there to
|
||||
// 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;
|
||||
}
|
||||
|
||||
|
|
@ -2261,7 +2249,7 @@ class ConfluenceChart {
|
|||
const last = this.bars[this.bars.length - 1];
|
||||
if (fallbackT > last.t) {
|
||||
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
|
||||
// rather than measuring against undefined.
|
||||
|
|
@ -2290,7 +2278,7 @@ class ConfluenceChart {
|
|||
if (!this.gesture) return;
|
||||
const { start, end } = this.gesture;
|
||||
if (this.armedTool === 'level') {
|
||||
const price = this.snapPrice(end.p);
|
||||
const price = ConfluenceChart.snapToTick(end.p);
|
||||
const y = this.candles.priceToCoordinate(price);
|
||||
if (y == null) return;
|
||||
this.previewLine.removeAttribute('hidden');
|
||||
|
|
@ -3189,7 +3177,7 @@ class ConfluenceChart {
|
|||
? sourcePoint - sourceStart
|
||||
: 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
|
||||
? this.shiftLineTime(original, original.anchor_t, indexShift)
|
||||
: this.timeAtIndex(this.indexAt(original.anchor_t) + indexShift);
|
||||
|
|
@ -3255,7 +3243,7 @@ class ConfluenceChart {
|
|||
}
|
||||
const snapped = this.snapToBars
|
||||
? this.snapPoint(point)
|
||||
: { ...point, p: this.snapPrice(point.p) };
|
||||
: { ...point, p: ConfluenceChart.snapToTick(point.p) };
|
||||
if (snapped.t == null || snapped.p == null) return;
|
||||
const time = snapped.t;
|
||||
const price = snapped.p;
|
||||
|
|
@ -3383,8 +3371,7 @@ class ConfluenceChart {
|
|||
|
||||
syncRthLines() {
|
||||
if (!this.rthPrimitive) return;
|
||||
if (!this.rthEnabled || this.rthMode !== 'spy_rth'
|
||||
|| !this.bars.length || this.bars[0]?.tf === '1d') {
|
||||
if (!this.rthEnabled || !this.bars.length || this.bars[0]?.tf === '1d') {
|
||||
this.rthPrimitive.setMarks([]);
|
||||
return;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -138,9 +138,9 @@
|
|||
<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>
|
||||
</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">
|
||||
<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>
|
||||
</form>
|
||||
<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 === 'fibonacci'">#{{ item.number }} · FIB · {{ item.label }}</span>
|
||||
<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"
|
||||
@keydown.enter="$event.target.blur()"
|
||||
@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"
|
||||
@keydown.enter="$event.target.blur()"
|
||||
@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
|
||||
|
||||
|
||||
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):
|
||||
engine(tmp_path).evaluate(zone(), 100, 1, 42, "/ES")
|
||||
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.
|
||||
assert len(payload["fired"]) == 1
|
||||
assert payload["fired"][0]["at"] == 42
|
||||
assert payload["fired"][0]["symbol"] == "/ES"
|
||||
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 more is True
|
||||
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):
|
||||
|
|
@ -57,6 +56,3 @@ def test_dispatched_alerts_land_in_the_event_log(tmp_path):
|
|||
assert "#27" 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]["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)
|
||||
|
||||
|
||||
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):
|
||||
path = tmp_path / "manual_lines.json"
|
||||
store = ManualLineStore(path)
|
||||
|
|
@ -178,7 +159,6 @@ def test_a_null_cutoff_clears_an_ended_line(tmp_path):
|
|||
"cutoff_t": 150,
|
||||
}).json()
|
||||
assert created["cutoff_t"] == 150
|
||||
assert created["symbol"] == "/ES"
|
||||
|
||||
cleared = client.patch(f"/api/lines/{created['id']}", json={"cutoff_t": None}).json()
|
||||
assert cleared["cutoff_t"] is None
|
||||
|
|
@ -218,7 +198,6 @@ def test_restoring_a_deleted_line_keeps_its_id_and_number(tmp_path):
|
|||
}).json()
|
||||
assert restored["id"] == line_id
|
||||
assert restored["number"] == number
|
||||
assert restored["symbol"] == "/ES"
|
||||
assert client.post("/api/lines/restore", json={
|
||||
"id": line_id,
|
||||
"tf": "1m",
|
||||
|
|
@ -229,30 +208,3 @@ def test_restoring_a_deleted_line_keeps_its_id_and_number(tmp_path):
|
|||
"last_t": 200,
|
||||
"number": number,
|
||||
}).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]
|
||||
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)
|
||||
|
||||
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):
|
||||
|
|
|
|||
Loading…
Reference in a new issue