Compare commits

..
38 changed files with 4522 additions and 3484 deletions
+4 -88
View File
@@ -1,91 +1,7 @@
from __future__ import annotations """Compatibility alias for the canonical review-assistant implementation."""
import json import sys
import urllib.error
import urllib.request
from collections.abc import Iterator
from typing import Any
from llm_stream import OpenAIStreamAccumulator from backend.features.review import agent as _implementation
sys.modules[__name__] = _implementation
class ReviewAssistantError(RuntimeError):
pass
def stream_review_assistant(
context: dict[str, Any],
question: str,
history: list[dict[str, str]],
api_key: str,
base_url: str,
model: str,
timeout: int = 120,
) -> Iterator[str]:
if not api_key or not model:
raise ReviewAssistantError("智能解读服务尚未配置。")
messages = [{"role": "system", "content": _system_prompt(context)}]
messages.extend(history[-12:])
messages.append({"role": "user", "content": question})
request = urllib.request.Request(
f"{base_url.rstrip('/')}/chat/completions",
data=json.dumps(
{"model": model, "messages": messages, "stream": True}, ensure_ascii=False
).encode("utf-8"),
headers={
"Content-Type": "application/json",
"Authorization": f"Bearer {api_key}",
"User-Agent": "XiaobaiReviewWeb/1.0",
"Accept": "text/event-stream",
},
method="POST",
)
try:
with urllib.request.urlopen(request, timeout=timeout) as response:
yielded = False
accumulator = OpenAIStreamAccumulator()
for raw_line in response:
line = raw_line.decode("utf-8", errors="replace").strip()
if not line or line.startswith(":"):
continue
if line.startswith("data:"):
line = line[5:].strip()
if line == "[DONE]":
break
try:
payload = json.loads(line)
except json.JSONDecodeError:
continue
choices = payload.get("choices") or []
if not choices:
continue
choice = choices[0] or {}
content = accumulator.feed(choice)
if content:
yielded = True
yield str(content)
if not yielded:
raise ReviewAssistantError("智能解读未返回有效内容。")
except urllib.error.HTTPError as exc:
raise ReviewAssistantError(f"智能解读服务暂不可用({exc.code})。") from exc
except (urllib.error.URLError, TimeoutError, OSError) as exc:
raise ReviewAssistantError("智能解读连接中断,请稍后重试。") from exc
def _system_prompt(context: dict[str, Any]) -> str:
context_json = json.dumps(context, ensure_ascii=False, separators=(",", ":"))
return f"""
你是“小白复盘”的统一复盘助手。你负责把网页中已经存在的市场统计、策略跟踪、提醒、复盘笔记和手工交易日志连接起来,帮助用户复盘和形成下一步观察计划。
最高优先级规则:
1. 只能使用下方“网页复盘数据”,数据缺失就明确说明,不得补造行情、交易或胜率。
2. 不自动下单,不声称已执行任何操作,不修改策略、提醒、笔记或交易日志。
3. 不承诺收益,不给无条件买卖指令。建议必须写成条件、失效条件和风险边界。
4. 区分市场事实、用户记录和你的推断。引用数字时写明数据日期。
5. 优先结合用户自己的策略跟踪与交易日志寻找可验证的重复模式;样本不足时明确标注。
6. 使用中文,先直接回答,再给数据依据和下一步观察。避免空泛口号,不展示模型、接口或内部工程信息。
7. 控制在 800 个中文字符以内,除非用户明确要求展开。
网页复盘数据:
{context_json}
""".strip()
+10 -1605
View File
File diff suppressed because it is too large Load Diff
+16 -1
View File
@@ -1,3 +1,18 @@
from .facade import AlertServiceMixin
from .http import AlertHttpMixin
from .repository import AlertRepositoryMixin
__all__ = [
"AlertHttpMixin",
"AlertRepositoryMixin",
"AlertService",
"AlertServiceMixin",
]
def __getattr__(name: str):
if name == "AlertService":
from .service import AlertService from .service import AlertService
__all__ = ["AlertService"] return AlertService
raise AttributeError(name)
+30
View File
@@ -0,0 +1,30 @@
from __future__ import annotations
from datetime import date
from typing import Any
class AlertServiceMixin:
def alert_center(self, status: str = "all", as_of: str = "") -> dict[str, Any]:
tracking = self.strategy_tracking.list_tracking(self.current_user_id, 12)
self.alert_service.sync_strategy_tracking(self.current_user_id, tracking)
return self.alert_service.list_alerts(
self.current_user_id, status, as_of
)
def create_alert(self, payload: dict[str, Any]) -> dict[str, Any]:
alert_id = self.alert_service.create_manual(self.current_user_id, payload)
return {"id": alert_id, **self.alert_center()}
def mark_alert_read(self, alert_id: int) -> dict[str, Any]:
self.alert_service.mark_read(self.current_user_id, alert_id)
return self.alert_center()
def mark_all_alerts_read(self, as_of: str = "") -> dict[str, Any]:
compact_date = self.alert_service.calendar_date(as_of or date.today().isoformat())
self.alert_service.mark_all_read(self.current_user_id, compact_date)
return self.alert_center(as_of=compact_date)
def delete_alert(self, alert_id: int) -> dict[str, Any]:
deleted = self.alert_service.delete(self.current_user_id, alert_id)
return {"deleted": deleted, **self.alert_center()}
+16
View File
@@ -0,0 +1,16 @@
from __future__ import annotations
import json
from http import HTTPStatus
class AlertHttpMixin:
def save_alert(self) -> None:
try:
body = self.read_json_body()
self.send_json(
{"ok": True, **self.application_service.create_alert(body)},
HTTPStatus.CREATED,
)
except (ValueError, json.JSONDecodeError) as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
+110
View File
@@ -0,0 +1,110 @@
from __future__ import annotations
from datetime import datetime
from typing import Any
class AlertRepositoryMixin:
def save_alert(
self,
user_id: int,
kind: str,
title: str,
content: str,
available_date: str,
code: str,
dedupe_key: str,
) -> int:
now = datetime.now().astimezone().isoformat(timespec="seconds")
with self.connect() as connection:
connection.execute(
"""
INSERT INTO alerts
(user_id, kind, title, content, available_date, code, dedupe_key,
is_read, created_at, updated_at)
VALUES (?, ?, ?, ?, ?, ?, ?, 0, ?, ?)
ON CONFLICT(user_id, dedupe_key) DO UPDATE SET
title=excluded.title, content=excluded.content,
available_date=excluded.available_date, updated_at=excluded.updated_at
""",
(
int(user_id), kind, title, content, available_date, code,
dedupe_key, now, now,
),
)
row = connection.execute(
"SELECT id FROM alerts WHERE user_id = ? AND dedupe_key = ?",
(int(user_id), dedupe_key),
).fetchone()
return int(row["id"])
def list_alerts(
self, user_id: int, as_of: str, unread_only: bool = False, limit: int = 100
) -> list[dict[str, Any]]:
with self.connect() as connection:
if unread_only:
rows = connection.execute(
"""
SELECT id, kind, title, content, available_date, code, is_read,
created_at, updated_at, read_at
FROM alerts
WHERE user_id = ? AND available_date <= ? AND is_read = 0
ORDER BY available_date DESC, id DESC LIMIT ?
""",
(int(user_id), as_of, max(1, min(300, int(limit)))),
).fetchall()
else:
rows = connection.execute(
"""
SELECT id, kind, title, content, available_date, code, is_read,
created_at, updated_at, read_at
FROM alerts WHERE user_id = ?
ORDER BY CASE WHEN available_date > ? THEN 0 ELSE 1 END,
is_read, available_date, id DESC LIMIT ?
""",
(int(user_id), as_of, max(1, min(300, int(limit)))),
).fetchall()
return [{**dict(row), "is_read": bool(row["is_read"])} for row in rows]
def count_unread_alerts(self, user_id: int, as_of: str) -> int:
with self.connect() as connection:
row = connection.execute(
"""
SELECT COUNT(*) AS total FROM alerts
WHERE user_id = ? AND available_date <= ? AND is_read = 0
""",
(int(user_id), as_of),
).fetchone()
return int(row["total"] if row else 0)
def mark_alert_read(self, user_id: int, alert_id: int) -> bool:
now = datetime.now().astimezone().isoformat(timespec="seconds")
with self.connect() as connection:
cursor = connection.execute(
"""
UPDATE alerts SET is_read = 1, read_at = ?, updated_at = ?
WHERE id = ? AND user_id = ?
""",
(now, now, int(alert_id), int(user_id)),
)
return cursor.rowcount > 0
def mark_all_alerts_read(self, user_id: int, as_of: str) -> int:
now = datetime.now().astimezone().isoformat(timespec="seconds")
with self.connect() as connection:
cursor = connection.execute(
"""
UPDATE alerts SET is_read = 1, read_at = ?, updated_at = ?
WHERE user_id = ? AND available_date <= ? AND is_read = 0
""",
(now, now, int(user_id), as_of),
)
return int(cursor.rowcount)
def delete_alert(self, user_id: int, alert_id: int) -> bool:
with self.connect() as connection:
cursor = connection.execute(
"DELETE FROM alerts WHERE id = ? AND user_id = ?",
(int(alert_id), int(user_id)),
)
return cursor.rowcount > 0
+24
View File
@@ -0,0 +1,24 @@
from .agent import HeavenAgentError, interpret_heaven
from .engine import (
build_five_phase_field,
build_market_hexagram,
build_manual_market_hexagram,
build_personal_field,
hexagram_from_lines,
)
from .http import HeavenHttpMixin
from .repository import HeavenRepositoryMixin
from .service import HeavenServiceMixin
__all__ = [
"HeavenAgentError",
"HeavenHttpMixin",
"HeavenRepositoryMixin",
"HeavenServiceMixin",
"build_five_phase_field",
"build_manual_market_hexagram",
"build_market_hexagram",
"build_personal_field",
"hexagram_from_lines",
"interpret_heaven",
]
+118
View File
@@ -0,0 +1,118 @@
from __future__ import annotations
import json
import time
import urllib.error
import urllib.request
from typing import Any
class HeavenAgentError(RuntimeError):
pass
def interpret_heaven(
mode: str,
context: dict[str, Any],
api_key: str,
base_url: str,
model: str,
timeout: int = 90,
) -> dict[str, Any]:
if mode not in {"trend", "fortune", "heart"}:
raise HeavenAgentError("不支持的问天解读模式。")
if not api_key or not model:
raise HeavenAgentError("LLM API Key 或模型尚未配置。")
system_prompt = _system_prompt(mode)
payload = json.dumps(
{
"model": model,
"messages": [
{"role": "system", "content": system_prompt},
{
"role": "user",
"content": json.dumps(context, ensure_ascii=False, separators=(",", ":")),
},
],
"stream": False,
},
ensure_ascii=False,
).encode("utf-8")
request = urllib.request.Request(
f"{base_url.rstrip('/')}/chat/completions",
data=payload,
headers={
"Content-Type": "application/json",
"Authorization": f"Bearer {api_key}",
"User-Agent": "XiaobaiReviewWeb/0.7",
},
method="POST",
)
started = time.perf_counter()
try:
with urllib.request.urlopen(request, timeout=timeout) as response:
result = json.loads(response.read().decode("utf-8"))
answer = str(result["choices"][0]["message"]["content"]).strip()
if not answer:
raise KeyError("empty response")
except urllib.error.HTTPError as exc:
raise HeavenAgentError(_http_error_message(exc)) from exc
except (urllib.error.URLError, TimeoutError, json.JSONDecodeError, KeyError, IndexError) as exc:
raise HeavenAgentError(f"问天模型调用失败:{exc}") from exc
return {
"answer": answer,
"model": model,
"latency_ms": round((time.perf_counter() - started) * 1000),
}
def _system_prompt(mode: str) -> str:
common = """
你是“小白复盘”的问天解读器。所有历法、卦象、爻位和市场指标已经由确定性程序计算,你只能解释提供的数据,不得改卦、改爻、改干支或编造行情。
问天属于传统文化与娱乐化观察,不是预测模型,不承诺应验,不输出无条件买卖指令,不用神秘话术制造确定性。
使用中文,先给核心判断,再解释结构。引用市场数字时标明数据日期。输出纯文本,可使用简短标题。
""".strip()
if mode == "trend":
return common + """
当前任务是“观势·解势”。六爻从初爻到上爻依次是个股内核、个股外显、板块内核、板块外显、指数内核、指数外显;初二为地、三四为人、五上为天。
行情数据只负责生成六爻,本次解势必须以卦象本身为主,不得根据指数涨跌、板块强弱、涨停家数、成交量或个股表现直接推演方向。context中不会提供这些数字,也不会提供爻位对应的市场角色。
先解释本卦卦名的核心义、上下卦组合及大象;再只解释实际动爻所代表的转折,并说明本卦如何走向之卦;最后可把这一组卦势翻译成克制的市场语言。
重点是“本卦为当下之势,动爻为变化关节,之卦为所趋之势”。不要说明某一动爻对应指数、板块或个股,也不要输出“一看指数、二看涨停家数”一类行情观察条件。
全文控制在300至450个中文字符,最多四小段。卦理约占九成,市场翻译最多一句,只能落到节制、等待、守信、辨伪等行为态度,不得据此预测市场下一阶段、涨跌方向或动能变化。不直接荐股,不使用Markdown表格。
不要使用“必然、确定、必涨、必跌、后续将、进入某阶段”等断语;天机只点出势的性质与变化关系,不替用户宣布结果。
""".strip()
if mode == "fortune":
return common + """
当前任务是“观气·解运”。严格区分五运、六气、节气、月令和日干,不把丙午简单解释为火年。
严格服从five_phase_field.framework提供的确定性结构,不自行重新计算五行:年纲由中运与司天在泉构成;岁半以前司天为主、在泉为辅,岁半以后在泉为主、司天为辅;当前六气层以客气加临主气为核心;日辰只负责触发。节气只用于定位当前六气阶段,不得再次叠加为独立力量。
重点解释framework.relations中的客主同气、客生主、主生客、客克主或主克客,以及客胜为从、主胜为逆、司天在泉同位、天符岁会等已经判定的关系。不得把司天、在泉、主气、客气视为彼此独立的证据重复计权,也不得自行增删传统格局。
首要解释当日气场容易放大参与者的哪些情绪、判断偏差和操作冲动,例如急躁、恐惧、迟疑、追涨、过早止损或路径依赖;再给出一至两个调节动作。
如有personal_profile,结合其日主、十神、五行平衡倾向说明当日对该用户主观状态的影响,但不得把简化平衡倾向说成唯一喜用神,也不得复述或猜测出生日期。
不得引用市场上涨下跌家数、涨跌停数量、成交额、板块强度或个股表现来证明气场。industry_affinity只是五行行业取象示例,不是行情旁证;行业契合度最多在末尾用一句话说明,不得写“当日共振”或暗示相关行业必然涨跌。
全文控制在420至600个中文字符,按“三层气机、人的状态、操作偏向、个人影响(如有)、制衡动作”组织,标题必须写“三层气机”。明确这些是传统历法框架下的观察语言,不宣称气候或五行直接导致股价。
""".strip()
return common + """
当前任务是“观心·解卦”。用户的问题始终只在心中,没有输入给你,因此你不能猜测问题内容,也不能替用户作具体决定。
全文控制在180至350个中文字符。只写一句卦意;一小段动爻与之卦;最后三句极短的问心句。
不要重述六条爻辞,不猜用户未说出口的问题,不以吉凶二字替代思考,不给出股票涨跌预测。语气安静、克制,越短越有余味。
""".strip()
def _http_error_message(exc: urllib.error.HTTPError) -> str:
detail = ""
try:
payload = json.loads(exc.read().decode("utf-8", errors="replace"))
error = payload.get("error")
if isinstance(error, dict):
detail = str(error.get("message") or error.get("code") or "")
elif error:
detail = str(error)
elif payload.get("message"):
detail = str(payload["message"])
except (json.JSONDecodeError, OSError):
detail = ""
suffix = f"{detail[:300]}" if detail else ""
return f"问天模型调用失败(HTTP {exc.code}{suffix}"
File diff suppressed because it is too large Load Diff
+30
View File
@@ -0,0 +1,30 @@
from __future__ import annotations
import json
from http import HTTPStatus
class HeavenHttpMixin:
def heaven_hexagram(self) -> None:
try:
body = self.read_json_body()
result = self.application_service.heaven_hexagram(body.get("lines"))
self.send_json({"ok": True, "hexagram": result})
except (ValueError, json.JSONDecodeError) as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
def heaven_personal(self) -> None:
try:
body = self.read_json_body()
result = self.application_service.heaven_personal(body)
self.send_json({"ok": True, "personal": result})
except (ValueError, json.JSONDecodeError) as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
def heaven_interpret(self) -> None:
try:
body = self.read_json_body()
result = self.application_service.heaven_interpret(body)
self.send_json({"ok": True, **result})
except (ValueError, json.JSONDecodeError) as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
+111
View File
@@ -0,0 +1,111 @@
from __future__ import annotations
import json
import sqlite3
from datetime import datetime
from typing import Any
class HeavenRepositoryMixin:
@staticmethod
def _heaven_reading_dict(row: sqlite3.Row | None) -> dict[str, Any] | None:
if not row:
return None
return {
"id": int(row["id"]),
"mode": str(row["mode"]),
"context_date": str(row["context_date"]),
"subject": str(row["subject"]),
"subject_detail": str(row["subject_detail"]),
"answer": str(row["answer"]),
"created_at": str(row["created_at"]),
}
def save_heaven_reading(
self,
user_id: int,
mode: str,
context_date: str,
subject: str,
subject_detail: str,
answer: str,
context_snapshot: dict[str, Any],
dedupe_key: str,
) -> dict[str, Any]:
now = datetime.now().astimezone().isoformat(timespec="seconds")
snapshot_json = json.dumps(
context_snapshot, ensure_ascii=False, separators=(",", ":")
)
with self.connect() as connection:
connection.execute(
"""
INSERT INTO heaven_readings
(user_id, mode, context_date, subject, subject_detail, answer,
context_snapshot, dedupe_key, created_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(user_id, dedupe_key) DO NOTHING
""",
(
int(user_id), mode, context_date, subject, subject_detail,
answer, snapshot_json, dedupe_key, now,
),
)
row = connection.execute(
"""
SELECT id, mode, context_date, subject, subject_detail, answer, created_at
FROM heaven_readings WHERE user_id = ? AND dedupe_key = ?
""",
(int(user_id), dedupe_key),
).fetchone()
connection.execute(
"""
DELETE FROM heaven_readings
WHERE user_id = ? AND mode = ? AND id NOT IN (
SELECT id FROM heaven_readings
WHERE user_id = ? AND mode = ? ORDER BY id DESC LIMIT 100
)
""",
(int(user_id), mode, int(user_id), mode),
)
result = self._heaven_reading_dict(row)
if not result:
raise ValueError("解读记录保存失败。")
return result
def list_heaven_readings(
self,
user_id: int,
mode: str,
context_date: str = "",
limit: int = 100,
) -> list[dict[str, Any]]:
clauses = ["user_id = ?", "mode = ?"]
parameters: list[Any] = [int(user_id), mode]
if context_date:
clauses.append("context_date = ?")
parameters.append(context_date)
parameters.append(max(1, min(100, int(limit))))
with self.connect() as connection:
rows = connection.execute(
f"""
SELECT id, mode, context_date, subject, subject_detail, answer, created_at
FROM heaven_readings WHERE {' AND '.join(clauses)}
ORDER BY context_date DESC, id DESC LIMIT ?
""",
parameters,
).fetchall()
return [self._heaven_reading_dict(row) for row in rows if row]
def latest_heaven_reading(
self, user_id: int, mode: str, context_date: str = ""
) -> dict[str, Any] | None:
items = self.list_heaven_readings(user_id, mode, context_date, 1)
return items[0] if items else None
def delete_heaven_reading(self, user_id: int, reading_id: int) -> bool:
with self.connect() as connection:
cursor = connection.execute(
"DELETE FROM heaven_readings WHERE id = ? AND user_id = ?",
(int(reading_id), int(user_id)),
)
return cursor.rowcount > 0
File diff suppressed because it is too large Load Diff
+22 -2
View File
@@ -1,3 +1,23 @@
from .trade_journal import EMOTIONS, TRADE_ACTIONS, TradeJournalService from .agent import ReviewAssistantError, stream_review_assistant
from .http import ReviewHttpMixin
from .repository import ReviewRepositoryMixin
from .service import ReviewServiceMixin
__all__ = ["EMOTIONS", "TRADE_ACTIONS", "TradeJournalService"] __all__ = [
"EMOTIONS",
"ReviewAssistantError",
"ReviewHttpMixin",
"ReviewRepositoryMixin",
"ReviewServiceMixin",
"TRADE_ACTIONS",
"TradeJournalService",
"stream_review_assistant",
]
def __getattr__(name: str):
if name in {"EMOTIONS", "TRADE_ACTIONS", "TradeJournalService"}:
from . import trade_journal
return getattr(trade_journal, name)
raise AttributeError(name)
+91
View File
@@ -0,0 +1,91 @@
from __future__ import annotations
import json
import urllib.error
import urllib.request
from collections.abc import Iterator
from typing import Any
from llm_stream import OpenAIStreamAccumulator
class ReviewAssistantError(RuntimeError):
pass
def stream_review_assistant(
context: dict[str, Any],
question: str,
history: list[dict[str, str]],
api_key: str,
base_url: str,
model: str,
timeout: int = 120,
) -> Iterator[str]:
if not api_key or not model:
raise ReviewAssistantError("智能解读服务尚未配置。")
messages = [{"role": "system", "content": _system_prompt(context)}]
messages.extend(history[-12:])
messages.append({"role": "user", "content": question})
request = urllib.request.Request(
f"{base_url.rstrip('/')}/chat/completions",
data=json.dumps(
{"model": model, "messages": messages, "stream": True}, ensure_ascii=False
).encode("utf-8"),
headers={
"Content-Type": "application/json",
"Authorization": f"Bearer {api_key}",
"User-Agent": "XiaobaiReviewWeb/1.0",
"Accept": "text/event-stream",
},
method="POST",
)
try:
with urllib.request.urlopen(request, timeout=timeout) as response:
yielded = False
accumulator = OpenAIStreamAccumulator()
for raw_line in response:
line = raw_line.decode("utf-8", errors="replace").strip()
if not line or line.startswith(":"):
continue
if line.startswith("data:"):
line = line[5:].strip()
if line == "[DONE]":
break
try:
payload = json.loads(line)
except json.JSONDecodeError:
continue
choices = payload.get("choices") or []
if not choices:
continue
choice = choices[0] or {}
content = accumulator.feed(choice)
if content:
yielded = True
yield str(content)
if not yielded:
raise ReviewAssistantError("智能解读未返回有效内容。")
except urllib.error.HTTPError as exc:
raise ReviewAssistantError(f"智能解读服务暂不可用({exc.code})。") from exc
except (urllib.error.URLError, TimeoutError, OSError) as exc:
raise ReviewAssistantError("智能解读连接中断,请稍后重试。") from exc
def _system_prompt(context: dict[str, Any]) -> str:
context_json = json.dumps(context, ensure_ascii=False, separators=(",", ":"))
return f"""
你是“小白复盘”的统一复盘助手。你负责把网页中已经存在的市场统计、策略跟踪、提醒、复盘笔记和手工交易日志连接起来,帮助用户复盘和形成下一步观察计划。
最高优先级规则:
1. 只能使用下方“网页复盘数据”,数据缺失就明确说明,不得补造行情、交易或胜率。
2. 不自动下单,不声称已执行任何操作,不修改策略、提醒、笔记或交易日志。
3. 不承诺收益,不给无条件买卖指令。建议必须写成条件、失效条件和风险边界。
4. 区分市场事实、用户记录和你的推断。引用数字时写明数据日期。
5. 优先结合用户自己的策略跟踪与交易日志寻找可验证的重复模式;样本不足时明确标注。
6. 使用中文,先直接回答,再给数据依据和下一步观察。避免空泛口号,不展示模型、接口或内部工程信息。
7. 控制在 800 个中文字符以内,除非用户明确要求展开。
网页复盘数据:
{context_json}
""".strip()
+97
View File
@@ -0,0 +1,97 @@
from __future__ import annotations
import json
from datetime import date
from http import HTTPStatus
from backend.bootstrap.config import normalize_date, validate_stock_code, validate_text
from backend.features.review.agent import ReviewAssistantError
class ReviewHttpMixin:
def save_trade_entry(self) -> None:
try:
body = self.read_json_body()
self.send_json(
{"ok": True, **self.application_service.save_trade_entry(body)},
HTTPStatus.CREATED,
)
except (ValueError, json.JSONDecodeError) as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
def stream_assistant_chat(self) -> None:
try:
body = self.read_json_body()
stream = self.application_service.assistant_stream(body)
except (ValueError, json.JSONDecodeError) as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
return
self.send_response(HTTPStatus.OK)
self.send_header("Content-Type", "application/x-ndjson; charset=utf-8")
self.send_header("Cache-Control", "no-cache, no-transform")
self.send_header("X-Accel-Buffering", "no")
self.send_header("Connection", "close")
self.end_headers()
try:
for chunk in stream:
self._write_stream_event({"type": "delta", "content": chunk})
self._write_stream_event({"type": "done"})
except (ValueError, ReviewAssistantError) as exc:
self._write_stream_event({"type": "error", "error": str(exc)})
except (BrokenPipeError, ConnectionResetError):
pass
finally:
self.close_connection = True
def save_watchlist(self) -> None:
try:
body = self.read_json_body()
code = validate_stock_code(str(body.get("code", "")))
name = validate_text(body.get("name"), "股票名称", 30, required=True)
sector = validate_text(body.get("sector"), "所属板块", 50)
color = str(body.get("color") or "red")
if color not in {"red", "blue", "green", "amber"}:
raise ValueError("标记颜色不支持。")
remark = validate_text(body.get("remark"), "跟踪备注", 240)
service = self.application_service
service.database.save_watchlist(
service.current_user_id, code, name, sector, color, remark
)
self.send_json(
{
"ok": True,
"items": service.database.list_watchlist(service.current_user_id),
}
)
except (ValueError, json.JSONDecodeError) as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
def save_note(self) -> None:
try:
body = self.read_json_body()
code = str(body.get("code") or "").strip()
if code:
code = validate_stock_code(code)
stock_name = validate_text(body.get("stock_name"), "股票名称", 30)
trade_date = normalize_date(str(body.get("trade_date") or date.today().isoformat()))
summary = validate_text(body.get("summary"), "盘面摘要", 500)
content = validate_text(body.get("content"), "复盘内容", 5000)
plan = validate_text(body.get("plan"), "明日计划", 2000)
if not summary and not content and not plan:
raise ValueError("每日复盘内容不能全部为空。")
raw_id = body.get("id")
note_id = int(raw_id) if raw_id else None
service = self.application_service
saved_id = service.database.save_note(
service.current_user_id,
code,
stock_name,
trade_date,
content,
plan,
note_id,
summary=summary,
)
self.send_json({"ok": True, "id": saved_id})
except (ValueError, TypeError, json.JSONDecodeError) as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
+281
View File
@@ -0,0 +1,281 @@
from __future__ import annotations
import json
from datetime import datetime
from typing import Any
class ReviewRepositoryMixin:
def list_watchlist(self, user_id: int) -> list[dict[str, Any]]:
with self.connect() as connection:
rows = connection.execute(
"""
SELECT code, name, sector, color, remark, created_at, updated_at
FROM watchlist WHERE user_id = ? ORDER BY updated_at DESC, code
""",
(int(user_id),),
).fetchall()
return [dict(row) for row in rows]
def save_watchlist(
self, user_id: int, code: str, name: str, sector: str, color: str,
remark: str | None = None,
) -> None:
now = datetime.now().astimezone().isoformat(timespec="seconds")
with self.connect() as connection:
existing = connection.execute(
"SELECT remark FROM watchlist WHERE user_id = ? AND code = ?",
(int(user_id), code),
).fetchone()
saved_remark = (
str(existing["remark"] or "") if remark is None and existing else str(remark or "")
)
connection.execute(
"""
INSERT INTO watchlist
(user_id, code, name, sector, color, remark, created_at, updated_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(user_id, code) DO UPDATE SET
name = excluded.name,
sector = excluded.sector,
color = excluded.color,
remark = excluded.remark,
updated_at = excluded.updated_at
""",
(int(user_id), code, name, sector, color, saved_remark, now, now),
)
def watchlist_price_history(
self, codes: list[str], end_date: str, limit_per_code: int = 6
) -> dict[str, list[dict[str, Any]]]:
result: dict[str, list[dict[str, Any]]] = {}
if not codes:
return result
with self.connect() as connection:
for code in codes:
rows = connection.execute(
"""
SELECT trade_date, ts_code, close, pct_chg
FROM daily_bars
WHERE substr(ts_code, 1, 6) = ? AND trade_date <= ?
ORDER BY trade_date DESC LIMIT ?
""",
(str(code), end_date, int(limit_per_code)),
).fetchall()
result[str(code)] = [dict(row) for row in reversed(rows)]
return result
def delete_watchlist(self, user_id: int, code: str) -> bool:
with self.connect() as connection:
cursor = connection.execute(
"DELETE FROM watchlist WHERE user_id = ? AND code = ?",
(int(user_id), code),
)
return cursor.rowcount > 0
def list_notes(
self,
user_id: int,
code: str = "",
trade_date: str = "",
scope: str = "all",
) -> list[dict[str, Any]]:
clauses: list[str] = ["user_id = ?"]
parameters: list[Any] = [int(user_id)]
if scope == "daily":
clauses.append("code = ''")
elif scope == "stock":
clauses.append("code <> ''")
if code:
clauses.append("code = ?")
parameters.append(code)
if trade_date:
clauses.append("trade_date = ?")
parameters.append(trade_date)
where = f"WHERE {' AND '.join(clauses)}" if clauses else ""
with self.connect() as connection:
rows = connection.execute(
f"""
SELECT id, code, stock_name, trade_date, summary, content, plan, created_at, updated_at
FROM review_notes {where}
ORDER BY trade_date DESC, updated_at DESC, id DESC LIMIT 200
""",
parameters,
).fetchall()
return [dict(row) for row in rows]
def save_note(
self,
user_id: int,
code: str,
stock_name: str,
trade_date: str,
content: str,
plan: str,
note_id: int | None = None,
summary: str = "",
) -> int:
now = datetime.now().astimezone().isoformat(timespec="seconds")
with self.connect() as connection:
if note_id:
cursor = connection.execute(
"""
UPDATE review_notes
SET code = ?, stock_name = ?, trade_date = ?, summary = ?, content = ?, plan = ?, updated_at = ?
WHERE id = ? AND user_id = ?
""",
(code, stock_name, trade_date, summary, content, plan, now, note_id, int(user_id)),
)
if cursor.rowcount == 0:
raise ValueError("复盘笔记不存在。")
return note_id
cursor = connection.execute(
"""
INSERT INTO review_notes
(user_id, code, stock_name, trade_date, summary, content, plan, created_at, updated_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)
""",
(int(user_id), code, stock_name, trade_date, summary, content, plan, now, now),
)
return int(cursor.lastrowid)
def delete_note(self, user_id: int, note_id: int) -> bool:
with self.connect() as connection:
cursor = connection.execute(
"DELETE FROM review_notes WHERE id = ? AND user_id = ?",
(note_id, int(user_id)),
)
return cursor.rowcount > 0
def save_trade_entry(
self,
user_id: int,
trade_date: str,
code: str,
name: str,
action: str,
price: float,
quantity: int,
position_pct: float,
pnl_amount: float | None,
pnl_pct: float | None,
thesis: str,
execution: str,
emotion: str,
tags: list[str],
trade_id: int | None = None,
) -> int:
now = datetime.now().astimezone().isoformat(timespec="seconds")
tags_json = json.dumps(tags, ensure_ascii=False, separators=(",", ":"))
with self.connect() as connection:
if trade_id:
cursor = connection.execute(
"""
UPDATE trade_entries SET
trade_date=?, code=?, name=?, action=?, price=?, quantity=?,
position_pct=?, pnl_amount=?, pnl_pct=?, thesis=?, execution=?,
emotion=?, tags=?, updated_at=?
WHERE id=? AND user_id=?
""",
(
trade_date, code, name, action, price, quantity, position_pct,
pnl_amount, pnl_pct, thesis, execution, emotion, tags_json, now,
int(trade_id), int(user_id),
),
)
if cursor.rowcount == 0:
raise ValueError("交易记录不存在或无权修改。")
return int(trade_id)
cursor = connection.execute(
"""
INSERT INTO trade_entries
(user_id, trade_date, code, name, action, price, quantity,
position_pct, pnl_amount, pnl_pct, thesis, execution, emotion,
tags, created_at, updated_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
""",
(
int(user_id), trade_date, code, name, action, price, quantity,
position_pct, pnl_amount, pnl_pct, thesis, execution, emotion,
tags_json, now, now,
),
)
return int(cursor.lastrowid)
def list_trade_entries(
self, user_id: int, start_date: str = "", end_date: str = "", code: str = "",
limit: int = 300,
) -> list[dict[str, Any]]:
clauses = ["user_id = ?"]
parameters: list[Any] = [int(user_id)]
if start_date:
clauses.append("trade_date >= ?")
parameters.append(start_date)
if end_date:
clauses.append("trade_date <= ?")
parameters.append(end_date)
if code:
clauses.append("code = ?")
parameters.append(code)
parameters.append(max(1, min(1000, int(limit))))
with self.connect() as connection:
rows = connection.execute(
f"""
SELECT * FROM trade_entries WHERE {' AND '.join(clauses)}
ORDER BY trade_date DESC, id DESC LIMIT ?
""",
parameters,
).fetchall()
return [dict(row) for row in rows]
def delete_trade_entry(self, user_id: int, trade_id: int) -> bool:
with self.connect() as connection:
cursor = connection.execute(
"DELETE FROM trade_entries WHERE id = ? AND user_id = ?",
(int(trade_id), int(user_id)),
)
return cursor.rowcount > 0
def save_assistant_exchange(
self, user_id: int, question: str, answer: str, context_date: str
) -> None:
now = datetime.now().astimezone().isoformat(timespec="seconds")
with self.connect() as connection:
connection.executemany(
"""
INSERT INTO assistant_messages
(user_id, role, content, context_date, created_at)
VALUES (?, ?, ?, ?, ?)
""",
[
(int(user_id), "user", question, context_date, now),
(int(user_id), "assistant", answer, context_date, now),
],
)
connection.execute(
"""
DELETE FROM assistant_messages WHERE user_id = ? AND id NOT IN (
SELECT id FROM assistant_messages
WHERE user_id = ? ORDER BY id DESC LIMIT 200
)
""",
(int(user_id), int(user_id)),
)
def list_assistant_messages(self, user_id: int, limit: int = 100) -> list[dict[str, Any]]:
with self.connect() as connection:
rows = connection.execute(
"""
SELECT role, content, context_date, created_at FROM assistant_messages
WHERE user_id = ? ORDER BY id DESC LIMIT ?
""",
(int(user_id), max(1, min(200, int(limit)))),
).fetchall()
return [dict(row) for row in reversed(rows)]
def delete_assistant_messages(self, user_id: int) -> int:
with self.connect() as connection:
cursor = connection.execute(
"DELETE FROM assistant_messages WHERE user_id = ?", (int(user_id),)
)
return int(cursor.rowcount)
+188
View File
@@ -0,0 +1,188 @@
from __future__ import annotations
from datetime import date, datetime, timedelta
from typing import Any
from backend.bootstrap.config import normalize_date, tushare_code, validate_text
from backend.data.providers.tushare_client import TushareError
from backend.features.review.agent import ReviewAssistantError, stream_review_assistant
class ReviewServiceMixin:
def trade_entries(
self, start_date: str = "", end_date: str = "", code: str = ""
) -> dict[str, Any]:
return self.trade_journal.list_entries(
self.current_user_id, start_date, end_date, code
)
def review_watchlist(self, trade_date: str) -> dict[str, Any]:
normalized_date = normalize_date(trade_date)
items = self.database.list_watchlist(self.current_user_id)
if not items:
return {"items": [], "trade_date": normalized_date}
resolved_date = normalized_date
if self.configured:
try:
client = self._tushare_client()
resolved_date, _ = client.resolve_trade_context(normalized_date)
history = self.database.watchlist_price_history(
[str(item["code"]) for item in items], resolved_date
)
missing_codes = [
str(item["code"]) for item in items
if len(history.get(str(item["code"])) or []) < 6
]
start_date = (
datetime.strptime(resolved_date, "%Y%m%d") - timedelta(days=24)
).strftime("%Y%m%d")
for code in missing_codes:
rows = client.query(
"daily",
{
"ts_code": tushare_code(code),
"start_date": start_date,
"end_date": resolved_date,
},
"ts_code,trade_date,open,high,low,close,pct_chg,vol,amount",
)
if rows:
self.database.upsert_daily_bars(rows)
if missing_codes:
history = self.database.watchlist_price_history(
[str(item["code"]) for item in items], resolved_date
)
except (TushareError, ValueError):
history = self.database.watchlist_price_history(
[str(item["code"]) for item in items], resolved_date
)
else:
history = self.database.watchlist_price_history(
[str(item["code"]) for item in items], resolved_date
)
auction_scores: dict[str, Any] = {}
try:
auction = self.auction_center(normalized_date, False)
auction_scores = {
str(row.get("code") or ""): row.get("attention_score")
for row in (auction.get("watchlist_rows") or [])
if row.get("available", True)
}
except (TushareError, ValueError):
pass
enriched = []
for item in items:
code = str(item.get("code") or "")
bars = history.get(code) or []
latest = bars[-1] if bars else {}
close = float(latest.get("close") or 0)
base_close = float(bars[-6].get("close") or 0) if len(bars) >= 6 else 0
enriched.append(
{
**item,
"change": (
round(float(latest.get("pct_chg") or 0), 2) if latest else None
),
"return_5d": (
round((close / base_close - 1) * 100, 2)
if close > 0 and base_close > 0 else None
),
"attention_score": auction_scores.get(code),
"market_date": str(latest.get("trade_date") or ""),
}
)
return {"items": enriched, "trade_date": resolved_date}
def save_trade_entry(self, payload: dict[str, Any]) -> dict[str, Any]:
trade_id = self.trade_journal.save(self.current_user_id, payload)
return {"id": trade_id, **self.trade_entries()}
def delete_trade_entry(self, trade_id: int) -> dict[str, Any]:
deleted = self.trade_journal.delete(self.current_user_id, trade_id)
return {"deleted": deleted, **self.trade_entries()}
def assistant_messages(self) -> list[dict[str, Any]]:
return self.database.list_assistant_messages(self.current_user_id)
def clear_assistant_messages(self) -> int:
return self.database.delete_assistant_messages(self.current_user_id)
def assistant_stream(self, payload: dict[str, Any]):
question = validate_text(payload.get("question"), "问题", 2000, required=True)
trade_date = normalize_date(
str(payload.get("trade_date") or date.today().isoformat())
)
context = self._assistant_context(trade_date)
history = [
{"role": item["role"], "content": str(item["content"])[:4000]}
for item in self.assistant_messages()[-12:]
if item.get("role") in {"user", "assistant"}
]
def generate():
answer_parts: list[str] = []
events = self.llm_gateway.stream(
"assistant",
"review-assistant-v1",
lambda profile: stream_review_assistant(
context,
question,
history,
profile.api_key,
profile.base_url,
profile.model,
),
(ReviewAssistantError,),
)
for event in events:
if event.kind == "delta":
chunk = str(event.value or "")
answer_parts.append(chunk)
yield chunk
elif event.kind == "complete":
self.database.save_assistant_exchange(
self.current_user_id,
question,
"".join(answer_parts).strip(),
trade_date,
)
return generate()
def _assistant_context(self, trade_date: str) -> dict[str, Any]:
dashboard = self.get_dashboard(trade_date)
actual_date = normalize_date(
str((dashboard.get("meta") or {}).get("trade_date") or trade_date)
)
sentiment = self.sentiment_history(actual_date, 10)
tracking = self.strategy_tracking.list_tracking(self.current_user_id, 5)
alerts = self.alert_service.list_alerts(
self.current_user_id, "all", date.today().isoformat()
)
trades = self.trade_journal.list_entries(
self.current_user_id, end_date=actual_date
)
return {
"data_date": actual_date,
"market": {
"overview": dashboard.get("overview") or {},
"top_sectors": (dashboard.get("sectors") or [])[:8],
"limit_performance": dashboard.get("limit_performance") or {},
"sentiment_history": (sentiment.get("rows") or [])[-10:],
},
"personal": {
"watchlist": self.database.list_watchlist(self.current_user_id)[:30],
"review_notes": self.database.list_notes(
self.current_user_id, scope="daily"
)[:10],
"strategy_tracking": {
"summary": tracking.get("summary") or {},
"batches": (tracking.get("batches") or [])[:5],
},
"alerts": (alerts.get("items") or [])[:20],
"trade_summary": trades.get("summary") or {},
"trade_entries": (trades.get("items") or [])[:30],
},
}
+6 -486
View File
@@ -8,12 +8,15 @@ from typing import Any
from backend.database import MIGRATIONS, MigrationRunner, SQLiteConnectionFactory from backend.database import MIGRATIONS, MigrationRunner, SQLiteConnectionFactory
from backend.features.accounts.repository import AccountRepositoryMixin from backend.features.accounts.repository import AccountRepositoryMixin
from backend.features.alerts.repository import AlertRepositoryMixin
from backend.features.auction.repository import AuctionRepositoryMixin from backend.features.auction.repository import AuctionRepositoryMixin
from backend.features.dragon_tiger.repository import DragonTigerRepositoryMixin from backend.features.dragon_tiger.repository import DragonTigerRepositoryMixin
from backend.features.heaven.repository import HeavenRepositoryMixin
from backend.features.market.repository import MarketRepositoryMixin from backend.features.market.repository import MarketRepositoryMixin
from backend.features.mentor.repository import MentorRepositoryMixin from backend.features.mentor.repository import MentorRepositoryMixin
from backend.features.pools.repository import PoolRepositoryMixin from backend.features.pools.repository import PoolRepositoryMixin
from backend.features.popularity.repository import PopularityRepositoryMixin from backend.features.popularity.repository import PopularityRepositoryMixin
from backend.features.review.repository import ReviewRepositoryMixin
from backend.features.screener.repository import ScreenerRepositoryMixin from backend.features.screener.repository import ScreenerRepositoryMixin
from backend.features.system.repository import SystemSettingsRepositoryMixin from backend.features.system.repository import SystemSettingsRepositoryMixin
from backend.llm.repository import LLMAuditRepositoryMixin from backend.llm.repository import LLMAuditRepositoryMixin
@@ -21,12 +24,15 @@ from backend.llm.repository import LLMAuditRepositoryMixin
class ReviewDatabase( class ReviewDatabase(
AccountRepositoryMixin, AccountRepositoryMixin,
AlertRepositoryMixin,
AuctionRepositoryMixin, AuctionRepositoryMixin,
DragonTigerRepositoryMixin, DragonTigerRepositoryMixin,
HeavenRepositoryMixin,
MarketRepositoryMixin, MarketRepositoryMixin,
MentorRepositoryMixin, MentorRepositoryMixin,
PoolRepositoryMixin, PoolRepositoryMixin,
PopularityRepositoryMixin, PopularityRepositoryMixin,
ReviewRepositoryMixin,
ScreenerRepositoryMixin, ScreenerRepositoryMixin,
SystemSettingsRepositoryMixin, SystemSettingsRepositoryMixin,
LLMAuditRepositoryMixin, LLMAuditRepositoryMixin,
@@ -660,147 +666,6 @@ class ReviewDatabase(
MigrationRunner().apply(connection, MIGRATIONS) MigrationRunner().apply(connection, MIGRATIONS)
def list_watchlist(self, user_id: int) -> list[dict[str, Any]]:
with self.connect() as connection:
rows = connection.execute(
"""
SELECT code, name, sector, color, remark, created_at, updated_at
FROM watchlist WHERE user_id = ? ORDER BY updated_at DESC, code
""",
(int(user_id),),
).fetchall()
return [dict(row) for row in rows]
def save_watchlist(
self, user_id: int, code: str, name: str, sector: str, color: str,
remark: str | None = None,
) -> None:
now = datetime.now().astimezone().isoformat(timespec="seconds")
with self.connect() as connection:
existing = connection.execute(
"SELECT remark FROM watchlist WHERE user_id = ? AND code = ?",
(int(user_id), code),
).fetchone()
saved_remark = (
str(existing["remark"] or "") if remark is None and existing else str(remark or "")
)
connection.execute(
"""
INSERT INTO watchlist
(user_id, code, name, sector, color, remark, created_at, updated_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(user_id, code) DO UPDATE SET
name = excluded.name,
sector = excluded.sector,
color = excluded.color,
remark = excluded.remark,
updated_at = excluded.updated_at
""",
(int(user_id), code, name, sector, color, saved_remark, now, now),
)
def watchlist_price_history(
self, codes: list[str], end_date: str, limit_per_code: int = 6
) -> dict[str, list[dict[str, Any]]]:
result: dict[str, list[dict[str, Any]]] = {}
if not codes:
return result
with self.connect() as connection:
for code in codes:
rows = connection.execute(
"""
SELECT trade_date, ts_code, close, pct_chg
FROM daily_bars
WHERE substr(ts_code, 1, 6) = ? AND trade_date <= ?
ORDER BY trade_date DESC LIMIT ?
""",
(str(code), end_date, int(limit_per_code)),
).fetchall()
result[str(code)] = [dict(row) for row in reversed(rows)]
return result
def delete_watchlist(self, user_id: int, code: str) -> bool:
with self.connect() as connection:
cursor = connection.execute(
"DELETE FROM watchlist WHERE user_id = ? AND code = ?",
(int(user_id), code),
)
return cursor.rowcount > 0
def list_notes(
self,
user_id: int,
code: str = "",
trade_date: str = "",
scope: str = "all",
) -> list[dict[str, Any]]:
clauses: list[str] = ["user_id = ?"]
parameters: list[Any] = [int(user_id)]
if scope == "daily":
clauses.append("code = ''")
elif scope == "stock":
clauses.append("code <> ''")
if code:
clauses.append("code = ?")
parameters.append(code)
if trade_date:
clauses.append("trade_date = ?")
parameters.append(trade_date)
where = f"WHERE {' AND '.join(clauses)}" if clauses else ""
with self.connect() as connection:
rows = connection.execute(
f"""
SELECT id, code, stock_name, trade_date, summary, content, plan, created_at, updated_at
FROM review_notes {where}
ORDER BY trade_date DESC, updated_at DESC, id DESC LIMIT 200
""",
parameters,
).fetchall()
return [dict(row) for row in rows]
def save_note(
self,
user_id: int,
code: str,
stock_name: str,
trade_date: str,
content: str,
plan: str,
note_id: int | None = None,
summary: str = "",
) -> int:
now = datetime.now().astimezone().isoformat(timespec="seconds")
with self.connect() as connection:
if note_id:
cursor = connection.execute(
"""
UPDATE review_notes
SET code = ?, stock_name = ?, trade_date = ?, summary = ?, content = ?, plan = ?, updated_at = ?
WHERE id = ? AND user_id = ?
""",
(code, stock_name, trade_date, summary, content, plan, now, note_id, int(user_id)),
)
if cursor.rowcount == 0:
raise ValueError("复盘笔记不存在。")
return note_id
cursor = connection.execute(
"""
INSERT INTO review_notes
(user_id, code, stock_name, trade_date, summary, content, plan, created_at, updated_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)
""",
(int(user_id), code, stock_name, trade_date, summary, content, plan, now, now),
)
return int(cursor.lastrowid)
def delete_note(self, user_id: int, note_id: int) -> bool:
with self.connect() as connection:
cursor = connection.execute(
"DELETE FROM review_notes WHERE id = ? AND user_id = ?",
(note_id, int(user_id)),
)
return cursor.rowcount > 0
def list_sector_phase_overrides(self) -> dict[str, str]: def list_sector_phase_overrides(self) -> dict[str, str]:
with self.connect() as connection: with self.connect() as connection:
@@ -830,10 +695,6 @@ class ReviewDatabase(
(name,), (name,),
) )
return cursor.rowcount > 0 return cursor.rowcount > 0
def list_wencai_saved_queries( def list_wencai_saved_queries(
self, user_id: int, limit: int = 30 self, user_id: int, limit: int = 30
) -> list[dict[str, Any]]: ) -> list[dict[str, Any]]:
@@ -883,344 +744,3 @@ class ReviewDatabase(
(int(query_id), int(user_id)), (int(query_id), int(user_id)),
) )
return cursor.rowcount > 0 return cursor.rowcount > 0
def save_alert(
self,
user_id: int,
kind: str,
title: str,
content: str,
available_date: str,
code: str,
dedupe_key: str,
) -> int:
now = datetime.now().astimezone().isoformat(timespec="seconds")
with self.connect() as connection:
connection.execute(
"""
INSERT INTO alerts
(user_id, kind, title, content, available_date, code, dedupe_key,
is_read, created_at, updated_at)
VALUES (?, ?, ?, ?, ?, ?, ?, 0, ?, ?)
ON CONFLICT(user_id, dedupe_key) DO UPDATE SET
title=excluded.title, content=excluded.content,
available_date=excluded.available_date, updated_at=excluded.updated_at
""",
(
int(user_id), kind, title, content, available_date, code,
dedupe_key, now, now,
),
)
row = connection.execute(
"SELECT id FROM alerts WHERE user_id = ? AND dedupe_key = ?",
(int(user_id), dedupe_key),
).fetchone()
return int(row["id"])
def list_alerts(
self, user_id: int, as_of: str, unread_only: bool = False, limit: int = 100
) -> list[dict[str, Any]]:
with self.connect() as connection:
if unread_only:
rows = connection.execute(
"""
SELECT id, kind, title, content, available_date, code, is_read,
created_at, updated_at, read_at
FROM alerts
WHERE user_id = ? AND available_date <= ? AND is_read = 0
ORDER BY available_date DESC, id DESC LIMIT ?
""",
(int(user_id), as_of, max(1, min(300, int(limit)))),
).fetchall()
else:
rows = connection.execute(
"""
SELECT id, kind, title, content, available_date, code, is_read,
created_at, updated_at, read_at
FROM alerts WHERE user_id = ?
ORDER BY CASE WHEN available_date > ? THEN 0 ELSE 1 END,
is_read, available_date, id DESC LIMIT ?
""",
(int(user_id), as_of, max(1, min(300, int(limit)))),
).fetchall()
return [{**dict(row), "is_read": bool(row["is_read"])} for row in rows]
def count_unread_alerts(self, user_id: int, as_of: str) -> int:
with self.connect() as connection:
row = connection.execute(
"""
SELECT COUNT(*) AS total FROM alerts
WHERE user_id = ? AND available_date <= ? AND is_read = 0
""",
(int(user_id), as_of),
).fetchone()
return int(row["total"] if row else 0)
def mark_alert_read(self, user_id: int, alert_id: int) -> bool:
now = datetime.now().astimezone().isoformat(timespec="seconds")
with self.connect() as connection:
cursor = connection.execute(
"""
UPDATE alerts SET is_read = 1, read_at = ?, updated_at = ?
WHERE id = ? AND user_id = ?
""",
(now, now, int(alert_id), int(user_id)),
)
return cursor.rowcount > 0
def mark_all_alerts_read(self, user_id: int, as_of: str) -> int:
now = datetime.now().astimezone().isoformat(timespec="seconds")
with self.connect() as connection:
cursor = connection.execute(
"""
UPDATE alerts SET is_read = 1, read_at = ?, updated_at = ?
WHERE user_id = ? AND available_date <= ? AND is_read = 0
""",
(now, now, int(user_id), as_of),
)
return int(cursor.rowcount)
def delete_alert(self, user_id: int, alert_id: int) -> bool:
with self.connect() as connection:
cursor = connection.execute(
"DELETE FROM alerts WHERE id = ? AND user_id = ?",
(int(alert_id), int(user_id)),
)
return cursor.rowcount > 0
def save_trade_entry(
self,
user_id: int,
trade_date: str,
code: str,
name: str,
action: str,
price: float,
quantity: int,
position_pct: float,
pnl_amount: float | None,
pnl_pct: float | None,
thesis: str,
execution: str,
emotion: str,
tags: list[str],
trade_id: int | None = None,
) -> int:
now = datetime.now().astimezone().isoformat(timespec="seconds")
tags_json = json.dumps(tags, ensure_ascii=False, separators=(",", ":"))
with self.connect() as connection:
if trade_id:
cursor = connection.execute(
"""
UPDATE trade_entries SET
trade_date=?, code=?, name=?, action=?, price=?, quantity=?,
position_pct=?, pnl_amount=?, pnl_pct=?, thesis=?, execution=?,
emotion=?, tags=?, updated_at=?
WHERE id=? AND user_id=?
""",
(
trade_date, code, name, action, price, quantity, position_pct,
pnl_amount, pnl_pct, thesis, execution, emotion, tags_json, now,
int(trade_id), int(user_id),
),
)
if cursor.rowcount == 0:
raise ValueError("交易记录不存在或无权修改。")
return int(trade_id)
cursor = connection.execute(
"""
INSERT INTO trade_entries
(user_id, trade_date, code, name, action, price, quantity,
position_pct, pnl_amount, pnl_pct, thesis, execution, emotion,
tags, created_at, updated_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
""",
(
int(user_id), trade_date, code, name, action, price, quantity,
position_pct, pnl_amount, pnl_pct, thesis, execution, emotion,
tags_json, now, now,
),
)
return int(cursor.lastrowid)
def list_trade_entries(
self, user_id: int, start_date: str = "", end_date: str = "", code: str = "",
limit: int = 300,
) -> list[dict[str, Any]]:
clauses = ["user_id = ?"]
parameters: list[Any] = [int(user_id)]
if start_date:
clauses.append("trade_date >= ?")
parameters.append(start_date)
if end_date:
clauses.append("trade_date <= ?")
parameters.append(end_date)
if code:
clauses.append("code = ?")
parameters.append(code)
parameters.append(max(1, min(1000, int(limit))))
with self.connect() as connection:
rows = connection.execute(
f"""
SELECT * FROM trade_entries WHERE {' AND '.join(clauses)}
ORDER BY trade_date DESC, id DESC LIMIT ?
""",
parameters,
).fetchall()
return [dict(row) for row in rows]
def delete_trade_entry(self, user_id: int, trade_id: int) -> bool:
with self.connect() as connection:
cursor = connection.execute(
"DELETE FROM trade_entries WHERE id = ? AND user_id = ?",
(int(trade_id), int(user_id)),
)
return cursor.rowcount > 0
def save_assistant_exchange(
self, user_id: int, question: str, answer: str, context_date: str
) -> None:
now = datetime.now().astimezone().isoformat(timespec="seconds")
with self.connect() as connection:
connection.executemany(
"""
INSERT INTO assistant_messages
(user_id, role, content, context_date, created_at)
VALUES (?, ?, ?, ?, ?)
""",
[
(int(user_id), "user", question, context_date, now),
(int(user_id), "assistant", answer, context_date, now),
],
)
connection.execute(
"""
DELETE FROM assistant_messages WHERE user_id = ? AND id NOT IN (
SELECT id FROM assistant_messages
WHERE user_id = ? ORDER BY id DESC LIMIT 200
)
""",
(int(user_id), int(user_id)),
)
def list_assistant_messages(self, user_id: int, limit: int = 100) -> list[dict[str, Any]]:
with self.connect() as connection:
rows = connection.execute(
"""
SELECT role, content, context_date, created_at FROM assistant_messages
WHERE user_id = ? ORDER BY id DESC LIMIT ?
""",
(int(user_id), max(1, min(200, int(limit)))),
).fetchall()
return [dict(row) for row in reversed(rows)]
def delete_assistant_messages(self, user_id: int) -> int:
with self.connect() as connection:
cursor = connection.execute(
"DELETE FROM assistant_messages WHERE user_id = ?", (int(user_id),)
)
return int(cursor.rowcount)
@staticmethod
def _heaven_reading_dict(row: sqlite3.Row | None) -> dict[str, Any] | None:
if not row:
return None
return {
"id": int(row["id"]),
"mode": str(row["mode"]),
"context_date": str(row["context_date"]),
"subject": str(row["subject"]),
"subject_detail": str(row["subject_detail"]),
"answer": str(row["answer"]),
"created_at": str(row["created_at"]),
}
def save_heaven_reading(
self,
user_id: int,
mode: str,
context_date: str,
subject: str,
subject_detail: str,
answer: str,
context_snapshot: dict[str, Any],
dedupe_key: str,
) -> dict[str, Any]:
now = datetime.now().astimezone().isoformat(timespec="seconds")
snapshot_json = json.dumps(
context_snapshot, ensure_ascii=False, separators=(",", ":")
)
with self.connect() as connection:
connection.execute(
"""
INSERT INTO heaven_readings
(user_id, mode, context_date, subject, subject_detail, answer,
context_snapshot, dedupe_key, created_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(user_id, dedupe_key) DO NOTHING
""",
(
int(user_id), mode, context_date, subject, subject_detail,
answer, snapshot_json, dedupe_key, now,
),
)
row = connection.execute(
"""
SELECT id, mode, context_date, subject, subject_detail, answer, created_at
FROM heaven_readings WHERE user_id = ? AND dedupe_key = ?
""",
(int(user_id), dedupe_key),
).fetchone()
connection.execute(
"""
DELETE FROM heaven_readings
WHERE user_id = ? AND mode = ? AND id NOT IN (
SELECT id FROM heaven_readings
WHERE user_id = ? AND mode = ? ORDER BY id DESC LIMIT 100
)
""",
(int(user_id), mode, int(user_id), mode),
)
result = self._heaven_reading_dict(row)
if not result:
raise ValueError("解读记录保存失败。")
return result
def list_heaven_readings(
self,
user_id: int,
mode: str,
context_date: str = "",
limit: int = 100,
) -> list[dict[str, Any]]:
clauses = ["user_id = ?", "mode = ?"]
parameters: list[Any] = [int(user_id), mode]
if context_date:
clauses.append("context_date = ?")
parameters.append(context_date)
parameters.append(max(1, min(100, int(limit))))
with self.connect() as connection:
rows = connection.execute(
f"""
SELECT id, mode, context_date, subject, subject_detail, answer, created_at
FROM heaven_readings WHERE {' AND '.join(clauses)}
ORDER BY context_date DESC, id DESC LIMIT ?
""",
parameters,
).fetchall()
return [self._heaven_reading_dict(row) for row in rows if row]
def latest_heaven_reading(
self, user_id: int, mode: str, context_date: str = ""
) -> dict[str, Any] | None:
items = self.list_heaven_readings(user_id, mode, context_date, 1)
return items[0] if items else None
def delete_heaven_reading(self, user_id: int, reading_id: int) -> bool:
with self.connect() as connection:
cursor = connection.execute(
"DELETE FROM heaven_readings WHERE id = ? AND user_id = ?",
(int(reading_id), int(user_id)),
)
return cursor.rowcount > 0
+4 -115
View File
@@ -1,118 +1,7 @@
from __future__ import annotations """Compatibility alias for the canonical heaven agent implementation."""
import json import sys
import time
import urllib.error
import urllib.request
from typing import Any
from backend.features.heaven import agent as _implementation
class HeavenAgentError(RuntimeError): sys.modules[__name__] = _implementation
pass
def interpret_heaven(
mode: str,
context: dict[str, Any],
api_key: str,
base_url: str,
model: str,
timeout: int = 90,
) -> dict[str, Any]:
if mode not in {"trend", "fortune", "heart"}:
raise HeavenAgentError("不支持的问天解读模式。")
if not api_key or not model:
raise HeavenAgentError("LLM API Key 或模型尚未配置。")
system_prompt = _system_prompt(mode)
payload = json.dumps(
{
"model": model,
"messages": [
{"role": "system", "content": system_prompt},
{
"role": "user",
"content": json.dumps(context, ensure_ascii=False, separators=(",", ":")),
},
],
"stream": False,
},
ensure_ascii=False,
).encode("utf-8")
request = urllib.request.Request(
f"{base_url.rstrip('/')}/chat/completions",
data=payload,
headers={
"Content-Type": "application/json",
"Authorization": f"Bearer {api_key}",
"User-Agent": "XiaobaiReviewWeb/0.7",
},
method="POST",
)
started = time.perf_counter()
try:
with urllib.request.urlopen(request, timeout=timeout) as response:
result = json.loads(response.read().decode("utf-8"))
answer = str(result["choices"][0]["message"]["content"]).strip()
if not answer:
raise KeyError("empty response")
except urllib.error.HTTPError as exc:
raise HeavenAgentError(_http_error_message(exc)) from exc
except (urllib.error.URLError, TimeoutError, json.JSONDecodeError, KeyError, IndexError) as exc:
raise HeavenAgentError(f"问天模型调用失败:{exc}") from exc
return {
"answer": answer,
"model": model,
"latency_ms": round((time.perf_counter() - started) * 1000),
}
def _system_prompt(mode: str) -> str:
common = """
你是“小白复盘”的问天解读器。所有历法、卦象、爻位和市场指标已经由确定性程序计算,你只能解释提供的数据,不得改卦、改爻、改干支或编造行情。
问天属于传统文化与娱乐化观察,不是预测模型,不承诺应验,不输出无条件买卖指令,不用神秘话术制造确定性。
使用中文,先给核心判断,再解释结构。引用市场数字时标明数据日期。输出纯文本,可使用简短标题。
""".strip()
if mode == "trend":
return common + """
当前任务是“观势·解势”。六爻从初爻到上爻依次是个股内核、个股外显、板块内核、板块外显、指数内核、指数外显;初二为地、三四为人、五上为天。
行情数据只负责生成六爻,本次解势必须以卦象本身为主,不得根据指数涨跌、板块强弱、涨停家数、成交量或个股表现直接推演方向。context中不会提供这些数字,也不会提供爻位对应的市场角色。
先解释本卦卦名的核心义、上下卦组合及大象;再只解释实际动爻所代表的转折,并说明本卦如何走向之卦;最后可把这一组卦势翻译成克制的市场语言。
重点是“本卦为当下之势,动爻为变化关节,之卦为所趋之势”。不要说明某一动爻对应指数、板块或个股,也不要输出“一看指数、二看涨停家数”一类行情观察条件。
全文控制在300至450个中文字符,最多四小段。卦理约占九成,市场翻译最多一句,只能落到节制、等待、守信、辨伪等行为态度,不得据此预测市场下一阶段、涨跌方向或动能变化。不直接荐股,不使用Markdown表格。
不要使用“必然、确定、必涨、必跌、后续将、进入某阶段”等断语;天机只点出势的性质与变化关系,不替用户宣布结果。
""".strip()
if mode == "fortune":
return common + """
当前任务是“观气·解运”。严格区分五运、六气、节气、月令和日干,不把丙午简单解释为火年。
严格服从five_phase_field.framework提供的确定性结构,不自行重新计算五行:年纲由中运与司天在泉构成;岁半以前司天为主、在泉为辅,岁半以后在泉为主、司天为辅;当前六气层以客气加临主气为核心;日辰只负责触发。节气只用于定位当前六气阶段,不得再次叠加为独立力量。
重点解释framework.relations中的客主同气、客生主、主生客、客克主或主克客,以及客胜为从、主胜为逆、司天在泉同位、天符岁会等已经判定的关系。不得把司天、在泉、主气、客气视为彼此独立的证据重复计权,也不得自行增删传统格局。
首要解释当日气场容易放大参与者的哪些情绪、判断偏差和操作冲动,例如急躁、恐惧、迟疑、追涨、过早止损或路径依赖;再给出一至两个调节动作。
如有personal_profile,结合其日主、十神、五行平衡倾向说明当日对该用户主观状态的影响,但不得把简化平衡倾向说成唯一喜用神,也不得复述或猜测出生日期。
不得引用市场上涨下跌家数、涨跌停数量、成交额、板块强度或个股表现来证明气场。industry_affinity只是五行行业取象示例,不是行情旁证;行业契合度最多在末尾用一句话说明,不得写“当日共振”或暗示相关行业必然涨跌。
全文控制在420至600个中文字符,按“三层气机、人的状态、操作偏向、个人影响(如有)、制衡动作”组织,标题必须写“三层气机”。明确这些是传统历法框架下的观察语言,不宣称气候或五行直接导致股价。
""".strip()
return common + """
当前任务是“观心·解卦”。用户的问题始终只在心中,没有输入给你,因此你不能猜测问题内容,也不能替用户作具体决定。
全文控制在180至350个中文字符。只写一句卦意;一小段动爻与之卦;最后三句极短的问心句。
不要重述六条爻辞,不猜用户未说出口的问题,不以吉凶二字替代思考,不给出股票涨跌预测。语气安静、克制,越短越有余味。
""".strip()
def _http_error_message(exc: urllib.error.HTTPError) -> str:
detail = ""
try:
payload = json.loads(exc.read().decode("utf-8", errors="replace"))
error = payload.get("error")
if isinstance(error, dict):
detail = str(error.get("message") or error.get("code") or "")
elif error:
detail = str(error)
elif payload.get("message"):
detail = str(payload["message"])
except (json.JSONDecodeError, OSError):
detail = ""
suffix = f"{detail[:300]}" if detail else ""
return f"问天模型调用失败(HTTP {exc.code}{suffix}"
+3 -1178
View File
File diff suppressed because it is too large Load Diff
+2 -2
View File
@@ -9,11 +9,11 @@ from tushare_client import _sector_coverage_issue
def load_method(name: str): def load_method(name: str):
source = Path("backend/application.py").read_text(encoding="utf-8") source = Path("backend/features/heaven/service.py").read_text(encoding="utf-8")
tree = ast.parse(source) tree = ast.parse(source)
dashboard_service = next( dashboard_service = next(
node for node in tree.body node for node in tree.body
if isinstance(node, ast.ClassDef) and node.name == "DashboardService" if isinstance(node, ast.ClassDef) and node.name == "HeavenServiceMixin"
) )
method = next( method = next(
node for node in dashboard_service.body node for node in dashboard_service.body
+163
View File
@@ -0,0 +1,163 @@
from __future__ import annotations
import ast
import hashlib
import unittest
from pathlib import Path
import heaven_agent
import heaven_engine
from backend.features.heaven import agent as canonical_agent
from backend.features.heaven import engine as canonical_engine
APP_ROOT = Path(__file__).resolve().parents[1]
ORIGINAL_ROOT = APP_ROOT.parent
HEAVEN_SERVICE_METHODS = {
"_heaven_manual_schema",
"_validate_heaven_manual_data",
"_apply_heaven_manual_data",
"_heaven_line_checks",
"_resolve_heaven_stock_code",
"heaven_setup",
"_heaven_stock_context",
"_heaven_market_mode",
"_heaven_trend_sources",
"_heaven_trend_quality_issues",
"heaven_personal",
"heaven_hexagram",
"heaven_readings",
"_heaven_reading_identity",
"heaven_interpret",
"_legacy_truncated_heaven_reading",
"_call_heaven_agent",
"_heaven_index_context",
"_aggregate_index_context",
"_heaven_sector_context",
}
HEAVEN_REPOSITORY_METHODS = {
"_heaven_reading_dict",
"save_heaven_reading",
"list_heaven_readings",
"latest_heaven_reading",
"delete_heaven_reading",
}
HEAVEN_HTTP_METHODS = {"heaven_hexagram", "heaven_personal", "heaven_interpret"}
def sha256(path: Path) -> str:
return hashlib.sha256(path.read_bytes()).hexdigest()
def class_methods(path: Path, class_name: str) -> dict[str, str]:
tree = ast.parse(path.read_text(encoding="utf-8"), filename=str(path))
owner = next(
node
for node in tree.body
if isinstance(node, ast.ClassDef) and node.name == class_name
)
return {
node.name: ast.dump(node, include_attributes=False)
for node in owner.body
if isinstance(node, (ast.FunctionDef, ast.AsyncFunctionDef))
}
def top_level_definitions(path: Path) -> dict[str, str]:
tree = ast.parse(path.read_text(encoding="utf-8"), filename=str(path))
return {
node.name: ast.dump(node, include_attributes=False)
for node in tree.body
if isinstance(node, (ast.FunctionDef, ast.AsyncFunctionDef, ast.ClassDef))
}
class HeavenSliceSourceEquivalenceTests(unittest.TestCase):
def assert_methods_equal(
self,
original_path: Path,
original_class: str,
migrated_path: Path,
migrated_class: str,
expected: set[str],
adapted: set[str] | None = None,
) -> None:
original = class_methods(original_path, original_class)
migrated = class_methods(migrated_path, migrated_class)
self.assertEqual(set(migrated), expected)
for name in sorted(expected - (adapted or set())):
self.assertEqual(migrated[name], original[name], name)
def test_heaven_agent_is_an_exact_file(self) -> None:
self.assertEqual(
sha256(ORIGINAL_ROOT / "heaven_agent.py"),
sha256(APP_ROOT / "backend" / "features" / "heaven" / "agent.py"),
)
def test_heaven_engine_definitions_are_exact_original_ast(self) -> None:
self.assertEqual(
top_level_definitions(ORIGINAL_ROOT / "heaven_engine.py"),
top_level_definitions(
APP_ROOT / "backend" / "features" / "heaven" / "engine.py"
),
)
def test_compatibility_modules_are_canonical_module_objects(self) -> None:
self.assertIs(heaven_agent, canonical_agent)
self.assertIs(heaven_engine, canonical_engine)
def test_heaven_service_methods_are_exact_original_ast(self) -> None:
self.assert_methods_equal(
ORIGINAL_ROOT / "server.py",
"DashboardService",
APP_ROOT / "backend" / "features" / "heaven" / "service.py",
"HeavenServiceMixin",
HEAVEN_SERVICE_METHODS,
{"_heaven_reading_identity"},
)
source = (
APP_ROOT / "backend" / "features" / "heaven" / "service.py"
).read_text(encoding="utf-8")
self.assertIn(
"MarketServiceMixin._display_compact_date(context_date)", source
)
def test_heaven_repository_methods_are_exact_original_ast(self) -> None:
self.assert_methods_equal(
ORIGINAL_ROOT / "database.py",
"ReviewDatabase",
APP_ROOT / "backend" / "features" / "heaven" / "repository.py",
"HeavenRepositoryMixin",
HEAVEN_REPOSITORY_METHODS,
)
def test_original_classes_no_longer_duplicate_moved_methods(self) -> None:
remaining_service = class_methods(
APP_ROOT / "backend" / "application.py", "DashboardService"
)
remaining_database = class_methods(APP_ROOT / "database.py", "ReviewDatabase")
remaining_http = class_methods(
APP_ROOT / "backend" / "application.py", "RequestHandler"
)
self.assertTrue(HEAVEN_SERVICE_METHODS.isdisjoint(remaining_service))
self.assertTrue(HEAVEN_REPOSITORY_METHODS.isdisjoint(remaining_database))
self.assertTrue(HEAVEN_HTTP_METHODS.isdisjoint(remaining_http))
def test_http_mixin_preserves_all_heaven_endpoints(self) -> None:
methods = class_methods(
APP_ROOT / "backend" / "features" / "heaven" / "http.py",
"HeavenHttpMixin",
)
self.assertEqual(set(methods), HEAVEN_HTTP_METHODS)
source = (
APP_ROOT / "backend" / "features" / "heaven" / "http.py"
).read_text(encoding="utf-8")
self.assertNotIn("SERVICE.", source)
self.assertEqual(source.count("self.application_service.heaven_"), 3)
if __name__ == "__main__":
unittest.main()
@@ -0,0 +1,211 @@
from __future__ import annotations
import ast
import hashlib
import unittest
from pathlib import Path
import assistant_agent
from backend.features.review import agent as canonical_agent
APP_ROOT = Path(__file__).resolve().parents[1]
ORIGINAL_ROOT = APP_ROOT.parent
ALERT_SERVICE_METHODS = {
"alert_center",
"create_alert",
"mark_alert_read",
"mark_all_alerts_read",
"delete_alert",
}
REVIEW_SERVICE_METHODS = {
"trade_entries",
"review_watchlist",
"save_trade_entry",
"delete_trade_entry",
"assistant_messages",
"clear_assistant_messages",
"assistant_stream",
"_assistant_context",
}
ALERT_REPOSITORY_METHODS = {
"save_alert",
"list_alerts",
"count_unread_alerts",
"mark_alert_read",
"mark_all_alerts_read",
"delete_alert",
}
REVIEW_REPOSITORY_METHODS = {
"list_watchlist",
"save_watchlist",
"watchlist_price_history",
"delete_watchlist",
"list_notes",
"save_note",
"delete_note",
"save_trade_entry",
"list_trade_entries",
"delete_trade_entry",
"save_assistant_exchange",
"list_assistant_messages",
"delete_assistant_messages",
}
ALERT_HTTP_METHODS = {"save_alert"}
REVIEW_HTTP_METHODS = {
"save_trade_entry",
"stream_assistant_chat",
"save_watchlist",
"save_note",
}
def sha256(path: Path) -> str:
return hashlib.sha256(path.read_bytes()).hexdigest()
def class_methods(path: Path, class_name: str) -> dict[str, str]:
tree = ast.parse(path.read_text(encoding="utf-8"), filename=str(path))
owner = next(
node
for node in tree.body
if isinstance(node, ast.ClassDef) and node.name == class_name
)
return {
node.name: ast.dump(node, include_attributes=False)
for node in owner.body
if isinstance(node, (ast.FunctionDef, ast.AsyncFunctionDef))
}
def top_level_definition(path: Path, name: str) -> str:
tree = ast.parse(path.read_text(encoding="utf-8"), filename=str(path))
node = next(
item
for item in tree.body
if isinstance(item, (ast.FunctionDef, ast.AsyncFunctionDef, ast.ClassDef))
and item.name == name
)
return ast.dump(node, include_attributes=False)
class ReviewAlertsSliceSourceEquivalenceTests(unittest.TestCase):
def assert_methods_equal(
self,
original_path: Path,
original_class: str,
migrated_path: Path,
migrated_class: str,
expected: set[str],
) -> None:
original = class_methods(original_path, original_class)
migrated = class_methods(migrated_path, migrated_class)
self.assertEqual(set(migrated), expected)
for name in sorted(expected):
self.assertEqual(migrated[name], original[name], name)
def test_review_assistant_agent_is_an_exact_file(self) -> None:
self.assertEqual(
sha256(ORIGINAL_ROOT / "assistant_agent.py"),
sha256(APP_ROOT / "backend" / "features" / "review" / "agent.py"),
)
def test_review_assistant_compatibility_module_is_canonical(self) -> None:
self.assertIs(assistant_agent, canonical_agent)
def test_alert_and_review_service_methods_are_exact_original_ast(self) -> None:
self.assert_methods_equal(
ORIGINAL_ROOT / "server.py",
"DashboardService",
APP_ROOT / "backend" / "features" / "alerts" / "facade.py",
"AlertServiceMixin",
ALERT_SERVICE_METHODS,
)
self.assert_methods_equal(
ORIGINAL_ROOT / "server.py",
"DashboardService",
APP_ROOT / "backend" / "features" / "review" / "service.py",
"ReviewServiceMixin",
REVIEW_SERVICE_METHODS,
)
def test_alert_and_review_repository_methods_are_exact_original_ast(self) -> None:
self.assert_methods_equal(
ORIGINAL_ROOT / "database.py",
"ReviewDatabase",
APP_ROOT / "backend" / "features" / "alerts" / "repository.py",
"AlertRepositoryMixin",
ALERT_REPOSITORY_METHODS,
)
self.assert_methods_equal(
ORIGINAL_ROOT / "database.py",
"ReviewDatabase",
APP_ROOT / "backend" / "features" / "review" / "repository.py",
"ReviewRepositoryMixin",
REVIEW_REPOSITORY_METHODS,
)
def test_existing_alert_and_trade_journal_services_preserve_original_classes(self) -> None:
self.assertEqual(
top_level_definition(
ORIGINAL_ROOT / "backend" / "features" / "alerts" / "service.py",
"AlertService",
),
top_level_definition(
APP_ROOT / "backend" / "features" / "alerts" / "service.py",
"AlertService",
),
)
self.assertEqual(
top_level_definition(
ORIGINAL_ROOT / "backend" / "features" / "review" / "trade_journal.py",
"TradeJournalService",
),
top_level_definition(
APP_ROOT / "backend" / "features" / "review" / "trade_journal.py",
"TradeJournalService",
),
)
def test_original_classes_no_longer_duplicate_moved_methods(self) -> None:
remaining_service = class_methods(
APP_ROOT / "backend" / "application.py", "DashboardService"
)
remaining_database = class_methods(APP_ROOT / "database.py", "ReviewDatabase")
remaining_http = class_methods(
APP_ROOT / "backend" / "application.py", "RequestHandler"
)
self.assertTrue(ALERT_SERVICE_METHODS.isdisjoint(remaining_service))
self.assertTrue(REVIEW_SERVICE_METHODS.isdisjoint(remaining_service))
self.assertTrue(ALERT_REPOSITORY_METHODS.isdisjoint(remaining_database))
self.assertTrue(REVIEW_REPOSITORY_METHODS.isdisjoint(remaining_database))
self.assertTrue(ALERT_HTTP_METHODS.isdisjoint(remaining_http))
self.assertTrue(REVIEW_HTTP_METHODS.isdisjoint(remaining_http))
def test_http_mixins_preserve_all_endpoints_without_global_service(self) -> None:
alerts = class_methods(
APP_ROOT / "backend" / "features" / "alerts" / "http.py",
"AlertHttpMixin",
)
review = class_methods(
APP_ROOT / "backend" / "features" / "review" / "http.py",
"ReviewHttpMixin",
)
self.assertEqual(set(alerts), ALERT_HTTP_METHODS)
self.assertEqual(set(review), REVIEW_HTTP_METHODS)
for path in (
APP_ROOT / "backend" / "features" / "alerts" / "http.py",
APP_ROOT / "backend" / "features" / "review" / "http.py",
):
source = path.read_text(encoding="utf-8")
self.assertNotIn("SERVICE.", source)
self.assertIn("self.application_service.", source)
if __name__ == "__main__":
unittest.main()
@@ -0,0 +1,69 @@
# 切片 08:问天、观势、观气与观心
> 基线:`2919229`(切片 07
> 回档标签:`xiaobai-preservation-slice-08-20260731`
> 结论:源码、API、数据库、真实页面、动画与全量回归通过;最终视觉仍等待全站人工验收
## 1. 原实现归位
本切片只移动原版问天 Agent、历法/卦象引擎、服务、持久化与三个 HTTP 入口,没有从
`next/`取用代码,也没有改写六爻、安全门、五运六气、个人合参、铜钱起卦、历史去重或
LLM 解读逻辑。
| 原位置 | 新的唯一实现位置 | 兼容方式 |
|---|---|---|
| `app/heaven_agent.py` | `app/backend/features/heaven/agent.py` | 根级模块指向同一模块对象 |
| `app/heaven_engine.py` | `app/backend/features/heaven/engine.py` | 根级模块指向同一模块对象 |
| `DashboardService`问天方法 | `app/backend/features/heaven/service.py` | `HeavenServiceMixin` |
| 问天历史持久化 | `app/backend/features/heaven/repository.py` | `HeavenRepositoryMixin` |
| 三个问天 POST 入口 | `app/backend/features/heaven/http.py` | `HeavenHttpMixin` |
历法引擎搬迁后仍从`app/vendor``app/data/iching_zh.json`读取原资源;这是唯一资源路径适配,
顶层函数与类定义未变。旧测试依赖的根级模块继续有效。
## 2. 源码等价
- `heaven_agent.py`与原版文件 SHA-256 一致。
- `heaven_engine.py`所有顶层函数和类与原版无位置信息 AST 一致。
- 20个问天服务方法中19个与原版 AST 一致;`_heaven_reading_identity`仅将旧巨型类名
`DashboardService`替换为已归位的`MarketServiceMixin`静态日期格式化方法。
- 5个问天 Repository 方法与原版 AST 一致。
- 三个 HTTP 方法只把全局`SERVICE`改为 Mixin 的`self.application_service`,状态码、字段和异常语义不变。
- `DashboardService``ReviewDatabase``RequestHandler`不再重复保留已迁移实现。
## 3. API 与数据库差分
- 原版`8788`和迁移版`8789`使用同一正式数据库的两个临时副本。
- 固定日期的问天初始化、三类历史读取、六爻成卦和个人五行计算共6个接口,状态码与业务 JSON 完全一致。
- 未调用真实`/api/heaven/interpret`,避免外部模型波动及测试记录写入;提示词和调用链由源码等价与单元测试覆盖。
- 两版均为62个 schema 对象;`heaven_readings``users``system_settings``llm_usage`
`user_birth_profiles``sector_phase_overrides`逐行一致。
- 差分只在系统临时目录副本运行,正式数据库未写入问天记录或 LLM 用量。
## 4. 真实浏览器与动画检查
- 1920×1080视口逐页检查观势、观气、观心,无横向溢出,页面滚动边界可用。
- 观势保留星空、三才六爻与八卦待载入场景;观气生成100个星点,五运六气环动画名为`wt-spin`
五行行业默认折叠;观心呼吸波纹、香火燃烧和阶段文字均持续运行。
- 夜间模式问天背景为`rgb(11, 17, 32)`,全页无控制台错误。
- 没有触发真实解势、解运或解卦调用,加载动画的静态资产及脚本未改动。
- 截图 SHA-256
- `heaven-trend-1080.png``9e70d9cbfd741e74d3c4511c0e76ef8e625a6a31fd5d6bde014103f807bcc76e`
- `heaven-fortune-1080.png``e4cdf0e0fc1d9325724654a781febed1d509e49ddccecb3355803b5e9f1d2f5c`
- `heaven-heart-breathing-1080.png``3493a59b9032c1a8a1ba7da609b499c736170d2c7947b3aa0ac4c8b95d981ff5`
- `heaven-heart-dark-1080.png``3d346fea6b593fa4e43037f5fc68eda1a3f3446a223c2022d9084922f646a088`
## 5. 自动验证与保留边界
| 验证 | 结果 |
|---|---:|
| 原版`python -m unittest discover -s tests -q` | 231项通过 |
| 迁移版`python -m unittest discover -s tests -q` | 280项通过 |
| `python -m unittest tests.test_preservation_slice_heaven -q` | 7项通过 |
| `npx.cmd playwright test --reporter=dot` | 45项通过(2.2分钟) |
| `git diff --check` | 通过 |
- 复盘助手、交易日志、自选、笔记与提醒属于切片09,本切片不提前移动。
- 前端 DOM、页面 JS、CSS和移动端行为未改动,统一归档延至切片10。
- `static/heaven-loading.js`是否失效仍不确定,继续保留到切片11试删。
- 没有删除待定代码、没有修改正式数据库、没有测试或切换 Docker/NAS。
@@ -0,0 +1,71 @@
{
"all_equal": true,
"endpoints": [
{
"name": "heaven setup awaiting stock selection",
"method": "GET",
"endpoint": "/api/heaven/setup?trade_date=2026-07-29",
"original_status": 200,
"migrated_status": 200,
"original_sha256": "b72f2e972c3993717d639d1ccc68a9f1fb716a5846823a485a05458103465b92",
"migrated_sha256": "b72f2e972c3993717d639d1ccc68a9f1fb716a5846823a485a05458103465b92",
"equal": true,
"first_difference": null
},
{
"name": "trend interpretation history",
"method": "GET",
"endpoint": "/api/heaven/readings?mode=trend&limit=20",
"original_status": 200,
"migrated_status": 200,
"original_sha256": "2aa83ea177bd5cb2351ac4f3a60fdcd17e02cfc14b5753e157e515f03fc2c6c1",
"migrated_sha256": "2aa83ea177bd5cb2351ac4f3a60fdcd17e02cfc14b5753e157e515f03fc2c6c1",
"equal": true,
"first_difference": null
},
{
"name": "fortune interpretation history",
"method": "GET",
"endpoint": "/api/heaven/readings?mode=fortune&limit=20",
"original_status": 200,
"migrated_status": 200,
"original_sha256": "e60debaad54c58ad34fc7f025bfeaba7d2349aee2531b3f507c6e1968dbdcbf0",
"migrated_sha256": "e60debaad54c58ad34fc7f025bfeaba7d2349aee2531b3f507c6e1968dbdcbf0",
"equal": true,
"first_difference": null
},
{
"name": "heart interpretation history",
"method": "GET",
"endpoint": "/api/heaven/readings?mode=heart&limit=20",
"original_status": 200,
"migrated_status": 200,
"original_sha256": "1fe13531909c83c645b4fbb6740f46f2590adb346ab6e43414023967bb51b072",
"migrated_sha256": "1fe13531909c83c645b4fbb6740f46f2590adb346ab6e43414023967bb51b072",
"equal": true,
"first_difference": null
},
{
"name": "six-line hexagram calculation",
"method": "POST",
"endpoint": "/api/heaven/hexagram",
"original_status": 200,
"migrated_status": 200,
"original_sha256": "9773a9f288ea5fbe063b778fc47c613972fd7747f8fc6f80333316f33923c960",
"migrated_sha256": "9773a9f288ea5fbe063b778fc47c613972fd7747f8fc6f80333316f33923c960",
"equal": true,
"first_difference": null
},
{
"name": "personal five-phase calculation",
"method": "POST",
"endpoint": "/api/heaven/personal",
"original_status": 200,
"migrated_status": 200,
"original_sha256": "3b179ea4e3e0b43d12e929b4057a51264772fadcf41e5c87573dc5d18873bc9e",
"migrated_sha256": "3b179ea4e3e0b43d12e929b4057a51264772fadcf41e5c87573dc5d18873bc9e",
"equal": true,
"first_difference": null
}
]
}
@@ -0,0 +1,42 @@
[
{
"name": "heaven setup awaiting stock selection",
"method": "GET",
"endpoint": "/api/heaven/setup?trade_date=2026-07-29",
"payload": null
},
{
"name": "trend interpretation history",
"method": "GET",
"endpoint": "/api/heaven/readings?mode=trend&limit=20",
"payload": null
},
{
"name": "fortune interpretation history",
"method": "GET",
"endpoint": "/api/heaven/readings?mode=fortune&limit=20",
"payload": null
},
{
"name": "heart interpretation history",
"method": "GET",
"endpoint": "/api/heaven/readings?mode=heart&limit=20",
"payload": null
},
{
"name": "six-line hexagram calculation",
"method": "POST",
"endpoint": "/api/heaven/hexagram",
"payload": {
"lines": [7, 8, 9, 6, 7, 8]
}
},
{
"name": "personal five-phase calculation",
"method": "POST",
"endpoint": "/api/heaven/personal",
"payload": {
"trade_date": "2026-07-29"
}
}
]
@@ -0,0 +1,59 @@
{
"all_equal": true,
"schema": {
"object_count": 62,
"original_sha256": "60a4e044f0ddb6502e44b45f65bedc1a0a31d4bd596f386d4ccfe153a6c8ddd1",
"migrated_sha256": "60a4e044f0ddb6502e44b45f65bedc1a0a31d4bd596f386d4ccfe153a6c8ddd1",
"equal": true
},
"tables": [
{
"table": "heaven_readings",
"original_count": 31,
"migrated_count": 31,
"original_sha256": "e6cfc86f010fb182d7522bb99e5d268e826f93896035172120df51a011590951",
"migrated_sha256": "e6cfc86f010fb182d7522bb99e5d268e826f93896035172120df51a011590951",
"equal": true
},
{
"table": "users",
"original_count": 3,
"migrated_count": 3,
"original_sha256": "4a0f135bef8ebae454693d3f40e1157d84d814d18d849f33e2e760ed1d8a7f89",
"migrated_sha256": "4a0f135bef8ebae454693d3f40e1157d84d814d18d849f33e2e760ed1d8a7f89",
"equal": true
},
{
"table": "system_settings",
"original_count": 1,
"migrated_count": 1,
"original_sha256": "86d333a5c7feaf7111cd79b2db29513326a9b0d3781607f78bf633579e2ca8e9",
"migrated_sha256": "86d333a5c7feaf7111cd79b2db29513326a9b0d3781607f78bf633579e2ca8e9",
"equal": true
},
{
"table": "llm_usage",
"original_count": 72,
"migrated_count": 72,
"original_sha256": "03a070316a125cef904bbfb2bb06b06e792242d5713142e70354402c741393c3",
"migrated_sha256": "03a070316a125cef904bbfb2bb06b06e792242d5713142e70354402c741393c3",
"equal": true
},
{
"table": "user_birth_profiles",
"original_count": 1,
"migrated_count": 1,
"original_sha256": "2a74d4b2edd3a730a46f17fb6ab3926575c42662c1d968643b0d54586ab68b13",
"migrated_sha256": "2a74d4b2edd3a730a46f17fb6ab3926575c42662c1d968643b0d54586ab68b13",
"equal": true
},
{
"table": "sector_phase_overrides",
"original_count": 0,
"migrated_count": 0,
"original_sha256": "4f53cda18c2baa0c0354bb5f9a3ecbe5ed12ab4d8e11ba873c2f11161202b945",
"migrated_sha256": "4f53cda18c2baa0c0354bb5f9a3ecbe5ed12ab4d8e11ba873c2f11161202b945",
"equal": true
}
]
}
Binary file not shown.

After

Width:  |  Height:  |  Size: 105 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 72 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 71 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 76 KiB

@@ -0,0 +1,60 @@
# 切片 09:我的复盘、自选、笔记、交易日志、提醒与复盘助手
> 基线:`b3df070`(切片 08
> 回档标签:`xiaobai-preservation-slice-09-20260731`
> 结论:源码、API、数据库、账户隔离、真实页面及全量回归通过;最终视觉仍等待全站人工验收
## 1. 原实现归位
本切片只机械移动原版复盘、提醒和复盘助手实现,没有从`next/`取用代码,也没有修改自选、笔记、交易日志、提醒、会员灰化、LLM 上下文或流式协议。
| 原位置 | 新的唯一实现位置 | 兼容方式 |
|---|---|---|
| `assistant_agent.py` | `app/backend/features/review/agent.py` | 根级模块指向同一模块对象 |
| `DashboardService`提醒方法 | `app/backend/features/alerts/facade.py` | `AlertServiceMixin` |
| `DashboardService`复盘方法 | `app/backend/features/review/service.py` | `ReviewServiceMixin` |
| `ReviewDatabase`提醒方法 | `app/backend/features/alerts/repository.py` | `AlertRepositoryMixin` |
| `ReviewDatabase`复盘方法 | `app/backend/features/review/repository.py` | `ReviewRepositoryMixin` |
| 五个 POST/流式入口 | `app/backend/features/alerts/http.py``review/http.py` | HTTP Mixin |
原版已经存在的`AlertService``TradeJournalService`、Repository Port 与 SQLite Adapter 继续保留为唯一底层实现。本切片没有制造第二套提醒或交易日志服务。
## 2. 源码等价与必要适配
- `assistant_agent.py`与正式 Agent 文件 SHA-256 一致,根级兼容模块与正式 Agent 是同一模块对象。
- 5 个提醒编排方法、8 个复盘编排方法、6 个提醒 Repository 方法和 13 个复盘 Repository 方法均与原版无位置信息 AST 一致。
- 五个 HTTP 方法只把全局`SERVICE`改为 Mixin 的`self.application_service`;状态码、字段、异常和流式事件写法不变。
- `backend.features.alerts``backend.features.review`仅对旧底层服务采用延迟导出,避免`database -> feature package -> sqlite adapter -> database`初始化循环;外部`from ... import AlertService/TradeJournalService`接口不变。
- 公共`RequestHandler._write_stream_event`继续保留,因为问师和复盘助手共同使用,且现有测试直接覆盖;统一传输层留到切片 10。
## 3. API 与数据库差分
- 原版`8790`和迁移版`8791`使用正式数据库的两个临时副本,未写正式数据库。
- 固定日期下自选、全部复盘笔记、交易日志、提醒中心及复盘助手历史共 5 个鉴权接口,状态码与业务 JSON 完全一致。
- 两版均为 62 个 schema 对象;`watchlist``review_notes``trade_entries``alerts``assistant_messages``users``system_settings``llm_usage`逐行一致。
- 真实浏览器在迁移版临时副本中新增一条交易记录,弹窗正确关闭并显示短提示“交易记录已保存”;没有写正式数据库。
- 未调用真实`/api/assistant/chat`,避免外部模型波动和 LLM 用量写入;Agent、流式累积与权限路径由源码等价和单元测试覆盖。
## 4. 真实浏览器检查
- 1920×1080 检查我的复盘完整布局,自选 6 只、交易日志、每日复盘三个独立输入框和最近复盘均可访问,无横向溢出。
- 添加自选与交易日志弹窗居中可用;交易日志保存后没有出现超长空弹窗。
- 提醒中心正常展示筛选、新建提醒和提醒记录。
- 管理员会员账号显示完整可用的复盘助手;临时普通账号显示相同界面、顶部“复盘助手仅对会员开放”,下方快捷问题、输入框和发送按钮均灰化禁用。
- 临时普通账号自选和交易日志均为空,证明用户私有数据没有串到管理员账号。
- 夜间模式背景为`rgb(18, 20, 22)`,1920×1080 页面无横向溢出,控制台无错误。
- 截图`review-workspace-night-1920x1080.png` SHA-256`bc803305a391930a0aafd68062bab6de86eca14612799eebfda220cf2ad0798e`
## 5. 自动验证与保留边界
| 验证 | 结果 |
|---|---:|
| 原版`python -m unittest discover -s tests -q` | 231 项通过 |
| 迁移版`python -m unittest discover -s tests -q` | 287 项通过 |
| `python -m unittest tests.test_preservation_slice_review_alerts -q` | 7 项通过 |
| `npx.cmd playwright test --reporter=dot` | 45 项通过(2.2 分钟) |
| `git diff --check` | 通过 |
- 前端 DOM、页面 JS、CSS 与移动端行为未改动,统一归位延至切片 10。
- `demo_data.py``static/heaven-loading.js`、五个疑似无引用前端函数和`wencai_saved_queries`继续保留到切片 11 试删。
- 没有修改`next/`、正式数据库、Docker/NAS 或`8765`服务。
@@ -0,0 +1,60 @@
{
"all_equal": true,
"endpoints": [
{
"name": "review watchlist",
"method": "GET",
"endpoint": "/api/watchlist?trade_date=2026-07-30",
"original_status": 200,
"migrated_status": 200,
"original_sha256": "29c001cc6859312fd90082d0fb83818d263827e28ec21157b94c415b71c45078",
"migrated_sha256": "29c001cc6859312fd90082d0fb83818d263827e28ec21157b94c415b71c45078",
"equal": true,
"first_difference": null
},
{
"name": "all review notes",
"method": "GET",
"endpoint": "/api/notes?scope=all",
"original_status": 200,
"migrated_status": 200,
"original_sha256": "a2a1f08610d017de0e46080a1268f72b734eeb80f9cada311d7e03eb1416656a",
"migrated_sha256": "a2a1f08610d017de0e46080a1268f72b734eeb80f9cada311d7e03eb1416656a",
"equal": true,
"first_difference": null
},
{
"name": "trade journal through fixed date",
"method": "GET",
"endpoint": "/api/trades?end_date=2026-07-30",
"original_status": 200,
"migrated_status": 200,
"original_sha256": "77287b34046e22748e94e61ead841e157f6b63f08de1b2fe61f5d715b2e989a4",
"migrated_sha256": "77287b34046e22748e94e61ead841e157f6b63f08de1b2fe61f5d715b2e989a4",
"equal": true,
"first_difference": null
},
{
"name": "alert center through fixed date",
"method": "GET",
"endpoint": "/api/alerts?status=all&as_of=2026-07-30",
"original_status": 200,
"migrated_status": 200,
"original_sha256": "8807772f2b611ea8b95033aa19cae84f08b08e92d76fff8e3cd5b79966ea04f4",
"migrated_sha256": "8807772f2b611ea8b95033aa19cae84f08b08e92d76fff8e3cd5b79966ea04f4",
"equal": true,
"first_difference": null
},
{
"name": "review assistant history",
"method": "GET",
"endpoint": "/api/assistant/messages",
"original_status": 200,
"migrated_status": 200,
"original_sha256": "eef46741adfc3a9f76294d3b78f37a45f113092ac9d44ee77c7a038a88ff09a1",
"migrated_sha256": "eef46741adfc3a9f76294d3b78f37a45f113092ac9d44ee77c7a038a88ff09a1",
"equal": true,
"first_difference": null
}
]
}
@@ -0,0 +1,32 @@
[
{
"name": "review watchlist",
"method": "GET",
"endpoint": "/api/watchlist?trade_date=2026-07-30",
"payload": null
},
{
"name": "all review notes",
"method": "GET",
"endpoint": "/api/notes?scope=all",
"payload": null
},
{
"name": "trade journal through fixed date",
"method": "GET",
"endpoint": "/api/trades?end_date=2026-07-30",
"payload": null
},
{
"name": "alert center through fixed date",
"method": "GET",
"endpoint": "/api/alerts?status=all&as_of=2026-07-30",
"payload": null
},
{
"name": "review assistant history",
"method": "GET",
"endpoint": "/api/assistant/messages",
"payload": null
}
]
@@ -0,0 +1,75 @@
{
"all_equal": true,
"schema": {
"object_count": 62,
"original_sha256": "60a4e044f0ddb6502e44b45f65bedc1a0a31d4bd596f386d4ccfe153a6c8ddd1",
"migrated_sha256": "60a4e044f0ddb6502e44b45f65bedc1a0a31d4bd596f386d4ccfe153a6c8ddd1",
"equal": true
},
"tables": [
{
"table": "watchlist",
"original_count": 6,
"migrated_count": 6,
"original_sha256": "070161f5cdcc33addf289f11cb0ae54bfc0d669a39426fc33398b64512dba2fe",
"migrated_sha256": "070161f5cdcc33addf289f11cb0ae54bfc0d669a39426fc33398b64512dba2fe",
"equal": true
},
{
"table": "review_notes",
"original_count": 3,
"migrated_count": 3,
"original_sha256": "86ed485c2ce922e3434f8ac742d38431ef430b5c22024cacb6f64f0e8e71ee98",
"migrated_sha256": "86ed485c2ce922e3434f8ac742d38431ef430b5c22024cacb6f64f0e8e71ee98",
"equal": true
},
{
"table": "trade_entries",
"original_count": 0,
"migrated_count": 0,
"original_sha256": "4f53cda18c2baa0c0354bb5f9a3ecbe5ed12ab4d8e11ba873c2f11161202b945",
"migrated_sha256": "4f53cda18c2baa0c0354bb5f9a3ecbe5ed12ab4d8e11ba873c2f11161202b945",
"equal": true
},
{
"table": "alerts",
"original_count": 2,
"migrated_count": 2,
"original_sha256": "37b46fb4a754cf9c6179b1be88eed7f5a5398ddfc87160f41ded622ffc313c88",
"migrated_sha256": "37b46fb4a754cf9c6179b1be88eed7f5a5398ddfc87160f41ded622ffc313c88",
"equal": true
},
{
"table": "assistant_messages",
"original_count": 0,
"migrated_count": 0,
"original_sha256": "4f53cda18c2baa0c0354bb5f9a3ecbe5ed12ab4d8e11ba873c2f11161202b945",
"migrated_sha256": "4f53cda18c2baa0c0354bb5f9a3ecbe5ed12ab4d8e11ba873c2f11161202b945",
"equal": true
},
{
"table": "users",
"original_count": 3,
"migrated_count": 3,
"original_sha256": "4a0f135bef8ebae454693d3f40e1157d84d814d18d849f33e2e760ed1d8a7f89",
"migrated_sha256": "4a0f135bef8ebae454693d3f40e1157d84d814d18d849f33e2e760ed1d8a7f89",
"equal": true
},
{
"table": "system_settings",
"original_count": 1,
"migrated_count": 1,
"original_sha256": "86d333a5c7feaf7111cd79b2db29513326a9b0d3781607f78bf633579e2ca8e9",
"migrated_sha256": "86d333a5c7feaf7111cd79b2db29513326a9b0d3781607f78bf633579e2ca8e9",
"equal": true
},
{
"table": "llm_usage",
"original_count": 72,
"migrated_count": 72,
"original_sha256": "03a070316a125cef904bbfb2bb06b06e792242d5713142e70354402c741393c3",
"migrated_sha256": "03a070316a125cef904bbfb2bb06b06e792242d5713142e70354402c741393c3",
"equal": true
}
]
}
Binary file not shown.

After

Width:  |  Height:  |  Size: 98 KiB

+5 -5
View File
@@ -1,6 +1,6 @@
{ {
"schema_version": 1, "schema_version": 1,
"updated_at": "2026-07-31T04:17:23+08:00", "updated_at": "2026-07-31T11:52:29+08:00",
"status": "active", "status": "active",
"migration_mode": "behavior_preserving_source_migration", "migration_mode": "behavior_preserving_source_migration",
"source_of_truth": "current_original_webapp_runtime_and_source", "source_of_truth": "current_original_webapp_runtime_and_source",
@@ -9,10 +9,10 @@
"failed_roots": [ "failed_roots": [
"next" "next"
], ],
"current_slice": "slice-08-heaven-trend-fortune-heart", "current_slice": "slice-10-frontend-shell-pages-components-css-mobile",
"last_completed_slice": "slice-07-mentor-skills-llm-streaming", "last_completed_slice": "slice-09-review-watchlist-notes-journal-alerts-assistant",
"last_checkpoint": "xiaobai-preservation-slice-07-20260731", "last_checkpoint": "xiaobai-preservation-slice-09-20260731",
"next_action": "capture_slice-08_heaven_trend_fortune_heart_and_animation_contracts_then_move_original_implementations", "next_action": "capture_frontend_asset_dependency_map_then_relocate_original_shell_pages_components_css_and_mobile_responsibilities_without_runtime_changes",
"authoritative_documents": [ "authoritative_documents": [
"AGENTS.md", "AGENTS.md",
"docs/migration/原版保真迁移总纲.md", "docs/migration/原版保真迁移总纲.md",
+27 -1
View File
@@ -1,6 +1,6 @@
# 小白复盘保真迁移账本 # 小白复盘保真迁移账本
> 当前状态:正式迁移,切片07“问师、模型 Skill 与 LLM 流式链路”已完成 > 当前状态:正式迁移,切片09“复盘、自选、交易日志、提醒与复盘助手”已完成
本账本是上下文恢复和人工审计的连续记录。任何迁移提交必须在同一提交中更新本文件及 本账本是上下文恢复和人工审计的连续记录。任何迁移提交必须在同一提交中更新本文件及
`保真迁移状态.json` `保真迁移状态.json`
@@ -27,6 +27,7 @@
| 2026-07-31 | `xiaobai-preservation-slice-05-20260731` | 集合竞价、题材库、人气热榜与龙虎榜原实现归位 | 自动、API、数据库与浏览器差分通过,进入切片06 | | 2026-07-31 | `xiaobai-preservation-slice-05-20260731` | 集合竞价、题材库、人气热榜与龙虎榜原实现归位 | 自动、API、数据库与浏览器差分通过,进入切片06 |
| 2026-07-31 | `xiaobai-preservation-slice-06-20260731` | 智能选股、自定义选股与策略持续跟踪原实现归位 | 自动、API、数据库与浏览器差分通过,进入切片07 | | 2026-07-31 | `xiaobai-preservation-slice-06-20260731` | 智能选股、自定义选股与策略持续跟踪原实现归位 | 自动、API、数据库与浏览器差分通过,进入切片07 |
| 2026-07-31 | `xiaobai-preservation-slice-07-20260731` | 问师、模型Skill与LLM流式链路原实现归位 | 自动、API、数据库、Skill与浏览器差分通过,进入切片08 | | 2026-07-31 | `xiaobai-preservation-slice-07-20260731` | 问师、模型Skill与LLM流式链路原实现归位 | 自动、API、数据库、Skill与浏览器差分通过,进入切片08 |
| 2026-07-31 | `xiaobai-preservation-slice-08-20260731` | 问天、观势、观气与观心原实现归位 | 自动、API、数据库、动画与浏览器差分通过,进入切片09 |
## 资产处置登记 ## 资产处置登记
@@ -58,6 +59,11 @@
| 问师消息与偏好方法 | 持久化 | 用户对话、置顶和排序 | 按职责机械移动并保持Mixin原接口 | `app/backend/features/mentor/repository.py` | 5个方法AST一致;相关表逐行一致 | 已移动 | | 问师消息与偏好方法 | 持久化 | 用户对话、置顶和排序 | 按职责机械移动并保持Mixin原接口 | `app/backend/features/mentor/repository.py` | 5个方法AST一致;相关表逐行一致 | 已移动 |
| `llm_stream.py`与模型访问方法 | 公共模型能力 | 问师、问天、复盘助手与策略编译 | 移入唯一模型边界并保留兼容别名 | `app/backend/llm/` | 流式文件哈希一致;21个服务方法与既有用户边界一致 | 已移动 | | `llm_stream.py`与模型访问方法 | 公共模型能力 | 问师、问天、复盘助手与策略编译 | 移入唯一模型边界并保留兼容别名 | `app/backend/llm/` | 流式文件哈希一致;21个服务方法与既有用户边界一致 | 已移动 |
| 公开`游资skills` | 运行资产 | 问师模型库 | 原样保留 | `app/游资skills/` | 190个文件逐路径和SHA-256一致 | 已复制 | | 公开`游资skills` | 运行资产 | 问师模型库 | 原样保留 | `app/游资skills/` | 190个文件逐路径和SHA-256一致 | 已复制 |
| `heaven_agent.py``heaven_engine.py`与问天服务 | 业务计算 | 观势、观气、观心 | 机械移动并保留兼容别名 | `app/backend/features/heaven/` | Agent文件哈希、引擎定义AST、20个服务方法、6个API及真实动画页面等价 | 已移动 |
| 问天历史方法 | 持久化 | 三类解读历史与当日解运复用 | 按职责机械移动 | `app/backend/features/heaven/repository.py` | 5个方法AST一致;62个schema对象及6张关键表逐行一致 | 已移动 |
| `assistant_agent.py`与复盘服务 | 业务服务 | 我的复盘、自选、交易日志与复盘助手 | 机械移动并保留兼容别名 | `app/backend/features/review/` | Agent文件哈希、8个服务方法AST、5个真实API及会员/非会员页面等价 | 已移动 |
| 复盘、自选、交易日志与助手消息方法 | 持久化 | 用户私有复盘数据 | 按职责机械移动并保持Mixin原接口 | `app/backend/features/review/repository.py` | 13个方法AST一致;相关表逐行一致;浏览器账户隔离通过 | 已移动 |
| 提醒编排与提醒持久化方法 | 业务服务/持久化 | 提醒中心与策略跟踪提醒 | 按职责机械移动 | `app/backend/features/alerts/` | 5个服务方法、6个Repository方法AST及真实API一致 | 已移动 |
处置只允许:`原样保留``移动``合并重复``待定``确认废弃` 处置只允许:`原样保留``移动``合并重复``待定``确认废弃`
@@ -145,6 +151,26 @@
- 回档:标签`xiaobai-preservation-slice-07-20260731` - 回档:标签`xiaobai-preservation-slice-07-20260731`
- 完整证据:`docs/migration/evidence/slice-07/README.md` - 完整证据:`docs/migration/evidence/slice-07/README.md`
已完成切片:`slice-08-heaven-trend-fortune-heart`
- 原版基线:提交`2919229`,即切片07回档点。
- 迁移范围:问天Agent、历法/卦象引擎、20个服务方法、5个持久化方法及3个HTTP入口。
- 兼容边界:根级`heaven_agent.py``heaven_engine.py`指向正式模块对象;原动画、DOM、JS与CSS未改动。
- API与数据库:6个固定输入真实API完全一致;62个schema对象及6张关键表逐行一致。
- 验收:原版231项、迁移版280项Python测试、7项切片源码等价测试、45项Playwright及1080P三模式动画页面通过。
- 回档:标签`xiaobai-preservation-slice-08-20260731`
- 完整证据:`docs/migration/evidence/slice-08/README.md`
已完成切片:`slice-09-review-watchlist-notes-journal-alerts-assistant`
- 原版基线:提交`b3df070`,即切片08回档点。
- 迁移范围:复盘助手Agent、5个提醒编排方法、8个复盘编排方法、19个持久化方法及5个HTTP入口。
- 兼容边界:根级`assistant_agent.py`指向正式模块对象;提醒与交易日志旧服务继续作为唯一底层实现;包导出采用延迟加载以切断初始化循环。
- API与数据库:5个固定输入真实API完全一致;62个schema对象及8张关键表逐行一致。
- 验收:原版231项、迁移版287项Python测试、7项切片源码等价测试、45项Playwright及1920×1080会员/非会员真实页面通过。
- 回档:标签`xiaobai-preservation-slice-09-20260731`
- 完整证据:`docs/migration/evidence/slice-09/README.md`
## 决策记录 ## 决策记录
| 日期 | 决策 | 原因 | | 日期 | 决策 | 原因 |