import time import re import traceback from config import runing_task, runing_shop,DELETE_BRAND_API_BASE from datetime import datetime from amazon.amazon_base import AmamzonBase, kill_process,TaskBase from amazon.tool import show_notification import requests class AmzoneMatchAction(AmamzonBase): mark_name = "匹配价格" def __init__(self, user_info: dict, socket_port: int = 19890): super().__init__(user_info,socket_port) self.already_asin = set() def reset_already_asin(self): self.already_asin = set() self.log("already_asin 已重置") def wait_loaded(self): # 等待加载完成 try: load_ele = self.tab.ele( 'xpath://kat-panel[@data-testid="kat-panel-ActionPanelContent"]//div[contains(@class,"Loader-module__loader")]|//kat-panel[@data-testid="kat-panel-ActionPanelContent"]/div[@data-f1-component]//div[contains(@class,"==")]/div[contains(@class,"==")]/span', timeout=3) load_ele.wait.deleted(timeout=5, raise_err=False) except Exception as e: print(f"【{self.mark_name}】等待加载中消失出错", e) def run_page_action(self): self.log(f"==================={self.mark_name}=======================") self.log("开始执行...") num = 0 retry_num = 0 get_page_faild = 0 while retry_num < 3: # 最多重试3次 # if num > 3: #测试 # return # 等待加载完成 try: load_ele = self.tab.eles("xpath://div[contains(@class,'Loader-module__loader')]",timeout=5) if len(load_ele) > 0: load_ele[0].wait.deleted(timeout=3, raise_err=False) time.sleep(0.5) # 获取当前页码 try: page_pamel = self.tab.eles('xpath://kat-pagination',timeout=5) if len(page_pamel) > 0: current_page = page_pamel[0].sr('xpath:.//ul[@class="pages"]//li[@aria-current="true"]').text # 测试 # if current_page > 1: # break # 总页数 total_page = page_pamel[0].sr.eles('xpath:.//ul[@class="pages"]//span[@class="page__inner"][last()]') if len(total_page) > 0: total_page = total_page[-1].text else: total_page = 0 self.log(f"【{self.mark_name}】当前页码: {current_page} / 总页数: {total_page}") get_page_faild = 0 except Exception as e: self.log(f"【{self.mark_name}】获取页码失败", e) get_page_faild+= 1 if get_page_faild > 2: show_notification(f"【{self.mark_name}】获取页码失败超3次停止任务!") break # 保存当前URL,用于失败重试时访问 try: self.url = self.tab.url self.log(f"获取当前的链接:{self.url}") except Exception as e: self.url = None self.log(f"获取当前的链接失败:{e}") sku_ls = self.tab.eles("xpath://div[@data-sku]",timeout=10) self.log(f"【{self.mark_name}】获取到 {len(sku_ls)}") # for sku_ele in sku_ls[0:2]: for sku_ele in sku_ls: # solve_problem = sku_ele.eles('xpath:.//kat-link[@label="解决商品信息问题"]') asin = sku_ele.ele('xpath:.//div[contains(@class,"JanusSplitBox-module__container")]//div[contains(@class,"JanusSplitBox-module__panel--") and contains(string(.),"ASIN")]/..//div[last()]',timeout=3).text self.log(f"【{self.mark_name}】ASIN {asin} 找到....") if asin in self.already_asin: self.log(f"【{self.mark_name}】{asin} 已经处理过了,跳过") continue price_match = sku_ele.eles('xpath:.//div[@data-test-id="FeaturedOfferPrice"]//a[text()="匹配"]',timeout=2) if len(price_match) == 0: self.log(f"【{self.mark_name}】{asin},没有推荐价格匹配按钮") yield (asin,"无需处理") self.already_asin.add(asin) continue try: price_match[0].click() save_all_btn = self.tab.ele('xpath://kat-button[@label="保存所有"]',timeout=10) save_all_btn.wait.displayed(timeout=10, raise_err=False) save_all_btn.click() save_all_btn.wait.deleted(timeout=10, raise_err=False) yield (asin,"处理完成") self.already_asin.add(asin) except Exception as e: print(f"{asin} 处理失败", e) yield (asin,"处理失败") self.already_asin.add(asin) continue # 判断是否存在需要翻页的情况 page_pamel = self.tab.eles('xpath://kat-pagination',timeout=5) if len(page_pamel) == 0: break next_page_btn = page_pamel[0].sr('xpath:.//span[@part="pagination-nav-right"]') class_str = next_page_btn.attr('class') if "end" in class_str: break next_page_btn.click() num += 1 self.log(f"【{self.mark_name}】【程序计算】正在翻页,已翻 {num} 页...") except Exception as e: self.log(f"【{self.mark_name}】处理匹配操作异常", e) traceback.print_exc() retry_num += 1 self.tab.refresh() self.tab.wait.doc_loaded(raise_err=False,timeout=120) class MatchTak(TaskBase): task_name = "匹配价格-TASK" def process_task(self, task_data: dict): """处理审批任务主入口 Args: task_data: 任务数据 """ try: data = task_data.get("data", {}) task_id = data.get("taskId") items = data.get("items", []) country_codes = data.get("country_codes", []) risk_listing_filter = data.get("risk_listing_filter", "All") user_id = data.get("user_id") stage_index = data.get("stage_index") final_stage = bool(data.get("final_stage", True)) # 用于测试 limit = data.get("limit",None) if not task_id or not items or not country_codes: self.log(f"任务ID为空{task_id}/店铺列表为空{items}/国家列表为空{country_codes},跳过", "WARNING") return self.log(f"开始处理审批任务 {task_id},共 {len(items)} 个店铺,{len(country_codes)} 个国家") runing_task[task_id] = { "status": "running", "start_time": datetime.now().strftime("%Y-%m-%d %H:%M:%S"), "total_shops": len(items), "processed_shops": 0, "total_countries": len(country_codes) * len(items), "processed_countries": 0, "total_asins": 0, "processed_asins": 0, "success_count": 0, "failed_count": 0, "stop_requested": False } # 遍历处理每个店铺 for idx, shop_item in enumerate(items, 1): # 检查是否收到暂停请求 if task_id in runing_task and runing_task[task_id].get("stop_requested", False): self.log(f"检测到任务 {task_id} 的暂停请求,停止处理", "WARNING") runing_task[task_id]["status"] = "stopped" return shop_name = shop_item.get("shopName", "未知店铺") self.log(f"[{idx}/{len(items)}] 开始处理店铺: {shop_name}") show_notification(f"开始处理店铺: {shop_name}", "info") try: self.process_shop(shop_item, country_codes, task_id, risk_listing_filter, user_id, stage_index, final_stage,limit=limit) # self.process_shop(shop_item, country_codes, task_id,risk_listing_filter) # 更新已处理店铺数 if task_id in runing_task: runing_task[task_id]["processed_shops"] += 1 except Exception as e: self.log(f"处理店铺 {shop_name} 失败: {str(e)}", "ERROR") self.log(traceback.format_exc(), "ERROR") # 更新任务状态 if task_id in runing_task: if runing_task[task_id].get("stop_requested", False): runing_task[task_id]["status"] = "stopped" self.log(f"任务 {task_id} 已被暂停!") else: runing_task[task_id]["status"] = "completed" self.log(f"任务 {task_id} 处理完成!") except Exception as e: self.log(f"任务处理失败: {traceback.format_exc()}", "ERROR") if task_id: if task_id in runing_task: runing_task[task_id]["status"] = "failed" runing_task[task_id]["error"] = str(e) def process_shop(self, shop_item: dict, country_codes: list, task_id: int, risk_listing_filter: str, user_id=None, stage_index=None, final_stage: bool = True,limit:str=None): """处理单个店铺 Args: shop_item: 店铺信息 country_codes: 国家代码列表 task_id: 任务ID risk_listing_filter: 风险商品筛选条件 """ shop_name = shop_item.get("shopName", "未知店铺") company_name = shop_item.get("companyName", "") if not company_name: self.log(f"店铺 {shop_name} 的公司名称为空,跳过", "WARNING") return if task_id in runing_task: runing_task[task_id]["current_shop"] = shop_name # 将店铺添加到正在执行中的店铺列表 start_time = datetime.now().strftime("%Y-%m-%d %H:%M:%S") runing_shop[shop_name] = start_time self.log(f"店铺 {shop_name} 已添加到执行列表,账号: {company_name},开始时间: {start_time}") # 店铺打开重试最多3次 driver = None max_retries = 3 iskill = False # 处理每个国家 for country_code in country_codes: # 检查是否收到暂停请求 if task_id in runing_task and runing_task[task_id].get("stop_requested", False): self.log(f"检测到任务 {task_id} 的暂停请求,停止处理国家", "WARNING") break # 打开店铺 driver = self.open_shop(cls=AmzoneMatchAction,max_retries=max_retries, company_name=company_name, shop_name=shop_name, iskill=iskill) if driver is None: self.log(f"任务 {task_id} 启动店铺失败,结束任务", "ERROR") for country_code in country_codes: self.post_result(task_id, shop_name, country_code, "", "", is_done=True) return current_url = None for _ in range(max_retries): try: self.process_country(driver, country_code, task_id, shop_name,risk_listing_filter,limit, target_url=current_url) driver.reset_already_asin() driver.close_store() break except Exception as e: self.log(f"处理国家 {country_code} 失败: {str(e)}", "ERROR") self.log(traceback.format_exc(), "ERROR") if "与页面的连接已断开" in str(e): iskill = True current_url = driver.url # 更新已处理国家数 if task_id in runing_task: runing_task[task_id]["processed_countries"] += 1 self.log(f"{task_id}任务处理完成") # 最后回传,标记完成 try: # self.post_result(task_id, shop_name, country_code, "", "", is_done=True) if final_stage: self.post_result(task_id, shop_name, country_code, "", "", is_done=True) else: self.post_stage_finished(task_id, user_id, stage_index) except Exception as e: self.log(f"回传结果失败: {str(e)}", "ERROR") # 从正在执行中的店铺列表中移除 if shop_name in runing_shop: del runing_shop[shop_name] self.log(f"店铺 {shop_name} 已从执行列表中移除") def process_country(self, driver, country_code, task_id, shop_name, risk_listing_filter,limit=None,target_url=None): """处理单个国家的审批任务 """ # 转换国家代码为中文名称 country_name = self.country_info.get(country_code, country_code) info_mes = f"开始处理国家: {country_name} ({country_code})" self.log(info_mes) show_notification(info_mes, "info") # 更新当前处理的国家 if task_id in runing_task: runing_task[task_id]["current_country"] = country_name max_retries = 5 switch_success, switch_success_pg, search_success, sku_ls = self.action_init(driver, country_name, shop_name, risk_listing_filter) # 如果切换失败,直接返回 if not switch_success or not switch_success_pg or not search_success: error_message = f"切换到国家({switch_success})/库存页面({switch_success_pg})/搜索筛选到指定选项({search_success}) {country_name} 失败,已重试 {max_retries} 次,跳过该国家" self.log(error_message, "ERROR") if task_id in runing_task: runing_task[task_id]["processed_countries"] += 1 return # # 如果没有需要审批的商品,直接返回 if len(sku_ls) == 0: self.log(f"国家 {country_name} 没有需要审批的商品") if task_id in runing_task: runing_task[task_id]["processed_countries"] += 1 return self.log(f"国家 {country_name} 搜索出 {len(sku_ls)} 商品,开始处理...") if target_url is not None: # 从失败的链接继续 self.log(f"开始访问:{target_url}") driver.tab.get(url=target_url) driver.tab.wait.doc_loaded(timeout=60, raise_err=False) result = [] # 处理所有需要审批的商品(通过yield获取结果) for asin, status in driver.run_page_action(): # 检查是否收到暂停请求 if task_id in runing_task and runing_task[task_id].get("stop_requested", False): self.log(f"检测到任务 {task_id} 的暂停请求,停止处理ASIN", "WARNING") break self.log(f"ASIN {asin} 处理结果: {status}") # 更新任务状态 if task_id in runing_task: runing_task[task_id]["current_asin"] = asin runing_task[task_id]["processed_asins"] += 1 runing_task[task_id]["failed_count"] += 1 result.append({ "asin": asin, "status": status, "done": False }) if len(result) > 20: self.post_result_batch(task_id, shop_name, country_code,result) result = [] if len(result) > 0: self.post_result_batch(task_id, shop_name, country_code, result) self.log(f"国家 {country_name} 处理完成") def post_stage_finished(self, task_id: int, user_id, stage_index): if user_id in (None, "", 0): raise ValueError("user_id is required for stage completion callback") if stage_index is None: raise ValueError("stage_index is required for stage completion callback") url = f"{DELETE_BRAND_API_BASE}/api/shop-match/tasks/{task_id}/stage-finished" payload = {"stage_index": stage_index} params = {"user_id": user_id} max_retries = 3 for retry in range(max_retries): try: self.log(f"Attempting stage completion callback ({retry + 1}/{max_retries})") response = requests.post( url, params=params, json=payload, headers={"Content-Type": "application/json"}, timeout=30, verify=False, ) self.log(f"Stage completion callback response: {response.text}") data = response.json() if response.text else {} if response.status_code == 200 and isinstance(data, dict) and data.get("success"): self.log(f"Stage completion callback succeeded: task={task_id}, stage={stage_index}") return self.log(f"Stage completion callback failed, status={response.status_code}", "WARNING") except Exception as e: self.log(f"Stage completion callback exception: {str(e)}", "ERROR") if retry < max_retries - 1: time.sleep(2) raise RuntimeError(f"Stage completion callback failed after retries: task={task_id}, stage={stage_index}") def post_result(self, task_id: int, shop_name: str, country_code: str, asin: str, status: str,is_done: bool = False): """回传处理结果到API Args: task_id: 任务ID shop_name: 店铺名称 country_code: 国家代码 asin: ASIN status: 处理状态 """ url = f"{DELETE_BRAND_API_BASE}/api/shop-match/tasks/{task_id}/result" payload = { "shops": [ { "error": "", "countries": { country_code: [ { "asin": asin, "status": status, "done": is_done } ] }, "shopName": shop_name } ] } max_retries = 3 for retry in range(max_retries): try: print("================【匹配】=====================") self.log(f"尝试回传结果 (第 {retry + 1}/{max_retries} 次)") self.log(f"回传URL: {url}") self.log(f"回传数据: {payload}") response = requests.post( url, json=payload, headers={"Content-Type": "application/json"}, timeout=30, verify=False ) self.log(f"回传结果: {response.text}") data = response.json() if response.text else {} if response.status_code == 200 and isinstance(data, dict) and data.get("success"): self.log(f"结果回传成功: {asin} - {status}") return else: self.log(f"结果回传失败,状态码: {response.status_code}", "WARNING") print("=====================================") except Exception as e: self.log(f"调用API异常: {str(e)}", "ERROR") print("=====================================") # 如果还有重试机会,等待后继续 if retry < max_retries - 1: time.sleep(2) self.log(f"已达到最大重试次数,结果回传最终失败", "ERROR") raise RuntimeError("已达到最大重试次数,结果回传最终失败") def post_result_batch(self, task_id: int, shop_name: str, country_code: str, country_data:list, is_done: bool = False): """回传处理结果到API Args: task_id: 任务ID shop_name: 店铺名称 country_code: 国家代码 asin: ASIN status: 处理状态 """ url = f"{DELETE_BRAND_API_BASE}/api/shop-match/tasks/{task_id}/result" payload = { "shops": [ { "error": "", "countries": { country_code: country_data }, "shopName": shop_name } ] } max_retries = 3 for retry in range(max_retries): try: print("================【匹配】=====================") self.log(f"尝试回传结果 (第 {retry + 1}/{max_retries} 次)") self.log(f"回传URL: {url}") self.log(f"回传数据: {payload}") response = requests.post( url, json=payload, headers={"Content-Type": "application/json"}, timeout=30, verify=False ) self.log(f"回传结果: {response.text}") data = response.json() if response.text else {} if response.status_code == 200 and isinstance(data, dict) and data.get("success"): self.log(f"结果回传成功") return else: self.log(f"结果回传失败,状态码: {response.status_code}", "WARNING") print("=====================================") except Exception as e: self.log(f"调用API异常: {str(e)}", "ERROR") print("=====================================") # 如果还有重试机会,等待后继续 if retry < max_retries - 1: time.sleep(2) self.log(f"已达到最大重试次数,结果回传最终失败", "ERROR") raise RuntimeError("已达到最大重试次数,结果回传最终失败") if __name__ == "__main__": user_info = { "company": "rongchuang123", "username": "自动化_Robot", "password": "#20zsg25" } shop_name = "魏振峰" country = "德国" kill_process('v6') driver = AmzoneMatchAction(user_info) browser = driver.open_shop(shop_name) sw_suc = driver.SwitchingCountries(country) driver.SwitchPage() risk_listing_filter = "All" for _ in range(3): try: sku_ls = driver.search(filter_type=risk_listing_filter) break except Exception as e: print(e) driver.tab.refresh() if len(sku_ls) > 0: print("有数据,开始操作") for asin, status in driver.run_page_action(): print(f"ASIN {asin} 的处理结果: {status}") print("已完成操作")