yumoqing d9a1903d7a feat: pipeline-app独立产线后端服务
3个业务模块:
- pipeline_core: 产线定义(pipelines/steps/versions)
- pipeline_ops: 运营(定价/供应量/使用记录)
- pipeline_dist: 分销(分销商/独立定价/API密钥)

- ahserver独立部署(端口9090)
- 独立数据库pipeline
- 80个文件, 符合module/db-table/crud三规范
2026-08-03 16:09:55 +08:00

280 lines
9.9 KiB
Python

"""pipeline_core 模块 - 产线管理核心模块"""
import json
from appPublic.uniqueID import getID
from sqlor.dbpools import DBPools
from ahserver.serverenv import ServerEnv
from appPublic.log import debug
MODULE_NAME = 'pipeline_core'
DBNAME = 'pipeline'
def _get_sor():
"""Get database pool and dbname for pipeline_core module."""
return DBPools(), DBNAME
async def create_pipeline(params_kw):
"""创建产线记录"""
result = {'success': False, 'message': ''}
try:
db, dbname = _get_sor()
async with db.sqlorContext(dbname) as sor:
data = params_kw.copy()
data.pop('page', None)
data.pop('rows', None)
data.pop('data_filter', None)
data['id'] = getID()
await sor.C('pipelines', data)
result['success'] = True
result['message'] = '创建成功'
result['id'] = data['id']
except Exception as e:
result['message'] = str(e)
return json.dumps(result, ensure_ascii=False, default=str)
async def update_pipeline(params_kw):
"""更新产线记录"""
result = {'success': False, 'message': ''}
try:
db, dbname = _get_sor()
async with db.sqlorContext(dbname) as sor:
data = params_kw.copy()
data.pop('page', None)
data.pop('rows', None)
data.pop('data_filter', None)
record_id = data.pop('id', None)
if not record_id:
result['message'] = '缺少id'
else:
await sor.U('pipelines', data, {'id': record_id})
result['success'] = True
result['message'] = '更新成功'
except Exception as e:
result['message'] = str(e)
return json.dumps(result, ensure_ascii=False, default=str)
async def delete_pipeline(params_kw):
"""删除产线记录"""
result = {'success': False, 'message': ''}
try:
db, dbname = _get_sor()
async with db.sqlorContext(dbname) as sor:
record_id = params_kw.get('id')
if not record_id:
result['message'] = '缺少id'
else:
await sor.D('pipelines', {'id': record_id})
result['success'] = True
result['message'] = '删除成功'
except Exception as e:
result['message'] = str(e)
return json.dumps(result, ensure_ascii=False, default=str)
async def create_pipeline_step(params_kw):
"""创建产线步骤"""
result = {'success': False, 'message': ''}
try:
db, dbname = _get_sor()
async with db.sqlorContext(dbname) as sor:
data = params_kw.copy()
data.pop('page', None)
data.pop('rows', None)
data.pop('data_filter', None)
data['id'] = getID()
await sor.C('pipeline_steps', data)
result['success'] = True
result['message'] = '创建成功'
result['id'] = data['id']
except Exception as e:
result['message'] = str(e)
return json.dumps(result, ensure_ascii=False, default=str)
async def update_pipeline_step(params_kw):
"""更新产线步骤"""
result = {'success': False, 'message': ''}
try:
db, dbname = _get_sor()
async with db.sqlorContext(dbname) as sor:
data = params_kw.copy()
data.pop('page', None)
data.pop('rows', None)
data.pop('data_filter', None)
record_id = data.pop('id', None)
if not record_id:
result['message'] = '缺少id'
else:
await sor.U('pipeline_steps', data, {'id': record_id})
result['success'] = True
result['message'] = '更新成功'
except Exception as e:
result['message'] = str(e)
return json.dumps(result, ensure_ascii=False, default=str)
async def delete_pipeline_step(params_kw):
"""删除产线步骤"""
result = {'success': False, 'message': ''}
try:
db, dbname = _get_sor()
async with db.sqlorContext(dbname) as sor:
record_id = params_kw.get('id')
if not record_id:
result['message'] = '缺少id'
else:
await sor.D('pipeline_steps', {'id': record_id})
result['success'] = True
result['message'] = '删除成功'
except Exception as e:
result['message'] = str(e)
return json.dumps(result, ensure_ascii=False, default=str)
async def create_pipeline_version(params_kw):
"""创建产线发布记录"""
result = {'success': False, 'message': ''}
try:
db, dbname = _get_sor()
async with db.sqlorContext(dbname) as sor:
data = params_kw.copy()
data.pop('page', None)
data.pop('rows', None)
data.pop('data_filter', None)
data['id'] = getID()
await sor.C('pipeline_versions', data)
result['success'] = True
result['message'] = '创建成功'
result['id'] = data['id']
except Exception as e:
result['message'] = str(e)
return json.dumps(result, ensure_ascii=False, default=str)
async def update_pipeline_version(params_kw):
"""更新产线发布记录"""
result = {'success': False, 'message': ''}
try:
db, dbname = _get_sor()
async with db.sqlorContext(dbname) as sor:
data = params_kw.copy()
data.pop('page', None)
data.pop('rows', None)
data.pop('data_filter', None)
record_id = data.pop('id', None)
if not record_id:
result['message'] = '缺少id'
else:
await sor.U('pipeline_versions', data, {'id': record_id})
result['success'] = True
result['message'] = '更新成功'
except Exception as e:
result['message'] = str(e)
return json.dumps(result, ensure_ascii=False, default=str)
async def delete_pipeline_version(params_kw):
"""删除产线发布记录"""
result = {'success': False, 'message': ''}
try:
db, dbname = _get_sor()
async with db.sqlorContext(dbname) as sor:
record_id = params_kw.get('id')
if not record_id:
result['message'] = '缺少id'
else:
await sor.D('pipeline_versions', {'id': record_id})
result['success'] = True
result['message'] = '删除成功'
except Exception as e:
result['message'] = str(e)
return json.dumps(result, ensure_ascii=False, default=str)
async def publish_pipeline(params_kw):
"""发布产线 - 创建发布记录并更新产线状态"""
result = {'success': False, 'message': ''}
try:
pipeline_id = params_kw.get('pipeline_id')
if not pipeline_id:
result['message'] = '缺少pipeline_id'
return json.dumps(result, ensure_ascii=False, default=str)
changelog = params_kw.get('changelog', '')
published_by = params_kw.get('published_by', '')
db, dbname = _get_sor()
async with db.sqlorContext(dbname) as sor:
# Get current pipeline
recs = await sor.R('pipelines', {'id': pipeline_id})
if not recs:
result['message'] = f'产线不存在: {pipeline_id}'
return json.dumps(result, ensure_ascii=False, default=str)
pipeline = recs[0]
current_version = pipeline.get('version', '1.0.0') if hasattr(pipeline, 'get') else getattr(pipeline, 'version', '1.0.0')
# Build config snapshot from current pipeline data
config_snapshot = json.dumps({
'pipeline_config': getattr(pipeline, 'pipeline_config', None),
'model_api_url': getattr(pipeline, 'model_api_url', None),
'version': current_version,
}, ensure_ascii=False, default=str)
# Create version record
version_id = getID()
version_data = {
'id': version_id,
'pipeline_id': pipeline_id,
'version': current_version,
'publish_status': 'published',
'published_by': published_by,
'changelog': changelog,
'config_snapshot': config_snapshot,
}
await sor.C('pipeline_versions', version_data)
# Update pipeline status to published
await sor.U('pipelines', {'status': 'published'}, {'id': pipeline_id})
result['success'] = True
result['message'] = '发布成功'
result['version_id'] = version_id
except Exception as e:
result['message'] = str(e)
return json.dumps(result, ensure_ascii=False, default=str)
def load_pipeline_core():
"""注册函数到 ServerEnv"""
env = ServerEnv()
# Pipelines
env.create_pipeline = create_pipeline
env.create_pipelines = create_pipeline
env.update_pipeline = update_pipeline
env.update_pipelines = update_pipeline
env.delete_pipeline = delete_pipeline
env.delete_pipelines = delete_pipeline
# Pipeline Steps
env.create_pipeline_step = create_pipeline_step
env.create_pipeline_steps = create_pipeline_step
env.update_pipeline_step = update_pipeline_step
env.update_pipeline_steps = update_pipeline_step
env.delete_pipeline_step = delete_pipeline_step
env.delete_pipeline_steps = delete_pipeline_step
# Pipeline Versions
env.create_pipeline_version = create_pipeline_version
env.create_pipeline_versions = create_pipeline_version
env.update_pipeline_version = update_pipeline_version
env.update_pipeline_versions = update_pipeline_version
env.delete_pipeline_version = delete_pipeline_version
env.delete_pipeline_versions = delete_pipeline_version
# Publish
env.publish_pipeline = publish_pipeline
debug(f'[{MODULE_NAME}] module loaded')
return True