8 Commits

2 changed files with 614 additions and 82 deletions

View File

@@ -4,6 +4,8 @@ ENV PYTHONUNBUFFERED=1
WORKDIR /app WORKDIR /app
RUN mkdir -p /app/data
COPY requirements.txt . COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt RUN pip install --no-cache-dir -r requirements.txt

696
main.py
View File

@@ -1,18 +1,26 @@
""" """
ModelHub Cancel-All 智能体 ModelHub 全云端提交智能体
启动后自动取消当前账号所有 waiting/running 状态的验证任务 部署在平台容器中,直接从内网提交验证任务
用法:推送此版本到平台,选择对应 tag 运行即可
功能:
- POST /run → 触发一次完整流程(搜索→筛选→提交)
- GET /status → 查看当前状态
- GET /health → 健康检查
""" """
import json import json
import os import os
import signal import signal
import time import sqlite3
import threading import threading
import time
import traceback
from datetime import datetime from datetime import datetime
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from urllib.parse import urlparse, parse_qs
import requests import requests
import yaml
# ============================================================ # ============================================================
# 配置 # 配置
@@ -20,19 +28,37 @@ import requests
HOST = "0.0.0.0" HOST = "0.0.0.0"
PORT = 8080 PORT = 8080
STRATEGY_ID = os.getenv("STRATEGY_ID", "")
# 目标GPU和Token跟提交版本保持一致 # 目标GPU
TARGET_GPU = "Iluvatar_bi-150" TARGET_GPU = "Iluvatar_bi-100"
# 账号Token
TARGET_TOKEN = "f45f1aae2c094426be237c88b1085015" TARGET_TOKEN = "f45f1aae2c094426be237c88b1085015"
# 架构白名单
SUPPORTED_ARCH_KEYWORDS = ['Qwen', 'Qwen2', 'Qwen3']
SUPPORTED_MODEL_TYPES = [
'qwen', 'qwen2', 'qwen2_vl', 'qwen2_5_vl', 'qwen2_audio',
'qwen3', 'qwen3_vl', 'qwen3_5', 'qwen3_5_moe',
]
SUPPORTED_SPECIAL_ARCHS = ['Eagle3Speculator', 'LlamaForCausalLMEagle3']
MODELHUB_API = "https://modelhub.org.cn/api" MODELHUB_API = "https://modelhub.org.cn/api"
BATCH_SIZE = 50 # 每批取消数量API上限
# 搜索关键词
SEARCH_KEYWORDS = ['Llama-3', 'Llama-3.1', 'Llama-3.2', 'Meta-Llama', 'Llama-4']
# ============================================================ # ============================================================
# 全局状态 # 全局状态
# ============================================================ # ============================================================
state = {'running': False, 'last_run': None, 'logs': []} state = {
'running': False,
'last_run': None,
'last_result': None,
'logs': [],
}
state_lock = threading.Lock() state_lock = threading.Lock()
@@ -46,76 +72,375 @@ def log(msg: str):
state['logs'] = state['logs'][-300:] state['logs'] = state['logs'][-300:]
def cancel_all_tasks(): # ============================================================
"""取消所有 waiting 和 running 状态的任务""" # 数据库(内存 SQLite
# ============================================================
db_conn = None
def init_db():
global db_conn
os.makedirs('/app/data', exist_ok=True)
db_conn = sqlite3.connect('/app/data/submit_history.db', check_same_thread=False)
db_conn.execute('''CREATE TABLE IF NOT EXISTS submitted (
model_id TEXT, gpu TEXT, task_id TEXT, status TEXT,
submitted_at TEXT, checked_at TEXT,
PRIMARY KEY(model_id, gpu)
)''')
db_conn.execute('''CREATE TABLE IF NOT EXISTS failed (
model_id TEXT, gpu TEXT, reason TEXT, failed_at TEXT,
PRIMARY KEY(model_id, gpu)
)''')
db_conn.commit()
def is_model_failed(model_id: str) -> bool:
"""检查模型是否已知失败"""
if db_conn:
row = db_conn.execute(
'SELECT 1 FROM failed WHERE model_id=? AND gpu=?',
(model_id, TARGET_GPU)
).fetchone()
return row is not None
return False
def record_failed(model_id: str, reason: str):
"""记录失败的模型"""
if db_conn:
db_conn.execute(
'INSERT OR REPLACE INTO failed VALUES (?,?,?,?)',
(model_id, TARGET_GPU, reason, datetime.now().isoformat())
)
db_conn.commit()
# ============================================================
# ModelScope 搜索
# ============================================================
MODELSCOPE_API = "https://modelscope.cn/api/v1"
DOWNLOAD_MIN = 50
DOWNLOAD_MAX = 5000
SEARCH_PAGES = 5 # 每个关键词搜5页50*5=250个结果
def search_models(keyword: str) -> list:
"""从 ModelScope 搜索模型多页筛选下载量50-5000的冷门模型"""
url = "https://modelscope.cn/openapi/v1/models"
models = []
for page in range(1, SEARCH_PAGES + 1):
params = {
'search': keyword,
'page_size': 50,
'page_number': page,
'sort': 'downloads',
}
try:
resp = requests.get(url, params=params, timeout=20,
headers={'User-Agent': 'Mozilla/5.0'})
data = resp.json()
if data.get('success'):
page_models = data.get('data', {}).get('models', [])
for m in page_models:
dl = m.get('downloads', 0)
if DOWNLOAD_MIN <= dl <= DOWNLOAD_MAX:
models.append({'id': m.get('id'), 'downloads': dl})
if len(page_models) < 50:
break # 最后一页,不继续
else:
break
except Exception as e:
log(f" [{keyword}] page={page}: {e}")
break
time.sleep(0.3)
log(f" [{keyword}]: {len(models)} 个 (50<={DOWNLOAD_MAX})")
return models
def check_architecture(model_id: str) -> tuple:
"""检查模型架构"""
try:
cfg_url = f"{MODELSCOPE_API}/models/{model_id}/repo?Revision=master&FilePath=config.json"
resp = requests.get(cfg_url, timeout=10)
if resp.status_code == 200:
cfg = resp.json()
archs = cfg.get('architectures', [])
mtype = cfg.get('model_type', '')
arch_str = str(archs)
for special in SUPPORTED_SPECIAL_ARCHS:
if special in arch_str:
return True, special
for kw in SUPPORTED_ARCH_KEYWORDS:
if kw in arch_str:
return True, kw
if mtype in SUPPORTED_MODEL_TYPES:
return True, mtype
return False, f"arch={archs} type={mtype}"
mid_upper = model_id.upper()
if 'QWEN3' in mid_upper or 'QWEN2' in mid_upper:
return True, "GGUF"
return False, "no config, not Qwen"
except Exception as e:
return True, f"check error: {e}"
def normalize_model_url(model_url: str) -> str:
"""标准化 URL 格式"""
if '/models/' in model_url:
return model_url
if 'modelscope.cn/' in model_url:
parts = model_url.split('modelscope.cn/')
if len(parts) == 2:
return f"https://www.modelscope.cn/models/{parts[1]}"
return model_url
# ============================================================
# 平台 API
# ============================================================
def check_platform_verify(model_id: str) -> dict:
"""查询全平台验证状态"""
headers = {'Xc-Token': TARGET_TOKEN, 'Accept': 'application/json'}
url = f"{MODELHUB_API}/computility/models/search-by-model-id"
try:
resp = requests.get(url, headers=headers, params={'modelId': model_id}, timeout=10)
data = resp.json()
if data.get('code') == 0:
result = data.get('data', {}).get('verifyResult', {})
return result if isinstance(result, dict) else {}
except Exception:
pass
return {}
def check_my_submitted(model_id: str) -> bool:
"""检查自己是否已提交"""
headers = {'Xc-Token': TARGET_TOKEN, 'Accept': 'application/json'}
url = f"{MODELHUB_API}/adapt/task/page"
try:
resp = requests.get(url, headers=headers, params={
'current': 1, 'pageSize': 100, 'onlyMine': 'true',
'gpuType': TARGET_GPU, 'modelId': model_id,
}, timeout=10)
data = resp.json()
if data.get('code') == 0:
records = data.get('data', {}).get('records', [])
for r in records:
if r.get('modelId') == model_id and r.get('gpuType') == gpu:
status = r.get('status', '')
if status in ('waiting', 'running', 'success'):
return True
except Exception:
pass
return False
def check_queue_available() -> int:
"""查询队列可用位置"""
headers = {'Xc-Token': TARGET_TOKEN, 'Accept': 'application/json'}
url = f"{MODELHUB_API}/adapt/task/page"
try:
resp = requests.get(url, headers=headers, params={
'current': 1, 'pageSize': 1, 'onlyMine': 'true',
'gpuType': TARGET_GPU, 'status': 'waiting',
}, timeout=10)
data = resp.json()
if data.get('code') == 0:
waiting = int(data['data'].get('total', 0))
return max(0, 100 - waiting)
except Exception:
pass
return -1
def build_config_params() -> str:
"""构建 YAML 配置 - vllm"""
params = {
'framework': 'vllm',
'nv_framework': 'vllm',
'api': 'completion',
'max_tokens': 1024,
'temperature': 0.7,
'repetition_penalty': 1.2,
'top_p': 0.9,
'lang': 'zh',
'max_model_len': 2048,
'sut_config': {
'gpu_num': 1,
'values': {
'command': [
'vllm', 'serve', '/model', '--port', '8000',
'--served-model-name', 'llm', '--max-model-len', '2048',
'--dtype', 'auto', '--gpu-memory-utilization', '0.95',
'-tp', '1', '--enforce-eager', '--trust-remote-code',
]
}
},
'ref_config': {
'gpu_num': 1,
'values': {
'command': [
'vllm', 'serve', '/model', '--port', '80',
'--served-model-name', 'llm', '--max-model-len', '4096',
'--enforce-eager', '--trust-remote-code', '-tp', '1',
]
}
},
}
return yaml.dump(params, default_flow_style=False, allow_unicode=True, width=1000)
def submit_model(model_url: str) -> tuple:
"""提交单个模型"""
headers = { headers = {
'Xc-Token': TARGET_TOKEN, 'Xc-Token': TARGET_TOKEN,
'Accept': 'application/json', 'Accept': 'application/json',
'Content-Type': 'application/json', 'Content-Type': 'application/json',
} }
url = f"{MODELHUB_API}/adapt/task/add"
log("=" * 50) payload = {
log("开始取消所有任务") 'modelAddress': normalize_model_url(model_url),
log(f"目标GPU: {TARGET_GPU}") 'taskType': 'text-generation',
'targetGpu': TARGET_GPU,
# 1. 查询所有任务 'framework': 'vllm',
all_tasks = [] 'strategyId': STRATEGY_ID,
page = 1 'configParams': build_config_params(),
while True: }
try: try:
resp = requests.get( resp = requests.post(url, headers=headers, json=payload, timeout=30)
f"{MODELHUB_API}/adapt/task/page",
headers={'Xc-Token': TARGET_TOKEN, 'Accept': 'application/json'},
params={'current': page, 'pageSize': 100, 'onlyMine': 'true',
'gpuType': TARGET_GPU},
timeout=15
)
data = resp.json()
if data.get('code') != 0:
log(f"查询失败: code={data.get('code')}")
break
records = data['data'].get('records', [])
if not records:
break
all_tasks.extend(records)
total = data['data'].get('total', 0)
log(f"{page}页: {len(records)} 条, 累计 {len(all_tasks)} / {total}")
if len(all_tasks) >= int(total):
break
page += 1
except Exception as e:
log(f" 查询异常: {e}")
return
# 2. 筛选 waiting/running 状态
cancellable = [t for t in all_tasks if t.get('status') in ('waiting', 'running')]
if not cancellable:
log(f" 没有可取消的任务 (总任务 {len(all_tasks)} 个)")
log("=" * 50)
return
task_ids = [int(t.get('taskId')) for t in cancellable if t.get('taskId')]
log(f" 总任务 {len(all_tasks)} 个, 可取消 {len(task_ids)}")
# 3. 分批取消
url = f"{MODELHUB_API}/async/task/stop-create-contest-task"
cancelled = 0
for i in range(0, len(task_ids), BATCH_SIZE):
batch = task_ids[i:i + BATCH_SIZE]
try:
resp = requests.put(url, headers=headers, json={'taskIds': batch}, timeout=15)
data = resp.json() data = resp.json()
if data.get('code') == 0: if data.get('code') == 0:
cancelled += len(batch) task_id = data.get('data', {}).get('id')
log(f" 批次 {i // BATCH_SIZE + 1}: 取消 {len(batch)} 个 OK") return True, task_id, 'success'
else: return False, None, data.get('message', 'unknown error')
log(f" 批次 {i // BATCH_SIZE + 1}: 失败 code={data.get('code')} msg={data.get('message', '')[:80]}")
except Exception as e: except Exception as e:
log(f" 批次 {i // BATCH_SIZE + 1}: 异常 {e}") return False, None, str(e)
time.sleep(1)
log(f" 取消完成: {cancelled} / {len(task_ids)}")
# ============================================================
# 主流程
# ============================================================
def run_pipeline(submit_limit: int = 2):
"""完整流程搜索→筛选→提交只针对目标GPU"""
init_db()
log("=" * 50) log("=" * 50)
log("开始执行流程")
log(f"目标GPU: {TARGET_GPU}")
log(f"提交限制: {submit_limit}")
# 1. 搜索
log("\n--- 阶段1: 搜索 ModelScope ---")
seen = set()
all_models = []
for kw in SEARCH_KEYWORDS:
models = search_models(kw)
for m in models:
mid = m.get('id', '')
if mid and mid not in seen:
seen.add(mid)
all_models.append({
'model_id': mid,
'url': f"https://modelscope.cn/{mid}",
'downloads': m.get('downloads', 0),
})
time.sleep(0.3)
log(f"搜索完成: {len(seen)} 个唯一模型, {len(all_models)} 个下载量{DOWNLOAD_MIN}-{DOWNLOAD_MAX}")
# 2. 格式筛选(排除 GPTQ/AWQ保留 GGUF 和 HuggingFace
log("\n--- 阶段2: 格式筛选 ---")
hf_models = []
format_skipped = 0
SKIP_FORMATS = ['GPTQ', 'AWQ']
for m in all_models:
mid_upper = m['model_id'].upper()
skip = False
for fmt in SKIP_FORMATS:
if fmt in mid_upper:
format_skipped += 1
log(f" x {m['model_id']}: {fmt}格式,跳过")
skip = True
break
if not skip:
hf_models.append(m)
log(f"格式筛选: {len(hf_models)} 通过, {format_skipped} 跳过 (GPTQ/AWQ)")
# 3. 架构筛选只保留有标准config.json的模型排除无效格式
log("\n--- 阶段3: 架构检查 ---")
arch_passed = []
arch_rejected = 0
for m in hf_models:
ok, reason = check_architecture(m['model_id'])
if not ok:
arch_rejected += 1
log(f" x {m['model_id']}: {reason}")
elif reason == 'GGUF':
arch_rejected += 1
log(f" x {m['model_id']}: GGUF(无config)")
else:
arch_passed.append(m)
time.sleep(0.1)
log(f"架构检查: {len(arch_passed)} 通过, {arch_rejected} 拒绝")
# 4. 筛选并提交只针对目标GPU
log(f"\n--- 阶段4: 筛选并提交 [{TARGET_GPU}] ---")
# 检查队列
available = check_queue_available()
if available <= 0:
log(f" 队列满,跳过")
return 0
log(f" 队列可用: {available}")
# 筛选
to_submit = []
for m in arch_passed:
model_id = m['model_id']
# 检查全平台验证状态
verify = check_platform_verify(model_id)
if TARGET_GPU in verify:
continue # 已有记录,跳过
# 检查自己是否已提交
if check_my_submitted(model_id):
continue
# 检查是否已知失败(避免重复提交)
if is_model_failed(model_id):
continue
to_submit.append(m)
if len(to_submit) >= min(submit_limit, available):
break
time.sleep(0.2)
log(f" 待提交: {len(to_submit)}")
# 提交
submitted = 0
for m in to_submit:
ok, task_id, msg = submit_model(m['url'])
if ok:
submitted += 1
log(f"{m['model_id']}")
db_conn.execute(
'INSERT OR REPLACE INTO submitted VALUES (?,?,?,?)',
(m['model_id'], TARGET_GPU, str(task_id), datetime.now().isoformat())
)
else:
log(f"{m['model_id']}: {msg}")
# 永久失败类型记录到 failed 表
if any(kw in str(msg) for kw in ['保护期', '白名单', '唯一性']):
record_failed(m['model_id'], msg)
time.sleep(0.5)
log(f" 提交完成: {submitted}/{len(to_submit)}")
db_conn.commit()
log(f"\n{'=' * 50}")
log(f"流程完成,共提交 {submitted} 个模型")
return submitted
# ============================================================ # ============================================================
@@ -124,21 +449,148 @@ def cancel_all_tasks():
class AgentHandler(BaseHTTPRequestHandler): class AgentHandler(BaseHTTPRequestHandler):
def do_GET(self): def do_GET(self):
if self.path == '/health': parsed = urlparse(self.path)
path = parsed.path
if path == '/health':
self._json({'status': 'ok'}) self._json({'status': 'ok'})
elif self.path == '/': elif path == '/':
self._json({ self._json({
'name': 'modelhub-cancel-all', 'name': 'modelhub-submit-agent',
'gpu': TARGET_GPU, 'strategy_id': STRATEGY_ID,
'status': 'running' if state['running'] else 'idle', 'status': 'running' if state['running'] else 'idle',
'last_result': state.get('last_result'), 'last_run': state['last_run'],
'last_result': state['last_result'],
}) })
elif self.path.startswith('/logs'): elif path == '/test':
lines = 100 self._json(self._run_connectivity_test())
elif path == '/status':
self._json({
'running': state['running'],
'last_run': state['last_run'],
'last_result': state['last_result'],
'queue_count': db_conn.execute('SELECT COUNT(*) FROM queue').fetchone()[0] if db_conn else 0,
'submitted_count': db_conn.execute('SELECT COUNT(*) FROM submitted').fetchone()[0] if db_conn else 0,
})
elif path == '/logs':
lines = int(parse_qs(parsed.query).get('lines', ['50'])[0])
self._json({'logs': state['logs'][-lines:]}) self._json({'logs': state['logs'][-lines:]})
else: else:
self._json({'error': 'not found'}, 404) self._json({'error': 'not found'}, 404)
def do_POST(self):
parsed = urlparse(self.path)
path = parsed.path
if path == '/run':
if state['running']:
self._json({'error': 'already running'}, 409)
return
# 读取请求体
content_len = int(self.headers.get('Content-Length', 0))
body = {}
if content_len > 0:
body = json.loads(self.rfile.read(content_len))
limit = body.get('limit', 2)
self._json({'status': 'started', 'gpu': TARGET_GPU, 'limit': limit})
# 后台运行
def _run():
try:
state['running'] = True
state['last_run'] = datetime.now().isoformat()
count = run_pipeline(submit_limit=limit)
state['last_result'] = {'submitted': count, 'success': True}
except Exception as e:
log(f"流程异常: {traceback.format_exc()}")
state['last_result'] = {'error': str(e), 'success': False}
finally:
state['running'] = False
threading.Thread(target=_run, daemon=True).start()
else:
self._json({'error': 'not found'}, 404)
def _run_connectivity_test(self) -> dict:
"""测试各 API 连通性"""
results = {}
# 1. ModelScope 搜索 API
try:
resp = requests.get(
'https://modelscope.cn/openapi/v1/models',
params={'search': 'qwen', 'page_size': 2, 'sort': 'downloads'},
timeout=10,
headers={'User-Agent': 'Mozilla/5.0'}
)
results['modelscope_search'] = {
'status': resp.status_code,
'ok': resp.status_code == 200,
'body_preview': resp.text[:200] if resp.status_code == 200 else resp.text[:100],
}
except Exception as e:
results['modelscope_search'] = {'ok': False, 'error': str(e)}
# 2. ModelScope config.json API
try:
resp = requests.get(
'https://modelscope.cn/api/v1/models/Qwen/Qwen3-8B/repo?Revision=master&FilePath=config.json',
timeout=10,
)
results['modelscope_config'] = {
'status': resp.status_code,
'ok': resp.status_code == 200,
'body_preview': resp.text[:200] if resp.status_code == 200 else resp.text[:100],
}
except Exception as e:
results['modelscope_config'] = {'ok': False, 'error': str(e)}
# 3. ModelHub 查询 API
try:
resp = requests.get(
'https://modelhub.org.cn/api/adapt/task/page',
headers={'Xc-Token': TARGET_TOKEN, 'Accept': 'application/json'},
params={'current': 1, 'pageSize': 1, 'onlyMine': 'true'},
timeout=10,
)
data = resp.json()
results['modelhub_query'] = {
'ok': data.get('code') == 0,
'code': data.get('code'),
'total': data.get('data', {}).get('total'),
}
except Exception as e:
results['modelhub_query'] = {'ok': False, 'error': str(e)}
# 4. ModelHub 提交 API (dry test)
try:
resp = requests.post(
'https://modelhub.org.cn/api/adapt/task/add',
headers={'Xc-Token': TARGET_TOKEN, 'Accept': 'application/json', 'Content-Type': 'application/json'},
json={
'modelAddress': 'https://www.modelscope.cn/models/Qwen/Qwen3-8B',
'taskType': 'text-generation',
'targetGpu': TARGET_GPU,
'framework': 'vllm',
'strategyId': STRATEGY_ID,
'configParams': build_config_params(),
},
timeout=10,
)
data = resp.json()
results['modelhub_submit'] = {
'ok': data.get('code') == 0,
'code': data.get('code'),
'message': data.get('message', '')[:100],
}
except Exception as e:
results['modelhub_submit'] = {'ok': False, 'error': str(e)}
return results
def _json(self, body: dict, status: int = 200): def _json(self, body: dict, status: int = 200):
payload = json.dumps(body, ensure_ascii=False).encode() payload = json.dumps(body, ensure_ascii=False).encode()
self.send_response(status) self.send_response(status)
@@ -148,7 +600,7 @@ class AgentHandler(BaseHTTPRequestHandler):
self.wfile.write(payload) self.wfile.write(payload)
def log_message(self, fmt, *args): def log_message(self, fmt, *args):
pass pass # 静默 HTTP 日志
# ============================================================ # ============================================================
@@ -168,27 +620,105 @@ def main():
signal.signal(signal.SIGTERM, _handle_signal) signal.signal(signal.SIGTERM, _handle_signal)
signal.signal(signal.SIGINT, _handle_signal) signal.signal(signal.SIGINT, _handle_signal)
init_db()
server = ThreadingHTTPServer((HOST, PORT), AgentHandler) server = ThreadingHTTPServer((HOST, PORT), AgentHandler)
server.timeout = 1 server.timeout = 1
log(f"Cancel-All 智能体启动 | {HOST}:{PORT}") log(f"智能体启动 | {HOST}:{PORT}")
log(f"STRATEGY_ID: {STRATEGY_ID}")
log(f"目标GPU: {TARGET_GPU}") log(f"目标GPU: {TARGET_GPU}")
# 启动后自动取消所有任务 # 启动后自动运行连通性测试
def _auto_cancel(): def _startup_test():
time.sleep(2) time.sleep(2)
log("=" * 50)
log("连通性测试开始")
log("=" * 50)
# 1. ModelScope 搜索
try:
resp = requests.get(
'https://modelscope.cn/openapi/v1/models',
params={'search': 'qwen', 'page_size': 2, 'sort': 'downloads'},
timeout=10, headers={'User-Agent': 'Mozilla/5.0'}
)
if resp.status_code == 200:
data = resp.json()
models = data.get('data', {}).get('models', [])
log(f"[ModelScope搜索] ✅ status={resp.status_code} models={len(models)}")
for m in models[:2]:
log(f" {m.get('id')} downloads={m.get('downloads')}")
else:
log(f"[ModelScope搜索] ❌ status={resp.status_code} body={resp.text[:100]}")
except Exception as e:
log(f"[ModelScope搜索] ❌ error={e}")
# 2. ModelScope config.json
try:
resp = requests.get(
'https://modelscope.cn/api/v1/models/Qwen/Qwen3-8B/repo?Revision=master&FilePath=config.json',
timeout=10
)
if resp.status_code == 200:
cfg = resp.json()
log(f"[ModelScope配置] ✅ arch={cfg.get('architectures')} type={cfg.get('model_type')}")
else:
log(f"[ModelScope配置] ❌ status={resp.status_code}")
except Exception as e:
log(f"[ModelScope配置] ❌ error={e}")
# 3. ModelHub 查询
try:
resp = requests.get(
'https://modelhub.org.cn/api/adapt/task/page',
headers={'Xc-Token': TARGET_TOKEN, 'Accept': 'application/json'},
params={'current': 1, 'pageSize': 1, 'onlyMine': 'true'},
timeout=10
)
data = resp.json()
if data.get('code') == 0:
log(f"[ModelHub查询] ✅ total={data['data'].get('total')}")
else:
log(f"[ModelHub查询] ❌ code={data.get('code')} msg={data.get('message','')[:80]}")
except Exception as e:
log(f"[ModelHub查询] ❌ error={e}")
# 4. ModelHub 提交
try:
resp = requests.post(
'https://modelhub.org.cn/api/adapt/task/add',
headers={'Xc-Token': TARGET_TOKEN, 'Accept': 'application/json', 'Content-Type': 'application/json'},
json={
'modelAddress': 'https://www.modelscope.cn/models/Qwen/Qwen3-8B',
'taskType': 'text-generation', 'targetGpu': TARGET_GPU,
'framework': 'vllm', 'strategyId': STRATEGY_ID,
'configParams': build_config_params(),
}, timeout=10
)
data = resp.json()
log(f"[ModelHub提交] code={data.get('code')} msg={data.get('message','')[:80]}")
except Exception as e:
log(f"[ModelHub提交] ❌ error={e}")
log("=" * 50)
log("连通性测试完成")
log("=" * 50)
# 开始正式提交流程
log("\n自动触发提交流程...")
try: try:
state['running'] = True state['running'] = True
state['last_run'] = datetime.now().isoformat() state['last_run'] = datetime.now().isoformat()
cancel_all_tasks() count = run_pipeline(submit_limit=2)
state['last_result'] = {'success': True} state['last_result'] = {'submitted': count, 'success': True}
except Exception as e: except Exception as e:
log(f"异常: {e}") log(f"流程异常: {traceback.format_exc()}")
state['last_result'] = {'error': str(e)} state['last_result'] = {'error': str(e), 'success': False}
finally: finally:
state['running'] = False state['running'] = False
threading.Thread(target=_auto_cancel, daemon=True).start() threading.Thread(target=_startup_test, daemon=True).start()
while not shutdown_requested: while not shutdown_requested:
server.handle_request() server.handle_request()