diff --git a/README.md b/README.md index 6f84e44..2ed546e 100644 --- a/README.md +++ b/README.md @@ -1,24 +1,23 @@ -# xc_validation_strategy_gguf +# xc_validation_strategy_vllm_submit -GGUF 模型下载 + 验证任务提交流水线策略服务:批量创建 GGUF 模型下载任务(最大并发 8), -每个模型下载成功后立即提交 hygon / bi150 两个验证任务,之后保持 HTTP 服务存活供平台探活。 +批量向 ModelHub XC 平台提交 vLLM 模型适配任务的策略服务(fanyi 账号), +之后保持 HTTP 服务存活供平台探活。 ## 功能 -- 自动登录 ModelHub 获取 Token(失败时回退到预设 Token) -- 按流水线批量创建 GGUF 模型下载任务(HuggingFace 源,最大并发 8) -- 每个模型下载成功后,立即提交 hygon_k100-ai 与 Iluvatar_bi-150 两个验证任务(llamacpp 框架) -- 下载成功的模型 ID 写入 `downloaded_success_models.txt` +- 通过 `/api/adapt/task/add` 接口(xc-Token 认证)批量提交模型适配任务 +- vLLM 框架,ppu_zw_810e GPU +- 提交成功的模型写入 `submitted_adapt_tasks.txt` - 暴露 `/health` 和 `/status` 接口满足平台运行时契约 ## 项目结构 ``` . -├── main.py # 主入口:HTTP 服务 + 下载/提交流水线 +├── main.py # 主入口:HTTP 服务 + 提交逻辑 ├── Dockerfile # 平台镜像构建配置 ├── requirements.txt # Python 依赖 -└── downloaded_success_models.txt # 运行后自动生成,记录下载成功的模型 +└── submitted_adapt_tasks.txt # 运行后自动生成,记录提交结果 ``` ## 平台契约说明 diff --git a/main.py b/main.py index 5ee3ecb..c7f4f6a 100644 --- a/main.py +++ b/main.py @@ -1,47 +1,41 @@ """ -xc_validation_strategy_gguf — 主入口 +xc_validation_strategy_vllm_submit — 主入口 -GGUF 模型下载 + 验证任务提交流水线(部署框架与 xc_validation_strategy 一致)。 - -启动后运行流水线:批量创建 GGUF 模型下载任务(最大并发 8), -每个模型下载成功后立即提交 hygon / bi150 两个验证任务; +启动后通过 /api/adapt/task/add 接口(xc-Token 认证)批量提交 +vLLM 模型适配任务(ppu_zw_810e),之后保持 HTTP 服务存活。 同时暴露 /health(K8s 探活)和 /status(运行状态)。 + +部署框架与 xc_validation_strategy 一致。 """ import json import os -import re import signal import threading from datetime import datetime from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer -from typing import Set +from typing import List 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" +BASE_URL = os.environ.get("BASE_URL", "https://modelhub.org.cn") +ADD_TASK_ENDPOINT = "/api/adapt/task/add" -# 登录账号:启动时优先用账号密码换取新 token(流水线运行时间长,预设 token 可能过期) -USER_ACCOUNT = "zhoushasha@4paradigm.com" -USER_PASSWORD = "ganshenme0" +# fanyi 账号的 xc-Token(该接口使用 xc-Token 认证,无需登录) +USER_ACCOUNT = "fanyi" +XC_TOKEN = "f2d501c9ae6543a589cd6cb789108c41" -# 通过 curl -X POST https://modelhub.org.cn/adminApi/user/login 获取后填入(登录失败时的回退) -AUTH_TOKEN = "eyJ0eXAiOiJKV1QiLCJhbGciOiJIUzI1NiJ9.eyJ1c2VyQWNjb3VudCI6Inpob3VzaGFzaGEiLCJpZCI6MTQsInVzZXJSb2xlIjoibGVhZGVyYm9hcmQiLCJleHAiOjE3ODQ1NDc1NDYsImlhdCI6MTc4Mzk0Mjc0Nn0.ZcOqcrfI22LPi4mGMnt164nZGhi61ZxtJGYsoO7fZdM" +GPU_TYPE = "ppu_zw_810e" +TASK_TYPE = "text-generation" +STRATEGY_ID = os.environ.get("STRATEGY_ID", "") # 平台自动注入,无需修改 -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 # 下载状态轮询间隔(秒) +HEADERS = { + "Content-Type": "application/json", + "xc-Token": XC_TOKEN, +} HTTP_HOST = "0.0.0.0" HTTP_PORT = 8080 @@ -51,75 +45,83 @@ HTTP_PORT = 8080 # ══════════════════════════════════════════════════════════ ALL_MODEL_IDS = [ - - - "mradermacher/Qwen2.5-3B-instruct-argus-v3-GGUF", - "mradermacher/NeuralsirkrishnaShadow_OgnoExperiment27-GGUF", - "mradermacher/GREEN-RadLlama2-7b-GGUF", - "mradermacher/snakmodel-7b-instruct-GGUF", - "mradermacher/BioMistral-7B-Starling-SLERP-GGUF", - "mradermacher/Marco-Llama-3.2-3B-GGUF", - "mradermacher/Qwen2.5-1.5B-Instruct-Open-R1-GRPO-GGUF", - "mradermacher/Qwen2.5-0.5B-Distill-Fast-GGUF", - "mradermacher/chomsky_16_bit_model-GGUF", - "mradermacher/Mistral-7B-DFT-GGUF", - "mradermacher/mergekit-task_arithmetic-qjeuqjw-GGUF", - "mradermacher/M7Yamshadowexperiment28_Experiment27Inex12-GGUF", - "mradermacher/YamshadowStrangemerges_32_Experiment28Inex12-GGUF", - "mradermacher/MeliodasPercival_01_Experiment29Pastiche-GGUF", - "mradermacher/Llama-3-6B-v0-GGUF", - "mradermacher/Excalibur-7b-DPO-GGUF", - "mradermacher/Llama-2-7b-Indian-Law-GGUF", - "mradermacher/StarlingHermes-2.5-Mistral-7B-slerp-GGUF", - "mradermacher/DeepThinker-7B-Sce-v1-GGUF", - "mradermacher/Alif-Llama-EXP2-GGUF", - "mradermacher/qwen-2.5-1.5B-Rasa-GGUF", - "mradermacher/German_RAG-PHI-3.5-MINI-4B-MERGED-HESSIAN-AI-GGUF", - "mradermacher/MFANN-phigments-slerp-V3.2-GGUF", - "mradermacher/AdityaGPT-GGUF", - "mradermacher/Jaja-small-v4-GGUF", - "mradermacher/Jaja-medium-v1-GGUF", - "mradermacher/astrollama-2-7b-base_abstract-GGUF", - "mradermacher/Llasagna-v0.1-GGUF", - "mradermacher/DeepThinker-v-GGUF", - "mradermacher/Jaja-small-v3-GGUF", - "mradermacher/Llama-3.2-1B-FC-v1.2-think-GGUF", - "mradermacher/Deepseek-Qwen2.5-1.5B-Redistil-GGUF", - - - - - - + "Hyeji0101/qwen2_5_1_5b_demo", + "GenueAI/geode-onyx", + "GM77/qwen3-4b-verilog-grpo", + "ChuGyouk/F_R13_T2", + "ChuGyouk/R17", + "Ingingdo/bit-0.5b-final-logic", + "beomi/Llama-3-Open-Ko-8B", + "Fiscus/trinitite_safe_rl_base_model", + "ChuGyouk/F_R12_T3", + "ChuGyouk/F_R12_T2", + "xw1234gan/cnk12_Main_fixed_SFTanchor_3B_step_9", + "xw1234gan/cnk12_Main_fixed_SFTanchor_3B_step_10", + "xw1234gan/cnk12_Main_fixed_BaseAnchor_3B_step_1", + "kmseong/llama3_2_3b-instruct-math-safedelta-scale0.99", + "opencompass/anah-v2", + "ChuGyouk/R14", + "trishajean/qwen-math-cebuano-1.5b-merged", + "GyanAISystems/Gyan-AI-G1-Official", + "Divij/Qwen2.5-3B-Instruct-sft-without-thoughts", + "Divij/Qwen2.5-3B-Instruct-sft-with-thoughts", + "ChuGyouk/R5_1", + "ChuGyouk/R18_1", + "ChuGyouk/R19_1", + "ChuGyouk/R12", + "ChuGyouk/F_R11_T4", + "ChuGyouk/F_R12", + "ChuGyouk/F_R11_T2", + "ChuGyouk/F_R11_T3", + "ChuGyouk/F_R13_1_T1", + "ChuGyouk/F_R12_T4", + "automerger/T3qm7xNeuralsirkrishna-7B", + "Ford91/clifford-ai-v2", + "ChuGyouk/R16_1", + "ChuGyouk/R15_1", + "nkatara/gita-text-generation-gpt2", + "HINT-lab/Qwen2.5-7B-Instruct-Self-Calibration", + "thirdeyeai/Qwen2.5-1.5B-Instruct-uncensored", + "karaselerm/qwen2.5-1.5b-instruct-ru-abliterated-hw6", + "xw1234gan/cnk12_Main_fixed_BaseAnchor_3B_step_2", + "ontocord/wide_3b_sft_stage1.1-ss1-with_intr_math.no_issue", + "mncai/Foundation_Law_epoch4", + "gauri0508/med-record-audit-qwen2.5-3b-grpo", + "unsloth/Phi-4-mini-instruct", + "E-motionAssistant/qwen-2.5-3b-tamil-therapy-merged", + "EscapeJeju/qwen2_5_1_5b_demo", + "AgPerry/Qwen3-8B-fim-v2v3pt-swe-lego-posttrain", + "ChuGyouk/F_R11", + "ChuGyouk/F_R11_1_T1", + "LorenaYannnnn/general_reward-Qwen3-0.6B-OURS_self-seed_1", + "Vortex5/Crimson-Constellation-12B", + "cloudyu/mistral_11B_instruct_v0.1", + "pkupie/Qwen2.5-3B-ug-cpt", + "iproskurina/qwen-hf-fewshot-iter-np-iter3", + "ontocord/wide_3b_sft_stage1.2-ss1-expert_wiki", + "kmseong/llama3_2_3b-instruct-math-safedelta-scale2", + "Thrillcrazyer/Qwen-2.5-1.5B_TAC_Teacher_Qwen32B", + "nyu-dice-lab/VeriThoughts-Reasoning-7B", + "ontocord/wide_3b", + "silvercoder67/Mistral-7b-instruct-v0.2-summ-sft-e2m", + "lihaoxin2020/qwen3-4b-sft-gpt54-ep2-instance-rubric-gpt54-step200", + "lihaoxin2020/qwen3-4b-sft-gpt54-ep2-instance-rubric-gpt54-step150", + "lihaoxin2020/qwen3-4b-sft-gpt54-ep2-evolving-rubric-gem3-flash-step150", + "Guilherme34/Firefly-V3", + ] - -# 去重(保持原有顺序) -_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, + "strategy_id": STRATEGY_ID, + "phase": "starting", # starting | submitting | done | error + "total": len(ALL_MODEL_IDS), + "submitted": 0, + "failed": 0, + "started_at": None, + "finished_at": None, } _shutdown = threading.Event() @@ -157,336 +159,119 @@ def _run_http(): print("[http] 已关闭", flush=True) # ══════════════════════════════════════════════════════════ -# 工具函数:生成模型文件名 +# 业务逻辑 # ══════════════════════════════════════════════════════════ -def get_model_filename(model_id: str) -> str: - """ - 从 model_id 生成标准 GGUF 模型文件名。 +def submit_task(model_id: str) -> bool: + config_content = f""" +gpu_type: ppu_zw_810e +framework: vllm +docker_image: harbor.4pd.io/hardcore-tech/asllm:1.10.1-pytorch2.10.0-ubuntu24.04-sail2.1.0-cuda13.0-sglang0.5.10-vllm0.19.0-py312 +nv_docker_image: harbor-contest.4pd.io/sunruoxi/vllm-openai-fix-tokenizer:v0.11.0 - 规则: - - 移除组织名(/ 前部分) - - 处理 '_-_' 分割(保留原有逻辑) - - 移除末尾 '-GGUF'(不区分大小写) - - 若移除后以 -i1, -i2, ..., -i99 结尾: - → 替换为 .i1, .i2, ... 并添加 '-Q4_0.gguf' - 否则: - → 直接添加 '.f16.gguf' +sut_config: + values: + gpu_num: 1 + env: + - name: test + value: fp16 + command: + - bash + - /opt/t-head/entrypoint.sh + - python3 + - -m + - asllm.entrypoints.api_server + - --model + - /model + - --port + - '30000' + - --host + - 0.0.0.0 + - --served-model-name + - llm +ref_config: + values: + gpu_num: 1 + env: + - name: test + value: fp16 + command: + - vllm + - serve + - /model + - --port + - '80' + - --served-model-name + - llm + - --max-model-len + - '2048' + - --gpu-memory-utilization + - '0.9' + - --enforce-eager + - --trust-remote-code + - -tp + - '1' +""" - 示例: - 'mradermacher/Qwen3-8B-makisu-v2.0.1-i1-GGUF' - → 'Qwen3-8B-makisu-v2.0.1.i1-Q4_0.gguf' - - 'QuantFactory/Apollo2-9B-GGUF' - → 'Apollo2-9B.f16.gguf' - """ - # 1. 提取模型名部分(/ 后) - base_name = model_id.split("/")[-1] - - # 2. 处理 '_-_' 分割 - if '_-_' in base_name: - base_name = base_name.split('_-_')[-1] - - # 3. 移除末尾的 -GGUF(不区分大小写) - if base_name.lower().endswith("-gguf"): - base_name = base_name[:-5] - - # 4. 检查是否以 -i<数字> 结尾(支持 i1~i99 等) - 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(失败时回退到预设 AUTH_TOKEN) -# ══════════════════════════════════════════════════════════ -def login() -> str: - payload = {"userAccount": USER_ACCOUNT, "userPassword": USER_PASSWORD} - print("[login] 正在登录...", flush=True) - try: - resp = requests.post(BASE_URL + LOGIN_ENDPOINT, headers=HEADERS, json=payload, timeout=15) - data = resp.json() - if resp.status_code == 200 and data.get("code") == 0: - print("[login] 登录成功", flush=True) - return data["data"]["token"] - print(f"[login] 登录失败: {data.get('message')},回退使用预设 Token", flush=True) - except Exception as e: - print(f"[login] 登录异常: {e},回退使用预设 Token", flush=True) - return AUTH_TOKEN - -# ══════════════════════════════════════════════════════════ -# 创建单个模型的下载任务 -# ══════════════════════════════════════════════════════════ -def create_download_task(token: str, model_id: str) -> bool: - 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 - else: - print(f"⚠️ 下载任务业务失败 ({model_id}): {data.get('message')}", flush=True) - return False - else: - print(f"❌ HTTP 错误 ({model_id}): {resp.status_code} - {resp.text}", flush=True) - return False - except Exception as e: - print(f"💥 创建下载任务异常 ({model_id}): {e}", flush=True) - return False - -# ══════════════════════════════════════════════════════════ -# 查询单个模型的最新下载任务状态 -# ══════════════════════════════════════════════════════════ -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], + "configParams": config_content, + "framework": "vllm", + "modelAddress": f"https://huggingface.co/{model_id}", + "targetGpu": GPU_TYPE, "taskType": TASK_TYPE, - "modelId": model_id, - "strategyId": STRATEGY_ID, # 平台要求 - "submissionConfig": [{ - "config": config_content, - "gpuType": gpu_type, - "taskType": TASK_TYPE - }] + "strategyId": STRATEGY_ID, # 平台要求;若接口不支持该字段会被忽略 } - print(f"📤 提交测试任务 (hygon): {model_id}", flush=True) + + print(f"📤 提交任务: {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 + resp = requests.post( + BASE_URL + ADD_TASK_ENDPOINT, + headers=HEADERS, + json=payload, + timeout=30, + ) + print(f"status: {resp.status_code}", flush=True) + result = resp.json() + print(result, flush=True) + if result.get("code") == 0: + print(f"✅ 提交成功: {model_id}", flush=True) + return True else: - print(f"❌ 测试任务 HTTP 错误: {resp.status_code} - {resp.text}", flush=True) + print(f"❌ 提交失败: {result.get('message')}", flush=True) return False except Exception as e: - print(f"💥 提交测试任务异常 ({model_id}): {e}", flush=True) + 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" + _state["phase"] = "submitting" + successful: List[str] = [] + for model_id in ALL_MODEL_IDS: + if _shutdown.is_set(): + break + if submit_task(model_id): + _state["submitted"] += 1 + successful.append(model_id) + else: + _state["failed"] += 1 + + # 写入结果文件 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"): - completed_results[model_id] = status - active_models.remove(model_id) - print(f"⏹️ {model_id} 完成,状态: {status}", 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. 补充新任务(最多补到 MAX_CONCURRENT_DOWNLOADS 个) - while len(active_models) < MAX_CONCURRENT_DOWNLOADS and pending_models and not _shutdown.is_set(): - next_model = pending_models.pop(0) - if create_download_task(token, next_model): - active_models.add(next_model) - print(f"▶️ 启动下载: {next_model} (当前活跃: {len(active_models)})", flush=True) - else: - # 创建失败也视为完成(避免卡住) - completed_results[next_model] = "CREATE_FAILED" - _state["download_failed"] += 1 - print(f"❌ 创建失败: {next_model}", flush=True) - - _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: + with open("submitted_adapt_tasks.txt", "w", encoding="utf-8") as f: + for mid in successful: 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']}", + f"[worker] 完成 submitted={_state['submitted']} failed={_state['failed']} total={_state['total']}", flush=True, ) - # 流水线完成后继续保持进程存活,等待平台停止 + # 提交完成后继续保持进程存活,等待平台停止 # ══════════════════════════════════════════════════════════ # 入口 @@ -504,7 +289,7 @@ def main(): http_thread = threading.Thread(target=_run_http, daemon=False) http_thread.start() - # 流水线线程 + # 提交任务线程 worker_thread = threading.Thread(target=_run_worker, daemon=True) worker_thread.start()