modified: app/amazon/__pycache__/approve.cpython-39.pyc

modified:   app/amazon/__pycache__/del_brand.cpython-39.pyc
	modified:   app/amazon/__pycache__/main.cpython-39.pyc
	new file:   app/amazon/__pycache__/match_action.cpython-39.pyc
	new file:   app/amazon/__pycache__/tool.cpython-39.pyc
	new file:   "app/amazon/approve - \345\211\257\346\234\254.py"
	modified:   app/amazon/approve.py
	modified:   app/amazon/del_brand.py
	modified:   app/amazon/main.py
	new file:   app/amazon/match_action.py
	new file:   app/amazon/tool.py
	modified:   app/main.py
	modified:   app/web_source/brand.html
This commit is contained in:
铭坤
2026-04-10 16:51:52 +08:00
parent ccaec4a984
commit 2d0d1d3461
13 changed files with 2125 additions and 239 deletions

View File

@@ -7,6 +7,7 @@ from concurrent.futures import ThreadPoolExecutor, as_completed
from config import JSON_TASK_QUEUE, runing_task, runing_shop, DELETE_BRAND_API_BASE, ZN_COMPANY, ZN_USERNAME, ZN_PASSWORD
from amazon.del_brand import AmazoneDriver, kill_process
from amazon.approve import ApproveTask
from amazon.match_action import MatchTak
class TaskMonitor:
@@ -47,6 +48,10 @@ class TaskMonitor:
self.executor = ThreadPoolExecutor(max_workers=self.max_workers)
futures = [] # 保存所有提交的任务Future对象
task_type_info = {
"product-risk-resolve-run" : "产品风险审批",
"shop-match-run" : "匹配价格"
}
try:
while self.running:
try:
@@ -63,10 +68,10 @@ class TaskMonitor:
future = self.executor.submit(self._process_task_wrapper, task_data)
futures.append(future)
elif task_type == "product-risk-resolve-run":
elif task_type in task_type_info:
# 提交产品风险审批任务到线程池
self.log(f"接收到产品风险审批任务,提交到线程池处理...")
future = self.executor.submit(self._process_approve_task_wrapper, task_data)
future = self.executor.submit(self._process_approve_task_wrapper, task_data, task_type_info[task_type])
futures.append(future)
else:
@@ -102,16 +107,21 @@ class TaskMonitor:
except Exception as e:
self.log(f"线程 {id(task_data)} 删除品牌任务处理异常: {traceback.format_exc()}", "ERROR")
def _process_approve_task_wrapper(self, task_data: Dict[str, Any]):
def _process_approve_task_wrapper(self, task_data: Dict[str, Any],TASK_TYPE:str):
"""审批任务处理包装器(用于线程池调用)
Args:
task_data: 任务数据
"""
try:
TASK_INFO = {
"产品风险审批" : ApproveTask,
"匹配价格" : MatchTak
}
self.log(f"线程 {id(task_data)} 开始处理产品风险审批任务...")
# 创建ApproveTask实例并处理任务
approve_task = ApproveTask(user_info=self.user_info)
TASK_CLS = TASK_INFO[TASK_TYPE] # 根据任务类型选择处理类默认为ApproveTask
approve_task = TASK_CLS(user_info=self.user_info)
approve_task.process_task(task_data)
self.log(f"线程 {id(task_data)} 产品风险审批任务处理完成")
except Exception as e:
@@ -232,24 +242,24 @@ class TaskMonitor:
driver = None
# 判断是否需要登录
# need_login = driver.need_login()
# if need_login:
# self.log(f"店铺 {shop_name} 需要登录,正在登录...")
# password = ""
need_login = driver.need_login()
if need_login:
self.log(f"店铺 {shop_name} 需要登录,正在登录...")
password = ""
# login_success = driver.login(password)
# if login_success:
# self.log(f"店铺 {shop_name} 登录成功,正在重新打开店铺...")
# browser = driver.open_shop(shop_name)
# if browser and browser != "店铺不存在":
# self.log(f"成功打开店铺 {shop_name} 登录后")
# break
# else:
# self.log(f"登录后打开店铺失败: {browser}", "WARNING")
# driver = None
# else:
# self.log(f"店铺 {shop_name} 登录失败", "WARNING")
# driver = None
login_success = driver.login(password)
if login_success:
self.log(f"店铺 {shop_name} 登录成功,正在重新打开店铺...")
browser = driver.open_shop(shop_name)
if browser and browser != "店铺不存在":
self.log(f"成功打开店铺 {shop_name} 登录后")
break
else:
self.log(f"登录后打开店铺失败: {browser}", "WARNING")
driver = None
else:
self.log(f"店铺 {shop_name} 登录失败", "WARNING")
driver = None
except Exception as e: