from .customer import weixin_customer_controller from .user import weixin_user_controller from app.controllers.msg import msg_controller from app.controllers.action import action_controller from app.schemas.weixin import WeixinGroupChatEvent, BindOrderResultEvent, CustomerRepeatPurchaseEvent, CustomerAssignOrderEvent, CustomerRefundOrderEvent from app.utils.common import transform_pydantic_to_list from datetime import datetime, timedelta from app.utils.event_task import event_manager, EventType import os import asyncio from typing import Any from apscheduler.schedulers.asyncio import AsyncIOScheduler as BackgroundScheduler from app.controllers.automation.scenario import automation_scenario_controller from app.controllers.automation.task import task_controller sync_lock = asyncio.Lock() import logging logger = logging.getLogger(__name__) async def sync_xingyun_contact_info_async(event_type, *args, **kwargs): if os.getenv("APP_ENV") != "prod": logger.info(f'非生产环境,不实际同步客户数据,{args} {kwargs}') return logger.info(f'同步客户数据,{args} {kwargs}') now_time = datetime.now() start_time = (now_time - timedelta(days=1)).replace(hour=0, minute=0, second=0, microsecond=0) end_time = now_time + timedelta(days=1) await weixin_customer_controller.load_user_from_xingyun(add_time_start=start_time, add_time_end=end_time) async def load_all_user_from_xingyun(event_type, *args, **kwargs): if os.getenv("APP_ENV") != "prod": logger.info(f'非生产环境,不实际从星云加载全量客户数据,{args} {kwargs}') return logger.info(f'从星云加载全量客户数据,{args} {kwargs}') start_time = datetime(2025, 6, 1) end_time = datetime.now() + timedelta(days=1) return await weixin_customer_controller.load_all_user_from_xingyun(task_id='load_all_user_from_xingyun', add_time_start=start_time, add_time_end=end_time) async def monitor_erp_order(event_type, *args, **kwargs): if os.getenv("APP_ENV") != "prod": logger.info(f'非生产环境,不实际从星云加载全量客户数据,{args} {kwargs}') return now_time = datetime.now() logger.info(f'监控ERP订单,{args} {kwargs}') new_order_list = [1] new_order_list = await weixin_customer_controller.monitor_erp_order() if new_order_list: logger.info(f'监控到新订单,{new_order_list}') await msg_controller.sync_msg() logger.debug(f'耗时:{datetime.now() - now_time}') async def update_erp_order(*args, **kwargs): logger.info(f'更新ERP订单,{args} {kwargs}') await msg_controller.get_order_user_days_before(2) async def check_remark_action_is_done(*args, **kwargs): if os.getenv("APP_ENV") != "prod": logger.info(f'非生产环境,不实际检查备注操作是否完成,{args} {kwargs}') return logger.info(f'检查备注操作是否完成,{args} {kwargs}') await action_controller.check_remark_action_is_done() # 每天到点执行一次 today = datetime.now() todayStart = today.replace(hour=0, minute=0, second=0, microsecond=0) if today.hour in [23] and today.minute >= 50: todayStart += timedelta(days=1) class_type = '早班' elif today.hour in [0, 1] and today.minute <= 10: class_type = '早班' elif today.hour in [16] and today.minute >= 50 or today.hour in [17,18] and today.minute <= 10: class_type = '晚班' else: return date = todayStart.strftime('%Y-%m-%d') title = f'设置 {date} {class_type} 的员工活码' success_title = title+' 成功' fail_title = title+' 失败' async with sync_lock: has_set = await msg_controller.model.filter(title__in=[success_title]).first() if has_set: logger.info(f'【已操作】{title}') return logger.info(f'【开始】{title}') try: result, next_class = await weixin_user_controller.set_user_online(class_type=class_type, date=todayStart) except Exception as e: logger.exception(f'设置 {date} {class_type} 的员工活码失败,{e}') result, next_class = {}, None is_failed = bool(list(filter(lambda x: x.get('fail') != 0, result.values()))) or not result logger.info(f'【完成】{title},结果:{result} {"失败" if is_failed else "成功"}') result_str = '\n'.join([f'{k}:{v}' for k, v in result.items()]) logger.info(f'{success_title if not is_failed else fail_title},{result_str}') content = f'\n\n{result_str}' if result and not next_class: content += f'\n\n⚠️ 未设置下一班值班人员,请及时设置' logger.info(f'消息通知:{content}') await msg_controller.new_system_msg(title=success_title if not is_failed else fail_title, content=content) await msg_controller.sync_msg() async def group_chat_update(event_type, group_info, *args): if os.getenv("APP_ENV") != "prod": logger.info(f'非生产环境,不实际检查群聊更新,{group_info} {args}') return group_event = WeixinGroupChatEvent(**group_info) await sync_xingyun_contact_info_async() async def test_event(event_type, event_data=None, *args, **kwargs): # logger.info(f'测试事件,{args} {kwargs}') # scenario_list = await automation_scenario_controller.find_active_scenarios(event_type, event_data=event_data) # logger.info(f'触发事件,{len(scenario_list)}') # for scenario in scenario_list: # task = await task_controller.create_from_scenario(scenario=scenario, event_data=event_data) # logger.info(f'创建任务,{task}') pass async def trigger_event(event_type: str, event_data: Any): """解析事件数据""" logger.info(f'测试事件,{event_type} {event_data}') scenario_list = await automation_scenario_controller.find_active_scenarios(event_type, event_data=event_data) logger.info(f'触发事件,{len(scenario_list)}') for scenario, reason in scenario_list: logger.info(f'触发场景,{scenario},{reason}') try: task = await task_controller.create_from_scenario(scenario=scenario, event_data=event_data, reason=reason) logger.info(f'创建任务,{task}') except Exception as e: logger.error(f'创建任务失败,{e}') # print(f'事件订阅:{transform_pydantic_to_list(WeixinGroupChatEvent)}') event_manager.subscribe(EventType.SYNC_XINGYUN_CONTACT_INFO, load_all_user_from_xingyun) event_manager.subscribe(EventType.OPEN_PERSONAL_CHAT, sync_xingyun_contact_info_async) event_manager.subscribe(EventType.GROUP_CHAT_UPDATED, group_chat_update, input_model=transform_pydantic_to_list(WeixinGroupChatEvent)) event_manager.subscribe(EventType.BIND_ORDER_FOR_USER, sync_xingyun_contact_info_async, input_model=transform_pydantic_to_list(BindOrderResultEvent)) event_manager.subscribe(EventType.OLD_CUSTOMER_REPEAT_PURCHASE, test_event, input_model=transform_pydantic_to_list(CustomerRepeatPurchaseEvent)) event_manager.subscribe(EventType.CUSTOMER_SERVICE_ASSIGN_ORDER, test_event, input_model=transform_pydantic_to_list(CustomerAssignOrderEvent)) event_manager.subscribe(EventType.DESIGNER_ASSIGN_ORDER, test_event, input_model=transform_pydantic_to_list(CustomerAssignOrderEvent)) event_manager.subscribe(EventType.CUSTOMER_RETURN_ORDER, test_event, input_model=transform_pydantic_to_list(CustomerRefundOrderEvent)) event_manager.subscribe(EventType.DESIGNER_UPLOAD_DESIGN, test_event) event_manager.register_pretask(trigger_event) if os.getenv("APP_ENV") == "prod": logger.info('生产环境,启动定时任务') # 创建调度器 scheduler = BackgroundScheduler() # scheduler.add_job(monitor_erp_order, "interval", seconds=60) # 每90秒执行一次 scheduler.add_job(update_erp_order, "interval", seconds=240) # 每90秒执行一次 scheduler.add_job(check_remark_action_is_done, "interval", seconds=180) # 每90秒执行一次 scheduler.start()