Compare commits

..
28 changed files with 617 additions and 158 deletions
+5 -2
View File
@@ -28,8 +28,11 @@ background scheduler
- `server.py` is the stable command/import facade. Runtime composition lives in
`backend/application.py` and `backend/bootstrap/`.
- `backend/http/` owns common authentication, request IDs, responses, static delivery, and
error normalization. Feature-specific transport handlers live beside their feature.
- `backend/bootstrap/` owns process configuration, dependency construction, startup, and
shared input/display-format contracts. It does not own feature behavior.
- `backend/http/` owns common authentication, request IDs, JSON/NDJSON responses, static
delivery, streaming connection lifecycle, and error normalization. Feature-specific
transport handlers live beside their feature.
Exact POST endpoints that only delegate to one of those handlers use the explicit maps in
`backend/application.py`; endpoints with path parameters, body handling, or special error
semantics remain visible control flow in `RequestHandler`.
+8 -45
View File
@@ -128,14 +128,19 @@ class DashboardService(
profile_supplier=self._resolved_llm_profile,
)
self.screener.ensure_builtin_strategies()
self._background_stop = threading.Event()
self._background_thread = self.jobs.start_scheduler(
def start_background_jobs(self) -> threading.Thread:
return self.jobs.start_scheduler(
self._background_refresh_tick,
self._background_stop,
interval_seconds=5,
initial_delay_seconds=3,
)
def stop_background_jobs(self, timeout_seconds: float = 5) -> bool:
scheduler_stopped = self.jobs.stop_scheduler(timeout_seconds)
workers_stopped = self.jobs.wait_for_idle(timeout_seconds)
return scheduler_stopped and workers_stopped
def _load_system_credentials(self, environment: dict[str, str]) -> dict[str, Any]:
encrypted = self.database.get_system_setting("credentials")
@@ -472,38 +477,6 @@ class DashboardService(
**self.database.status(),
}
@staticmethod
def _ifind_field(row: dict[str, Any], tokens: tuple[str, ...]) -> Any:
for key, value in row.items():
label = str(key or "")
if any(token.casefold() == label.casefold() for token in tokens):
return value
for key, value in row.items():
label = str(key or "")
if any(token in label for token in tokens):
return value
return None
@classmethod
def _ifind_row_code(cls, row: dict[str, Any]) -> str:
value = cls._ifind_field(row, ("股票代码", "证券代码", "代码", "thscode"))
match = re.search(r"(?<!\d)(\d{6})(?!\d)", str(value or ""))
if match:
return match.group(1)
for value in row.values():
match = re.search(r"(?<!\d)(\d{6})\.(?:SH|SZ|BJ)(?![A-Z])", str(value or ""), re.I)
if match:
return match.group(1)
return ""
SERVICE = DashboardService()
@@ -1020,16 +993,6 @@ class RequestHandler(
return
self.send_json({"error": "Not found"}, HTTPStatus.NOT_FOUND)
def _write_stream_event(self, payload: dict[str, Any]) -> None:
self.wfile.write(
(json.dumps(payload, ensure_ascii=False, separators=(",", ":")) + "\n").encode("utf-8")
)
self.wfile.flush()
def save_reason(self) -> None:
try:
body = self.read_json_body()
+4
View File
@@ -71,6 +71,10 @@ def normalize_date(value: str) -> str:
return parsed.strftime("%Y%m%d")
def display_compact_date(value: str) -> str:
return f"{value[:4]}-{value[4:6]}-{value[6:8]}" if len(value) == 8 else value
def validate_stock_code(value: str) -> str:
code = value.strip()
if not re.fullmatch(r"\d{6}", code):
+4 -3
View File
@@ -16,12 +16,13 @@ def main(handler_class: type[Any] | None = None, service: Any | None = None) ->
parser.add_argument("--port", type=int, default=8765)
args = parser.parse_args()
server = ThreadingHTTPServer((args.host, args.port), handler_class)
print(f"Xiaobai Review Web is running at http://{args.host}:{args.port}")
print("Press Ctrl+C to stop.")
try:
service.start_background_jobs()
print(f"Xiaobai Review Web is running at http://{args.host}:{args.port}")
print("Press Ctrl+C to stop.")
server.serve_forever()
except KeyboardInterrupt:
pass
finally:
service._background_stop.set()
service.stop_background_jobs()
server.server_close()
+1 -4
View File
@@ -11,6 +11,7 @@ from datetime import datetime, time as dt_time, timedelta
from threading import Lock
from typing import Any, ClassVar
from backend.bootstrap.config import display_compact_date as _display_date
from backend.data.numbers import finite_number as _number
from backend.features.sentiment.engine import apply_sentiment_to_dashboard
@@ -2007,10 +2008,6 @@ def _display_time(value: Any) -> str:
return f"{raw[:2]}:{raw[2:4]}:{raw[4:6]}"
def _display_date(value: str) -> str:
return f"{value[:4]}-{value[4:6]}-{value[6:8]}" if len(value) == 8 else value
def _realtime_market_status(current_time: dt_time) -> str:
if current_time < dt_time(9, 25):
return "pre_open"
+29
View File
@@ -7,6 +7,35 @@ from typing import Any
class HeavenRepositoryMixin:
def list_sector_phase_overrides(self) -> dict[str, str]:
with self.connect() as connection:
rows = connection.execute(
"SELECT name, element FROM sector_phase_overrides ORDER BY updated_at DESC, name"
).fetchall()
return {row["name"]: row["element"] for row in rows}
def save_sector_phase_override(self, name: str, element: str) -> None:
now = datetime.now().astimezone().isoformat(timespec="seconds")
with self.connect() as connection:
connection.execute(
"""
INSERT INTO sector_phase_overrides (name, element, updated_at)
VALUES (?, ?, ?)
ON CONFLICT(name) DO UPDATE SET
element = excluded.element,
updated_at = excluded.updated_at
""",
(name, element, now),
)
def delete_sector_phase_override(self, name: str) -> bool:
with self.connect() as connection:
cursor = connection.execute(
"DELETE FROM sector_phase_overrides WHERE name = ?",
(name,),
)
return cursor.rowcount > 0
@staticmethod
def _heaven_reading_dict(row: sqlite3.Row | None) -> dict[str, Any] | None:
if not row:
+1 -16
View File
@@ -14,19 +14,4 @@ class MentorHttpMixin:
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 event in stream:
self._write_stream_event(event)
self._write_stream_event({"type": "done"})
except (ValueError, MentorAgentError) as exc:
self._write_stream_event({"type": "error", "error": str(exc)})
except (BrokenPipeError, ConnectionResetError):
pass
finally:
self.close_connection = True
self.send_ndjson_stream(stream, (ValueError, MentorAgentError))
+24
View File
@@ -117,6 +117,30 @@ class PoolServiceMixin:
finally:
self._ifind_event_lock.release()
@staticmethod
def _ifind_field(row: dict[str, Any], tokens: tuple[str, ...]) -> Any:
for key, value in row.items():
label = str(key or "")
if any(token.casefold() == label.casefold() for token in tokens):
return value
for key, value in row.items():
label = str(key or "")
if any(token in label for token in tokens):
return value
return None
@classmethod
def _ifind_row_code(cls, row: dict[str, Any]) -> str:
value = cls._ifind_field(row, ("股票代码", "证券代码", "代码", "thscode"))
match = re.search(r"(?<!\d)(\d{6})(?!\d)", str(value or ""))
if match:
return match.group(1)
for value in row.values():
match = re.search(r"(?<!\d)(\d{6})\.(?:SH|SZ|BJ)(?![A-Z])", str(value or ""), re.I)
if match:
return match.group(1)
return ""
@staticmethod
def _normalize_ifind_event_time(value: Any) -> str:
text = str(value or "").strip()
+2 -16
View File
@@ -26,22 +26,8 @@ class ReviewHttpMixin:
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
events = ({"type": "delta", "content": chunk} for chunk in stream)
self.send_ndjson_stream(events, (ValueError, ReviewAssistantError))
def save_watchlist(self) -> None:
try:
+1 -4
View File
@@ -8,6 +8,7 @@ from collections import defaultdict
from datetime import datetime, timedelta
from typing import Any
from backend.bootstrap.config import display_compact_date as _display_date
from backend.data.numbers import finite_number as _number
from backend.data.providers.tushare_client import TushareClient, TushareError
from backend.features.screener.strategies import ADVANCED_CURATED_STRATEGIES
@@ -2200,7 +2201,3 @@ def _regime_reason(regime: str) -> str:
"divergence": "指数或核心仍强,但广度、封板质量开始分化。",
"retreat": "情绪指标继续走弱,应提高筛选门槛并接受无候选结果。",
}.get(regime, "市场阶段待确认。")
def _display_date(value: str) -> str:
return f"{value[:4]}-{value[4:6]}-{value[6:8]}" if len(value) == 8 else value
+29
View File
@@ -3,6 +3,7 @@ from __future__ import annotations
import json
import mimetypes
import secrets
from collections.abc import Iterable
from http import HTTPStatus
from http.cookies import SimpleCookie
from typing import Any
@@ -142,5 +143,33 @@ class HttpTransportMixin:
self.end_headers()
self.wfile.write(content)
def _write_stream_event(self, payload: dict[str, Any]) -> None:
self.wfile.write(
(json.dumps(payload, ensure_ascii=False, separators=(",", ":")) + "\n").encode("utf-8")
)
self.wfile.flush()
def send_ndjson_stream(
self,
events: Iterable[dict[str, Any]],
error_types: tuple[type[Exception], ...],
) -> None:
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 event in events:
self._write_stream_event(event)
self._write_stream_event({"type": "done"})
except error_types as exc:
self._write_stream_event({"type": "error", "error": str(exc)})
except (BrokenPipeError, ConnectionResetError):
pass
finally:
self.close_connection = True
def log_message(self, format_string: str, *args: Any) -> None:
print(f"[{self.log_date_time_string()}] {format_string % args}")
+37 -19
View File
@@ -18,6 +18,9 @@ class InProcessJobRunner:
self.repository = repository
self._locks: dict[str, threading.Lock] = {}
self._locks_guard = threading.Lock()
self._scheduler_guard = threading.Lock()
self._scheduler_stop = threading.Event()
self._scheduler_thread: threading.Thread | None = None
def submit(
self, job_id: str, idempotency_key: str, action: JobAction,
@@ -52,27 +55,42 @@ class InProcessJobRunner:
return True
def start_scheduler(
self, callback: Callable[[], None], stop_event: threading.Event,
interval_seconds: float, initial_delay_seconds: float = 0,
self, callback: Callable[[], None], interval_seconds: float,
initial_delay_seconds: float = 0,
) -> threading.Thread:
def schedule_loop() -> None:
if stop_event.wait(initial_delay_seconds):
return
while not stop_event.is_set():
try:
callback()
except Exception:
# Submitted jobs persist their own failures; the scheduler must stay alive.
pass
stop_event.wait(interval_seconds)
with self._scheduler_guard:
current = self._scheduler_thread
if current is not None and current.is_alive():
return current
self._scheduler_stop.clear()
thread = threading.Thread(
target=schedule_loop,
name="background-job-scheduler",
daemon=True,
)
thread.start()
return thread
def schedule_loop() -> None:
if self._scheduler_stop.wait(initial_delay_seconds):
return
while not self._scheduler_stop.is_set():
try:
callback()
except Exception:
# Submitted jobs persist failures; the scheduler must stay alive.
pass
self._scheduler_stop.wait(interval_seconds)
thread = threading.Thread(
target=schedule_loop,
name="background-job-scheduler",
daemon=True,
)
self._scheduler_thread = thread
thread.start()
return thread
def stop_scheduler(self, timeout_seconds: float = 5) -> bool:
with self._scheduler_guard:
thread = self._scheduler_thread
self._scheduler_stop.set()
if thread is not None and thread is not threading.current_thread():
thread.join(max(0, timeout_seconds))
return thread is None or not thread.is_alive()
def wait_for_idle(self, timeout_seconds: float = 5) -> bool:
deadline = time.monotonic() + max(0, timeout_seconds)
+1
View File
@@ -4,6 +4,7 @@ from datetime import datetime, timezone
from typing import Any
from urllib.parse import urlparse
from backend.bootstrap.config import validate_text
from backend.features.screener.compiler import LLMCompilerError, test_llm_connection
+43 -8
View File
@@ -220,6 +220,25 @@
"runtime_role": "index observation fallback"
}
],
"provider_construction": [
{
"client": "TushareClient",
"owner": "backend/data/providers/tushare.py",
"compatibility_fallback": "backend/features/market/service.py"
},
{
"client": "IfindHttpClient",
"owner": "backend/data/gateway.py"
},
{
"client": "MarketChartClient",
"owner": "backend/data/gateway.py"
},
{
"client": "WebRealtimeAggregator",
"owner": "backend/data/gateway.py"
}
],
"numeric_normalization": [
{
"function": "finite_number",
@@ -230,6 +249,12 @@
"path": "backend/data/numbers.py"
}
],
"date_formatting": [
{
"function": "display_compact_date",
"path": "backend/bootstrap/config.py"
}
],
"llm_entrypoints": [
{
"function": "stream_with_mentor",
@@ -262,6 +287,16 @@
"path": "backend/llm/transport.py"
}
],
"http_transport": [
{
"function": "send_json",
"path": "backend/http/handler.py"
},
{
"function": "send_ndjson_stream",
"path": "backend/http/handler.py"
}
],
"css_layers": [
"/shared/tokens.css?v=20260729-1",
"/styles/styles.css",
@@ -289,13 +324,13 @@
},
{
"path": "backend/features/screener/engine.py",
"bytes": 108434,
"lines": 2206
"bytes": 108387,
"lines": 2203
},
{
"path": "backend/data/providers/tushare_client.py",
"bytes": 94171,
"lines": 2168
"bytes": 94124,
"lines": 2165
},
{
"path": "frontend/app.js",
@@ -339,8 +374,8 @@
},
{
"path": "backend/application.py",
"bytes": 48749,
"lines": 1129
"bytes": 47769,
"lines": 1092
},
{
"path": "frontend/styles/theme.css",
@@ -349,8 +384,8 @@
},
{
"path": "database.py",
"bytes": 33284,
"lines": 746
"bytes": 32073,
"lines": 716
}
]
}
+1 -31
View File
@@ -2,7 +2,7 @@ from __future__ import annotations
import json
import sqlite3
from datetime import datetime, timezone
from datetime import datetime
from pathlib import Path
from typing import Any
@@ -665,36 +665,6 @@ class ReviewDatabase(
)
MigrationRunner().apply(connection, MIGRATIONS)
def list_sector_phase_overrides(self) -> dict[str, str]:
with self.connect() as connection:
rows = connection.execute(
"SELECT name, element FROM sector_phase_overrides ORDER BY updated_at DESC, name"
).fetchall()
return {row["name"]: row["element"] for row in rows}
def save_sector_phase_override(self, name: str, element: str) -> None:
now = datetime.now().astimezone().isoformat(timespec="seconds")
with self.connect() as connection:
connection.execute(
"""
INSERT INTO sector_phase_overrides (name, element, updated_at)
VALUES (?, ?, ?)
ON CONFLICT(name) DO UPDATE SET
element = excluded.element,
updated_at = excluded.updated_at
""",
(name, element, now),
)
def delete_sector_phase_override(self, name: str) -> bool:
with self.connect() as connection:
cursor = connection.execute(
"DELETE FROM sector_phase_overrides WHERE name = ?",
(name,),
)
return cursor.rowcount > 0
def list_wencai_saved_queries(
self, user_id: int, limit: int = 30
) -> list[dict[str, Any]]:
+65
View File
@@ -0,0 +1,65 @@
from __future__ import annotations
import ast
import sys
import unittest
from pathlib import Path
from unittest.mock import patch
from backend.bootstrap import runtime
class BackgroundLifecycleTests(unittest.TestCase):
def test_service_construction_does_not_start_the_scheduler(self) -> None:
source = (Path(__file__).parents[1] / "backend" / "application.py").read_text(
encoding="utf-8"
)
tree = ast.parse(source)
service_class = next(
node for node in tree.body
if isinstance(node, ast.ClassDef) and node.name == "DashboardService"
)
constructor = next(
node for node in service_class.body
if isinstance(node, ast.FunctionDef) and node.name == "__init__"
)
called_methods = {
node.func.attr for node in ast.walk(constructor)
if isinstance(node, ast.Call) and isinstance(node.func, ast.Attribute)
}
self.assertNotIn("start_scheduler", called_methods)
def test_runtime_starts_jobs_after_bind_and_stops_before_close(self) -> None:
events: list[str] = []
class Service:
def start_background_jobs(self) -> None:
events.append("start-jobs")
def stop_background_jobs(self) -> None:
events.append("stop-jobs")
class Server:
def __init__(self, address: tuple[str, int], handler: object) -> None:
events.append("bind")
def serve_forever(self) -> None:
events.append("serve")
raise KeyboardInterrupt
def server_close(self) -> None:
events.append("close")
with (
patch.object(runtime, "ThreadingHTTPServer", Server),
patch.object(sys, "argv", ["server.py", "--port", "8797"]),
):
runtime.main(object, Service())
self.assertEqual(
events, ["bind", "start-jobs", "serve", "stop-jobs", "close"]
)
if __name__ == "__main__":
unittest.main()
+29 -2
View File
@@ -1,7 +1,9 @@
from __future__ import annotations
import ast
import unittest
from datetime import datetime, timedelta
from pathlib import Path
from backend.data import (
DataPolicyError,
@@ -45,8 +47,6 @@ class DataGatewayTests(unittest.TestCase):
self.assertIs(gateway.chart_data.ifind, gateway.ifind)
def test_server_has_no_direct_runtime_tushare_construction(self) -> None:
from pathlib import Path
source = (
Path(__file__).resolve().parents[1]
/ "backend"
@@ -57,6 +57,33 @@ class DataGatewayTests(unittest.TestCase):
self.assertEqual(source.count("TushareClient(self.token)"), 1)
self.assertIn("return gateway.tushare()", source)
def test_provider_construction_has_unique_declared_owners(self) -> None:
root = Path(__file__).resolve().parents[1]
owners = {
"EastmoneyChartClient": {"backend/data/gateway.py"},
"IfindHttpClient": {"backend/data/gateway.py"},
"IfindProvider": {"backend/data/gateway.py"},
"MarketChartClient": {"backend/data/gateway.py"},
"TushareClient": {"backend/features/market/service.py"},
"TushareProvider": {"backend/data/gateway.py"},
"WebRealtimeAggregator": {"backend/data/gateway.py"},
}
found = {name: set() for name in owners}
for path in (root / "backend").rglob("*.py"):
relative = path.relative_to(root).as_posix()
tree = ast.parse(path.read_text(encoding="utf-8"), filename=str(path))
for node in ast.walk(tree):
if not isinstance(node, ast.Call):
continue
name = getattr(node.func, "id", None) or getattr(node.func, "attr", None)
if name in found:
found[name].add(relative)
self.assertEqual(found, owners)
provider_source = (root / "backend/data/providers/tushare.py").read_text(
encoding="utf-8"
)
self.assertIn("client_factory: Callable[[str], TushareClient] = TushareClient", provider_source)
def test_quality_gate_accepts_matching_daily_evidence(self) -> None:
timezone = market_timezone()
now = datetime(2026, 7, 29, 16, 0, tzinfo=timezone)
+85
View File
@@ -0,0 +1,85 @@
from __future__ import annotations
import io
import json
import unittest
from backend.http.handler import HttpTransportMixin
class StreamError(RuntimeError):
pass
class TransportStub(HttpTransportMixin):
def __init__(self) -> None:
self.statuses: list[int] = []
self.response_headers: list[tuple[str, str]] = []
self.wfile = io.BytesIO()
self.close_connection = False
def send_response(self, status: int) -> None:
self.statuses.append(int(status))
def send_header(self, name: str, value: str) -> None:
self.response_headers.append((name, value))
def end_headers(self) -> None:
pass
class HttpStreamingTests(unittest.TestCase):
def test_stream_transport_preserves_headers_events_and_completion(self) -> None:
handler = TransportStub()
handler.send_ndjson_stream(
({"type": "delta", "content": value} for value in ("", "")),
(StreamError,),
)
self.assertEqual(handler.statuses, [200])
self.assertEqual(
dict(handler.response_headers),
{
"Content-Type": "application/x-ndjson; charset=utf-8",
"Cache-Control": "no-cache, no-transform",
"X-Accel-Buffering": "no",
"Connection": "close",
},
)
events = [
json.loads(line)
for line in handler.wfile.getvalue().decode("utf-8").splitlines()
]
self.assertEqual(
events,
[
{"type": "delta", "content": ""},
{"type": "delta", "content": ""},
{"type": "done"},
],
)
self.assertTrue(handler.close_connection)
def test_stream_transport_preserves_feature_error_event(self) -> None:
def events():
yield {"type": "delta", "content": "partial"}
raise StreamError("stream failed")
handler = TransportStub()
handler.send_ndjson_stream(events(), (StreamError,))
payloads = [
json.loads(line)
for line in handler.wfile.getvalue().decode("utf-8").splitlines()
]
self.assertEqual(
payloads,
[
{"type": "delta", "content": "partial"},
{"type": "error", "error": "stream failed"},
],
)
self.assertTrue(handler.close_connection)
if __name__ == "__main__":
unittest.main()
+23
View File
@@ -68,6 +68,29 @@ class JobRunnerTests(unittest.TestCase):
)
self.assertEqual(self.repository.recent(1)[0]["status"], "failed")
def test_scheduler_start_is_idempotent_and_stop_waits_for_exit(self) -> None:
first = self.runner.start_scheduler(
lambda: None, interval_seconds=60, initial_delay_seconds=60
)
second = self.runner.start_scheduler(
lambda: None, interval_seconds=60, initial_delay_seconds=60
)
self.assertIs(first, second)
self.assertTrue(first.is_alive())
self.assertTrue(self.runner.stop_scheduler())
self.assertFalse(first.is_alive())
def test_scheduler_can_restart_after_an_orderly_stop(self) -> None:
first = self.runner.start_scheduler(
lambda: None, interval_seconds=60, initial_delay_seconds=60
)
self.assertTrue(self.runner.stop_scheduler())
second = self.runner.start_scheduler(
lambda: None, interval_seconds=60, initial_delay_seconds=60
)
self.assertIsNot(first, second)
self.assertTrue(self.runner.stop_scheduler())
if __name__ == "__main__":
unittest.main()
+27
View File
@@ -1,8 +1,10 @@
from __future__ import annotations
import unittest
from unittest.mock import patch
from backend.llm import LLMGateway, LLMGatewayError
from backend.llm.service import LLMServiceMixin
class ProviderFailure(RuntimeError):
@@ -113,6 +115,31 @@ class LLMGatewayTests(unittest.TestCase):
with self.assertRaisesRegex(LLMGatewayError, "额度已用完"):
gateway.call("mentor", "mentor-v1", lambda model: "unused", (ProviderFailure,))
def test_saved_system_model_can_reach_the_connection_probe(self) -> None:
service = LLMServiceMixin()
service._system_credentials = {
"llm_models": [
{
"id": "primary-model",
"name": "主模型",
"api_key": "secret",
"base_url": "https://model.example/v1",
"model": "model-name",
}
]
}
service.llm_gateway = self.gateway()
expected = {"ok": True, "reply": "OK"}
with patch(
"backend.llm.service.test_llm_connection", return_value=expected
) as connection_probe:
result = service.test_system_llm_profile("primary-model", {})
self.assertEqual(result, expected)
connection_probe.assert_called_once_with(
"secret", "https://model.example/v1", "model-name"
)
if __name__ == "__main__":
unittest.main()
@@ -40,6 +40,8 @@ class AccountSliceStructureTests(unittest.TestCase):
"require_access",
"serve_static",
"send_json",
"_write_stream_event",
"send_ndjson_stream",
):
self.assertNotIn(method, RequestHandler.__dict__)
self.assertIn(method, HttpTransportMixin.__dict__)
@@ -38,6 +38,9 @@ HEAVEN_SERVICE_METHODS = {
}
HEAVEN_REPOSITORY_METHODS = {
"list_sector_phase_overrides",
"save_sector_phase_override",
"delete_sector_phase_override",
"_heaven_reading_dict",
"save_heaven_reading",
"list_heaven_readings",
@@ -146,6 +146,7 @@ class MarketSliceSourceEquivalenceTests(unittest.TestCase):
self.assertEqual(sha256(ORIGINAL_ROOT / original), sha256(APP_ROOT / migrated))
original_tushare = top_level_definitions(ORIGINAL_ROOT / "tushare_client.py")
original_tushare.pop("_number")
original_tushare.pop("_display_date")
self.assertEqual(
original_tushare,
top_level_definitions(APP_ROOT / "backend/data/providers/tushare_client.py"),
@@ -155,6 +156,13 @@ class MarketSliceSourceEquivalenceTests(unittest.TestCase):
function_contract(APP_ROOT / "backend/data/numbers.py", "finite_number"),
)
self.assertIs(canonical_tushare._number, finite_number)
self.assertEqual(
function_contract(ORIGINAL_ROOT / "tushare_client.py", "_display_date"),
function_contract(
APP_ROOT / "backend/bootstrap/config.py", "display_compact_date"
),
)
self.assertIs(canonical_tushare._display_date, bootstrap_config.display_compact_date)
original_charts = top_level_definitions(ORIGINAL_ROOT / "chart_data_provider.py")
original_charts.pop("_stock_market_code")
self.assertEqual(
+10 -2
View File
@@ -9,6 +9,7 @@ import advanced_strategies
import llm_strategy
import screener
import strategy_tracking
from backend.bootstrap import config as bootstrap_config
from backend.data.numbers import finite_number
from backend.features.screener import compiler, engine, strategies, tracking
from backend.features.screener import service as screener_service
@@ -151,12 +152,12 @@ class ScreenerSliceSourceEquivalenceTests(unittest.TestCase):
self.assertEqual(
module_contract(
ORIGINAL_ROOT / "screener.py",
excluded_definitions={"_number"},
excluded_definitions={"_display_date", "_number"},
exclude_imports=True,
),
module_contract(
APP_ROOT / "backend" / "features" / "screener" / "engine.py",
excluded_definitions={"_number"},
excluded_definitions={"_display_date", "_number"},
exclude_imports=True,
),
)
@@ -165,6 +166,13 @@ class ScreenerSliceSourceEquivalenceTests(unittest.TestCase):
function_contract(APP_ROOT / "backend/data/numbers.py", "finite_number"),
)
self.assertIs(engine._number, finite_number)
self.assertEqual(
function_contract(ORIGINAL_ROOT / "screener.py", "_display_date"),
function_contract(
APP_ROOT / "backend/bootstrap/config.py", "display_compact_date"
),
)
self.assertIs(engine._display_date, bootstrap_config.display_compact_date)
self.assertEqual(
class_methods(
ORIGINAL_ROOT / "backend" / "features" / "screener" / "tracking.py",
@@ -29,6 +29,8 @@ POOL_METHODS = {
"_apply_reason_overrides",
"_schedule_ifind_event_enrichment",
"_refresh_ifind_event_enrichment",
"_ifind_field",
"_ifind_row_code",
"_normalize_ifind_event_time",
"_merge_ifind_event_enrichment",
}
+13
View File
@@ -155,10 +155,19 @@ def build() -> dict[str, Any]:
{"provider": "eastmoney", "path": "backend/data/realtime.py", "runtime_role": "isolated realtime observation"},
{"provider": "tencent", "path": "backend/data/realtime.py", "runtime_role": "index observation fallback"},
],
"provider_construction": [
{"client": "TushareClient", "owner": "backend/data/providers/tushare.py", "compatibility_fallback": "backend/features/market/service.py"},
{"client": "IfindHttpClient", "owner": "backend/data/gateway.py"},
{"client": "MarketChartClient", "owner": "backend/data/gateway.py"},
{"client": "WebRealtimeAggregator", "owner": "backend/data/gateway.py"},
],
"numeric_normalization": [
{"function": "finite_number", "path": "backend/data/numbers.py"},
{"function": "non_nan_number", "path": "backend/data/numbers.py"},
],
"date_formatting": [
{"function": "display_compact_date", "path": "backend/bootstrap/config.py"},
],
"llm_entrypoints": [
{"function": "stream_with_mentor", "path": "backend/features/mentor/agent.py"},
{"function": "interpret_heaven", "path": "backend/features/heaven/agent.py"},
@@ -170,6 +179,10 @@ def build() -> dict[str, Any]:
{"function": "chat_completion", "path": "backend/llm/transport.py"},
{"function": "stream_chat_completion", "path": "backend/llm/transport.py"},
],
"http_transport": [
{"function": "send_json", "path": "backend/http/handler.py"},
{"function": "send_ndjson_stream", "path": "backend/http/handler.py"},
],
"css_layers": css_layers(html),
"code_hotspots": code_hotspots(),
}
+11 -6
View File
@@ -32,17 +32,22 @@ def main() -> None:
from server import RequestHandler, SERVICE
server = ThreadingHTTPServer(("127.0.0.1", args.port), RequestHandler)
print(
f"Preservation runtime is running at http://127.0.0.1:{args.port} "
f"with database {SERVICE.database.path}",
flush=True,
)
try:
if hasattr(SERVICE, "start_background_jobs"):
SERVICE.start_background_jobs()
print(
f"Preservation runtime is running at http://127.0.0.1:{args.port} "
f"with database {SERVICE.database.path}",
flush=True,
)
server.serve_forever()
except KeyboardInterrupt:
pass
finally:
SERVICE._background_stop.set()
if hasattr(SERVICE, "stop_background_jobs"):
SERVICE.stop_background_jobs()
else:
SERVICE._background_stop.set()
server.server_close()
+149
View File
@@ -23,6 +23,12 @@
| CR-03 | 股票市场后缀转换 | Tushare业务与iFinD图表各保留一份完全相同的沪深京代码转换函数 | 图表复用`bootstrap/config.py::tushare_code`,只保留一份函数体 | 已完成 |
| CR-04 | 数值归一化策略 | 四个业务模块分别保留两组完全相同的数值转换函数体 | 由`backend/data/numbers.py`集中拥有两种既有语义,消费者保留原局部别名 | 已完成 |
| CR-05 | 根级兼容入口 | 正式后端仍有五处通过迁移兼容模块反向导入规范实现 | 正式代码改用规范路径;兼容入口只服务原公开导入契约 | 已完成 |
| CR-06 | 紧凑日期显示 | Tushare与选股引擎各保留一份完全相同的`YYYYMMDD`显示转换 | 由`bootstrap/config.py`拥有唯一格式策略,消费者保留原局部别名 | 已完成 |
| CR-07 | 数据Provider组装 | 核查网关、容器和业务服务是否重复创建外部数据客户端 | 固化唯一创建位置及兼容例外,不改数据源语义 | 已完成 |
| CR-08 | NDJSON流式传输 | 问师与复盘助手重复维护响应头、事件写入、完成、断线及关闭流程 | HTTP共享层拥有唯一流式连接生命周期 | 已完成 |
| CR-09 | 应用服务门面 | 两个仅服务股池iFinD补全的方法仍错位在全局`DashboardService` | 原函数体机械归位到`PoolServiceMixin` | 已完成 |
| CR-10 | Repository所有权 | 五行行业阶段覆盖的三项持久化方法仍错位在根级数据库门面 | 原函数体机械归位到问天Repository,兼容数据继续保留 | 已完成 |
| CR-11 | 后台任务生命周期 | 导入应用即启动调度线程,启动早于端口绑定,停止只置位但不等待 | 运行时显式启停,每个Runner只拥有一个可等待的调度线程 | 已完成 |
## CR-01验收口径
@@ -139,6 +145,149 @@
本批基线为`xiaobai-reduction-04-numeric-normalization-20260801`;检查点为
`xiaobai-reduction-05-compatibility-boundaries-20260801`
## CR-06验收口径
- 只合并参数、函数体和异常行为完全相同的日期/文本转换;名称相似但空值、未来日期、错误文案或
输入格式不同的函数不得合并。
- Tushare与选股引擎继续暴露局部`_display_date`名称,并分别指向唯一共享实现。
- 原版两个`_display_date`函数必须分别与共享实现AST相等,所有原调用结果保持不变。
- 市场洞察的日期显示函数会清理连字符并容忍空值,语义不同,必须继续独立保留。
## CR-06结果
- 删除Tushare与选股引擎内两个重复日期函数体,新增`display_compact_date`唯一策略;生产代码总
行数不增加,重复函数体由两份降为一份。
- 架构清单登记日期格式唯一所有权;保持性测试改为未改范围AST相等、共享函数AST相等和运行时
对象身份三重契约,没有放宽原迁移门禁。
- `normalize_date`、市场洞察日期显示、实时行情时间格式和会员日期边界因语义不同均原样保留。
- 候选321项、纯`app/`导出258项、24个JavaScript文件、API/架构注册表、Git空白检查和SQLite
完整性检查通过;本批不涉及页面、CSS或浏览器行为。
- 本批不修改日期输入规则、业务计算、选股结果、接口、数据库、数据源、LLM、权限或部署。
本批基线为`xiaobai-reduction-05-compatibility-boundaries-20260801`;检查点为
`xiaobai-reduction-06-date-formatting-20260801`
## CR-07验收口径
- iFinD、图表、实时观察器和Provider适配器必须只在`build_data_gateway`创建,并由容器共享。
- Tushare必须继续通过实时Token供应器按需创建,不能为了减少对象数量缓存过期Token。
- 市场服务中为原版隔离测试桩保留的一处`TushareClient(self.token)`是明确兼容例外,不得被误判为
第二条正式数据链路。
- 不得合并Tushare、iFinD、东方财富和腾讯的传输、缓存、重试或降级逻辑。
## CR-07结果
- 全后端构造点扫描确认iFinD、MarketChart、东方财富图表和实时观察器均只有网关一个创建位置;
`ApplicationContainer`暴露的是同一对象引用,没有第二份客户端。
- Tushare Provider使用动态Token供应器,市场服务只有一处已登记测试兼容回退;本批没有发现可安全
删除的生产实现,因此不为追求行数强行修改运行代码。
- 架构清单新增Provider创建所有权,自动测试会在未来出现第二个未登记构造点时失败。
- 候选322项、纯`app/`导出259项、24个JavaScript文件、API/架构注册表、Git空白检查和SQLite
完整性检查通过;本批不涉及页面、CSS或浏览器行为。
- 本批不修改请求频率、缓存、重试、Token更新、数据源选择、计算口径、API、数据库或前端。
本批基线为`xiaobai-reduction-06-date-formatting-20260801`;检查点为
`xiaobai-reduction-07-provider-ownership-20260801`
## CR-08验收口径
- 只合并NDJSON响应头、事件序列化、完成事件、业务错误事件、客户端断线和连接关闭这些传输行为。
- 问师继续直接输出原事件字典;复盘助手继续把文本分片包装为`delta/content`事件。
- 两个功能各自的业务异常类型、请求体错误状态码、错误正文和流创建时机保持不变。
- `_write_stream_event`与连接生命周期必须只由`backend/http/handler.py`拥有,应用大类和功能HTTP
模块不再保留第二份实现。
## CR-08结果
- 删除问师和复盘助手各16行重复流式控制流,并将应用大类中的10行事件写入方法归入HTTP共享层;
新共享实现29行、两个调用适配共3行,生产代码净减少10行。
- 新增专项测试固定四个响应头、中文NDJSON序列化、增量顺序、完成事件、业务错误事件和关闭状态。
- 候选324项、纯`app/`导出261项及45项Playwright通过;24个JavaScript文件、API/架构注册表、
Git空白检查和SQLite完整性检查通过。
- 本批不修改提示词、模型选择、会员计次、流式正文、前端解析、API路径、数据库或数据源。
本批基线为`xiaobai-reduction-07-provider-ownership-20260801`;检查点为
`xiaobai-reduction-08-ndjson-transport-20260801`
## CR-09验收口径
- 只有消费者全部属于单一领域、且能够按原函数体机械移动的方法才从应用门面移出。
- `_ifind_field``_ifind_row_code`继续保持静态/类方法签名、字段优先级、大小写规则和代码正则。
- 账号委托属于稳定公开门面;系统设置属于跨领域协调;后台刷新留到CR-11,本批均不得删除或重写。
- 移动后`DashboardService`必须继续通过Mixin解析同名方法,调用点和返回值不变。
## CR-09结果
- 将iFinD字段匹配和股票代码提取两个方法从应用大类机械移动到股池服务,原版与迁移方法AST逐项
相等;应用大类不再直接拥有股池专属实现。
- 连同迁移期遗留空行,`backend/application.py`减少32行,股池服务增加24行,生产代码净减少8行。
- 候选324项、纯`app/`导出261项、24个JavaScript文件、API/架构注册表、Git空白检查和SQLite
完整性检查通过;本批不涉及页面、CSS或浏览器行为。
- 本批不修改字段匹配、涨跌停原因补全、接口、数据源、缓存、数据库、权限或前端。
本批基线为`xiaobai-reduction-08-ndjson-transport-20260801`;检查点为
`xiaobai-reduction-09-service-facade-20260801`
## CR-10验收口径
- 只移动调用方、数据表和业务含义均明确属于单一领域的方法;数据库连接、事务和返回值必须保持不变。
- `list_sector_phase_overrides``save_sector_phase_override``delete_sector_phase_override`必须由
`backend/features/heaven/repository.py`拥有,并继续通过`ReviewDatabase`的Mixin解析。
- 不修改表结构、迁移顺序、时间格式、排序、冲突更新或删除结果语义。
- `wencai_saved_queries`及其三个方法属于已登记的账户隔离兼容数据;即使前端入口已取消,也必须保留。
## CR-10结果
- 将五行行业阶段覆盖的查询、保存和删除三个方法从根级`database.py`机械移动到问天Repository
调用名称、SQL、事务边界、时间值和返回结果均未改变。
- 根级数据库门面不再直接拥有问天领域的持久化实现,问财历史兼容表和方法完整保留,未扩大删除范围。
- 34项Repository、问天、账户隔离、清理契约及迁移定向测试通过;候选324项、纯`app/`导出261项、
24个JavaScript文件、API/架构注册表、Git空白检查和SQLite完整性检查通过。
- 本批不修改页面、CSS、API、数据库结构、行情、数据源、业务计算、LLM、权限或后台任务。
本批基线为`xiaobai-reduction-09-service-facade-20260801`;检查点为
`xiaobai-reduction-10-repository-ownership-20260801`
## CR-11验收口径
- 导入`backend.application`或构造`DashboardService`不得启动后台调度;必须先成功绑定HTTP端口,
再由运行时显式启动。
- 同一`InProcessJobRunner`重复启动调度器必须返回同一活动线程,停止必须置位并在限定时间内等待退出,
有序停止后允许重新启动。
- 运行时关闭顺序固定为:停止调度、等待已提交任务、关闭HTTP服务器;迁移对比工具继续兼容原版入口。
- 三个任务的注册定义、5秒刷新频率、3秒初始延迟、幂等键、锁、重试、业务函数和结果不得改变。
- 声明的超时继续是目标与审计字段;Python线程不能安全强杀,本批不伪造硬取消能力。
## CR-11结果
- 删除`DashboardService`构造阶段的调度副作用,端口占用、模块导入和单元测试不再提前创建后台写线程;
`backend/bootstrap/runtime.py`成为正式启动与停止所有者。
- `InProcessJobRunner`集中持有调度停止事件和线程引用;重复启动幂等,停止可等待,原任务锁、持久化运行
状态、成功幂等、失败记录和后续重试逻辑保持不变。
- 新增语法树与运行顺序门禁,固定“构造不启动”“绑定后启动”“停止后关服”,并补齐重复启动与重启测试。
- 候选328项、纯`app/`导出265项和45项Playwright通过;24个JavaScript文件、API/架构注册表、
Git空白检查和SQLite完整性检查通过。
- 本批没有可安全删除的重复任务实现;为补齐原先缺失的生命周期,生产代码净增加24行。增加内容仅为
调度状态、幂等启停和运行时委托,不新增业务层、任务或兼容包装。
- 本批不修改页面、CSS、API、数据源、刷新计算、自动选股条件、数据库结构、权限或LLM。
本批基线为`xiaobai-reduction-10-repository-ownership-20260801`;检查点为
`xiaobai-reduction-11-background-jobs-20260801`
## CR-11人工验收修正
- 2026-08-02人工验收发现本地页面可访问,但行情与LLM同时无法连接。第一原因是验收服务由受限
自动化会话启动,子进程继承了禁止外部网络访问的权限;重新在主机正常网络权限下启动后恢复。
- 随后的主模型连接测试暴露`LLMServiceMixin`机械迁移时遗漏`validate_text`导入,导致请求在真正
访问模型前抛出`NameError`并关闭HTTP连接;恢复原依赖并增加保存模型连接探测的运行契约测试。
- 网站自身实测`000001`返回Tushare日K 60根、分时242点;主模型
`MiniMax-M2.7-highspeed`在2236毫秒内回复`OK`,证明服务进程的数据与LLM出网链路均已恢复。
- 修正后候选329项、纯`app/`导出266项通过;本次只恢复缺失导入和测试,不修改模型配置、额度、
提示词、回退策略、行情来源或计算逻辑。
本修正基线为`xiaobai-reduction-11-background-jobs-20260801`;检查点为
`xiaobai-reduction-11-runtime-connectivity-fix-20260802`
## 人工验收记录
- 2026-08-01:用户检查CR-02与CR-03运行结果,确认未发现明显异常。本记录仅表示本轮可见功能与