feat: 算力分配回收API — pool统计/可用节点查询/集群节点/分配/回收

This commit is contained in:
yumoqing 2026-08-07 13:06:12 +08:00
parent db3f13b7ec
commit 56726bc679
5 changed files with 100 additions and 5 deletions

View File

@ -0,0 +1,26 @@
# 查询集群已分配节点
# params: cluster_id (必填)
cluster_id = params_kw.get('cluster_id', '')
if not cluster_id:
return {'status': 'error', 'message': 'Missing cluster_id'}
env = request._run_ns
dbname = env.get_module_dbname('pccs')
async with DBPools().sqlorContext(dbname) as sor:
rows = await sor.sqlExe(
"SELECT cn.id as node_id, cn.name, cn.ip_address, cn.cpu_cores, cn.memory_gb, "
"cn.gpu_count, cn.gpu_type, cnn.role, cnn.status as node_status, cnn.assigned_at "
"FROM cluster_node cnn "
"JOIN compute_node cn ON cn.id = cnn.node_id "
"WHERE cnn.cluster_id = ${cid}$ AND cnn.status != 'removed' "
"ORDER BY cnn.role, cn.name",
{'cid': cluster_id}
)
result = [{
'node_id': r.node_id, 'name': r.name, 'ip_address': r.ip_address,
'cpu_cores': r.cpu_cores, 'memory_gb': r.memory_gb,
'gpu_count': r.gpu_count, 'gpu_type': r.gpu_type,
'role': r.role, 'status': r.node_status,
'assigned_at': str(r.assigned_at)
} for r in rows]
return {'status': 'ok', 'data': result, 'total': len(result)}

View File

@ -1,4 +1,28 @@
# 从算力池分配节点到集群
# params: pool_id, cluster_id, role(control/compute), count, cpu, memory, gpu
# 从算力池分配节点到集群 (算力单元分配)
# params:
# pool_id - 算力池 ID (必填)
# cluster_id - 目标集群 ID
# role - 节点角色: control/compute (默认 compute)
# count - 分配数量 (默认 1)
# cpu - 每节点 CPU 核数要求 (可选0 表示不限制)
# memory - 每节点内存 GB 要求 (可选)
# gpu - 每节点 GPU 数量要求 (可选)
# return: {status, allocated_nodes: [node_id, ...]}
result = await allocate_nodes(request, params_kw)
return result if isinstance(result, dict) else {'status': 'error', 'message': str(result)}
if isinstance(result, dict) and result.get('status') == 'ok':
return {
'widgettype': 'Message',
'options': {
'title': '分配成功',
'message': '已分配 {} 个节点:\n{}'.format(
len(result.get('allocated_nodes', [])),
', '.join(result.get('allocated_nodes', []))
),
'type': 'success',
'timeout': 5
}
}
return result if isinstance(result, dict) else {
'widgettype': 'Error',
'options': {'title': '分配失败', 'message': str(result), 'cwidth': 16, 'cheight': 9, 'timeout': 3}
}

View File

@ -0,0 +1,24 @@
# 查询可用算力节点
# params: pool_id (必填), role (可选 control/compute)
pool_id = params_kw.get('pool_id', '')
role = params_kw.get('role', '')
if not pool_id:
return {'status': 'error', 'message': 'Missing pool_id'}
env = request._run_ns
dbname = env.get_module_dbname('pccs')
async with DBPools().sqlorContext(dbname) as sor:
sql = "SELECT * FROM compute_node WHERE pool_id = ${pid}$ AND status = 'available'"
params = {'pid': pool_id}
if role:
sql = sql + " AND node_type = ${role}$"
params['role'] = role
sql = sql + " ORDER BY cpu_cores DESC, memory_gb DESC"
rows = await sor.sqlExe(sql, params)
result = [{
'id': r.id, 'name': r.name, 'ip_address': r.ip_address,
'cpu_cores': r.cpu_cores, 'memory_gb': r.memory_gb,
'gpu_count': r.gpu_count, 'gpu_type': r.gpu_type,
'node_type': r.node_type
} for r in rows]
return {'status': 'ok', 'data': result, 'total': len(result)}

View File

@ -1,3 +1,20 @@
# 回收节点回算力池
# 回收节点回算力池 (算力单元回收)
# params:
# node_ids - 节点 ID 列表,逗号分隔 (必填)
# pool_id - 算力池 ID (可选,不传则自动从节点关联获取)
# return: {status, released: count}
result = await release_nodes(request, params_kw)
return result if isinstance(result, dict) else {'status': 'error', 'message': str(result)}
if isinstance(result, dict) and result.get('status') == 'ok':
return {
'widgettype': 'Message',
'options': {
'title': '回收成功',
'message': '已回收 {} 个节点到算力池'.format(result.get('released', 0)),
'type': 'success',
'timeout': 5
}
}
return result if isinstance(result, dict) else {
'widgettype': 'Error',
'options': {'title': '回收失败', 'message': str(result), 'cwidth': 16, 'cheight': 9, 'timeout': 3}
}

View File

@ -0,0 +1,4 @@
# 算力池资源统计
# params: pool_id (可选,不传则返回全部池)
result = await pool_stats(request, params_kw)
return result