From e2b928273c06df663f6c61041231690bf6eae24e Mon Sep 17 00:00:00 2001 From: stefanfeng Date: Thu, 20 Aug 2026 17:10:54 +0800 Subject: [PATCH] fix(avatar): acknowledge BOXIM messages as read --- .../backend/services/boxim_client.py | 21 +++++++++ .../backend/services/takeover_service.py | 15 ++++++- .../backend/tests/test_boxim_client.py | 12 ++++++ .../backend/tests/test_takeover_service.py | 43 +++++++++++++++++++ 4 files changed, 90 insertions(+), 1 deletion(-) diff --git a/digital-avatar-app/backend/services/boxim_client.py b/digital-avatar-app/backend/services/boxim_client.py index 9c50af0..f3a0fdf 100644 --- a/digital-avatar-app/backend/services/boxim_client.py +++ b/digital-avatar-app/backend/services/boxim_client.py @@ -162,6 +162,27 @@ class BoxIMClient: raise BoxIMError("BOXIM 私聊消息格式不正确") return [item for item in data if isinstance(item, dict)] + async def mark_private_messages_read( + self, + access_token: str, + friend_id: int | str, + message_id: int | str, + ) -> None: + """Mark one private conversation read through its latest received message.""" + friend_id_text = str(friend_id).strip() + message_id_text = str(message_id).strip() + if not friend_id_text.isdigit() or not message_id_text.isdigit(): + raise BoxIMError("BOXIM 已读回执参数不正确") + await self._request( + "PUT", + "/message/private/readed", + access_token, + params={ + "friendId": int(friend_id_text), + "messageId": int(message_id_text), + }, + ) + async def send_private_message( self, access_token: str, diff --git a/digital-avatar-app/backend/services/takeover_service.py b/digital-avatar-app/backend/services/takeover_service.py index 8fc46eb..d803345 100644 --- a/digital-avatar-app/backend/services/takeover_service.py +++ b/digital-avatar-app/backend/services/takeover_service.py @@ -268,6 +268,7 @@ class TakeoverService: messages.sort(key=lambda item: (_numeric_id(item.get("id")), item.get("sendTime") or 0)) priming = not bool(cursor.initialized) max_message_id = _numeric_id(cursor.last_message_id) + read_receipts: dict[str, int] = {} for message in messages: self._record_message( db, @@ -276,7 +277,19 @@ class TakeoverService: message, schedule_reply=not priming, ) - max_message_id = max(max_message_id, _numeric_id(message.get("id"))) + message_id = _numeric_id(message.get("id")) + max_message_id = max(max_message_id, message_id) + send_id = str(message.get("sendId") or "") + recv_id = str(message.get("recvId") or "") + if recv_id == cursor.boxim_owner_id and send_id and message_id: + read_receipts[send_id] = max(read_receipts.get(send_id, 0), message_id) + + # BOXIM publishes this HTTP state change to connected socket clients. + # Do it before advancing the cursor so a failed receipt is retried. + for peer_id, message_id in read_receipts.items(): + await self.boxim.mark_private_messages_read( + session["access_token"], peer_id, message_id + ) cursor.last_message_id = str(max_message_id) cursor.initialized = True diff --git a/digital-avatar-app/backend/tests/test_boxim_client.py b/digital-avatar-app/backend/tests/test_boxim_client.py index 6a65b23..31ef6cb 100644 --- a/digital-avatar-app/backend/tests/test_boxim_client.py +++ b/digital-avatar-app/backend/tests/test_boxim_client.py @@ -96,6 +96,18 @@ async def test_send_private_message_matches_boxim_payload(config): } +@pytest.mark.asyncio +async def test_mark_private_messages_read_uses_latest_message_id(config): + mocked, client = _client_patch(request_payload={"code": 200, "data": None}) + with mocked: + await BoxIMClient(config).mark_private_messages_read("box-token", "77", "101") + + call = client.request.await_args + assert call.args[:2] == ("PUT", "https://im.example/api/message/private/readed") + assert call.kwargs["headers"] == {"accessToken": "box-token"} + assert call.kwargs["params"] == {"friendId": 77, "messageId": 101} + + @pytest.mark.asyncio async def test_boxim_auth_error_is_explicit(config): mocked, _ = _client_patch( diff --git a/digital-avatar-app/backend/tests/test_takeover_service.py b/digital-avatar-app/backend/tests/test_takeover_service.py index 4260f7d..2786a25 100644 --- a/digital-avatar-app/backend/tests/test_takeover_service.py +++ b/digital-avatar-app/backend/tests/test_takeover_service.py @@ -32,6 +32,7 @@ class FakeBoxIM: def __init__(self): self.messages = [] self.sent = [] + self.read_receipts = [] async def exchange_access_token(self, huihui_token): assert huihui_token == "prod-huihui-token" @@ -45,6 +46,12 @@ class FakeBoxIM: assert access_token == "box-token" return [item.copy() for item in self.messages if int(item["id"]) > int(min_id)] + async def mark_private_messages_read(self, access_token, friend_id, message_id): + assert access_token == "box-token" + self.read_receipts.append( + {"friendId": str(friend_id), "messageId": str(message_id)} + ) + async def send_private_message(self, access_token, peer_id, content, *, local_id=None): self.sent.append({"peerId": str(peer_id), "content": content, "localId": str(local_id)}) return {"id": 900 + len(self.sent), "localId": int(local_id)} @@ -101,6 +108,7 @@ async def test_first_sync_primes_cursor_without_replying_to_history(service_cont assert db.query(TakeoverMessage).count() == 1 assert db.query(TakeoverReplyTask).count() == 0 assert boxim.sent == [] + assert boxim.read_receipts == [{"friendId": "200", "messageId": "10"}] finally: db.close() @@ -116,6 +124,7 @@ async def test_incoming_message_is_prepared_then_sent_at_three_seconds(service_c with patch("routers.chat._resolve_reply", return_value={"answer": "**你好**\n\n很高兴见到你"}): await service.poll_and_process_messages() assert boxim.sent == [] + assert boxim.read_receipts == [{"friendId": "200", "messageId": "11"}] clock.advance(2) await service.poll_and_process_messages() @@ -134,6 +143,40 @@ async def test_incoming_message_is_prepared_then_sent_at_three_seconds(service_c db.close() +@pytest.mark.asyncio +async def test_read_receipt_failure_does_not_advance_cursor(service_context): + session_factory, service, boxim, clock = service_context + await service.poll_and_process_messages() + boxim.messages.append( + {"id": 12, "localId": 3, "sendId": 200, "recvId": 100, "sendTime": clock.millis(), "type": 0, "content": "未读消息"} + ) + boxim.mark_private_messages_read = AsyncMock(side_effect=BoxIMError("回执失败")) + + with patch("routers.chat._resolve_reply", return_value={"answer": "稍后回复"}): + await service.poll_and_process_messages() + + db = session_factory() + try: + cursor = db.query(TakeoverCursor).one() + assert cursor.last_message_id == "0" + assert db.query(TakeoverMessage).count() == 0 + assert db.query(TakeoverReplyTask).count() == 0 + finally: + db.close() + + boxim.mark_private_messages_read = AsyncMock(return_value=None) + with patch("routers.chat._resolve_reply", return_value={"answer": "稍后回复"}): + await service.poll_and_process_messages() + + db = session_factory() + try: + assert db.query(TakeoverCursor).one().last_message_id == "12" + assert db.query(TakeoverMessage).count() == 1 + assert db.query(TakeoverReplyTask).count() == 1 + finally: + db.close() + + @pytest.mark.asyncio async def test_owner_message_cancels_pending_reply(service_context): session_factory, service, boxim, clock = service_context