From ec75ee0aad23f4e7f8f0f18f4ed9d6b025a0cf6e Mon Sep 17 00:00:00 2001 From: ymq Date: Sun, 16 Aug 2026 23:45:08 +0800 Subject: [PATCH] =?UTF-8?q?fix(storage/executor):=20=E4=BF=AE=E5=A4=8D=20d?= =?UTF-8?q?ir(rec)=20=E5=8F=96=E5=88=97=20bug(DictObject=20=E5=88=97?= =?UTF-8?q?=E5=90=8D=E4=B8=8D=E6=9A=B4=E9=9C=B2=E4=B8=BA=E5=AE=9E=E4=BE=8B?= =?UTF-8?q?=E5=B1=9E=E6=80=A7)=EF=BC=8C=E7=BB=9F=E4=B8=80=E7=94=A8=20dict(?= =?UTF-8?q?rec)=20=E5=8F=96=E5=88=97?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit get_task/get_task_steps/list_tasks/list_human_tasks/get_human_task/get_pipeline_steps/_get_task_raw 全部改用 _rec_to_dict(dict(rec)),修复返回 dict 缺列名(如 state)导致 KeyError --- pipeline_service/executor.py | 5 ++++- pipeline_service/storage.py | 39 +++++++++++++++--------------------- 2 files changed, 20 insertions(+), 24 deletions(-) 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: