diff --git a/app/ARCHITECTURE.md b/app/ARCHITECTURE.md index 8c6cb78..542357d 100644 --- a/app/ARCHITECTURE.md +++ b/app/ARCHITECTURE.md @@ -30,8 +30,9 @@ background scheduler `backend/application.py` and `backend/bootstrap/`. - `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, responses, static delivery, and - error normalization. Feature-specific transport handlers live beside their feature. +- `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`. diff --git a/app/backend/application.py b/app/backend/application.py index efea17b..f0abec4 100644 --- a/app/backend/application.py +++ b/app/backend/application.py @@ -1020,16 +1020,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() diff --git a/app/backend/features/mentor/http.py b/app/backend/features/mentor/http.py index 241b28f..cee3f0e 100644 --- a/app/backend/features/mentor/http.py +++ b/app/backend/features/mentor/http.py @@ -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)) diff --git a/app/backend/features/review/http.py b/app/backend/features/review/http.py index ad66e52..251bfeb 100644 --- a/app/backend/features/review/http.py +++ b/app/backend/features/review/http.py @@ -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: diff --git a/app/backend/http/handler.py b/app/backend/http/handler.py index 186fb7e..59b3643 100644 --- a/app/backend/http/handler.py +++ b/app/backend/http/handler.py @@ -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}") diff --git a/app/config/architecture-inventory.json b/app/config/architecture-inventory.json index 625fdeb..413b671 100644 --- a/app/config/architecture-inventory.json +++ b/app/config/architecture-inventory.json @@ -287,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", @@ -364,8 +374,8 @@ }, { "path": "backend/application.py", - "bytes": 48749, - "lines": 1129 + "bytes": 48513, + "lines": 1119 }, { "path": "frontend/styles/theme.css", diff --git a/app/tests/test_http_streaming.py b/app/tests/test_http_streaming.py new file mode 100644 index 0000000..1d24823 --- /dev/null +++ b/app/tests/test_http_streaming.py @@ -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() diff --git a/app/tests/test_preservation_slice_accounts.py b/app/tests/test_preservation_slice_accounts.py index f2c0eed..2c4170f 100644 --- a/app/tests/test_preservation_slice_accounts.py +++ b/app/tests/test_preservation_slice_accounts.py @@ -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__) diff --git a/app/tools/build_architecture_inventory.py b/app/tools/build_architecture_inventory.py index f5b84c6..948e651 100644 --- a/app/tools/build_architecture_inventory.py +++ b/app/tools/build_architecture_inventory.py @@ -179,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(), } diff --git a/docs/governance/code-reduction.md b/docs/governance/code-reduction.md index cbc57db..46b851f 100644 --- a/docs/governance/code-reduction.md +++ b/docs/governance/code-reduction.md @@ -25,6 +25,7 @@ | CR-05 | 根级兼容入口 | 正式后端仍有五处通过迁移兼容模块反向导入规范实现 | 正式代码改用规范路径;兼容入口只服务原公开导入契约 | 已完成 | | CR-06 | 紧凑日期显示 | Tushare与选股引擎各保留一份完全相同的`YYYYMMDD`显示转换 | 由`bootstrap/config.py`拥有唯一格式策略,消费者保留原局部别名 | 已完成 | | CR-07 | 数据Provider组装 | 核查网关、容器和业务服务是否重复创建外部数据客户端 | 固化唯一创建位置及兼容例外,不改数据源语义 | 已完成 | +| CR-08 | NDJSON流式传输 | 问师与复盘助手重复维护响应头、事件写入、完成、断线及关闭流程 | HTTP共享层拥有唯一流式连接生命周期 | 已完成 | ## CR-01验收口径 @@ -185,6 +186,26 @@ 本批基线为`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`。 + ## 人工验收记录 - 2026-08-01:用户检查CR-02与CR-03运行结果,确认未发现明显异常。本记录仅表示本轮可见功能与