diff --git a/pipeline_service/executor.py b/pipeline_service/executor.py index 59c7820..b845a31 100644 --- a/pipeline_service/executor.py +++ b/pipeline_service/executor.py @@ -146,7 +146,10 @@ async def _get_task_raw(task_id: str) -> dict: if not recs: return None rec = recs[0] - return {k: getattr(rec, k) for k in dir(rec) if not k.startswith('_')} if hasattr(rec, '__dict__') else dict(rec) + try: + return dict(rec) + except (TypeError, ValueError): + return {k: getattr(rec, k) for k in dir(rec) if not k.startswith('_')} async def _execute_step(task_id: str, step_name: str, step_graph: dict, task_info: dict): diff --git a/pipeline_service/storage.py b/pipeline_service/storage.py index 122ec86..48190c2 100644 --- a/pipeline_service/storage.py +++ b/pipeline_service/storage.py @@ -19,6 +19,16 @@ def _get_db(): return db, DBNAME +def _rec_to_dict(rec): + """sqlor 返回 DictObject,列名不暴露为实例属性(dir() 只列 dict 方法),须用 dict(rec) 取列。""" + if isinstance(rec, dict): + return dict(rec) + try: + return dict(rec) + except (TypeError, ValueError): + return {k: getattr(rec, k) for k in dir(rec) if not k.startswith('_')} + + async def get_pipeline_steps(pipeline_id: str) -> list: """Read step definitions from pipeline_steps table (defined by pipeline_core). @@ -32,10 +42,7 @@ async def get_pipeline_steps(pipeline_id: str) -> list: return [] result = [] for rec in recs: - if hasattr(rec, '__dict__'): - d = {k: getattr(rec, k) for k in dir(rec) if not k.startswith('_')} - else: - d = dict(rec) + d = _rec_to_dict(rec) # Extract deps from step_config JSON cfg_raw = d.get('step_config', '{}') try: @@ -99,9 +106,7 @@ async def get_task(tenant_id: str, task_id: str) -> Optional[dict]: if not recs: return None rec = recs[0] - if hasattr(rec, '__dict__'): - return {k: getattr(rec, k) for k in dir(rec) if not k.startswith('_')} - return dict(rec) + return _rec_to_dict(rec) async def get_task_steps(task_id: str) -> list: @@ -111,10 +116,7 @@ async def get_task_steps(task_id: str) -> list: 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__'): - result.append({k: getattr(rec, k) for k in dir(rec) if not k.startswith('_')}) - else: - result.append(dict(rec)) + result.append(_rec_to_dict(rec)) return result @@ -226,10 +228,7 @@ async def list_tasks(tenant_id: str, pipeline_id: str = None, limit: int = 100) recs = await sor.sqlExe(sql, {'tenant_id': tenant_id}) result = [] for rec in (recs or [])[:limit]: - if hasattr(rec, '__dict__'): - result.append({k: getattr(rec, k) for k in dir(rec) if not k.startswith('_')}) - else: - result.append(dict(rec)) + result.append(_rec_to_dict(rec)) return result @@ -318,10 +317,7 @@ async def list_human_tasks(tenant_id=None, assignee_role=None, assignee_id=None, recs = await sor.sqlExe(sql, params) result = [] for rec in (recs or []): - if hasattr(rec, '__dict__'): - d = {k: getattr(rec, k) for k in dir(rec) if not k.startswith('_')} - else: - d = dict(rec) + d = _rec_to_dict(rec) # Parse JSON fields for field in ('form_schema', 'result_data'): if d.get(field) and isinstance(d[field], str): @@ -341,10 +337,7 @@ async def get_human_task(task_id, step_name): if not recs: return None rec = recs[0] - if hasattr(rec, '__dict__'): - d = {k: getattr(rec, k) for k in dir(rec) if not k.startswith('_')} - else: - d = dict(rec) + d = _rec_to_dict(rec) for field in ('form_schema', 'result_data'): if d.get(field) and isinstance(d[field], str): try: