fix(storage/executor): 修复 dir(rec) 取列 bug(DictObject 列名不暴露为实例属性),统一用 dict(rec) 取列

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
This commit is contained in:
ymq 2026-08-16 23:45:08 +08:00
parent a3a15ae666
commit ec75ee0aad
2 changed files with 20 additions and 24 deletions

View File

@ -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):

View File

@ -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: