Compare commits

..
20 changed files with 436 additions and 142 deletions
+3 -2
View File
@@ -30,8 +30,9 @@ background scheduler
`backend/application.py` and `backend/bootstrap/`. `backend/application.py` and `backend/bootstrap/`.
- `backend/bootstrap/` owns process configuration, dependency construction, startup, and - `backend/bootstrap/` owns process configuration, dependency construction, startup, and
shared input/display-format contracts. It does not own feature behavior. shared input/display-format contracts. It does not own feature behavior.
- `backend/http/` owns common authentication, request IDs, responses, static delivery, and - `backend/http/` owns common authentication, request IDs, JSON/NDJSON responses, static
error normalization. Feature-specific transport handlers live beside their feature. 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 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 `backend/application.py`; endpoints with path parameters, body handling, or special error
semantics remain visible control flow in `RequestHandler`. semantics remain visible control flow in `RequestHandler`.
+8 -45
View File
@@ -128,14 +128,19 @@ class DashboardService(
profile_supplier=self._resolved_llm_profile, profile_supplier=self._resolved_llm_profile,
) )
self.screener.ensure_builtin_strategies() 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_refresh_tick,
self._background_stop,
interval_seconds=5, interval_seconds=5,
initial_delay_seconds=3, 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]: def _load_system_credentials(self, environment: dict[str, str]) -> dict[str, Any]:
encrypted = self.database.get_system_setting("credentials") encrypted = self.database.get_system_setting("credentials")
@@ -472,38 +477,6 @@ class DashboardService(
**self.database.status(), **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() SERVICE = DashboardService()
@@ -1020,16 +993,6 @@ class RequestHandler(
return return
self.send_json({"error": "Not found"}, HTTPStatus.NOT_FOUND) 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: def save_reason(self) -> None:
try: try:
body = self.read_json_body() body = self.read_json_body()
+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) parser.add_argument("--port", type=int, default=8765)
args = parser.parse_args() args = parser.parse_args()
server = ThreadingHTTPServer((args.host, args.port), handler_class) 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: 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() server.serve_forever()
except KeyboardInterrupt: except KeyboardInterrupt:
pass pass
finally: finally:
service._background_stop.set() service.stop_background_jobs()
server.server_close() server.server_close()
+29
View File
@@ -7,6 +7,35 @@ from typing import Any
class HeavenRepositoryMixin: 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 @staticmethod
def _heaven_reading_dict(row: sqlite3.Row | None) -> dict[str, Any] | None: def _heaven_reading_dict(row: sqlite3.Row | None) -> dict[str, Any] | None:
if not row: if not row:
+1 -16
View File
@@ -14,19 +14,4 @@ class MentorHttpMixin:
except (ValueError, json.JSONDecodeError) as exc: except (ValueError, json.JSONDecodeError) as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST) self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
return return
self.send_response(HTTPStatus.OK) self.send_ndjson_stream(stream, (ValueError, MentorAgentError))
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
+24
View File
@@ -117,6 +117,30 @@ class PoolServiceMixin:
finally: finally:
self._ifind_event_lock.release() 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 @staticmethod
def _normalize_ifind_event_time(value: Any) -> str: def _normalize_ifind_event_time(value: Any) -> str:
text = str(value or "").strip() text = str(value or "").strip()
+2 -16
View File
@@ -26,22 +26,8 @@ class ReviewHttpMixin:
except (ValueError, json.JSONDecodeError) as exc: except (ValueError, json.JSONDecodeError) as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST) self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
return return
self.send_response(HTTPStatus.OK) events = ({"type": "delta", "content": chunk} for chunk in stream)
self.send_header("Content-Type", "application/x-ndjson; charset=utf-8") self.send_ndjson_stream(events, (ValueError, ReviewAssistantError))
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: def save_watchlist(self) -> None:
try: try:
+29
View File
@@ -3,6 +3,7 @@ from __future__ import annotations
import json import json
import mimetypes import mimetypes
import secrets import secrets
from collections.abc import Iterable
from http import HTTPStatus from http import HTTPStatus
from http.cookies import SimpleCookie from http.cookies import SimpleCookie
from typing import Any from typing import Any
@@ -142,5 +143,33 @@ class HttpTransportMixin:
self.end_headers() self.end_headers()
self.wfile.write(content) 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: def log_message(self, format_string: str, *args: Any) -> None:
print(f"[{self.log_date_time_string()}] {format_string % args}") print(f"[{self.log_date_time_string()}] {format_string % args}")
+37 -19
View File
@@ -18,6 +18,9 @@ class InProcessJobRunner:
self.repository = repository self.repository = repository
self._locks: dict[str, threading.Lock] = {} self._locks: dict[str, threading.Lock] = {}
self._locks_guard = 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( def submit(
self, job_id: str, idempotency_key: str, action: JobAction, self, job_id: str, idempotency_key: str, action: JobAction,
@@ -52,27 +55,42 @@ class InProcessJobRunner:
return True return True
def start_scheduler( def start_scheduler(
self, callback: Callable[[], None], stop_event: threading.Event, self, callback: Callable[[], None], interval_seconds: float,
interval_seconds: float, initial_delay_seconds: float = 0, initial_delay_seconds: float = 0,
) -> threading.Thread: ) -> threading.Thread:
def schedule_loop() -> None: with self._scheduler_guard:
if stop_event.wait(initial_delay_seconds): current = self._scheduler_thread
return if current is not None and current.is_alive():
while not stop_event.is_set(): return current
try: self._scheduler_stop.clear()
callback()
except Exception:
# Submitted jobs persist their own failures; the scheduler must stay alive.
pass
stop_event.wait(interval_seconds)
thread = threading.Thread( def schedule_loop() -> None:
target=schedule_loop, if self._scheduler_stop.wait(initial_delay_seconds):
name="background-job-scheduler", return
daemon=True, while not self._scheduler_stop.is_set():
) try:
thread.start() callback()
return thread 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: def wait_for_idle(self, timeout_seconds: float = 5) -> bool:
deadline = time.monotonic() + max(0, timeout_seconds) deadline = time.monotonic() + max(0, timeout_seconds)
+14 -4
View File
@@ -287,6 +287,16 @@
"path": "backend/llm/transport.py" "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": [
"/shared/tokens.css?v=20260729-1", "/shared/tokens.css?v=20260729-1",
"/styles/styles.css", "/styles/styles.css",
@@ -364,8 +374,8 @@
}, },
{ {
"path": "backend/application.py", "path": "backend/application.py",
"bytes": 48749, "bytes": 47769,
"lines": 1129 "lines": 1092
}, },
{ {
"path": "frontend/styles/theme.css", "path": "frontend/styles/theme.css",
@@ -374,8 +384,8 @@
}, },
{ {
"path": "database.py", "path": "database.py",
"bytes": 33284, "bytes": 32073,
"lines": 746 "lines": 716
} }
] ]
} }
+1 -31
View File
@@ -2,7 +2,7 @@ from __future__ import annotations
import json import json
import sqlite3 import sqlite3
from datetime import datetime, timezone from datetime import datetime
from pathlib import Path from pathlib import Path
from typing import Any from typing import Any
@@ -665,36 +665,6 @@ class ReviewDatabase(
) )
MigrationRunner().apply(connection, MIGRATIONS) 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( 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]]:
+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()
+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") 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__": if __name__ == "__main__":
unittest.main() unittest.main()
@@ -40,6 +40,8 @@ class AccountSliceStructureTests(unittest.TestCase):
"require_access", "require_access",
"serve_static", "serve_static",
"send_json", "send_json",
"_write_stream_event",
"send_ndjson_stream",
): ):
self.assertNotIn(method, RequestHandler.__dict__) self.assertNotIn(method, RequestHandler.__dict__)
self.assertIn(method, HttpTransportMixin.__dict__) self.assertIn(method, HttpTransportMixin.__dict__)
@@ -38,6 +38,9 @@ HEAVEN_SERVICE_METHODS = {
} }
HEAVEN_REPOSITORY_METHODS = { HEAVEN_REPOSITORY_METHODS = {
"list_sector_phase_overrides",
"save_sector_phase_override",
"delete_sector_phase_override",
"_heaven_reading_dict", "_heaven_reading_dict",
"save_heaven_reading", "save_heaven_reading",
"list_heaven_readings", "list_heaven_readings",
@@ -29,6 +29,8 @@ POOL_METHODS = {
"_apply_reason_overrides", "_apply_reason_overrides",
"_schedule_ifind_event_enrichment", "_schedule_ifind_event_enrichment",
"_refresh_ifind_event_enrichment", "_refresh_ifind_event_enrichment",
"_ifind_field",
"_ifind_row_code",
"_normalize_ifind_event_time", "_normalize_ifind_event_time",
"_merge_ifind_event_enrichment", "_merge_ifind_event_enrichment",
} }
@@ -179,6 +179,10 @@ def build() -> dict[str, Any]:
{"function": "chat_completion", "path": "backend/llm/transport.py"}, {"function": "chat_completion", "path": "backend/llm/transport.py"},
{"function": "stream_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), "css_layers": css_layers(html),
"code_hotspots": code_hotspots(), "code_hotspots": code_hotspots(),
} }
+11 -6
View File
@@ -32,17 +32,22 @@ def main() -> None:
from server import RequestHandler, SERVICE from server import RequestHandler, SERVICE
server = ThreadingHTTPServer(("127.0.0.1", args.port), RequestHandler) 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: 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() server.serve_forever()
except KeyboardInterrupt: except KeyboardInterrupt:
pass pass
finally: finally:
SERVICE._background_stop.set() if hasattr(SERVICE, "stop_background_jobs"):
SERVICE.stop_background_jobs()
else:
SERVICE._background_stop.set()
server.server_close() server.server_close()
+89
View File
@@ -25,6 +25,10 @@
| CR-05 | 根级兼容入口 | 正式后端仍有五处通过迁移兼容模块反向导入规范实现 | 正式代码改用规范路径;兼容入口只服务原公开导入契约 | 已完成 | | CR-05 | 根级兼容入口 | 正式后端仍有五处通过迁移兼容模块反向导入规范实现 | 正式代码改用规范路径;兼容入口只服务原公开导入契约 | 已完成 |
| CR-06 | 紧凑日期显示 | Tushare与选股引擎各保留一份完全相同的`YYYYMMDD`显示转换 | 由`bootstrap/config.py`拥有唯一格式策略,消费者保留原局部别名 | 已完成 | | CR-06 | 紧凑日期显示 | Tushare与选股引擎各保留一份完全相同的`YYYYMMDD`显示转换 | 由`bootstrap/config.py`拥有唯一格式策略,消费者保留原局部别名 | 已完成 |
| CR-07 | 数据Provider组装 | 核查网关、容器和业务服务是否重复创建外部数据客户端 | 固化唯一创建位置及兼容例外,不改数据源语义 | 已完成 | | CR-07 | 数据Provider组装 | 核查网关、容器和业务服务是否重复创建外部数据客户端 | 固化唯一创建位置及兼容例外,不改数据源语义 | 已完成 |
| CR-08 | NDJSON流式传输 | 问师与复盘助手重复维护响应头、事件写入、完成、断线及关闭流程 | HTTP共享层拥有唯一流式连接生命周期 | 已完成 |
| CR-09 | 应用服务门面 | 两个仅服务股池iFinD补全的方法仍错位在全局`DashboardService` | 原函数体机械归位到`PoolServiceMixin` | 已完成 |
| CR-10 | Repository所有权 | 五行行业阶段覆盖的三项持久化方法仍错位在根级数据库门面 | 原函数体机械归位到问天Repository,兼容数据继续保留 | 已完成 |
| CR-11 | 后台任务生命周期 | 导入应用即启动调度线程,启动早于端口绑定,停止只置位但不等待 | 运行时显式启停,每个Runner只拥有一个可等待的调度线程 | 已完成 |
## CR-01验收口径 ## CR-01验收口径
@@ -185,6 +189,91 @@
本批基线为`xiaobai-reduction-06-date-formatting-20260801`;检查点为 本批基线为`xiaobai-reduction-06-date-formatting-20260801`;检查点为
`xiaobai-reduction-07-provider-ownership-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`
## 人工验收记录 ## 人工验收记录
- 2026-08-01:用户检查CR-02与CR-03运行结果,确认未发现明显异常。本记录仅表示本轮可见功能与 - 2026-08-01:用户检查CR-02与CR-03运行结果,确认未发现明显异常。本记录仅表示本轮可见功能与