Files

627 lines
29 KiB
Python
Raw Permalink Normal View History

"""
xc_validation_strategy_gguf_hygon 主入口
GGUF 模型下载 + 验证任务提交流水线部署框架与 xc_validation_strategy 一致
启动后运行流水线批量创建 GGUF 模型下载任务最大并发 8
遇到并发限制自动等待重试每个模型下载成功后立即提交
hygon / bi150 两个验证任务
同时暴露 /healthK8s 探活 /status运行状态
"""
import json
import os
import re
import signal
import threading
from datetime import datetime
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from typing import Set
import requests
# ══════════════════════════════════════════════════════════
# 配置
# ══════════════════════════════════════════════════════════
BASE_URL = os.environ.get("BASE_URL", "https://modelhub.org.cn")
LOGIN_ENDPOINT = "/adminApi/user/login"
CREATE_DOWNLOAD_TASK_ENDPOINT = "/adminApi/async/task/model-download-task"
SUBMIT_TEST_TASK_ENDPOINT = "/adminApi/async/task/create-contest-task"
# 流水线内部登录账号(创建下载任务 / 提交验证任务用)
USER_ACCOUNT = "zhangyuanxi"
USER_PASSWORD = "jianghanjiao0"
CONTEST_API_TOKEN = "ef1ef82f3c9efee413d602345fbe224d"
HF_TOKEN = "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"
CONTRIBUTORS = "zhoushasha"
TASK_TYPE = "text-generation"
STRATEGY_ID = os.environ.get("STRATEGY_ID", "") # 平台自动注入,无需修改
MAX_CONCURRENT_DOWNLOADS = 8 # 同时下载的模型数上限
CHECK_INTERVAL_SECONDS = 10 # 下载状态轮询间隔(秒)
RETRY_WAIT_SECONDS = 20 # 命中并发限制时的等待时间(秒)
HTTP_HOST = "0.0.0.0"
HTTP_PORT = 8080
# ══════════════════════════════════════════════════════════
# 模型列表
# ══════════════════════════════════════════════════════════
ALL_MODEL_IDS = [
"mradermacher/Ice0.80-03.02-RP-GGUF",
"mradermacher/NeuralLemonKuno-GGUF",
"mradermacher/DavidAU-Dark_Mistress-8B-GGUF",
"mradermacher/dermai-v2-GGUF",
"mradermacher/LemonNalyvkaRP-GGUF",
"mradermacher/distill_70b_infra_together-GGUF",
"mradermacher/nsfw-i-hate-my-life-v1-GGUF",
"mradermacher/M7Yamshadowexperiment28_Ognoexperiment27Multi_verse_model-GGUF",
"mradermacher/nsfw-i-hate-my-life-v2-GGUF",
"mradermacher/Experiment28T3q_YamOgno-GGUF",
"mradermacher/Ice0.81-03.02-RP-GGUF",
"mradermacher/Qwen2-7B-sft-hhrlhf-spin-GGUF",
"mradermacher/WillowGPT-GGUF",
"mradermacher/Dark_Master-v0.1-8B-GGUF",
"mradermacher/YarnGPTII-GGUF",
"mradermacher/kyara-lite-sft-1M-GGUF",
"mradermacher/text-arco-GGUF",
"mradermacher/Mistral-7B-v0.3-sft-ultrachat-GGUF",
"mradermacher/Saba2-Preview-GGUF",
"mradermacher/Anemoi-4B-GGUF",
"mradermacher/exaone-3.5-2.4b-instruct-dacon-llm-GGUF",
"mradermacher/light-3B-beta-GGUF",
"mradermacher/bellman-mistral-7b-instruct-v0.3-GGUF",
"mradermacher/Kosmos-EVAA-Franken-Immersive-v40-8B-GGUF",
"mradermacher/Qwen2.5-Interpreter-GGUF",
"mradermacher/Llama-3.1-8B-Gemma-2-9B-mix-GGUF",
"mradermacher/Llama-3.1-8B-sft-ultrachat-GGUF",
"mradermacher/t5-base-english-japanese-GGUF",
"mradermacher/tourist-GGUF",
"mradermacher/t5-base-japanese-question-generation-GGUF",
"mradermacher/t5-qiita-title-generation-GGUF",
"mradermacher/Not_Even_My_Final_Form-8B-Model_Stock-GGUF",
"mradermacher/Morphing-8B-Model_Stock-GGUF",
"mradermacher/Llamaverse-3.1-8B-Instruct-GGUF",
"mradermacher/Gemma-2-9B-sft-ultrachat-GGUF",
"mradermacher/Italy-10B-GGUF",
"mradermacher/Qwen2.5-7B-sft-ultrachat-safeRLHF-GGUF",
"mradermacher/Qwen2-7B-sft-ultrachat-safeRLHF-GGUF",
"mradermacher/qwen-2.5-0.5b-unlabeled-pile-GGUF",
"mradermacher/HQwen-3B-00-fp16-GGUF",
"mradermacher/AI_Pet_code-GGUF",
"mradermacher/codellama-7b-betterv-params-c2v-chisel-GGUF",
"mradermacher/LongRAG-Qwen2.5-7B-Instruct-GGUF",
"mradermacher/granite-3.0-3b-a800m-base-GGUF",
"mradermacher/granite-3.0-8b-base-GGUF",
"mradermacher/moxin-llm-7b-GGUF",
"mradermacher/llama-3.2-3b-merged-base-3K-updated-GGUF",
"mradermacher/II-Tulu-8B-SFT-GGUF",
"mradermacher/HebQwen-json-2025-GGUF",
"mradermacher/WiNGPT-Babel-GGUF",
"mradermacher/oh-dcft-v3.1-llama-3.1-405b-v2dummytesting-GGUF",
"mradermacher/greesychat-turbo-GGUF",
"mradermacher/MFANN-abliterated-phi2-merge-unretrained-GGUF",
"mradermacher/oh-dcft-v3.1-llama-3.1-405b-qwen-v2dummytesting-GGUF",
"mradermacher/Llama-3.1-8B-sft-ultrachat-safeRLHF-GGUF",
"mradermacher/Falcon-3B-8k-data-new-classes-GGUF",
"mradermacher/Lunar-4B-GGUF",
"mradermacher/Zelus-8B-Model_Stock-GGUF",
"mradermacher/mpn_mistral7bv3_sft-GGUF",
"mradermacher/Qwen2-7B-sft-ultrachat-GGUF",
"mradermacher/Llama-3.1-8B-Instruct-GSM8k-GGUF",
"mradermacher/Mistral-7B-v0.3-sft-ultrachat-safeRLHF-GGUF",
"mradermacher/QwENDEAVR-VL-2B-GGUF",
"mradermacher/EXAONE-3.5-7.8B-it-DACON-LLM-GGUF",
"mradermacher/Gemma-2-9B-sft-ultrachat-safeRLHF-GGUF",
"mradermacher/Eunoia-GwQ-Gemma-9B-GGUF",
"mradermacher/periquito-3B-GGUF",
"mradermacher/T3qm7xNeuralsirkrishna_Strangemerges_30Ogno-GGUF",
"mradermacher/Chalice1-GGUF",
"mradermacher/Yamshadowexperiment28M70.5-0.45-0.64-0.12-0.94-0.72-7B-GGUF",
"mradermacher/M7Yamshadowexperiment28_Strangemerges_32Yamshadow-GGUF",
"mradermacher/Experiment26Yam_Strangemerges_32Ogno-GGUF",
"mradermacher/Experiment28T3q_NeuralsirkrishnaExperiment26-GGUF",
"mradermacher/YamshadowStrangemerges_32_Strangemerges_32Experiment28-GGUF",
"mradermacher/Experiment28M7_Experiment28Neuralsirkrishna-GGUF",
"mradermacher/YamshadowStrangemerges_32_Neuralsirkrishnaexperiment26Experiment24-GGUF",
"mradermacher/YamshadowInex12_Experiment29Pastiche-GGUF",
"mradermacher/tsukasa-7b-qlora-GGUF",
"mradermacher/Qwen2.5-3B-Model-Stock-v2-GGUF",
"mradermacher/LongRAG_Ministral-8B-Instruct-GGUF",
"mradermacher/LongRAG_Llama3.1-8B-Instruct-GGUF",
"mradermacher/inter-play-sim-assistant-test-GGUF",
"mradermacher/Llama-3.2-1B-magnitude-0.1-GGUF",
"mradermacher/Howdy-8B-LINEAR-GGUF",
"mradermacher/Qwen2.5-3B-Model-Stock-GGUF",
"mradermacher/Ice0.51-16.01-RP-GGUF",
"mradermacher/Qwen2.5-3B-Model-Stock-v3-GGUF",
"mradermacher/SalvadAI-GGUF",
"mradermacher/Yi-6B-Chat-Llama-3.3-70B-Instruct-honest_lying-GGUF",
"mradermacher/Funny-Chatbot-Llama-3.1-8B-Instruct-GGUF",
"mradermacher/Aira-2-355M-GGUF",
"mradermacher/Qmerft-GGUF",
"mradermacher/Qwen2.5-1.5B-Instruct-DPO-bad-boy-GGUF",
"mradermacher/NxMobileLM-1.5B-SFT-GGUF",
"mradermacher/Llama-3-ELYZA-JP-8B-normal-chosen-GGUF",
"mradermacher/HebQwen-json-2025-meta-GGUF",
"mradermacher/JapMed-GGUF",
"mradermacher/What_A_Thrill-8B-Model_Stock-GGUF",
"mradermacher/phi-3-clinical-GGUF",
"mradermacher/ORPO-Llama-SQL-8B-Spider-GGUF",
"mradermacher/Gemma-Radiation-RP-9B-lorablated-GGUF",
"mradermacher/D.A.T.A.v1-GGUF",
"mradermacher/gemmarpp-GGUF",
"mradermacher/Meta-Llama-3-8B-wanda-0.1-GGUF",
"mradermacher/LwQ-10B-Instruct-GGUF",
"mradermacher/occiglot-7b-de-en-instruct-dpo1-GGUF",
"mradermacher/MCQ-qkvo-2-GGUF",
"mradermacher/YamshadowStrangemerges_32_Inex12Experiment28-GGUF",
"mradermacher/Bahasa-4b-GGUF",
"mradermacher/MCQ-head-16-GGUF",
"mradermacher/Experiment26Yam_Ognoexperiment27Multi_verse_model-GGUF",
"mradermacher/An4-7Bv2.4-GGUF",
"mradermacher/YamshadowStrangemerges_32_YamExperiment28-GGUF",
"mradermacher/bitmap-M7-alpaca-70m-GGUF",
"mradermacher/UpshotLlama-3-8B-GGUF",
"mradermacher/ChimeraLlama-3-8B-GGUF",
"mradermacher/Yamshadowexperiment28M70.8-0.85-0.86-0.36-0.42-0.16-7B-GGUF",
"mradermacher/microcoderfim-1B-GGUF",
"mradermacher/Medusa_OpenHermes-2.5-Mistral7B-DPO-GGUF",
"mradermacher/YamshadowInex12_Experiment24Ognoexperiment27-GGUF",
"mradermacher/gemma-2b-datascience-instruct-v4.5-GGUF",
"mradermacher/Josiefied-Qwen2-1.5B-Instruct-abliterated-GGUF",
"mradermacher/zoyllm-7b-slimorca-GGUF",
"mradermacher/M7Yamshadowexperiment28_T3qm7xExperiment28-GGUF",
"mradermacher/Aria-7b-128k-v4-GGUF",
"mradermacher/T3qm7xNeuralsirkrishna_CalmeInex12-GGUF",
"mradermacher/T3qm7xNeuralsirkrishna_YamStrangemerges_32-GGUF",
"mradermacher/YamshadowInex12_T3qm7xExperiment28-GGUF",
"mradermacher/Model_Stock_Mixv0.1-7B-GGUF",
"mradermacher/StrangeMerges_49-7B-dare_ties-GGUF",
"mradermacher/galland-mistral-7b-instruct-v0.2-v1-GGUF",
"mradermacher/Magiq-Translator-0-GGUF",
"mradermacher/occiglot-7b-de-en-instruct-GGUF",
"mradermacher/M7Yamshadowexperiment28_YamShadow-GGUF",
"mradermacher/StrangeMerges_17-7B-dare_ties-GGUF",
"mradermacher/occiglot-7b-es-en-instruct-hf-GGUF",
"mradermacher/ko-gemma-1.1-2b-GGUF",
"mradermacher/M7Yamshadowexperiment28_Percival_01Yamshadow-GGUF",
"mradermacher/TinyAlpaca-1.1B-GGUF",
"mradermacher/MeliodasT3q_M7Alloyingotneoy-GGUF",
"mradermacher/Vistral-7B-Chat-function-calling-GGUF",
]
# 去重(保持原有顺序)
_seen = set()
_deduplicated = []
for _mid in ALL_MODEL_IDS:
if _mid not in _seen:
_deduplicated.append(_mid)
_seen.add(_mid)
ALL_MODEL_IDS = _deduplicated
print(f"[INFO] 去重后模型数量: {len(ALL_MODEL_IDS)}", flush=True)
HEADERS = {"Content-Type": "application/json"}
# ══════════════════════════════════════════════════════════
# 全局状态(供 /status 展示)
# ══════════════════════════════════════════════════════════
_state = {
"strategy_id": STRATEGY_ID,
"phase": "starting", # starting | running | done | error
"total": len(ALL_MODEL_IDS),
"downloading": [], # 当前正在下载的模型
"download_success": 0,
"download_failed": 0,
"submitted": 0, # 成功提交的验证任务数hygon + bi150
"submit_failed": 0,
"started_at": None,
"finished_at": None,
}
_shutdown = threading.Event()
# ══════════════════════════════════════════════════════════
# HTTP 服务
# ══════════════════════════════════════════════════════════
class Handler(BaseHTTPRequestHandler):
def do_GET(self):
if self.path == "/health":
self._json({"status": "ok"})
elif self.path == "/status":
self._json(_state)
else:
self._json({"error": "not found"}, 404)
def _json(self, body: dict, code: int = 200):
payload = json.dumps(body, default=str).encode()
self.send_response(code)
self.send_header("Content-Type", "application/json")
self.send_header("Content-Length", str(len(payload)))
self.end_headers()
self.wfile.write(payload)
def log_message(self, fmt, *args):
print(f"[http] {self.address_string()} {fmt % args}", flush=True)
def _run_http():
server = ThreadingHTTPServer((HTTP_HOST, HTTP_PORT), Handler)
server.timeout = 1
print(f"[http] 监听 {HTTP_HOST}:{HTTP_PORT}", flush=True)
while not _shutdown.is_set():
server.handle_request()
server.server_close()
print("[http] 已关闭", flush=True)
# ══════════════════════════════════════════════════════════
# 工具函数:生成模型文件名
# ══════════════════════════════════════════════════════════
def get_model_filename(model_id: str) -> str:
"""
model_id 生成标准 GGUF 模型文件名
规则
- 移除组织名/ 前部分
- 处理 '_-_' 分割保留原有逻辑
- 移除末尾 '-GGUF'不区分大小写
- 若移除后以 -i1, -i2, ..., -i99 结尾
替换为 .i1, .i2, ... 并添加 '-Q4_0.gguf'
否则
直接添加 '.f16.gguf'
"""
base_name = model_id.split("/")[-1]
if '_-_' in base_name:
base_name = base_name.split('_-_')[-1]
if base_name.lower().endswith("-gguf"):
base_name = base_name[:-5]
match = re.search(r'-i(\d+)$', base_name)
if match:
number = match.group(1)
base_name = base_name[:match.start()] + f".i{number}-Q4_0.gguf"
else:
base_name = base_name + ".f16.gguf"
return base_name
# ══════════════════════════════════════════════════════════
# 登录获取 token
# ══════════════════════════════════════════════════════════
def login(account: str = USER_ACCOUNT, password: str = USER_PASSWORD) -> str:
payload = {"userAccount": account, "userPassword": password}
print(f"🔑 正在登录: {account} ...", flush=True)
resp = requests.post(BASE_URL + LOGIN_ENDPOINT, headers=HEADERS, json=payload, timeout=15)
if resp.status_code != 200:
raise Exception(f"HTTP 登录失败: {resp.text}")
data = resp.json()
if data.get("code") != 0:
raise Exception(f"业务登录失败: {data.get('message')}")
token = data["data"]["token"]
print(f"✅ 登录成功: {account}", flush=True)
return token
# ══════════════════════════════════════════════════════════
# 创建单个模型的下载任务(只发一次请求,由外部控制重试)
# ══════════════════════════════════════════════════════════
def create_download_task_once(token: str, model_id: str):
filename = get_model_filename(model_id)
auth_headers = {**HEADERS, "Authorization": f"Bearer {token}"}
payload = {
"allowPatterns": [filename],
"hfToken": HF_TOKEN,
"modelId": model_id,
"source": "HUGGING_FACE",
"stillDownloadAlreadySuccessDownloadedModel": False
}
print(f"📥 创建下载任务: {model_id}{filename}", flush=True)
try:
resp = requests.post(BASE_URL + CREATE_DOWNLOAD_TASK_ENDPOINT, headers=auth_headers, json=payload, timeout=15)
if resp.status_code == 200:
data = resp.json()
if data.get("code") == 0:
print(f"✅ 下载任务已提交: {model_id}", flush=True)
return True, "SUCCESS"
else:
message = data.get("message", "Unknown Error")
print(f"⚠️ 下载任务业务失败 ({model_id}): {message}", flush=True)
return False, message
else:
print(f"❌ HTTP 错误 ({model_id}): {resp.status_code} - {resp.text}", flush=True)
return False, "HTTP_ERROR"
except Exception as e:
print(f"💥 创建下载任务异常 ({model_id}): {e}", flush=True)
return False, "EXCEPTION"
# ══════════════════════════════════════════════════════════
# 查询单个模型的最新下载任务状态
# ══════════════════════════════════════════════════════════
def check_model_status(token: str, model_id: str) -> str:
"""
返回状态: 'WAITING', 'RUNNING', 'SUCCESS', 'FAILED', 'UNKNOWN'
"""
url = BASE_URL + CREATE_DOWNLOAD_TASK_ENDPOINT
auth_headers = {**HEADERS, "Authorization": f"Bearer {token}"}
params = {"modelId": model_id, "current": 1, "pageSize": 1}
try:
resp = requests.get(url, headers=auth_headers, params=params, timeout=10)
if resp.status_code != 200:
return "UNKNOWN"
data = resp.json()
if data.get("code") != 0:
return "UNKNOWN"
records = data.get("data", {}).get("records", [])
if not records:
return "UNKNOWN"
status = records[0].get("status", "UNKNOWN").upper()
return status
except Exception as e:
print(f"⚠️ 查询状态异常 ({model_id}): {e}", flush=True)
return "UNKNOWN"
# ══════════════════════════════════════════════════════════
# 提交单个模型的测试任务 hygon
# ══════════════════════════════════════════════════════════
def submit_test_task(token: str, model_id: str) -> bool:
auth_headers = {**HEADERS, "Authorization": f"Bearer {token}"}
model_filename = get_model_filename(model_id)
gpu_type = "hygon_k100-ai"
config_content = f"""docker_image: git.modelhub.org.cn:9443/enginex-hygon/hygon-llama.cpp:b7516
nv_docker_image: harbor-contest.4pd.io/luxinlong02/llama-cpp:b7003-cuda-full-12.3
framework: llamacpp
storage: gpfs
modelhub_options:
srcRelativePath: leaderboard/modelHubXC/{model_id}
mountPoint: /model
api: completion
temperature: 0
repetition_penalty: 1.1
top_p: 0.9
max_model_len: 4096
sut_config:
gpu_num: 1
values:
command: ['/app/llama-server','--model', '/model/{model_filename}', '--alias', 'llm', '--threads', '20','--n-gpu-layers', '999', '--prio', '3', '--min_p', '0.01', '--ctx-size', '4096', '--host', '0.0.0.0', '--port', '8000', '--jinja', '--flash-attn', 'off']
ref_config:
gpu_num: 1
values:
command: ['/workspace/llama.cpp/build/bin/llama-server','--model', '/model/{model_filename}', '--alias', 'llm', '--threads', '20', '--n-gpu-layers', '999', '--prio', '3', '--min_p', '0.01', '--ctx-size', '4096', '--host', '0.0.0.0', '--port', '8000', '--jinja', '--flash-attn', 'off']
"""
task_data = {
"contestApiToken": CONTEST_API_TOKEN,
"contributors": CONTRIBUTORS,
"gpuTypes": [gpu_type],
"taskType": TASK_TYPE,
"modelId": model_id,
"strategyId": STRATEGY_ID, # 平台要求
"submissionConfig": [{
"config": config_content,
"gpuType": gpu_type,
"taskType": TASK_TYPE
}]
}
print(f"📤 提交测试任务 (hygon): {model_id}", flush=True)
try:
resp = requests.post(BASE_URL + SUBMIT_TEST_TASK_ENDPOINT, json=task_data, headers=auth_headers, timeout=15)
if resp.status_code == 200:
result = resp.json()
if result.get("code") == 0:
task_id = result.get("data", {}).get("taskId")
print(f"✅ 测试任务提交成功! Task ID: {task_id}", flush=True)
return True
else:
print(f"❌ 测试任务业务错误: {result.get('message')}", flush=True)
return False
else:
print(f"❌ 测试任务 HTTP 错误: {resp.status_code} - {resp.text}", flush=True)
return False
except Exception as e:
print(f"💥 提交测试任务异常 ({model_id}): {e}", flush=True)
return False
# ══════════════════════════════════════════════════════════
# 提交单个模型的测试任务 bi150
# ══════════════════════════════════════════════════════════
def submit_test_task_bi150(token: str, model_id: str) -> bool:
auth_headers = {**HEADERS, "Authorization": f"Bearer {token}"}
model_filename = get_model_filename(model_id)
gpu_type = "Iluvatar_bi-150"
config_content = f"""docker_image: git.modelhub.org.cn:9443/enginex-iluvatar/iluvatar-llama.cpp:b7516-bi150
nv_docker_image: harbor-contest.4pd.io/luxinlong02/llama-cpp:b7003-cuda-full-12.3
framework: llamacpp
storage: gpfs
modelhub_options:
srcRelativePath: leaderboard/modelHubXC/{model_id}
mountPoint: /model
max_model_len: 4096
sut_config:
gpu_num: 1
values:
command: ['/app/llama-server','--model', '/model/{model_filename}', '--alias', 'llm', '--threads', '20','--n-gpu-layers','128', '--ctx-size', '4096', '--host', '0.0.0.0', '--port', '8000', '--jinja', '--flash-attn', 'off', '--no-mmap', '--sync-to-temp']
ref_config:
gpu_num: 1
values:
command: ['/workspace/llama.cpp/build/bin/llama-server','--model', '/model/{model_filename}', '--alias', 'llm', '--threads', '20','--n-gpu-layers','128', '--ctx-size', '4096', '--host', '0.0.0.0', '--port', '8000', '--jinja', '--flash-attn', 'off']
"""
task_data = {
"contestApiToken": CONTEST_API_TOKEN,
"contributors": CONTRIBUTORS,
"gpuTypes": [gpu_type],
"taskType": TASK_TYPE,
"modelId": model_id,
"strategyId": STRATEGY_ID, # 平台要求
"submissionConfig": [{
"config": config_content,
"gpuType": gpu_type,
"taskType": TASK_TYPE
}]
}
print(f"📤 提交测试任务 (bi150): {model_id}", flush=True)
try:
resp = requests.post(BASE_URL + SUBMIT_TEST_TASK_ENDPOINT, json=task_data, headers=auth_headers, timeout=15)
if resp.status_code == 200:
result = resp.json()
if result.get("code") == 0:
task_id = result.get("data", {}).get("taskId")
print(f"✅ 测试任务提交成功! Task ID: {task_id}", flush=True)
return True
else:
print(f"❌ 测试任务业务错误: {result.get('message')}", flush=True)
return False
else:
print(f"❌ 测试任务 HTTP 错误: {resp.status_code} - {resp.text}", flush=True)
return False
except Exception as e:
print(f"💥 提交测试任务异常 ({model_id}): {e}", flush=True)
return False
# ══════════════════════════════════════════════════════════
# 业务逻辑:动态流水线(下载 → 提交验证任务,带并发限制重试)
# ══════════════════════════════════════════════════════════
def _run_worker():
_state["started_at"] = datetime.utcnow().isoformat()
_state["phase"] = "running"
try:
token = login()
except Exception as e:
print(f"[worker] 登录失败: {e}", flush=True)
_state["phase"] = "error"
return
pending_models = list(ALL_MODEL_IDS) # 尚未开始下载的模型
active_models: Set[str] = set() # 当前正在下载的模型
completed_results = {} # model_id -> status
print(f"🚀 总共 {len(pending_models)} 个模型待下载。最大并发数: {MAX_CONCURRENT_DOWNLOADS}\n", flush=True)
while (pending_models or active_models) and not _shutdown.is_set():
# 1. 检查并更新活跃任务的状态
for model_id in list(active_models):
if _shutdown.is_set():
break
status = check_model_status(token, model_id)
if status in ("SUCCESS", "FAILED"):
active_models.remove(model_id)
completed_results[model_id] = status
print(f"⏹️ {model_id} 完成,状态: {status} | 活跃数: {len(active_models)}", flush=True)
if status == "SUCCESS":
_state["download_success"] += 1
# 下载成功后立即提交该模型的验证任务
if submit_test_task(token, model_id):
_state["submitted"] += 1
print(f"🧪 已为 {model_id} 提交 hygon 验证任务", flush=True)
else:
_state["submit_failed"] += 1
print(f"⚠️ {model_id} hygon 验证任务提交失败", flush=True)
if submit_test_task_bi150(token, model_id):
_state["submitted"] += 1
print(f"🧪 已为 {model_id} 提交 bi150 验证任务", flush=True)
else:
_state["submit_failed"] += 1
print(f"⚠️ {model_id} bi150 验证任务提交失败", flush=True)
else:
_state["download_failed"] += 1
# 2. 提交新任务(遇到并发限制原地等待重试,其他错误丢弃)
while len(active_models) < MAX_CONCURRENT_DOWNLOADS and pending_models and not _shutdown.is_set():
next_model = pending_models[0] # 只看队首,不先弹出
success = False
while not success and next_model in pending_models and not _shutdown.is_set():
success, message = create_download_task_once(token, next_model)
if success:
pending_models.pop(0)
active_models.add(next_model)
print(f"▶️ 启动下载: {next_model} (当前活跃: {len(active_models)})", flush=True)
break
else:
if "用户最多同时运行 8 个下载任务" in message:
print(f"⏳ 检测到并发限制,等待 {RETRY_WAIT_SECONDS} 秒后重试该模型...", flush=True)
_shutdown.wait(RETRY_WAIT_SECONDS)
else:
print(f"❌ 非并发错误,丢弃该模型: {next_model}", flush=True)
pending_models.pop(0)
completed_results[next_model] = "DISCARDED"
_state["download_failed"] += 1
break
_state["downloading"] = sorted(active_models)
# 3. 稍作等待,避免频繁查询(可被 shutdown 信号打断)
if active_models or pending_models:
_shutdown.wait(CHECK_INTERVAL_SECONDS)
print("\n✅ 所有模型处理完毕!\n", flush=True)
# 4. 收集所有成功下载的模型ID并写入结果文件
success_models = [
mid for mid, status in completed_results.items()
if status == "SUCCESS"
]
try:
with open("downloaded_success_models.txt", "w", encoding="utf-8") as f:
for mid in success_models:
f.write(f"{mid}\n")
except Exception:
pass
print("🎉 下载成功的模型ID列表", flush=True)
print("[", flush=True)
for mid in success_models:
print(f' "{mid}",', flush=True)
print("]", flush=True)
_state["downloading"] = []
_state["finished_at"] = datetime.utcnow().isoformat()
_state["phase"] = "done"
print(
f"[worker] 完成 download_success={_state['download_success']} "
f"download_failed={_state['download_failed']} "
f"submitted={_state['submitted']} submit_failed={_state['submit_failed']}",
flush=True,
)
# 流水线完成后继续保持进程存活,等待平台停止
# ══════════════════════════════════════════════════════════
# 入口
# ══════════════════════════════════════════════════════════
def _handle_signal(signum, _frame):
print(f"[main] 收到信号 {signum},正在关闭...", flush=True)
_shutdown.set()
def main():
signal.signal(signal.SIGTERM, _handle_signal)
signal.signal(signal.SIGINT, _handle_signal)
# HTTP 服务线程
http_thread = threading.Thread(target=_run_http, daemon=False)
http_thread.start()
# 流水线线程
worker_thread = threading.Thread(target=_run_worker, daemon=True)
worker_thread.start()
# 主线程等待 shutdown
_shutdown.wait()
print("[main] 等待 HTTP 服务关闭...", flush=True)
http_thread.join(timeout=5)
print("[main] 退出", flush=True)
if __name__ == "__main__":
main()