Files
2026-06-25 17:41:06 +08:00

409 lines
18 KiB
Python

from tortoise.expressions import Q
from app.core.crud import CRUDBase
from app.schemas.weixin import (
WeixinGroupChatCreate,
WeixinGroupChatUpdate,
WeixinGroupChatXingyunCreate,
)
from some_sdk.services.binder import wk_client, lintao_client, xy_client
from some_sdk.wk_weixin_sdk.apis import corp_group
from some_sdk.wk_weixin_sdk.apis.extern_user import get_external_user_chat_info
from some_sdk.xingyun_sdk.apis.work_user import list_work_group
from some_sdk.xingyun_sdk.apis.customer import list_user_join_group
from some_sdk.lintao_sdk.biz.by_order import get_order_relative_user
from some_sdk.lintao_sdk.apis.orderlist import list_product_raw, get_order_log
from app.models.weixin import WeixinGroupChat, RoleType, CustomerGroup
from app.controllers.weixin.user import weixin_user_controller
from app.controllers.weixin.customer import weixin_customer_controller
from app.utils.common import async_split_generator, async_generator_to_list
from app.core.cache import cache_if, invalidate_cache
import re
import logging
logger = logging.getLogger(__name__)
ROLE_LIST = [role.value for role in RoleType]
class WeixinGroupChatController(CRUDBase[WeixinGroupChat, WeixinGroupChatCreate, WeixinGroupChatUpdate]):
def __init__(self):
super().__init__(model=WeixinGroupChat)
async def get_external_group_chat_info(self, chat_id: str):
group_info = await corp_group.get_external_group_chat_info(wk_client, chat_id)
internal_member_list = group_info['internal_member_list']
external_member_list = group_info['external_member_list']
group_name = group_info.get('name')
external_member_list.sort(key=lambda user: user.get('join_time'))
unique_dict = {
'name': group_name,
'external_user_count': len(external_member_list),
'internal_member_count': len(internal_member_list),
}
group_info.update(unique_dict)
# 同步更新数据库
status, obj = await self.create_or_update(group_info, query_kwargs={'chat_id': chat_id}, update_kwargs=unique_dict)
if not status:
logger.debug(f'群聊数据不变,不用更新,chat_id: {chat_id}')
# elif status == 'create':
# await self.create_customer_group(chat_id)
else:
# 群更新:触发数据同步事件
# event_manager.subscribe("group_chat_updated", update_cache_handler, EventPriority.MEDIUM)
# 只保存准确无误的数据
# buyer_nick = get_buyer_nick_from_group_name(group_name)
# logger.debug(f'group_name: {group_name}; buyer_nick: {buyer_nick}')
# if not buyer_nick:
# logger.error(f'群聊未命名,无法绑定订单,chat_id: {chat_id}')
# return group_info
# self.load_erp_order_info_by_group(buyer_nick, internal_member_list, external_member_list)
update_result1 = await weixin_user_controller.update_from_group([*internal_member_list])
update_result2 = await weixin_customer_controller.update_from_group([*external_member_list])
logger.debug(f'群聊用户数据更新结果,chat_id: {chat_id}, {update_result1}')
logger.debug(f'群聊客户数据更新结果,chat_id: {chat_id}, {update_result2}')
# 刷新用户缓存
for user in [*internal_member_list, *external_member_list]:
await self.expire_user_cache(user.get('userid'))
group_info['is_create'] = status == 'create'
group_info['detect_update'] = status != False
return group_info
# 从星云有客中加载群聊信息
async def load_xingyun_group_info(self, chat_name: str=''):
# 检索出最新的记录,只加载最新的记录
newest_group = await self.model.filter(xingyun_chat_id__isnull=False).order_by('-create_time').first()
newest_time = newest_group.create_time if newest_group else 1735660800 # 2025-01-01 00:00:00
load_result = {"update_count": 0, "add_count": 0}
group_iter = list_work_group(xy_client, chatName=chat_name)
async for group_list in async_split_generator(group_iter):
if group_list[-1].get('createTime') < newest_time:
break
logger.debug(f'处理群聊信息,chat_name: {chat_name}, {len(group_list)}')
result = await self.save_xingyun_group_info(group_list)
load_result['update_count'] += result['update_count']
load_result['add_count'] += result['add_count']
logger.debug(f'群聊数据更新结果,chat_name: {chat_name}, {result}')
if load_result['add_count'] + load_result['update_count'] > 300:
break
await self.refresh_xingyun_group_info()
return load_result
async def save_xingyun_group_info(self, group_list: list):
group_dict = {group.get('groupChatId'): group for group in group_list}
group_orm_list = await self.model.filter(chat_id__in=group_dict.keys())
update_list = []
history = {}
for group_orm in group_orm_list:
group = group_dict.pop(group_orm.chat_id, history.get(group_orm.chat_id, {}))
if not group: continue
history[group_orm.chat_id] = group
try:
group_in = WeixinGroupChatXingyunCreate(**group)
except Exception as e:
logger.error(f'群聊数据转换错误,chat_id: {group}, {e}')
continue
updated = False
if group_in.avatars:
group_orm.avatars = group_in.avatars
updated = True
if group_orm.xingyun_chat_id != group_in.xingyun_chat_id:
group_orm.xingyun_chat_id = group_in.xingyun_chat_id
updated = True
if updated:
update_list.append(group_orm)
if update_list:
await self.model.bulk_update(update_list, ['avatars', 'xingyun_chat_id'])
add_list = []
for chat_id, group in group_dict.items():
group_in = WeixinGroupChatXingyunCreate(**group)
add_list.append(self.model(**group_in.model_dump(exclude_unset=True)))
if add_list:
await self.model.bulk_create(add_list)
return {
'update_count': len(update_list),
'add_count': len(add_list),
}
# 刷新数据库中的群聊信息
async def refresh_xingyun_group_info(self):
for _ in range(10):
group_orm_list = await self.model.filter(xingyun_chat_id__isnull=False, external_user_count__isnull=True).order_by('-create_time').limit(20).all()
if not group_orm_list: break
logger.debug(f'刷新群聊数据,{len(group_orm_list)}')
for group_orm in group_orm_list:
try:
await self.get_external_group_chat_info(group_orm.chat_id)
except Exception as e:
logger.error(f'刷新群聊数据错误,chat_id: {group_orm.chat_id}, {e}')
continue
logger.debug(f'刷新群聊数据结果,chat_id: {group_orm.name} 已更新')
# 同步erp中的用户和客户信息到数据库
async def load_erp_order_info_by_group(self, buyer_nick: str, internal_member_list: list, external_member_list: list):
if not buyer_nick:
logger.error(f'群聊未命名,无法绑定订单,buyer_nick: {buyer_nick}')
return
orders = await async_generator_to_list(get_order_relative_user(lintao_client, buyer_nick=buyer_nick))
assert orders, f'未找到订单,buyer_nick: {buyer_nick}'
internal_member_dict = {user.get('name'): user for user in internal_member_list}
order = orders[0]
order_relative_user = order.get("users", [])
order_relative_user_dict = {user.get('name'): user for user in order_relative_user}
for name, user in order_relative_user_dict.items():
erp_id = user.get('id', '')
if not erp_id: continue
if re.match(r'^\d+$', erp_id) is None:
# 客户
if not external_member_list:
logger.error(f'客户群聊用户数为0,无法绑定客户, name: {name}, erp_id: {erp_id}')
continue
if name != buyer_nick:
logger.error(f'客户群聊用户与ERP卖家名称不一致,无法绑定客户, name: {name}, erp_id: {erp_id}')
continue
if len(external_member_list) != 1:
external_member_list.sort(key=lambda user: user.get('join_time'))
logger.warning(f'客户群聊用户数不是1个,将选取最先入群的用户作为客户,name: {name}, erp_id: {erp_id}')
for member in external_member_list:
member['order_id'] = order.get('trade_no')
weixin_user = external_member_list[0]
weixin_user.update({
'order_id': order.get('trade_no'),
'taobao_id': user.get('id', ''),
'taobao_name': user.get('name', ''),
'erp_message': user.get('message', ''),
'erp_remark': user.get('memo', ''),
'need_confirm': len(external_member_list) > 1,
})
continue
else:
# 员工
erp_id = int(erp_id)
weixin_user = internal_member_dict.get(name)
logger.debug(f'name: {name}, erp_id: {erp_id}, weixin_user: {weixin_user}')
if not weixin_user:
logger.error(f'{internal_member_dict.keys()} {name}')
logger.error(f'客户群聊用户与ERP员工名称不一致,无法绑定员工, name: {name}, erp_id: {erp_id}')
continue
weixin_user.update({
'erp_id': erp_id,
'erp_name': name,
'role': user.get('role', ''),
})
continue
# 通过客户ID获取客户的所有订单
async def get_order_relative_user_list_by_weixin_userid(self, user_id: str):
customers = await weixin_customer_controller.model.filter(weixin_id=user_id)
if not customers:
await weixin_customer_controller.create({
"weixin_id": user_id,
})
# 新用户:触发数据同步事件
return False, []
order_list = []
buyer_set = set()
for customer in customers:
buyer_nick = customer.taobao_name
buyer_id = customer.taobao_id
if buyer_nick:
logger.debug(f'name: {user_id}; buyer_nick: {buyer_nick}')
if buyer_nick in buyer_set: continue
buyer_set.add(buyer_nick)
orders = await async_generator_to_list(get_order_relative_user(lintao_client, buyer_nick=buyer_nick))
elif buyer_id:
if buyer_id in buyer_set: continue
buyer_set.add(buyer_id)
orders = await async_generator_to_list(get_order_relative_user(lintao_client, buyer_id=buyer_id))
else:
assert buyer_nick, "当前客户未绑定淘宝账号"
logger.debug(f'orders.length: {len(orders)}')
order_list.extend(orders)
# for order in orders:
# 原地按照id排序
order_list.sort(key=lambda order: order.get('id'), reverse=True)
if order_list:
await self.get_user_detail_by_order(order=order_list[0])
return True, order_list
# 根据客户名字,获取群聊中客户的所有订单
async def get_order_relative_user_list_by_weixin_group_name(self, buyer_nick: str):
async with cache_if(f'order:relative_user:{buyer_nick}', ttl=10) as cache:
if cache.hit:
order_list = cache.value
logger.debug(f'{cache.key} hit')
else:
logger.debug(f'{cache.key} not hit')
order_list = await async_generator_to_list(get_order_relative_user(lintao_client, buyer_nick=buyer_nick))
cache.set(order_list)
logger.debug(f'order_list.length: {len(order_list)}')
# for order in order_list:
order_list.sort(key=lambda order: order.get('id'), reverse=True)
if order_list:
await self.get_user_detail_by_order(order=order_list[0])
order_list[0]['__type__'] = 'load_order_by_weixin_group_name'
return order_list
# 根据客户ID,获取群聊中客户的所有订单
async def get_order_relative_user_list_by_weixin_group_buyer_id(self, buyer_id: str):
orders = await async_generator_to_list(get_order_relative_user(lintao_client, buyer_id=buyer_id))
logger.debug(f'orders.length: {len(orders)}')
# for order in orders:
if orders:
await self.get_user_detail_by_order(order=orders[0])
return orders
async def get_user_detail_by_order(self, order_id: str='', order: dict = None):
if not order:
orders = await async_generator_to_list(get_order_relative_user(lintao_client, trade_no=order_id))
if not orders: return
_orders = [o for o in orders if o.get('ctid') == order_id]
# assert len(orders) == 1, f'订单号 {order_id} 对应多个订单'
if not _orders:
for order in orders:
logger.warning(f'订单号 {order_id} 对应多个订单,{order}')
# return {}
_orders = orders
order = _orders[0]
# 缓存订单产品信息
tid = order.get("ctid")
async with cache_if(f'order:product:{tid}', ttl=3600*24) as cache:
if cache.hit:
product_resp = cache.value
else:
product_resp = (await list_product_raw(lintao_client, tid=tid)).get('data', {})
if product_resp: cache.set(product_resp)
if product_resp:
order['product_title'] = product_resp[0].get("title")
order['product_pic'] = product_resp[0].get("pic_path")
order['has_fetched_user_detail'] = True
user_list = order.get("users", [])
async with cache_if(f'user:auto_relate', ttl=600) as cache:
if cache.hit:
default_user_list = cache.value
else:
default_user_list = await weixin_user_controller.all(Q(auto_relate=True))
default_user_list = [user.to_dict() for user in default_user_list]
cache.set(default_user_list)
default_user_set = {str(user.get('erp_id')) for user in default_user_list}
for user in user_list:
erp_id = user.get('id', '')
if erp_id in default_user_set: continue
await self.get_order_user_info(user)
order['users'].extend(default_user_list)
# 根据 role_list 中的顺序进行排序
role_order = {role: index for index, role in enumerate(ROLE_LIST)}
order['users'].sort(key=lambda user: role_order.get(user['role'], len(ROLE_LIST)))
return order
async def get_order_user_info(self, user: dict=None):
if not user: return
erp_id = user.get('id', '')
if not erp_id: return
if re.match(r'^\d+$', erp_id) is None:
# logger.debug(f'淘宝客户, erp_id: {erp_id}, user: {user}')
# 客户
taobao_id = user.get('id', '')
query = Q(taobao_id=taobao_id)
weixin_customer_list = await weixin_customer_controller.all(query)
if not weixin_customer_list:
logger.error(f'客户未绑定订单, taobao_id: {taobao_id}, user: {user}')
return
weixin_customer = weixin_customer_list[0]
buyer = weixin_customer.to_dict()
if not buyer.get('name') and user.get('name'):
weixin_customer.taobao_name = user.get('name')
await weixin_customer.save()
user.update({**buyer, **user})
else:
# 员工
erp_id = int(erp_id)
# logger.debug(f'员工, erp_id: {erp_id}, user: {user}')
query = Q(erp_id=erp_id)
weixin_user_list = await weixin_user_controller.all(query)
if not weixin_user_list:
user['erp_id'] = erp_id
logger.error(f'员工未绑定ERP, erp_id: {erp_id}, user: {user}')
return
weixin_user = weixin_user_list[0]
user.update(weixin_user.to_dict())
# 获取erp日志
async def get_erp_order_log(self, order_id: str):
return await get_order_log(lintao_client, order_id)
async def expire_user_cache(self, weixin_id: str):
await invalidate_cache(f'customer:user_datail:{weixin_id}')
await invalidate_cache(f'customer:user_datail_0:{weixin_id}')
# async def list_user_join_group(self, xingyun_id: int):
# assert xingyun_id, 'id 不能为空'
# result = await async_generator_to_list(list_user_join_group(xy_client, cid=xingyun_id))
# xingyun_chat_ids = [item.get('id') for item in result]
# all_group_chat = await self.model.filter(xingyun_chat_id__in=xingyun_chat_ids)
# group_chat_dict = {item.xingyun_chat_id: item for item in all_group_chat}
# logger.debug(f'list_user_join_group, xingyun_chat_ids: {xingyun_chat_ids}')
# remain_chat_ids = set()
# for item in result:
# group_chat_id = item.get('id')
# if group_chat_id in group_chat_dict: continue
# remain_chat_ids.add(group_chat_id)
# chat_name = item.get('chatName')
# if not chat_name: continue
# await self.load_xingyun_group_info(chat_name=chat_name)
# remain_chat_list = await self.model.filter(xingyun_chat_id__in=remain_chat_ids)
# all_group_chat.extend(remain_chat_list)
# return all_group_chat
weixin_group_chat_controller = WeixinGroupChatController()