modified: amazon/approve.py modified: amazon/asin_status.py new file: amazon/chrome_base.py modified: amazon/del_brand.py modified: amazon/detail_spider.py new file: "amazon/detail_spider_\345\244\207\344\273\275.py" modified: amazon/main.py modified: amazon/match_action.py modified: amazon/price_match.py new file: amazon/similar_asin.py modified: app.py new file: assets/appearance-patent-Br1dtmol.css new file: assets/appearance-patent-D1DgeYhe.css new file: assets/delete-brand-BfMLVSQU.css new file: assets/patrol-delete-BAzbYtgc.css new file: assets/price-track-Hozgt_em.css new file: assets/product-risk-CbX7bwZi.css new file: assets/query-asin-B4RsOiza.css new file: assets/shop-match-B2HWAgBY.css modified: main.py
552 lines
22 KiB
Python
552 lines
22 KiB
Python
import traceback
|
||
import winreg
|
||
import subprocess
|
||
import time
|
||
import uuid
|
||
import requests
|
||
import json
|
||
import os
|
||
from datetime import datetime
|
||
|
||
from amazon.amazon_base import AmamzonBase,kill_process,TaskBase
|
||
from amazon.tool import show_notification,get_shop_info
|
||
|
||
|
||
from config import runing_task,runing_shop,DELETE_BRAND_API_BASE
|
||
|
||
|
||
class AmazoneDriver(AmamzonBase):
|
||
"""亚马逊专用驱动类,继承自ZiniaoDriver"""
|
||
mark_name = "删除品牌"
|
||
|
||
def SwitchPage(self):
|
||
"""
|
||
切换至 管理所有库存页面
|
||
1、等待 //navigation-favorites-bar[@class="hydrated"] 出现
|
||
"""
|
||
navigation = self.tab.ele('xpath://navigation-favorites-bar[@class="hydrated"]')
|
||
navigation.wait.displayed(raise_err=False)
|
||
page_btn = navigation.sr('xpath://internal-fav-bar-links[@data-internal="navigation"]').sr(
|
||
'xpath://a[@data-page-id="ezdpc-gui-inventory-mons"]')
|
||
page_btn.wait.displayed(raise_err=False)
|
||
page_btn.click(timeout=5)
|
||
|
||
self.tab.wait.doc_loaded()
|
||
# 等待搜索框出现
|
||
search_region = self.tab.ele('xpath://div[@id="searchBoxContainer"]//kat-input-group')
|
||
search_region.wait.displayed(raise_err=False)
|
||
|
||
def search(self, asin):
|
||
search_region = self.tab.ele('xpath://div[@id="searchBoxContainer"]//kat-input-group')
|
||
search_region.wait.displayed(raise_err=False)
|
||
time.sleep(0.6)
|
||
search_input = self.tab.ele("xpath://kat-input[contains(@class,'SearchBox-module__searchInput')]").sr(
|
||
'xpath://span[@class="container"]//input[@part="input"]')
|
||
search_input.input(asin,clear=True)
|
||
sku_ls = []
|
||
for _ in range(3):
|
||
search_btn = self.tab.ele("xpath://kat-icon[@name='search']")
|
||
search_btn.click()
|
||
|
||
load_ele = self.tab.ele("xpath://div[contains(@class,'Loader-module__loader')]")
|
||
# load_ele.wait.hidden(timeout=3, raise_err=False)
|
||
load_ele.wait.deleted(timeout=3, raise_err=False)
|
||
time.sleep(0.5)
|
||
|
||
sku_ls = self.tab.eles("xpath://div[@data-sku]",timeout=3)
|
||
if len(sku_ls) > 0:
|
||
break
|
||
return sku_ls
|
||
|
||
def del_action(self, sku_ele):
|
||
dropdown = sku_ele.ele('xpath:.//kat-dropdown-button[@variant="secondary"]')
|
||
dropdown.click()
|
||
time.sleep(1)
|
||
|
||
del_btn = dropdown.sr("xpath:.//button[@role='menuitem' and @data-action='DeleteListing']")
|
||
del_btn.wait.displayed(raise_err=False)
|
||
del_btn.click()
|
||
|
||
confirm_btn = self.tab.ele('xpath://kat-modal[@data-testid="action-modal"]//kat-button[@variant="primary"]')
|
||
confirm_btn.wait.displayed(raise_err=False)
|
||
confirm_btn.wait.enabled()
|
||
confirm_btn.click()
|
||
|
||
suc_alert = self.tab.ele("xpath://kat-alert[@variant='success']")
|
||
suc = suc_alert.wait.displayed(raise_err=False, timeout=5)
|
||
return suc
|
||
|
||
|
||
class DelbramdTask(TaskBase):
|
||
task_name = "删除品牌-TASK"
|
||
|
||
def update_task_status(self, task_id: int, **kwargs):
|
||
"""更新任务状态
|
||
|
||
Args:
|
||
task_id: 任务ID
|
||
**kwargs: 要更新的字段
|
||
"""
|
||
if task_id not in runing_task:
|
||
runing_task[task_id] = {}
|
||
runing_task[task_id].update(kwargs)
|
||
|
||
def process_task(self,task_data):
|
||
try:
|
||
data = task_data.get("data", {})
|
||
task_id = data.get("taskId")
|
||
items = data.get("items", [])
|
||
|
||
if not task_id:
|
||
self.log("任务ID为空,跳过", "WARNING")
|
||
return
|
||
|
||
# 初始化任务状态
|
||
self.update_task_status(
|
||
task_id,
|
||
status="running",
|
||
start_time=datetime.now().strftime("%Y-%m-%d %H:%M:%S"),
|
||
total_shops=len(items),
|
||
processed_shops=0,
|
||
total_asins=sum(shop.get("totalRows", 0) for shop in items),
|
||
processed_asins=0,
|
||
success_count=0,
|
||
failed_count=0,
|
||
stop_requested=False # 暂停请求标志
|
||
)
|
||
|
||
self.log(f"开始处理任务 {task_id},共 {len(items)} 个店铺")
|
||
|
||
# 遍历处理每个店铺
|
||
for idx, shop_data in enumerate(items, 1):
|
||
shop_name = shop_data.get("shopName", "未知店铺")
|
||
self.log(f"[{idx}/{len(items)}] 开始处理店铺: {shop_name}")
|
||
|
||
try:
|
||
self.process_shop(shop_data, task_id)
|
||
# 更新已处理店铺数
|
||
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 and runing_task[task_id].get("stop_requested", False):
|
||
self.update_task_status(task_id, status="stopped")
|
||
self.log(f"任务 {task_id} 已被暂停!")
|
||
else:
|
||
self.update_task_status(task_id, status="completed")
|
||
self.log(f"任务 {task_id} 处理完成!")
|
||
|
||
except Exception as e:
|
||
self.log(f"任务处理失败: {traceback.format_exc()}", "ERROR")
|
||
if task_id:
|
||
self.update_task_status(task_id, status="failed", error=str(e))
|
||
|
||
def process_shop(self, shop_data, task_id):
|
||
"""处理单个店铺(包含重试逻辑)
|
||
Args:
|
||
shop_data: 店铺数据
|
||
task_id: 任务ID
|
||
"""
|
||
shop_name = shop_data.get("shopName", "未知店铺")
|
||
result_id = shop_data.get("resultId")
|
||
countries = shop_data.get("countries", [])
|
||
company_name = shop_data.get("companyName", "未知公司")
|
||
|
||
if not countries:
|
||
self.log(f"店铺 {shop_name} 没有国家数据,跳过", "WARNING")
|
||
return
|
||
|
||
# 更新当前处理的店铺
|
||
self.update_task_status(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"账号:{company_name},店铺 {shop_name} 已添加到执行列表,开始时间: {start_time}")
|
||
|
||
# 店铺打开重试最多3次
|
||
max_retries = 5
|
||
iskill = False
|
||
|
||
|
||
# 处理每个国家
|
||
chunk_index = 1
|
||
for country_data in countries:
|
||
self.log(f"开始处理国家:{country_data}")
|
||
# 检查是否收到暂停请求
|
||
if task_id in runing_task and runing_task[task_id].get("stop_requested", False):
|
||
self.log(f"检测到任务 {task_id} 的暂停请求,停止处理国家", "WARNING")
|
||
break # 跳出循环,进入finally关闭店铺
|
||
|
||
# 打开店铺
|
||
driver = self.open_shop(AmazoneDriver,max_retries=max_retries, company_name=company_name,
|
||
shop_name=shop_name, iskill=iskill)
|
||
|
||
# 检查是否成功打开
|
||
if not driver:
|
||
self.log(f"店铺 {shop_name} 打开失败,已重试 {max_retries} 次,跳过该店铺", "ERROR")
|
||
return
|
||
|
||
try:
|
||
chunk_index = self.process_country(driver, country_data, task_id, result_id, shop_data, chunk_index)
|
||
except Exception as e:
|
||
country_name = country_data.get("country", "未知")
|
||
self.log(f"处理国家 {country_name} 失败: {str(e)}", "ERROR")
|
||
self.log(traceback.format_exc(), "ERROR")
|
||
self.log(traceback.format_exc(), "ERROR")
|
||
if "与页面的连接已断开" in str(e):
|
||
iskill = True
|
||
|
||
try:
|
||
driver.close_store()
|
||
except Exception as e:
|
||
iskill = True
|
||
self.log(f"关闭国家失败,{e}")
|
||
|
||
# 从正在执行中的店铺列表中移除
|
||
if shop_name in runing_shop:
|
||
del runing_shop[shop_name]
|
||
self.log(f"店铺 {shop_name} 已从执行列表中移除")
|
||
|
||
def process_country(self, driver: AmazoneDriver, country_data, task_id: int,
|
||
result_id: int, shop_data, chunk_index: int):
|
||
"""处理单个国家的所有ASIN
|
||
|
||
Args:
|
||
driver: 亚马逊驱动实例
|
||
country_data: 国家数据
|
||
task_id: 任务ID
|
||
result_id: 结果ID
|
||
shop_data: 店铺数据
|
||
"""
|
||
country = country_data.get("country", "未知")
|
||
items = country_data.get("items", [])
|
||
shop_name = shop_data.get("shopName", "未知店铺")
|
||
|
||
|
||
self.log(f"开始处理国家: {country},共 {len(items)} 个ASIN")
|
||
self.update_task_status(task_id, current_country=country)
|
||
|
||
max_retries = 5
|
||
switch_success, switch_success_pg, search_success, sku_ls = self.action_init(driver, country, shop_name,
|
||
None)
|
||
# 如果切换失败,直接返回
|
||
if not switch_success or not switch_success_pg or not search_success:
|
||
error_message = f"切换到国家({switch_success})/库存页面({switch_success_pg})/搜索筛选到指定选项({search_success}) {country} 失败,已重试 {max_retries} 次,跳过该国家"
|
||
self.log(error_message, "ERROR")
|
||
if task_id in runing_task:
|
||
runing_task[task_id]["processed_countries"] += 1
|
||
return chunk_index
|
||
|
||
# 处理每个ASIN
|
||
file_key = shop_data.get("fileKey", "")
|
||
source_filename = shop_data.get("sourceFilename", "")
|
||
total_rows = shop_data.get("totalRows", 0)
|
||
|
||
for idx, asin_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")
|
||
return # 返回到上层,会触发店铺关闭
|
||
|
||
try:
|
||
self.log(f"[{idx}/{len(items)}] 处理ASIN: {asin_item.get('asin', '')}")
|
||
self.process_asin(driver, asin_item, country, task_id, result_id,
|
||
file_key, source_filename, total_rows, chunk_index)
|
||
except Exception as e:
|
||
asin = asin_item.get("asin", "未知")
|
||
self.log(f"处理ASIN {asin} 失败: {str(e)}", "ERROR")
|
||
# 继续处理下一个ASIN
|
||
chunk_index += 1
|
||
|
||
return chunk_index
|
||
|
||
def process_asin(self, driver: AmazoneDriver, asin_item,
|
||
country: str, task_id: int, result_id: int, file_key: str,
|
||
source_filename: str, total_rows: int, chunk_index: int):
|
||
"""处理单个ASIN并回传结果
|
||
|
||
Args:
|
||
driver: 亚马逊驱动实例
|
||
asin_item: ASIN数据项
|
||
country: 国家名称
|
||
task_id: 任务ID
|
||
result_id: 结果ID
|
||
file_key: 文件KEY
|
||
source_filename: 源文件名
|
||
total_rows: 总行数
|
||
chunk_index: 当前处理索引
|
||
"""
|
||
asin = asin_item.get("asin", "")
|
||
row_index = asin_item.get("rowIndex", 0)
|
||
|
||
# 更新当前处理的ASIN
|
||
self.update_task_status(task_id, current_asin=asin)
|
||
|
||
status = "失败"
|
||
max_retries = 3 # 最多重试3次
|
||
|
||
for retry in range(max_retries):
|
||
try:
|
||
self.log(f"处理ASIN {asin} (第 {retry + 1}/{max_retries} 次)")
|
||
|
||
# 如果不是第一次尝试,先刷新页面
|
||
if retry > 0:
|
||
self.log("重试前刷新页面...")
|
||
try:
|
||
driver.tab.refresh()
|
||
time.sleep(3)
|
||
except Exception as e:
|
||
self.log(f"刷新页面失败: {str(e)}", "WARNING")
|
||
|
||
# 搜索ASIN
|
||
sku_ls = driver.search(asin=asin)
|
||
self.log(f"搜索到 {len(sku_ls)} 个SKU")
|
||
|
||
if len(sku_ls) == 0:
|
||
status = "查询不到"
|
||
self.log(f"ASIN {asin} 未找到商品", "WARNING")
|
||
break # 查询不到商品,无需重试
|
||
else:
|
||
# 删除所有找到的SKU,如果任何一个失败则重新开始整个流程
|
||
success_count = 0
|
||
all_success = True # 标记是否所有SKU都删除成功
|
||
total_sku_count = len(sku_ls)
|
||
|
||
for sku in sku_ls:
|
||
try:
|
||
suc = driver.del_action(sku)
|
||
if suc:
|
||
success_count += 1
|
||
self.log(f"SKU 删除成功 ({success_count}/{total_sku_count})")
|
||
else:
|
||
self.log(f"SKU 删除失败", "WARNING")
|
||
all_success = False
|
||
break # 任何一个失败,退出循环,准备重试整个流程
|
||
except Exception as e:
|
||
self.log(f"删除SKU异常: {str(e)}", "ERROR")
|
||
all_success = False
|
||
break # 发生异常,退出循环,准备重试整个流程
|
||
|
||
# 判断删除结果
|
||
if all_success and success_count == total_sku_count:
|
||
status = "成功"
|
||
self.log(f"ASIN {asin} 所有SKU删除成功 ({success_count}/{total_sku_count})")
|
||
# 更新成功计数
|
||
if task_id in runing_task:
|
||
runing_task[task_id]["success_count"] += 1
|
||
break # 全部成功,跳出重试循环
|
||
|
||
else:
|
||
# 有失败的SKU
|
||
if retry < max_retries - 1:
|
||
self.log(f"有SKU删除失败,准备重试整个流程... ({retry + 1}/{max_retries})")
|
||
time.sleep(2)
|
||
continue # 继续下一次重试
|
||
else:
|
||
# 所有重试都用完了
|
||
if success_count > 0:
|
||
status = "部分成功"
|
||
self.log(f"ASIN {asin} 部分SKU删除成功 ({success_count}/{total_sku_count})", "WARNING")
|
||
if task_id in runing_task:
|
||
runing_task[task_id]["success_count"] += 1
|
||
else:
|
||
status = "失败"
|
||
self.log(f"ASIN {asin} 所有SKU删除失败", "ERROR")
|
||
if task_id in runing_task:
|
||
runing_task[task_id]["failed_count"] += 1
|
||
|
||
except Exception as e:
|
||
status = "删除异常"
|
||
self.log(f"处理ASIN {asin} 异常: {str(e)}", "ERROR")
|
||
|
||
# 如果还有重试机会,继续重试
|
||
if retry < max_retries - 1:
|
||
self.log(f"发生异常,准备重试... ({retry + 1}/{max_retries})")
|
||
time.sleep(2)
|
||
else:
|
||
# 所有重试都失败了,更新失败计数
|
||
if task_id in runing_task:
|
||
runing_task[task_id]["failed_count"] += 1
|
||
|
||
# 更新已处理ASIN计数
|
||
if task_id in runing_task:
|
||
runing_task[task_id]["processed_asins"] += 1
|
||
|
||
# 回传结果到API
|
||
try:
|
||
payload = {
|
||
"submissionId": "",
|
||
"files": [{
|
||
"fileKey": file_key,
|
||
"sourceFilename": source_filename,
|
||
"chunkIndex": chunk_index,
|
||
"chunkTotal": total_rows,
|
||
"processedRows": chunk_index,
|
||
"totalRows": total_rows,
|
||
"currentCountry": country,
|
||
"currentAsin": asin,
|
||
"countries": [{
|
||
"country": country,
|
||
"items": [{
|
||
"asin": asin,
|
||
"status": status
|
||
}]
|
||
}]
|
||
}]
|
||
}
|
||
|
||
self.post_result(task_id, payload)
|
||
self.log(f"ASIN {asin} 结果已回传,状态: {status}")
|
||
|
||
except Exception as e:
|
||
self.log(f"回传结果失败: {str(e)}", "ERROR")
|
||
|
||
def _report_all_asins_failed(self, country: str, items, task_id, shop_data, chunk_index):
|
||
"""将国家下所有ASIN标记为失败并回传
|
||
|
||
Args:
|
||
country: 国家名称
|
||
items: ASIN列表
|
||
task_id: 任务ID
|
||
shop_data: 店铺数据
|
||
"""
|
||
self.log(f"开始回传国家 {country} 下的 {len(items)} 个ASIN失败状态")
|
||
|
||
file_key = shop_data.get("fileKey", "")
|
||
source_filename = shop_data.get("sourceFilename", "")
|
||
total_rows = shop_data.get("totalRows", 0)
|
||
|
||
for asin_item in items:
|
||
asin = asin_item.get("asin", "")
|
||
|
||
try:
|
||
# 更新失败计数
|
||
if task_id in runing_task:
|
||
runing_task[task_id]["failed_count"] += 1
|
||
runing_task[task_id]["processed_asins"] += 1
|
||
|
||
# 回传失败状态
|
||
payload = {
|
||
"submissionId": "",
|
||
"files": [{
|
||
"fileKey": file_key,
|
||
"sourceFilename": source_filename,
|
||
"chunkIndex": chunk_index,
|
||
"chunkTotal": total_rows,
|
||
"processedRows": chunk_index,
|
||
"totalRows": total_rows,
|
||
"currentCountry": country,
|
||
"currentAsin": asin,
|
||
"countries": [{
|
||
"country": country,
|
||
"items": [{
|
||
"asin": asin,
|
||
"status": "失败"
|
||
}]
|
||
}]
|
||
}]
|
||
}
|
||
|
||
self.post_result(task_id, payload)
|
||
chunk_index += 1
|
||
self.log(f"ASIN {asin} 失败状态已回传")
|
||
|
||
except Exception as e:
|
||
self.log(f"回传ASIN {asin} 失败状态时出错: {str(e)}", "ERROR")
|
||
|
||
return chunk_index
|
||
|
||
def post_result(self, task_id: int, payload):
|
||
"""回传结果到API(带重试机制)
|
||
|
||
Args:
|
||
task_id: 任务ID
|
||
payload: 结果数据
|
||
"""
|
||
url = f"{DELETE_BRAND_API_BASE}/api/delete-brand/tasks/{task_id}/result"
|
||
max_retries = 3 # 最多重试3次
|
||
|
||
for retry in range(max_retries):
|
||
try:
|
||
self.log(f"尝试回传结果 (第 {retry + 1}/{max_retries} 次)")
|
||
|
||
response = requests.post(
|
||
url,
|
||
json=payload,
|
||
headers={"Content-Type": "application/json"},
|
||
timeout=30,
|
||
verify=False # 忽略SSL证书验证
|
||
)
|
||
print("【结果提交】:", payload)
|
||
print("【结果提交返回】:", response.text)
|
||
|
||
if response.status_code == 200:
|
||
self.log(f"结果回传成功: {url}")
|
||
return # 成功后直接返回,不再重试
|
||
else:
|
||
self.log(f"结果回传失败,状态码: {response.status_code}", "WARNING")
|
||
# 如果还有重试机会,继续重试
|
||
if retry < max_retries - 1:
|
||
self.log(f"准备重试... ({retry + 1}/{max_retries})")
|
||
time.sleep(2) # 等待2秒后重试
|
||
else:
|
||
self.log(f"已达到最大重试次数,结果回传最终失败", "ERROR")
|
||
|
||
except Exception as e:
|
||
self.log(f"调用API异常: {str(e)}", "ERROR")
|
||
# 如果还有重试机会,继续重试
|
||
if retry < max_retries - 1:
|
||
self.log(f"发生异常,准备重试... ({retry + 1}/{max_retries})")
|
||
time.sleep(2) # 等待2秒后重试
|
||
else:
|
||
self.log(f"已达到最大重试次数,结果回传最终失败", "ERROR")
|
||
|
||
def update_task_status(self, task_id: int, **kwargs):
|
||
"""更新任务状态
|
||
|
||
Args:
|
||
task_id: 任务ID
|
||
**kwargs: 要更新的字段
|
||
"""
|
||
if task_id not in runing_task:
|
||
runing_task[task_id] = {}
|
||
|
||
runing_task[task_id].update(kwargs)
|
||
|
||
|
||
|
||
if __name__ == '__main__':
|
||
# 使用示例
|
||
user_info = {
|
||
"company": "rongchuang123",
|
||
"username": "自动化_Robot",
|
||
"password": "#20zsg25"
|
||
}
|
||
shop_name = "郭亚芳"
|
||
kill_process("v6")
|
||
|
||
# 创建驱动实例
|
||
driver = AmazoneDriver(user_info)
|
||
|
||
# 打开店铺并获取浏览器实例
|
||
browser = driver.open_shop(shop_name)
|
||
|
||
country = "西班牙"
|
||
asin = "Voanos"
|
||
|
||
sw_suc = driver.SwitchingCountries(country)
|
||
driver.SwitchPage()
|
||
sku_ls = driver.search(asin=asin)
|
||
print("查询结果有:", len(sku_ls))
|
||
|
||
for sku in sku_ls:
|
||
suc = driver.del_action(sku)
|
||
print(sku, "删除结果", suc)
|
||
|
||
|
||
|