P0 from docs/async_refactor.md. The mutating routes are sync `def`, so FastAPI runs them in a threadpool, and they reach Runtime.broadcast through rebuild_levels — writing asyncio.Queue directly from there. That queue is not thread-safe: it wakes a consumer by resolving a Future, which only the loop thread may do. A dropped wakeup means a drawing made in one browser does not reach another until the next market tick. broadcast now posts through call_soon_threadsafe when it is off the loop, and publishes directly when it is on it, so the stream's own path pays nothing. Worth being straight about the tests: the race is timing-dependent and did not reproduce in twenty attempts — a foreign-thread put_nowait usually lands in the ready queue before the loop sleeps, and a tick every second covers the rest. Even asyncio's debug thread-affinity check stays quiet unless a consumer is parked on the Future at that instant. So the tests assert the contract rather than provoke the failure: a broadcast from a worker thread must go through call_soon_threadsafe, one from the loop must deliver synchronously, and both must arrive. Also adds the loop-lag probe, which reports scheduling drift as loop_lag_ms on /api/status. It found P1 on its first run: 19,441ms worst against 1.5ms in steady state, which is seeding blocking the loop. "The chart feels laggy" is now a number. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
283 lines
9.8 KiB
Python
283 lines
9.8 KiB
Python
import logging
|
|
import time
|
|
import uuid
|
|
|
|
from fastapi import APIRouter, Depends, HTTPException, Query, Request, Response
|
|
from pydantic import BaseModel, Field
|
|
|
|
from app.bars.models import Timeframe
|
|
from app.analysis.levels import Side
|
|
from app.analysis.manual_lines import ManualLine
|
|
from app.api.deps import require_token
|
|
|
|
# Everything here needs the token when CHART_AUTH_TOKEN is set. /health and
|
|
# /version live in app.api.meta and stay open on purpose.
|
|
logger = logging.getLogger(__name__)
|
|
|
|
router = APIRouter(prefix="/api", dependencies=[Depends(require_token)])
|
|
|
|
|
|
class LineCreate(BaseModel):
|
|
tf: Timeframe
|
|
side: Side
|
|
anchor_t: int
|
|
anchor_p: float
|
|
end_t: int
|
|
end_p: float
|
|
note: str = ""
|
|
hidden: bool = False
|
|
color: str = Field("#65b7cf", pattern=r"^#[0-9a-fA-F]{6}$")
|
|
line_width: int = Field(2, ge=1, le=4)
|
|
cutoff_t: int | None = None
|
|
armed: bool = True
|
|
|
|
|
|
class PriceAlertCreate(BaseModel):
|
|
"""A horizontal level typed in rather than drawn.
|
|
|
|
Structurally just a manual line with zero slope, so it inherits persistence,
|
|
editing, clustering and — importantly — the rule that a hand-placed level
|
|
alerts regardless of confluence score.
|
|
"""
|
|
|
|
price: float = Field(gt=0)
|
|
note: str = ""
|
|
tf: Timeframe = Timeframe.D1
|
|
color: str = Field("#e0a34a", pattern=r"^#[0-9a-fA-F]{6}$")
|
|
line_width: int = Field(2, ge=1, le=4)
|
|
|
|
|
|
class LinePatch(BaseModel):
|
|
side: Side | None = None
|
|
note: str | None = None
|
|
hidden: bool | None = None
|
|
color: str | None = Field(None, pattern=r"^#[0-9a-fA-F]{6}$")
|
|
line_width: int | None = Field(None, ge=1, le=4)
|
|
anchor_t: int | None = None
|
|
anchor_p: float | None = None
|
|
slope: float | None = None
|
|
last_t: int | None = None
|
|
cutoff_t: int | None = None
|
|
armed: bool | None = None
|
|
pinned: bool | None = None
|
|
x: float | None = Field(None, ge=0.0, le=1.0)
|
|
y: float | None = Field(None, ge=0.0, le=1.0)
|
|
collapsed: bool | None = None
|
|
|
|
|
|
@router.get("/status")
|
|
def status(request: Request):
|
|
runtime = request.app.state.runtime
|
|
return {
|
|
"stream": runtime.stream.status,
|
|
"source": runtime.stream.source.name,
|
|
"delay_minutes": runtime.stream.source.delay_minutes,
|
|
"symbol": runtime.stream.symbol,
|
|
"last_bar_t": runtime.stream.last_bar_t,
|
|
"bars_held": runtime.store.counts(),
|
|
"warm": {tf.value: bool(runtime.store.get(tf)) for tf in Timeframe},
|
|
# How late the event loop is running. Rising numbers mean something is
|
|
# blocking it — see docs/async_refactor.md.
|
|
"loop_lag_ms": {
|
|
"recent": round(runtime.loop_lag_recent * 1000, 1),
|
|
"worst": round(runtime.loop_lag_worst * 1000, 1),
|
|
},
|
|
}
|
|
|
|
|
|
@router.get("/bars")
|
|
def bars(request: Request, tf: str = "1m", limit: int = Query(500, ge=1, le=5000)):
|
|
try:
|
|
timeframe = Timeframe(tf)
|
|
except ValueError as exc:
|
|
raise HTTPException(400, "Unknown timeframe") from exc
|
|
values = request.app.state.runtime.store.get(timeframe, limit)
|
|
return {"tf": timeframe.value, "bars": [bar.to_dict() for bar in values]}
|
|
|
|
|
|
@router.get("/levels")
|
|
def levels(request: Request, tf: str = "all"):
|
|
values = request.app.state.runtime.levels
|
|
if tf != "all":
|
|
try:
|
|
timeframe = Timeframe(tf)
|
|
except ValueError as exc:
|
|
raise HTTPException(400, "Unknown timeframe") from exc
|
|
values = [level for level in values if level.tf is timeframe]
|
|
return {"levels": [level.to_dict() for level in values]}
|
|
|
|
|
|
@router.get("/confluence")
|
|
def confluence(request: Request):
|
|
runtime = request.app.state.runtime
|
|
return {
|
|
"price": runtime.price,
|
|
"clusters": [cluster.to_dict() for cluster in runtime.clusters],
|
|
}
|
|
|
|
|
|
@router.post("/lines", status_code=201)
|
|
def create_line(request: Request, payload: LineCreate):
|
|
if payload.end_t == payload.anchor_t:
|
|
raise HTTPException(400, "Line endpoints must have different times")
|
|
line = ManualLine(
|
|
id=f"ml_{uuid.uuid4().hex}",
|
|
tf=payload.tf,
|
|
side=payload.side,
|
|
anchor_t=payload.anchor_t,
|
|
anchor_p=payload.anchor_p,
|
|
slope=(payload.end_p - payload.anchor_p) / (payload.end_t - payload.anchor_t),
|
|
last_t=payload.end_t,
|
|
created_at=int(time.time()),
|
|
note=payload.note,
|
|
hidden=payload.hidden,
|
|
color=payload.color,
|
|
line_width=payload.line_width,
|
|
cutoff_t=payload.cutoff_t,
|
|
armed=payload.armed,
|
|
)
|
|
runtime = request.app.state.runtime
|
|
line = runtime.manual_lines.add(line)
|
|
runtime.rebuild_levels()
|
|
return line.to_level().to_dict()
|
|
|
|
|
|
@router.post("/lines/price", status_code=201)
|
|
def create_price_alert(request: Request, payload: PriceAlertCreate):
|
|
runtime = request.app.state.runtime
|
|
now = int(time.time())
|
|
# Side is only used for the label; clustering derives it positionally.
|
|
reference = runtime.price if runtime.price is not None else payload.price
|
|
line = ManualLine(
|
|
id=f"ml_{uuid.uuid4().hex}",
|
|
tf=payload.tf,
|
|
side=Side.RESISTANCE if payload.price >= reference else Side.SUPPORT,
|
|
anchor_t=now,
|
|
anchor_p=payload.price,
|
|
slope=0.0,
|
|
# A horizontal level has no natural end. The span only matters to the
|
|
# fallback geometry; price_at() is constant either way.
|
|
last_t=now + 3600,
|
|
created_at=now,
|
|
note=payload.note,
|
|
color=payload.color,
|
|
line_width=payload.line_width,
|
|
)
|
|
line = runtime.manual_lines.add(line)
|
|
runtime.rebuild_levels()
|
|
return line.to_level().to_dict()
|
|
|
|
|
|
class CommentCreate(BaseModel):
|
|
text: str = Field(min_length=1, max_length=2000)
|
|
# Pinned to a moment on the chart, or floating and always on screen.
|
|
pinned: bool = True
|
|
anchor_t: int | None = None
|
|
anchor_p: float | None = None
|
|
x: float = Field(0.72, ge=0.0, le=1.0)
|
|
y: float = Field(0.12, ge=0.0, le=1.0)
|
|
color: str = Field("#c8992f", pattern=r"^#[0-9a-fA-F]{6}$")
|
|
tf: Timeframe = Timeframe.M1
|
|
|
|
|
|
@router.post("/comments", status_code=201)
|
|
def create_comment(request: Request, payload: CommentCreate):
|
|
"""A note on the chart. Stored with the lines so it shares their numbering,
|
|
filtering and deletion, but it is never a level — see ManualLineStore.levels.
|
|
"""
|
|
runtime = request.app.state.runtime
|
|
now = int(time.time())
|
|
anchored = payload.anchor_t if payload.anchor_t is not None else now
|
|
line = ManualLine(
|
|
id=f"ml_{uuid.uuid4().hex}",
|
|
tf=payload.tf,
|
|
# Side is meaningless for a comment; clustering never sees it.
|
|
side=Side.SUPPORT,
|
|
anchor_t=anchored,
|
|
anchor_p=payload.anchor_p if payload.anchor_p is not None else (runtime.price or 0.0),
|
|
slope=0.0,
|
|
last_t=anchored,
|
|
created_at=now,
|
|
note=payload.text,
|
|
color=payload.color,
|
|
kind="comment",
|
|
pinned=payload.pinned,
|
|
x=payload.x,
|
|
y=payload.y,
|
|
# A comment must never alert, whatever else changes around it.
|
|
armed=False,
|
|
)
|
|
line = runtime.manual_lines.add(line)
|
|
return line.to_dict()
|
|
|
|
|
|
class SnapReport(BaseModel):
|
|
"""What the browser computed for one snap, for diagnosing chart geometry."""
|
|
cursor_t: int | None = None
|
|
cursor_p: float | None = None
|
|
cursor_x: float | None = None
|
|
snapped_t: int | None = None
|
|
snapped_p: float | None = None
|
|
bars_held: int | None = None
|
|
first_bar_t: int | None = None
|
|
last_bar_t: int | None = None
|
|
tf: str | None = None
|
|
chart_w: float | None = None
|
|
chart_h: float | None = None
|
|
cursor_y: float | None = None
|
|
# Where the indicator actually landed versus where the price says it should
|
|
# — the only way to tell a wrong answer from a correctly-computed one drawn
|
|
# in the wrong place.
|
|
dot_y: float | None = None
|
|
expected_y: float | None = None
|
|
bar_low_y: float | None = None
|
|
bar_high_y: float | None = None
|
|
note: str | None = None
|
|
|
|
|
|
@router.post("/debug/snap", status_code=204)
|
|
def debug_snap(payload: SnapReport):
|
|
"""Record one snap sample from a browser running diagnostic mode.
|
|
|
|
Chart geometry bugs live in the client, and the browser is usually on a
|
|
different machine from whoever is debugging it — so its own numbers cannot
|
|
be read any other way. Off unless the page is opened with ?diag=1; see
|
|
AGENTS.md. Logged at warning level so it appears without reconfiguring
|
|
uvicorn's log levels.
|
|
"""
|
|
logger.warning("SNAPDBG %s", payload.model_dump())
|
|
return Response(status_code=204)
|
|
|
|
|
|
@router.get("/drawings")
|
|
def drawings(request: Request):
|
|
"""Every drawing, comments included, for the sidebar list."""
|
|
store = request.app.state.runtime.manual_lines
|
|
return {"drawings": [
|
|
{**line.to_dict(), "kind": line.drawing_kind} for line in store.drawings()
|
|
]}
|
|
|
|
|
|
@router.patch("/lines/{line_id}")
|
|
def patch_line(request: Request, line_id: str, payload: LinePatch):
|
|
changes = payload.model_dump(exclude_none=True)
|
|
if changes.get("anchor_t") == changes.get("last_t") and "anchor_t" in changes:
|
|
raise HTTPException(400, "Line endpoints must have different times")
|
|
try:
|
|
line = request.app.state.runtime.manual_lines.update(line_id, changes)
|
|
except KeyError as exc:
|
|
raise HTTPException(404, "Line not found") from exc
|
|
request.app.state.runtime.rebuild_levels()
|
|
# A comment has no level form — returning one would hand the caller a shape
|
|
# that looks like something the confluence engine tracks.
|
|
return line.to_dict() if line.is_comment else line.to_level().to_dict()
|
|
|
|
|
|
@router.delete("/lines/{line_id}", status_code=204)
|
|
def delete_line(request: Request, line_id: str):
|
|
try:
|
|
request.app.state.runtime.manual_lines.delete(line_id)
|
|
except KeyError as exc:
|
|
raise HTTPException(404, "Line not found") from exc
|
|
request.app.state.runtime.rebuild_levels()
|
|
return Response(status_code=204)
|