History was only fetched by the startup seed, so a Schwab outage stayed a hole until the next deploy — Sept 9 to 24 after a refresh token expired. A reconnect more than two minutes past the last bar now fetches the gap from Yahoo, fills empty buckets only, refolds the live forming buckets, rebuilds levels without alerting, and resyncs every socket. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
589 lines
26 KiB
Python
589 lines
26 KiB
Python
import asyncio
|
|
import logging
|
|
import time
|
|
from dataclasses import dataclass, field, replace
|
|
from typing import ClassVar
|
|
|
|
from app.analysis.alerts import Alert, AlertEngine
|
|
from app.analysis.event_log import EventLog
|
|
from app.bars.models import Bar, Timeframe
|
|
from app.bars.aggregator import Aggregator
|
|
from app.bars.session import bucket_start
|
|
from app.analysis.bar_space import price_in_bar_space, price_in_timeframe_space
|
|
from app.analysis.horizontals import build_prior_day_levels
|
|
from app.analysis.levels import Level, LevelKind
|
|
from app.analysis.moving_averages import build_ma_levels
|
|
from app.analysis.vwap import build_vwap_level
|
|
from app.analysis.confluence import Cluster, cluster_levels
|
|
from app.analysis.indicators import atr
|
|
from app.analysis.manual_lines import ManualLineStore
|
|
from app.analysis.user_prefs import UserPrefStore
|
|
from app.bars.store import InMemoryBarStore
|
|
from app.config import Settings
|
|
from app.market.factory import live_source, seed_source
|
|
from app.market.stream import StreamService
|
|
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
|
|
store: InMemoryBarStore = field(init=False)
|
|
stream: StreamService = field(init=False)
|
|
subscribers: set[asyncio.Queue[dict]] = field(default_factory=set)
|
|
aggregator: Aggregator = field(init=False)
|
|
levels: list[Level] = field(default_factory=list)
|
|
clusters: list[Cluster] = field(default_factory=list)
|
|
price: float | None = None
|
|
atr15: float = 0.0
|
|
manual_lines: ManualLineStore = field(init=False)
|
|
user_prefs: UserPrefStore = field(init=False)
|
|
ma_levels: list[Level] = field(default_factory=list)
|
|
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
|
|
# True while history is being replayed, so on_bar accumulates quietly.
|
|
seeding: bool = False
|
|
# Coalescing window for level rebuilds; see request_rebuild.
|
|
REBUILD_INTERVAL: ClassVar[float] = 0.25
|
|
_last_rebuild: float = 0.0
|
|
_rebuild_pending: bool = False
|
|
# Seconds the loop ran late, worst since start and most recent sample.
|
|
loop_lag_worst: float = 0.0
|
|
loop_lag_recent: float = 0.0
|
|
_lag_task: asyncio.Task | None = None
|
|
_rebuild_task: asyncio.Task | None = None
|
|
_token_task: asyncio.Task | None = None
|
|
schwab_login: object | None = None
|
|
events: EventLog = field(init=False)
|
|
|
|
def __post_init__(self) -> None:
|
|
self.store = InMemoryBarStore(self.settings.max_bars_per_tf)
|
|
self.aggregator = Aggregator(self.settings.enabled_timeframes)
|
|
self.manual_lines = ManualLineStore(self.settings.manual_lines_path)
|
|
self.user_prefs = UserPrefStore(self.settings.user_prefs_path)
|
|
self.events = EventLog(self.settings.events_path)
|
|
# One engine for the process, not one per browser connection. Cooldowns
|
|
# are only meaningful if they outlive a page reload, and a phone push
|
|
# must not depend on a tab being open to produce it.
|
|
self.alert_engine = AlertEngine(
|
|
self.settings.confluence_min_score,
|
|
self.settings.alert_cooldown_seconds,
|
|
self.settings.alert_state_path,
|
|
self.settings.alert_timezone,
|
|
)
|
|
self.levels = self.manual_lines.levels()
|
|
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
|
|
# updates the last candle and the live price, and stops there.
|
|
#
|
|
# It must not reach the aggregator: that accumulates volume with
|
|
# `current.v += incoming.v`, so re-sending the same forming minute would
|
|
# add its volume to every higher timeframe again on each update. Alerts
|
|
# stay on closed bars for the same reason they always were — a level is
|
|
# judged on a settled bar, not on a price that may not last the minute.
|
|
if not bar.closed:
|
|
self.store.put(bar)
|
|
self.price = bar.c
|
|
self.broadcast({"type": "bar", "bar": bar})
|
|
for provisional in self.provisional_higher(bar):
|
|
self.store.put(provisional)
|
|
self.broadcast({"type": "bar", "bar": provisional})
|
|
return
|
|
|
|
alert_bar: Bar | None = None
|
|
for aggregated in self.aggregator.update(bar):
|
|
self.store.put(aggregated)
|
|
if self.seeding:
|
|
# Seeding replays years of history through this method. Every
|
|
# bar used to broadcast to nobody and rebuild every level from
|
|
# scratch, which is the whole of the startup stall — 19.5
|
|
# seconds of measured loop lag and two minutes before the port
|
|
# opened. Bars are accumulated here and the levels are built
|
|
# once, at the end, from the finished store.
|
|
continue
|
|
self.broadcast({"type": "bar", "bar": aggregated})
|
|
if self.settings.ma_sets.get(aggregated.tf):
|
|
self.request_rebuild()
|
|
if aggregated.tf is Timeframe.M1 and aggregated.closed:
|
|
self.price = aggregated.c
|
|
alert_bar = aggregated
|
|
if alert_bar is not None:
|
|
values = atr(self.store.get(Timeframe.M15), 14)
|
|
self.atr15 = next((value for value in reversed(values) if value is not None), 0.0)
|
|
# VWAP re-prices every minute, so levels are rebuilt here too. The
|
|
# broadcast is a delta, which is what keeps that affordable.
|
|
self.request_rebuild()
|
|
self.rebuild_clusters(
|
|
evaluate_alerts=True,
|
|
alert_t=alert_bar.t,
|
|
alert_price=alert_bar.c,
|
|
)
|
|
|
|
def provisional_higher(self, bar: Bar) -> list[Bar]:
|
|
"""Higher-timeframe bars including the minute still being traded.
|
|
|
|
The aggregator cannot be asked for these: it accumulates with
|
|
``current.v += incoming.v``, so re-feeding the same forming minute on
|
|
every tick would add its volume to each higher timeframe again and
|
|
again. Its committed state — every minute that has actually closed — is
|
|
combined with the live minute here instead, without mutating it. The
|
|
next closed minute goes through the aggregator normally and replaces
|
|
what this produced, because the store keys on the bucket's timestamp.
|
|
"""
|
|
out: list[Bar] = []
|
|
for tf in self.settings.enabled_timeframes:
|
|
if tf is Timeframe.M1:
|
|
continue
|
|
start = bucket_start(bar.t, tf)
|
|
base = self.aggregator.forming.get(tf)
|
|
if base is None or start > base.t:
|
|
# The live minute opens a bucket the aggregator has not started.
|
|
out.append(replace(bar, tf=tf, t=start, closed=False))
|
|
continue
|
|
if start < base.t:
|
|
continue
|
|
out.append(
|
|
replace(
|
|
base,
|
|
h=max(base.h, bar.h),
|
|
l=min(base.l, bar.l),
|
|
c=bar.c,
|
|
v=base.v + bar.v,
|
|
closed=False,
|
|
)
|
|
)
|
|
return out
|
|
|
|
def broadcast(self, event: dict) -> None:
|
|
"""Publish an event to every subscriber, from any thread.
|
|
|
|
asyncio.Queue is not thread-safe: it wakes a waiting consumer by
|
|
resolving a Future, which only the loop thread may do. The mutating
|
|
routes are sync `def`, so FastAPI runs them in a threadpool, and they
|
|
reach here through rebuild_levels — writing the queue directly from
|
|
there can drop a socket's wakeup. The visible symptom is a drawing made
|
|
in one browser not reaching another until the next market tick, which
|
|
is why it has gone unnoticed: the stream ticks about once a second and
|
|
covers it over.
|
|
"""
|
|
loop = self._loop
|
|
if loop is None or self._on_loop_thread(loop):
|
|
self._publish(event)
|
|
return
|
|
loop.call_soon_threadsafe(self._publish, event)
|
|
|
|
@staticmethod
|
|
def _on_loop_thread(loop: asyncio.AbstractEventLoop) -> bool:
|
|
try:
|
|
return asyncio.get_running_loop() is loop
|
|
except RuntimeError:
|
|
# No loop in this thread at all, so certainly not that one.
|
|
return False
|
|
|
|
def _publish(self, event: dict) -> None:
|
|
for queue in self.subscribers.copy():
|
|
if queue.full():
|
|
queue.get_nowait()
|
|
queue.put_nowait(event)
|
|
|
|
def rebuild_levels(self) -> None:
|
|
self.ma_levels = build_ma_levels(
|
|
{tf: self.store.get(tf) for tf in self.settings.ma_sets},
|
|
self.settings.ma_sets,
|
|
)
|
|
minute_bars = self.store.get(Timeframe.M1)
|
|
manual = self.manual_lines.levels()
|
|
self.position_manual_levels(manual, minute_bars)
|
|
self.levels = (
|
|
self.ma_levels
|
|
+ build_prior_day_levels(self.store.get(Timeframe.D1), self.price)
|
|
+ build_vwap_level(minute_bars)
|
|
+ manual
|
|
)
|
|
self.broadcast_level_delta()
|
|
self.rebuild_clusters()
|
|
|
|
def position_manual_levels(
|
|
self, levels: list[Level], bars: list[Bar], now: int | None = None,
|
|
) -> None:
|
|
"""Price sloped lines across bars rather than seconds.
|
|
|
|
The chart spaces bars evenly, so the line a person drew advances per bar.
|
|
Pricing it per second instead put the alert somewhere the line visibly
|
|
was not — 147 points out across a weekend.
|
|
"""
|
|
if not bars:
|
|
return
|
|
minute_times = [bar.t for bar in bars]
|
|
now = minute_times[-1] if now is None else now
|
|
for level in levels:
|
|
if not level.slope:
|
|
continue
|
|
if not self.settings.trendline_source_geometry:
|
|
level.current_p = price_in_bar_space(level, minute_times, now)
|
|
continue
|
|
source_times = [bar.t for bar in self.store.get(level.tf)]
|
|
level.current_p = price_in_timeframe_space(level, source_times, level.tf, now)
|
|
level.geometry_resolved = level.current_p is not None
|
|
|
|
def broadcast_level_delta(self) -> None:
|
|
"""Send only levels whose serialised form actually changed.
|
|
|
|
A daily moving average carries hundreds of points and changes once a
|
|
session; VWAP changes every minute. Broadcasting the whole set on the
|
|
VWAP cadence would push the entire history every minute, so subscribers
|
|
get a delta and merge it by id.
|
|
"""
|
|
current = {level.id: level.to_dict() for level in self.levels}
|
|
changed = [value for id_, value in current.items() if self._sent_levels.get(id_) != value]
|
|
removed = [id_ for id_ in self._sent_levels if id_ not in current]
|
|
self._sent_levels = current
|
|
if changed or removed:
|
|
self.broadcast({"type": "levels", "changed": changed, "removed": removed})
|
|
|
|
def rebuild_clusters(
|
|
self,
|
|
evaluate_alerts: bool = False,
|
|
alert_t: int | None = None,
|
|
alert_price: float | None = None,
|
|
) -> None:
|
|
if self.price is None or self.stream.last_bar_t is None:
|
|
return
|
|
self.clusters = cluster_levels(self.levels, self.stream.last_bar_t, self.price, self.atr15)
|
|
self.broadcast({"type": "clusters", "price": self.price, "clusters": self.clusters})
|
|
if evaluate_alerts:
|
|
# Evaluated over every level, deliberately ignoring per-connection
|
|
# layer preferences: those are a display choice made in one browser,
|
|
# and a push notification has no business depending on them.
|
|
evaluation_t = self.stream.last_bar_t if alert_t is None else alert_t
|
|
evaluation_price = self.price if alert_price is None else alert_price
|
|
alert_levels = [replace(level) for level in self.levels]
|
|
self.position_manual_levels(
|
|
alert_levels,
|
|
self.store.get(Timeframe.M1),
|
|
evaluation_t,
|
|
)
|
|
alert_clusters = cluster_levels(
|
|
alert_levels,
|
|
evaluation_t,
|
|
evaluation_price,
|
|
self.atr15,
|
|
)
|
|
self.dispatch_alerts(
|
|
self.alert_engine.evaluate(
|
|
alert_clusters,
|
|
evaluation_price,
|
|
self.atr15,
|
|
evaluation_t,
|
|
self.stream.symbol,
|
|
self.watched_ma_levels(),
|
|
confluence=self.confluence_alerts_enabled(),
|
|
)
|
|
)
|
|
|
|
def confluence_alerts_enabled(self) -> bool:
|
|
return bool(self.user_prefs.get("confluence_alerts"))
|
|
|
|
def set_confluence_alerts(self, enabled: bool) -> dict:
|
|
self.user_prefs.put("confluence_alerts", bool(enabled))
|
|
return {"enabled": self.confluence_alerts_enabled()}
|
|
|
|
def ma_alerts(self) -> dict:
|
|
stored = self.user_prefs.get("ma_alerts") or {}
|
|
return {"1d": list(stored.get("1d") or [])}
|
|
|
|
def set_ma_alerts(self, value: dict) -> dict:
|
|
self.user_prefs.put("ma_alerts", value)
|
|
return self.ma_alerts()
|
|
|
|
def watched_ma_levels(self) -> list[Level]:
|
|
armed = set(self.ma_alerts().get("1d", []))
|
|
return [
|
|
level for level in self.levels
|
|
if level.kind is LevelKind.MA
|
|
and level.tf is Timeframe.D1
|
|
and level.period in armed
|
|
]
|
|
|
|
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))
|
|
|
|
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.broadcast({
|
|
"type": "alert",
|
|
"cluster": alert.cluster,
|
|
"message": alert.message,
|
|
# Same number the push carries, so a phone and a screen agree.
|
|
"number": alert.number,
|
|
"at": alert.at,
|
|
})
|
|
task = asyncio.create_task(self.notify(alert.push or alert.message))
|
|
# Held so the task is not garbage collected mid-flight.
|
|
self._notify_tasks.add(task)
|
|
task.add_done_callback(self._notify_tasks.discard)
|
|
tripped.update(alert.tripped)
|
|
if tripped:
|
|
self.disarm(tripped)
|
|
|
|
def disarm(self, ids: set[str]) -> None:
|
|
"""A hand-placed level fires once, then waits to be re-armed."""
|
|
changed = False
|
|
for line_id in ids:
|
|
try:
|
|
self.manual_lines.update(line_id, {"armed": False})
|
|
changed = True
|
|
except KeyError:
|
|
continue
|
|
# Rebuilt after the loop, not inside it: rebuild_levels re-enters
|
|
# rebuild_clusters, and doing that mid-dispatch would rewrite the very
|
|
# clusters being iterated.
|
|
if changed:
|
|
self.rebuild_levels()
|
|
|
|
async def notify(self, message: str) -> None:
|
|
try:
|
|
await send_ntfy(self.settings.ntfy_server, self.settings.ntfy_topic, message)
|
|
except Exception:
|
|
# A push outage must not take down the stream or the sockets.
|
|
logger.warning("ntfy delivery failed", exc_info=True)
|
|
|
|
def request_rebuild(self) -> None:
|
|
"""Rebuild levels soon, at most REBUILD_INTERVAL apart.
|
|
|
|
A rebuild costs a pass over every moving average and a diff of their
|
|
points, so doing one per bar is fine at one bar a minute and ruinous in
|
|
a burst. Yahoo's first poll emits a whole day of minutes at once, which
|
|
measured as twenty seconds of loop lag. Coalescing turns that into a
|
|
handful of rebuilds without changing what subscribers eventually see:
|
|
the pending flag guarantees a trailing rebuild, so the last bar of a
|
|
burst is never the one that gets dropped.
|
|
"""
|
|
now = time.monotonic()
|
|
if now - self._last_rebuild >= self.REBUILD_INTERVAL:
|
|
self._last_rebuild = now
|
|
self._rebuild_pending = False
|
|
self.rebuild_levels()
|
|
return
|
|
self._rebuild_pending = True
|
|
|
|
async def rebuild_watch(self, interval: float = 0.25) -> None:
|
|
"""Flush a rebuild that request_rebuild deferred during a burst."""
|
|
while True:
|
|
await asyncio.sleep(interval)
|
|
if self._rebuild_pending:
|
|
self._rebuild_pending = False
|
|
self._last_rebuild = time.monotonic()
|
|
self.rebuild_levels()
|
|
|
|
def settle_after_seed(self) -> None:
|
|
"""Derive everything the replay deliberately skipped, once.
|
|
|
|
on_bar normally maintains price, ATR and the level set as each bar
|
|
arrives. During a seed it only fills the store, so the same state is
|
|
computed here from the finished history — one pass instead of one per
|
|
bar. Alerts are not evaluated: a level touched two years ago is not news,
|
|
and firing on replayed history is how a deploy used to re-alert.
|
|
"""
|
|
minute_bars = self.store.get(Timeframe.M1)
|
|
if minute_bars:
|
|
self.price = minute_bars[-1].c
|
|
values = atr(self.store.get(Timeframe.M15), 14)
|
|
self.atr15 = next((value for value in reversed(values) if value is not None), 0.0)
|
|
self.rebuild_levels()
|
|
|
|
async def loop_lag_watch(self, interval: float = 0.1) -> None:
|
|
"""Measure how late the event loop is running its own timers.
|
|
|
|
The loop is single-threaded and everything shares it: the market
|
|
stream, every WebSocket, and any CPU work that has strayed onto it.
|
|
When something blocks, the symptom reaching a person is "the chart
|
|
feels laggy" — unfalsifiable. This turns it into a number.
|
|
|
|
Scheduling drift is the honest measure: sleep for a known interval and
|
|
see how much longer it actually took.
|
|
"""
|
|
while True:
|
|
before = time.perf_counter()
|
|
await asyncio.sleep(interval)
|
|
lag = (time.perf_counter() - before) - interval
|
|
if lag > self.loop_lag_worst:
|
|
self.loop_lag_worst = lag
|
|
self.loop_lag_recent = lag
|
|
|
|
async def start(self) -> asyncio.Task:
|
|
# Captured here so a threadpool route can post events back to the loop
|
|
# that owns the queues, rather than touching them across threads.
|
|
self._loop = asyncio.get_running_loop()
|
|
self.seeding = True
|
|
try:
|
|
source = seed_source(self.settings)
|
|
# Always the Yahoo symbol: Schwab has no history to seed from.
|
|
seed_symbol = self.settings.yahoo_symbol
|
|
await self.stream.seed(
|
|
source, Timeframe.H1, self.settings.seed_1h_range, seed_symbol
|
|
)
|
|
# One-hour history cannot reconstruct a 30-minute candle. Yahoo
|
|
# supplies roughly 60 days natively; the newer overlap is replaced
|
|
# below by bars aggregated from the finer 1m seed.
|
|
await self.stream.seed(
|
|
source, Timeframe.M30, self.settings.seed_30m_range, seed_symbol
|
|
)
|
|
await self.stream.seed(
|
|
source, Timeframe.M1, self.settings.seed_1m_range, seed_symbol
|
|
)
|
|
except Exception:
|
|
# A transient seed failure must not prevent the live stream or UI starting.
|
|
pass
|
|
finally:
|
|
self.seeding = False
|
|
self.settle_after_seed()
|
|
self._lag_task = asyncio.create_task(self.loop_lag_watch(), name="loop-lag")
|
|
self._rebuild_task = asyncio.create_task(self.rebuild_watch(), name="level-rebuild")
|
|
keep_alive = getattr(self.stream.source, "keep_alive", None)
|
|
if keep_alive is not None:
|
|
self._token_task = asyncio.create_task(keep_alive(), name="token-keepalive")
|
|
return asyncio.create_task(self.stream.run(), name="market-stream")
|
|
|
|
def needs_login(self) -> bool:
|
|
error = self.stream.last_error or ""
|
|
return "invalid_grant" in error or "Refresh token" in error
|
|
|
|
def start_schwab_login(self) -> str:
|
|
from app.api.schwab_auth import begin_login
|
|
|
|
context = begin_login(self.settings)
|
|
self.schwab_login = context
|
|
return context.authorization_url
|
|
|
|
def finish_schwab_login(self, redirect_url: str) -> None:
|
|
from app.api.schwab_auth import complete_login
|
|
|
|
context = self.schwab_login
|
|
if context is None:
|
|
raise RuntimeError("No login in progress")
|
|
complete_login(self.settings, context, redirect_url)
|
|
self.schwab_login = None
|
|
reset = getattr(self.stream.source, "reset_client", None)
|
|
if reset is not None:
|
|
reset()
|