fix(avatar): acknowledge BOXIM messages as read
This commit is contained in:
@@ -162,6 +162,27 @@ class BoxIMClient:
|
|||||||
raise BoxIMError("BOXIM 私聊消息格式不正确")
|
raise BoxIMError("BOXIM 私聊消息格式不正确")
|
||||||
return [item for item in data if isinstance(item, dict)]
|
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(
|
async def send_private_message(
|
||||||
self,
|
self,
|
||||||
access_token: str,
|
access_token: str,
|
||||||
|
|||||||
@@ -268,6 +268,7 @@ class TakeoverService:
|
|||||||
messages.sort(key=lambda item: (_numeric_id(item.get("id")), item.get("sendTime") or 0))
|
messages.sort(key=lambda item: (_numeric_id(item.get("id")), item.get("sendTime") or 0))
|
||||||
priming = not bool(cursor.initialized)
|
priming = not bool(cursor.initialized)
|
||||||
max_message_id = _numeric_id(cursor.last_message_id)
|
max_message_id = _numeric_id(cursor.last_message_id)
|
||||||
|
read_receipts: dict[str, int] = {}
|
||||||
for message in messages:
|
for message in messages:
|
||||||
self._record_message(
|
self._record_message(
|
||||||
db,
|
db,
|
||||||
@@ -276,7 +277,19 @@ class TakeoverService:
|
|||||||
message,
|
message,
|
||||||
schedule_reply=not priming,
|
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.last_message_id = str(max_message_id)
|
||||||
cursor.initialized = True
|
cursor.initialized = True
|
||||||
|
|||||||
@@ -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
|
@pytest.mark.asyncio
|
||||||
async def test_boxim_auth_error_is_explicit(config):
|
async def test_boxim_auth_error_is_explicit(config):
|
||||||
mocked, _ = _client_patch(
|
mocked, _ = _client_patch(
|
||||||
|
|||||||
@@ -32,6 +32,7 @@ class FakeBoxIM:
|
|||||||
def __init__(self):
|
def __init__(self):
|
||||||
self.messages = []
|
self.messages = []
|
||||||
self.sent = []
|
self.sent = []
|
||||||
|
self.read_receipts = []
|
||||||
|
|
||||||
async def exchange_access_token(self, huihui_token):
|
async def exchange_access_token(self, huihui_token):
|
||||||
assert huihui_token == "prod-huihui-token"
|
assert huihui_token == "prod-huihui-token"
|
||||||
@@ -45,6 +46,12 @@ class FakeBoxIM:
|
|||||||
assert access_token == "box-token"
|
assert access_token == "box-token"
|
||||||
return [item.copy() for item in self.messages if int(item["id"]) > int(min_id)]
|
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):
|
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)})
|
self.sent.append({"peerId": str(peer_id), "content": content, "localId": str(local_id)})
|
||||||
return {"id": 900 + len(self.sent), "localId": int(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(TakeoverMessage).count() == 1
|
||||||
assert db.query(TakeoverReplyTask).count() == 0
|
assert db.query(TakeoverReplyTask).count() == 0
|
||||||
assert boxim.sent == []
|
assert boxim.sent == []
|
||||||
|
assert boxim.read_receipts == [{"friendId": "200", "messageId": "10"}]
|
||||||
finally:
|
finally:
|
||||||
db.close()
|
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很高兴见到你"}):
|
with patch("routers.chat._resolve_reply", return_value={"answer": "**你好**\n\n很高兴见到你"}):
|
||||||
await service.poll_and_process_messages()
|
await service.poll_and_process_messages()
|
||||||
assert boxim.sent == []
|
assert boxim.sent == []
|
||||||
|
assert boxim.read_receipts == [{"friendId": "200", "messageId": "11"}]
|
||||||
|
|
||||||
clock.advance(2)
|
clock.advance(2)
|
||||||
await service.poll_and_process_messages()
|
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()
|
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
|
@pytest.mark.asyncio
|
||||||
async def test_owner_message_cancels_pending_reply(service_context):
|
async def test_owner_message_cancels_pending_reply(service_context):
|
||||||
session_factory, service, boxim, clock = service_context
|
session_factory, service, boxim, clock = service_context
|
||||||
|
|||||||
Reference in New Issue
Block a user