From 83e9fd1724b7a210f85ac15e6210911aa1c7cf2e Mon Sep 17 00:00:00 2001 From: Kris <2893855659@qq.com> Date: Sat, 18 Jul 2026 18:44:39 +0800 Subject: [PATCH] =?UTF-8?q?fix(channels):=20=E4=BF=AE=E5=A4=8D=E5=87=BA?= =?UTF-8?q?=E7=AB=99=E4=BC=9A=E8=AF=9D=E9=94=81=E5=BF=99=E6=97=B6=E7=9A=84?= =?UTF-8?q?=E5=BC=82=E5=B8=B8=E5=A4=84=E7=90=86=E9=80=BB=E8=BE=91?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 1. 移除对ConflictError的抛出,改为记录INFO日志并调用_handleOutboundFailure 2. 新增单元测试覆盖锁忙场景的正确行为:不崩溃、触发幂等回退、允许渠道重试 3. 补充中文注释说明锁忙场景的处理逻辑与设计意图 --- .../usecase/inbound_message_service.py | 30 +++- .../usecase/test_inbound_message_service.py | 134 ++++++++++++++++++ 2 files changed, 159 insertions(+), 5 deletions(-) diff --git a/backend/package/yuxi/channels/application/messaging/usecase/inbound_message_service.py b/backend/package/yuxi/channels/application/messaging/usecase/inbound_message_service.py index 90267beb..de32c33c 100644 --- a/backend/package/yuxi/channels/application/messaging/usecase/inbound_message_service.py +++ b/backend/package/yuxi/channels/application/messaging/usecase/inbound_message_service.py @@ -51,7 +51,6 @@ from yuxi.channels.contract.dtos.messaging.inbound import ( from yuxi.channels.contract.errors import ( AgentRunCancelledError, AgentRunInterruptedError, - ConflictError, Error, IdempotencyConflictError, InternalError, @@ -659,11 +658,22 @@ class InboundMessageService: lock_key = f"outbound:{outbound_ctx.conversation_id}" lock_token = await self._cache_port.acquireAdvisoryLock(lock_key, ttl_seconds=60) if lock_token is None: + # 锁忙是可预期的并发冲突(同会话前一条出站未完成,或进程 + # 重启遗留 stale lock),不作为 crash 处理。走 + # _handleOutboundFailure 触发 NACK + 幂等回退,允许渠道侧 + # 重试,保持 ACK 互斥语义不被绕过(与 pipeline.run 失败 + # 路径行为一致)。 error_occurred = True - raise ConflictError( - "outbound pipeline busy for conversation", + await self._logger.info( + "outbound pipeline busy for conversation, skip delivery", trace_id=ctx.trace_id, + agent_run_id=ctx.agent_run_id, + channel_type=ctx.channel_type, + account_id=ctx.account_id, + conversation_id=outbound_ctx.conversation_id, ) + await self._handleOutboundFailure(ctx, "outbound pipeline busy for conversation") + return try: ok, err = await self._outbound_pipeline.run(outbound_ctx) except Exception as e: @@ -784,11 +794,21 @@ class InboundMessageService: lock_key = f"outbound:{outbound_ctx.conversation_id}" lock_token = await self._cache_port.acquireAdvisoryLock(lock_key, ttl_seconds=60) if lock_token is None: + # 锁忙是可预期的并发冲突(同会话前一条出站未完成,或进程 + # 重启遗留 stale lock),不作为 crash 处理。走 + # _handleOutboundFailure 触发 NACK + 幂等回退,允许渠道侧 + # 重试,保持 ACK 互斥语义不被绕过(与 pipeline.run 失败 + # 路径行为一致)。 error_occurred = True - raise ConflictError( - "outbound pipeline busy for conversation", + await self._logger.info( + "outbound pipeline busy for conversation, skip delivery", trace_id=ctx.trace_id, + channel_type=ctx.channel_type, + account_id=ctx.account_id, + conversation_id=outbound_ctx.conversation_id, ) + await self._handleOutboundFailure(ctx, "outbound pipeline busy for conversation") + return try: ok, err = await self._outbound_pipeline.run(outbound_ctx) except Exception as e: diff --git a/backend/test/unit/channels/application/usecase/test_inbound_message_service.py b/backend/test/unit/channels/application/usecase/test_inbound_message_service.py index a031c668..fd1d4b13 100644 --- a/backend/test/unit/channels/application/usecase/test_inbound_message_service.py +++ b/backend/test/unit/channels/application/usecase/test_inbound_message_service.py @@ -1109,6 +1109,140 @@ class TestOutboundSkipOnAgentRunTerminalState: assert result.ack_decision == "pending" +@pytest.mark.unit +class TestOutboundLockBusy: + """H-O5:会话级出站锁忙时按 delivery failure 处理,不作为 crash。 + + 覆盖 ``_deliverAgentResponse`` / ``_deliverSilentCommandResponse`` 在 + ``acquireAdvisoryLock`` 返回 None(同会话前一条出站未完成,或进程重启 + 遗留 stale lock)时的行为: + - 不抛 ConflictError 穿透到 ``_deliverOutboundAsync`` 兜底 except + - 记录 INFO 日志(不是 ERROR crashed) + - 调用 ``_handleOutboundFailure`` 触发 NACK + 幂等回退 + - 不执行 ``outbound_pipeline.run``(未拿到锁,直接跳过) + - 不调用 ``releaseAdvisoryLock``(未拿到锁无需释放) + """ + + @pytest.mark.asyncio + async def test_agent_response_lock_busy_triggers_outbound_failure( + self, + fake_logger, + fake_config, + fake_tracer, + fake_transaction_port, + fake_idempotency_repo, + ): + # Arrange Agent 响应出站时锁忙 → 走 _handleOutboundFailure + inbound_pipeline = AsyncMock() + + async def _set_agent_run_with_idempotency(ctx): + ctx.agent_run_id = "run-1" + ctx.idempotency_record_id = 42 + return (True, None) + + inbound_pipeline.run.side_effect = _set_agent_run_with_idempotency + outbound_pipeline = AsyncMock() + streaming_config = MagicMock() + streaming_config.streaming_enabled = False + fake_config.getStreamingConfig.return_value = streaming_config + # 锁忙:acquireAdvisoryLock 返回 None + cache_port = AsyncMock() + cache_port.acquireAdvisoryLock.return_value = None + service = _make_service( + fake_logger, + fake_config, + fake_tracer, + fake_transaction_port, + inbound_pipeline=inbound_pipeline, + outbound_pipeline=outbound_pipeline, + idempotency_repository=fake_idempotency_repo, + cache_port=cache_port, + ) + + # Act + result = await service.receiveInbound(_make_inbound_cmd()) + await _drain_outbound_tasks(service) + + # Assert 幂等记录回退为 failed(走 _handleOutboundFailure) + last_call = fake_idempotency_repo.updateIdempotencyStatus.call_args_list[-1] + assert last_call.args[0] == 42 + assert last_call.args[1] == "failed" + # 未执行出站管道(锁忙直接 return) + assert outbound_pipeline.run.await_count == 0 + # 未调用 releaseAdvisoryLock(未拿到锁) + assert cache_port.releaseAdvisoryLock.await_count == 0 + # INFO 日志记录锁忙(不是 ERROR crashed) + busy_logs = [ + call for call in fake_logger.info.call_args_list + if "outbound pipeline busy" in str(call) + ] + assert busy_logs, "应记录 'outbound pipeline busy' INFO 日志" + # 不应有 crash 级 ERROR 日志 + crash_logs = [ + call for call in fake_logger.exception.call_args_list + if "outbound delivery task crashed" in str(call) + ] + assert not crash_logs, "锁忙不应记为 crashed" + assert result.ack_decision == "pending" + + @pytest.mark.asyncio + async def test_silent_command_lock_busy_triggers_outbound_failure( + self, + fake_logger, + fake_config, + fake_tracer, + fake_transaction_port, + fake_idempotency_repo, + ): + # Arrange 静默命令出站时锁忙 → 走 _handleOutboundFailure + inbound_pipeline = AsyncMock() + + async def _set_silent_command_with_idempotency(ctx): + ctx.is_silent = True + ctx.command = CommandResult( + command="/help", + is_silent=True, + response=CommandResponse(content="help"), + ) + ctx.idempotency_record_id = 42 + return (True, None) + + inbound_pipeline.run.side_effect = _set_silent_command_with_idempotency + outbound_pipeline = AsyncMock() + streaming_config = MagicMock() + streaming_config.streaming_enabled = False + fake_config.getStreamingConfig.return_value = streaming_config + cache_port = AsyncMock() + cache_port.acquireAdvisoryLock.return_value = None + service = _make_service( + fake_logger, + fake_config, + fake_tracer, + fake_transaction_port, + inbound_pipeline=inbound_pipeline, + outbound_pipeline=outbound_pipeline, + idempotency_repository=fake_idempotency_repo, + cache_port=cache_port, + ) + + # Act + result = await service.receiveInbound(_make_inbound_cmd()) + await _drain_outbound_tasks(service) + + # Assert 幂等记录回退为 failed + last_call = fake_idempotency_repo.updateIdempotencyStatus.call_args_list[-1] + assert last_call.args[1] == "failed" + # 未执行出站管道 + assert outbound_pipeline.run.await_count == 0 + # INFO 日志记录锁忙 + busy_logs = [ + call for call in fake_logger.info.call_args_list + if "outbound pipeline busy" in str(call) + ] + assert busy_logs, "应记录 'outbound pipeline busy' INFO 日志" + assert result.ack_decision == "pending" + + @pytest.mark.unit class TestReceiveWebhook: """receiveWebhook 行为测试。"""