Author SHA1 Message Date
jackyu66git 97aa61705d fix: pipeline MACD 参数统一为标准 12/26/9(与 web/交易所一致) 2026-09-12 02:15:17 +08:00
jackyu66gitandCursor 340676bfbd fix(web): 分型框竖边 canvas 绘制,换币对强制全量刷新
LWC 折线无法画真竖线;增量刷新时用坐标采样补刷竖边,避免与横边脱节。

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-11 17:32:18 +08:00
jackyu66gitandCursor 9cf625c413 fix(web): 小周期切换时对齐标记,避免 LWC Value is null
主周期笔/KLC 分型标记在切到 1m/2m 主图时未对齐 K 线 time;过滤均线无效点并钳制视窗恢复。顺带统一 BI 中枢计算路径。

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-08 17:06:31 +08:00
jackyu66gitandCursor 18a7f485e6 feat(web): 增量自动刷新、结构区修复与默认指标/周期
自动刷新常态只拉 recent 尾部 K,每 1 分钟全量重算缠论;修复结构区缓存导入;默认指标/4h·1h·15m/近30天;同步 ECR-009 screener 相关改动。

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-08 15:45:40 +08:00
jackyu66gitandCursor 0f6eb92a1f test(ECR-009): 补页面/API 路由冒烟与 TEST_REPORT
交付前缺 Flask 常驻与路由断言;现补齐 pytest 与报告。

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-07 16:03:37 +08:00
jackyu66gitandCursor 9880e236a5 docs(ECR-009): record implementation commit in TRACEABILITY
Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-07 15:46:35 +08:00
jackyu66gitandCursor ec08de098e feat(ECR-009): Crypto Wyckoff Screener 独立页(D/W/M)
移植 A_Share_DP 引擎;本地缓存与 60s tip;月线由日线 UTC 聚合;不碰主站 analyze/缠论叠层。

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-07 15:46:35 +08:00
54 changed files with 5521 additions and 190 deletions
+3
View File
@@ -44,3 +44,6 @@ data_provider/._config.json
# ESS gate / engineering-loop working dirs(归档进 docs/runs/ # ESS gate / engineering-loop working dirs(归档进 docs/runs/
.gates/ .gates/
loop/ loop/
# Crypto Wyckoff Screener local cache
data/crypto_wyckoff/
+2 -2
View File
@@ -55,8 +55,8 @@ class IndicatorsBuilderMixin:
return None return None
def add_indicators(self, df): def add_indicators(self, df):
fast = 26 fast = 12
slow = 52 slow = 26
period = 9 period = 9
macd = ta.MACD(df, fastperiod=fast, slowperiod=slow, signalperiod=period) macd = ta.MACD(df, fastperiod=fast, slowperiod=slow, signalperiod=period)
bb365 = ta.BBANDS(df, timeperiod=365, nbdevup=3.0, nbdevdn=3.0, matype=0) bb365 = ta.BBANDS(df, timeperiod=365, nbdevup=3.0, nbdevdn=3.0, matype=0)
+5
View File
@@ -0,0 +1,5 @@
"""crypto_wyckoff — multi-TF screener for crypto (ported from A_Share_DP Architecture v1.0)."""
from crypto_wyckoff.version import ARCHITECTURE_VERSION, WYCKOFF_ENGINE_VERSION
__all__ = ["WYCKOFF_ENGINE_VERSION", "ARCHITECTURE_VERSION"]
+342
View File
@@ -0,0 +1,342 @@
"""Walk-forward Wyckoff phase/event annotations for chart overlay."""
from __future__ import annotations
from datetime import date
from crypto_wyckoff.domain_models import OHLCVFrame, WyckoffCycle, WyckoffEvent, WyckoffPhase
from crypto_wyckoff.cycle import CycleEngine
from crypto_wyckoff.event import EventEngine
from crypto_wyckoff.features import FeatureEngine
from crypto_wyckoff.phase import PhaseEngine
_MIN_BARS = {"1d": 40, "1w": 26, "1M": 18}
_NOTABLE_EVENTS = {
WyckoffEvent.PS.value,
WyckoffEvent.SC.value,
WyckoffEvent.AR.value,
WyckoffEvent.ST.value,
WyckoffEvent.SPRING.value,
WyckoffEvent.TEST.value,
WyckoffEvent.SOS.value,
WyckoffEvent.LPS.value,
WyckoffEvent.JUMP.value,
WyckoffEvent.BACKUP.value,
WyckoffEvent.BC.value,
WyckoffEvent.UTAD.value,
WyckoffEvent.SOW.value,
WyckoffEvent.LPSY.value,
}
def _slice_frame(frame: OHLCVFrame, end_idx: int) -> OHLCVFrame:
n = end_idx + 1
return OHLCVFrame(
ts_code=frame.ts_code,
timeframe=frame.timeframe,
trade_dates=frame.trade_dates[:n],
open=frame.open[:n],
high=frame.high[:n],
low=frame.low[:n],
close=frame.close[:n],
volume=frame.volume[:n],
amount=frame.amount[:n] if frame.amount else [],
)
def _compress_phases(points: list[tuple[str, str]]) -> list[dict]:
"""points: [(date_iso, phase), ...] → segments."""
if not points:
return []
segs: list[dict] = []
start, phase = points[0]
prev = start
for d, p in points[1:]:
if p != phase:
segs.append({"start": start, "end": prev, "phase": phase})
start, phase = d, p
prev = d
segs.append({"start": start, "end": prev, "phase": phase})
return segs
def annotate_frame(
frame: OHLCVFrame,
step: int | None = None,
*,
role: str | None = None,
) -> dict:
"""Pure annotation: phase bands + event markers + latest levels.
``role`` is the D/W/M rule alias (1d/1w/1M). Defaults to frame.timeframe.
``step`` defaults by role to keep interactive charts snappy.
"""
tf = role or frame.timeframe
min_bars = _MIN_BARS.get(tf, 30)
if step is None:
step = {"1d": 2, "1w": 1, "1M": 1}.get(tf, 2)
empty = {
"phases": [],
"events": [],
"levels": {},
"bars": len(frame),
"timeframe": tf,
}
if frame.empty or len(frame) < min_bars:
return empty
feat_eng = FeatureEngine()
cycle_eng = CycleEngine()
phase_eng = PhaseEngine()
event_eng = EventEngine()
phase_points: list[tuple[str, str]] = []
events: list[dict] = []
last_event: str | None = None
levels: dict = {}
# Ensure last bar is always evaluated
indices = list(range(min_bars - 1, len(frame), step))
if indices[-1] != len(frame) - 1:
indices.append(len(frame) - 1)
for i in indices:
sub = _slice_frame(frame, i)
f = feat_eng.run(sub, tf)
c = cycle_eng.run(f, tf)
p = phase_eng.run(c, f, tf)
e = event_eng.run(c, p, f, tf)
d = str(frame.trade_dates[i])[:10]
phase = p.payload.get("phase") or WyckoffPhase.NONE.value
phase_points.append((d, phase))
cur = e.payload.get("current_event") or WyckoffEvent.NONE.value
if cur in _NOTABLE_EVENTS and cur != last_event:
events.append({
"date": d,
"event": cur,
"price": float(frame.close[i]),
"low": float(frame.low[i]),
"high": float(frame.high[i]),
})
last_event = cur
elif cur == WyckoffEvent.NONE.value:
last_event = None
if i == len(frame) - 1 and not f.payload.get("insufficient"):
levels = {
k: f.payload.get(k)
for k in (
"range_high", "range_low", "ma20", "ma60",
"swing_high", "swing_low", "close",
)
if f.payload.get(k) is not None
}
levels["phase"] = phase
levels["cycle"] = c.payload.get("cycle")
levels["current_event"] = cur
return {
"phases": _compress_phases(phase_points),
"events": events,
"levels": levels,
"bars": len(frame),
"timeframe": tf,
}
_RANGE_CYCLES = {
WyckoffCycle.ACCUMULATION.value,
WyckoffCycle.RE_ACCUMULATION.value,
WyckoffCycle.DISTRIBUTION.value,
WyckoffCycle.RE_DISTRIBUTION.value,
}
def _build_range_zones(
price_frame: OHLCVFrame,
cycle_segs: list[dict],
levels: dict | None = None,
) -> list[dict]:
"""Build price boxes (high/low × date span) for accum/distrib ranges."""
if price_frame.empty:
return []
dates = [str(d)[:10] for d in price_frame.trade_dates]
highs = price_frame.high
lows = price_frame.low
zones: list[dict] = []
for seg in cycle_segs or []:
cy = seg.get("cycle")
if cy not in _RANGE_CYCLES:
continue
start, end = seg["start"], seg["end"]
idxs = [i for i, d in enumerate(dates) if start <= d <= end]
if not idxs:
# weekly bar date may sit between daily bars — take nearest window
i0 = next((i for i, d in enumerate(dates) if d >= start), None)
if i0 is None:
continue
i1 = next((i for i, d in enumerate(dates) if d > end), len(dates)) - 1
idxs = list(range(i0, max(i0, i1) + 1))
if not idxs:
continue
# pad short weekly hits to at least ~1 week of dailies for visibility
if len(idxs) < 5 and idxs[-1] + 1 < len(dates):
extra = min(5 - len(idxs), len(dates) - 1 - idxs[-1])
idxs = list(range(idxs[0], idxs[-1] + 1 + max(0, extra)))
hi = max(highs[i] for i in idxs)
lo = min(lows[i] for i in idxs)
if hi <= lo:
continue
zones.append({
"kind": cy,
"start": dates[idxs[0]],
"end": dates[idxs[-1]],
"high": float(hi),
"low": float(lo),
"current": False,
})
# Always expose the latest trading-range box from feature snapshot
levels = levels or {}
rh, rl = levels.get("range_high"), levels.get("range_low")
if rh is not None and rl is not None and float(rh) > float(rl):
look = min(60, len(dates))
cy = levels.get("cycle") or "Unknown"
if cy not in _RANGE_CYCLES:
# Phase B/C in a range → treat as accumulation-style TR for display
ph = levels.get("phase") or ""
if ph in ("A", "B", "C"):
cy = WyckoffCycle.ACCUMULATION.value
elif ph in ("D", "E") and float(levels.get("close") or 0) < float(rh):
cy = WyckoffCycle.ACCUMULATION.value
else:
cy = "Range"
zones.append({
"kind": cy,
"start": dates[-look],
"end": dates[-1],
"high": float(rh),
"low": float(rl),
"current": True,
})
return zones
def annotate_symbol(
ts_code: str,
freq: str,
end_date: date | None = None,
lookback: int = 180,
*,
combo_id: str | None = None,
) -> dict:
"""IO + annotate for one symbol (used by API).
For the combo *low* chart, phase bands come from **mid** structure,
while event markers / levels come from the low TF.
"""
from crypto_wyckoff.combos import ROLE_HIGH, ROLE_LOW, ROLE_MID, get_combo
from crypto_wyckoff.io import load_frame
combo = get_combo(combo_id)
allowed = {combo["low"], combo["mid"], combo["high"]}
if freq not in allowed:
raise ValueError(f"freq {freq} not in combo {combo['id']} ({combo['label']})")
empty = {
"ts_code": ts_code,
"freq": freq,
"phases": [],
"events": [],
"levels": {},
"zones": [],
"bars": 0,
"phase_source": freq,
"cycles": [],
"combo_id": combo["id"],
}
_ = end_date
if freq == combo["low"]:
low = load_frame(ts_code, combo["low"], lookback)
mid = load_frame(ts_code, combo["mid"], max(60, lookback // 3))
if low is None:
return empty
d_ann = annotate_frame(low, role=ROLE_LOW)
w_ann = annotate_frame(mid, role=ROLE_MID) if mid is not None else {"phases": []}
cycles = _cycle_segments(mid, role=ROLE_MID) if mid is not None else []
levels = d_ann.get("levels") or {}
if cycles:
levels = {**levels, "cycle": cycles[-1].get("cycle") or levels.get("cycle")}
for p in reversed(w_ann.get("phases") or []):
if p.get("phase") not in (None, "None"):
levels = {**levels, "phase": p["phase"]}
break
return {
"ts_code": ts_code,
"freq": freq,
"end_date": low.trade_dates[-1].isoformat() if low.trade_dates else None,
"phases": w_ann.get("phases") or [],
"events": d_ann.get("events") or [],
"levels": d_ann.get("levels") or {},
"zones": _build_range_zones(low, cycles, levels),
"bars": d_ann.get("bars", 0),
"phase_source": combo["mid"],
"cycles": cycles,
"combo_id": combo["id"],
}
role = ROLE_MID if freq == combo["mid"] else ROLE_HIGH
frame = load_frame(ts_code, freq, lookback)
if frame is None:
return empty
out = annotate_frame(frame, role=role)
out["ts_code"] = ts_code
out["freq"] = freq
out["end_date"] = frame.trade_dates[-1].isoformat() if frame.trade_dates else None
out["phase_source"] = freq
out["cycles"] = _cycle_segments(frame, role=ROLE_HIGH if role == ROLE_HIGH else ROLE_MID)
out["zones"] = _build_range_zones(frame, out["cycles"], out.get("levels") or {})
out["combo_id"] = combo["id"]
if role == ROLE_HIGH:
if not any(p.get("phase") not in (None, "None") for p in out["phases"]):
out["phases"] = [
{"start": c["start"], "end": c["end"], "phase": c["cycle"]}
for c in out["cycles"]
if c.get("cycle") and c["cycle"] != "Unknown"
]
return out
def _cycle_segments(
frame: OHLCVFrame,
step: int | None = None,
*,
role: str | None = None,
) -> list[dict]:
"""Walk-forward cycle labels compressed to segments."""
tf = role or frame.timeframe
min_bars = _MIN_BARS.get(tf, 30)
if step is None:
step = {"1d": 3, "1w": 1, "1M": 1}.get(tf, 2)
if frame.empty or len(frame) < min_bars:
return []
feat_eng = FeatureEngine()
cycle_eng = CycleEngine()
points: list[tuple[str, str]] = []
indices = list(range(min_bars - 1, len(frame), step))
if indices[-1] != len(frame) - 1:
indices.append(len(frame) - 1)
for i in indices:
sub = _slice_frame(frame, i)
f = feat_eng.run(sub, tf)
c = cycle_eng.run(f, tf)
points.append((str(frame.trade_dates[i])[:10], c.payload.get("cycle") or "Unknown"))
segs = _compress_phases(points)
return [{"start": s["start"], "end": s["end"], "cycle": s["phase"]} for s in segs]
+248
View File
@@ -0,0 +1,248 @@
"""Multi-timeframe combo presets for Crypto Wyckoff Screener.
Roles (engine rule aliases stay D/W/M):
high → Cycle (rules as 1M)
mid → Phase (rules as 1w)
low → Event (rules as 1d)
Actual bar TFs come from the combo (e.g. 8h/4h/1h).
"""
from __future__ import annotations
import json
import re
import threading
from copy import deepcopy
from pathlib import Path
from typing import Any
from crypto_wyckoff.io import DATA_DIR, ensure_dirs
ROLE_LOW = "1d"
ROLE_MID = "1w"
ROLE_HIGH = "1M"
# Minutes for ordering / validation (provider labels)
_TF_MINUTES: dict[str, int] = {
"1m": 1, "2m": 2, "3m": 3, "4m": 4, "5m": 5,
"10m": 10, "15m": 15, "20m": 20, "25m": 25, "30m": 30, "45m": 45,
"1h": 60, "2h": 120, "3h": 180, "4h": 240, "5h": 300,
"6h": 360, "7h": 420, "8h": 480, "9h": 540, "10h": 600,
"11h": 660, "12h": 720, "16h": 960, "20h": 1200,
"1d": 1440, "2d": 2880, "3d": 4320, "4d": 5760, "5d": 7200, "6d": 8640,
"1w": 10080, "2w": 20160, "3w": 30240,
"1M": 43200,
}
# TFs we allow in custom combos (provider-backed + local 1M)
ALLOWED_TFS: tuple[str, ...] = (
"1h", "2h", "3h", "4h", "6h", "8h", "12h",
"1d", "2d", "3d", "1w", "1M",
)
BUILTIN: list[dict[str, Any]] = [
{
"id": "h8_4_1",
"label": "8h / 4h / 1h",
"high": "8h",
"mid": "4h",
"low": "1h",
"builtin": True,
},
{
"id": "d_w_m",
"label": "1d / 1w / 1M",
"high": "1M",
"mid": "1w",
"low": "1d",
"builtin": True,
},
]
_COMBOS_FILE = DATA_DIR / "combos.json"
_lock = threading.Lock()
_cache: list[dict[str, Any]] | None = None
def tf_minutes(tf: str) -> int | None:
if tf in _TF_MINUTES:
return _TF_MINUTES[tf]
# tolerate provider typo "10" → skip
m = re.fullmatch(r"(\d+)([mhdwM])", tf)
if not m:
return None
n, u = int(m.group(1)), m.group(2)
mult = {"m": 1, "h": 60, "d": 1440, "w": 10080, "M": 43200}[u]
return n * mult
def combo_id_for(high: str, mid: str, low: str) -> str:
def _tok(t: str) -> str:
return t.replace("/", "_")
return f"{_tok(high)}_{_tok(mid)}_{_tok(low)}"
def validate_combo(high: str, mid: str, low: str) -> str | None:
"""Return error message or None if ok."""
for tf in (high, mid, low):
if tf not in ALLOWED_TFS:
return f"不支持的周期: {tf}"
if len({high, mid, low}) < 3:
return "高/中/低周期必须互不相同"
hm, mm, lm = tf_minutes(high), tf_minutes(mid), tf_minutes(low)
if hm is None or mm is None or lm is None:
return "无法解析周期长度"
if not (hm > mm > lm):
return "须满足 高 > 中 > 低(例如 8h > 4h > 1h"
return None
def _normalize(row: dict[str, Any]) -> dict[str, Any] | None:
high, mid, low = row.get("high"), row.get("mid"), row.get("low")
if not high or not mid or not low:
return None
err = validate_combo(str(high), str(mid), str(low))
if err:
return None
cid = str(row.get("id") or combo_id_for(high, mid, low))
label = str(row.get("label") or f"{high} / {mid} / {low}")
return {
"id": cid,
"label": label,
"high": str(high),
"mid": str(mid),
"low": str(low),
"builtin": bool(row.get("builtin", False)),
}
def _load_raw() -> list[dict[str, Any]]:
ensure_dirs()
if not _COMBOS_FILE.exists():
return deepcopy(BUILTIN)
try:
data = json.loads(_COMBOS_FILE.read_text(encoding="utf-8"))
items = data.get("combos") if isinstance(data, dict) else data
if not isinstance(items, list):
return deepcopy(BUILTIN)
except (OSError, json.JSONDecodeError):
return deepcopy(BUILTIN)
out: list[dict[str, Any]] = []
seen: set[str] = set()
for b in BUILTIN:
out.append(deepcopy(b))
seen.add(b["id"])
for row in items:
if not isinstance(row, dict):
continue
norm = _normalize(row)
if not norm or norm["id"] in seen:
continue
if norm["id"] in {b["id"] for b in BUILTIN}:
continue
norm["builtin"] = False
out.append(norm)
seen.add(norm["id"])
return out
def _save(combos: list[dict[str, Any]]) -> None:
ensure_dirs()
custom = [c for c in combos if not c.get("builtin")]
payload = {"combos": custom}
tmp = _COMBOS_FILE.with_suffix(".tmp")
tmp.write_text(json.dumps(payload, ensure_ascii=False, indent=2), encoding="utf-8")
tmp.replace(_COMBOS_FILE)
def list_combos() -> list[dict[str, Any]]:
global _cache
with _lock:
if _cache is None:
_cache = _load_raw()
return deepcopy(_cache)
def get_combo(combo_id: str | None) -> dict[str, Any]:
combos = list_combos()
if combo_id:
for c in combos:
if c["id"] == combo_id:
return deepcopy(c)
return deepcopy(combos[0])
def add_combo(high: str, mid: str, low: str, label: str | None = None) -> dict[str, Any]:
err = validate_combo(high, mid, low)
if err:
raise ValueError(err)
cid = combo_id_for(high, mid, low)
row = {
"id": cid,
"label": label or f"{high} / {mid} / {low}",
"high": high,
"mid": mid,
"low": low,
"builtin": False,
}
with _lock:
combos = _load_raw()
for c in combos:
if c["id"] == cid or (c["high"], c["mid"], c["low"]) == (high, mid, low):
_cache = combos
return deepcopy(c)
combos.append(row)
_save(combos)
_cache = combos
return deepcopy(row)
def delete_combo(combo_id: str) -> bool:
with _lock:
combos = _load_raw()
kept: list[dict[str, Any]] = []
removed = False
for c in combos:
if c["id"] == combo_id:
if c.get("builtin"):
raise ValueError("内置组合不可删除")
removed = True
continue
kept.append(c)
if removed:
_save(kept)
_cache = kept
return removed
def all_tfs_for_combos(combos: list[dict[str, Any]] | None = None) -> list[str]:
"""Unique TFs needed by active combos (stable order)."""
rows = combos if combos is not None else list_combos()
seen: list[str] = []
for c in rows:
for k in ("low", "mid", "high"):
tf = c[k]
if tf not in seen:
seen.append(tf)
return seen
def lookback_for(tf: str) -> int:
defaults = {
"1h": 500,
"2h": 400,
"3h": 350,
"4h": 300,
"6h": 280,
"8h": 250,
"12h": 220,
"1d": 250,
"2d": 200,
"3d": 180,
"1w": 104,
"1M": 60,
}
return defaults.get(tf, 200)
+102
View File
@@ -0,0 +1,102 @@
"""Cycle Engine — monthly/weekly macro cycle via Rule Registry."""
from __future__ import annotations
from crypto_wyckoff.domain_models import EngineResult, WyckoffCycle
from crypto_wyckoff.rules.base import RuleHit
from crypto_wyckoff.rules.registry import rule_registry
def _resolve_range_conflict(hits: list[RuleHit], features: dict) -> list[RuleHit]:
"""Accumulation vs Distribution overlap → mutually exclusive by MA120 position."""
accum = [h for h in hits if h.cycle == WyckoffCycle.ACCUMULATION.value]
dist = [h for h in hits if h.cycle == WyckoffCycle.DISTRIBUTION.value]
if not (accum and dist):
return hits
close = float(features.get("close") or 0)
ma120 = float(features.get("ma120") or close) or close
others = [
h for h in hits
if h.cycle not in (WyckoffCycle.ACCUMULATION.value, WyckoffCycle.DISTRIBUTION.value)
]
# Below MA120 → accumulation; above → distribution; equal band uses relative position
if close < ma120 * 0.995:
return others + accum
if close > ma120 * 1.005:
return others + dist
# Tight band: keep higher confidence only
best_a = max(accum, key=lambda h: h.confidence)
best_d = max(dist, key=lambda h: h.confidence)
return others + ([best_a] if best_a.confidence >= best_d.confidence else [best_d])
class CycleEngine:
name = "Cycle"
version = "1.0.0"
def run(self, feature: EngineResult, timeframe: str) -> EngineResult:
features = feature.payload
if features.get("insufficient"):
return EngineResult(
name=self.name,
version=self.version,
confidence=15.0,
score=40.0,
reasons=[f"{timeframe} 数据不足,Cycle=Unknown"],
warnings=["insufficient_features"],
payload={
"cycle": WyckoffCycle.UNKNOWN.value,
"timeframe": timeframe,
"trend_score": 40.0,
},
)
context = {"features": features, "timeframe": timeframe}
hits: list[RuleHit] = []
for rule in rule_registry.by_category("cycle", timeframe):
hit = rule.evaluate(context)
if hit and hit.cycle:
hits.append(hit)
hits = _resolve_range_conflict(hits, features)
if not hits:
return EngineResult(
name=self.name,
version=self.version,
confidence=30.0,
score=40.0,
reasons=["无匹配周期规则,标记 Unknown"],
payload={
"cycle": WyckoffCycle.UNKNOWN.value,
"timeframe": timeframe,
"trend_score": 40.0,
},
)
best = max(hits, key=lambda h: h.confidence)
trend_score = best.score
if best.cycle == WyckoffCycle.MARKUP.value:
trend_score = max(trend_score, 75.0)
elif best.cycle == WyckoffCycle.ACCUMULATION.value:
trend_score = max(60.0, trend_score * 0.9)
elif best.cycle == WyckoffCycle.DISTRIBUTION.value:
trend_score = min(45.0, 100 - trend_score * 0.5)
elif best.cycle == WyckoffCycle.MARKDOWN.value:
trend_score = min(30.0, 100 - trend_score)
return EngineResult(
name=self.name,
version=self.version,
confidence=best.confidence,
score=trend_score,
reasons=best.reasons,
metrics=best.metrics,
payload={
"cycle": best.cycle,
"timeframe": timeframe,
"rule_id": best.rule_id,
"trend_score": trend_score,
},
)
+195
View File
@@ -0,0 +1,195 @@
"""Decision Engine — multi-timeframe fusion and tradability (Architecture v1.0)."""
from __future__ import annotations
from crypto_wyckoff.domain_models import (
DecisionSignal,
EngineResult,
RiskLevel,
WyckoffCycle,
WyckoffEvent,
WyckoffPhase,
)
BULL_CYCLES = {
WyckoffCycle.ACCUMULATION.value,
WyckoffCycle.RE_ACCUMULATION.value,
WyckoffCycle.MARKUP.value,
}
BEAR_CYCLES = {
WyckoffCycle.DISTRIBUTION.value,
WyckoffCycle.RE_DISTRIBUTION.value,
WyckoffCycle.MARKDOWN.value,
}
class DecisionEngine:
name = "Decision"
version = "1.0.0"
def run(
self,
monthly_cycle: EngineResult,
weekly_cycle: EngineResult,
weekly_phase: EngineResult,
weekly_event: EngineResult,
daily_event: EngineResult,
daily_signal: EngineResult,
) -> EngineResult:
m_cycle = monthly_cycle.payload.get("cycle", WyckoffCycle.UNKNOWN.value)
w_cycle = weekly_cycle.payload.get("cycle", WyckoffCycle.UNKNOWN.value)
w_phase = weekly_phase.payload.get("phase", WyckoffPhase.NONE.value)
w_event = weekly_event.payload.get("current_event", WyckoffEvent.NONE.value)
d_event = daily_event.payload.get("current_event", WyckoffEvent.NONE.value)
trend_score = float(monthly_cycle.payload.get("trend_score", monthly_cycle.score))
structure_score = float(weekly_phase.payload.get("structure_score", weekly_phase.score))
entry_score = float(daily_event.payload.get("entry_score", daily_event.score))
overall_score = 0.30 * trend_score + 0.30 * structure_score + 0.40 * entry_score
reasons: list[str] = []
warnings: list[str] = []
alignment = 50.0
m_bull = m_cycle in BULL_CYCLES
m_bear = m_cycle in BEAR_CYCLES
w_bull = w_cycle in BULL_CYCLES
d_bullish_event = d_event in {
WyckoffEvent.SPRING.value,
WyckoffEvent.TEST.value,
WyckoffEvent.SOS.value,
WyckoffEvent.LPS.value,
WyckoffEvent.JUMP.value,
WyckoffEvent.BACKUP.value,
}
d_bearish_event = d_event in {
WyckoffEvent.UTAD.value,
WyckoffEvent.SOW.value,
WyckoffEvent.LPSY.value,
}
# Alignment scoring
if m_bull and w_bull and d_bullish_event:
alignment = 92.0
reasons.append("✓ 月/周多头结构与日线多头事件一致")
elif m_bull and d_bullish_event:
alignment = 78.0
reasons.append("✓ 月线支持,日线有入场事件")
if not w_bull:
warnings.append("周线结构未完全确认")
alignment -= 8
elif m_bear and d_bullish_event:
alignment = 35.0
reasons.append("✗ 月线派发/下跌,日线弹簧可能只是反弹")
elif m_bear and d_bearish_event:
alignment = 85.0
reasons.append("✓ 空头多周期一致")
else:
alignment = 55.0
reasons.append("○ 多周期部分一致,需观察")
if w_phase in (WyckoffPhase.D.value, WyckoffPhase.E.value) and m_bull:
alignment = min(98.0, alignment + 6)
reasons.append(f"✓ 周线阶段 {w_phase} 结构成熟({w_event}")
active = daily_event.payload.get("active_events") or daily_event.payload.get("recent_events") or []
if d_event == WyckoffEvent.SPRING.value and len(active) >= 3:
alignment = min(98.0, alignment + 4)
reasons.append("✓ 日线多重事件同时确认")
# Decision signal — hard gate on monthly bear + daily spring
decision = DecisionSignal.WATCH.value
risk = RiskLevel.MEDIUM.value
if m_bear and d_event == WyckoffEvent.SPRING.value:
decision = DecisionSignal.WATCH.value
risk = RiskLevel.HIGH.value
overall_score = min(overall_score, 55.0)
reasons.append("→ 决策:观察(月线不支持,禁止追日线弹簧)")
elif m_bear and d_bullish_event:
decision = DecisionSignal.AVOID.value
risk = RiskLevel.HIGH.value
overall_score = min(overall_score, 48.0)
reasons.append("→ 决策:回避(逆大周期多头事件)")
elif (
m_bull
and w_phase in (WyckoffPhase.D.value, WyckoffPhase.E.value, WyckoffPhase.C.value)
and d_event in (WyckoffEvent.SPRING.value, WyckoffEvent.LPS.value, WyckoffEvent.SOS.value)
and alignment >= 85
and overall_score >= 80
):
decision = DecisionSignal.STRONG_BUY.value
risk = RiskLevel.LOW.value
reasons.append("→ 决策:强烈买入(三级共振)")
elif m_bull and d_bullish_event and overall_score >= 68 and alignment >= 70:
decision = DecisionSignal.BUY.value
risk = RiskLevel.LOW.value if alignment >= 80 else RiskLevel.MEDIUM.value
reasons.append("→ 决策:买入")
elif m_bear and d_bearish_event and overall_score >= 65:
decision = DecisionSignal.SELL.value
risk = RiskLevel.MEDIUM.value
reasons.append("→ 决策:卖出")
else:
decision = DecisionSignal.WATCH.value
reasons.append("→ 决策:观察")
# Stars from score + alignment
combo = 0.6 * overall_score + 0.4 * alignment
if combo >= 90:
stars = 5
elif combo >= 80:
stars = 4
elif combo >= 65:
stars = 3
elif combo >= 50:
stars = 2
else:
stars = 1
overall_confidence = (
0.25 * monthly_cycle.confidence
+ 0.25 * weekly_phase.confidence
+ 0.25 * daily_event.confidence
+ 0.25 * daily_signal.confidence
)
# Weak event pulls overall down
if daily_event.confidence < 60:
overall_confidence = min(overall_confidence, daily_event.confidence + 15)
return EngineResult(
name=self.name,
version=self.version,
confidence=overall_confidence,
score=overall_score,
reasons=reasons,
warnings=warnings,
metrics={
"trend_score": trend_score,
"structure_score": structure_score,
"entry_score": entry_score,
"alignment": alignment,
"stars": stars,
},
payload={
"decision_signal": decision,
"alignment": alignment,
"stars": stars,
"risk": risk,
"overall_score": overall_score,
"overall_confidence": overall_confidence,
"trend_score": trend_score,
"structure_score": structure_score,
"entry_score": entry_score,
"m_cycle": m_cycle,
"w_cycle": w_cycle,
"w_phase": w_phase,
"w_event": w_event,
"d_event": d_event,
# Facts preserved — never overwritten
"facts": {
"monthly": {"cycle": m_cycle},
"weekly": {"cycle": w_cycle, "phase": w_phase, "event": w_event},
"daily": {"event": d_event},
},
},
)
+154
View File
@@ -0,0 +1,154 @@
"""Wyckoff Screener domain models — Architecture v1.0 frozen contracts."""
from __future__ import annotations
from dataclasses import dataclass, field
from datetime import date, datetime
from enum import Enum
from typing import Any, Optional
class WyckoffCycle(str, Enum):
ACCUMULATION = "Accumulation"
RE_ACCUMULATION = "ReAccumulation"
MARKUP = "Markup"
DISTRIBUTION = "Distribution"
RE_DISTRIBUTION = "ReDistribution"
MARKDOWN = "Markdown"
UNKNOWN = "Unknown"
class WyckoffPhase(str, Enum):
A = "A"
B = "B"
C = "C"
D = "D"
E = "E"
NONE = "None"
class WyckoffEvent(str, Enum):
PS = "PS"
SC = "SC"
AR = "AR"
ST = "ST"
SPRING = "Spring"
TEST = "Test"
SOS = "SOS"
LPS = "LPS"
JUMP = "Jump"
BACKUP = "Backup"
BC = "BC"
UTAD = "UTAD"
SOW = "SOW"
LPSY = "LPSY"
NONE = "None"
class DecisionSignal(str, Enum):
STRONG_BUY = "StrongBuy"
BUY = "Buy"
WATCH = "Watch"
AVOID = "Avoid"
SELL = "Sell"
class RiskLevel(str, Enum):
LOW = "Low"
MEDIUM = "Medium"
HIGH = "High"
@dataclass
class EngineResult:
"""Unified result envelope for every Wyckoff engine (v1.0 contract)."""
name: str
version: str = "1.0.0"
confidence: float = 0.0
score: float = 0.0
reasons: list[str] = field(default_factory=list)
warnings: list[str] = field(default_factory=list)
metrics: dict[str, Any] = field(default_factory=dict)
payload: dict[str, Any] = field(default_factory=dict)
def to_dict(self) -> dict[str, Any]:
return {
"name": self.name,
"version": self.version,
"confidence": self.confidence,
"score": self.score,
"reasons": self.reasons,
"warnings": self.warnings,
"metrics": self.metrics,
"payload": self.payload,
}
@dataclass
class OHLCVFrame:
"""In-memory OHLCV for one symbol one timeframe. Engines never touch DB."""
ts_code: str
timeframe: str # "1d" | "1w" | "1M"
trade_dates: list[date]
open: list[float]
high: list[float]
low: list[float]
close: list[float]
volume: list[float]
amount: list[float] = field(default_factory=list)
def __len__(self) -> int:
return len(self.close)
@property
def empty(self) -> bool:
return len(self.close) == 0
@dataclass
class WyckoffScanRow:
"""Persisted scan row for wyckoff_scan table."""
trade_date: date
ts_code: str
name: str = ""
industry: str = ""
engine_version: str = "v1.0.0"
combo_id: str = "d_w_m"
m_cycle: str = WyckoffCycle.UNKNOWN.value
cycle_confidence: float = 0.0
trend_score: float = 0.0
w_cycle: str = WyckoffCycle.UNKNOWN.value
w_phase: str = WyckoffPhase.NONE.value
w_current_event: str = WyckoffEvent.NONE.value
w_recent_events_json: str = "[]"
phase_confidence: float = 0.0
structure_score: float = 0.0
d_current_event: str = WyckoffEvent.NONE.value
d_recent_events_json: str = "[]"
event_confidence: float = 0.0
entry_score: float = 0.0
entry: Optional[float] = None
stop: Optional[float] = None
target1: Optional[float] = None
target2: Optional[float] = None
rr: Optional[float] = None
alignment: float = 0.0
stars: int = 1
decision_signal: str = DecisionSignal.WATCH.value
signal_confidence: float = 0.0
overall_confidence: float = 0.0
overall_score: float = 0.0
risk: str = RiskLevel.MEDIUM.value
reasons_json: str = "[]"
feature_snapshot_json: str = "{}"
markers_json: str = "[]"
scanned_at: datetime = field(default_factory=datetime.now)
+149
View File
@@ -0,0 +1,149 @@
"""Event Engine — active concurrent events via Rule Registry.
Note: `active_events` are rules that fire on the latest bar snapshot,
NOT a historical SC→AR→ST timeline. Do not present as chronological chain.
"""
from __future__ import annotations
from crypto_wyckoff.domain_models import EngineResult, WyckoffEvent
from crypto_wyckoff.rules.registry import rule_registry
# Display order only (not temporal history)
_DISPLAY_ORDER = [
WyckoffEvent.PS.value,
WyckoffEvent.SC.value,
WyckoffEvent.AR.value,
WyckoffEvent.ST.value,
WyckoffEvent.SPRING.value,
WyckoffEvent.TEST.value,
WyckoffEvent.SOS.value,
WyckoffEvent.LPS.value,
WyckoffEvent.JUMP.value,
WyckoffEvent.BACKUP.value,
WyckoffEvent.BC.value,
WyckoffEvent.UTAD.value,
WyckoffEvent.SOW.value,
WyckoffEvent.LPSY.value,
]
# Dominant event: highest confidence wins; ties broken by this priority
_DOMINANCE_PRIORITY = [
WyckoffEvent.SOS.value,
WyckoffEvent.LPS.value,
WyckoffEvent.UTAD.value,
WyckoffEvent.SPRING.value,
WyckoffEvent.JUMP.value,
WyckoffEvent.BACKUP.value,
WyckoffEvent.TEST.value,
WyckoffEvent.SC.value,
WyckoffEvent.SOW.value,
WyckoffEvent.AR.value,
WyckoffEvent.ST.value,
]
class EventEngine:
name = "Event"
version = "1.0.0"
def run(
self,
cycle: EngineResult,
phase: EngineResult,
feature: EngineResult,
timeframe: str,
) -> EngineResult:
if feature.payload.get("insufficient"):
return EngineResult(
name=self.name,
version=self.version,
confidence=20.0,
score=30.0,
reasons=["特征不足,跳过事件识别"],
warnings=["insufficient_features"],
payload={
"current_event": WyckoffEvent.NONE.value,
"active_events": [],
"recent_events": [], # alias for DB/API compat; same as active_events
"timeframe": timeframe,
"entry_score": 30.0,
},
)
context = {
"features": feature.payload,
"cycle": cycle.payload,
"phase": phase.payload,
"timeframe": timeframe,
}
hits = []
for rule in rule_registry.by_category("event", timeframe):
hit = rule.evaluate(context)
if hit and hit.event:
hits.append(hit)
if not hits:
return EngineResult(
name=self.name,
version=self.version,
confidence=35.0,
score=40.0,
reasons=["无显著事件"],
payload={
"current_event": WyckoffEvent.NONE.value,
"active_events": [],
"recent_events": [],
"timeframe": timeframe,
"entry_score": 40.0,
},
)
by_event: dict[str, float] = {}
reasons: list[str] = []
metrics: dict = {}
for h in hits:
prev = by_event.get(h.event, -1.0)
if h.confidence >= prev:
by_event[h.event] = h.confidence
reasons.extend(h.reasons)
metrics.update(h.metrics)
active = [e for e in _DISPLAY_ORDER if e in by_event]
for e in by_event:
if e not in active:
active.append(e)
# Dominant = max confidence; tie-break by dominance priority index
def _dom_key(ev: str) -> tuple:
conf = by_event[ev]
try:
prio = _DOMINANCE_PRIORITY.index(ev)
except ValueError:
prio = 99
return (conf, -prio)
current = max(by_event.keys(), key=_dom_key)
event_conf = by_event[current]
co_bonus = min(12.0, max(0, len(active) - 1) * 3)
entry_score = min(98.0, event_conf + co_bonus)
if current == WyckoffEvent.SPRING.value and WyckoffEvent.TEST.value in by_event:
entry_score = min(98.0, entry_score + 5)
return EngineResult(
name=self.name,
version=self.version,
confidence=event_conf,
score=entry_score,
reasons=list(dict.fromkeys(reasons))[:8],
warnings=["active_events_are_concurrent_not_timeline"],
metrics=metrics,
payload={
"current_event": current,
"active_events": active,
"recent_events": active, # persisted column name; semantic = active
"event_scores": by_event,
"timeframe": timeframe,
"entry_score": entry_score,
},
)
+206
View File
@@ -0,0 +1,206 @@
"""Feature Engine — pure function over OHLCVFrame → EngineResult(FeatureSnapshot)."""
from __future__ import annotations
from typing import Any
import numpy as np
from crypto_wyckoff.domain_models import EngineResult, OHLCVFrame
def _sma(arr: np.ndarray, n: int) -> float:
if len(arr) < n:
return float(arr[-1]) if len(arr) else 0.0
return float(np.mean(arr[-n:]))
def _atr(high: np.ndarray, low: np.ndarray, close: np.ndarray, n: int = 14) -> float:
if len(close) < 2:
return 0.0
prev_close = close[:-1]
tr = np.maximum(high[1:] - low[1:], np.maximum(np.abs(high[1:] - prev_close), np.abs(low[1:] - prev_close)))
if len(tr) < n:
return float(np.mean(tr)) if len(tr) else 0.0
return float(np.mean(tr[-n:]))
def _adx(high: np.ndarray, low: np.ndarray, close: np.ndarray, n: int = 14) -> float:
"""Simplified ADX approximation."""
if len(close) < n + 2:
return 15.0
up = high[1:] - high[:-1]
down = low[:-1] - low[1:]
plus_dm = np.where((up > down) & (up > 0), up, 0.0)
minus_dm = np.where((down > up) & (down > 0), down, 0.0)
tr = np.maximum(high[1:] - low[1:], np.maximum(np.abs(high[1:] - close[:-1]), np.abs(low[1:] - close[:-1])))
atr = np.mean(tr[-n:]) or 1e-9
plus_di = 100 * np.mean(plus_dm[-n:]) / atr
minus_di = 100 * np.mean(minus_dm[-n:]) / atr
denom = plus_di + minus_di
if denom < 1e-9:
return 10.0
dx = 100 * abs(plus_di - minus_di) / denom
return float(min(60.0, dx))
def compute_feature_snapshot(frame: OHLCVFrame) -> dict[str, Any]:
"""Compute technical snapshot dict from OHLCV (no I/O)."""
if frame.empty or len(frame) < 5:
return {"ts_code": frame.ts_code, "timeframe": frame.timeframe, "bars": len(frame)}
close = np.asarray(frame.close, dtype=float)
high = np.asarray(frame.high, dtype=float)
low = np.asarray(frame.low, dtype=float)
volume = np.asarray(frame.volume, dtype=float)
open_ = np.asarray(frame.open, dtype=float)
ma20 = _sma(close, 20)
ma60 = _sma(close, 60)
ma120 = _sma(close, min(120, len(close)))
atr = _atr(high, low, close, 14)
vol_ma20 = _sma(volume, 20) or 1e-9
volume_ratio = float(volume[-1] / vol_ma20)
look = min(60, len(close))
window_h = high[-look:]
window_l = low[-look:]
range_high = float(np.max(window_h))
range_low = float(np.min(window_l))
rng = max(range_high - range_low, 1e-9)
range_pct_60 = float(rng / close[-1]) if close[-1] else 0.0
range_position = float((close[-1] - range_low) / rng)
# Spring / UTAD hints
pierce_below = max(0.0, (range_low - low[-1]) / close[-1]) if close[-1] else 0.0
# if previous bars broke below and last close back in range
prior_low = float(np.min(low[-6:-1])) if len(low) >= 6 else float(low[-2])
pierce_below = max(pierce_below, max(0.0, (range_low - prior_low) / close[-1]))
close_back_in_range = 1.0 if close[-1] >= range_low else 0.0
reclaim_speed = 0.0
if pierce_below > 0 and close[-1] >= range_low:
reclaim_speed = min(1.0, (close[-1] - low[-1]) / max(atr, 1e-9) / 2)
pierce_above = max(0.0, (high[-1] - range_high) / close[-1])
fail_back = 1.0 if pierce_above > 0 and close[-1] <= range_high else 0.0
breakout_above = 1.0 if close[-1] > range_high and volume_ratio >= 1.0 else -1.0
# pullback hold: close near ma20 from above after being higher
pullback_hold = 0.0
if len(close) >= 5 and close[-1] > ma20 and close[-3] > close[-1] and (close[-1] - ma20) / max(atr, 1e-9) < 1.5:
pullback_hold = 0.8
ma60_prev = _sma(close[:-5], 60) if len(close) > 65 else ma60
ma60_slope = (ma60 - ma60_prev) / max(abs(ma60_prev), 1e-9)
# volume trend: recent 10 vs prior 10
if len(volume) >= 20:
volume_trend = float(np.mean(volume[-10:]) / (np.mean(volume[-20:-10]) + 1e-9) - 1.0)
else:
volume_trend = 0.0
bar_range_atr = float((high[-1] - low[-1]) / max(atr, 1e-9))
bounce_from_low = float((close[-1] - float(np.min(low[-10:]))) / close[-1]) if close[-1] else 0.0
gap_up_pct = float((open_[-1] - close[-2]) / close[-2]) if len(close) >= 2 and close[-2] else 0.0
after_strength = 0.0
if len(close) >= 4 and close[-3] > close[-4]:
after_strength = 0.7
spring_score_hint = 0.0
if pierce_below >= 0.002 and close_back_in_range:
spring_score_hint = min(90.0, 50 + pierce_below * 1500 + reclaim_speed * 20)
utad_score_hint = min(90.0, 50 + pierce_above * 1500) if pierce_above >= 0.002 and fail_back else 0.0
# swing
swing_high = float(np.max(high[-20:])) if len(high) >= 5 else float(high[-1])
swing_low = float(np.min(low[-20:])) if len(low) >= 5 else float(low[-1])
return {
"ts_code": frame.ts_code,
"timeframe": frame.timeframe,
"bars": len(frame),
"close": float(close[-1]),
"open": float(open_[-1]),
"high": float(high[-1]),
"low": float(low[-1]),
"volume": float(volume[-1]),
"ma20": ma20,
"ma60": ma60,
"ma120": ma120,
"ma60_slope": float(ma60_slope),
"atr": atr,
"adx": _adx(high, low, close),
"volume_ma20": float(vol_ma20),
"volume_ratio": volume_ratio,
"volume_trend": volume_trend,
"range_high": range_high,
"range_low": range_low,
"range_pct_60": range_pct_60,
"range_position": range_position,
"pierce_below_range": pierce_below,
"pierce_above_range": pierce_above,
"close_back_in_range": close_back_in_range,
"reclaim_speed": reclaim_speed,
"fail_back_into_range": fail_back,
"breakout_above_range": breakout_above,
"pullback_hold": pullback_hold,
"bar_range_atr": bar_range_atr,
"bounce_from_low": bounce_from_low,
"gap_up_pct": gap_up_pct,
"after_strength": after_strength,
"spring_score_hint": spring_score_hint,
"utad_score_hint": utad_score_hint,
"swing_high": swing_high,
"swing_low": swing_low,
"trade_date": str(frame.trade_dates[-1]) if frame.trade_dates else None,
}
# Minimum bars before a timeframe is considered usable (no cross-TF borrow)
_MIN_BARS = {"1d": 40, "1w": 26, "1M": 18}
class FeatureEngine:
"""Pure Feature Engine — no database access."""
name = "Feature"
version = "1.0.0"
def run(self, frame: OHLCVFrame | None, timeframe: str | None = None) -> EngineResult:
tf = timeframe or (frame.timeframe if frame else "1d")
min_bars = _MIN_BARS.get(tf, 30)
if frame is None or frame.empty or len(frame) < min_bars:
bars = 0 if frame is None or frame.empty else len(frame)
return EngineResult(
name=self.name,
version=self.version,
confidence=10.0,
score=10.0,
reasons=[f"{tf} bars={bars} < min={min_bars},标记 insufficient"],
warnings=["insufficient_features"],
metrics={"bars": bars, "min_bars": min_bars},
payload={
"ts_code": getattr(frame, "ts_code", ""),
"timeframe": tf,
"bars": bars,
"insufficient": True,
},
)
snap = compute_feature_snapshot(frame)
snap["insufficient"] = False
conf = 90.0 if snap.get("bars", 0) >= 60 else 50.0 + min(40.0, snap.get("bars", 0) * 0.5)
warnings = []
if snap.get("bars", 0) < 60:
warnings.append("bars偏少,特征可靠性中等")
return EngineResult(
name=self.name,
version=self.version,
confidence=conf,
score=conf,
reasons=[f"computed {snap.get('bars', 0)} bars {tf}"],
warnings=warnings,
metrics={"bars": snap.get("bars", 0)},
payload=snap,
)
+363
View File
@@ -0,0 +1,363 @@
"""Paths + OHLCV cache + DATA_SERVICE fetch (crypto continuous calendar)."""
from __future__ import annotations
import json
import logging
import os
import sqlite3
import time
from datetime import date, datetime, timezone
from pathlib import Path
from typing import Iterable
import requests
from crypto_wyckoff.domain_models import OHLCVFrame
logger = logging.getLogger(__name__)
_REPO_ROOT = Path(__file__).resolve().parents[1]
DATA_DIR = Path(os.environ.get("CRYPTO_WYCKOFF_DATA", str(_REPO_ROOT / "data" / "crypto_wyckoff")))
BARS_DB = DATA_DIR / "bars.sqlite"
SCAN_DB = DATA_DIR / "scan.sqlite"
DATA_SERVICE_URL = os.environ.get(
"DATA_SERVICE_URL",
os.environ.get("DATASVC_URL", "https://provider.jackyu66.com"),
).rstrip("/")
# Continuous crypto: bar counts (not A-share weekend-padded calendar multipliers)
# Provider has many TFs; 1M is resampled locally from daily UTC months.
LOOKBACK = {
"1h": 500,
"2h": 400,
"4h": 300,
"6h": 280,
"8h": 250,
"12h": 220,
"1d": 250,
"1w": 104,
"1M": 60,
}
# Default D/W/M stack (kept for compat); combos may request more TFs from provider.
TF_PROVIDER = ("1h", "4h", "8h", "1d", "1w")
TF_LIST = ("1d", "1w", "1M")
LOCAL_ONLY_TFS = frozenset({"1M"})
def ensure_dirs() -> None:
DATA_DIR.mkdir(parents=True, exist_ok=True)
def _symbol_key(symbol: str) -> str:
return symbol.replace("/", "_").replace(":", "_")
def _bars_conn() -> sqlite3.Connection:
ensure_dirs()
conn = sqlite3.connect(str(BARS_DB), timeout=60)
conn.execute(
"""
CREATE TABLE IF NOT EXISTS bars (
symbol TEXT NOT NULL,
tf TEXT NOT NULL,
ts INTEGER NOT NULL,
open REAL, high REAL, low REAL, close REAL, volume REAL,
PRIMARY KEY (symbol, tf, ts)
)
"""
)
conn.execute("CREATE INDEX IF NOT EXISTS idx_bars_sym_tf ON bars(symbol, tf)")
return conn
def fetch_candles(
symbol: str,
tf: str,
*,
limit: int | None = None,
start_ms: int | None = None,
end_ms: int | None = None,
timeout: float = 15.0,
) -> list[dict]:
params: dict = {"symbol": symbol, "tf": tf}
if limit is not None:
params["limit"] = int(limit)
if start_ms is not None:
params["start"] = int(start_ms)
if end_ms is not None:
params["end"] = int(end_ms)
resp = requests.get(f"{DATA_SERVICE_URL}/api/candles", params=params, timeout=timeout)
resp.raise_for_status()
data = resp.json()
if not isinstance(data, list):
return []
out = []
for row in data:
try:
ts = int(float(row["timestamp"]))
out.append(
{
"ts": ts,
"open": float(row["open"]),
"high": float(row["high"]),
"low": float(row["low"]),
"close": float(row["close"]),
"volume": float(row.get("volume") or 0),
}
)
except (KeyError, TypeError, ValueError):
continue
out.sort(key=lambda r: r["ts"])
return out
def upsert_bars(symbol: str, tf: str, rows: list[dict]) -> int:
if not rows:
return 0
conn = _bars_conn()
try:
conn.executemany(
"""
INSERT INTO bars(symbol, tf, ts, open, high, low, close, volume)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(symbol, tf, ts) DO UPDATE SET
open=excluded.open, high=excluded.high, low=excluded.low,
close=excluded.close, volume=excluded.volume
""",
[
(symbol, tf, r["ts"], r["open"], r["high"], r["low"], r["close"], r["volume"])
for r in rows
],
)
conn.commit()
return len(rows)
finally:
conn.close()
def is_intraday_tf(tf: str) -> bool:
"""True for minute/hour TFs that need clock time on charts."""
t = (tf or "").strip()
return t.endswith("m") or t.endswith("h")
def load_bars_with_ts(
symbol: str, tf: str, lookback: int | None = None
) -> list[dict]:
"""Return OHLCV rows with UTC ms ts (for chart labels).
``datetime`` is wall-clock in Asia/Shanghai (UTC+8) for display.
"""
from zoneinfo import ZoneInfo
tz_cn = ZoneInfo("Asia/Shanghai")
if lookback is None:
try:
from crypto_wyckoff.combos import lookback_for
lookback = lookback_for(tf)
except Exception:
lookback = LOOKBACK.get(tf, 100)
lookback = lookback or LOOKBACK.get(tf, 100)
conn = _bars_conn()
try:
cur = conn.execute(
"""
SELECT ts, open, high, low, close, volume FROM bars
WHERE symbol=? AND tf=?
ORDER BY ts DESC LIMIT ?
""",
(symbol, tf, lookback),
)
rows = list(reversed(cur.fetchall()))
finally:
conn.close()
out = []
for ts, o, h, l, c, v in rows:
dt_utc = datetime.fromtimestamp(ts / 1000.0, tz=timezone.utc)
dt_cn = dt_utc.astimezone(tz_cn)
out.append(
{
"ts": int(ts),
"datetime": dt_cn.strftime("%Y-%m-%dT%H:%M:%S+08:00"),
"date": dt_cn.strftime("%Y-%m-%d"),
"open": o,
"high": h,
"low": l,
"close": c,
"volume": v,
}
)
return out
def load_frame(symbol: str, tf: str, lookback: int | None = None) -> OHLCVFrame | None:
rows = load_bars_with_ts(symbol, tf, lookback)
if not rows:
return None
return OHLCVFrame(
ts_code=symbol,
timeframe=tf,
trade_dates=[
datetime.fromtimestamp(r["ts"] / 1000.0, tz=timezone.utc).date() for r in rows
],
open=[r["open"] for r in rows],
high=[r["high"] for r in rows],
low=[r["low"] for r in rows],
close=[r["close"] for r in rows],
volume=[r["volume"] for r in rows],
)
def bar_count(symbol: str, tf: str) -> int:
conn = _bars_conn()
try:
cur = conn.execute(
"SELECT COUNT(*) FROM bars WHERE symbol=? AND tf=?", (symbol, tf)
)
return int(cur.fetchone()[0])
finally:
conn.close()
def rebuild_monthly_from_daily(symbol: str) -> int:
"""Aggregate UTC calendar-month OHLCV from local daily bars (provider has no 1M)."""
conn = _bars_conn()
try:
cur = conn.execute(
"""
SELECT ts, open, high, low, close, volume FROM bars
WHERE symbol=? AND tf='1d' ORDER BY ts ASC
""",
(symbol,),
)
daily = cur.fetchall()
finally:
conn.close()
if not daily:
return 0
months: dict[tuple[int, int], dict] = {}
for ts, o, h, l, c, v in daily:
dt = datetime.fromtimestamp(ts / 1000.0, tz=timezone.utc)
key = (dt.year, dt.month)
# month bar open timestamp = first day 00:00 UTC
month_ts = int(datetime(dt.year, dt.month, 1, tzinfo=timezone.utc).timestamp() * 1000)
if key not in months:
months[key] = {
"ts": month_ts,
"open": o,
"high": h,
"low": l,
"close": c,
"volume": v or 0.0,
}
else:
m = months[key]
m["high"] = max(m["high"], h)
m["low"] = min(m["low"], l)
m["close"] = c
m["volume"] = (m["volume"] or 0) + (v or 0)
rows = sorted(months.values(), key=lambda r: r["ts"])
# drop stale months then upsert
conn = _bars_conn()
try:
conn.execute("DELETE FROM bars WHERE symbol=? AND tf='1M'", (symbol,))
conn.commit()
finally:
conn.close()
return upsert_bars(symbol, "1M", rows)
def backfill_symbol(symbol: str, tfs: Iterable[str] = TF_LIST) -> dict:
"""Pull history for requested TFs; monthly derived from daily when needed."""
wanted = list(dict.fromkeys(tfs))
stats: dict = {}
need_monthly = "1M" in wanted
if need_monthly and "1d" not in wanted:
wanted = ["1d", *wanted]
for tf in wanted:
if tf in LOCAL_ONLY_TFS:
continue
need = LOOKBACK.get(tf, 100)
if tf == "1d" and need_monthly:
need = max(need, LOOKBACK["1M"] * 31)
try:
rows = fetch_candles(symbol, tf, limit=need)
n = upsert_bars(symbol, tf, rows)
stats[tf] = n
except Exception as e:
logger.warning("backfill %s %s failed: %s", symbol, tf, e)
stats[tf] = 0
time.sleep(0.05)
if need_monthly:
try:
stats["1M"] = rebuild_monthly_from_daily(symbol)
except Exception as e:
logger.warning("monthly rebuild %s failed: %s", symbol, e)
stats["1M"] = 0
return stats
def tip_update_symbol(symbol: str, tfs: Iterable[str] = TF_LIST) -> bool:
"""Update forming tip bars (limit=3). Returns True if any bar changed."""
wanted = list(dict.fromkeys(tfs))
changed = False
for tf in wanted:
if tf in LOCAL_ONLY_TFS:
continue
try:
rows = fetch_candles(symbol, tf, limit=3)
if not rows:
continue
before = _tip_fingerprint(symbol, tf)
upsert_bars(symbol, tf, rows)
after = _tip_fingerprint(symbol, tf)
if before != after:
changed = True
except Exception as e:
logger.debug("tip %s %s: %s", symbol, tf, e)
time.sleep(0.02)
if "1M" in wanted:
before_m = _tip_fingerprint(symbol, "1M")
try:
rebuild_monthly_from_daily(symbol)
except Exception as e:
logger.debug("monthly tip %s: %s", symbol, e)
after_m = _tip_fingerprint(symbol, "1M")
if before_m != after_m:
changed = True
return changed
def _tip_fingerprint(symbol: str, tf: str) -> tuple | None:
conn = _bars_conn()
try:
cur = conn.execute(
"""
SELECT ts, open, high, low, close, volume FROM bars
WHERE symbol=? AND tf=? ORDER BY ts DESC LIMIT 1
""",
(symbol, tf),
)
row = cur.fetchone()
return tuple(row) if row else None
finally:
conn.close()
def fetch_symbols_from_provider() -> list[str]:
try:
resp = requests.get(f"{DATA_SERVICE_URL}/health", timeout=8)
resp.raise_for_status()
payload = resp.json()
symbols = payload.get("symbols") or payload.get("symbol_list") or []
return [s for s in symbols if isinstance(s, str)]
except Exception as e:
logger.warning("health symbols failed: %s", e)
return []
+78
View File
@@ -0,0 +1,78 @@
"""Phase Engine — Phase AE via Rule Registry."""
from __future__ import annotations
from crypto_wyckoff.domain_models import EngineResult, WyckoffPhase
from crypto_wyckoff.rules.registry import rule_registry
class PhaseEngine:
name = "Phase"
version = "1.0.0"
def run(self, cycle: EngineResult, feature: EngineResult, timeframe: str) -> EngineResult:
if feature.payload.get("insufficient") or cycle.payload.get("cycle") == "Unknown":
return EngineResult(
name=self.name,
version=self.version,
confidence=20.0,
score=30.0,
reasons=["数据/周期不足,Phase=None"],
warnings=["insufficient_features"],
payload={
"phase": WyckoffPhase.NONE.value,
"timeframe": timeframe,
"cycle": cycle.payload.get("cycle"),
"structure_score": 30.0,
},
)
context = {
"features": feature.payload,
"cycle": cycle.payload,
"timeframe": timeframe,
}
hits = []
for rule in rule_registry.by_category("phase", timeframe):
hit = rule.evaluate(context)
if hit and hit.phase:
hits.append(hit)
if not hits:
return EngineResult(
name=self.name,
version=self.version,
confidence=40.0,
score=cycle.score * 0.5,
reasons=["未识别明确 Phase"],
payload={
"phase": WyckoffPhase.NONE.value,
"timeframe": timeframe,
"cycle": cycle.payload.get("cycle"),
"structure_score": cycle.score * 0.5,
},
)
best = max(hits, key=lambda h: h.confidence)
structure_score = best.score
# Phase D/E stronger structure
if best.phase in (WyckoffPhase.D.value, WyckoffPhase.E.value):
structure_score = max(structure_score, 80.0)
elif best.phase == WyckoffPhase.C.value:
structure_score = max(structure_score, 72.0)
return EngineResult(
name=self.name,
version=self.version,
confidence=best.confidence,
score=structure_score,
reasons=best.reasons,
metrics=best.metrics,
payload={
"phase": best.phase,
"timeframe": timeframe,
"cycle": cycle.payload.get("cycle"),
"rule_id": best.rule_id,
"structure_score": structure_score,
},
)
+181
View File
@@ -0,0 +1,181 @@
"""Scan pipeline: load local frames → engines → store (per TF combo)."""
from __future__ import annotations
import json
import logging
from datetime import date, datetime, timezone
from crypto_wyckoff.combos import ROLE_HIGH, ROLE_LOW, ROLE_MID, get_combo, lookback_for
from crypto_wyckoff.cycle import CycleEngine
from crypto_wyckoff.decision import DecisionEngine
from crypto_wyckoff.domain_models import WyckoffScanRow
from crypto_wyckoff.event import EventEngine
from crypto_wyckoff.features import FeatureEngine
from crypto_wyckoff.io import load_frame
from crypto_wyckoff.phase import PhaseEngine
from crypto_wyckoff.plan import PlanEngine
from crypto_wyckoff.signal import SignalEngine
from crypto_wyckoff.store import upsert_row
from crypto_wyckoff.symbols_cn import display_name_cn
from crypto_wyckoff.version import WYCKOFF_ENGINE_VERSION
logger = logging.getLogger(__name__)
def analyze_symbol(
low_frame,
mid_frame,
high_frame,
*,
feature_eng: FeatureEngine,
cycle_eng: CycleEngine,
phase_eng: PhaseEngine,
event_eng: EventEngine,
signal_eng: SignalEngine,
decision_eng: DecisionEngine,
plan_eng: PlanEngine,
) -> dict:
"""Run engines with D/W/M *role* aliases so existing rules match.
Frames may be any TF combo (e.g. 1h/4h/8h); rules still see 1d/1w/1M roles.
"""
f_d = feature_eng.run(low_frame, ROLE_LOW)
f_w = feature_eng.run(mid_frame, ROLE_MID)
f_m = feature_eng.run(high_frame, ROLE_HIGH)
c_m = cycle_eng.run(f_m, ROLE_HIGH)
c_w = cycle_eng.run(f_w, ROLE_MID)
p_w = phase_eng.run(c_w, f_w, ROLE_MID)
p_d = phase_eng.run(c_w, f_d, ROLE_LOW)
e_w = event_eng.run(c_w, p_w, f_w, ROLE_MID)
e_d = event_eng.run(c_w, p_d, f_d, ROLE_LOW)
s_d = signal_eng.run(e_d, p_d)
decision = decision_eng.run(c_m, c_w, p_w, e_w, e_d, s_d)
plan = plan_eng.run(f_d, decision)
return {
"f_d": f_d, "f_w": f_w, "f_m": f_m,
"c_m": c_m, "c_w": c_w, "p_w": p_w,
"e_w": e_w, "e_d": e_d, "s_d": s_d,
"decision": decision, "plan": plan,
}
def _to_row(
trade_date: date,
symbol: str,
result: dict,
*,
combo_id: str,
combo_label: str,
) -> WyckoffScanRow:
d = result["decision"]
p = result["plan"]
c_m, c_w, p_w = result["c_m"], result["c_w"], result["p_w"]
e_w, e_d, s_d = result["e_w"], result["e_d"], result["s_d"]
f_d, f_w, f_m = result["f_d"], result["f_w"], result["f_m"]
snapshot = {
"combo_id": combo_id,
"combo_label": combo_label,
"daily": {k: f_d.payload.get(k) for k in (
"ma20", "ma60", "ma120", "atr", "adx", "volume_ratio",
"range_high", "range_low", "swing_high", "swing_low", "close",
)},
"weekly": {k: f_w.payload.get(k) for k in ("ma20", "ma60", "adx", "close")},
"monthly": {k: f_m.payload.get(k) for k in ("ma20", "ma60", "adx", "close")},
}
markers = []
for key, typ in (("entry", "entry"), ("stop", "stop"), ("target1", "target1"), ("target2", "target2")):
if p.payload.get(key) is not None:
markers.append({"type": typ, "price": p.payload[key]})
return WyckoffScanRow(
trade_date=trade_date,
ts_code=symbol,
name=display_name_cn(symbol),
industry="crypto",
engine_version=WYCKOFF_ENGINE_VERSION,
m_cycle=c_m.payload.get("cycle", "Unknown"),
cycle_confidence=c_m.confidence,
trend_score=float(d.payload.get("trend_score", c_m.score)),
w_cycle=c_w.payload.get("cycle", "Unknown"),
w_phase=p_w.payload.get("phase", "None"),
w_current_event=e_w.payload.get("current_event", "None"),
w_recent_events_json=json.dumps(
e_w.payload.get("active_events") or e_w.payload.get("recent_events") or [],
ensure_ascii=False,
),
phase_confidence=p_w.confidence,
structure_score=float(d.payload.get("structure_score", p_w.score)),
d_current_event=e_d.payload.get("current_event", "None"),
d_recent_events_json=json.dumps(
e_d.payload.get("active_events") or e_d.payload.get("recent_events") or [],
ensure_ascii=False,
),
event_confidence=e_d.confidence,
entry_score=float(d.payload.get("entry_score", e_d.score)),
entry=p.payload.get("entry"),
stop=p.payload.get("stop"),
target1=p.payload.get("target1"),
target2=p.payload.get("target2"),
rr=p.payload.get("rr"),
alignment=float(d.payload.get("alignment", 0)),
stars=int(d.payload.get("stars", 1)),
decision_signal=d.payload.get("decision_signal", "Watch"),
signal_confidence=s_d.confidence,
overall_confidence=float(d.payload.get("overall_confidence", d.confidence)),
overall_score=float(d.payload.get("overall_score", d.score)),
risk=d.payload.get("risk", "Medium"),
reasons_json=json.dumps(d.reasons + d.warnings, ensure_ascii=False),
feature_snapshot_json=json.dumps(snapshot, ensure_ascii=False),
markers_json=json.dumps(markers, ensure_ascii=False),
scanned_at=datetime.now(timezone.utc),
combo_id=combo_id,
)
_ENGINES = None
def _engines():
global _ENGINES
if _ENGINES is None:
_ENGINES = {
"feature_eng": FeatureEngine(),
"cycle_eng": CycleEngine(),
"phase_eng": PhaseEngine(),
"event_eng": EventEngine(),
"signal_eng": SignalEngine(),
"decision_eng": DecisionEngine(),
"plan_eng": PlanEngine(),
}
return _ENGINES
def analyze_and_store(
symbol: str,
trade_date: date | None = None,
*,
combo_id: str | None = None,
) -> WyckoffScanRow | None:
eng = _engines()
combo = get_combo(combo_id)
low_tf, mid_tf, high_tf = combo["low"], combo["mid"], combo["high"]
low = load_frame(symbol, low_tf, lookback_for(low_tf))
mid = load_frame(symbol, mid_tf, lookback_for(mid_tf))
high = load_frame(symbol, high_tf, lookback_for(high_tf))
if low is None or len(low) < 40:
return None
result = analyze_symbol(low, mid, high, **eng)
td = trade_date or (
low.trade_dates[-1] if low.trade_dates else datetime.now(timezone.utc).date()
)
row = _to_row(td, symbol, result, combo_id=combo["id"], combo_label=combo["label"])
upsert_row(row)
return row
+78
View File
@@ -0,0 +1,78 @@
"""Plan Engine — Entry / Stop / Target / RR only when Decision is tradable."""
from __future__ import annotations
from crypto_wyckoff.domain_models import DecisionSignal, EngineResult
_TRADABLE = {
DecisionSignal.STRONG_BUY.value,
DecisionSignal.BUY.value,
DecisionSignal.SELL.value,
}
class PlanEngine:
name = "Plan"
version = "1.0.0"
def run(self, daily_feature: EngineResult, decision: EngineResult) -> EngineResult:
f = daily_feature.payload
close = float(f.get("close") or 0)
atr = float(f.get("atr") or 0) or close * 0.02
swing_low = float(f.get("swing_low") or close - 2 * atr)
swing_high = float(f.get("swing_high") or close + 2 * atr)
range_high = float(f.get("range_high") or swing_high)
signal = decision.payload.get("decision_signal", DecisionSignal.WATCH.value)
entry = stop = t1 = t2 = rr = None
reasons: list[str] = []
if signal not in _TRADABLE or close <= 0:
reasons.append(f"无交易计划(信号={signal}")
return EngineResult(
name=self.name,
version=self.version,
confidence=decision.confidence,
score=decision.score,
reasons=reasons,
payload={
"entry": None,
"stop": None,
"target1": None,
"target2": None,
"rr": None,
},
)
if signal in (DecisionSignal.STRONG_BUY.value, DecisionSignal.BUY.value):
entry = round(close, 4)
stop = round(min(swing_low, close - 1.5 * atr), 4)
risk = max(entry - stop, 1e-6)
t1 = round(entry + 2.0 * risk, 4)
t2 = round(max(range_high, entry + 3.0 * risk), 4)
rr = round((t1 - entry) / risk, 2)
reasons.append(f"入场={entry} 止损={stop} 目标一={t1} 盈亏比={rr}")
else: # Sell
entry = round(close, 4)
stop = round(max(swing_high, close + 1.5 * atr), 4)
risk = max(stop - entry, 1e-6)
t1 = round(entry - 2.0 * risk, 4)
t2 = round(entry - 3.0 * risk, 4)
rr = round((entry - t1) / risk, 2)
reasons.append(f"做空计划 入场={entry} 止损={stop} 目标一={t1}")
return EngineResult(
name=self.name,
version=self.version,
confidence=decision.confidence,
score=decision.score,
reasons=reasons,
payload={
"entry": entry,
"stop": stop,
"target1": t1,
"target2": t2,
"rr": rr,
},
)
+3
View File
@@ -0,0 +1,3 @@
from crypto_wyckoff.rules.registry import rule_registry
__all__ = ["rule_registry"]
+33
View File
@@ -0,0 +1,33 @@
"""Rule protocol for Wyckoff Rule Registry."""
from __future__ import annotations
from abc import ABC, abstractmethod
from dataclasses import dataclass, field
from typing import Any
@dataclass
class RuleHit:
"""A single rule match."""
rule_id: str
event: str | None = None
phase: str | None = None
cycle: str | None = None
confidence: float = 0.0
score: float = 0.0
reasons: list[str] = field(default_factory=list)
metrics: dict[str, Any] = field(default_factory=dict)
class WyckoffRule(ABC):
"""Pluggable rule. Engines iterate registry; never hardcode rule lists."""
rule_id: str
category: str # cycle | phase | event
timeframes: tuple[str, ...] = ("1d", "1w", "1M")
@abstractmethod
def evaluate(self, context: dict[str, Any]) -> RuleHit | None:
"""Return RuleHit if matched, else None. Pure — no I/O."""
+126
View File
@@ -0,0 +1,126 @@
"""Cycle classification rules (monthly / weekly)."""
from __future__ import annotations
from typing import Any
from crypto_wyckoff.domain_models import WyckoffCycle
from crypto_wyckoff.rules.base import RuleHit, WyckoffRule
def _f(ctx: dict[str, Any], key: str, default: float = 0.0) -> float:
v = ctx.get("features", {}).get(key, default)
try:
return float(v) if v is not None else default
except (TypeError, ValueError):
return default
class MarkupCycleRule(WyckoffRule):
rule_id = "cycle_markup"
category = "cycle"
timeframes = ("1M", "1w")
def evaluate(self, context: dict[str, Any]) -> RuleHit | None:
close = _f(context, "close")
ma20 = _f(context, "ma20")
ma60 = _f(context, "ma60")
ma120 = _f(context, "ma120")
adx = _f(context, "adx")
slope = _f(context, "ma60_slope")
if close > ma20 > ma60 and (ma60 >= ma120 or slope > 0) and adx >= 18:
conf = min(95.0, 55 + adx + (10 if close > ma120 else 0))
return RuleHit(
rule_id=self.rule_id,
cycle=WyckoffCycle.MARKUP.value,
confidence=conf,
score=conf,
reasons=["价格位于均线多头排列", f"ADX={adx:.1f}"],
metrics={"adx": adx, "slope": slope},
)
return None
class MarkdownCycleRule(WyckoffRule):
rule_id = "cycle_markdown"
category = "cycle"
timeframes = ("1M", "1w")
def evaluate(self, context: dict[str, Any]) -> RuleHit | None:
close = _f(context, "close")
ma20 = _f(context, "ma20")
ma60 = _f(context, "ma60")
ma120 = _f(context, "ma120")
adx = _f(context, "adx")
slope = _f(context, "ma60_slope")
if close < ma20 < ma60 and (ma60 <= ma120 or slope < 0) and adx >= 18:
conf = min(95.0, 55 + adx + (10 if close < ma120 else 0))
return RuleHit(
rule_id=self.rule_id,
cycle=WyckoffCycle.MARKDOWN.value,
confidence=conf,
score=conf,
reasons=["价格位于均线空头排列", f"ADX={adx:.1f}"],
metrics={"adx": adx},
)
return None
class AccumulationCycleRule(WyckoffRule):
rule_id = "cycle_accumulation"
category = "cycle"
timeframes = ("1M", "1w")
def evaluate(self, context: dict[str, Any]) -> RuleHit | None:
adx = _f(context, "adx")
range_pct = _f(context, "range_pct_60")
close = _f(context, "close")
ma120 = _f(context, "ma120")
vol_trend = _f(context, "volume_trend")
# Range-bound after decline: strictly at/below MA120 (mutually exclusive vs Distribution)
if adx < 22 and range_pct < 0.28 and close <= ma120:
conf = 60 + (10 if vol_trend > 0 else 0) + (10 if close < ma120 else 0)
return RuleHit(
rule_id=self.rule_id,
cycle=WyckoffCycle.ACCUMULATION.value,
confidence=min(90.0, conf),
score=min(90.0, conf),
reasons=["低趋势强度区间震荡", "疑似吸筹区间"],
metrics={"adx": adx, "range_pct_60": range_pct},
)
return None
class DistributionCycleRule(WyckoffRule):
rule_id = "cycle_distribution"
category = "cycle"
timeframes = ("1M", "1w")
def evaluate(self, context: dict[str, Any]) -> RuleHit | None:
adx = _f(context, "adx")
range_pct = _f(context, "range_pct_60")
close = _f(context, "close")
ma120 = _f(context, "ma120")
vol_trend = _f(context, "volume_trend")
# Range-bound near highs: strictly above MA120 (mutually exclusive vs Accumulation)
if adx < 22 and range_pct < 0.28 and close > ma120:
conf = 60 + (10 if vol_trend < 0 else 0) + (10 if close > ma120 else 0)
return RuleHit(
rule_id=self.rule_id,
cycle=WyckoffCycle.DISTRIBUTION.value,
confidence=min(90.0, conf),
score=min(90.0, conf),
reasons=["高位低趋势震荡", "疑似派发区间"],
metrics={"adx": adx, "range_pct_60": range_pct},
)
return None
def build_rules() -> list[WyckoffRule]:
# Order: trend cycles first (more decisive), then range cycles
return [
MarkupCycleRule(),
MarkdownCycleRule(),
AccumulationCycleRule(),
DistributionCycleRule(),
]
+254
View File
@@ -0,0 +1,254 @@
"""Event rules: Spring/SOS/LPS/UTAD/SC/AR/ST/..."""
from __future__ import annotations
from typing import Any
from crypto_wyckoff.domain_models import WyckoffCycle, WyckoffEvent, WyckoffPhase
from crypto_wyckoff.rules.base import RuleHit, WyckoffRule
def _f(ctx: dict[str, Any], key: str, default: float = 0.0) -> float:
v = ctx.get("features", {}).get(key, default)
try:
return float(v) if v is not None else default
except (TypeError, ValueError):
return default
def _cycle(ctx: dict[str, Any]) -> str:
return (ctx.get("cycle") or {}).get("cycle") or ""
def _phase(ctx: dict[str, Any]) -> str:
return (ctx.get("phase") or {}).get("phase") or ""
class SpringRule(WyckoffRule):
rule_id = "event_spring"
category = "event"
timeframes = ("1d",)
def evaluate(self, context: dict[str, Any]) -> RuleHit | None:
cycle = _cycle(context)
if cycle not in (WyckoffCycle.ACCUMULATION.value, WyckoffCycle.RE_ACCUMULATION.value,
WyckoffCycle.MARKUP.value):
# Allow spring only in accumulative contexts; Decision will filter MTF
if cycle == WyckoffCycle.DISTRIBUTION.value:
pass # still detect for facts but lower confidence
pierce = _f(context, "pierce_below_range")
reclaim = _f(context, "reclaim_speed")
vol_ratio = _f(context, "volume_ratio")
close_in_range = _f(context, "close_back_in_range")
if pierce >= 0.002 and close_in_range >= 0.5 and reclaim >= 0.3:
strength = min(98.0, 50 + pierce * 2000 + reclaim * 20 + (15 if vol_ratio < 1.2 else 5))
return RuleHit(
rule_id=self.rule_id,
event=WyckoffEvent.SPRING.value,
confidence=strength,
score=strength,
reasons=[
f"跌破区间后收回 (pierce={pierce:.3%})",
f"回收速度={reclaim:.2f}",
f"量比={vol_ratio:.2f}",
],
metrics={"pierce": pierce, "reclaim": reclaim, "volume_ratio": vol_ratio},
)
return None
class TestRule(WyckoffRule):
rule_id = "event_test"
category = "event"
timeframes = ("1d", "1w")
def evaluate(self, context: dict[str, Any]) -> RuleHit | None:
pos = _f(context, "range_position")
vol_ratio = _f(context, "volume_ratio")
near_low = pos < 0.2
if near_low and vol_ratio < 0.85:
return RuleHit(
rule_id=self.rule_id,
event=WyckoffEvent.TEST.value,
confidence=68.0,
score=65.0,
reasons=["低位缩量回测"],
)
return None
class SOSRule(WyckoffRule):
rule_id = "event_sos"
category = "event"
timeframes = ("1d", "1w")
def evaluate(self, context: dict[str, Any]) -> RuleHit | None:
breakout = _f(context, "breakout_above_range")
vol_ratio = _f(context, "volume_ratio")
close = _f(context, "close")
ma20 = _f(context, "ma20")
if breakout >= 0.0 and vol_ratio >= 1.2 and close > ma20:
conf = min(95.0, 70 + vol_ratio * 8)
return RuleHit(
rule_id=self.rule_id,
event=WyckoffEvent.SOS.value,
confidence=conf,
score=conf,
reasons=["放量突破区间上沿 (SOS)"],
metrics={"vol_ratio": vol_ratio},
)
return None
class LPSRule(WyckoffRule):
rule_id = "event_lps"
category = "event"
timeframes = ("1d", "1w")
def evaluate(self, context: dict[str, Any]) -> RuleHit | None:
# Pullback hold above broken range / MA20 after prior strength
pullback = _f(context, "pullback_hold")
vol_ratio = _f(context, "volume_ratio")
above_ma = _f(context, "close") > _f(context, "ma20")
if pullback >= 0.5 and above_ma and vol_ratio <= 1.1:
return RuleHit(
rule_id=self.rule_id,
event=WyckoffEvent.LPS.value,
confidence=74.0,
score=76.0,
reasons=["突破后缩量回踩支撑 (LPS)"],
)
return None
class SCRule(WyckoffRule):
rule_id = "event_sc"
category = "event"
timeframes = ("1w", "1d")
def evaluate(self, context: dict[str, Any]) -> RuleHit | None:
vol_ratio = _f(context, "volume_ratio")
bar_range = _f(context, "bar_range_atr")
pos = _f(context, "range_position")
if vol_ratio >= 1.8 and bar_range >= 1.5 and pos < 0.35:
return RuleHit(
rule_id=self.rule_id,
event=WyckoffEvent.SC.value,
confidence=72.0,
score=70.0,
reasons=["低位放量宽幅,疑似 Selling Climax"],
)
return None
class ARRule(WyckoffRule):
rule_id = "event_ar"
category = "event"
timeframes = ("1w", "1d")
def evaluate(self, context: dict[str, Any]) -> RuleHit | None:
# Automatic rally: bounce from lows
bounce = _f(context, "bounce_from_low")
if bounce >= 0.04:
return RuleHit(
rule_id=self.rule_id,
event=WyckoffEvent.AR.value,
confidence=65.0,
score=62.0,
reasons=["低点后自动反弹 (AR)"],
)
return None
class STRule(WyckoffRule):
rule_id = "event_st"
category = "event"
timeframes = ("1w", "1d")
def evaluate(self, context: dict[str, Any]) -> RuleHit | None:
pos = _f(context, "range_position")
vol_ratio = _f(context, "volume_ratio")
if 0.15 < pos < 0.45 and vol_ratio < 1.0:
return RuleHit(
rule_id=self.rule_id,
event=WyckoffEvent.ST.value,
confidence=60.0,
score=58.0,
reasons=["次级测试 (ST)"],
)
return None
class UTADRule(WyckoffRule):
rule_id = "event_utad"
category = "event"
timeframes = ("1w", "1d")
def evaluate(self, context: dict[str, Any]) -> RuleHit | None:
cycle = _cycle(context)
pierce_up = _f(context, "pierce_above_range")
fail = _f(context, "fail_back_into_range")
if cycle in (WyckoffCycle.DISTRIBUTION.value, WyckoffCycle.RE_DISTRIBUTION.value,
WyckoffCycle.MARKUP.value):
if pierce_up >= 0.002 and fail >= 0.5:
return RuleHit(
rule_id=self.rule_id,
event=WyckoffEvent.UTAD.value,
confidence=76.0,
score=74.0,
reasons=["冲高失败回到区间 (UTAD)"],
)
return None
class JumpRule(WyckoffRule):
rule_id = "event_jump"
category = "event"
timeframes = ("1d",)
def evaluate(self, context: dict[str, Any]) -> RuleHit | None:
gap = _f(context, "gap_up_pct")
vol_ratio = _f(context, "volume_ratio")
if gap >= 0.03 and vol_ratio >= 1.3:
return RuleHit(
rule_id=self.rule_id,
event=WyckoffEvent.JUMP.value,
confidence=70.0,
score=72.0,
reasons=["放量向上跳跃 (Jump)"],
)
return None
class BackupRule(WyckoffRule):
rule_id = "event_backup"
category = "event"
timeframes = ("1d",)
def evaluate(self, context: dict[str, Any]) -> RuleHit | None:
pullback = _f(context, "pullback_hold")
after_jump = _f(context, "after_strength")
if after_jump >= 0.5 and pullback >= 0.5:
return RuleHit(
rule_id=self.rule_id,
event=WyckoffEvent.BACKUP.value,
confidence=68.0,
score=70.0,
reasons=["跳跃后回踩 (Backup)"],
)
return None
def build_rules() -> list[WyckoffRule]:
return [
SpringRule(),
UTADRule(),
SOSRule(),
LPSRule(),
SCRule(),
JumpRule(),
BackupRule(),
TestRule(),
ARRule(),
STRule(),
]
+163
View File
@@ -0,0 +1,163 @@
"""Phase AE rules (primarily weekly)."""
from __future__ import annotations
from typing import Any
from crypto_wyckoff.domain_models import WyckoffCycle, WyckoffPhase
from crypto_wyckoff.rules.base import RuleHit, WyckoffRule
def _f(ctx: dict[str, Any], key: str, default: float = 0.0) -> float:
v = ctx.get("features", {}).get(key, default)
try:
return float(v) if v is not None else default
except (TypeError, ValueError):
return default
def _cycle(ctx: dict[str, Any]) -> str:
return (ctx.get("cycle") or {}).get("cycle") or WyckoffCycle.UNKNOWN.value
class PhaseARule(WyckoffRule):
rule_id = "phase_a"
category = "phase"
timeframes = ("1w", "1d")
def evaluate(self, context: dict[str, Any]) -> RuleHit | None:
cycle = _cycle(context)
if cycle not in (WyckoffCycle.ACCUMULATION.value, WyckoffCycle.DISTRIBUTION.value,
WyckoffCycle.RE_ACCUMULATION.value, WyckoffCycle.RE_DISTRIBUTION.value):
return None
# Stopping action: high vol + large range recently, still range-bound
vol_ratio = _f(context, "volume_ratio")
range_last = _f(context, "bar_range_atr")
if vol_ratio >= 1.4 and range_last >= 1.2:
return RuleHit(
rule_id=self.rule_id,
phase=WyckoffPhase.A.value,
confidence=70.0,
score=65.0,
reasons=["放量宽幅波动,疑似 Phase A 停止行为"],
)
return None
class PhaseBRule(WyckoffRule):
rule_id = "phase_b"
category = "phase"
timeframes = ("1w", "1d")
def evaluate(self, context: dict[str, Any]) -> RuleHit | None:
cycle = _cycle(context)
if cycle not in (WyckoffCycle.ACCUMULATION.value, WyckoffCycle.DISTRIBUTION.value):
return None
adx = _f(context, "adx")
range_pct = _f(context, "range_pct_60")
pos = _f(context, "range_position") # 0=low 1=high of range
if adx < 20 and 0.25 < pos < 0.75 and range_pct < 0.30:
return RuleHit(
rule_id=self.rule_id,
phase=WyckoffPhase.B.value,
confidence=72.0,
score=68.0,
reasons=["区间中部震荡,疑似 Phase B 建仓/派发"],
)
return None
class PhaseCRule(WyckoffRule):
rule_id = "phase_c"
category = "phase"
timeframes = ("1w", "1d")
def evaluate(self, context: dict[str, Any]) -> RuleHit | None:
cycle = _cycle(context)
pos = _f(context, "range_position")
spring_like = _f(context, "spring_score_hint")
utad_like = _f(context, "utad_score_hint")
if cycle in (WyckoffCycle.ACCUMULATION.value, WyckoffCycle.RE_ACCUMULATION.value):
if pos < 0.25 or spring_like >= 50:
return RuleHit(
rule_id=self.rule_id,
phase=WyckoffPhase.C.value,
confidence=75.0 + min(15.0, spring_like * 0.15),
score=78.0,
reasons=["区间低位测试,疑似 Phase C (Spring/Test)"],
)
if cycle in (WyckoffCycle.DISTRIBUTION.value, WyckoffCycle.RE_DISTRIBUTION.value):
if pos > 0.75 or utad_like >= 50:
return RuleHit(
rule_id=self.rule_id,
phase=WyckoffPhase.C.value,
confidence=75.0,
score=78.0,
reasons=["区间高位测试,疑似 Phase C (UTAD)"],
)
return None
class PhaseDRule(WyckoffRule):
rule_id = "phase_d"
category = "phase"
timeframes = ("1w", "1d")
def evaluate(self, context: dict[str, Any]) -> RuleHit | None:
cycle = _cycle(context)
close = _f(context, "close")
ma20 = _f(context, "ma20")
range_high = _f(context, "range_high")
range_low = _f(context, "range_low")
vol_ratio = _f(context, "volume_ratio")
if cycle in (WyckoffCycle.ACCUMULATION.value, WyckoffCycle.RE_ACCUMULATION.value):
if close > ma20 and range_high > 0 and close >= range_high * 0.98 and vol_ratio >= 1.1:
return RuleHit(
rule_id=self.rule_id,
phase=WyckoffPhase.D.value,
confidence=80.0,
score=82.0,
reasons=["突破区间上沿放量,疑似 Phase D SOS"],
)
if cycle in (WyckoffCycle.DISTRIBUTION.value, WyckoffCycle.RE_DISTRIBUTION.value):
if close < ma20 and range_low > 0 and close <= range_low * 1.02:
return RuleHit(
rule_id=self.rule_id,
phase=WyckoffPhase.D.value,
confidence=80.0,
score=82.0,
reasons=["跌破区间下沿,疑似 Phase D SOW"],
)
return None
class PhaseERule(WyckoffRule):
rule_id = "phase_e"
category = "phase"
timeframes = ("1w", "1d")
def evaluate(self, context: dict[str, Any]) -> RuleHit | None:
cycle = _cycle(context)
# Markup/Markdown already imply trend continuation (Phase E of prior structure)
if cycle == WyckoffCycle.MARKUP.value:
return RuleHit(
rule_id=self.rule_id,
phase=WyckoffPhase.E.value,
confidence=78.0,
score=80.0,
reasons=["趋势上行,对应 Phase E Markup"],
)
if cycle == WyckoffCycle.MARKDOWN.value:
return RuleHit(
rule_id=self.rule_id,
phase=WyckoffPhase.E.value,
confidence=78.0,
score=80.0,
reasons=["趋势下行,对应 Phase E Markdown"],
)
return None
def build_rules() -> list[WyckoffRule]:
# More specific phases first
return [PhaseDRule(), PhaseCRule(), PhaseARule(), PhaseBRule(), PhaseERule()]
+39
View File
@@ -0,0 +1,39 @@
"""Rule Registry — register Wyckoff rules without modifying engines."""
from __future__ import annotations
from crypto_wyckoff.rules.base import WyckoffRule
class RuleRegistry:
def __init__(self) -> None:
self._rules: dict[str, WyckoffRule] = {}
def register(self, rule: WyckoffRule) -> None:
self._rules[rule.rule_id] = rule
def get(self, rule_id: str) -> WyckoffRule | None:
return self._rules.get(rule_id)
def by_category(self, category: str, timeframe: str | None = None) -> list[WyckoffRule]:
out = [r for r in self._rules.values() if r.category == category]
if timeframe:
out = [r for r in out if timeframe in r.timeframes]
return out
def all(self) -> list[WyckoffRule]:
return list(self._rules.values())
rule_registry = RuleRegistry()
def _register_defaults() -> None:
from crypto_wyckoff.rules import cycle_rules, event_rules, phase_rules
for mod in (cycle_rules, phase_rules, event_rules):
for rule in mod.build_rules():
rule_registry.register(rule)
_register_defaults()
+127
View File
@@ -0,0 +1,127 @@
"""Background tip + scan scheduler for crypto wyckoff (all enabled combos)."""
from __future__ import annotations
import logging
import threading
from datetime import datetime, timezone
from crypto_wyckoff.combos import all_tfs_for_combos, list_combos
from crypto_wyckoff.io import (
backfill_symbol,
bar_count,
fetch_symbols_from_provider,
tip_update_symbol,
)
from crypto_wyckoff.pipeline import analyze_and_store
logger = logging.getLogger(__name__)
_thread: threading.Thread | None = None
_stop = threading.Event()
_status: dict = {
"running": False,
"last_tick_at": None,
"last_error": None,
"symbols_total": 0,
"symbols_scanned": 0,
"tick_interval_sec": 60,
"backfill_done": False,
}
_status_lock = threading.Lock()
def _set(**kwargs):
with _status_lock:
_status.update(kwargs)
def get_status() -> dict:
with _status_lock:
return dict(_status)
def run_tick(max_symbols: int | None = None, force_rescan: bool = False) -> dict:
"""One cycle: refresh symbols, tip-update, analyze each combo."""
symbols = fetch_symbols_from_provider()
if max_symbols:
symbols = symbols[:max_symbols]
combos = list_combos()
tfs = all_tfs_for_combos(combos)
_set(symbols_total=len(symbols), running=True, last_error=None)
scanned = 0
errors = 0
changed_n = 0
for i, sym in enumerate(symbols):
try:
# Prefer low-TF of first combo for "enough history" gate
low0 = combos[0]["low"] if combos else "1d"
if bar_count(sym, low0) < 40:
backfill_symbol(sym, tfs)
tip_changed = tip_update_symbol(sym, tfs)
if tip_changed:
changed_n += 1
if force_rescan or tip_changed:
for combo in combos:
row = analyze_and_store(sym, combo_id=combo["id"])
if row:
scanned += 1
except Exception as e:
errors += 1
if errors <= 5:
logger.warning("tick %s: %s", sym, e)
_set(last_error=str(e))
if (i + 1) % 25 == 0:
_set(symbols_scanned=scanned)
logger.info("wyckoff tick progress %s/%s scanned=%s", i + 1, len(symbols), scanned)
_set(
running=False,
symbols_scanned=scanned,
last_tick_at=datetime.now(timezone.utc).isoformat(),
backfill_done=True,
)
return {
"symbols": len(symbols),
"scanned": scanned,
"changed_tips": changed_n,
"errors": errors,
"combos": [c["id"] for c in combos],
"tfs": tfs,
}
def _loop(interval: int, max_symbols: int | None):
try:
run_tick(max_symbols=max_symbols, force_rescan=True)
except Exception as e:
logger.exception("initial tick failed: %s", e)
_set(last_error=str(e), running=False)
while not _stop.wait(interval):
try:
# Tip-driven: only force full rescan when tips change is handled inside
run_tick(max_symbols=max_symbols, force_rescan=False)
except Exception as e:
logger.exception("tick failed: %s", e)
_set(last_error=str(e), running=False)
def start_scheduler(interval_sec: int = 60, max_symbols: int | None = None) -> None:
global _thread
if _thread and _thread.is_alive():
return
_stop.clear()
_set(tick_interval_sec=interval_sec)
_thread = threading.Thread(
target=_loop,
args=(interval_sec, max_symbols),
name="crypto-wyckoff-scheduler",
daemon=True,
)
_thread.start()
logger.info("crypto wyckoff scheduler started interval=%ss", interval_sec)
def stop_scheduler() -> None:
_stop.set()
+35
View File
@@ -0,0 +1,35 @@
"""Signal Engine — timeframe-local status labels only (not tradability)."""
from __future__ import annotations
from crypto_wyckoff.domain_models import EngineResult, WyckoffEvent
class SignalEngine:
"""Maps local Event/Phase into a status label. Decision decides tradability."""
name = "Signal"
version = "1.0.0"
def run(self, event: EngineResult, phase: EngineResult | None = None) -> EngineResult:
current = event.payload.get("current_event", WyckoffEvent.NONE.value)
conf = event.confidence
label = current # status label mirrors event for V1
reasons = [f"本地事件标签: {label}"]
if phase and phase.payload.get("phase"):
reasons.append(f"本地阶段: {phase.payload.get('phase')}")
return EngineResult(
name=self.name,
version=self.version,
confidence=conf,
score=event.score,
reasons=reasons,
payload={
"signal_label": label,
"current_event": current,
"phase": (phase.payload.get("phase") if phase else None),
"active_events": event.payload.get("active_events")
or event.payload.get("recent_events", []),
},
)
+236
View File
@@ -0,0 +1,236 @@
"""SQLite persistence for crypto wyckoff scan rows (per combo)."""
from __future__ import annotations
import sqlite3
from datetime import datetime
from typing import Any
from crypto_wyckoff.domain_models import WyckoffScanRow
from crypto_wyckoff.io import SCAN_DB, ensure_dirs
_COLS = [
"trade_date", "combo_id", "ts_code", "name", "industry", "engine_version",
"m_cycle", "cycle_confidence", "trend_score",
"w_cycle", "w_phase", "w_current_event", "w_recent_events_json",
"phase_confidence", "structure_score",
"d_current_event", "d_recent_events_json", "event_confidence", "entry_score",
"entry", "stop", "target1", "target2", "rr",
"alignment", "stars", "decision_signal", "signal_confidence",
"overall_confidence", "overall_score", "risk", "reasons_json",
"feature_snapshot_json", "markers_json", "scanned_at",
]
_CREATE_SQL = """
CREATE TABLE IF NOT EXISTS wyckoff_scan (
trade_date TEXT NOT NULL,
combo_id TEXT NOT NULL DEFAULT 'd_w_m',
ts_code TEXT NOT NULL,
name TEXT DEFAULT '',
industry TEXT DEFAULT '',
engine_version TEXT,
m_cycle TEXT, cycle_confidence REAL, trend_score REAL,
w_cycle TEXT, w_phase TEXT, w_current_event TEXT, w_recent_events_json TEXT,
phase_confidence REAL, structure_score REAL,
d_current_event TEXT, d_recent_events_json TEXT, event_confidence REAL, entry_score REAL,
entry REAL, stop REAL, target1 REAL, target2 REAL, rr REAL,
alignment REAL, stars INTEGER, decision_signal TEXT, signal_confidence REAL,
overall_confidence REAL, overall_score REAL, risk TEXT, reasons_json TEXT,
feature_snapshot_json TEXT, markers_json TEXT, scanned_at TEXT,
PRIMARY KEY (trade_date, combo_id, ts_code)
)
"""
def _migrate(c: sqlite3.Connection) -> None:
cur = c.execute(
"SELECT name FROM sqlite_master WHERE type='table' AND name='wyckoff_scan'"
)
if not cur.fetchone():
c.execute(_CREATE_SQL)
c.execute(
"CREATE INDEX IF NOT EXISTS idx_cw_score "
"ON wyckoff_scan(trade_date, combo_id, overall_score DESC)"
)
return
cols = {r[1] for r in c.execute("PRAGMA table_info(wyckoff_scan)")}
if "combo_id" in cols:
c.execute(
"CREATE INDEX IF NOT EXISTS idx_cw_score "
"ON wyckoff_scan(trade_date, combo_id, overall_score DESC)"
)
return
# Legacy PK (trade_date, ts_code) → add combo_id via table rebuild
c.execute("ALTER TABLE wyckoff_scan RENAME TO wyckoff_scan_old")
c.execute(_CREATE_SQL)
old_cols = [r[1] for r in c.execute("PRAGMA table_info(wyckoff_scan_old)")]
shared = [col for col in _COLS if col != "combo_id" and col in old_cols]
col_sql = ",".join(shared)
c.execute(
f"""
INSERT INTO wyckoff_scan (combo_id, {col_sql})
SELECT 'd_w_m', {col_sql} FROM wyckoff_scan_old
"""
)
c.execute("DROP TABLE wyckoff_scan_old")
c.execute(
"CREATE INDEX IF NOT EXISTS idx_cw_score "
"ON wyckoff_scan(trade_date, combo_id, overall_score DESC)"
)
def _conn() -> sqlite3.Connection:
ensure_dirs()
c = sqlite3.connect(str(SCAN_DB), timeout=60)
c.row_factory = sqlite3.Row
_migrate(c)
c.commit()
return c
def upsert_row(row: WyckoffScanRow) -> None:
combo_id = getattr(row, "combo_id", None) or "d_w_m"
vals = (
row.trade_date.isoformat() if hasattr(row.trade_date, "isoformat") else str(row.trade_date),
combo_id,
row.ts_code, row.name, row.industry, row.engine_version,
row.m_cycle, row.cycle_confidence, row.trend_score,
row.w_cycle, row.w_phase, row.w_current_event, row.w_recent_events_json,
row.phase_confidence, row.structure_score,
row.d_current_event, row.d_recent_events_json, row.event_confidence, row.entry_score,
row.entry, row.stop, row.target1, row.target2, row.rr,
row.alignment, row.stars, row.decision_signal, row.signal_confidence,
row.overall_confidence, row.overall_score, row.risk, row.reasons_json,
row.feature_snapshot_json, row.markers_json,
row.scanned_at.isoformat() if isinstance(row.scanned_at, datetime) else str(row.scanned_at),
)
c = _conn()
try:
placeholders = ",".join("?" * len(_COLS))
col_sql = ",".join(_COLS)
updates = ",".join(
f"{col}=excluded.{col}"
for col in _COLS
if col not in ("trade_date", "combo_id", "ts_code")
)
c.execute(
f"""
INSERT INTO wyckoff_scan ({col_sql}) VALUES ({placeholders})
ON CONFLICT(trade_date, combo_id, ts_code) DO UPDATE SET {updates}
""",
vals,
)
c.commit()
finally:
c.close()
def latest_trade_date(combo_id: str | None = None) -> str | None:
c = _conn()
try:
if combo_id:
cur = c.execute(
"SELECT MAX(trade_date) FROM wyckoff_scan WHERE combo_id=?",
(combo_id,),
)
else:
cur = c.execute("SELECT MAX(trade_date) FROM wyckoff_scan")
row = cur.fetchone()
return row[0] if row and row[0] else None
finally:
c.close()
def count_for_date(trade_date: str | None = None, combo_id: str | None = None) -> int:
td = trade_date or latest_trade_date(combo_id)
if not td:
return 0
c = _conn()
try:
if combo_id:
cur = c.execute(
"SELECT COUNT(*) FROM wyckoff_scan WHERE trade_date=? AND combo_id=?",
(td, combo_id),
)
else:
cur = c.execute("SELECT COUNT(*) FROM wyckoff_scan WHERE trade_date=?", (td,))
return int(cur.fetchone()[0])
finally:
c.close()
def query_scan(
*,
trade_date: str | None = None,
combo_id: str | None = None,
m_cycle: str | None = None,
w_phase: str | None = None,
d_event: str | None = None,
decision_signal: str | None = None,
min_overall_score: float | None = None,
min_alignment: float | None = None,
sort: str = "overall_score",
limit: int = 100,
offset: int = 0,
) -> list[dict[str, Any]]:
cid = combo_id or "d_w_m"
td = trade_date or latest_trade_date(cid)
if not td:
return []
sort_col = sort if sort in {
"overall_score", "alignment", "entry_score", "trend_score", "structure_score", "stars"
} else "overall_score"
clauses = ["trade_date=?", "combo_id=?"]
args: list[Any] = [td, cid]
if m_cycle:
clauses.append("m_cycle=?")
args.append(m_cycle)
if w_phase:
clauses.append("w_phase=?")
args.append(w_phase)
if d_event:
clauses.append("d_current_event=?")
args.append(d_event)
if decision_signal:
clauses.append("decision_signal=?")
args.append(decision_signal)
if min_overall_score is not None:
clauses.append("overall_score>=?")
args.append(min_overall_score)
if min_alignment is not None:
clauses.append("alignment>=?")
args.append(min_alignment)
where = " AND ".join(clauses)
args.extend([limit, offset])
c = _conn()
try:
cur = c.execute(
f"SELECT * FROM wyckoff_scan WHERE {where} ORDER BY {sort_col} DESC LIMIT ? OFFSET ?",
args,
)
return [dict(r) for r in cur.fetchall()]
finally:
c.close()
def get_symbol(
ts_code: str,
trade_date: str | None = None,
combo_id: str | None = None,
) -> dict[str, Any] | None:
cid = combo_id or "d_w_m"
td = trade_date or latest_trade_date(cid)
if not td:
return None
c = _conn()
try:
cur = c.execute(
"SELECT * FROM wyckoff_scan WHERE trade_date=? AND combo_id=? AND ts_code=?",
(td, cid, ts_code),
)
row = cur.fetchone()
return dict(row) if row else None
finally:
c.close()
+51
View File
@@ -0,0 +1,51 @@
"""Crypto symbol → Chinese display name for screener UI."""
from __future__ import annotations
# Base asset → 中文名(覆盖 provider 当前币对;未知则回退 base)
_BASE_CN: dict[str, str] = {
"BTC": "比特币",
"ETH": "以太坊",
"SOL": "索拉纳",
"XAU": "黄金",
"XAG": "白银",
"SAGA": "Saga",
"CL": "原油",
"ZEC": "大零币",
"XRP": "瑞波币",
"DOGE": "狗狗币",
"BNB": "币安币",
"SUI": "Sui",
"BILL": "Bill",
"BZ": "BZ",
"LAB": "Lab",
"TON": "通联币",
"CRCL": "Circle",
"SNDK": "SNDK",
"1000PEPE": "千倍佩佩",
"PEPE": "佩佩",
"CHIP": "CHIP",
"WIF": "狗帽子",
}
def base_asset(symbol: str) -> str:
"""BTC/USDT:USDT → BTC1000PEPE/USDT:USDT → 1000PEPE."""
s = (symbol or "").strip()
if not s:
return ""
head = s.split(":")[0]
return head.split("/")[0].upper() if "/" in head else head.upper()
def display_name_cn(symbol: str) -> str:
base = base_asset(symbol)
if not base:
return symbol or ""
return _BASE_CN.get(base, base)
def symbol_name_map(symbols: list[str] | None = None) -> dict[str, str]:
if not symbols:
return {f"{k}/USDT:USDT": v for k, v in _BASE_CN.items()}
return {s: display_name_cn(s) for s in symbols}
+4
View File
@@ -0,0 +1,4 @@
"""Wyckoff Screener engine version — bump when rules change."""
WYCKOFF_ENGINE_VERSION = "v1.0.0"
ARCHITECTURE_VERSION = "1.0"
+3 -1
View File
@@ -24,6 +24,7 @@
- ECR-004 ReviewedTR 评分硬化 + VP 少系列 + 阶段/门闩/单测(无币种参数) - ECR-004 ReviewedTR 评分硬化 + VP 少系列 + 阶段/门闩/单测(无币种参数)
- ECR-007 Final Approval / `276481e`Wyckoff Live Structure`live.py`);Confirmed ≠ Liveexecution 仅 confirmed - ECR-007 Final Approval / `276481e`Wyckoff Live Structure`live.py`);Confirmed ≠ Liveexecution 仅 confirmed
- ECR-008 Reviewed:主站 `chart_tv.js``chart_tv_{lifecycle,shell,indicators,chan,overlays,finalize}.js` + 薄门面 - ECR-008 Reviewed:主站 `chart_tv.js``chart_tv_{lifecycle,shell,indicators,chan,overlays,finalize}.js` + 薄门面
- ECR-009 Implementing`/wyckoff_crypto` 独立选股页(`crypto_wyckoff/`);D/W + 本地月线;60s tip
- 威科夫数据随主 analyze 默认返回;UI 开关仅显隐叠层 - 威科夫数据随主 analyze 默认返回;UI 开关仅显隐叠层
- Live 观察:主图左下角 Cycle Summary(「形成中」= FORMING);无单独 Live 图层 - Live 观察:主图左下角 Cycle Summary(「形成中」= FORMING);无单独 Live 图层
@@ -31,7 +32,7 @@
- `/api/analyze` 字段可增不可删 - `/api/analyze` 字段可增不可删
- 无 ADR 不改笔/段/中枢/买卖点语义 - 无 ADR 不改笔/段/中枢/买卖点语义
- 威科夫为独立叠层(ECR-003/007);勿借机改缠论算法 - 威科夫为独立叠层(ECR-003/007);Crypto Screener 为独立页(ECR-009),勿混进缠论引擎
- Live candidate **不得**进入 execution;交易 L2+ → RISK_REVIEW + EXPLive 须 Human - Live candidate **不得**进入 execution;交易 L2+ → RISK_REVIEW + EXPLive 须 Human
## 已知债务 ## 已知债务
@@ -42,3 +43,4 @@
- 威科夫启发式参数未做 UI 调参 - 威科夫启发式参数未做 UI 调参
- ECR-007 待 Human 在 Gitea 开 PR 合入 `dev` - ECR-007 待 Human 在 Gitea 开 PR 合入 `dev`
- `chart_tv_overlays.js` 仍偏大,可后续再拆 - `chart_tv_overlays.js` 仍偏大,可后续再拆
- ECR-009:月线历史受日线深度限制;Cycle 规则在 crypto 上可能偏 Unknown,看效果再调参
+6
View File
@@ -2,6 +2,12 @@
## Unreleased — 2026-08-07 ## Unreleased — 2026-08-07
### ECR-009L2,进行中)
- 独立页 `/wyckoff_crypto`:移植 A_Share_DP D/W/M 威科夫选股引擎至数字货币
- 本地 `data/crypto_wyckoff/`60s tip;月线由日线 UTC 自然月聚合(provider 无 1M
- API`/api/wyckoff_crypto/*`;不碰主站 analyze / 缠论叠层
### ECR-008L3Reviewed ### ECR-008L3Reviewed
- 主站 `chart_tv.js` 拆为 lifecycle / shell / indicators / chan / overlays / finalize + 薄门面 - 主站 `chart_tv.js` 拆为 lifecycle / shell / indicators / chan / overlays / finalize + 薄门面
@@ -0,0 +1,25 @@
# ECR-009
**Title:** Crypto Wyckoff Screener 独立页(D/W/M
**Status:** Implementing
**Date:** 2026-08-07
**Change Level:** L2
## Change
新增 `crypto_wyckoff/` 包(移植 A_Share_DP 引擎)+ `/wyckoff_crypto` 页 + `/api/wyckoff_crypto/*`;本地缓存 K 线;60s tip 更新。
周期组合:内置 `8h/4h/1h`(默认)与 `1d/1w/1M`;UI 下拉切换;可添加自定义高/中/低组合(规则引擎仍按 D/W/M 角色映射)。
## Forbidden
- 改缠论算法、主站叠层、`/api/analyze``config/`/`strategies/`
- 自动下单
## Acceptance
- [ ] 页面可列出扫描结果(decision/cycle/phase/event
- [ ] 本地 `data/crypto_wyckoff/` 有 K 线与 scan
- [ ] 调度可跑 tip 更新
- [ ] Decision 门闩单测通过
- [ ] 下拉可选 `8h/4h/1h`,可添加新组合
@@ -0,0 +1,31 @@
# ENGINEERING_SPEC — ECR-009 Crypto Wyckoff Screener
**Level:** L2 · 独立页
**Date:** 2026-08-07
## Goal
数字货币 D/W/M 威科夫选股观察页(A_Share_DP 引擎语义);24/7 tip 每分钟更新。
## Package
`crypto_wyckoff/`features → cycle/phase/event/signal → decision → plan;本地 `data/crypto_wyckoff/`
## API
- `GET /wyckoff_crypto`
- `GET /api/wyckoff_crypto/meta|status|scan`
- `GET /api/wyckoff_crypto/symbol/<symbol>`
- `POST /api/wyckoff_crypto/tick`
## Env
- `CRYPTO_WYCKOFF_DISABLE=1` 关闭调度
- `CRYPTO_WYCKOFF_INTERVAL=60`
- `CRYPTO_WYCKOFF_MAX_SYMBOLS=N` 小样本调试
- `DATA_SERVICE_URL` 默认 provider.jackyu66.com
## Crypto calendar
UTC 连续盘;回填不做 A 股周末放大。
**月线**provider 无 `1M`,由本地日线按 **UTC 自然月** OHLCV 聚合;日/周直接拉 `1d`/`1w`
@@ -0,0 +1,13 @@
# Idea: Crypto Wyckoff Screener(独立页)
## Problem
主站威科夫是图叠层;需要 A_Share_DP 式 D/W/M 多周期选股/决策观察,用于数字货币。
## Hypothesis
独立包 + 独立页,币对来自 DATA_SERVICE,本地缓存 1d/1w/1M,每分钟 tip 更新,不碰缠论主链路。
## Change Level Guess
**L2**(新行为面;不改 strategies
+8 -8
View File
@@ -1,8 +1,8 @@
# STATE # STATE
**owner:** idle **owner:** engineer
**active_ecr:** noneECR-008 ReviewedECR-007 待合入 `dev` **active_ecr:** ECR-009crypto wyckoff screener
**phase:** post-review **phase:** implementing
**system_version:** v1.0.0 **system_version:** v1.0.0
**strategy_version:** unchanged **strategy_version:** unchanged
**updated:** 2026-08-07 **updated:** 2026-08-07
@@ -16,12 +16,12 @@
| ECR-002 | L3 | Done (Reviewed) | runtime 包拆分 | | ECR-002 | L3 | Done (Reviewed) | runtime 包拆分 |
| ECR-003 | L2 | Done (Reviewed) | `081a57a` 主站威科夫 | | ECR-003 | L2 | Done (Reviewed) | `081a57a` 主站威科夫 |
| ECR-004 | L2 | Done (Reviewed) | 威科夫硬化 / VP 减负 | | ECR-004 | L2 | Done (Reviewed) | 威科夫硬化 / VP 减负 |
| ECR-007 | L2 | Done (Final Approval) | Live Structure · `276481e` · 待 Gitea PR → `dev` | | ECR-007 | L2 | Done (Final Approval) | Live Structure · 待合入 `dev` |
| ECR-008 | L3 | Done (Reviewed) | chart_tv 拆分 · 本分支 | | ECR-008 | L3 | Done (Reviewed) | chart_tv 拆分 |
| ECR-009 | L2 | Implementing | `/wyckoff_crypto` · D/W/M |
## Notes ## Notes
- ECR-007**FINAL_APPROVAL** · 已 push;开 PRhttps://git.jackyu66.com/jack/Chan/pulls/new/feature/ECR-007-wyckoff-live-structure base `dev` - ECR-009:打开 http://localhost:8128/wyckoff_crypto ;默认组合 `8h/4h/1h`,可下拉切 `1d/1w/1M` 或「添加组合」
- ECR-008**Approve** · `node --check` 绿;请硬刷新 `?v=20260807f` 目测 - 可用 `CRYPTO_WYCKOFF_MAX_SYMBOLS` 限流;月线仍由日线 UTC 聚合
- 归档:`docs/runs/LOOP-RUN-005/`
- 未请求新 system tag - 未请求新 system tag
+5
View File
@@ -0,0 +1,5 @@
ecr: ECR-009
owner: engineer
phase: implementing
updated: 2026-08-07
notes: crypto wyckoff screener · D/W/M · 24/7 tip
+25
View File
@@ -0,0 +1,25 @@
# TEST_REPORT — ECR-009
**Date:** 2026-08-07
## Commands
```bash
PYTHONPATH=. python -m pytest tests/test_crypto_wyckoff_decision.py -q
CRYPTO_WYCKOFF_DISABLE=1 PYTHONPATH=.:web python -m pytest web/tests/test_wyckoff_crypto_routes.py -q
# Manual / live:
# cd web && CRYPTO_WYCKOFF_MAX_SYMBOLS=5 PYTHONPATH=..:. python app.py
# curl -I http://127.0.0.1:8128/wyckoff_crypto
```
## Result
| Check | Result |
|-------|--------|
| Decision gate unit | 2 passed |
| Route page/meta/scan | 补测(本文件) |
| Live HTTP 2026-08-07 | `GET /wyckoff_crypto` → 200(需先启动 web |
## Note
此前冒烟只做了引擎 tick,**未**在交付前保持 Flask 常驻并给浏览器 URL——属 ESS 测试缺口,已补路由测试与本报告。
+8
View File
@@ -61,3 +61,11 @@
| ECR-008 | 拆分 chart_tv 单体 | ENG-008 | `chart_tv_*.js` + 薄门面 | `node --check` | dbb6202 | | ECR-008 | 拆分 chart_tv 单体 | ENG-008 | `chart_tv_*.js` + 薄门面 | `node --check` | dbb6202 |
| ECR-008 | 对外 API 不变 | ENG-008 | `initTradingView` / `disposeTradingViewCharts` | ui.js 调用点 | dbb6202 | | ECR-008 | 对外 API 不变 | ENG-008 | `initTradingView` / `disposeTradingViewCharts` | ui.js 调用点 | dbb6202 |
| ECR-008 | 无打包器 | PROFILE | `index.html` script 顺序 | 人工 | dbb6202 | | ECR-008 | 无打包器 | PROFILE | `index.html` script 顺序 | 人工 | dbb6202 |
## ECR-009
| ECR | Requirement | Spec | Code | Test | Commit |
|-----|-------------|------|------|------|--------|
| ECR-009 | Crypto D/W/M screener 独立页 | ENG-009 | `crypto_wyckoff/` + `/wyckoff_crypto` | `test_crypto_wyckoff_decision` | ec08de0 |
| ECR-009 | 月线本地聚合 | ENG-009 | `io.rebuild_monthly_from_daily` | smoke tip | ec08de0 |
| ECR-009 | 不碰 analyze/缠论 | ECR-009 Forbidden | 新 API 前缀 | 人工 | ec08de0 |
+31
View File
@@ -0,0 +1,31 @@
"""Unit tests for TF combo validation."""
from __future__ import annotations
import pytest
from crypto_wyckoff.combos import (
add_combo,
delete_combo,
get_combo,
list_combos,
validate_combo,
)
def test_builtin_default_is_h8_4_1():
c = get_combo(None)
assert c["id"] == "h8_4_1"
assert (c["high"], c["mid"], c["low"]) == ("8h", "4h", "1h")
def test_validate_order():
assert validate_combo("8h", "4h", "1h") is None
assert validate_combo("1h", "4h", "8h") is not None
assert validate_combo("8h", "8h", "1h") is not None
def test_list_includes_dwm():
ids = {c["id"] for c in list_combos()}
assert "h8_4_1" in ids
assert "d_w_m" in ids
+66
View File
@@ -0,0 +1,66 @@
"""Decision engine MTF gate tests (ported semantics)."""
from crypto_wyckoff.domain_models import (
DecisionSignal,
EngineResult,
WyckoffCycle,
WyckoffEvent,
WyckoffPhase,
)
from crypto_wyckoff.decision import DecisionEngine
def _er(name, payload, score=70, confidence=70):
return EngineResult(name=name, score=score, confidence=confidence, payload=payload)
def test_monthly_distribution_daily_spring_is_watch():
eng = DecisionEngine()
monthly = _er("Cycle", {"cycle": WyckoffCycle.DISTRIBUTION.value, "trend_score": 40}, score=40)
weekly_c = _er("Cycle", {"cycle": WyckoffCycle.ACCUMULATION.value, "trend_score": 70}, score=70)
weekly_p = _er(
"Phase",
{"phase": WyckoffPhase.B.value, "cycle": WyckoffCycle.ACCUMULATION.value, "structure_score": 65},
score=65,
)
weekly_e = _er("Event", {"current_event": WyckoffEvent.ST.value, "recent_events": ["SC", "AR", "ST"]}, score=60)
daily_e = _er(
"Event",
{"current_event": WyckoffEvent.SPRING.value, "recent_events": ["SC", "AR", "ST", "Spring"], "entry_score": 92},
score=92,
confidence=92,
)
daily_s = _er("Signal", {"signal_label": "Spring", "current_event": "Spring"}, confidence=92, score=92)
out = eng.run(monthly, weekly_c, weekly_p, weekly_e, daily_e, daily_s)
assert out.payload["decision_signal"] == DecisionSignal.WATCH.value
assert out.payload["d_event"] == WyckoffEvent.SPRING.value
def test_bull_alignment_can_strong_buy():
eng = DecisionEngine()
monthly = _er("Cycle", {"cycle": WyckoffCycle.MARKUP.value, "trend_score": 90}, score=90, confidence=90)
weekly_c = _er("Cycle", {"cycle": WyckoffCycle.ACCUMULATION.value, "trend_score": 85}, score=85, confidence=85)
weekly_p = _er(
"Phase",
{"phase": WyckoffPhase.D.value, "cycle": WyckoffCycle.ACCUMULATION.value, "structure_score": 88},
score=88,
confidence=88,
)
weekly_e = _er("Event", {"current_event": WyckoffEvent.SOS.value, "recent_events": ["SOS"]}, score=85, confidence=85)
daily_e = _er(
"Event",
{
"current_event": WyckoffEvent.SPRING.value,
"recent_events": ["SC", "AR", "ST", "Spring", "Test"],
"active_events": ["SC", "AR", "ST", "Spring"],
"entry_score": 92,
},
score=92,
confidence=92,
)
daily_s = _er("Signal", {"signal_label": "Spring"}, confidence=92, score=92)
out = eng.run(monthly, weekly_c, weekly_p, weekly_e, daily_e, daily_s)
assert out.payload["decision_signal"] in (
DecisionSignal.STRONG_BUY.value,
DecisionSignal.BUY.value,
)
+77
View File
@@ -2,6 +2,9 @@
from flask import Blueprint, jsonify, request from flask import Blueprint, jsonify, request
from services.runtime import * # noqa: F403 from services.runtime import * # noqa: F403
from services import runtime as R from services import runtime as R
# import * 不会带出下划线私有名;结构区缓存需显式导入
from services.runtime.state import _zone_cache
from services.runtime.timeframes import _zone_cache_ttl
bp = Blueprint("analyze", __name__) bp = Blueprint("analyze", __name__)
@@ -785,3 +788,77 @@ def analyze():
return jsonify(result) return jsonify(result)
def _serialize_kl_tail(df, limit: int):
"""只序列化最近 limit 根,供自动刷新增量合并。"""
if df is None or getattr(df, "empty", True):
return []
tail = df.tail(limit)
clean = clean_dataframe_for_json(tail)
records = clean.to_dict("records")
for row in records:
d = row.get("date")
if hasattr(d, "isoformat"):
try:
row["date"] = d.isoformat()
except Exception:
row["date"] = str(d)
# timestamp 统一成 int ms,便于前端按 key 合并
ts = row.get("timestamp")
if ts is not None:
try:
row["timestamp"] = int(ts)
except (TypeError, ValueError):
pass
elif hasattr(d, "timestamp"):
try:
row["timestamp"] = int(d.timestamp() * 1000)
except Exception:
pass
return records
@bp.route("/api/klines/recent")
def klines_recent():
"""轻量拉取最近 N 根 K 线(不做缠论/威科夫),供主站自动刷新增量。"""
symbol = (request.args.get("symbol") or "").strip()
if not symbol:
return jsonify({"error": "交易对不能为空"}), 400
timeframe = request.args.get("timeframe", "5m")
try:
limit = int(request.args.get("limit", 2))
except (TypeError, ValueError):
limit = 2
limit = max(1, min(limit, 20))
element_timeframe = request.args.get("element_timeframe") or None
sub_sub_timeframe = request.args.get("sub_sub_timeframe") or None
# 只取尾部:不传 start/end,避免全量窗口回拉
df = get_kl_data(symbol, timeframe, limit=limit)
if df is None:
return jsonify({"error": "获取数据失败"}), 502
if len(df) == 0:
return jsonify({"error": "没有数据"}), 404
result = {
"partial": True,
"symbol": symbol,
"timeframe": timeframe,
"limit": limit,
"kline_data": _serialize_kl_tail(df, limit),
}
if element_timeframe:
edf = get_kl_data(symbol, element_timeframe, limit=limit)
result["element_timeframe"] = element_timeframe
result["element_kline_data"] = _serialize_kl_tail(edf, limit) if edf is not None else []
if sub_sub_timeframe:
sdf = get_kl_data(symbol, sub_sub_timeframe, limit=limit)
result["sub_sub_timeframe"] = sub_sub_timeframe
result["sub_sub_kline_data"] = _serialize_kl_tail(sdf, limit) if sdf is not None else []
return jsonify(result)
+236
View File
@@ -0,0 +1,236 @@
"""Crypto Wyckoff Screener API + page (independent of /api/analyze)."""
from __future__ import annotations
import os
import threading
from flask import Blueprint, jsonify, render_template, request
from crypto_wyckoff.combos import (
ALLOWED_TFS,
add_combo,
delete_combo,
get_combo,
list_combos,
)
from crypto_wyckoff.domain_models import DecisionSignal, WyckoffCycle, WyckoffEvent, WyckoffPhase
from crypto_wyckoff.scheduler import get_status, run_tick, start_scheduler
from crypto_wyckoff import store as wyckoff_store
from crypto_wyckoff.symbols_cn import display_name_cn, symbol_name_map
from crypto_wyckoff.version import ARCHITECTURE_VERSION, WYCKOFF_ENGINE_VERSION
bp = Blueprint("wyckoff_crypto", __name__)
_scheduler_started = False
_sched_lock = threading.Lock()
def ensure_scheduler() -> None:
global _scheduler_started
with _sched_lock:
if _scheduler_started:
return
if os.environ.get("CRYPTO_WYCKOFF_DISABLE", "").lower() in ("1", "true", "yes"):
return
interval = int(os.environ.get("CRYPTO_WYCKOFF_INTERVAL", "60"))
max_sym = os.environ.get("CRYPTO_WYCKOFF_MAX_SYMBOLS")
max_symbols = int(max_sym) if max_sym else None
start_scheduler(interval_sec=interval, max_symbols=max_symbols)
_scheduler_started = True
def _safe_int(raw, default: int, *, lo: int | None = None, hi: int | None = None) -> int:
try:
v = int(raw)
except (TypeError, ValueError):
v = default
if lo is not None:
v = max(lo, v)
if hi is not None:
v = min(hi, v)
return v
@bp.route("/wyckoff_crypto")
def page():
ensure_scheduler()
return render_template("wyckoff_crypto.html")
@bp.route("/api/wyckoff_crypto/meta")
def meta():
ensure_scheduler()
combo_id = request.args.get("combo_id")
combo = get_combo(combo_id)
latest = wyckoff_store.latest_trade_date(combo["id"])
return jsonify(
{
"architecture_version": ARCHITECTURE_VERSION,
"engine_version": WYCKOFF_ENGINE_VERSION,
"latest_trade_date": latest,
"scan_count": wyckoff_store.count_for_date(latest, combo["id"]),
"cycles": [c.value for c in WyckoffCycle],
"phases": [p.value for p in WyckoffPhase],
"events": [e.value for e in WyckoffEvent],
"decision_signals": [s.value for s in DecisionSignal],
"timezone": "Asia/Shanghai",
"utc_offset": "+08:00",
"timeframes": [combo["low"], combo["mid"], combo["high"]],
"combo": combo,
"combos": list_combos(),
"allowed_tfs": list(ALLOWED_TFS),
"symbol_names": symbol_name_map(),
"default_symbol": "BTC/USDT:USDT",
"status": get_status(),
}
)
@bp.route("/api/wyckoff_crypto/combos", methods=["GET"])
def combos_list():
ensure_scheduler()
return jsonify({"combos": list_combos(), "allowed_tfs": list(ALLOWED_TFS)})
@bp.route("/api/wyckoff_crypto/combos", methods=["POST"])
def combos_add():
ensure_scheduler()
body = request.get_json(silent=True) or {}
high = (body.get("high") or request.args.get("high") or "").strip()
mid = (body.get("mid") or request.args.get("mid") or "").strip()
low = (body.get("low") or request.args.get("low") or "").strip()
label = (body.get("label") or request.args.get("label") or "").strip() or None
try:
row = add_combo(high, mid, low, label=label)
except ValueError as e:
return jsonify({"error": str(e)}), 400
return jsonify({"ok": True, "combo": row, "combos": list_combos()})
@bp.route("/api/wyckoff_crypto/combos/<combo_id>", methods=["DELETE"])
def combos_delete(combo_id: str):
ensure_scheduler()
try:
removed = delete_combo(combo_id)
except ValueError as e:
return jsonify({"error": str(e)}), 400
if not removed:
return jsonify({"error": "not_found"}), 404
return jsonify({"ok": True, "combos": list_combos()})
@bp.route("/api/wyckoff_crypto/status")
def status():
ensure_scheduler()
return jsonify(get_status())
@bp.route("/api/wyckoff_crypto/scan")
def scan():
ensure_scheduler()
combo = get_combo(request.args.get("combo_id"))
rows = wyckoff_store.query_scan(
trade_date=request.args.get("trade_date"),
combo_id=combo["id"],
m_cycle=request.args.get("m_cycle"),
w_phase=request.args.get("w_phase"),
d_event=request.args.get("d_event"),
decision_signal=request.args.get("decision_signal"),
min_overall_score=_float_or_none(request.args.get("min_overall_score")),
min_alignment=_float_or_none(request.args.get("min_alignment")),
sort=request.args.get("sort") or "overall_score",
limit=_safe_int(request.args.get("limit"), 100, lo=1, hi=500),
offset=_safe_int(request.args.get("offset"), 0, lo=0),
)
for row in rows:
row["name"] = display_name_cn(row.get("ts_code") or "")
return jsonify({"rows": rows, "count": len(rows), "combo": combo})
@bp.route("/api/wyckoff_crypto/symbol/<path:symbol>")
def symbol_detail(symbol: str):
ensure_scheduler()
combo = get_combo(request.args.get("combo_id"))
row = wyckoff_store.get_symbol(symbol, request.args.get("trade_date"), combo["id"])
if not row:
return jsonify({"error": "not_found"}), 404
return jsonify(row)
@bp.route("/api/wyckoff_crypto/tick", methods=["POST"])
def manual_tick():
"""Manual one-shot tick (debug). Optional JSON/query max_symbols."""
ensure_scheduler()
body = request.get_json(silent=True) or {}
max_sym = request.args.get("max_symbols") or body.get("max_symbols")
max_symbols = int(max_sym) if max_sym not in (None, "") else None
def _job():
try:
run_tick(max_symbols=max_symbols, force_rescan=True)
except Exception:
pass
threading.Thread(target=_job, daemon=True).start()
return jsonify({"ok": True, "started": True})
@bp.route("/api/wyckoff_crypto/klines")
def klines():
"""Local cached OHLCV for chart (combo TFs)."""
ensure_scheduler()
from crypto_wyckoff.io import is_intraday_tf, load_bars_with_ts
symbol = request.args.get("symbol") or ""
combo = get_combo(request.args.get("combo_id"))
allowed = {combo["low"], combo["mid"], combo["high"]}
tf = request.args.get("tf") or combo["low"]
limit = _safe_int(request.args.get("limit"), 180, lo=1, hi=500)
if not symbol or tf not in allowed:
return jsonify({"error": "bad_request", "allowed": sorted(allowed)}), 400
items = load_bars_with_ts(symbol, tf, lookback=limit)
return jsonify({
"items": items,
"symbol": symbol,
"tf": tf,
"count": len(items),
"intraday": is_intraday_tf(tf),
"combo": combo,
})
@bp.route("/api/wyckoff_crypto/overlay")
def overlay():
"""Phase/event overlay for chart."""
ensure_scheduler()
from crypto_wyckoff.annotate import annotate_symbol
symbol = request.args.get("symbol") or ""
combo = get_combo(request.args.get("combo_id"))
allowed = {combo["low"], combo["mid"], combo["high"]}
tf = request.args.get("tf") or combo["low"]
bars = _safe_int(request.args.get("bars"), 180, lo=20, hi=400)
if not symbol or tf not in allowed:
return jsonify({"error": "bad_request", "allowed": sorted(allowed)}), 400
try:
data = annotate_symbol(symbol, freq=tf, lookback=bars, combo_id=combo["id"])
except Exception:
return jsonify({
"error": "overlay_failed",
"phases": [],
"events": [],
"levels": {},
"zones": [],
"combo_id": combo["id"],
}), 500
return jsonify(data)
def _float_or_none(v):
if v in (None, ""):
return None
try:
return float(v)
except (TypeError, ValueError):
return None
+7
View File
@@ -15,6 +15,7 @@ from api.analyze import bp as analyze_bp
from api.pages import bp as pages_bp from api.pages import bp as pages_bp
from api.symbols import bp as symbols_bp from api.symbols import bp as symbols_bp
from api.trend import bp as trend_bp from api.trend import bp as trend_bp
from api.wyckoff_crypto import bp as wyckoff_crypto_bp, ensure_scheduler
def create_app() -> Flask: def create_app() -> Flask:
@@ -23,6 +24,12 @@ def create_app() -> Flask:
app.register_blueprint(analyze_bp) app.register_blueprint(analyze_bp)
app.register_blueprint(symbols_bp) app.register_blueprint(symbols_bp)
app.register_blueprint(trend_bp) app.register_blueprint(trend_bp)
app.register_blueprint(wyckoff_crypto_bp)
# Start crypto wyckoff tip scheduler (daemon); disable with CRYPTO_WYCKOFF_DISABLE=1
try:
ensure_scheduler()
except Exception:
pass
return app return app
+2 -2
View File
@@ -25,8 +25,8 @@ def analyze_chan(df, symbol=None, timeframe=None):
zs_list = chan.calculate_seg_zs(seg_list) zs_list = chan.calculate_seg_zs(seg_list)
# 计算笔中枢(BI中枢)并拍平成列表 # 计算笔中枢(BI中枢)并拍平成列表
#bi_zs_list = chan.cal_bi_zs_list_pure(bi_list) bi_zs_list = chan.cal_bi_zs_list_pure(bi_list)
bi_zs_list = chan.cal_bi_zs(seg_list) #bi_zs_list = chan.cal_bi_zs(seg_list)
bsp_list = [] bsp_list = []
if len(bi_zs_list) > 0: if len(bi_zs_list) > 0:
bsp_list = chan.find_all_bsp(bi_list, bi_zs_list) bsp_list = chan.find_all_bsp(bi_list, bi_zs_list)
+4 -4
View File
@@ -94,19 +94,19 @@ def _prefer_smaller(candidates, labels_ordered, ceiling_tf, timeframe_keys):
def compute_timeframe_defaults(labels_ordered): def compute_timeframe_defaults(labels_ordered):
""" """
根据已排序的周期 中文标签映射计算主 / / 次次周期默认值 根据已排序的周期 中文标签映射计算主 / / 次次周期默认值
默认偏好 4h 2h次次 1h威科夫与结构在小时级更可读 默认偏好 4h 1h次次 15m
labels_ordered: OrderedDict 或按插入顺序排列的 dict labels_ordered: OrderedDict 或按插入顺序排列的 dict
""" """
if not labels_ordered: if not labels_ordered:
labels_ordered = DEFAULT_TIMEFRAME_LABELS.copy() labels_ordered = DEFAULT_TIMEFRAME_LABELS.copy()
timeframe_keys = list(labels_ordered.keys()) timeframe_keys = list(labels_ordered.keys())
preferred_main = next((tf for tf in ['4h', '2h', '1h'] if tf in labels_ordered), None) preferred_main = next((tf for tf in ['4h', '1h', '15m'] if tf in labels_ordered), None)
default_main = preferred_main or (timeframe_keys[0] if timeframe_keys else '1m') default_main = preferred_main or (timeframe_keys[0] if timeframe_keys else '1m')
if default_main not in labels_ordered and timeframe_keys: if default_main not in labels_ordered and timeframe_keys:
default_main = timeframe_keys[0] default_main = timeframe_keys[0]
default_element = _prefer_smaller(['2h', '1h'], labels_ordered, default_main, timeframe_keys) default_element = _prefer_smaller(['1h', '15m'], labels_ordered, default_main, timeframe_keys)
default_sub_sub = _prefer_smaller(['1h'], labels_ordered, default_element, timeframe_keys) default_sub_sub = _prefer_smaller(['15m', '5m'], labels_ordered, default_element, timeframe_keys)
return default_main, default_element, default_sub_sub, timeframe_keys return default_main, default_element, default_sub_sub, timeframe_keys
+106 -19
View File
@@ -9,11 +9,26 @@ function updateTradingViewData() {
return; return;
} }
// 保存当前的可视范围 // 优先用请求前冻结的视窗;否则现场拍(自动刷新短间隔 delta≈0,两种都稳)
const frozen = window._preserveViewOnRefresh;
const oldBarCount = window._preserveViewBarCount || 0;
let savedScrollPosition = null;
if (tvWidget.mainChart) { if (tvWidget.mainChart) {
tvWidget.state.visibleRange = tvWidget.mainChart.timeScale().getVisibleRange(); const ts = tvWidget.mainChart.timeScale();
tvWidget.state.logicalRange = tvWidget.mainChart.timeScale().getVisibleLogicalRange(); if (frozen) {
tvWidget.state.visibleRange = frozen.visibleRange;
tvWidget.state.logicalRange = frozen.logicalRange;
savedScrollPosition = (typeof frozen.scrollPosition === 'number') ? frozen.scrollPosition : null;
} else {
tvWidget.state.visibleRange = ts.getVisibleRange();
tvWidget.state.logicalRange = ts.getVisibleLogicalRange();
try {
savedScrollPosition = ts.scrollPosition ? ts.scrollPosition() : null;
} catch (e) {}
}
} }
window._preserveViewOnRefresh = null;
window._preserveViewBarCount = 0;
// 检查是否显示原始K线 // 检查是否显示原始K线
const showOriginalKline = $('#showOriginalKline').is(':checked'); const showOriginalKline = $('#showOriginalKline').is(':checked');
@@ -71,6 +86,24 @@ function updateTradingViewData() {
}; };
}); });
} }
// LWC 不允许 null/NaN;时间用整秒,避免 Line 渲染抛 Value is null
candles = (candles || []).filter(function (c) {
return c && c.time != null &&
isFinite(Number(c.open)) && isFinite(Number(c.high)) &&
isFinite(Number(c.low)) && isFinite(Number(c.close));
}).map(function (c) {
return {
time: Math.floor(Number(c.time)),
open: Number(c.open),
high: Number(c.high),
low: Number(c.low),
close: Number(c.close)
};
});
const newBarCount = candles.length;
const barDelta = (oldBarCount > 0 && newBarCount > 0) ? (newBarCount - oldBarCount) : 0;
// 更新主系列数据(根据klineType) // 更新主系列数据(根据klineType)
const klineType = ($('#klineType').val() || (showOriginalKline ? 'candlestick' : 'line')); const klineType = ($('#klineType').val() || (showOriginalKline ? 'candlestick' : 'line'));
@@ -270,23 +303,77 @@ function updateTradingViewData() {
// 更新EMA52显示 // 更新EMA52显示
updateEMA52Display(currentData); updateEMA52Display(currentData);
// 恢复之前的可视范围 - 优先使用visibleRange以确保时间轴对齐 // 与自动刷新一致:增量更新绝不碰 barSpacing(缩放本来就留在图表实例上)。
// 一写 barSpacing,LWC 会按右边缘重锚 → 放大往右、缩小往左。
// 这里只在 setData 之后把位置扳回刷新前的 logical / time 窗口。
if (tvWidget.mainChart) { if (tvWidget.mainChart) {
if (tvWidget.state.visibleRange) { const charts = [
console.log('🔄 恢复可见范围:', tvWidget.state.visibleRange); tvWidget.mainChart,
tvWidget.mainChart.timeScale().setVisibleRange(tvWidget.state.visibleRange); tvWidget.volumeChart,
if (tvWidget.volumeChart) tvWidget.volumeChart.timeScale().setVisibleRange(tvWidget.state.visibleRange); tvWidget.atrChart,
if (tvWidget.atrChart) tvWidget.atrChart.timeScale().setVisibleRange(tvWidget.state.visibleRange); tvWidget.macdChart,
if (tvWidget.macdChart) tvWidget.macdChart.timeScale().setVisibleRange(tvWidget.state.visibleRange); tvWidget.chanMacdChart
if (tvWidget.chanMacdChart) tvWidget.chanMacdChart.timeScale().setVisibleRange(tvWidget.state.visibleRange); ].filter(Boolean);
} else if (tvWidget.state.logicalRange) {
console.log('🔄 恢复逻辑范围:', tvWidget.state.logicalRange); const vr = tvWidget.state.visibleRange;
tvWidget.mainChart.timeScale().setVisibleLogicalRange(tvWidget.state.logicalRange); const lr = tvWidget.state.logicalRange;
if (tvWidget.volumeChart) tvWidget.volumeChart.timeScale().setVisibleLogicalRange(tvWidget.state.logicalRange); const savedScroll = savedScrollPosition;
if (tvWidget.atrChart) tvWidget.atrChart.timeScale().setVisibleLogicalRange(tvWidget.state.logicalRange);
if (tvWidget.macdChart) tvWidget.macdChart.timeScale().setVisibleLogicalRange(tvWidget.state.logicalRange); const applyPosition = function (tag) {
if (tvWidget.chanMacdChart) tvWidget.chanMacdChart.timeScale().setVisibleLogicalRange(tvWidget.state.logicalRange); let ok = false;
} if (lr && lr.from !== undefined && lr.to !== undefined && newBarCount > 0) {
// 视窗超出当前 K 线数量时,LWC Line 绘制会抛 Value is null
const span = Math.max(1, lr.to - lr.from);
let to = lr.to;
let from = lr.from;
const maxTo = newBarCount - 1 + 8;
if (to > maxTo) {
to = maxTo;
from = to - span;
}
if (from < -8) {
from = -8;
to = from + span;
}
const clamped = { from: from, to: to };
charts.forEach(c => {
try {
c.timeScale().setVisibleLogicalRange(clamped);
ok = true;
} catch (e) {}
});
if (ok) console.log('🔄 恢复位置 logical' + (tag || '') + ':', clamped);
}
if (!ok && vr && vr.from !== undefined && vr.to !== undefined) {
charts.forEach(c => {
try {
c.timeScale().setVisibleRange(vr);
ok = true;
} catch (e) {}
});
if (ok) console.log('🔄 恢复位置 time' + (tag || '') + ':', vr);
}
if (!ok && typeof savedScroll === 'number') {
const pos = savedScroll + (barDelta || 0);
charts.forEach(c => {
try { c.timeScale().scrollToPosition(pos, false); } catch (e) {}
});
console.log('🔄 恢复位置 scroll' + (tag || '') + ':', pos);
}
};
applyPosition('');
setTimeout(function () { applyPosition('@0'); }, 0);
setTimeout(function () { applyPosition('@50'); }, 50);
// 增量 setData 常不触发可见时间范围回调,但价格轴会变:补刷分型竖边
var bumpFxVert = function () {
if (typeof window._redrawFxBoxVerticalOverlay === 'function') {
window._redrawFxBoxVerticalOverlay();
}
};
bumpFxVert();
setTimeout(bumpFxVert, 0);
setTimeout(bumpFxVert, 50);
} }
console.log('增量更新图表完成'); console.log('增量更新图表完成');
+25 -14
View File
@@ -24,15 +24,13 @@ function chartTvFinalize(ctx) {
var chanMacdChart = ctx.chanMacdChart; var chanMacdChart = ctx.chanMacdChart;
var createChartOptions = ctx.createChartOptions; var createChartOptions = ctx.createChartOptions;
// 同步所有图表的时间轴配置 // 同步所有图表的时间轴配置
const hasPendingRestoreView = !!window._pendingRestoreView;
const pendingView = window._pendingRestoreView;
const syncTimeScaleSettings = () => { const syncTimeScaleSettings = () => {
// 获取主图表的时间轴设置
const mainTimeScale = mainChart.timeScale();
const baseOptions = { const baseOptions = {
timeVisible: true, timeVisible: true,
secondsVisible: false, secondsVisible: false,
borderColor: '#ddd', borderColor: '#ddd',
barSpacing: symbolConfig.type === 'a_stock' ? 6 : 10,
rightOffset: 12,
lockVisibleTimeRangeOnResize: true, lockVisibleTimeRangeOnResize: true,
// 关键:确保所有图表边缘行为完全一致 // 关键:确保所有图表边缘行为完全一致
fixLeftEdge: false, fixLeftEdge: false,
@@ -41,6 +39,12 @@ function chartTvFinalize(ctx) {
ticksVisible: true, ticksVisible: true,
minimumHeight: 0, minimumHeight: 0,
}; };
// 有待恢复视图时不要先写 barSpacing/rightOffset(会钉右缘导致图往右偏),
// 交给后面 setVisibleRange 一次锁定位置+缩放。
if (!pendingView) {
baseOptions.barSpacing = symbolConfig.type === 'a_stock' ? 6 : 10;
baseOptions.rightOffset = 12;
}
console.log('🔧 同步时间轴设置:', baseOptions); console.log('🔧 同步时间轴设置:', baseOptions);
@@ -59,8 +63,12 @@ function chartTvFinalize(ctx) {
// 仅在没有待恢复视图时,设置默认可见范围 // 仅在没有待恢复视图时,设置默认可见范围
const totalBars = candles ? candles.length : 0; const totalBars = candles ? candles.length : 0;
const visibleBarsCount = 200; const visibleBarsCount = 200;
const hasPendingRestoreView = !!window._pendingRestoreView; const allChartsNow = [mainChart, volumeChart, atrChart]
if (!hasPendingRestoreView) { .concat(showMacd && macdChart ? [macdChart] : [])
.concat(showMacd && chanMacdChart ? [chanMacdChart] : []);
if (hasPendingRestoreView && pendingView) {
restoreChartViewState(allChartsNow, pendingView, { preferTime: true });
} else {
// 显示最近 200 根K线而非全部挤压(避免K线过多时重叠) // 显示最近 200 根K线而非全部挤压(避免K线过多时重叠)
if (totalBars > visibleBarsCount) { if (totalBars > visibleBarsCount) {
const rangeFrom = totalBars - visibleBarsCount; const rangeFrom = totalBars - visibleBarsCount;
@@ -71,8 +79,12 @@ function chartTvFinalize(ctx) {
} }
} }
// 立即同步其他图表到主图表的范围 // 立即同步其他图表到主图表的范围(无 pending 时)
setTimeout(() => { setTimeout(() => {
if (window._pendingRestoreView) {
restoreChartViewState(allChartsNow, window._pendingRestoreView, { preferTime: true });
return;
}
const logRange = mainChart.timeScale().getVisibleLogicalRange(); const logRange = mainChart.timeScale().getVisibleLogicalRange();
if (logRange) { if (logRange) {
console.log('🔧 同步可见范围:', logRange); console.log('🔧 同步可见范围:', logRange);
@@ -120,11 +132,10 @@ function chartTvFinalize(ctx) {
} }
const defaultMAs = [ const defaultMAs = [
{ type: 'EMA', length: 13, color: '#800080', name: 'EMA13', visible: true }, // { type: 'EMA', length: 26, color: '#FF8C00', name: 'EMA26', visible: false }, //
{ type: 'EMA', length: 26, color: '#FF8C00', name: 'EMA26', visible: true }, // 橙色 { type: 'EMA', length: 52, color: '#000000', name: 'EMA52', visible: true }, // 黑色 · 默认开
{ type: 'EMA', length: 52, color: '#000000', name: 'EMA52', visible: false }, // 黑色 { type: 'SMA', length: 30, color: '#1E90FF', name: 'MA30', visible: true }, // 蓝色 · 默认开
{ type: 'EMA', length: 104, color: '#1E90FF', name: 'EMA104', visible: false }, // 蓝色 { type: 'SMA', length: 250, color: '#800080', name: 'MA250', visible: true } // 紫色 · 默认开
{ type: 'EMA', length: 156, color: '#F700FF', name: 'EMA156', visible: false } // 粉色
]; ];
defaultMAs.forEach(ma => { defaultMAs.forEach(ma => {
@@ -191,9 +202,9 @@ function chartTvFinalize(ctx) {
window._pendingRestoreView = null; window._pendingRestoreView = null;
if (pending) { if (pending) {
// 恢复刷新前的缩放和位置(优先可见范围/逻辑范围,最后回退到滚动位置 // 恢复刷新前的缩放和位置(时间范围优先,避免数据滑动后逻辑索引错位
console.log('📌 恢复图表视图:', JSON.stringify(pending)); console.log('📌 恢复图表视图:', JSON.stringify(pending));
restoreChartViewState(allCharts, pending); restoreChartViewState(allCharts, pending, { preferTime: true });
} else { } else {
// 无保存视图,正常同步主图到子图 // 无保存视图,正常同步主图到子图
const visibleRange = mainChart.timeScale().getVisibleRange(); const visibleRange = mainChart.timeScale().getVisibleRange();
+234 -69
View File
@@ -1,6 +1,216 @@
/* chart_tv_overlays.js — structure zones / wyckoff / BSP / FX / bollinger */ /* chart_tv_overlays.js — structure zones / wyckoff / BSP / FX / bollinger */
/** 标记 time 必须落在主 series 的 K 线 time 上,否则 LWC 会抛 Value is null */
function alignMarkersToCandles(markers, candles) {
if (!Array.isArray(markers) || !markers.length) return [];
if (!Array.isArray(candles) || !candles.length) return [];
var times = [];
for (var i = 0; i < candles.length; i++) {
var ct = candles[i] && candles[i].time;
if (ct == null || !isFinite(Number(ct))) continue;
times.push(Math.floor(Number(ct)));
}
if (!times.length) return [];
var set = {};
for (var j = 0; j < times.length; j++) set[times[j]] = true;
var nearest = function (target) {
var best = times[0];
var bestDiff = Math.abs(best - target);
// 两端夹逼:大数据量时比全扫略好
var lo = 0, hi = times.length - 1;
while (lo <= hi) {
var mid = (lo + hi) >> 1;
var t = times[mid];
var d = Math.abs(t - target);
if (d < bestDiff) { best = t; bestDiff = d; }
if (t < target) lo = mid + 1;
else hi = mid - 1;
}
if (lo < times.length) {
var d2 = Math.abs(times[lo] - target);
if (d2 < bestDiff) best = times[lo];
}
if (hi >= 0) {
var d3 = Math.abs(times[hi] - target);
if (d3 < bestDiff) best = times[hi];
}
return best;
};
var out = [];
for (var k = 0; k < markers.length; k++) {
var m = markers[k];
if (!m || m.time == null || !isFinite(Number(m.time))) continue;
var t0 = Math.floor(Number(m.time));
var aligned = set[t0] ? t0 : nearest(t0);
var copy = Object.assign({}, m, { time: aligned });
out.push(copy);
}
return out;
}
function safeOverlayLineSetData(series, points) {
if (!series || typeof series.setData !== 'function' || !Array.isArray(points) || points.length < 2) return;
try {
var a = points[0], b = points[1];
if (!a || !b || a.time == null || b.time == null) return;
var t0 = Math.floor(Number(a.time));
var t1 = Math.floor(Number(b.time));
var v0 = Number(a.value);
var v1 = Number(b.value);
if (!isFinite(t0) || !isFinite(t1) || !isFinite(v0) || !isFinite(v1)) return;
// 竖边不用折线(任意时间差都会斜),改走 canvas
if (t0 === t1) return;
if (t0 > t1) {
series.setData([{ time: t1, value: v1 }, { time: t0, value: v0 }]);
} else {
series.setData([{ time: t0, value: v0 }, { time: t1, value: v1 }]);
}
} catch (e) {
console.warn('叠层线 setData 跳过:', e && e.message ? e.message : e);
}
}
function pushFxBoxVertical(time, lo, hi, color) {
if (!window._fxBoxVerticals) window._fxBoxVerticals = [];
var t = Math.floor(Number(time));
var a = Number(lo), b = Number(hi);
if (!isFinite(t) || !isFinite(a) || !isFinite(b) || a === b) return;
window._fxBoxVerticals.push({
time: t,
lo: Math.min(a, b),
hi: Math.max(a, b),
color: color || '#888'
});
}
function getMainPriceSeries() {
if (!window.tvWidget || !tvWidget.series) return null;
var s = tvWidget.series;
return s.candleSeries || s.klcSeries || s.barSeries || s.heikinSeries || s.renkoSeries ||
s.lineSeries || s.areaSeries || s.baselineSeries || null;
}
function syncFxBoxVerticalOverlay(mainChart, mainChartContainer) {
if (!mainChart || !mainChartContainer) return;
if (typeof window._fxBoxOverlayCleanup === 'function') {
try { window._fxBoxOverlayCleanup(); } catch (e) {}
window._fxBoxOverlayCleanup = null;
}
var canvas = mainChartContainer.querySelector('.fx-box-vert-overlay');
if (!canvas) {
canvas = document.createElement('canvas');
canvas.className = 'fx-box-vert-overlay';
canvas.style.cssText = 'position:absolute;left:0;top:0;width:100%;height:100%;pointer-events:none;z-index:6;';
if (getComputedStyle(mainChartContainer).position === 'static') {
mainChartContainer.style.position = 'relative';
}
mainChartContainer.appendChild(canvas);
}
var lastSig = '';
var watchRaf = null;
var cleaned = false;
var redrawPending = false;
var quant = function (v) {
if (v == null || !isFinite(Number(v))) return 'n';
return String(Math.round(Number(v)));
};
// LWC 4 无 priceScale 订阅:采样坐标变化(含增量 setData 后自动缩放)
var sampleSig = function () {
var boxes = window._fxBoxVerticals || [];
var series = getMainPriceSeries();
if (!series || !boxes.length) return '0';
var ts = mainChart.timeScale();
var a = boxes[0];
var b = boxes[boxes.length - 1];
return [
boxes.length,
quant(ts.timeToCoordinate(a.time)),
quant(series.priceToCoordinate(a.hi)),
quant(series.priceToCoordinate(a.lo)),
quant(ts.timeToCoordinate(b.time)),
quant(series.priceToCoordinate(b.hi)),
quant(series.priceToCoordinate(b.lo))
].join('|');
};
var redraw = function () {
var boxes = window._fxBoxVerticals || [];
var series = getMainPriceSeries();
var rect = mainChartContainer.getBoundingClientRect();
var dpr = window.devicePixelRatio || 1;
canvas.width = Math.max(1, Math.floor(rect.width * dpr));
canvas.height = Math.max(1, Math.floor(rect.height * dpr));
canvas.style.width = rect.width + 'px';
canvas.style.height = rect.height + 'px';
var ctx2 = canvas.getContext('2d');
if (!ctx2) return;
ctx2.setTransform(dpr, 0, 0, dpr, 0, 0);
ctx2.clearRect(0, 0, rect.width, rect.height);
if (!series || !boxes.length) {
lastSig = sampleSig();
return;
}
var ts = mainChart.timeScale();
for (var i = 0; i < boxes.length; i++) {
var box = boxes[i];
var x = ts.timeToCoordinate(box.time);
var y1 = series.priceToCoordinate(box.hi);
var y2 = series.priceToCoordinate(box.lo);
if (x == null || y1 == null || y2 == null) continue;
ctx2.beginPath();
ctx2.strokeStyle = box.color;
ctx2.lineWidth = 1;
ctx2.setLineDash([4, 3]);
ctx2.moveTo(Math.round(x) + 0.5, y1);
ctx2.lineTo(Math.round(x) + 0.5, y2);
ctx2.stroke();
}
ctx2.setLineDash([]);
lastSig = sampleSig();
};
var scheduleRedraw = function () {
if (cleaned || redrawPending) return;
redrawPending = true;
requestAnimationFrame(function () {
redrawPending = false;
if (!cleaned) redraw();
});
};
var watch = function () {
if (cleaned) return;
watchRaf = requestAnimationFrame(watch);
var sig = sampleSig();
if (sig !== lastSig) scheduleRedraw();
};
try { mainChart.timeScale().subscribeVisibleLogicalRangeChange(scheduleRedraw); } catch (e) {}
try { mainChart.timeScale().subscribeVisibleTimeRangeChange(scheduleRedraw); } catch (e) {}
var ro = null;
if (typeof ResizeObserver !== 'undefined') {
ro = new ResizeObserver(scheduleRedraw);
ro.observe(mainChartContainer);
}
window._redrawFxBoxVerticalOverlay = scheduleRedraw;
window._fxBoxOverlayCleanup = function () {
if (cleaned) return;
cleaned = true;
if (watchRaf != null) {
try { cancelAnimationFrame(watchRaf); } catch (e) {}
watchRaf = null;
}
window._redrawFxBoxVerticalOverlay = null;
try { mainChart.timeScale().unsubscribeVisibleLogicalRangeChange(scheduleRedraw); } catch (e) {}
try { mainChart.timeScale().unsubscribeVisibleTimeRangeChange(scheduleRedraw); } catch (e) {}
if (ro) try { ro.disconnect(); } catch (e) {}
try { if (canvas && canvas.parentNode) canvas.parentNode.removeChild(canvas); } catch (e) {}
};
if (!window._tvInitCleanups) window._tvInitCleanups = [];
window._tvInitCleanups.push(window._fxBoxOverlayCleanup);
scheduleRedraw();
setTimeout(scheduleRedraw, 50);
watchRaf = requestAnimationFrame(watch);
}
function chartTvRenderOverlays(ctx) { function chartTvRenderOverlays(ctx) {
window._fxBoxVerticals = [];
var symbol = ctx.symbol; var symbol = ctx.symbol;
var timeframe = ctx.timeframe; var timeframe = ctx.timeframe;
var symbolConfig = ctx.symbolConfig; var symbolConfig = ctx.symbolConfig;
@@ -1642,7 +1852,7 @@ function chartTvRenderOverlays(ctx) {
priceLineVisible: false, priceLineVisible: false,
crosshairMarkerVisible: false, crosshairMarkerVisible: false,
}); });
topSeries.setData([{ time: startTs, value: boxHigh }, { time: endTs, value: boxHigh }]); safeOverlayLineSetData(topSeries, [{ time: startTs, value: boxHigh }, { time: endTs, value: boxHigh }]);
const bottomSeries = mainChart.addLineSeries({ const bottomSeries = mainChart.addLineSeries({
color: boxColor, color: boxColor,
@@ -1652,31 +1862,13 @@ function chartTvRenderOverlays(ctx) {
priceLineVisible: false, priceLineVisible: false,
crosshairMarkerVisible: false, crosshairMarkerVisible: false,
}); });
bottomSeries.setData([{ time: startTs, value: boxLow }, { time: endTs, value: boxLow }]); safeOverlayLineSetData(bottomSeries, [{ time: startTs, value: boxLow }, { time: endTs, value: boxLow }]);
const leftSeries = mainChart.addLineSeries({ pushFxBoxVertical(startTs, boxLow, boxHigh, boxColor);
color: boxColor, pushFxBoxVertical(endTs, boxLow, boxHigh, boxColor);
lineWidth: 1,
lineStyle: 2, // 虚线
lastValueVisible: false,
priceLineVisible: false,
crosshairMarkerVisible: false,
});
// 左边竖线:同一 time 上下两个点(和你已有ZS绘制写法保持一致)
leftSeries.setData([{ time: startTs, value: boxLow }, { time: startTs, value: boxHigh }]);
const rightSeries = mainChart.addLineSeries({
color: boxColor,
lineWidth: 1,
lineStyle: 2, // 虚线
lastValueVisible: false,
priceLineVisible: false,
crosshairMarkerVisible: false,
});
rightSeries.setData([{ time: endTs, value: boxLow }, { time: endTs, value: boxHigh }]);
if (!tvWidget.series.mainKlcFxBoxSeries) tvWidget.series.mainKlcFxBoxSeries = []; if (!tvWidget.series.mainKlcFxBoxSeries) tvWidget.series.mainKlcFxBoxSeries = [];
tvWidget.series.mainKlcFxBoxSeries.push(topSeries, bottomSeries, leftSeries, rightSeries); tvWidget.series.mainKlcFxBoxSeries.push(topSeries, bottomSeries);
} }
} }
@@ -1849,7 +2041,7 @@ function chartTvRenderOverlays(ctx) {
priceLineVisible: false, priceLineVisible: false,
crosshairMarkerVisible: false, crosshairMarkerVisible: false,
}); });
topSeries.setData([{ time: startTs, value: boxHigh }, { time: endTs, value: boxHigh }]); safeOverlayLineSetData(topSeries, [{ time: startTs, value: boxHigh }, { time: endTs, value: boxHigh }]);
const bottomSeries = mainChart.addLineSeries({ const bottomSeries = mainChart.addLineSeries({
color: boxColor, color: boxColor,
@@ -1859,30 +2051,13 @@ function chartTvRenderOverlays(ctx) {
priceLineVisible: false, priceLineVisible: false,
crosshairMarkerVisible: false, crosshairMarkerVisible: false,
}); });
bottomSeries.setData([{ time: startTs, value: boxLow }, { time: endTs, value: boxLow }]); safeOverlayLineSetData(bottomSeries, [{ time: startTs, value: boxLow }, { time: endTs, value: boxLow }]);
const leftSeries = mainChart.addLineSeries({ pushFxBoxVertical(startTs, boxLow, boxHigh, boxColor);
color: boxColor, pushFxBoxVertical(endTs, boxLow, boxHigh, boxColor);
lineWidth: 1,
lineStyle: 2,
lastValueVisible: false,
priceLineVisible: false,
crosshairMarkerVisible: false,
});
leftSeries.setData([{ time: startTs, value: boxLow }, { time: startTs, value: boxHigh }]);
const rightSeries = mainChart.addLineSeries({
color: boxColor,
lineWidth: 1,
lineStyle: 2,
lastValueVisible: false,
priceLineVisible: false,
crosshairMarkerVisible: false,
});
rightSeries.setData([{ time: endTs, value: boxLow }, { time: endTs, value: boxHigh }]);
if (!tvWidget.series.elementKlcFxBoxSeries) tvWidget.series.elementKlcFxBoxSeries = []; if (!tvWidget.series.elementKlcFxBoxSeries) tvWidget.series.elementKlcFxBoxSeries = [];
tvWidget.series.elementKlcFxBoxSeries.push(topSeries, bottomSeries, leftSeries, rightSeries); tvWidget.series.elementKlcFxBoxSeries.push(topSeries, bottomSeries);
} }
} }
@@ -2002,7 +2177,7 @@ function chartTvRenderOverlays(ctx) {
priceLineVisible: false, priceLineVisible: false,
crosshairMarkerVisible: false, crosshairMarkerVisible: false,
}); });
topSeries.setData([{ time: startTs, value: boxHigh }, { time: endTs, value: boxHigh }]); safeOverlayLineSetData(topSeries, [{ time: startTs, value: boxHigh }, { time: endTs, value: boxHigh }]);
const bottomSeries = mainChart.addLineSeries({ const bottomSeries = mainChart.addLineSeries({
color: boxColor, color: boxColor,
@@ -2012,30 +2187,13 @@ function chartTvRenderOverlays(ctx) {
priceLineVisible: false, priceLineVisible: false,
crosshairMarkerVisible: false, crosshairMarkerVisible: false,
}); });
bottomSeries.setData([{ time: startTs, value: boxLow }, { time: endTs, value: boxLow }]); safeOverlayLineSetData(bottomSeries, [{ time: startTs, value: boxLow }, { time: endTs, value: boxLow }]);
const leftSeries = mainChart.addLineSeries({ pushFxBoxVertical(startTs, boxLow, boxHigh, boxColor);
color: boxColor, pushFxBoxVertical(endTs, boxLow, boxHigh, boxColor);
lineWidth: 1,
lineStyle: 2,
lastValueVisible: false,
priceLineVisible: false,
crosshairMarkerVisible: false,
});
leftSeries.setData([{ time: startTs, value: boxLow }, { time: startTs, value: boxHigh }]);
const rightSeries = mainChart.addLineSeries({
color: boxColor,
lineWidth: 1,
lineStyle: 2,
lastValueVisible: false,
priceLineVisible: false,
crosshairMarkerVisible: false,
});
rightSeries.setData([{ time: endTs, value: boxLow }, { time: endTs, value: boxHigh }]);
if (!tvWidget.series.subSubKlcFxBoxSeries) tvWidget.series.subSubKlcFxBoxSeries = []; if (!tvWidget.series.subSubKlcFxBoxSeries) tvWidget.series.subSubKlcFxBoxSeries = [];
tvWidget.series.subSubKlcFxBoxSeries.push(topSeries, bottomSeries, leftSeries, rightSeries); tvWidget.series.subSubKlcFxBoxSeries.push(topSeries, bottomSeries);
} }
} }
} catch (e) { console.error('绘制次次周期KLC分型标记出错:', e); } } catch (e) { console.error('绘制次次周期KLC分型标记出错:', e); }
@@ -2180,7 +2338,7 @@ function chartTvRenderOverlays(ctx) {
else if (klineType === 'klc') targetSeries = tvWidget.series.klcSeries; else if (klineType === 'klc') targetSeries = tvWidget.series.klcSeries;
if (targetSeries) { if (targetSeries) {
try { try {
targetSeries.setMarkers(combinedMarkers); targetSeries.setMarkers(alignMarkersToCandles(combinedMarkers, candles));
} catch (e) { } catch (e) {
console.warn('设置主系列标记失败(可能series已释放):', e); console.warn('设置主系列标记失败(可能series已释放):', e);
} }
@@ -2309,7 +2467,7 @@ function chartTvRenderOverlays(ctx) {
else if (klineType2 === 'klc') targetSeries2 = tvWidget.series.klcSeries; else if (klineType2 === 'klc') targetSeries2 = tvWidget.series.klcSeries;
if (targetSeries2) { if (targetSeries2) {
try { try {
targetSeries2.setMarkers(onlyMainAndU); targetSeries2.setMarkers(alignMarkersToCandles(onlyMainAndU, candles));
} catch (e) { } catch (e) {
console.warn('设置主系列标记失败(可能series已释放):', e); console.warn('设置主系列标记失败(可能series已释放):', e);
} }
@@ -2338,4 +2496,11 @@ function chartTvRenderOverlays(ctx) {
} }
} }
} }
// KLC 分型框竖边:canvas 真竖线(LWC 折线做不到不斜)
try {
syncFxBoxVerticalOverlay(mainChart, mainChartContainer);
} catch (e) {
console.warn('分型竖边 overlay 失败:', e);
}
} }
+4 -4
View File
@@ -38,7 +38,7 @@ function chartTvBuildShell(ctx) {
} }
candles = klineDataSource.map((kline) => { candles = klineDataSource.map((kline) => {
const date = new Date(kline.date); const date = new Date(kline.date);
const timestamp = date.getTime() / 1000; const timestamp = Math.floor(date.getTime() / 1000);
return { return {
time: timestamp, time: timestamp,
open: parseFloat(kline.open), open: parseFloat(kline.open),
@@ -46,7 +46,7 @@ function chartTvBuildShell(ctx) {
low: parseFloat(kline.low), low: parseFloat(kline.low),
close: parseFloat(kline.close), close: parseFloat(kline.close),
}; };
}); }).filter((c) => isFinite(c.time) && isFinite(c.open) && isFinite(c.high) && isFinite(c.low) && isFinite(c.close));
} else { } else {
if (!currentData.kline_data || !Array.isArray(currentData.kline_data)) { if (!currentData.kline_data || !Array.isArray(currentData.kline_data)) {
console.error('主周期K线数据不存在或不是数组:', currentData.kline_data); console.error('主周期K线数据不存在或不是数组:', currentData.kline_data);
@@ -54,7 +54,7 @@ function chartTvBuildShell(ctx) {
} }
candles = currentData.kline_data.map((kline) => { candles = currentData.kline_data.map((kline) => {
const date = new Date(kline.date); const date = new Date(kline.date);
const timestamp = date.getTime() / 1000; const timestamp = Math.floor(date.getTime() / 1000);
return { return {
time: timestamp, time: timestamp,
open: parseFloat(kline.open), open: parseFloat(kline.open),
@@ -62,7 +62,7 @@ function chartTvBuildShell(ctx) {
low: parseFloat(kline.low), low: parseFloat(kline.low),
close: parseFloat(kline.close), close: parseFloat(kline.close),
}; };
}); }).filter((c) => isFinite(c.time) && isFinite(c.open) && isFinite(c.high) && isFinite(c.low) && isFinite(c.close));
} }
// 根据交易对类型过滤数据(仅用于显示优化) // 根据交易对类型过滤数据(仅用于显示优化)
+141 -15
View File
@@ -1,4 +1,46 @@
/* chart_view.js — split from chart.js */ /* chart_view.js — split from chart.js */
/** 用尾部 N 根合并进已有 K 线(同 timestamp 覆盖,更新则追加) */
function mergeKlineTail(existing, incoming) {
if (!Array.isArray(incoming) || !incoming.length) {
return Array.isArray(existing) ? existing : [];
}
if (!Array.isArray(existing) || !existing.length) {
return incoming.slice();
}
const out = existing.slice();
const barTs = (row) => {
if (row && row.timestamp != null && row.timestamp !== '') {
const n = Number(row.timestamp);
if (!Number.isNaN(n)) return n;
}
const t = row && row.date != null ? new Date(row.date).getTime() : NaN;
return Number.isNaN(t) ? null : t;
};
for (let i = 0; i < incoming.length; i++) {
const row = incoming[i];
const ts = barTs(row);
if (ts == null) continue;
let idx = -1;
const scanFrom = Math.max(0, out.length - 8);
for (let j = out.length - 1; j >= scanFrom; j--) {
if (barTs(out[j]) === ts) {
idx = j;
break;
}
}
if (idx >= 0) {
out[idx] = Object.assign({}, out[idx], row);
} else {
const lastTs = barTs(out[out.length - 1]);
if (lastTs == null || ts > lastTs) {
out.push(row);
}
}
}
return out;
}
function updateChart(options) { function updateChart(options) {
options = options || {}; options = options || {};
// 只显示旋转加载图标 // 只显示旋转加载图标
@@ -47,9 +89,88 @@ function updateChart(options) {
if (options.fromAutoRefresh && window._analyzeXhr && window._analyzeXhr.readyState !== 4) { if (options.fromAutoRefresh && window._analyzeXhr && window._analyzeXhr.readyState !== 4) {
try { window._analyzeXhr.abort(); } catch (e) {} try { window._analyzeXhr.abort(); } catch (e) {}
} }
// 请求发出前冻结视窗(与自动刷新同一套;避免等响应时/setData 后 logical 索引漂移)
try {
if (tvWidget && tvWidget.mainChart) {
window._preserveViewOnRefresh = captureChartViewState(tvWidget.mainChart);
const prev = currentData && (
($('#subSubPeriodKline').is(':checked') && currentData.sub_sub_kline_data) ||
($('#elementPeriodKline').is(':checked') && currentData.element_kline_data) ||
currentData.kline_data
);
window._preserveViewBarCount = Array.isArray(prev) ? prev.length : 0;
console.log('📌 刷新前冻结视窗 bars=', window._preserveViewBarCount, window._preserveViewOnRefresh);
}
} catch (e) {
window._preserveViewOnRefresh = null;
window._preserveViewBarCount = 0;
}
const requestId = ++lastRequestId;
const chartsReady = !!(tvWidget && tvWidget.state && tvWidget.state.isInitialized && tvWidget.mainChart);
const hasBaseline = !!(currentData && Array.isArray(currentData.kline_data) && currentData.kline_data.length);
const baselineSymbol = (currentData && currentData.symbol) || window._lastChartSymbol || '';
// 自动刷新常态:只拉最近 2 根;换币对后基线不一致则禁止尾部合并(否则会叠旧缠论)
// fullAnalyze(约每 1 分钟)走全量 analyze 更新缠论
const useRecentTail = !!(
options.fromAutoRefresh &&
!options.fullAnalyze &&
chartsReady &&
hasBaseline &&
baselineSymbol &&
baselineSymbol === symbol
);
if (useRecentTail) {
console.log('自动刷新 → /api/klines/recent limit=2');
window._analyzeXhr = $.ajax({
url: '/api/klines/recent',
data: {
symbol: symbol,
timeframe: timeframe,
limit: 2,
element_timeframe: elementTimeframe || undefined,
sub_sub_timeframe: subSubTimeframe || undefined
},
success: function(partial) {
$('#refreshLoadingSpinner').hide();
if (requestId !== lastRequestId) return;
if (!partial || !Array.isArray(partial.kline_data)) {
console.warn('recent 响应无效,回退全量 analyze');
updateChart({ incremental: true, reason: 'recent-fallback' });
return;
}
currentData.kline_data = mergeKlineTail(currentData.kline_data, partial.kline_data);
if (Array.isArray(partial.element_kline_data)) {
currentData.element_kline_data = mergeKlineTail(
currentData.element_kline_data, partial.element_kline_data
);
if (partial.element_timeframe) {
currentData.element_timeframe = partial.element_timeframe;
}
}
if (Array.isArray(partial.sub_sub_kline_data)) {
currentData.sub_sub_kline_data = mergeKlineTail(
currentData.sub_sub_kline_data, partial.sub_sub_kline_data
);
if (partial.sub_sub_timeframe) {
currentData.sub_sub_timeframe = partial.sub_sub_timeframe;
}
}
refreshChart(currentData, { incremental: true, skipTables: true });
},
error: function(jqXHR, textStatus, errorThrown) {
$('#refreshLoadingSpinner').hide();
if (textStatus === 'abort') return;
console.warn('recent 失败,回退全量 analyze:', errorThrown);
updateChart({ incremental: true, reason: 'recent-error-fallback' });
}
});
return;
}
// 发送请求 // 手动 / 首拉:全量 analyze
const requestId = ++lastRequestId; // 标记本次请求
window._analyzeXhr = $.ajax({ window._analyzeXhr = $.ajax({
url: '/api/analyze', url: '/api/analyze',
data: { data: {
@@ -75,21 +196,31 @@ function updateChart(options) {
} }
// 保存当前数据 // 保存当前数据
const prevSymbol = (currentData && currentData.symbol) || window._lastChartSymbol || '';
if (currentData) { if (currentData) {
// 覆盖前断开旧引用,帮助GC尽快回收 // 覆盖前断开旧引用,帮助GC尽快回收
delete currentData.original_kline_data; delete currentData.original_kline_data;
delete currentData.original_macd; delete currentData.original_macd;
} }
currentData = data; currentData = data;
window._lastChartSymbol = symbol;
window._lastFullAnalyzeAt = Date.now();
if (typeof renderWyckoffCycleSummary === 'function') { if (typeof renderWyckoffCycleSummary === 'function') {
renderWyckoffCycleSummary(); renderWyckoffCycleSummary();
} }
refreshChart(data, { // 有图则增量;笔/段/中枢/结构区只在全量 init 绘制
incremental: options.incremental !== undefined // 换币对 / 手动分析 / 结构区:必须全量重建,否则会残留旧币对叠层
? !!options.incremental const ready = !!(tvWidget && tvWidget.state && tvWidget.state.isInitialized && tvWidget.mainChart);
: !!options.fromAutoRefresh const structureZonesOn = $('#showMainStructureZone').is(':checked');
}); const symbolChanged = !!(prevSymbol && prevSymbol !== symbol);
let wantIncremental = options.incremental !== undefined
? !!options.incremental
: (ready || !!options.fromAutoRefresh);
if (structureZonesOn || options.fullAnalyze || symbolChanged || options.incremental === false) {
wantIncremental = false;
}
refreshChart(data, { incremental: wantIncremental });
}, },
error: function(jqXHR, textStatus, errorThrown) { error: function(jqXHR, textStatus, errorThrown) {
// 隐藏加载图标 // 隐藏加载图标
@@ -121,24 +252,21 @@ function captureChartViewState(chart) {
} }
function restoreChartViewState(charts, viewState) { function restoreChartViewState(charts, viewState) {
// 全量重建备用:先缩放,再位置;不要在位置前写 rightOffset(会右边缘锚定)
if (!viewState || !Array.isArray(charts) || charts.length === 0) return; if (!viewState || !Array.isArray(charts) || charts.length === 0) return;
const validCharts = charts.filter(c => c && c.timeScale); const validCharts = charts.filter(c => c && c.timeScale);
if (validCharts.length === 0) return; if (validCharts.length === 0) return;
validCharts.forEach(c => { validCharts.forEach(c => {
try { try {
const optionsPatch = {}; if (typeof viewState.barSpacing === 'number') {
if (typeof viewState.barSpacing === 'number') optionsPatch.barSpacing = viewState.barSpacing; c.timeScale().applyOptions({ barSpacing: viewState.barSpacing });
if (typeof viewState.rightOffset === 'number') optionsPatch.rightOffset = viewState.rightOffset;
if (Object.keys(optionsPatch).length) {
c.timeScale().applyOptions(optionsPatch);
} }
} catch (e) {} } catch (e) {}
}); });
let restored = false; let restored = false;
// 优先按逻辑范围恢复(对新数据更稳健)
if (viewState.logicalRange && viewState.logicalRange.from !== undefined && viewState.logicalRange.to !== undefined) { if (viewState.logicalRange && viewState.logicalRange.from !== undefined && viewState.logicalRange.to !== undefined) {
validCharts.forEach(c => { validCharts.forEach(c => {
try { try {
@@ -148,7 +276,6 @@ function restoreChartViewState(charts, viewState) {
}); });
} }
// 逻辑范围失败时,回退到时间可见范围
if (!restored && viewState.visibleRange && viewState.visibleRange.from !== undefined && viewState.visibleRange.to !== undefined) { if (!restored && viewState.visibleRange && viewState.visibleRange.from !== undefined && viewState.visibleRange.to !== undefined) {
validCharts.forEach(c => { validCharts.forEach(c => {
try { try {
@@ -158,7 +285,6 @@ function restoreChartViewState(charts, viewState) {
}); });
} }
// 最后回退到滚动位置
if (!restored && typeof viewState.scrollPosition === 'number') { if (!restored && typeof viewState.scrollPosition === 'number') {
validCharts.forEach(c => { validCharts.forEach(c => {
try { c.timeScale().scrollToPosition(viewState.scrollPosition, false); } catch (e) {} try { c.timeScale().scrollToPosition(viewState.scrollPosition, false); } catch (e) {}
+2 -2
View File
@@ -85,9 +85,9 @@ $(document).on('change', '#showMainBiZs', function() {
$(document).on('change', '#showMainStructureZone', function() { $(document).on('change', '#showMainStructureZone', function() {
const on = $('#showMainStructureZone').is(':checked'); const on = $('#showMainStructureZone').is(':checked');
console.log('结构区切换为:', on); console.log('结构区切换为:', on);
// 勾选后才向服务器请求多周期结构区数据;取消勾选仅重绘,不重复拉取 // 勾选后才向服务器请求多周期结构区数据;结构区叠层只在全量 init 里绘制,必须 incremental:false
if (on) { if (on) {
updateChart(); updateChart({ incremental: false });
} else { } else {
updateChartDisplay(); updateChartDisplay();
} }
+47 -21
View File
@@ -272,10 +272,10 @@ function loadSymbols() {
}); });
} }
// 设置默认时间范围(需覆盖威科夫 lookback;1 天在 4h/1h 上几乎检不出区间) // 设置默认时间范围:最近 1 个月
function setDefaultTimeRange() { function setDefaultTimeRange() {
const now = new Date(); const now = new Date();
const daysBack = 14; const daysBack = 30;
const start = new Date(now.getTime() - (daysBack * 24 * 60 * 60 * 1000)); const start = new Date(now.getTime() - (daysBack * 24 * 60 * 60 * 1000));
// 格式化为datetime-local输入框所需的格式 YYYY-MM-DDThh:mm // 格式化为datetime-local输入框所需的格式 YYYY-MM-DDThh:mm
@@ -493,6 +493,9 @@ $(document).ready(function() {
let autoRefreshTimer = null; let autoRefreshTimer = null;
let nextRefreshTime = null; let nextRefreshTime = null;
let autoRefreshTick = 0; let autoRefreshTick = 0;
/** 自动刷新时,缠论全量重算间隔(毫秒);时间戳见 window._lastFullAnalyzeAt */
const AUTO_FULL_ANALYZE_MS = 60 * 1000;
// 初始化自动刷新功能 // 初始化自动刷新功能
function initAutoRefresh() { function initAutoRefresh() {
// 监听自动刷新勾选框变化 // 监听自动刷新勾选框变化
@@ -519,10 +522,10 @@ function startAutoRefresh() {
stopAutoRefresh(); stopAutoRefresh();
// 获取刷新频率(分钟) // 获取刷新频率(分钟)
const interval = parseFloat($('#refreshInterval').val()) || 5; const interval = parseFloat($('#refreshInterval').val()) || (5 / 60);
const intervalMs = interval * 60 * 1000; const intervalMs = interval * 60 * 1000;
console.log(`开始自动刷新,频率: ${interval}分钟 (${intervalMs}毫秒)`); console.log(`开始自动刷新,频率: ${interval}分钟 (${intervalMs}毫秒);缠论全量每 ${AUTO_FULL_ANALYZE_MS / 1000}s`);
// 计算下次刷新时间 // 计算下次刷新时间
nextRefreshTime = new Date(Date.now() + intervalMs); nextRefreshTime = new Date(Date.now() + intervalMs);
@@ -531,16 +534,36 @@ function startAutoRefresh() {
// 启动定时器 // 启动定时器
autoRefreshTick = 0; autoRefreshTick = 0;
autoRefreshTimer = setInterval(function() { autoRefreshTimer = setInterval(function() {
// 更新结束时间为当前时间 // 刷新前先钉住当前缩放/位置(updateEndTime / 请求返回前都可能被改写)
if (tvWidget && tvWidget.mainChart && typeof captureChartViewState === 'function') {
try {
window._pendingRestoreView = captureChartViewState(tvWidget.mainChart);
} catch (e) {
window._pendingRestoreView = null;
}
}
// 更新结束时间显示(仅 UI
updateEndTimeToNow(); updateEndTimeToNow();
// 多数周期增量更新;每隔若干次全量重建以刷新笔/段/中枢(dispose 已防泄漏)
autoRefreshTick += 1; autoRefreshTick += 1;
const fullRebuild = (autoRefreshTick % 6) === 0; const now = Date.now();
updateChart({ const lastFull = window._lastFullAnalyzeAt || 0;
fromAutoRefresh: true, const needFullAnalyze = !lastFull || (now - lastFull >= AUTO_FULL_ANALYZE_MS);
incremental: !fullRebuild // 常态:/api/klines/recent 合并尾部 K;满 1 分钟:全量 /api/analyze 刷新缠论
}); if (needFullAnalyze) {
console.log('自动刷新 → 全量缠论 analyze(距上次', lastFull ? Math.round((now - lastFull) / 1000) + 's' : '首次', '');
updateChart({
fromAutoRefresh: true,
fullAnalyze: true,
incremental: true
});
} else {
updateChart({
fromAutoRefresh: true,
incremental: true
});
}
// 更新下次刷新时间 // 更新下次刷新时间
nextRefreshTime = new Date(Date.now() + intervalMs); nextRefreshTime = new Date(Date.now() + intervalMs);
@@ -784,7 +807,8 @@ function refreshChart(data, options) {
// 自动刷新:增量更新,避免每次销毁/重建 Lightweight Charts // 自动刷新:增量更新,避免每次销毁/重建 Lightweight Charts
if (preferIncremental && chartsReady) { if (preferIncremental && chartsReady) {
try { try {
if (tvWidget.mainChart) { // 若定时器已捕获则保留;否则此刻再捕获一次
if (!window._pendingRestoreView && tvWidget.mainChart) {
try { try {
window._pendingRestoreView = captureChartViewState(tvWidget.mainChart); window._pendingRestoreView = captureChartViewState(tvWidget.mainChart);
} catch (e) { } catch (e) {
@@ -792,7 +816,10 @@ function refreshChart(data, options) {
} }
} }
updateTradingViewData(); updateTradingViewData();
updateTables(data); // recent-tail 刷新结构未变,跳过表格重绘以提速
if (!options.skipTables) {
updateTables(data);
}
if (currentData && currentData.ema52_dict) { if (currentData && currentData.ema52_dict) {
updateEMA52Display(currentData); updateEMA52Display(currentData);
} }
@@ -804,7 +831,7 @@ function refreshChart(data, options) {
// 保存当前缩放(barSpacing)和滚动位置(scrollPosition)到 window // 保存当前缩放(barSpacing)和滚动位置(scrollPosition)到 window
// tvWidget 会在 initTradingView 内被重建,所以必须存到 window 上 // tvWidget 会在 initTradingView 内被重建,所以必须存到 window 上
if (tvWidget && tvWidget.mainChart) { if (!window._pendingRestoreView && tvWidget && tvWidget.mainChart) {
try { try {
window._pendingRestoreView = captureChartViewState(tvWidget.mainChart); window._pendingRestoreView = captureChartViewState(tvWidget.mainChart);
console.log('📌 保存图表视图:', JSON.stringify(window._pendingRestoreView)); console.log('📌 保存图表视图:', JSON.stringify(window._pendingRestoreView));
@@ -812,6 +839,8 @@ function refreshChart(data, options) {
console.warn('保存图表视图失败:', e); console.warn('保存图表视图失败:', e);
window._pendingRestoreView = null; window._pendingRestoreView = null;
} }
} else if (window._pendingRestoreView) {
console.log('📌 使用已保存图表视图:', JSON.stringify(window._pendingRestoreView));
} }
initTradingView($('#symbol').val(), $('#timeframe').val()); initTradingView($('#symbol').val(), $('#timeframe').val());
@@ -845,14 +874,14 @@ $('#showElementMacdDiv').change(function() {
refreshChartOnly(); refreshChartOnly();
}); });
// 绑定分型类型显示开关 // 绑定分型类型显示开关(与笔一致:全量重建,避免增量路径标记未对齐)
$('#showKlcFxType').change(function() { $('#showKlcFxType').change(function() {
refreshChartOnly(); updateChartDisplay();
}); });
// 绑定小周期分型显示开关 // 绑定小周期分型显示开关
$('#showElementKlcFxType').change(function() { $('#showElementKlcFxType').change(function() {
refreshChart(currentData); updateChartDisplay();
}); });
@@ -865,10 +894,7 @@ $('#showElementBollinger').change(function() {
updateChartDisplay(); updateChartDisplay();
}); });
// 绑定K线周期切换 // K线周期切换由 macd_ui.js 统一走 updateChartDisplay(勿再绑 refreshChart,会重复且易漏对齐)
$('input[name="klinePeriod"]').change(function() {
refreshChart(currentData);
});
// 绑定主图U显示开关 // 绑定主图U显示开关
$('#toggleUOnMain').change(function() { $('#toggleUOnMain').change(function() {
+21 -5
View File
@@ -4,6 +4,16 @@ window.App.Charts = (function() {
// 依赖 Indicators // 依赖 Indicators
const Indicators = (window.App && window.App.Indicators) || {}; const Indicators = (window.App && window.App.Indicators) || {};
function sanitizeLinePoints(points) {
if (!Array.isArray(points)) return [];
return points.filter(function (p) {
return p && p.time != null && p.value != null &&
isFinite(Number(p.time)) && isFinite(Number(p.value));
}).map(function (p) {
return { time: Math.floor(Number(p.time)), value: Number(p.value) };
});
}
function addMovingAveragesToChart(candleData) { function addMovingAveragesToChart(candleData) {
if (!window.tvWidget || !tvWidget.mainChart || !candleData || candleData.length === 0) return; if (!window.tvWidget || !tvWidget.mainChart || !candleData || candleData.length === 0) return;
if (!window.movingAverages) return; if (!window.movingAverages) return;
@@ -21,6 +31,8 @@ window.App.Charts = (function() {
try { try {
const maData = Indicators.calculateMA(candleData, maConfig.type, maConfig.length, maConfig.source); const maData = Indicators.calculateMA(candleData, maConfig.type, maConfig.length, maConfig.source);
const smoothedData = maConfig.smoothType !== 'none' ? (window.applySmoothToMA ? window.applySmoothToMA(maData, maConfig.smoothType, maConfig.smoothLength) : maData) : maData; const smoothedData = maConfig.smoothType !== 'none' ? (window.applySmoothToMA ? window.applySmoothToMA(maData, maConfig.smoothType, maConfig.smoothLength) : maData) : maData;
const cleanData = sanitizeLinePoints(smoothedData);
if (!cleanData.length) return;
const maSeries = tvWidget.mainChart.addLineSeries({ const maSeries = tvWidget.mainChart.addLineSeries({
color: maConfig.color, color: maConfig.color,
lineWidth: maConfig.lineWidth || 2, lineWidth: maConfig.lineWidth || 2,
@@ -30,8 +42,8 @@ window.App.Charts = (function() {
priceLineVisible: false, priceLineVisible: false,
crosshairMarkerVisible: true, crosshairMarkerVisible: true,
}); });
maSeries.setData(smoothedData); maSeries.setData(cleanData);
maConfig.data = smoothedData; maConfig.data = cleanData;
tvWidget.series.maSeries.push(maSeries); tvWidget.series.maSeries.push(maSeries);
} catch(e) {} } catch(e) {}
}); });
@@ -51,12 +63,16 @@ window.App.Charts = (function() {
if (!bbConfig.visible) return; if (!bbConfig.visible) return;
try { try {
const bbData = Indicators.calculateBB(candleData, bbConfig.length, bbConfig.upperMultiplier, bbConfig.lowerMultiplier, bbConfig.source); const bbData = Indicators.calculateBB(candleData, bbConfig.length, bbConfig.upperMultiplier, bbConfig.lowerMultiplier, bbConfig.source);
const upper = sanitizeLinePoints(bbData.map(item => ({ time: item.time, value: item.upper })));
const middle = sanitizeLinePoints(bbData.map(item => ({ time: item.time, value: item.middle })));
const lower = sanitizeLinePoints(bbData.map(item => ({ time: item.time, value: item.lower })));
if (!upper.length || !middle.length || !lower.length) return;
const upperSeries = tvWidget.mainChart.addLineSeries({ color: bbConfig.upperColor, lineWidth: bbConfig.lineWidth || 2, lineStyle: bbConfig.lineStyle || 0, lastValueVisible: false, priceLineVisible: false, crosshairMarkerVisible: true }); const upperSeries = tvWidget.mainChart.addLineSeries({ color: bbConfig.upperColor, lineWidth: bbConfig.lineWidth || 2, lineStyle: bbConfig.lineStyle || 0, lastValueVisible: false, priceLineVisible: false, crosshairMarkerVisible: true });
const middleSeries = tvWidget.mainChart.addLineSeries({ color: bbConfig.middleColor, lineWidth: bbConfig.lineWidth || 2, lineStyle: bbConfig.lineStyle || 0, lastValueVisible: false, priceLineVisible: false, crosshairMarkerVisible: true }); const middleSeries = tvWidget.mainChart.addLineSeries({ color: bbConfig.middleColor, lineWidth: bbConfig.lineWidth || 2, lineStyle: bbConfig.lineStyle || 0, lastValueVisible: false, priceLineVisible: false, crosshairMarkerVisible: true });
const lowerSeries = tvWidget.mainChart.addLineSeries({ color: bbConfig.lowerColor, lineWidth: bbConfig.lineWidth || 2, lineStyle: bbConfig.lineStyle || 0, lastValueVisible: false, priceLineVisible: false, crosshairMarkerVisible: true }); const lowerSeries = tvWidget.mainChart.addLineSeries({ color: bbConfig.lowerColor, lineWidth: bbConfig.lineWidth || 2, lineStyle: bbConfig.lineStyle || 0, lastValueVisible: false, priceLineVisible: false, crosshairMarkerVisible: true });
upperSeries.setData(bbData.map(item => ({ time: item.time, value: item.upper }))); upperSeries.setData(upper);
middleSeries.setData(bbData.map(item => ({ time: item.time, value: item.middle }))); middleSeries.setData(middle);
lowerSeries.setData(bbData.map(item => ({ time: item.time, value: item.lower }))); lowerSeries.setData(lower);
bbConfig.data = bbData; bbConfig.data = bbData;
tvWidget.series.bbSeries.push(upperSeries, middleSeries, lowerSeries); tvWidget.series.bbSeries.push(upperSeries, middleSeries, lowerSeries);
} catch(e) {} } catch(e) {}
+4 -1
View File
@@ -63,7 +63,10 @@ window.App.Indicators = (function() {
default: default:
value = sourceData[i]; value = sourceData[i];
} }
result.push({ time: data[i].time, value }); if (value == null || !isFinite(value) || data[i].time == null || !isFinite(Number(data[i].time))) {
continue;
}
result.push({ time: Math.floor(Number(data[i].time)), value: Number(value) });
} }
return result; return result;
} }
+23 -23
View File
@@ -22,8 +22,8 @@
<script src="https://cdn.jsdelivr.net/npm/bootstrap@5.1.3/dist/js/bootstrap.bundle.min.js"></script> <script src="https://cdn.jsdelivr.net/npm/bootstrap@5.1.3/dist/js/bootstrap.bundle.min.js"></script>
<!-- TradingView Widget BEGIN --> <!-- TradingView Widget BEGIN -->
<script src="https://cdn.jsdelivr.net/npm/lightweight-charts@4.0.1/dist/lightweight-charts.standalone.production.js"></script> <script src="https://cdn.jsdelivr.net/npm/lightweight-charts@4.0.1/dist/lightweight-charts.standalone.production.js"></script>
<script defer src="{{ url_for('static', filename='js/indicators.js') }}"></script> <script defer src="{{ url_for('static', filename='js/indicators.js') }}?v=20260809i"></script>
<script defer src="{{ url_for('static', filename='js/charts.js') }}"></script> <script defer src="{{ url_for('static', filename='js/charts.js') }}?v=20260809i"></script>
<!-- TradingView Widget END --> <!-- TradingView Widget END -->
<script> <script>
window.AVAILABLE_TIMEFRAMES = JSON.parse('{{ timeframe_keys_json | safe }}'); window.AVAILABLE_TIMEFRAMES = JSON.parse('{{ timeframe_keys_json | safe }}');
@@ -982,7 +982,7 @@
<input type="datetime-local" id="end_time" class="form-control"> <input type="datetime-local" id="end_time" class="form-control">
</div> </div>
<div class="col-md-1"> <div class="col-md-1">
<button class="btn btn-primary w-100" onclick="updateChart()" style="padding: 8px 6px; font-size: 14px;"> <button class="btn btn-primary w-100" onclick="updateEndTimeToNow(); updateChart({ incremental: false, fullAnalyze: true })" style="padding: 8px 6px; font-size: 14px;">
分析 分析
</button> </button>
</div> </div>
@@ -1028,14 +1028,14 @@
<div class="d-flex align-items-center mb-2"> <div class="d-flex align-items-center mb-2">
<label for="refreshInterval" class="form-label me-2 mb-0">自动刷新:</label> <label for="refreshInterval" class="form-label me-2 mb-0">自动刷新:</label>
<select id="refreshInterval" class="form-select form-select-sm me-2" style="width: 80px;"> <select id="refreshInterval" class="form-select form-select-sm me-2" style="width: 80px;">
<option value="0.0833">5秒</option> <option value="0.0833" selected>5秒</option>
<option value="0.1667">10秒</option> <option value="0.1667">10秒</option>
<option value="0.25">15秒</option> <option value="0.25">15秒</option>
<option value="0.5">30秒</option> <option value="0.5">30秒</option>
<option value="1">1分钟</option> <option value="1">1分钟</option>
<option value="2">2分钟</option> <option value="2">2分钟</option>
<option value="3">3分钟</option> <option value="3">3分钟</option>
<option value="5" selected>5分钟</option> <option value="5">5分钟</option>
<option value="10">10分钟</option> <option value="10">10分钟</option>
</select> </select>
<div class="form-check form-check-inline me-2"> <div class="form-check form-check-inline me-2">
@@ -1383,24 +1383,24 @@
</div> </div>
<script src="https://cdn.jsdelivr.net/npm/bootstrap@5.1.3/dist/js/bootstrap.bundle.min.js"></script> <script src="https://cdn.jsdelivr.net/npm/bootstrap@5.1.3/dist/js/bootstrap.bundle.min.js"></script>
<script defer src="{{ url_for('static', filename='js/app/api_client.js') }}"></script> <script defer src="{{ url_for('static', filename='js/app/api_client.js') }}?v=20260808i"></script>
<script defer src="{{ url_for('static', filename='js/app/state.js') }}"></script> <script defer src="{{ url_for('static', filename='js/app/state.js') }}?v=20260808i"></script>
<script defer src="{{ url_for('static', filename='js/app/trend.js') }}"></script> <script defer src="{{ url_for('static', filename='js/app/trend.js') }}?v=20260808i"></script>
<script defer src="{{ url_for('static', filename='js/app/macd_ui.js') }}?v=20260807f"></script> <script defer src="{{ url_for('static', filename='js/app/macd_ui.js') }}?v=20260808j"></script>
<script defer src="{{ url_for('static', filename='js/app/chart_format.js') }}?v=20260807f"></script> <script defer src="{{ url_for('static', filename='js/app/chart_format.js') }}?v=20260808i"></script>
<script defer src="{{ url_for('static', filename='js/app/chart_view.js') }}?v=20260807f"></script> <script defer src="{{ url_for('static', filename='js/app/chart_view.js') }}?v=20260809q"></script>
<script defer src="{{ url_for('static', filename='js/app/chart_tv_lifecycle.js') }}?v=20260807f"></script> <script defer src="{{ url_for('static', filename='js/app/chart_tv_lifecycle.js') }}?v=20260808i"></script>
<script defer src="{{ url_for('static', filename='js/app/chart_tv_shell.js') }}?v=20260807f"></script> <script defer src="{{ url_for('static', filename='js/app/chart_tv_shell.js') }}?v=20260809j"></script>
<script defer src="{{ url_for('static', filename='js/app/chart_tv_indicators.js') }}?v=20260807f"></script> <script defer src="{{ url_for('static', filename='js/app/chart_tv_indicators.js') }}?v=20260808i"></script>
<script defer src="{{ url_for('static', filename='js/app/chart_tv_chan.js') }}?v=20260807f"></script> <script defer src="{{ url_for('static', filename='js/app/chart_tv_chan.js') }}?v=20260808i"></script>
<script defer src="{{ url_for('static', filename='js/app/chart_tv_overlays.js') }}?v=20260807f"></script> <script defer src="{{ url_for('static', filename='js/app/chart_tv_overlays.js') }}?v=20260809o"></script>
<script defer src="{{ url_for('static', filename='js/app/chart_tv_finalize.js') }}?v=20260807f"></script> <script defer src="{{ url_for('static', filename='js/app/chart_tv_finalize.js') }}?v=20260809d"></script>
<script defer src="{{ url_for('static', filename='js/app/chart_tv.js') }}?v=20260807f"></script> <script defer src="{{ url_for('static', filename='js/app/chart_tv.js') }}?v=20260808i"></script>
<script defer src="{{ url_for('static', filename='js/app/chart_sync.js') }}?v=20260807f"></script> <script defer src="{{ url_for('static', filename='js/app/chart_sync.js') }}?v=20260809j"></script>
<script defer src="{{ url_for('static', filename='js/app/chart_tables.js') }}?v=20260807f"></script> <script defer src="{{ url_for('static', filename='js/app/chart_tables.js') }}?v=20260808i"></script>
<script defer src="{{ url_for('static', filename='js/app/ui.js') }}?v=20260807f"></script> <script defer src="{{ url_for('static', filename='js/app/ui.js') }}?v=20260809j"></script>
<script defer src="{{ url_for('static', filename='js/app/overlays.js') }}"></script> <script defer src="{{ url_for('static', filename='js/app/overlays.js') }}?v=20260808i"></script>
<script defer src="{{ url_for('static', filename='js/app/main.js') }}"></script> <script defer src="{{ url_for('static', filename='js/app/main.js') }}?v=20260808i"></script>
<!-- 均线配置弹窗 --> <!-- 均线配置弹窗 -->
<div id="maConfigModal" class="ma-config-modal"> <div id="maConfigModal" class="ma-config-modal">
File diff suppressed because it is too large Load Diff
+22
View File
@@ -64,11 +64,33 @@ def test_analyze_route_registered():
rules = {r.rule for r in app.url_map.iter_rules()} rules = {r.rule for r in app.url_map.iter_rules()}
assert "/api/analyze" in rules assert "/api/analyze" in rules
assert "/api/klines/recent" in rules
assert "/api/chart_metadata" in rules assert "/api/chart_metadata" in rules
assert "/" in rules assert "/" in rules
assert "/chan_tv" in rules assert "/chan_tv" in rules
def test_klines_recent_returns_tail_only():
from app import app
df = make_ohlcv(n=30)
# analyze 蓝图 star-import 后绑定在 api.analyze 命名空间
with patch("api.analyze.get_kl_data", return_value=df):
client = app.test_client()
resp = client.get(
"/api/klines/recent",
query_string={"symbol": "BTC/USDT:USDT", "timeframe": "5m", "limit": 2},
)
assert resp.status_code == 200
body = resp.get_json()
assert body.get("partial") is True
assert body.get("limit") == 2
assert isinstance(body.get("kline_data"), list)
assert len(body["kline_data"]) == 2
assert "bi_list" not in body
assert "wyckoff" not in body
def test_contract_keys_stable(): def test_contract_keys_stable():
assert "bi_list" in CONTRACT_KEYS and "seg_list" in CONTRACT_KEYS assert "bi_list" in CONTRACT_KEYS and "seg_list" in CONTRACT_KEYS
for k in ("kline_data", "macd", "zs_list", "bsp_list", "chan_macd"): for k in ("kline_data", "macd", "zs_list", "bsp_list", "chan_macd"):
+123
View File
@@ -0,0 +1,123 @@
"""ECR-009: page/API smoke without requiring live provider during assert."""
from __future__ import annotations
import os
import sys
import pytest
# Ensure repo root + web on path like app.py
_ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
_WEB = os.path.join(_ROOT, "web")
for p in (_ROOT, _WEB):
if p not in sys.path:
sys.path.insert(0, p)
os.environ.setdefault("CRYPTO_WYCKOFF_DISABLE", "1")
@pytest.fixture()
def client():
from app import create_app
app = create_app()
app.config["TESTING"] = True
with app.test_client() as c:
yield c
def test_wyckoff_crypto_page_ok(client):
resp = client.get("/wyckoff_crypto")
assert resp.status_code == 200
assert b"Crypto Wyckoff Screener" in resp.data
assert b"fCombo" in resp.data
assert b"chartCanvas" in resp.data
def test_wyckoff_crypto_meta_ok(client):
resp = client.get("/api/wyckoff_crypto/meta")
assert resp.status_code == 200
data = resp.get_json()
assert "engine_version" in data
assert data.get("combo", {}).get("id") == "h8_4_1"
assert data["combo"]["low"] == "1h"
ids = {c["id"] for c in data.get("combos") or []}
assert "h8_4_1" in ids and "d_w_m" in ids
def test_wyckoff_crypto_scan_ok(client):
resp = client.get("/api/wyckoff_crypto/scan?limit=5&combo_id=h8_4_1")
assert resp.status_code == 200
data = resp.get_json()
assert "rows" in data
assert data.get("combo", {}).get("id") == "h8_4_1"
def test_wyckoff_crypto_klines_bad_request(client):
resp = client.get("/api/wyckoff_crypto/klines")
assert resp.status_code == 400
def test_wyckoff_crypto_klines_ok(client):
resp = client.get(
"/api/wyckoff_crypto/klines?symbol=BTC/USDT:USDT&tf=1h&limit=10&combo_id=h8_4_1"
)
assert resp.status_code == 200
data = resp.get_json()
assert "items" in data
assert data.get("tf") == "1h"
assert data.get("intraday") is True
if data["items"]:
assert "datetime" in data["items"][0]
assert "ts" in data["items"][0]
assert "T" in data["items"][0]["datetime"]
assert "+08:00" in data["items"][0]["datetime"]
def test_wyckoff_crypto_klines_bad_limit_ok(client):
resp = client.get(
"/api/wyckoff_crypto/klines?symbol=BTC/USDT:USDT&tf=1h&limit=abc&combo_id=h8_4_1"
)
assert resp.status_code == 200
def test_wyckoff_crypto_overlay_ok(client):
resp = client.get(
"/api/wyckoff_crypto/overlay?symbol=BTC/USDT:USDT&tf=1h&bars=60&combo_id=h8_4_1"
)
assert resp.status_code == 200
data = resp.get_json()
assert "phases" in data
assert "events" in data
def test_combos_add_and_list(client, tmp_path, monkeypatch):
from crypto_wyckoff import combos as cm
monkeypatch.setattr(cm, "_COMBOS_FILE", tmp_path / "combos.json")
monkeypatch.setattr(cm, "_cache", None)
resp = client.get("/api/wyckoff_crypto/combos")
assert resp.status_code == 200
assert len(resp.get_json()["combos"]) >= 2
bad = client.post(
"/api/wyckoff_crypto/combos",
json={"high": "1h", "mid": "4h", "low": "8h"},
)
assert bad.status_code == 400
ok = client.post(
"/api/wyckoff_crypto/combos",
json={"high": "12h", "mid": "4h", "low": "1h", "label": "12h/4h/1h"},
)
assert ok.status_code == 200
cid = ok.get_json()["combo"]["id"]
assert cid == "12h_4h_1h"
deleted = client.delete(f"/api/wyckoff_crypto/combos/{cid}")
assert deleted.status_code == 200
builtin = client.delete("/api/wyckoff_crypto/combos/h8_4_1")
assert builtin.status_code == 400