diff --git a/pipeline_service/storage.py b/pipeline_service/storage.py index d507da4..3cf6181 100644 --- a/pipeline_service/storage.py +++ b/pipeline_service/storage.py @@ -69,11 +69,11 @@ async def init_task_steps(task_id: str, step_records: list): data = { "id": step_id, "task_id": task_id, - "step_name": rec['step_name'], - "step_type": rec.get('step_type', rec['step_name']), - "display_name": rec.get('display_name', rec['step_name']), - "step_order": rec.get('step_order', 0), - "deps": rec.get('deps', '[]') if isinstance(rec.get('deps'), str) else json.dumps(rec.get('deps', [])), + "step_name": getattr(rec, 'step_name', ''), + "step_type": getattr(rec, 'step_type', '') or getattr(rec, 'step_name', ''), + "display_name": getattr(rec, 'display_name', '') or getattr(rec, 'step_name', ''), + "step_order": getattr(rec, 'step_order', 0), + "deps": getattr(rec, 'deps', '[]') if isinstance(getattr(rec, 'deps', '[]'), str) else json.dumps(getattr(rec, 'deps', [])), "state": "pending", } await sor.C('pipeline_task_steps', data) @@ -96,7 +96,7 @@ async def get_task_steps(task_id: str) -> list: """Get all step records for a task.""" db, dbname = _get_db() async with db.sqlorContext(dbname) as sor: - recs = await sor.R('pipeline_task_steps', {'task_id': task_id}, sort='step_order') + recs = await sor.sqlExe("SELECT * FROM pipeline_task_steps WHERE task_id=${tid}$ ORDER BY step_order", {'tid': task_id}) result = [] for rec in (recs or []): if hasattr(rec, '__dict__'):