群里按 BOT 通道激活当前工单、航线报价和成交跟进;没有 TMS 报价编号时用工单字段补齐仍调锁舱,避免被本地规则拦住。 Co-authored-by: Cursor <cursoragent@cursor.com>
346 lines
12 KiB
Python
346 lines
12 KiB
Python
"""
|
|
微盛会话存档:解析、跳过存档席、测服游标、入 inbox。
|
|
|
|
不连微盛外网、不连现网 Redis。
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import os
|
|
import sys
|
|
import unittest
|
|
|
|
ROOT = os.path.abspath(os.path.join(os.path.dirname(__file__), ".."))
|
|
if ROOT not in sys.path:
|
|
sys.path.insert(0, ROOT)
|
|
|
|
os.environ.setdefault("YTD_ENV", "test")
|
|
os.environ.setdefault("REDIS_BACKEND", "memory")
|
|
os.environ.setdefault("CHECKPOINT_BACKEND", "memory")
|
|
os.environ.setdefault("MESSAGE_STORE_BACKEND", "memory")
|
|
os.environ.setdefault("LLM_DATA_USAGE_CONFIRMED", "false")
|
|
os.environ.setdefault("LLM_ALLOW_NETWORK", "false")
|
|
os.environ.setdefault("LEDGER_BACKEND", "memory")
|
|
|
|
from agent.channel.queue import MemoryMessageStore
|
|
from agent.channel.wshoto.client import wshoto_ready
|
|
from agent.channel.wshoto.parse import archive_record_to_inbound
|
|
from agent.channel.wshoto.rooms import (
|
|
list_watched_rooms,
|
|
load_last_seq,
|
|
save_last_seq,
|
|
unwatch_collab_room,
|
|
watch_collab_room,
|
|
)
|
|
from agent.channel.wshoto.poll import ingest_records, poll_once
|
|
from agent.redis_coord.client import MemoryRedis, RedisClient
|
|
|
|
|
|
def _redis() -> RedisClient:
|
|
return RedisClient(raw=MemoryRedis(), key_prefix="inquiry_robot:", backend="memory")
|
|
|
|
|
|
class _Cfg:
|
|
wshoto_poll_enabled = True
|
|
wshoto_app_id = "app"
|
|
wshoto_app_secret = "secret"
|
|
wecom_archive_seat_user_ids = "AI001"
|
|
wshoto_bootstrap_lookback_sec = 3600
|
|
wshoto_page_size = 50
|
|
|
|
|
|
class _FakeApi:
|
|
def __init__(self, records: list[dict], *, error: WshotoError | None = None) -> None:
|
|
self.records = records
|
|
self.error = error
|
|
self.calls: list[dict] = []
|
|
|
|
def list_group_messages(self, **kwargs):
|
|
self.calls.append(kwargs)
|
|
if self.error is not None:
|
|
raise self.error
|
|
last = int(kwargs.get("last_seq") or 0)
|
|
return [r for r in self.records if int(r.get("seq") or 0) > last]
|
|
|
|
|
|
class WshotoParseTests(unittest.TestCase):
|
|
def test_member_text_becomes_group_inbound(self) -> None:
|
|
msg = archive_record_to_inbound(
|
|
{
|
|
"msgid": "m1",
|
|
"seq": 10,
|
|
"from": "WuJiLin",
|
|
"room_id": "wr_room_1",
|
|
"msgtype": "text",
|
|
"content": "运费 1200",
|
|
"msgtime": 1_700_000_000_000,
|
|
},
|
|
seat_userids=["AI001"],
|
|
)
|
|
self.assertIsNotNone(msg)
|
|
assert msg is not None
|
|
self.assertEqual(msg.sender_id, "WuJiLin")
|
|
self.assertEqual(msg.chat_id, "wr_room_1")
|
|
self.assertEqual(msg.chat_type, "group")
|
|
self.assertEqual(msg.content, "运费 1200")
|
|
self.assertEqual(msg.message_id, "wshoto:m1")
|
|
self.assertEqual(msg.create_time, 1_700_000_000)
|
|
|
|
def test_official_text_content_beats_short_preview(self) -> None:
|
|
"""管家 content 常是摘要或空,正文在企微官方 text.content。必须取更长的那份。"""
|
|
msg = archive_record_to_inbound(
|
|
{
|
|
"msgid": "m1b",
|
|
"seq": 11,
|
|
"from": "LuoXiaoHua",
|
|
"room_id": "wr_room_1",
|
|
"msgtype": "text",
|
|
"content": "码头操作费",
|
|
"msg_content": {
|
|
"text": {"content": "码头操作费100CNY\n报关费改为300CNY"}
|
|
},
|
|
},
|
|
seat_userids=["AI001"],
|
|
)
|
|
self.assertIsNotNone(msg)
|
|
assert msg is not None
|
|
self.assertEqual(msg.content, "码头操作费100CNY\n报关费改为300CNY")
|
|
|
|
def test_skip_archive_seat_and_display_name(self) -> None:
|
|
self.assertIsNone(
|
|
archive_record_to_inbound(
|
|
{
|
|
"msgid": "m2",
|
|
"from": "AI001",
|
|
"room_id": "wr_room_1",
|
|
"msgtype": "text",
|
|
"content": "群简介",
|
|
},
|
|
seat_userids=["AI001"],
|
|
)
|
|
)
|
|
self.assertIsNone(
|
|
archive_record_to_inbound(
|
|
{
|
|
"msgid": "m3",
|
|
"from": "询价机器人",
|
|
"room_id": "wr_room_1",
|
|
"msgtype": "text",
|
|
"content": "系统话术",
|
|
},
|
|
seat_userids=["AI001"],
|
|
)
|
|
)
|
|
|
|
def test_file_keeps_filename(self) -> None:
|
|
msg = archive_record_to_inbound(
|
|
{
|
|
"msgid": "m4",
|
|
"from": "LuoXiaoHua",
|
|
"room_id": "wr_room_1",
|
|
"msgtype": "file",
|
|
"content": "",
|
|
"msg_content": {"filename": "报价.xlsx", "sdkfileid": "sdk-1"},
|
|
},
|
|
seat_userids=["AI001"],
|
|
)
|
|
self.assertIsNotNone(msg)
|
|
assert msg is not None
|
|
self.assertEqual(msg.msg_type, "file")
|
|
self.assertEqual(msg.media.get("filename"), "报价.xlsx")
|
|
self.assertEqual(msg.media.get("media_id"), "sdk-1")
|
|
|
|
def test_skip_revoke_and_missing_room(self) -> None:
|
|
self.assertIsNone(
|
|
archive_record_to_inbound(
|
|
{"msgid": "m5", "from": "WuJiLin", "room_id": "", "msgtype": "text", "content": "x"},
|
|
seat_userids=["AI001"],
|
|
)
|
|
)
|
|
self.assertIsNone(
|
|
archive_record_to_inbound(
|
|
{
|
|
"msgid": "m6",
|
|
"from": "WuJiLin",
|
|
"room_id": "wr_room_1",
|
|
"msgtype": "revoke",
|
|
"content": "x",
|
|
},
|
|
seat_userids=["AI001"],
|
|
)
|
|
)
|
|
|
|
|
|
class WshotoPollTests(unittest.TestCase):
|
|
def test_ready_needs_keys(self) -> None:
|
|
cfg = _Cfg()
|
|
self.assertTrue(wshoto_ready(cfg))
|
|
cfg.wshoto_app_secret = ""
|
|
self.assertFalse(wshoto_ready(cfg))
|
|
|
|
def test_watch_rooms_and_cursor(self) -> None:
|
|
redis = _redis()
|
|
watch_collab_room("wr_a", redis=redis)
|
|
watch_collab_room("wr_b", redis=redis)
|
|
self.assertEqual(list_watched_rooms(redis=redis), ["wr_a", "wr_b"])
|
|
unwatch_collab_room("wr_a", redis=redis)
|
|
self.assertEqual(list_watched_rooms(redis=redis), ["wr_b"])
|
|
save_last_seq("seat:AI001", 88, redis=redis)
|
|
self.assertEqual(load_last_seq("seat:AI001", redis=redis), 88)
|
|
watch_collab_room("wr_c", work_order_no="WO1", redis=redis)
|
|
from agent.channel.wshoto.rooms import ticket_no_for_room
|
|
|
|
self.assertEqual(ticket_no_for_room("wr_c", redis=redis), "WO1")
|
|
|
|
def test_watch_room_enters_hot_window(self) -> None:
|
|
from agent.channel.wshoto.rooms import has_hot_room, wait_wshoto_poll
|
|
|
|
redis = _redis()
|
|
watch_collab_room("wr_hot", redis=redis)
|
|
self.assertTrue(has_hot_room())
|
|
self.assertTrue(wait_wshoto_poll(0.2))
|
|
|
|
def test_new_room_lookback_is_capped(self) -> None:
|
|
store = MemoryMessageStore()
|
|
redis = _redis()
|
|
watch_collab_room("wr_room_1", redis=redis)
|
|
api = _FakeApi(
|
|
[
|
|
{
|
|
"msgid": "fresh",
|
|
"seq": 7,
|
|
"from": "WuJiLin",
|
|
"room_id": "wr_room_1",
|
|
"msgtype": "text",
|
|
"content": "随时可提",
|
|
}
|
|
]
|
|
)
|
|
out = poll_once(
|
|
client=api,
|
|
settings=_Cfg(),
|
|
redis=redis,
|
|
store=store,
|
|
now=lambda: 1_800_000_000,
|
|
)
|
|
self.assertTrue(out["ok"])
|
|
call = api.calls[0]
|
|
self.assertEqual(call.get("msg_time_end"), 1_800_000_000)
|
|
self.assertEqual(call.get("msg_time_start"), 1_800_000_000 - 180)
|
|
|
|
def test_poll_enqueues_member_skips_seat_and_advances_seq(self) -> None:
|
|
store = MemoryMessageStore()
|
|
redis = _redis()
|
|
api = _FakeApi(
|
|
[
|
|
{
|
|
"msgid": "old",
|
|
"seq": 10,
|
|
"from": "WuJiLin",
|
|
"room_id": "wr_room_1",
|
|
"msgtype": "text",
|
|
"content": "旧消息",
|
|
},
|
|
{
|
|
"msgid": "seat",
|
|
"seq": 11,
|
|
"from": "AI001",
|
|
"room_id": "wr_room_1",
|
|
"msgtype": "text",
|
|
"content": "机器人自己说的",
|
|
},
|
|
{
|
|
"msgid": "new",
|
|
"seq": 12,
|
|
"from": "WuJiLin",
|
|
"room_id": "wr_room_1",
|
|
"msgtype": "text",
|
|
"content": "运费 1500",
|
|
},
|
|
]
|
|
)
|
|
save_last_seq("seat:AI001", 10, redis=redis)
|
|
out = poll_once(client=api, settings=_Cfg(), redis=redis, store=store)
|
|
self.assertTrue(out["ok"])
|
|
self.assertEqual(out["enqueued"], 1)
|
|
self.assertEqual(out["last_seq"], 12)
|
|
self.assertEqual(load_last_seq("seat:AI001", redis=redis), 12)
|
|
item = store.claim_inbound()
|
|
self.assertIsNotNone(item)
|
|
assert item is not None
|
|
self.assertEqual(item.message.content, "运费 1500")
|
|
self.assertEqual(item.message.chat_type, "group")
|
|
|
|
def test_first_poll_uses_time_window_not_seq_zero(self) -> None:
|
|
store = MemoryMessageStore()
|
|
redis = _redis()
|
|
api = _FakeApi(
|
|
[
|
|
{
|
|
"msgid": "recent",
|
|
"seq": 5,
|
|
"from": "WuJiLin",
|
|
"room_id": "wr_room_1",
|
|
"msgtype": "text",
|
|
"content": "刚说的",
|
|
}
|
|
]
|
|
)
|
|
out = poll_once(
|
|
client=api,
|
|
settings=_Cfg(),
|
|
redis=redis,
|
|
store=store,
|
|
now=lambda: 1_800_000_000,
|
|
)
|
|
self.assertTrue(out["ok"])
|
|
self.assertEqual(len(api.calls), 1)
|
|
call = api.calls[0]
|
|
self.assertEqual(call.get("last_seq"), 0)
|
|
self.assertGreater(call.get("msg_time_start") or 0, 0)
|
|
self.assertEqual(call.get("msg_time_end"), 1_800_000_000)
|
|
self.assertEqual(out["enqueued"], 1)
|
|
|
|
def test_ingest_dedupes_message_id(self) -> None:
|
|
store = MemoryMessageStore()
|
|
rec = {
|
|
"msgid": "same",
|
|
"seq": 3,
|
|
"from": "WuJiLin",
|
|
"room_id": "wr_room_1",
|
|
"msgtype": "text",
|
|
"content": "重复",
|
|
}
|
|
a, _ = ingest_records([rec], seat_userids=["AI001"], store=store)
|
|
b, skipped = ingest_records([rec], seat_userids=["AI001"], store=store)
|
|
self.assertEqual(a, 1)
|
|
self.assertEqual(b, 0)
|
|
self.assertEqual(skipped, 1)
|
|
|
|
def test_watched_rooms_preferred_without_seat_userid(self) -> None:
|
|
store = MemoryMessageStore()
|
|
redis = _redis()
|
|
watch_collab_room("wr_room_1", redis=redis)
|
|
api = _FakeApi(
|
|
[
|
|
{
|
|
"msgid": "room-msg",
|
|
"seq": 21,
|
|
"from": "LuoXiaoHua",
|
|
"room_id": "wr_room_1",
|
|
"msgtype": "text",
|
|
"content": "20GP 1",
|
|
}
|
|
]
|
|
)
|
|
out = poll_once(client=api, settings=_Cfg(), redis=redis, store=store)
|
|
self.assertTrue(out["ok"])
|
|
self.assertEqual(out["mode"], "rooms")
|
|
self.assertEqual(out["enqueued"], 1)
|
|
self.assertEqual(api.calls[0].get("userid"), "")
|
|
self.assertEqual(api.calls[0].get("target_id"), "wr_room_1")
|
|
|
|
|
|
if __name__ == "__main__":
|
|
unittest.main()
|