Browse Source

新增mcp排障中心

jackson 4 days ago
parent
commit
e9b818b2df

+ 13 - 0
.env.example

@@ -53,3 +53,16 @@ FMS_REDIS_PREFIX=fms:mcp:gateway:
 
 # tools/call rate limiting (default: 60 requests/minute per device session and tool)
 # initialize and tools/list are not rate limited
+
+# MCP diagnosis events are sent asynchronously to support. Keep disabled until
+# the support collector and matching HMAC key are deployed.
+MCP_DIAGNOSIS_ENABLED=false
+MCP_DIAGNOSIS_URL=https://support.example.com/internal/mcp-diagnostics/events
+MCP_DIAGNOSIS_KEY_ID=gateway-current
+MCP_DIAGNOSIS_SECRET=
+MCP_DIAGNOSIS_ALLOW_INSECURE_HTTP=false
+MCP_DIAGNOSIS_QUEUE_SIZE=1000
+MCP_DIAGNOSIS_BATCH_SIZE=100
+MCP_DIAGNOSIS_TIMEOUT_SECONDS=0.5
+MCP_DIAGNOSIS_INITIAL_BACKOFF_SECONDS=0.25
+MCP_DIAGNOSIS_MAX_BACKOFF_SECONDS=5.0

+ 5 - 5
AGENTS.md

@@ -32,11 +32,9 @@ Other repositories should link to this repository's README and project docs inst
 ## Documentation Routes
 
 - `README.md`: current tool catalog, routes, configuration, operations, and troubleshooting.
-- `project-docs/overview.md`: current system summary and cross-repository ownership map.
-- `project-docs/requirements.md`: current protocol, security, and output contracts.
-- `project-docs/tech-specs.md`: architecture, session, Presenter, logging, and rate-limit details.
-- `project-docs/user-structure.md`: user, request, and developer flows plus directory map.
-- `project-docs/timeline.md`: short milestones and current operational handoff only.
+- `Y:/all_project_docs/mcp/overview.md`: current system summary and cross-repository ownership map.
+- The same directory's `requirements.md`, `tech-specs.md`, `user-structure.md`, and `timeline.md`: current contracts, architecture, flows, and operational handoff.
+- `Y:/all_project_docs/mcp/guides/mcp-team-sharing.md`: current team onboarding and sharing guide.
 - `docs/superpowers/specs/`: dated design decisions and ADR-like context.
 
 The following are not authoritative and are not read by default:
@@ -45,6 +43,8 @@ The following are not authoritative and are not read by default:
 - Completed `docs/superpowers/plans/**`: execution history, not the current contract.
 When a temporary plan finishes, promote durable facts into README/project docs, record the milestone once, then retire the plan. Do not keep completed checklists as active instructions.
 
+If `Y:/all_project_docs` is unavailable, stop changes to contracts, architecture, databases, or deployment behavior and report the missing documentation repository instead of relying on stale information.
+
 ## Change Rules
 
 - Update tool metadata, both registries, CLI forwarding, Presenter mapping, and tests together when adding or changing a tool.

+ 47 - 7
README.md

@@ -58,9 +58,10 @@ Gateway 在 `tools/call` 最终边界处理展示字段,不改变 ThinkPHP 内
 - 三个筛选项工具保留“可传值、显示名称、业务编码”,确保返回值可继续传给查询或导出工具。
 - 两个导出工具返回 `files[].label + files[].url`,不暴露后端 `file_url` 键。
 - `request_id` 位于 MCP 结果 `_meta`;参数错误使用业务名称,未知异常不透传后端细节。
+- stdio 与公网 `tools/call` 缺少有效工具名时返回 JSON-RPC `-32602 Invalid params`,不包装为业务 `isError`。
 - stdio 与 public 模式共用 `services/output_presenter.py`;未知工具或畸形响应关闭失败。
 
-详细设计见 [技术规格的 Presenter 分类](project-docs/tech-specs.md#presenter-分类)。
+详细设计见集中文档仓的 [技术规格 Presenter 分类](../all_project_docs/mcp/tech-specs.md#presenter-分类)。
 
 省外进港资料导出使用 `outbound_numbers`、`container_codes`、`bl_numbers`、`so_numbers` 四个号码数组之一,并同时提供 `file_type=NB/SH/MS`。用户未明确号码类型时,AI 必须先让用户从排舱单号、柜号、提单号、SO 号中选择,确认前不得调用;提单号指后台排舱单列表的普通提单号,SO 号对应 `fms_booking_detail.so_number`。
 
@@ -70,11 +71,12 @@ Gateway 在 `tools/call` 最终边界处理展示字段,不改变 ThinkPHP 内
 
 - `AGENTS.md`:项目所有权、红线、协作规则和验证命令。
 - `README.md`:当前工具目录、接入配置、运行方式、运维与排障。
-- `project-docs/overview.md`:当前系统概况与跨仓职责。
-- `project-docs/requirements.md`:现行协议、安全、工具选择和输出合同。
-- `project-docs/tech-specs.md`:组件、会话、展示、限流、追踪与日志设计。
-- `project-docs/user-structure.md`:员工、调试、工具选择和跨仓开发流程。
-- `project-docs/timeline.md`:短里程碑和当前交接状态。
+- [`Y:/all_project_docs/mcp/guides/mcp-team-sharing.md`](../all_project_docs/mcp/guides/mcp-team-sharing.md):面向团队的从 0 到 1 架构、实践与问题复盘分享稿。
+- [`Y:/all_project_docs/mcp/overview.md`](../all_project_docs/mcp/overview.md):当前系统概况与跨仓职责。
+- [`requirements.md`](../all_project_docs/mcp/requirements.md):现行协议、安全、工具选择和输出合同。
+- [`tech-specs.md`](../all_project_docs/mcp/tech-specs.md):组件、会话、展示、限流、追踪与日志设计。
+- [`user-structure.md`](../all_project_docs/mcp/user-structure.md):员工、调试、工具选择和跨仓开发流程。
+- [`timeline.md`](../all_project_docs/mcp/timeline.md):短里程碑和当前交接状态。
 
 临时规划与已完成的执行计划不作为现行合同;交付后应把稳定事实并回上述入口,并从活动工作区清理过程文件。
 
@@ -342,6 +344,32 @@ FMS_REDIS_PASSWORD=change-me
 FMS_REDIS_PREFIX=fms:mcp:gateway:
 ```
 
+Gateway 诊断事件上报默认关闭。Support 的 internal collector、Mongo 索引和
+HMAC 密钥部署完成后,在 Gateway `.env` 增加:
+
+四项目统一发布时按集中文档仓的 [`MCP 排障中心全版本部署清单`](../all_project_docs/support/operations/mcp-diagnosis-full-deployment.md) 操作;Gateway 必须放在 Collector、Mongo 索引和两套 PHP Worker 之后灰度启用。
+
+```dotenv
+MCP_DIAGNOSIS_ENABLED=true
+MCP_DIAGNOSIS_URL=https://support.example.com/internal/mcp-diagnostics/events
+MCP_DIAGNOSIS_KEY_ID=gateway-current
+MCP_DIAGNOSIS_SECRET=replace-with-at-least-32-random-characters
+MCP_DIAGNOSIS_ALLOW_INSECURE_HTTP=false
+MCP_DIAGNOSIS_QUEUE_SIZE=1000
+MCP_DIAGNOSIS_BATCH_SIZE=100
+MCP_DIAGNOSIS_TIMEOUT_SECONDS=0.5
+MCP_DIAGNOSIS_INITIAL_BACKOFF_SECONDS=0.25
+MCP_DIAGNOSIS_MAX_BACKOFF_SECONDS=5.0
+```
+
+生产环境 `MCP_DIAGNOSIS_URL` 必须使用 HTTPS。只有隔离测试环境可显式设置 `MCP_DIAGNOSIS_ALLOW_INSECURE_HTTP=true` 使用 HTTP,默认值 `false` 会使 HTTP 配置回退为 `NullDiagnosticReporter`。
+
+`MCP_DIAGNOSIS_KEY_ID` 和 `MCP_DIAGNOSIS_SECRET` 必须与 Support
+`mcp_diagnosis.php` 的 Gateway current/previous key 对应。Reporter 使用内存有界队列、
+批量 HMAC 和短超时异步发送;Support 超时、拒绝、队列满或线程启动失败均不改变 MCP
+响应。进程异常退出时尚未发送的内存事件可能丢失,因此该链路用于排障观测,不作为业务
+审计唯一依据。
+
 兼容旧配置键:
 
 - `MCP_AUTH_BASE_URL`
@@ -406,7 +434,7 @@ ThinkPHP MCP 路由位于各后端项目的 `route/mcp/mcp_route.php`。
 
 Gateway 会调用以下路径:
 
-- Auth:`/mcp/auth/exchange`、`/mcp/auth/refresh`、`/mcp/auth/revoke`。
+- Auth:本地兼容会话只使用 `/mcp/auth/refresh`、`/mcp/auth/revoke`;正式公网设备配置由 base 的登录态设备接口创建,不调用已退役的 `/mcp/auth/exchange`。
 - 动态工具列表:`/mcp/tools/listEnabledTools`。
 - 查询:`/mcp/tools/queryOrder`、`/mcp/tools/queryTrack`、`/mcp/tools/queryOrderExact`、`/mcp/tools/queryOrderDetail`、`/mcp/tools/queryCustomsDeclarationFiles`、`/mcp/tools/queryOutboundList`、`/mcp/tools/queryOutboundDetail`。
 - 筛选项:`/mcp/tools/listOutboundFilterOptions`、`/mcp/tools/listOrderFilterOptions`、`/mcp/tools/listPendingOutboundExportFilterOptions`。
@@ -416,6 +444,18 @@ Gateway 会调用以下路径:
 
 ## 排障
 
+跨 Gateway、PHP 和数据库访问日志的排查顺序见集中文档仓的 [`mcp-team-sharing.md`](../all_project_docs/mcp/guides/mcp-team-sharing.md)“排错时怎么看日志”章节。
+
+若 Support 页面只能看到 fmsoperate 事件、看不到 Gateway 前置阶段,依次确认:
+
+1. Gateway `.env` 的 `MCP_DIAGNOSIS_ENABLED=true`,URL 指向已部署的 internal collector。
+2. Gateway 与 Support 使用相同的 key id 和至少 32 字符的 secret,且系统时钟误差在 Support 允许范围内。
+3. Support collector 已开启,Redis nonce 存储可用,Mongo 事件索引已在目标环境验证。
+4. 修改配置后已重启 Gateway;先灰度单个实例,再模拟设备失效或限流并从 Support 查询。
+
+不要在日志、命令历史、工单或群聊中打印 `MCP_DIAGNOSIS_SECRET`、请求体、`GWS_*`、
+`MT_*`、Authorization 或 Cookie。Gateway 事件只保存 `GWS_*` 的 SHA-256 前 12 位。
+
 ### query_order 返回 UTF-8 BOM 错误
 
 如果 Workbuddy 提示:

+ 42 - 13
app.py

@@ -11,6 +11,10 @@ from public_server import serve_public
 from services.api_client import ApiClient
 from services.auth_client import AuthClient
 from services.gateway_session_store import GatewaySessionStore
+from services.diagnostic_reporter import (
+    NullDiagnosticReporter,
+    diagnostic_reporter_from_config,
+)
 from services.scoped_api_client import ScopedApiClient
 from services.token_store import FileTokenStore, RedisSocketClient, RedisTokenStore
 from tools.list_order_filter_options import ListOrderFilterOptionsTool
@@ -49,10 +53,17 @@ def parse_string_list(value):
 
 
 class GatewayApp:
-    def __init__(self, auth_client=None, api_client=None, token_store=None):
+    def __init__(
+        self,
+        auth_client=None,
+        api_client=None,
+        token_store=None,
+        reporter=None,
+    ):
         self.auth_client = auth_client
         self.api_client = api_client
         self.token_store = token_store
+        self.reporter = reporter or NullDiagnosticReporter()
         self._tools = {
             'query_order': QueryOrderTool(api_client=api_client),
             'query_track': QueryTrackTool(api_client=api_client),
@@ -113,7 +124,12 @@ class GatewayApp:
             token_store=token_store,
             timeout=config.timeout_seconds,
         )
-        return cls(auth_client=auth_client, api_client=api_client, token_store=token_store)
+        return cls(
+            auth_client=auth_client,
+            api_client=api_client,
+            token_store=token_store,
+            reporter=diagnostic_reporter_from_config(config),
+        )
 
     def registered_tool_names(self):
         return tuple(self._tools.keys())
@@ -179,7 +195,7 @@ class GatewayApp:
         return tool.call(request_id=request_id, **arguments)
 
     def create_protocol_handler(self):
-        return McpProtocolHandler(self)
+        return McpProtocolHandler(self, reporter=self.reporter)
 
     def run_cli(self, argv=None, stdin=None, stdout=None):
         stdin = stdin or sys.stdin
@@ -248,9 +264,18 @@ class GatewayApp:
         if args.command == 'list-tools':
             payload = self.list_tools()
         elif args.command == 'serve-stdio':
-            return self.create_protocol_handler().run_stdio(stdin=stdin, stdout=stdout)
+            try:
+                return self.create_protocol_handler().run_stdio(
+                    stdin=stdin,
+                    stdout=stdout,
+                )
+            finally:
+                self.reporter.close()
         elif args.command == 'serve-public':
             config = GatewayConfig.from_env()
+            reporter = self.reporter
+            if isinstance(reporter, NullDiagnosticReporter):
+                reporter = diagnostic_reporter_from_config(config)
             redis = RedisSocketClient(
                 host=config.redis_host,
                 port=config.redis_port,
@@ -267,15 +292,19 @@ class GatewayApp:
                 session_store=session_store,
                 api_client=ScopedApiClient(config.tools_base_url, timeout=config.timeout_seconds),
             )
-            return serve_public(
-                public_app,
-                host=args.host,
-                port=args.port,
-                enable_rate_limit=config.rate_limit_enabled,
-                rate_limit_max_requests=config.rate_limit_max_requests,
-                rate_limit_window_seconds=config.rate_limit_window_seconds,
-                max_in_flight_per_tool=config.max_in_flight_per_tool,
-            )
+            try:
+                return serve_public(
+                    public_app,
+                    host=args.host,
+                    port=args.port,
+                    enable_rate_limit=config.rate_limit_enabled,
+                    rate_limit_max_requests=config.rate_limit_max_requests,
+                    rate_limit_window_seconds=config.rate_limit_window_seconds,
+                    max_in_flight_per_tool=config.max_in_flight_per_tool,
+                    reporter=reporter,
+                )
+            finally:
+                reporter.close()
         elif args.command == 'call':
             tool_args = {
                 'page': args.page,

+ 69 - 0
config.py

@@ -24,6 +24,16 @@ class GatewayConfig:
     rate_limit_max_requests: int = 60
     rate_limit_window_seconds: int = 60
     max_in_flight_per_tool: int = 2
+    diagnosis_enabled: bool = False
+    diagnosis_url: str = ''
+    diagnosis_key_id: str = ''
+    diagnosis_secret: str = ''
+    diagnosis_allow_insecure_http: bool = False
+    diagnosis_queue_size: int = 1000
+    diagnosis_batch_size: int = 100
+    diagnosis_timeout_seconds: float = 0.5
+    diagnosis_initial_backoff_seconds: float = 0.25
+    diagnosis_max_backoff_seconds: float = 5.0
 
     @classmethod
     def from_env(cls, env=None, dotenv_path=''):
@@ -75,6 +85,55 @@ class GatewayConfig:
             rate_limit_max_requests=cls._parse_int(cls._pick(dotenv_env, primary_env, 'FMS_RATE_LIMIT_MAX_REQUESTS', 'MCP_RATE_LIMIT_MAX_REQUESTS'), default=60),
             rate_limit_window_seconds=cls._parse_int(cls._pick(dotenv_env, primary_env, 'FMS_RATE_LIMIT_WINDOW_SECONDS', 'MCP_RATE_LIMIT_WINDOW_SECONDS'), default=60),
             max_in_flight_per_tool=cls._parse_int(cls._pick(dotenv_env, primary_env, 'FMS_MAX_IN_FLIGHT_PER_TOOL', 'MCP_MAX_IN_FLIGHT_PER_TOOL'), default=2),
+            diagnosis_enabled=cls._parse_bool(
+                cls._pick(dotenv_env, primary_env, 'MCP_DIAGNOSIS_ENABLED'),
+                default=False,
+            ),
+            diagnosis_url=cls._pick(
+                dotenv_env, primary_env, 'MCP_DIAGNOSIS_URL'
+            ).rstrip('/'),
+            diagnosis_key_id=cls._pick(
+                dotenv_env, primary_env, 'MCP_DIAGNOSIS_KEY_ID'
+            ),
+            diagnosis_secret=cls._pick(
+                dotenv_env, primary_env, 'MCP_DIAGNOSIS_SECRET'
+            ),
+            diagnosis_allow_insecure_http=cls._parse_bool(
+                cls._pick(
+                    dotenv_env,
+                    primary_env,
+                    'MCP_DIAGNOSIS_ALLOW_INSECURE_HTTP',
+                ),
+                default=False,
+            ),
+            diagnosis_queue_size=cls._parse_int(
+                cls._pick(dotenv_env, primary_env, 'MCP_DIAGNOSIS_QUEUE_SIZE'),
+                default=1000,
+            ),
+            diagnosis_batch_size=cls._parse_int(
+                cls._pick(dotenv_env, primary_env, 'MCP_DIAGNOSIS_BATCH_SIZE'),
+                default=100,
+            ),
+            diagnosis_timeout_seconds=cls._parse_float(
+                cls._pick(dotenv_env, primary_env, 'MCP_DIAGNOSIS_TIMEOUT_SECONDS'),
+                default=0.5,
+            ),
+            diagnosis_initial_backoff_seconds=cls._parse_float(
+                cls._pick(
+                    dotenv_env,
+                    primary_env,
+                    'MCP_DIAGNOSIS_INITIAL_BACKOFF_SECONDS',
+                ),
+                default=0.25,
+            ),
+            diagnosis_max_backoff_seconds=cls._parse_float(
+                cls._pick(
+                    dotenv_env,
+                    primary_env,
+                    'MCP_DIAGNOSIS_MAX_BACKOFF_SECONDS',
+                ),
+                default=5.0,
+            ),
         )
 
     @staticmethod
@@ -144,6 +203,16 @@ class GatewayConfig:
             return default
         return value not in ('0', 'false', 'no', 'off')
 
+    @staticmethod
+    def _parse_float(raw, default):
+        value = GatewayConfig._strip_comment(str(raw or '').strip())
+        if not value:
+            return default
+        try:
+            return float(value)
+        except ValueError:
+            return default
+
     @staticmethod
     def _build_default_session_key(env):
         computer = env.get('COMPUTERNAME') or env.get('HOSTNAME') or os.environ.get('COMPUTERNAME') or os.environ.get('HOSTNAME') or 'unknown-computer'

+ 114 - 12
mcp_protocol.py

@@ -4,6 +4,8 @@ import sys
 import uuid
 
 from services.output_presenter import OutputPresenter
+from services.diagnostic_event import RequestDiagnosticEmitter
+from services.diagnostic_reporter import NullDiagnosticReporter
 
 
 logger = logging.getLogger(__name__)
@@ -15,8 +17,9 @@ class McpProtocolHandler:
     server_version = '0.1.0'
     output_presenter = OutputPresenter()
 
-    def __init__(self, gateway_app):
+    def __init__(self, gateway_app, reporter=None):
         self.gateway_app = gateway_app
+        self.reporter = reporter or NullDiagnosticReporter()
         self.initialized = False
 
     def handle_message(self, message):
@@ -40,8 +43,28 @@ class McpProtocolHandler:
         method = str(request.get('method') or '').strip()
         tool_name = ''
         trace_request_id = self._build_trace_request_id()
+        emitter = RequestDiagnosticEmitter(self.reporter, trace_request_id)
+        emitter.emit(
+            stage='request_ingress',
+            status='started',
+            event_code='REQUEST_RECEIVED',
+            context={
+                'jsonrpc_method': method or 'unknown',
+                'transport': 'stdio',
+            },
+        )
+        backend_started = False
         try:
             if method == 'initialize':
+                emitter.emit(
+                    stage='protocol_validation',
+                    status='succeeded',
+                    event_code='PROTOCOL_VALIDATION_COMPLETED',
+                    context={
+                        'jsonrpc_method': method,
+                        'transport': 'stdio',
+                    },
+                )
                 self.initialized = True
                 return self._success_response(
                     request_id,
@@ -59,6 +82,15 @@ class McpProtocolHandler:
                     },
                 )
             if method == 'tools/list':
+                emitter.emit(
+                    stage='protocol_validation',
+                    status='succeeded',
+                    event_code='PROTOCOL_VALIDATION_COMPLETED',
+                    context={
+                        'jsonrpc_method': method,
+                        'transport': 'stdio',
+                    },
+                )
                 tools = [
                     self._normalize_tool(tool)
                     for tool in self.gateway_app.list_tools(
@@ -72,22 +104,83 @@ class McpProtocolHandler:
                     raise ValueError('tool parameters must be an object')
                 raw_tool_name = params.get('name')
                 if not isinstance(raw_tool_name, str) or not raw_tool_name.strip():
-                    raise ValueError('tool name must be a non-empty string')
+                    emitter.emit(
+                        stage='protocol_validation',
+                        status='failed',
+                        event_code='PARAM_VALIDATION_FAILED',
+                        context={
+                            'jsonrpc_method': method,
+                            'jsonrpc_code': -32602,
+                            'transport': 'stdio',
+                        },
+                    )
+                    return self._error_response(
+                        request_id,
+                        -32602,
+                        'Invalid params',
+                        trace_request_id,
+                    )
                 tool_name = raw_tool_name.strip()
                 arguments = params.get('arguments') or {}
-                if tool_name == 'query_order':
-                    tool_result = self.gateway_app.call_tool(tool_name, arguments)
-                else:
-                    tool_result = self.gateway_app.call_tool(
+                emitter.emit(
+                    stage='protocol_validation',
+                    status='succeeded',
+                    event_code='PROTOCOL_VALIDATION_COMPLETED',
+                    tool_code=tool_name,
+                    context={
+                        'jsonrpc_method': method,
+                        'transport': 'stdio',
+                    },
+                )
+                backend_started = True
+                emitter.emit(
+                    stage='backend_call',
+                    status='started',
+                    event_code='BACKEND_CALL_STARTED',
+                    tool_code=tool_name,
+                    context={'transport': 'stdio'},
+                )
+                tool_result = self.gateway_app.call_tool(
+                    tool_name,
+                    arguments,
+                    request_id=trace_request_id,
+                )
+                backend_started = False
+                emitter.emit(
+                    stage='backend_call',
+                    status='succeeded',
+                    event_code='BACKEND_CALL_COMPLETED',
+                    tool_code=tool_name,
+                    response_code=(
+                        tool_result.get('code')
+                        if isinstance(tool_result, dict)
+                        else None
+                    ),
+                    context={'transport': 'stdio'},
+                )
+                try:
+                    response = self._tool_call_response(
+                        request_id,
                         tool_name,
-                        arguments,
-                        request_id=trace_request_id,
+                        tool_result,
                     )
-                return self._tool_call_response(
-                    request_id,
-                    tool_name,
-                    tool_result,
+                except Exception:
+                    emitter.emit(
+                        stage='response_safety',
+                        status='failed',
+                        event_code='RESPONSE_SAFETY_REJECTED',
+                        tool_code=tool_name,
+                        context={'transport': 'stdio'},
+                    )
+                    raise
+                emitter.emit(
+                    stage='response_safety',
+                    status='succeeded',
+                    event_code='RESPONSE_SAFETY_COMPLETED',
+                    tool_code=tool_name,
+                    context={'transport': 'stdio'},
                 )
+                return response
             return self._error_response(
                 request_id,
                 -32601,
@@ -113,6 +206,15 @@ class McpProtocolHandler:
                         'exception_class': exc.__class__.__name__,
                     },
                 )
+                if backend_started:
+                    emitter.emit(
+                        stage='backend_call',
+                        status='failed',
+                        event_code=diagnostic_reason,
+                        tool_code=tool_name or None,
+                        response_code='MCP_9001',
+                        context={'transport': 'stdio'},
+                    )
                 return self._tool_exception_response(
                     request_id,
                     tool_name,

+ 0 - 44
project-docs/overview.md

@@ -1,44 +0,0 @@
-# 项目概述
-
-## 项目定位
-
-`mcp` 是物流货运系统面向 Workbuddy 的 Python MCP Gateway。它负责 MCP 工具 Schema 与注册、stdio/public 协议适配、设备会话转发、安全展示 DTO、追踪日志和运行手册,不拥有物流业务规则。
-
-Gateway 保持薄边界:不直连业务数据库,不在 Python 中复制员工、公司、菜单或数据权限。员工身份、工具权限和业务数据范围由 ThinkPHP 在每次请求中最终校验。
-
-## 当前能力
-
-本地 `GatewayApp` 与公网 `PublicGatewayApp` 注册同一组 12 个候选工具:
-
-| 类别 | 工具 |
-|---|---|
-| 订单与轨迹 | `query_order`、`query_order_exact`、`query_order_detail`、`query_track` |
-| 报关与排舱 | `query_customs_declaration_files`、`query_outbound_list`、`query_outbound_detail` |
-| 筛选项 | `list_outbound_filter_options`、`list_order_filter_options`、`list_pending_outbound_export_filter_options` |
-| 导出 | `export_pending_outbound_orders`、`export_out_of_province_port_data` |
-
-候选注册不代表员工一定可见。实际工具集合是本地注册集合与 fmsoperate 动态启用列表的交集,并继续受员工权限与数据范围约束。动态列表不可用或格式异常时关闭访问。
-
-`query_order` 保持旧响应兼容;其余 11 个工具由 `services/output_presenter.py` 转为中文白名单展示 DTO。订单详情按固定中文模块展示,图片附件保留文件名和安全链接,机器字段及内部 ID 关闭失败。所有工具遵守零推断边界:结果意图或号码业务类型不明确时先询问,禁止按格式猜测、并行试查或失败后跨字段、跨工具重试。
-
-## 跨仓职责
-
-| 仓库 | 所有权 |
-|---|---|
-| `Y:/mcp` | 工具 Schema、双注册表、协议、展示、Gateway 会话转发、运行与排障入口 |
-| `Y:/fmsoperate` | 工具路由、Validate/Logic/Model、动态注册表、业务权限、数据范围、导出和访问日志 |
-| `Y:/base` | 员工身份、设备会话、MCP token 生命周期和认证决策 |
-| `Y:/home` | 员工入口 UI 与薄代理,不承载认证或物流业务逻辑 |
-| `Y:/settlement_tests` | 跨仓 PHP 合同测试 |
-
-## 文档权威层
-
-1. 代码与运行时数据决定真实行为:代码给出候选注册和合同,fmsoperate 当前启用值决定实际可见工具。
-2. `README.md` 是接入、工具目录、配置、运维和排障入口。
-3. `project-docs/` 五件套记录现行产品、协议、架构、流程和短里程碑。
-4. `AGENTS.md` 只记录 AI 必须遵守的所有权、红线和协作规则。
-5. 临时规划和已完成的执行计划不属于权威层;交付后把稳定事实并回上述入口,并清理过程文件。
-
-## 当前交接状态
-
-用户已确认 `Y:/fmsoperate/sql` 原有 10 个 MCP SQL 与 `Y:/base/sql` 的 5 个 MCP SQL 均已执行。订单详情注册 SQL 尚待目标环境执行;仓库仍无法静态判断动态注册表当前启用值、运行中 Gateway 是否已重启、Workbuddy 是否已重新连接,这些状态必须从部署环境核验。

+ 0 - 67
project-docs/requirements.md

@@ -1,67 +0,0 @@
-# 需求与功能清单
-
-## 产品边界
-
-- Gateway 只处理 MCP 协议、会话转发、候选工具注册和展示适配,不直连业务数据库,不复制 ThinkPHP 业务规则。
-- 所有员工、公司、菜单、工具和数据权限由 ThinkPHP 最终判断;Gateway 不接受调用方提供的身份或权限覆盖参数。
-- 本地 stdio 与公网 HTTP 必须注册同一组候选工具,并共享同一套协议与展示合同。
-- 写入、审批、费用修改、订单状态流转或批量业务变更工具不得在未完成独立安全评审时开放。
-
-## 协议合同
-
-- 支持 MCP `initialize`、`tools/list` 和 `tools/call`;能力声明保持 `tools.listChanged=false`。
-- `tools/list` 返回本地 12 个候选工具与 fmsoperate 动态启用列表的交集。
-- `tools/call` 在调用前再次校验工具已注册、设备会话有效且后端仍允许当前员工使用该工具。
-- 动态列表缺失、格式错误或查询失败时关闭访问,不允许回退为本地全量工具。
-- 后端业务失败仍放在 JSON-RPC `result` 中并设置 `isError=true`,不得中断 stdio 或 HTTP 会话;协议、参数或未知方法错误使用 JSON-RPC `error`。
-- CLI `call` 保留工具原始 `code/msg/data/meta` 信封,不经过面向 Workbuddy 的 Presenter。
-
-## 候选工具
-
-| 类别 | 工具 | 核心要求 |
-|---|---|---|
-| 查询 | `query_order` | 保持旧展示协议兼容 |
-| 查询 | `query_order_exact`、`query_track` | `query_order_exact` 只返回订单列表/批量精准筛选结果;只在结果意图和号码类型明确后调用 |
-| 查询 | `query_order_detail` | 只按明确订单号查询单个订单详情和中文模块;多个订单需逐个调用;列表、批量筛选不得使用本工具;完整详情可使用“全部”一次查询概览及十个明细模块 |
-| 查询 | `query_customs_declaration_files` | 订单号数组与排舱单号数组必须二选一 |
-| 查询 | `query_outbound_list` | 排舱阶段必选;运输方式未指定时查询全部 |
-| 查询 | `query_outbound_detail` | 只按明确排舱单号查询,不接受内部排舱 ID |
-| 筛选 | `list_outbound_filter_options` | 统一返回排舱阶段、运输方式、集货仓库、是否直送柜、拖车、报关和清关七类选项 |
-| 筛选 | `list_order_filter_options` | 返回当前员工可用的精准订单筛选值 |
-| 筛选 | `list_pending_outbound_export_filter_options` | 返回当前员工可用的未排舱导出筛选值 |
-| 导出 | `export_pending_outbound_orders` | 复用筛选工具返回值,不猜测内部 ID |
-| 导出 | `export_out_of_province_port_data` | 排舱单号、柜号、提单号、SO 号四类数组严格四选一,并要求文件类型 |
-
-## 零推断与工具选择
-
-- AI 必须先确认用户要看的结果类别,再确认所给号码的业务类型;任一项不明确时先提问。
-- 禁止按号码外观猜测类型,禁止并行调用多个候选工具,禁止跨字段试查,禁止一次失败后自行换号码字段或换工具重试。
-- 多个筛选值命中时必须让用户选择;唯一命中才可把内部 `value` 传给目标工具。
-- 面向用户描述筛选条件时使用中文标签,不展示英文筛选字段、内部数字代码或技术参数文案。
-- 中文筛选展示约束不影响查询结果中的件数、重量、体积、日期和其他正常业务值。
-
-## 输出与错误安全
-
-- `query_order` 是唯一旧展示例外,继续使用 `columns + records` 和原文本行为。
-- 其余 11 个工具由 `OutputPresenter` 白名单处理:4 个表格工具、1 个排舱详情工具、1 个订单详情工具、3 个筛选项工具和 2 个导出工具。
-- 表格使用中文 `headers + rows + pagination`;排舱详情使用 `summary + details + pagination`;订单详情使用固定中文分组、明细和分页;导出只返回 `files[].label + files[].url`。
-- 订单详情未知模块、分组或字段必须关闭失败;“全部”响应必须严格包含概览和十个明细分组,各明细分组独立带分页;附件可展示分类、文件名、类型、图片标识、预览和下载链接,商品/入库/查验图片不展示 URL。
-- 筛选项保留可继续传参的 `value`、用户显示 `label` 和业务 `code`,但不得泄露未列入白名单的后端字段。
-- `request_id` 放在 MCP 结果 `_meta`;安全工具的业务错误使用固定中文消息,不透传后端 `msg`、异常数据、堆栈或原始响应。
-- Gateway 已知错误码与重试策略以 `OutputPresenter.ERROR_MESSAGES`、`NON_RETRYABLE_CODES` 为代码源;未知业务码对外归一为 `MCP_9001`。
-- stdio 与 public 共用同一个 Presenter;未知工具、畸形响应、未知字段和未知异常关闭失败。
-
-## 公网安全与运行合同
-
-- 公网请求使用 `GWS_xxx` 查找 Redis Gateway session,并从服务端会话取得 `mcp_token`;不得把 token 返回给客户端。
-- 公网入口生成可信 `rq_http_*` 追踪号;调用方 `X-Request-Id` 只允许记录短哈希,不得作为可信追踪号。
-- 日志不得记录明文 Gateway 凭据、MCP token、Cookie、Authorization、授权码或后端原始响应。
-- `initialize` 与 `tools/list` 不限流;`tools/call` 按已认证的 `gateway_session_id + tool_name` 使用可选滑动窗口配额和可释放的并发配额,超限返回 JSON-RPC `-32029`,不得关闭 MCP 连接。窗口配置为 `0` 时关闭累计次数限制;并发名额在请求结束或异常后立即释放。
-- 公网必须位于 HTTPS 反向代理后,Redis 只允许内网访问并启用生产密码。
-
-## 变更与验收
-
-- 工具变更必须同步更新工具 metadata、stdio/public 双注册表、CLI 转发、Presenter、ThinkPHP 路由/验证/Logic/Model、动态注册数据和相关测试。
-- 因 `tools.listChanged=false`,工具或 Schema 变化后必须重启 Gateway 并让客户端重新连接。
-- 用户已确认 fmsoperate 原有 10 个 MCP SQL 与 base 的 5 个 MCP SQL 均已执行;订单详情注册 SQL 尚待执行,动态启用值仍以运行环境为准。
-- Python 生产代码变更必须运行全量 unittest 和 `.coveragerc` 要求的语句、分支 100% 严格覆盖率;跨仓 PHP 变更必须运行对应合同测试与语法检查。

+ 0 - 90
project-docs/tech-specs.md

@@ -1,90 +0,0 @@
-# 技术规格
-
-## 技术栈与边界
-
-- Gateway:Python 3,核心运行模块只使用标准库;stdio 使用逐行 JSON-RPC,公网使用 `ThreadingHTTPServer`。
-- 后端:ThinkPHP 6.0 + PHP 7.1+,按 Controller / Logic / Model / Validate 分层实现业务工具。
-- 存储:Redis 保存本地 token 或公网 Gateway session;MySQL、MongoDB 和文件存储只由所属业务系统访问。
-- Gateway 不连接业务数据库;跨仓调用只通过明确的 HTTP 路由和服务端会话进行。
-
-## 组件
-
-| 组件 | 职责 |
-|---|---|
-| `app.py` | 本地 `GatewayApp`、stdio/CLI 入口、12 个工具注册与参数转发 |
-| `public_gateway.py` | 公网 `PublicGatewayApp`、设备会话解析后的 scoped 调用、同组 12 个工具注册 |
-| `public_server.py` | `/mcp` HTTP JSON-RPC、可信追踪号、审计上下文、限流和 `/health` |
-| `mcp_protocol.py` | MCP 握手、`tools/list`、`tools/call`、旧/新结果分流和协议错误 |
-| `services/api_client.py` | 本地 token 上下文下的 ThinkPHP 工具请求与动态列表请求 |
-| `services/scoped_api_client.py` | 公网请求级 `mcp_token` 转发,避免进程级身份串用 |
-| `services/gateway_session_store.py` | `GWS_xxx` 到服务端员工会话的 Redis 映射 |
-| `services/output_presenter.py` | 11 个安全工具的白名单 DTO、错误映射和文本渲染 |
-| `tools/*.py` | 工具名称、说明、输入 Schema、路由和调用封装 |
-
-## 本地与公网会话
-
-本地 stdio 流程:
-
-1. `GatewayApp` 从 Redis 或开发用文件 store 读取本机 session。
-2. token 临近过期时通过 base 刷新。
-3. `tools/list` 和 `tools/call` 携带该 token 请求 fmsoperate。
-
-公网流程:
-
-1. Workbuddy 通过 Header、Bearer 或 Cookie 提交 `GWS_xxx`。
-2. `RequestContextParser` 提取设备凭据,`GatewaySessionStore` 从 Redis 获取服务端 `mcp_token`。
-3. `PublicGatewayApp` 为单次请求创建 scoped API 调用上下文,不把 token 放进进程全局状态。
-4. fmsoperate 再校验 token、员工状态、公司、工具权限和业务数据范围。
-
-会话缺失、过期或后端动态列表不可用时关闭访问。公网不得使用文件 token store、固定 `FMS_SESSION_KEY` 或共享进程级员工 token。
-
-## 动态工具可见性
-
-- `GatewayApp` 与 `PublicGatewayApp` 的候选注册顺序和名称必须完全一致,当前各为 12 个。
-- `tools/list` 调用 `/mcp/tools/listEnabledTools`,只返回本地注册与后端启用代码的交集。
-- `tools/call` 重新读取启用集合,避免工具被禁用后继续调用。
-- MCP 返回 `tools.listChanged=false`,因此工具或 Schema 变更必须通过 Gateway 重启和客户端重连刷新。
-
-## Presenter 分类
-
-`query_order` 走旧协议分支,保留 `columns/records/meta` 和既有文本行为。`OutputPresenter.SAFE_TOOLS` 显式处理其余 11 个工具:
-
-| 类型 | 工具 | DTO |
-|---|---|---|
-| 表格 | `query_order_exact`、`query_track`、`query_customs_declaration_files`、`query_outbound_list` | `headers + rows + pagination` |
-| 排舱详情 | `query_outbound_detail` | `summary + details`,分页位于 `details.pagination` |
-| 订单详情 | `query_order_detail` | 中文订单号、详情模块、固定概览分组或明细及分页 |
-| 筛选项 | `list_order_filter_options`、`list_outbound_filter_options`、`list_pending_outbound_export_filter_options` | `value + label + code` 三列 |
-| 导出 | `export_pending_outbound_orders`、`export_out_of_province_port_data` | `files[].label + files[].url` |
-
-Presenter 对列定义、记录类型、详情汇总、附件结构和导出 URL 做白名单校验。未知工具、未允许列、畸形响应或未知异常返回安全 `isError=true`;真实异常仅记录在服务端日志。后端业务错误映射为固定消息,不透传原始 `msg/data`。
-
-订单详情 Presenter 不依赖原始请求参数,而是校验后端规范 `section`、固定分组及每行完整键集合。“全部”响应校验概览和十个明细分组,并将每个明细分组的分页信息独立展示。状态和时间类型只接受已知枚举;入库重量根据后端固定模式显示为单箱重量或总重量。
-
-## 追踪、日志与限流
-
-- stdio 自动生成 `rq_*`;公网入口始终生成可信 `rq_http_*`,不采信调用方 `X-Request-Id`。
-- 调用方 `X-Request-Id` 只记录 SHA-256 短哈希 `client_request_id_hash`,用于关联客户端反馈。
-- 追踪号贯穿动态工具列表、工具调用、MCP `_meta` 和结构化日志。
-- 审计日志只记录必要的 request/jsonrpc ID、协议方法、工具代码、员工/公司标识和脱敏客户端标识;禁止记录凭据与原始响应。
-- `initialize` 与 `tools/list` 不进入限流器;`tools/call` 使用 `gateway_session_id + tool_name` 分桶,并通过 `FMS_MAX_IN_FLIGHT_PER_TOOL` 控制可释放的同时执行数。`FMS_RATE_LIMIT_MAX_REQUESTS=0` 时关闭累计窗口限制,超限统一返回 JSON-RPC `-32029`。
-- 限流配置为 `FMS_RATE_LIMIT_ENABLED`、`FMS_RATE_LIMIT_MAX_REQUESTS`、`FMS_RATE_LIMIT_WINDOW_SECONDS`;生产环境还应在反向代理层限流。
-
-## 关键业务参数
-
-- `query_outbound_list.outbound_status` 必填,只暴露后台业务阶段;`shipping_method` 无默认值,未传表示全部运输方式。
-- `query_order_detail.order_number` 必填;`section` 接受 11 个中文模块或“全部”,页码和每页数量均为 1 至 100。
-- `query_outbound_list.warehouse_id` 接受正整数或 `-1`(客户仓),拒绝 0 与其他负数。
-- `list_outbound_filter_options.filter_type` 使用七个中文枚举;仓库分支由 fmsoperate 复用 `OrderModel::getSortWarehouse()`,所有分支先校验 `admin/Outbound/index`。
-- 排舱筛选响应保留 `value/label/code` 供跨工具传参,面向用户只展示中文标签。
-
-## 验证基线
-
-```powershell
-python -m unittest discover -s tests -p "test_*.py"
-python -m coverage run -m unittest discover -s tests -p "test_*.py"
-python -m coverage report -m --fail-under=100
-git diff --check
-```
-
-运行 coverage 会更新本地 `.coverage` 文件;只做文档治理时使用已有只读报告,不应制造覆盖率产物变更。

+ 0 - 43
project-docs/timeline.md

@@ -1,43 +0,0 @@
-# 时间线
-
-本文件只保留可帮助交接的里程碑与部署状态。测试过程细节由 Git 历史和当前测试代码承载,现行合同以 `README.md` 与其他四份 `project-docs` 为准。
-
-## 2026-07-06 至 2026-07-08:基础 Gateway 与公网安全
-
-- 建立项目文档与本地 stdio Gateway,随后增加公网 HTTP JSON-RPC、Redis 设备会话和健康检查。
-- 公网审计不记录授权码或 token,限流身份不再默认信任调用方可伪造的转发头。
-- 正式员工接入统一使用后台生成的 `GWS_xxx`,旧授权码绑定流程退出运行时。
-
-## 2026-07-13 至 2026-07-15:协议错误与安全展示
-
-- stdio/public 统一工具响应适配:后端业务失败返回 MCP `isError=true`,不再错误标记为成功,也不中断 JSON-RPC 会话。
-- 引入 `OutputPresenter`,将受支持工具的内部结果转换为中文白名单 DTO;`query_order` 保持旧协议兼容。
-- 增加未排舱订单导出、省外进港资料导出及对应权限筛选项。
-
-## 2026-07-16:排舱查询
-
-- 增加 `query_outbound_list` 与 `query_outbound_detail`,同步本地/公网注册、CLI 参数和 Presenter。
-- 排舱列表使用固定中文业务列;详情使用汇总与订单明细两层结构,附件只保留安全文件字段。
-
-## 2026-07-17:工具边界、筛选与限流
-
-- 12 个候选工具统一采用零推断说明:先确认结果意图与号码类型,禁止格式猜测、并行试查和跨字段/跨工具重试。
-- `list_outbound_filter_options` 统一七类排舱筛选项,仓库能力并入该工具,不保留重复入口。
-- 排舱阶段改为必选;运输方式未指定时查询全部运输方式。
-- `initialize` 与 `tools/list` 退出限流,`tools/call` 保持按设备会话与工具分桶。
-- 本地 `GatewayApp` 与公网 `PublicGatewayApp` 对齐为 12 个候选工具;`OutputPresenter` 处理其中 11 个安全工具。
-- 用户确认 `Y:/fmsoperate/sql` 的 10 个 MCP SQL 与 `Y:/base/sql` 的 5 个 MCP SQL 均已执行。
-- 文档治理建立 `AGENTS.md`,放行并重写 `project-docs` 五件套;历史计划与旧测试快照经校验备份后从活动工作区清理。
-
-## 2026-07-17:订单详情工具
-
-- 增加 `query_order_detail` 的中文模块 Schema、本地/公网双注册、CLI 转发和严格白名单 Presenter。
-- 支持 `section=全部` 一次返回概览和十个明细模块,各明细模块独立分页并保持中文白名单展示。
-- 概览固定展示状态节点及四组业务信息;明细覆盖箱单、DW、附件、入库、查验、轨迹、日志和派送,附件保留安全链接。
-- fmsoperate 注册 SQL 仍待目标环境执行;Gateway 重启、Workbuddy 重连和实际动态可见性仍需发布确认。
-
-## 当前交接
-
-- Python 测试继续执行 `.coveragerc` 的语句与分支 100% 门禁;发布前必须重新运行,不能引用历史报告代替。
-- 原有 SQL 执行状态已确认,订单详情注册 SQL尚待执行;工具当前动态启用值仍需通过实际员工会话的 `tools/list` 核验。
-- 运行中 Gateway 是否已加载当前代码、Workbuddy 是否已重新连接,无法由仓库静态判断,发布交接必须显式确认。

+ 0 - 68
project-docs/user-structure.md

@@ -1,68 +0,0 @@
-# 用户流程与项目结构
-
-## 员工公共使用流程
-
-1. 员工从后台或 home 的“连接 Workbuddy”入口创建设备配置。
-2. base 生成 `GWS_xxx`、服务端 MCP token session 和 Redis Gateway session;配置中不暴露 `mcp_token`。
-3. Workbuddy 连接公网 `/mcp`,通过 Header、Bearer 或 Cookie 稳定携带同一 `GWS_xxx`。
-4. `initialize` 完成握手,`tools/list` 返回 Gateway 候选集合与 fmsoperate 动态启用列表的交集。
-5. `tools/call` 解析设备会话、按会话和工具限流,再携带服务端 token 调用 fmsoperate。
-6. ThinkPHP 校验员工、公司、工具和业务数据权限;Gateway 将结果转为安全 MCP 展示。
-7. 设备撤销、token 失效或 Gateway session 过期后,后续调用关闭失败,员工重新创建设备配置。
-
-## 本地调试流程
-
-1. 开发者从 `Y:/mcp` 或员工可访问的 UNC 路径运行 `python app.py serve-stdio`。
-2. 本地 session 从 Redis 读取;文件 token store 只用于开发排障,不放在共享目录。
-3. 使用 `python app.py list-tools` 检查动态可见工具,使用 `python app.py call --tool ...` 查看原始 ThinkPHP 信封。
-4. 使用真实 MCP 客户端验证 `initialize`、`tools/list`、`tools/call` 和 Presenter 展示。
-5. 工具或 Schema 变化后重启 stdio 进程并重新连接客户端;`tools.listChanged=false` 不会主动推送变化。
-
-## AI 工具选择流程
-
-1. 确认结果意图:普通订单、精准订单、订单详情、轨迹、报关资料、排舱列表、排舱详情、筛选项或导出。
-2. 确认号码类型:订单号、排舱单号、物流单号、柜号、提单号或 SO 号。上下文不明确时先询问。
-3. 只调用一个已确认的目标工具;不得按号码格式猜测、并行试查、跨字段试查或失败后自行换工具。
-4. 需要筛选值时先调用对应筛选工具。唯一命中可传递内部 `value`,多条命中让用户选择。
-5. 排舱列表统一使用 `list_outbound_filter_options` 查询七类选项,包括集货仓库;不另设仓库工具。
-6. 完整订单详情先查订单概览,再逐一查询十个明细模块并跟随分页;某模块为空不能跳过其他模块。
-7. 面向用户展示中文标签和完整业务结果,内部值只用于结构化串联。
-
-## 跨仓开发流程
-
-1. 从 `README.md` 与 `project-docs/` 确认产品合同,再读取 `AGENTS.md` 的所有权和红线。
-2. 在 `Y:/mcp` 同步工具 metadata、输入 Schema、stdio/public 双注册、CLI、Presenter 与 Python 测试。
-3. 在 `Y:/fmsoperate` 按 Controller / Logic / Model / Validate 分层同步路由、业务权限、动态注册和 PHP 测试。
-4. 涉及员工身份、设备会话或 token 生命周期时,在 `Y:/base` 同步认证实现和合同测试;home 保持 UI/薄代理边界。
-5. 运行各仓定向与全量验证,核对双注册表和实际 `tools/list`。
-6. 发布后确认 Gateway 已重启、Workbuddy 已重连、动态启用值符合预期,再更新 README 和五件套中的稳定事实。
-
-## 请求与输出路径
-
-```text
-Workbuddy
-  -> /mcp (GWS_xxx)
-  -> PublicMcpHttpHandler
-  -> PublicGatewayApp + GatewaySessionStore
-  -> ScopedApiClient (server-side mcp_token)
-  -> fmsoperate MCP route
-  -> ThinkPHP permission/data checks
-  -> OutputPresenter
-  -> MCP content + structuredContent + _meta.request_id
-```
-
-本地 stdio 使用 `GatewayApp + ApiClient + TokenStore`,从协议处理开始与公网共用相同的工具合同和 Presenter。`query_order` 是旧展示例外。
-
-## 目录地图
-
-| 路径 | 内容 |
-|---|---|
-| `app.py` | 本地入口、CLI 和 `GatewayApp` |
-| `public_gateway.py`、`public_server.py` | 公网 Gateway 与 HTTP JSON-RPC 服务 |
-| `mcp_protocol.py` | MCP 协议处理与结果分流 |
-| `tools/` | 12 个工具的 Schema、说明和路由封装 |
-| `services/` | API、认证、session、request context、scoped client 与 Presenter |
-| `utils/` | 限流和安全散列工具 |
-| `tests/` | Python 单元、协议、展示、会话和安全回归测试 |
-| `project-docs/` | 现行项目五件套 |
-| `docs/superpowers/specs/` | 设计决策背景;不替代现行合同 |

+ 91 - 8
public_gateway.py

@@ -1,4 +1,5 @@
 import logging
+import time
 import uuid
 
 from constants import DEVICE_INVALID_MESSAGE
@@ -57,10 +58,30 @@ class PublicGatewayApp:
     def registered_tool_names(self):
         return tuple(self._tools.keys())
 
-    def _require_session(self, gateway_session_id):
+    def _require_session(self, gateway_session_id, diagnostic_emitter=None):
         session = self.session_store.get(gateway_session_id)
         if not session or not session.get('mcp_token'):
+            if diagnostic_emitter is not None:
+                diagnostic_emitter.emit(
+                    stage='gateway_session',
+                    status='failed',
+                    event_code='GATEWAY_SESSION_NOT_FOUND',
+                    session_credential=gateway_session_id,
+                    context={'transport': 'http'},
+                )
             raise RuntimeError(DEVICE_INVALID_MESSAGE)
+        if diagnostic_emitter is not None:
+            diagnostic_emitter.set_defaults(
+                session_credential=gateway_session_id,
+                admin_id=session.get('admin_id'),
+                company_id=session.get('company_id'),
+                context={'transport': 'http'},
+            )
+            diagnostic_emitter.emit(
+                stage='gateway_session',
+                status='succeeded',
+                event_code='GATEWAY_SESSION_RESOLVED',
+            )
         return session
 
     def _enabled_tool_names(self, response):
@@ -102,20 +123,58 @@ class PublicGatewayApp:
         request_id = str(request_id or '').strip()
         return request_id or 'rq_{0}'.format(uuid.uuid4().hex[:16])
 
-    def call_tool(self, gateway_session_id, name, arguments=None, request_id='', client_ip=''):
+    def call_tool(
+        self,
+        gateway_session_id,
+        name,
+        arguments=None,
+        request_id='',
+        client_ip='',
+        diagnostic_emitter=None,
+    ):
+        session = self._require_session(gateway_session_id, diagnostic_emitter)
+        request_id = self.build_request_id(request_id)
         if name not in self._tools:
+            if diagnostic_emitter is not None:
+                diagnostic_emitter.emit(
+                    stage='backend_call',
+                    status='failed',
+                    event_code='TOOL_NOT_REGISTERED',
+                    context={'transport': 'http'},
+                )
             raise KeyError('tool not registered: {0}'.format(name))
 
-        session = self._require_session(gateway_session_id)
-        request_id = self.build_request_id(request_id)
-        if name not in self._load_enabled_tool_names(
-            session['mcp_token'],
-            request_id,
-        ):
+        try:
+            enabled_tools = self._load_enabled_tool_names(
+                session['mcp_token'],
+                request_id,
+            )
+        except Exception:
+            if diagnostic_emitter is not None:
+                diagnostic_emitter.emit(
+                    stage='backend_call',
+                    status='failed',
+                    event_code='ENABLED_TOOL_LOOKUP_FAILED',
+                    tool_code=name,
+                    context={'transport': 'http'},
+                )
+            raise
+        if name not in enabled_tools:
+            if diagnostic_emitter is not None:
+                diagnostic_emitter.emit(
+                    stage='backend_call',
+                    status='failed',
+                    event_code='TOOL_DISABLED',
+                    tool_code=name,
+                    context={'transport': 'http'},
+                )
             raise RuntimeError('tool disabled: {0}'.format(name))
 
         tool = self._tools[name]
 
+        if diagnostic_emitter is not None:
+            diagnostic_emitter.set_defaults(tool_code=name)
+
         session_hash = hash_gateway_session_id(gateway_session_id)[:12]
         admin_id = session.get('admin_id')
         company_id = session.get('company_id')
@@ -131,6 +190,13 @@ class PublicGatewayApp:
         )
 
         try:
+            started_at = time.monotonic()
+            if diagnostic_emitter is not None:
+                diagnostic_emitter.emit(
+                    stage='backend_call',
+                    status='started',
+                    event_code='BACKEND_CALL_STARTED',
+                )
             result = self.api_client.call_tool(
                 token=session['mcp_token'],
                 tool_code=tool.name,
@@ -150,6 +216,16 @@ class PublicGatewayApp:
                     'response_code': result.get('code'),
                 },
             )
+            if diagnostic_emitter is not None:
+                diagnostic_emitter.emit(
+                    stage='backend_call',
+                    status='succeeded',
+                    event_code='BACKEND_CALL_COMPLETED',
+                    response_code=(
+                        result.get('code') if isinstance(result, dict) else None
+                    ),
+                    cost_ms=max(0, int((time.monotonic() - started_at) * 1000)),
+                )
             return result
         except Exception as e:
             logger.error(
@@ -163,4 +239,11 @@ class PublicGatewayApp:
                     'exception_class': e.__class__.__name__,
                 },
             )
+            if diagnostic_emitter is not None:
+                diagnostic_emitter.emit(
+                    stage='backend_call',
+                    status='failed',
+                    event_code='UNEXPECTED_EXCEPTION',
+                    response_code='MCP_9001',
+                )
             raise

+ 236 - 18
public_server.py

@@ -8,6 +8,8 @@ from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
 
 from constants import DEVICE_INVALID_MESSAGE
 from mcp_protocol import McpProtocolHandler
+from services.diagnostic_event import RequestDiagnosticEmitter
+from services.diagnostic_reporter import NullDiagnosticReporter
 from services.request_context import RequestContextParser
 from utils.rate_limiter import SimpleRateLimiter
 
@@ -21,10 +23,17 @@ def extract_client_ip(headers, client_address):
 
 
 class PublicMcpHttpHandler:
-    def __init__(self, gateway_app, context_parser=None, rate_limiter=None):
+    def __init__(
+        self,
+        gateway_app,
+        context_parser=None,
+        rate_limiter=None,
+        reporter=None,
+    ):
         self.gateway_app = gateway_app
         self.context_parser = context_parser or RequestContextParser()
         self.rate_limiter = rate_limiter
+        self.reporter = reporter or NullDiagnosticReporter()
         # Cache registered tool names at startup for rate-key validation;
         # unknown names fall back to the bare IP bucket, preventing bucket explosion
         self._known_tools = frozenset(gateway_app.registered_tool_names())
@@ -52,7 +61,8 @@ class PublicMcpHttpHandler:
         jsonrpc_id=None,
         tool_name='',
         *,
-        log_identity=None
+        log_identity=None,
+        diagnostic_emitter=None,
     ):
         """Returns an error response if rate limit exceeded, else None.
 
@@ -60,6 +70,17 @@ class PublicMcpHttpHandler:
         log_identity — optional string shown in warning logs (e.g. client_ip)
         """
         if self.rate_limiter and rate_key and not self.rate_limiter.is_allowed(rate_key):
+            if diagnostic_emitter is not None:
+                diagnostic_emitter.emit(
+                    stage='rate_limit',
+                    status='failed',
+                    event_code='RATE_LIMIT_EXCEEDED',
+                    tool_code=tool_name or None,
+                    context={
+                        'limit_type': 'request_window',
+                        'transport': 'http',
+                    },
+                )
             logger.warning(
                 "MCP public rate limit exceeded",
                 extra={
@@ -86,6 +107,7 @@ class PublicMcpHttpHandler:
         message,
         client_ip='',
         trace_request_id='',
+        diagnostic_emitter=None,
     ):
         request_id = message.get('id') if isinstance(message, dict) else None
         trace_request_id = str(trace_request_id or '').strip() \
@@ -94,9 +116,31 @@ class PublicMcpHttpHandler:
         message_params = (message or {}).get('params') or {}
         tool_name = str(message_params.get('name') or '').strip() \
             if isinstance(message_params, dict) else ''
+        emitter = diagnostic_emitter or RequestDiagnosticEmitter(
+            self.reporter,
+            trace_request_id,
+            defer_until_identity=True,
+        )
+        emitter.emit(
+            stage='request_ingress',
+            status='started',
+            event_code='REQUEST_RECEIVED',
+            context={'jsonrpc_method': method or 'unknown', 'transport': 'http'},
+        )
+        protocol_validated = False
 
         try:
             if method == 'initialize':
+                emitter.emit(
+                    stage='protocol_validation',
+                    status='succeeded',
+                    event_code='PROTOCOL_VALIDATION_COMPLETED',
+                    context={
+                        'jsonrpc_method': method,
+                        'transport': 'http',
+                    },
+                )
+                protocol_validated = True
                 # initialize is a stateless handshake that only returns server metadata;
                 # rate-limiting it would block clients from connecting at all, so we skip it.
                 return McpProtocolHandler._success_response(request_id, {
@@ -108,8 +152,24 @@ class PublicMcpHttpHandler:
                     },
                 })
             if method == 'tools/list':
+                emitter.emit(
+                    stage='protocol_validation',
+                    status='succeeded',
+                    event_code='PROTOCOL_VALIDATION_COMPLETED',
+                    context={
+                        'jsonrpc_method': method,
+                        'transport': 'http',
+                    },
+                )
+                protocol_validated = True
                 context = self.context_parser.parse(headers or {})
                 if not context.has_session():
+                    emitter.emit(
+                        stage='gateway_session',
+                        status='failed',
+                        event_code='GATEWAY_SESSION_NOT_FOUND',
+                        context={'transport': 'http'},
+                    )
                     logger.warning(
                         'MCP public device session unavailable',
                         extra={
@@ -140,9 +200,43 @@ class PublicMcpHttpHandler:
             if method == 'tools/call':
                 if not isinstance(message_params, dict):
                     raise ValueError('tool parameters must be an object')
+                if not tool_name:
+                    emitter.emit(
+                        stage='protocol_validation',
+                        status='failed',
+                        event_code='PARAM_VALIDATION_FAILED',
+                        context={
+                            'jsonrpc_method': method,
+                            'jsonrpc_code': -32602,
+                            'transport': 'http',
+                        },
+                    )
+                    return McpProtocolHandler._error_response(
+                        request_id,
+                        -32602,
+                        'Invalid params',
+                        trace_request_id,
+                    )
+                emitter.emit(
+                    stage='protocol_validation',
+                    status='succeeded',
+                    event_code='PROTOCOL_VALIDATION_COMPLETED',
+                    tool_code=tool_name,
+                    context={
+                        'jsonrpc_method': method,
+                        'transport': 'http',
+                    },
+                )
+                protocol_validated = True
                 # Parse context first — tools/call always requires a valid session
                 context = self.context_parser.parse(headers or {})
                 if not context.has_session():
+                    emitter.emit(
+                        stage='gateway_session',
+                        status='failed',
+                        event_code='GATEWAY_SESSION_NOT_FOUND',
+                        context={'transport': 'http'},
+                    )
                     logger.warning(
                         'MCP public device session unavailable',
                         extra={
@@ -171,6 +265,7 @@ class PublicMcpHttpHandler:
                     request_id,
                     tool_name,
                     log_identity=client_ip,
+                    diagnostic_emitter=emitter,
                 )
                 if blocked:
                     return blocked
@@ -178,6 +273,16 @@ class PublicMcpHttpHandler:
                 acquired = self.rate_limiter is None \
                     or self.rate_limiter.try_acquire(rate_key)
                 if not acquired:
+                    emitter.emit(
+                        stage='rate_limit',
+                        status='failed',
+                        event_code='CONCURRENCY_LIMIT_EXCEEDED',
+                        tool_code=tool_name or None,
+                        context={
+                            'limit_type': 'concurrency',
+                            'transport': 'http',
+                        },
+                    )
                     logger.warning(
                         'MCP public concurrency limit exceeded',
                         extra={
@@ -196,6 +301,16 @@ class PublicMcpHttpHandler:
                         'Too many requests in progress. Please try again later.',
                         trace_request_id,
                     )
+                emitter.emit(
+                    stage='rate_limit',
+                    status='succeeded',
+                    event_code='RATE_LIMIT_ALLOWED',
+                    tool_code=tool_name or None,
+                    context={
+                        'limit_type': 'tools_call',
+                        'transport': 'http',
+                    },
+                )
                 try:
                     result = self.gateway_app.call_tool(
                         session_id,
@@ -203,15 +318,42 @@ class PublicMcpHttpHandler:
                         params.get('arguments') or {},
                         request_id=trace_request_id,
                         client_ip=client_ip,
+                        diagnostic_emitter=emitter,
                     )
                 finally:
                     if self.rate_limiter is not None:
                         self.rate_limiter.release(rate_key)
-                return McpProtocolHandler._tool_call_response(
-                    request_id,
-                    tool_name,
-                    result,
+                try:
+                    response = McpProtocolHandler._tool_call_response(
+                        request_id,
+                        tool_name,
+                        result,
+                    )
+                except Exception:
+                    emitter.emit(
+                        stage='response_safety',
+                        status='failed',
+                        event_code='RESPONSE_SAFETY_REJECTED',
+                        context={'transport': 'http'},
+                    )
+                    raise
+                emitter.emit(
+                    stage='response_safety',
+                    status='succeeded',
+                    event_code='RESPONSE_SAFETY_COMPLETED',
+                    context={'transport': 'http'},
                 )
+                return response
+            emitter.emit(
+                stage='protocol_validation',
+                status='failed',
+                event_code='METHOD_NOT_FOUND',
+                context={
+                    'jsonrpc_method': method or 'unknown',
+                    'jsonrpc_code': -32601,
+                    'transport': 'http',
+                },
+            )
             return McpProtocolHandler._error_response(
                 request_id,
                 -32601,
@@ -219,6 +361,18 @@ class PublicMcpHttpHandler:
                 trace_request_id,
             )
         except Exception as exc:
+            if not protocol_validated:
+                emitter.emit(
+                    stage='protocol_validation',
+                    status='failed',
+                    event_code='PARAM_VALIDATION_FAILED',
+                    tool_code=tool_name or None,
+                    context={
+                        'jsonrpc_method': method or 'unknown',
+                        'jsonrpc_code': -32602,
+                        'transport': 'http',
+                    },
+                )
             if method == 'tools/call':
                 diagnostic_reason = (
                     'PARAM_VALIDATION_FAILED'
@@ -261,10 +415,16 @@ class PublicMcpHttpHandler:
                 'Gateway request failed. Please try again later.',
                 trace_request_id,
             )
+        finally:
+            emitter.flush()
 
 
-def create_http_handler(gateway_app, rate_limiter=None):
-    rpc_handler = PublicMcpHttpHandler(gateway_app, rate_limiter=rate_limiter)
+def create_http_handler(gateway_app, rate_limiter=None, reporter=None):
+    rpc_handler = PublicMcpHttpHandler(
+        gateway_app,
+        rate_limiter=rate_limiter,
+        reporter=reporter,
+    )
 
     class Handler(BaseHTTPRequestHandler):
         def log_message(self, format, *args):
@@ -282,6 +442,11 @@ def create_http_handler(gateway_app, rate_limiter=None):
             request_headers = dict(self.headers.items())
             client_ip = extract_client_ip(request_headers, self.client_address)
             trace_request_id = rpc_handler._build_trace_request_id(request_headers)
+            diagnostic_emitter = RequestDiagnosticEmitter(
+                rpc_handler.reporter,
+                trace_request_id,
+                defer_until_identity=True,
+            )
             client_request_id_hash = rpc_handler._client_request_id(
                 request_headers
             )
@@ -297,6 +462,24 @@ def create_http_handler(gateway_app, rate_limiter=None):
             try:
                 message = json.loads(body)
             except json.JSONDecodeError as exc:
+                diagnostic_emitter.emit(
+                    stage='request_ingress',
+                    status='started',
+                    event_code='REQUEST_RECEIVED',
+                    context={
+                        'jsonrpc_method': 'unknown',
+                        'transport': 'http',
+                    },
+                )
+                diagnostic_emitter.emit(
+                    stage='protocol_validation',
+                    status='failed',
+                    event_code='PARAM_VALIDATION_FAILED',
+                    context={
+                        'jsonrpc_code': -32700,
+                        'transport': 'http',
+                    },
+                )
                 logger.warning(
                     'MCP public invalid JSON',
                     extra={
@@ -310,12 +493,15 @@ def create_http_handler(gateway_app, rate_limiter=None):
                         'client_identity': client_ip,
                     },
                 )
-                self._write_json(McpProtocolHandler._error_response(
-                    None,
-                    -32700,
-                    'Parse error',
-                    trace_request_id,
-                ))
+                self._write_json(
+                    McpProtocolHandler._error_response(
+                        None,
+                        -32700,
+                        'Parse error',
+                        trace_request_id,
+                    ),
+                    diagnostic_emitter=diagnostic_emitter,
+                )
                 return
             method = message.get('method', '')
             request_id = message.get('id') if message.get('id') is not None else ''
@@ -334,10 +520,14 @@ def create_http_handler(gateway_app, rate_limiter=None):
                 message,
                 client_ip,
                 trace_request_id=trace_request_id,
+                diagnostic_emitter=diagnostic_emitter,
+            )
+            self._write_json(
+                response,
+                diagnostic_emitter=diagnostic_emitter,
             )
-            self._write_json(response)
 
-        def _write_json(self, payload):
+        def _write_json(self, payload, diagnostic_emitter=None):
             raw = json.dumps(payload, ensure_ascii=False).encode('utf-8')
             try:
                 self.send_response(200)
@@ -345,16 +535,41 @@ def create_http_handler(gateway_app, rate_limiter=None):
                 self.send_header('Content-Length', str(len(raw)))
                 self.end_headers()
                 self.wfile.write(raw)
+                if diagnostic_emitter is not None:
+                    diagnostic_emitter.emit(
+                        stage='response_write',
+                        status='succeeded',
+                        event_code='RESPONSE_WRITE_COMPLETED',
+                        context={
+                            'http_status': 200,
+                            'client_disconnected': False,
+                            'transport': 'http',
+                        },
+                    )
             except (BrokenPipeError, ConnectionResetError):
                 self.close_connection = True
+                if diagnostic_emitter is not None:
+                    diagnostic_emitter.emit(
+                        stage='response_write',
+                        status='failed',
+                        event_code='CLIENT_DISCONNECTED',
+                        context={
+                            'http_status': 200,
+                            'client_disconnected': True,
+                            'transport': 'http',
+                        },
+                    )
                 logger.info('MCP client disconnected before response was written')
+            finally:
+                if diagnostic_emitter is not None:
+                    diagnostic_emitter.flush()
 
     return Handler
 
 
 def serve_public(gateway_app, host='0.0.0.0', port=8765, enable_rate_limit=True,
                  rate_limit_max_requests=60, rate_limit_window_seconds=60,
-                 max_in_flight_per_tool=2):
+                 max_in_flight_per_tool=2, reporter=None):
     # Configure logging
     logging.basicConfig(
         level=logging.INFO,
@@ -373,7 +588,10 @@ def serve_public(gateway_app, host='0.0.0.0', port=8765, enable_rate_limit=True,
         logger.info(f"Rate limiting enabled: {rate_limit_max_requests} requests/{rate_limit_window_seconds}s and {max_in_flight_per_tool} in-flight per session and tool (tools/call only)")
 
     logger.info(f"Starting public MCP Gateway on {host}:{port}")
-    server = ThreadingHTTPServer((host, int(port)), create_http_handler(gateway_app, rate_limiter))
+    server = ThreadingHTTPServer(
+        (host, int(port)),
+        create_http_handler(gateway_app, rate_limiter, reporter=reporter),
+    )
 
     # Schedule periodic cleanup to prevent unbounded memory growth in the rate limiter
     if rate_limiter is not None:

+ 206 - 0
services/diagnostic_event.py

@@ -0,0 +1,206 @@
+import hashlib
+import re
+import uuid
+from datetime import datetime, timezone
+
+
+_STAGE_CONTEXT_FIELDS = {
+    'request_ingress': frozenset(('jsonrpc_method', 'http_status', 'transport')),
+    'protocol_validation': frozenset(('jsonrpc_method', 'jsonrpc_code', 'transport')),
+    'gateway_session': frozenset(('transport',)),
+    'rate_limit': frozenset(('limit_type', 'transport')),
+    'backend_call': frozenset(('http_status', 'transport')),
+    'response_safety': frozenset(('transport',)),
+    'response_write': frozenset(('http_status', 'client_disconnected', 'transport')),
+}
+_STATUSES = frozenset(('started', 'succeeded', 'failed', 'skipped'))
+_CODE_PATTERN = re.compile(r'^[A-Z][A-Z0-9_]*$')
+
+
+def _utc_timestamp():
+    return datetime.now(timezone.utc).isoformat(timespec='milliseconds').replace(
+        '+00:00', 'Z'
+    )
+
+
+def _optional_positive_integer(event, name, value):
+    if value is None:
+        return
+    if isinstance(value, bool) or not isinstance(value, int) or value <= 0:
+        raise ValueError('{0} must be a positive integer'.format(name))
+    event[name] = value
+
+
+def _optional_code(event, name, value, max_length):
+    if value is None or value == '':
+        return
+    if (
+        not isinstance(value, str)
+        or len(value) > max_length
+        or _CODE_PATTERN.fullmatch(value) is None
+    ):
+        raise ValueError('{0} is invalid'.format(name))
+    event[name] = value
+
+
+def _validated_context(stage, context):
+    context = {} if context is None else context
+    if not isinstance(context, dict):
+        raise ValueError('context must be an object')
+    unknown = set(context) - _STAGE_CONTEXT_FIELDS[stage]
+    if unknown:
+        raise ValueError('context contains unsupported fields')
+
+    result = {}
+    for key, value in context.items():
+        if key == 'http_status':
+            if isinstance(value, bool) or not isinstance(value, int) or not 100 <= value <= 599:
+                raise ValueError('http_status is invalid')
+        elif key == 'jsonrpc_code':
+            if isinstance(value, bool) or not isinstance(value, int):
+                raise ValueError('jsonrpc_code is invalid')
+        elif key == 'client_disconnected':
+            if not isinstance(value, bool):
+                raise ValueError('client_disconnected is invalid')
+        elif not isinstance(value, str) or not value or len(value) > 64:
+            raise ValueError('{0} is invalid'.format(key))
+        result[key] = value
+    return result
+
+
+def build_diagnostic_event(
+    request_id,
+    stage,
+    status,
+    event_code,
+    occurred_at=None,
+    session_credential='',
+    company_id=None,
+    admin_id=None,
+    tool_code=None,
+    response_code=None,
+    cost_ms=None,
+    summary_code=None,
+    context=None,
+):
+    if not isinstance(request_id, str) or not re.fullmatch(
+        r'rq_[A-Za-z0-9._:-]{1,97}', request_id
+    ):
+        raise ValueError('request_id is invalid')
+    if stage not in _STAGE_CONTEXT_FIELDS:
+        raise ValueError('stage is invalid')
+    if status not in _STATUSES:
+        raise ValueError('status is invalid')
+    if (
+        not isinstance(event_code, str)
+        or len(event_code) > 64
+        or _CODE_PATTERN.fullmatch(event_code) is None
+    ):
+        raise ValueError('event_code is invalid')
+
+    event = {
+        'schema_version': 1,
+        'event_id': 'evt_gateway_{0}'.format(uuid.uuid4().hex),
+        'request_id': request_id,
+        'source': 'gateway',
+        'stage': stage,
+        'status': status,
+        'event_code': event_code,
+        'occurred_at': occurred_at or _utc_timestamp(),
+    }
+    _optional_positive_integer(event, 'company_id', company_id)
+    _optional_positive_integer(event, 'admin_id', admin_id)
+    if tool_code is not None and tool_code != '':
+        if not isinstance(tool_code, str) or re.fullmatch(r'[a-z][a-z0-9_]{0,63}', tool_code) is None:
+            raise ValueError('tool_code is invalid')
+        event['tool_code'] = tool_code
+    if session_credential:
+        if not isinstance(session_credential, str):
+            raise ValueError('session_credential must be text')
+        event['session_hash'] = hashlib.sha256(
+            session_credential.encode('utf-8')
+        ).hexdigest()[:12]
+    _optional_code(event, 'response_code', response_code, 32)
+    _optional_code(event, 'summary_code', summary_code, 64)
+    if cost_ms is not None:
+        if (
+            isinstance(cost_ms, bool)
+            or not isinstance(cost_ms, int)
+            or not 0 <= cost_ms <= 3600000
+        ):
+            raise ValueError('cost_ms is invalid')
+        event['cost_ms'] = cost_ms
+    validated_context = _validated_context(stage, context)
+    if validated_context:
+        event['context'] = validated_context
+    return event
+
+
+class RequestDiagnosticEmitter:
+    MAX_EVENTS = 20
+
+    def __init__(
+        self,
+        reporter,
+        request_id,
+        defer_until_identity=False,
+        **event_defaults
+    ):
+        self.reporter = reporter
+        self.request_id = request_id
+        self.event_defaults = event_defaults
+        self.defer_until_identity = bool(defer_until_identity)
+        self._deferred = []
+        self.emitted_count = 0
+        self.dropped_count = 0
+
+    def set_defaults(self, **event_defaults):
+        self.event_defaults.update(event_defaults)
+        company_id = self.event_defaults.get('company_id')
+        if (
+            self.defer_until_identity
+            and isinstance(company_id, int)
+            and not isinstance(company_id, bool)
+            and company_id > 0
+        ):
+            self.flush()
+
+    def emit(self, **fields):
+        if self.emitted_count >= self.MAX_EVENTS:
+            self.dropped_count += 1
+            return False
+        event_fields = dict(self.event_defaults)
+        event_fields.update(fields)
+        event_fields['request_id'] = self.request_id
+        self.emitted_count += 1
+        if (
+            self.defer_until_identity
+            and not event_fields.get('company_id')
+        ):
+            self._deferred.append(event_fields)
+            return True
+        return self._deliver(event_fields)
+
+    def flush(self):
+        deferred = self._deferred
+        self._deferred = []
+        accepted = True
+        for fields in deferred:
+            enriched = dict(fields)
+            for name in ('company_id', 'admin_id', 'session_credential'):
+                if name in self.event_defaults:
+                    enriched[name] = self.event_defaults[name]
+            if not self._deliver(enriched):
+                accepted = False
+        return accepted
+
+    def _deliver(self, event_fields):
+        try:
+            event = build_diagnostic_event(**event_fields)
+            accepted = self.reporter.report(event)
+        except Exception:
+            self.dropped_count += 1
+            return False
+        if not accepted:
+            self.dropped_count += 1
+        return accepted

+ 217 - 0
services/diagnostic_reporter.py

@@ -0,0 +1,217 @@
+import hashlib
+import hmac
+import json
+import logging
+import queue
+import threading
+import time
+import urllib.request
+from urllib.parse import urlsplit
+import uuid
+
+
+logger = logging.getLogger(__name__)
+
+
+def diagnostic_reporter_from_config(config, **overrides):
+    if getattr(config, 'diagnosis_enabled', False) is not True:
+        return NullDiagnosticReporter()
+    url = str(getattr(config, 'diagnosis_url', '') or '').strip()
+    key_id = str(getattr(config, 'diagnosis_key_id', '') or '').strip()
+    secret = str(getattr(config, 'diagnosis_secret', '') or '')
+    allow_insecure_http = getattr(
+        config, 'diagnosis_allow_insecure_http', False
+    ) is True
+    parsed_url = urlsplit(url)
+    if (
+        (
+            parsed_url.scheme != 'https'
+            and not (parsed_url.scheme == 'http' and allow_insecure_http)
+        )
+        or not parsed_url.netloc
+        or not key_id
+        or len(secret) < 32
+    ):
+        logger.warning('MCP diagnostic reporter configuration is invalid')
+        return NullDiagnosticReporter()
+    options = {
+        'url': url,
+        'key_id': key_id,
+        'secret': secret,
+        'queue_size': getattr(config, 'diagnosis_queue_size', 1000),
+        'batch_size': getattr(config, 'diagnosis_batch_size', 100),
+        'timeout_seconds': getattr(config, 'diagnosis_timeout_seconds', 0.5),
+        'initial_backoff': getattr(
+            config, 'diagnosis_initial_backoff_seconds', 0.25
+        ),
+        'max_backoff': getattr(
+            config, 'diagnosis_max_backoff_seconds', 5.0
+        ),
+    }
+    options.update(overrides)
+    reporter = DiagnosticReporter(**options)
+    if reporter.start():
+        return reporter
+    reporter.close()
+    return NullDiagnosticReporter()
+
+
+def _http_transport(url, body, headers, timeout):
+    request = urllib.request.Request(url, data=body, headers=headers, method='POST')
+    with urllib.request.urlopen(request, timeout=timeout) as response:
+        return response.getcode(), response.read()
+
+
+class NullDiagnosticReporter:
+    drop_count = 0
+    pending_count = 0
+
+    def start(self):
+        return True
+
+    def report(self, _event):
+        return True
+
+    def close(self):
+        return None
+
+
+class DiagnosticReporter:
+    def __init__(
+        self,
+        url,
+        key_id,
+        secret,
+        queue_size=1000,
+        batch_size=100,
+        timeout_seconds=0.5,
+        initial_backoff=0.25,
+        max_backoff=5.0,
+        transport=None,
+        clock=None,
+        nonce_factory=None,
+        sleeper=None,
+        thread_factory=None,
+    ):
+        self.url = str(url)
+        self.key_id = str(key_id)
+        self.secret = str(secret)
+        self.batch_size = max(1, min(100, int(batch_size)))
+        self.timeout_seconds = float(timeout_seconds)
+        self.initial_backoff = max(0.0, float(initial_backoff))
+        self.max_backoff = max(self.initial_backoff, float(max_backoff))
+        self.transport = transport or _http_transport
+        self.clock = clock or time.time
+        self.nonce_factory = nonce_factory or (lambda: uuid.uuid4().hex)
+        self.sleeper = sleeper or time.sleep
+        self.thread_factory = thread_factory or threading.Thread
+        self._queue = queue.Queue(maxsize=max(1, int(queue_size)))
+        self._pending = []
+        self._backoff = self.initial_backoff
+        self._stop = threading.Event()
+        self._thread = None
+        self.drop_count = 0
+
+    @property
+    def pending_count(self):
+        return len(self._pending)
+
+    def report(self, event):
+        try:
+            self._queue.put_nowait(event)
+            return True
+        except queue.Full:
+            self.drop_count += 1
+            return False
+
+    def start(self):
+        if self._thread is not None:
+            return True
+        try:
+            thread = self.thread_factory(target=self._run, daemon=True)
+            thread.start()
+            self._thread = thread
+            return True
+        except Exception:
+            logger.warning('MCP diagnostic reporter thread unavailable')
+            return False
+
+    def close(self):
+        self._stop.set()
+        thread = self._thread
+        if thread is not None and thread.is_alive():
+            thread.join(timeout=self.timeout_seconds)
+
+    def process_once(self):
+        if not self._pending:
+            self._pending = self._take_batch()
+        if not self._pending:
+            return True
+
+        body = json.dumps(
+            {'events': self._pending},
+            ensure_ascii=False,
+            separators=(',', ':'),
+        ).encode('utf-8')
+        headers = self._signed_headers(body)
+        try:
+            status, response_body = self.transport(
+                self.url,
+                body,
+                headers,
+                self.timeout_seconds,
+            )
+            payload = json.loads(response_body.decode('utf-8'))
+            if status < 200 or status >= 300 or payload.get('code') != 'MCP_DIAG_INGEST_0000':
+                raise ValueError('collector rejected batch')
+        except Exception:
+            delay = self._backoff
+            self._backoff = min(
+                self.max_backoff,
+                max(self.initial_backoff, self._backoff * 2),
+            )
+            self.sleeper(delay)
+            return False
+
+        self._pending = []
+        self._backoff = self.initial_backoff
+        return True
+
+    def _take_batch(self):
+        events = []
+        while len(events) < self.batch_size:
+            try:
+                events.append(self._queue.get_nowait())
+            except queue.Empty:
+                break
+        return events
+
+    def _signed_headers(self, body):
+        timestamp = str(int(self.clock()))
+        nonce = str(self.nonce_factory())
+        path = urlsplit(self.url).path or '/'
+        canonical = 'POST\n{0}\n{1}\n{2}\n{3}'.format(
+            path,
+            timestamp,
+            nonce,
+            hashlib.sha256(body).hexdigest(),
+        )
+        signature = hmac.new(
+            self.secret.encode('utf-8'),
+            canonical.encode('utf-8'),
+            hashlib.sha256,
+        ).hexdigest()
+        return {
+            'Content-Type': 'application/json',
+            'X-MCP-Source': 'gateway',
+            'X-MCP-Timestamp': timestamp,
+            'X-MCP-Nonce': nonce,
+            'X-MCP-Key-Id': self.key_id,
+            'X-MCP-Signature': signature,
+        }
+
+    def _run(self):
+        while not self._stop.is_set():
+            if not self.process_once():
+                continue
+            self._stop.wait(0.1)

+ 41 - 0
tests/test_app_coverage.py

@@ -103,6 +103,12 @@ class DummyAuthClient:
 
 
 class GatewayAppBoundaryTest(unittest.TestCase):
+    def test_protocol_handler_uses_gateway_reporter(self):
+        reporter = MagicMock()
+        app = GatewayApp(reporter=reporter)
+
+        self.assertIs(reporter, app.create_protocol_handler().reporter)
+
     def test_registered_tool_names_returns_local_registry(self):
         app = GatewayApp(api_client=DummyApiClient())
 
@@ -286,6 +292,7 @@ class RunCliBranchCoverageTest(unittest.TestCase):
                 patch('app.GatewaySessionStore') as store_cls, \
                 patch('app.ScopedApiClient') as api_cls, \
                 patch('app.PublicGatewayApp') as public_app_cls, \
+                patch('app.diagnostic_reporter_from_config') as reporter_factory, \
                 patch('app.serve_public', return_value=17) as serve:
             result = GatewayApp().run_cli([
                 'serve-public', '--host', '127.0.0.1', '--port', '9000',
@@ -317,7 +324,41 @@ class RunCliBranchCoverageTest(unittest.TestCase):
             rate_limit_max_requests=12,
             rate_limit_window_seconds=34,
             max_in_flight_per_tool=2,
+            reporter=reporter_factory.return_value,
+        )
+        reporter_factory.return_value.close.assert_called_once_with()
+
+    def test_serve_public_reuses_existing_non_null_reporter(self):
+        config = SimpleNamespace(
+            redis_host='redis.test',
+            redis_port=6379,
+            redis_db=0,
+            redis_password='',
+            timeout_seconds=1,
+            redis_prefix='gateway:',
+            gateway_session_ttl_seconds=600,
+            tools_base_url='https://tools.test',
+            rate_limit_enabled=False,
+            rate_limit_max_requests=0,
+            rate_limit_window_seconds=60,
+            max_in_flight_per_tool=1,
         )
+        reporter = MagicMock()
+        with patch('app.GatewayConfig.from_env', return_value=config), \
+                patch('app.RedisSocketClient'), \
+                patch('app.GatewaySessionStore'), \
+                patch('app.ScopedApiClient'), \
+                patch('app.PublicGatewayApp'), \
+                patch('app.diagnostic_reporter_from_config') as factory, \
+                patch('app.serve_public', return_value=0) as serve:
+            result = GatewayApp(reporter=reporter).run_cli([
+                'serve-public', '--host', '127.0.0.1', '--port', '9000',
+            ])
+
+        self.assertEqual(0, result)
+        factory.assert_not_called()
+        self.assertIs(reporter, serve.call_args.kwargs['reporter'])
+        reporter.close.assert_called_once_with()
 
     def test_unsupported_command_raises_runtime_error(self):
         with patch(

+ 42 - 0
tests/test_config_compat.py

@@ -92,6 +92,40 @@ class GatewayConfigCompatTest(unittest.TestCase):
         self.assertEqual('public', config.gateway_mode)
         self.assertEqual('fms:mcp:gateway:', config.redis_prefix)
 
+    def test_diagnostic_reporter_config_reads_all_operational_limits(self):
+        config = GatewayConfig.from_env(env={
+            'FMS_API_BASE': 'https://base.example.com',
+            'MCP_DIAGNOSIS_ENABLED': 'true',
+            'MCP_DIAGNOSIS_URL': 'https://support.internal/internal/mcp-diagnostics/events',
+            'MCP_DIAGNOSIS_KEY_ID': 'gateway-current',
+            'MCP_DIAGNOSIS_SECRET': 's' * 32,
+            'MCP_DIAGNOSIS_QUEUE_SIZE': '500',
+            'MCP_DIAGNOSIS_BATCH_SIZE': '50',
+            'MCP_DIAGNOSIS_TIMEOUT_SECONDS': '0.4',
+            'MCP_DIAGNOSIS_INITIAL_BACKOFF_SECONDS': '0.2',
+            'MCP_DIAGNOSIS_MAX_BACKOFF_SECONDS': '3.0',
+            'MCP_DIAGNOSIS_ALLOW_INSECURE_HTTP': 'true',
+        }, dotenv_path='missing.env')
+
+        self.assertTrue(config.diagnosis_enabled)
+        self.assertEqual(
+            'https://support.internal/internal/mcp-diagnostics/events',
+            config.diagnosis_url,
+        )
+        self.assertEqual('gateway-current', config.diagnosis_key_id)
+        self.assertEqual('s' * 32, config.diagnosis_secret)
+        self.assertEqual(500, config.diagnosis_queue_size)
+        self.assertEqual(50, config.diagnosis_batch_size)
+        self.assertEqual(0.4, config.diagnosis_timeout_seconds)
+        self.assertEqual(0.2, config.diagnosis_initial_backoff_seconds)
+        self.assertEqual(3.0, config.diagnosis_max_backoff_seconds)
+        self.assertTrue(config.diagnosis_allow_insecure_http)
+
+        default_config = GatewayConfig.from_env(env={
+            'FMS_API_BASE': 'https://base.example.com',
+        }, dotenv_path='missing.env')
+        self.assertFalse(default_config.diagnosis_allow_insecure_http)
+
     def test_timeout_ms_conversion_in_preferred_env(self):
         """Test FMS_TIMEOUT_MS conversion in preferred env"""
         config = GatewayConfig.from_env(env={
@@ -110,6 +144,14 @@ class GatewayConfigCompatTest(unittest.TestCase):
 
         self.assertEqual(1, config.timeout_seconds)
 
+    def test_invalid_diagnostic_float_uses_default(self):
+        config = GatewayConfig.from_env(env={
+            'FMS_API_BASE': 'https://base.example.com',
+            'MCP_DIAGNOSIS_TIMEOUT_SECONDS': 'not-a-number',
+        }, dotenv_path='missing.env')
+
+        self.assertEqual(0.5, config.diagnosis_timeout_seconds)
+
     def test_timeout_resolution_covers_supported_env_keys(self):
         cases = (
             ({'MCP_TIMEOUT_SECONDS': '11'}, {}, 11),

+ 509 - 0
tests/test_diagnostic_reporter.py

@@ -0,0 +1,509 @@
+import hashlib
+import json
+import unittest
+from unittest.mock import MagicMock, patch
+
+from services.diagnostic_event import (
+    RequestDiagnosticEmitter,
+    build_diagnostic_event,
+)
+from services.diagnostic_reporter import (
+    DiagnosticReporter,
+    NullDiagnosticReporter,
+    _http_transport,
+    diagnostic_reporter_from_config,
+)
+
+
+class DiagnosticReporterTest(unittest.TestCase):
+    def test_event_builder_rejects_invalid_required_and_optional_fields(self):
+        base = {
+            'request_id': 'rq_valid',
+            'stage': 'request_ingress',
+            'status': 'started',
+            'event_code': 'REQUEST_RECEIVED',
+        }
+        cases = (
+            ('request_id', None),
+            ('request_id', 'bad'),
+            ('stage', 'unknown'),
+            ('status', 'unknown'),
+            ('event_code', 1),
+            ('event_code', 'A' * 65),
+            ('event_code', 'lowercase'),
+            ('company_id', True),
+            ('company_id', '1'),
+            ('company_id', 0),
+            ('tool_code', 1),
+            ('tool_code', 'UPPER'),
+            ('session_credential', 1),
+            ('response_code', 'bad-code'),
+            ('summary_code', 'A' * 65),
+            ('cost_ms', True),
+            ('cost_ms', '1'),
+            ('cost_ms', -1),
+            ('cost_ms', 3600001),
+        )
+        for field, value in cases:
+            with self.subTest(field=field, value=value):
+                values = dict(base)
+                values[field] = value
+                with self.assertRaises(ValueError):
+                    build_diagnostic_event(**values)
+
+    def test_event_builder_enforces_stage_context_types(self):
+        cases = (
+            ('request_ingress', 'not-object'),
+            ('request_ingress', {'secret': 'x'}),
+            ('request_ingress', {'http_status': True}),
+            ('request_ingress', {'http_status': '200'}),
+            ('request_ingress', {'http_status': 99}),
+            ('protocol_validation', {'jsonrpc_code': True}),
+            ('protocol_validation', {'jsonrpc_code': 'bad'}),
+            ('response_write', {'client_disconnected': 1}),
+            ('request_ingress', {'transport': 1}),
+            ('request_ingress', {'transport': ''}),
+            ('request_ingress', {'transport': 'x' * 65}),
+        )
+        for stage, context in cases:
+            with self.subTest(stage=stage, context=context):
+                with self.assertRaises(ValueError):
+                    build_diagnostic_event(
+                        request_id='rq_valid',
+                        stage=stage,
+                        status='failed',
+                        event_code='INVALID_EVENT',
+                        context=context,
+                    )
+
+    def test_event_builder_accepts_complete_whitelisted_event(self):
+        event = build_diagnostic_event(
+            request_id='rq_complete',
+            stage='response_write',
+            status='succeeded',
+            event_code='RESPONSE_WRITE_COMPLETED',
+            company_id=1,
+            admin_id=2,
+            tool_code='query_order',
+            response_code='MCP_0000',
+            summary_code='OK',
+            cost_ms=0,
+            context={
+                'http_status': 200,
+                'client_disconnected': False,
+                'transport': 'http',
+            },
+        )
+
+        self.assertTrue(event['occurred_at'].endswith('Z'))
+        self.assertEqual('OK', event['summary_code'])
+
+    def test_event_builder_hashes_session_and_keeps_whitelist(self):
+        event = build_diagnostic_event(
+            request_id='rq_http_test_1',
+            stage='gateway_session',
+            status='failed',
+            event_code='GATEWAY_SESSION_NOT_FOUND',
+            occurred_at='2026-07-20T02:00:00.000Z',
+            session_credential='GWS_super_secret',
+            company_id=1002,
+            admin_id=88,
+            tool_code='query_order_detail',
+            context={'transport': 'http'},
+        )
+
+        self.assertTrue(event['event_id'].startswith('evt_gateway_'))
+        self.assertEqual(
+            hashlib.sha256(b'GWS_super_secret').hexdigest()[:12],
+            event['session_hash'],
+        )
+        self.assertNotIn('session_credential', event)
+        self.assertNotIn('GWS_super_secret', json.dumps(event))
+        self.assertEqual('gateway', event['source'])
+
+    def test_request_emitter_caps_each_request_at_twenty_events(self):
+        reporter = RecordingReporter()
+        emitter = RequestDiagnosticEmitter(reporter, 'rq_http_test_1')
+
+        for index in range(25):
+            emitter.emit(
+                stage='request_ingress',
+                status='started',
+                event_code='REQUEST_RECEIVED',
+                context={'transport': 'http'},
+            )
+
+        self.assertEqual(20, len(reporter.events))
+        self.assertEqual(5, emitter.dropped_count)
+        self.assertEqual(20, len({item['event_id'] for item in reporter.events}))
+
+    def test_request_emitter_counts_reporter_rejection(self):
+        reporter = RecordingReporter(accepted=False)
+        emitter = RequestDiagnosticEmitter(reporter, 'rq_rejected')
+        emitter.set_defaults(context={'transport': 'http'})
+
+        self.assertFalse(emitter.emit(
+            stage='request_ingress',
+            status='started',
+            event_code='REQUEST_RECEIVED',
+        ))
+        self.assertEqual(1, emitter.dropped_count)
+
+    def test_deferred_emitter_enriches_prestage_events_after_identity_resolution(self):
+        reporter = RecordingReporter()
+        emitter = RequestDiagnosticEmitter(
+            reporter,
+            'rq_deferred',
+            defer_until_identity=True,
+        )
+        emitter.emit(
+            stage='request_ingress',
+            status='started',
+            event_code='REQUEST_RECEIVED',
+            context={'transport': 'http'},
+        )
+        emitter.emit(
+            stage='protocol_validation',
+            status='succeeded',
+            event_code='PROTOCOL_VALIDATION_COMPLETED',
+            context={'transport': 'http'},
+        )
+        self.assertEqual([], reporter.events)
+
+        emitter.set_defaults(
+            company_id=1002,
+            admin_id=88,
+            session_credential='GWS_private',
+        )
+
+        self.assertEqual(2, len(reporter.events))
+        self.assertTrue(all(
+            event['company_id'] == 1002 for event in reporter.events
+        ))
+        self.assertNotIn('GWS_private', str(reporter.events))
+
+    def test_deferred_emitter_flushes_unscoped_failures(self):
+        reporter = RecordingReporter()
+        emitter = RequestDiagnosticEmitter(
+            reporter,
+            'rq_unscoped',
+            defer_until_identity=True,
+        )
+        emitter.emit(
+            stage='gateway_session',
+            status='failed',
+            event_code='GATEWAY_SESSION_NOT_FOUND',
+            session_credential='GWS_missing',
+            context={'transport': 'http'},
+        )
+
+        self.assertTrue(emitter.flush())
+        self.assertEqual(1, len(reporter.events))
+        self.assertIn('session_hash', reporter.events[0])
+        self.assertNotIn('company_id', reporter.events[0])
+
+    def test_queue_full_only_increments_drop_count(self):
+        reporter = self.reporter(queue_size=1)
+        first = build_diagnostic_event(
+            request_id='rq_http_1', stage='request_ingress',
+            status='started', event_code='REQUEST_RECEIVED',
+        )
+        second = build_diagnostic_event(
+            request_id='rq_http_2', stage='request_ingress',
+            status='started', event_code='REQUEST_RECEIVED',
+        )
+
+        self.assertTrue(reporter.report(first))
+        self.assertFalse(reporter.report(second))
+        self.assertEqual(1, reporter.drop_count)
+
+    def test_batch_signs_exact_json_bytes_and_accepts_success(self):
+        calls = []
+
+        def transport(url, body, headers, timeout):
+            calls.append((url, body, headers, timeout))
+            return 200, json.dumps({
+                'code': 'MCP_DIAG_INGEST_0000',
+                'msg': 'success',
+                'data': {'accepted': 1, 'duplicate': 0, 'rejected': 0},
+            }).encode('utf-8')
+
+        reporter = self.reporter(transport=transport)
+        reporter.report(build_diagnostic_event(
+            request_id='rq_http_1', stage='request_ingress',
+            status='started', event_code='REQUEST_RECEIVED',
+        ))
+
+        self.assertTrue(reporter.process_once())
+        self.assertEqual(1, len(calls))
+        url, body, headers, timeout = calls[0]
+        self.assertEqual('https://support.internal/internal/mcp-diagnostics/events', url)
+        self.assertEqual(0.5, timeout)
+        self.assertEqual({'events'}, set(json.loads(body.decode('utf-8'))))
+        canonical = 'POST\n/internal/mcp-diagnostics/events\n{0}\n{1}\n{2}'.format(
+            headers['X-MCP-Timestamp'],
+            headers['X-MCP-Nonce'],
+            hashlib.sha256(body).hexdigest(),
+        )
+        expected = hashlib.pbkdf2_hmac(
+            'sha256', canonical.encode('utf-8'), b'x', 1
+        )
+        self.assertNotEqual(expected.hex(), headers['X-MCP-Signature'])
+        import hmac
+        self.assertEqual(
+            hmac.new(b's' * 32, canonical.encode('utf-8'), hashlib.sha256).hexdigest(),
+            headers['X-MCP-Signature'],
+        )
+
+    def test_failed_send_keeps_batch_and_uses_bounded_backoff(self):
+        attempts = []
+        sleeps = []
+
+        def transport(_url, _body, _headers, _timeout):
+            attempts.append(1)
+            if len(attempts) == 1:
+                raise OSError('support unavailable secret')
+            return 200, b'{"code":"MCP_DIAG_INGEST_0000","data":{}}'
+
+        reporter = self.reporter(
+            transport=transport,
+            sleeper=sleeps.append,
+            initial_backoff=0.25,
+            max_backoff=1.0,
+        )
+        reporter.report(build_diagnostic_event(
+            request_id='rq_http_1', stage='request_ingress',
+            status='started', event_code='REQUEST_RECEIVED',
+        ))
+
+        self.assertFalse(reporter.process_once())
+        self.assertEqual([0.25], sleeps)
+        self.assertTrue(reporter.process_once())
+        self.assertEqual(2, len(attempts))
+        self.assertEqual(0, reporter.pending_count)
+
+    def test_thread_start_failure_and_null_reporter_are_non_fatal(self):
+        class BrokenThread:
+            def start(self):
+                raise RuntimeError('thread unavailable')
+
+        reporter = self.reporter(thread_factory=lambda **_kwargs: BrokenThread())
+        self.assertFalse(reporter.start())
+        reporter.close()
+
+        null = NullDiagnosticReporter()
+        self.assertTrue(null.start())
+        self.assertTrue(null.report({'ignored': True}))
+        null.close()
+
+    def test_start_is_idempotent_and_close_joins_live_thread(self):
+        thread = AliveThread()
+        reporter = self.reporter(thread_factory=lambda **_kwargs: thread)
+
+        self.assertTrue(reporter.start())
+        self.assertTrue(reporter.start())
+        reporter.close()
+
+        self.assertEqual([0.5], thread.join_timeouts)
+
+    def test_process_once_accepts_empty_queue_and_rejects_bad_responses(self):
+        reporter = self.reporter()
+        self.assertTrue(reporter.process_once())
+
+        responses = (
+            (199, b'{"code":"MCP_DIAG_INGEST_0000"}'),
+            (300, b'{"code":"MCP_DIAG_INGEST_0000"}'),
+            (200, b'{"code":"MCP_DIAG_INGEST_INVALID"}'),
+        )
+        for status, body in responses:
+            with self.subTest(status=status, body=body):
+                reporter = self.reporter(
+                    transport=lambda *_args, result=(status, body): result,
+                )
+                reporter.report(build_diagnostic_event(
+                    request_id='rq_rejected',
+                    stage='request_ingress',
+                    status='started',
+                    event_code='REQUEST_RECEIVED',
+                ))
+                self.assertFalse(reporter.process_once())
+                self.assertEqual(1, reporter.pending_count)
+
+    def test_batch_size_caps_single_send_and_run_loop_handles_retry(self):
+        reporter = self.reporter(batch_size=1)
+        for request_id in ('rq_one', 'rq_two'):
+            reporter.report(build_diagnostic_event(
+                request_id=request_id,
+                stage='request_ingress',
+                status='started',
+                event_code='REQUEST_RECEIVED',
+            ))
+        self.assertTrue(reporter.process_once())
+        self.assertEqual(0, reporter.pending_count)
+        self.assertFalse(reporter._queue.empty())
+
+        reporter = self.reporter()
+        reporter._stop = MagicMock()
+        reporter._stop.is_set.side_effect = (False, False, True)
+        reporter.process_once = MagicMock(side_effect=(False, True))
+        reporter._run()
+        reporter._stop.wait.assert_called_once_with(0.1)
+
+    def test_default_http_transport_posts_and_reads_response(self):
+        response = MagicMock()
+        response.getcode.return_value = 200
+        response.read.return_value = b'{}'
+        response.__enter__.return_value = response
+        response.__exit__.return_value = False
+
+        with patch(
+            'services.diagnostic_reporter.urllib.request.urlopen',
+            return_value=response,
+        ) as urlopen:
+            result = _http_transport(
+                'https://support.test/path',
+                b'{}',
+                {'X-Test': '1'},
+                0.5,
+            )
+
+        self.assertEqual((200, b'{}'), result)
+        self.assertEqual(0.5, urlopen.call_args.kwargs['timeout'])
+
+    def test_config_factory_uses_null_when_disabled_or_invalid(self):
+        disabled = type('Config', (), {'diagnosis_enabled': False})()
+        invalid = type('Config', (), {
+            'diagnosis_enabled': True,
+            'diagnosis_url': '',
+            'diagnosis_key_id': '',
+            'diagnosis_secret': '',
+        })()
+
+        self.assertIsInstance(
+            diagnostic_reporter_from_config(disabled),
+            NullDiagnosticReporter,
+        )
+        self.assertIsInstance(
+            diagnostic_reporter_from_config(invalid),
+            NullDiagnosticReporter,
+        )
+
+    def test_config_factory_starts_enabled_reporter(self):
+        config = type('Config', (), {
+            'diagnosis_enabled': True,
+            'diagnosis_url': 'https://support.test/internal/mcp-diagnostics/events',
+            'diagnosis_key_id': 'gateway-current',
+            'diagnosis_secret': 's' * 32,
+            'diagnosis_queue_size': 10,
+            'diagnosis_batch_size': 20,
+            'diagnosis_timeout_seconds': 0.5,
+            'diagnosis_initial_backoff_seconds': 0.25,
+            'diagnosis_max_backoff_seconds': 2.0,
+        })()
+
+        reporter = diagnostic_reporter_from_config(
+            config,
+            thread_factory=lambda **_kwargs: StartedThread(),
+        )
+
+        self.assertIsInstance(reporter, DiagnosticReporter)
+        self.assertTrue(reporter._thread.started)
+        reporter.close()
+
+    def test_config_factory_requires_explicit_opt_in_for_http(self):
+        values = {
+            'diagnosis_enabled': True,
+            'diagnosis_url': 'http://support.test/internal/mcp-diagnostics/events',
+            'diagnosis_key_id': 'gateway-current',
+            'diagnosis_secret': 's' * 32,
+        }
+        default_config = type('Config', (), values)()
+        allowed_config = type(
+            'Config',
+            (),
+            dict(values, diagnosis_allow_insecure_http=True),
+        )()
+
+        self.assertIsInstance(
+            diagnostic_reporter_from_config(default_config),
+            NullDiagnosticReporter,
+        )
+        reporter = diagnostic_reporter_from_config(
+            allowed_config,
+            thread_factory=lambda **_kwargs: StartedThread(),
+        )
+        self.assertIsInstance(reporter, DiagnosticReporter)
+        reporter.close()
+
+    def test_config_factory_falls_back_when_thread_cannot_start(self):
+        config = type('Config', (), {
+            'diagnosis_enabled': True,
+            'diagnosis_url': 'https://support.test/events',
+            'diagnosis_key_id': 'gateway-current',
+            'diagnosis_secret': 's' * 32,
+        })()
+
+        reporter = diagnostic_reporter_from_config(
+            config,
+            thread_factory=lambda **_kwargs: BrokenThread(),
+        )
+
+        self.assertIsInstance(reporter, NullDiagnosticReporter)
+
+    def reporter(self, **overrides):
+        options = {
+            'url': 'https://support.internal/internal/mcp-diagnostics/events',
+            'key_id': 'gateway-current',
+            'secret': 's' * 32,
+            'queue_size': 10,
+            'batch_size': 100,
+            'timeout_seconds': 0.5,
+            'transport': lambda *_args: (200, b'{"code":"MCP_DIAG_INGEST_0000","data":{}}'),
+            'clock': lambda: 1784512800,
+            'nonce_factory': lambda: 'nonce-gateway-1234',
+            'sleeper': lambda _seconds: None,
+        }
+        options.update(overrides)
+        return DiagnosticReporter(**options)
+
+
+class RecordingReporter:
+    def __init__(self, accepted=True):
+        self.events = []
+        self.accepted = accepted
+
+    def report(self, event):
+        self.events.append(event)
+        return self.accepted
+
+
+class StartedThread:
+    def __init__(self):
+        self.started = False
+
+    def start(self):
+        self.started = True
+
+    def is_alive(self):
+        return False
+
+
+class AliveThread(StartedThread):
+    def __init__(self):
+        super().__init__()
+        self.join_timeouts = []
+
+    def is_alive(self):
+        return True
+
+    def join(self, timeout=None):
+        self.join_timeouts.append(timeout)
+
+
+class BrokenThread:
+    def start(self):
+        raise RuntimeError('thread unavailable')
+
+
+if __name__ == '__main__':
+    unittest.main()

+ 97 - 2
tests/test_mcp_protocol.py

@@ -119,7 +119,7 @@ class FullColumnsApiClient(DummyApiClient):
 
 
 class McpProtocolTest(unittest.TestCase):
-    def build_handler(self, api_client=None):
+    def build_handler(self, api_client=None, reporter=None):
         token_store = InMemoryTokenStore(refresh_skew_seconds=60)
         token_store.save('MT_demo', '2099-01-01T00:00:00')
         app = GatewayApp(
@@ -127,7 +127,93 @@ class McpProtocolTest(unittest.TestCase):
             api_client=api_client or DummyApiClient(),
             token_store=token_store,
         )
-        return McpProtocolHandler(app)
+        return McpProtocolHandler(app, reporter=reporter)
+
+    def test_stdio_tool_call_emits_correlated_diagnostic_stages(self):
+        reporter = RecordingReporter()
+        handler = self.build_handler(reporter=reporter)
+
+        response = handler.handle_request({
+            'jsonrpc': '2.0',
+            'id': 10,
+            'method': 'tools/call',
+            'params': {
+                'name': 'query_order',
+                'arguments': {'keyword': 'SO20260706001'},
+            },
+        })
+
+        self.assertFalse(response['result']['isError'])
+        self.assertEqual(
+            [
+                'request_ingress',
+                'protocol_validation',
+                'backend_call',
+                'backend_call',
+                'response_safety',
+            ],
+            [event['stage'] for event in reporter.events],
+        )
+        self.assertEqual(1, len({event['request_id'] for event in reporter.events}))
+
+    def test_stdio_missing_tool_name_returns_invalid_params(self):
+        reporter = RecordingReporter()
+        response = self.build_handler(reporter=reporter).handle_request({
+            'jsonrpc': '2.0',
+            'id': 11,
+            'method': 'tools/call',
+            'params': {'arguments': {}},
+        })
+
+        self.assertIn('error', response)
+        self.assertEqual(-32602, response['error']['code'])
+        self.assertNotIn('result', response)
+        self.assertEqual(
+            ['request_ingress', 'protocol_validation'],
+            [event['stage'] for event in reporter.events],
+        )
+        self.assertEqual('failed', reporter.events[-1]['status'])
+        self.assertEqual(
+            'PARAM_VALIDATION_FAILED',
+            reporter.events[-1]['event_code'],
+        )
+
+    def test_stdio_reporter_failure_does_not_change_response(self):
+        class BrokenReporter:
+            def report(self, _event):
+                raise RuntimeError('support unavailable')
+
+        response = self.build_handler(reporter=BrokenReporter()).handle_request({
+            'jsonrpc': '2.0',
+            'id': 1,
+            'method': 'initialize',
+            'params': {},
+        })
+
+        self.assertIn('result', response)
+
+    def test_stdio_invalid_tool_result_emits_response_safety_failure(self):
+        class InvalidGateway:
+            def call_tool(self, name, arguments, request_id=''):
+                return []
+
+        reporter = RecordingReporter()
+        response = McpProtocolHandler(
+            InvalidGateway(),
+            reporter=reporter,
+        ).handle_request({
+            'jsonrpc': '2.0',
+            'id': 1,
+            'method': 'tools/call',
+            'params': {'name': 'query_order', 'arguments': {}},
+        })
+
+        self.assertTrue(response['result']['isError'])
+        event = next(
+            item for item in reporter.events
+            if item['stage'] == 'response_safety'
+        )
+        self.assertEqual('failed', event['status'])
 
     def test_non_tool_exception_is_sanitized(self):
         class ExplodingGateway:
@@ -433,5 +519,14 @@ class McpProtocolTest(unittest.TestCase):
         self.assertEqual('2025-06-18', response['result']['protocolVersion'])
 
 
+class RecordingReporter:
+    def __init__(self):
+        self.events = []
+
+    def report(self, event):
+        self.events.append(event)
+        return True
+
+
 if __name__ == '__main__':
     unittest.main()

+ 2 - 5
tests/test_mcp_protocol_coverage.py

@@ -402,11 +402,8 @@ class HandleMessageEdgeCasesTest(unittest.TestCase):
                     'params': {'name': tool_name, 'arguments': {}},
                 })
 
-                self.assertTrue(response['result']['isError'])
-                self.assertEqual(
-                    '工具返回格式异常',
-                    response['result']['structuredContent']['message'],
-                )
+                self.assertEqual(-32602, response['error']['code'])
+                self.assertNotIn('result', response)
 
         handler.gateway_app.call_tool.assert_not_called()
 

+ 181 - 0
tests/test_public_gateway.py

@@ -1,6 +1,7 @@
 import unittest
 
 from public_gateway import PublicGatewayApp
+from services.diagnostic_event import RequestDiagnosticEmitter
 
 
 class FakeSessionStore:
@@ -38,6 +39,177 @@ class FakeApiClient:
 
 
 class PublicGatewayAppTest(unittest.TestCase):
+    def test_missing_redis_session_emits_gateway_session_failure(self):
+        reporter = RecordingReporter()
+        emitter = RequestDiagnosticEmitter(
+            reporter,
+            'rq_missing',
+            defer_until_identity=True,
+        )
+        app = PublicGatewayApp(FakeSessionStore(), FakeApiClient())
+
+        with self.assertRaises(RuntimeError):
+            app.call_tool(
+                'GWS_missing',
+                'query_order',
+                request_id='rq_missing',
+                diagnostic_emitter=emitter,
+            )
+        emitter.flush()
+
+        event = reporter.events[-1]
+        self.assertEqual('gateway_session', event['stage'])
+        self.assertEqual('failed', event['status'])
+        self.assertEqual('GATEWAY_SESSION_NOT_FOUND', event['event_code'])
+        self.assertIn('session_hash', event)
+
+    def test_enabled_tool_lookup_failure_emits_backend_failure(self):
+        class FailingListClient(FakeApiClient):
+            def list_enabled_tools(self, token, request_id=''):
+                raise OSError('registry unavailable secret')
+
+        store = FakeSessionStore()
+        store.sessions['GWS_A'] = {
+            'mcp_token': 'MT_A',
+            'admin_id': 88,
+            'company_id': 1002,
+        }
+        reporter = RecordingReporter()
+        emitter = RequestDiagnosticEmitter(
+            reporter,
+            'rq_lookup',
+            defer_until_identity=True,
+        )
+        emitter.emit(
+            stage='request_ingress',
+            status='started',
+            event_code='REQUEST_RECEIVED',
+            context={'transport': 'http'},
+        )
+        app = PublicGatewayApp(store, FailingListClient())
+
+        with self.assertRaises(OSError):
+            app.call_tool(
+                'GWS_A',
+                'query_order',
+                request_id='rq_lookup',
+                diagnostic_emitter=emitter,
+            )
+
+        self.assertTrue(all(
+            event['company_id'] == 1002 for event in reporter.events
+        ))
+        failed = reporter.events[-1]
+        self.assertEqual('backend_call', failed['stage'])
+        self.assertEqual('failed', failed['status'])
+        self.assertEqual('ENABLED_TOOL_LOOKUP_FAILED', failed['event_code'])
+
+    def test_unknown_and_disabled_tools_emit_backend_failures(self):
+        store = FakeSessionStore()
+        store.sessions['GWS_A'] = {
+            'mcp_token': 'MT_A',
+            'admin_id': 88,
+            'company_id': 1002,
+        }
+        app = PublicGatewayApp(store, FakeApiClient())
+
+        cases = (
+            ('not_registered', 'TOOL_NOT_REGISTERED', KeyError),
+            ('query_outbound_detail', 'TOOL_DISABLED', RuntimeError),
+        )
+        for tool_name, event_code, exception_class in cases:
+            with self.subTest(tool_name=tool_name):
+                reporter = RecordingReporter()
+                emitter = RequestDiagnosticEmitter(
+                    reporter,
+                    'rq_tool_check',
+                    defer_until_identity=True,
+                )
+                with self.assertRaises(exception_class):
+                    app.call_tool(
+                        'GWS_A',
+                        tool_name,
+                        request_id='rq_tool_check',
+                        diagnostic_emitter=emitter,
+                    )
+                self.assertEqual(event_code, reporter.events[-1]['event_code'])
+                self.assertEqual('failed', reporter.events[-1]['status'])
+
+    def test_enabled_tool_lookup_failure_without_emitter_keeps_exception(self):
+        class FailingListClient(FakeApiClient):
+            def list_enabled_tools(self, token, request_id=''):
+                raise OSError('registry unavailable')
+
+        store = FakeSessionStore()
+        store.sessions['GWS_A'] = {'mcp_token': 'MT_A'}
+        app = PublicGatewayApp(store, FailingListClient())
+
+        with self.assertRaises(OSError):
+            app.call_tool('GWS_A', 'query_order', request_id='rq_lookup')
+    def test_diagnostic_events_include_session_identity_and_backend_result(self):
+        store = FakeSessionStore()
+        store.sessions['GWS_A'] = {
+            'mcp_token': 'MT_A',
+            'admin_id': 88,
+            'company_id': 1002,
+        }
+        reporter = RecordingReporter()
+        emitter = RequestDiagnosticEmitter(reporter, 'rq_a')
+        app = PublicGatewayApp(
+            session_store=store,
+            api_client=FakeApiClient(),
+            auth_client=None,
+        )
+
+        app.call_tool(
+            'GWS_A',
+            'query_order',
+            {'keyword': 'A'},
+            request_id='rq_a',
+            diagnostic_emitter=emitter,
+        )
+
+        self.assertEqual(
+            ['gateway_session', 'backend_call', 'backend_call'],
+            [event['stage'] for event in reporter.events],
+        )
+        self.assertEqual(
+            ['succeeded', 'started', 'succeeded'],
+            [event['status'] for event in reporter.events],
+        )
+        self.assertTrue(all(event['company_id'] == 1002 for event in reporter.events))
+        self.assertNotIn('GWS_A', str(reporter.events))
+        self.assertNotIn('MT_A', str(reporter.events))
+
+    def test_backend_failure_emits_failed_event_without_changing_exception(self):
+        class FailingApiClient(FakeApiClient):
+            def call_tool(self, *args, **kwargs):
+                raise OSError('backend secret unavailable')
+
+        store = FakeSessionStore()
+        store.sessions['GWS_A'] = {
+            'mcp_token': 'MT_A',
+            'admin_id': 88,
+            'company_id': 1002,
+        }
+        reporter = RecordingReporter()
+        emitter = RequestDiagnosticEmitter(reporter, 'rq_a')
+        app = PublicGatewayApp(store, FailingApiClient(), auth_client=None)
+
+        with self.assertRaisesRegex(OSError, 'backend secret unavailable'):
+            app.call_tool(
+                'GWS_A',
+                'query_order',
+                request_id='rq_a',
+                diagnostic_emitter=emitter,
+            )
+
+        failed = reporter.events[-1]
+        self.assertEqual('backend_call', failed['stage'])
+        self.assertEqual('failed', failed['status'])
+        self.assertEqual('UNEXPECTED_EXCEPTION', failed['event_code'])
+        self.assertNotIn('backend secret', str(failed))
+
     def test_two_employees_use_isolated_tokens(self):
         store = FakeSessionStore()
         store.sessions['GWS_A'] = {'mcp_token': 'MT_A'}
@@ -94,5 +266,14 @@ class PublicGatewayAppTest(unittest.TestCase):
         self.assertEqual(['GWS_A'], store.touched)
 
 
+class RecordingReporter:
+    def __init__(self):
+        self.events = []
+
+    def report(self, event):
+        self.events.append(event)
+        return True
+
+
 if __name__ == '__main__':
     unittest.main()

+ 198 - 1
tests/test_public_server.py

@@ -3,6 +3,7 @@ import unittest
 from unittest.mock import Mock
 
 from public_server import PublicMcpHttpHandler, create_http_handler, extract_client_ip
+from services.diagnostic_event import RequestDiagnosticEmitter
 from utils.rate_limiter import SimpleRateLimiter
 
 
@@ -32,12 +33,136 @@ class FakeGateway:
         self.list_calls.append((gateway_session_id, request_id))
         return [{'name': 'query_track', 'description': 'query track', 'input_schema': {'type': 'object'}}]
 
-    def call_tool(self, gateway_session_id, name, arguments=None, request_id='', client_ip=''):
+    def call_tool(
+        self,
+        gateway_session_id,
+        name,
+        arguments=None,
+        request_id='',
+        client_ip='',
+        diagnostic_emitter=None,
+    ):
         self.calls.append((gateway_session_id, name, arguments, request_id, client_ip))
         return self.tool_result
 
 
 class PublicMcpHttpHandlerTest(unittest.TestCase):
+    def test_initialize_emits_successful_protocol_validation(self):
+        reporter = RecordingReporter()
+        handler = PublicMcpHttpHandler(FakeGateway(), reporter=reporter)
+
+        response = handler.handle_json_rpc(
+            headers={},
+            message={'jsonrpc': '2.0', 'id': 1, 'method': 'initialize'},
+        )
+
+        self.assertIn('result', response)
+        self.assertEqual(
+            ['request_ingress', 'protocol_validation'],
+            [event['stage'] for event in reporter.events],
+        )
+        self.assertEqual('succeeded', reporter.events[-1]['status'])
+
+    def test_missing_session_emits_safe_ingress_and_session_failure(self):
+        reporter = RecordingReporter()
+        handler = PublicMcpHttpHandler(
+            FakeGateway(),
+            context_parser=FakeParser(),
+            reporter=reporter,
+        )
+
+        response = handler.handle_json_rpc(
+            headers={},
+            message={
+                'jsonrpc': '2.0',
+                'id': 1,
+                'method': 'tools/call',
+                'params': {'name': 'query_order', 'arguments': {}},
+            },
+        )
+
+        self.assertEqual(-32001, response['error']['code'])
+        self.assertEqual(
+            ['request_ingress', 'protocol_validation', 'gateway_session'],
+            [event['stage'] for event in reporter.events],
+        )
+        self.assertEqual('failed', reporter.events[-1]['status'])
+
+    def test_reporter_failure_does_not_change_successful_response(self):
+        class BrokenReporter:
+            def report(self, _event):
+                raise RuntimeError('support unavailable')
+
+        handler = PublicMcpHttpHandler(
+            FakeGateway(),
+            context_parser=FakeParser(),
+            reporter=BrokenReporter(),
+        )
+
+        response = handler.handle_json_rpc(
+            headers={'X-Gateway-Session': 'GWS_A'},
+            message={'jsonrpc': '2.0', 'id': 1, 'method': 'initialize'},
+        )
+
+        self.assertIn('result', response)
+
+    def test_invalid_backend_result_emits_response_safety_failure(self):
+        reporter = RecordingReporter()
+        gateway = FakeGateway()
+        gateway.tool_result = []
+        handler = PublicMcpHttpHandler(
+            gateway,
+            context_parser=FakeParser(),
+            reporter=reporter,
+        )
+
+        response = handler.handle_json_rpc(
+            headers={'X-Gateway-Session': 'GWS_A'},
+            message={
+                'jsonrpc': '2.0',
+                'id': 1,
+                'method': 'tools/call',
+                'params': {'name': 'query_order', 'arguments': {}},
+            },
+        )
+
+        self.assertTrue(response['result']['isError'])
+        event = next(
+            item for item in reporter.events
+            if item['stage'] == 'response_safety'
+        )
+        self.assertEqual('failed', event['status'])
+        self.assertEqual('RESPONSE_SAFETY_REJECTED', event['event_code'])
+
+    def test_missing_tool_name_emits_protocol_validation_failure(self):
+        reporter = RecordingReporter()
+        handler = PublicMcpHttpHandler(
+            FakeGateway(),
+            context_parser=FakeParser(),
+            reporter=reporter,
+        )
+
+        response = handler.handle_json_rpc(
+            headers={'X-Gateway-Session': 'GWS_A'},
+            message={
+                'jsonrpc': '2.0',
+                'id': 1,
+                'method': 'tools/call',
+                'params': {'arguments': {}},
+            },
+        )
+
+        self.assertNotIn('result', response)
+        self.assertEqual(-32602, response['error']['code'])
+        self.assertTrue(response['error']['data']['request_id'].startswith('rq_http_'))
+        self.assertEqual(
+            'failed',
+            next(
+                event for event in reporter.events
+                if event['stage'] == 'protocol_validation'
+            )['status'],
+        )
+
     def test_constructor_uses_registered_names_without_loading_dynamic_list(self):
         gateway = FakeGateway()
 
@@ -308,6 +433,43 @@ class PublicMcpHttpHandlerTest(unittest.TestCase):
 
 
 class RateLimitTest(unittest.TestCase):
+    def test_rate_limit_helper_remains_usable_without_emitter(self):
+        handler = self._make_handler(max_requests=1)
+        self.assertIsNone(handler._check_rate_limit(
+            'GWS_A:query_order', 'tools/call', 'rq_http_first'
+        ))
+
+        response = handler._check_rate_limit(
+            'GWS_A:query_order', 'tools/call', 'rq_http_second'
+        )
+
+        self.assertEqual(-32029, response['error']['code'])
+
+    def test_rate_limit_rejection_emits_failed_diagnostic_event(self):
+        reporter = RecordingReporter()
+        limiter = SimpleRateLimiter(
+            max_requests=1,
+            window_seconds=60,
+            max_in_flight=1,
+        )
+        handler = PublicMcpHttpHandler(
+            FakeGateway(),
+            context_parser=FakeParser(),
+            rate_limiter=limiter,
+            reporter=reporter,
+        )
+        headers = {'X-Gateway-Session': 'GWS_A'}
+
+        handler.handle_json_rpc(headers, self._tools_call_msg())
+        reporter.events.clear()
+        handler.handle_json_rpc(headers, self._tools_call_msg())
+
+        event = next(
+            item for item in reporter.events if item['stage'] == 'rate_limit'
+        )
+        self.assertEqual('failed', event['status'])
+        self.assertEqual('RATE_LIMIT_EXCEEDED', event['event_code'])
+
     def _make_handler(self, max_requests=2, max_in_flight=2):
         gateway = FakeGateway()
         limiter = SimpleRateLimiter(
@@ -463,6 +625,32 @@ class RateLimitTest(unittest.TestCase):
 
 
 class HttpDisconnectTest(unittest.TestCase):
+    def test_broken_pipe_emits_response_write_failure(self):
+        reporter = RecordingReporter()
+        emitter = RequestDiagnosticEmitter(reporter, 'rq_http_write')
+        handler_class = create_http_handler(
+            FakeGateway(),
+            rate_limiter=None,
+            reporter=reporter,
+        )
+        handler = object.__new__(handler_class)
+        handler.send_response = Mock()
+        handler.send_header = Mock()
+        handler.end_headers = Mock()
+        handler.wfile = Mock()
+        handler.wfile.write.side_effect = BrokenPipeError()
+
+        handler._write_json(
+            {'jsonrpc': '2.0', 'id': 1, 'result': {}},
+            diagnostic_emitter=emitter,
+        )
+
+        self.assertEqual('response_write', reporter.events[-1]['stage'])
+        self.assertEqual('failed', reporter.events[-1]['status'])
+        self.assertTrue(
+            reporter.events[-1]['context']['client_disconnected']
+        )
+
     def test_broken_pipe_while_writing_response_is_handled(self):
         handler_class = create_http_handler(FakeGateway(), rate_limiter=None)
         handler = object.__new__(handler_class)
@@ -477,5 +665,14 @@ class HttpDisconnectTest(unittest.TestCase):
 
         self.assertIn('client disconnected before response', logs.output[0])
 
+
+class RecordingReporter:
+    def __init__(self):
+        self.events = []
+
+    def report(self, event):
+        self.events.append(event)
+        return True
+
 if __name__ == '__main__':
     unittest.main()