feat: improve workflow observability and lighten admin frontend

This commit is contained in:
2026-06-22 06:34:53 +08:00
parent 1b835a4aad
commit e1ef4c1a70
11 changed files with 360 additions and 416 deletions
@@ -226,6 +226,7 @@ def _make_execution_log_entry(
**extra: Any,
) -> dict:
extra.setdefault('timestamp', datetime.now().isoformat())
events = extra.pop('events', None)
metadata = dict(extra.pop('metadata', {}) or {})
event_type = _execution_log_event_type(extra.get('status'))
collaboration = _node_collaboration_metadata(node_map, node_id, node_type)
@@ -249,19 +250,106 @@ def _make_execution_log_entry(
error=extra.get('error'),
)
metadata['communication']['timestamp'] = extra.get('timestamp')
return {
entry = {
'node_id': node_id,
'node_type': node_type,
'node_label': _node_label_from_map(node_map, node_id),
'metadata': metadata,
**extra,
}
if events:
entry['events'] = events
return entry
def _event_log_status(event_type: str, event: Dict[str, Any]) -> str:
if event_type == 'error':
return 'failed'
if event_type == 'waiting_input':
return 'waiting'
if event_type in ('start', 'node_start', 'parallel_start'):
return 'running'
status = event.get('status') or 'completed'
if status == 'success':
return 'completed'
if status == 'failed':
return 'failed'
return status
def _make_stream_event_log_entry(event: Dict[str, Any]) -> Optional[dict]:
event_type = event.get('type')
if not event_type or event_type == 'llm_chunk':
return None
event_data = copy.deepcopy(event)
timestamp = event_data.get('timestamp') or datetime.now().isoformat()
event_data.setdefault('timestamp', timestamp)
communication = dict(event_data.get('communication') or {})
communication.setdefault('event', event_type)
communication.setdefault('timestamp', timestamp)
node_id = str(event_data.get('node_id') or event_data.get('run_id') or 'workflow')
node_type = str(event_data.get('node_type') or 'workflow')
node_label = (
event_data.get('node_label')
or event_data.get('workflow_name')
or communication.get('message')
or event_type
)
output = event_data.get('outputs')
if output is None:
output = event_data.get('content') or event_data.get('message')
error = event_data.get('error_message') or event_data.get('error')
return {
'type': event_type,
'node_id': node_id,
'node_type': node_type,
'node_label': str(node_label),
'status': _event_log_status(event_type, event_data),
'timestamp': timestamp,
'elapsed_time': event_data.get('elapsed_time'),
'tokens_used': event_data.get('tokens_used') or 0,
'output': output,
'error': error,
'event': event_data,
'events': [event_data],
'metadata': {
'stream_event': True,
'collaboration': event_data.get('collaboration') or {},
'communication': communication,
},
}
def _log_timestamp(log: Dict[str, Any]) -> str:
metadata = log.get('metadata') or {}
communication = metadata.get('communication') or {}
return log.get('timestamp') or communication.get('timestamp') or ''
def _merge_execution_logs(node_logs: List[dict], event_logs: List[dict]) -> List[dict]:
base_logs = [
log for log in (node_logs or [])
if not (log.get('metadata') or {}).get('stream_event')
]
merged = base_logs + list(event_logs or [])
return sorted(
merged,
key=lambda item: (_log_timestamp(item), 0 if item.get('metadata', {}).get('stream_event') else 1),
)
def _node_result_metadata(result: NodeResult) -> dict:
return dict(result.metadata or {})
def _node_result_events(result: NodeResult) -> List[Dict[str, Any]]:
events = getattr(result, 'events', None) or []
return copy.deepcopy(events) if events else []
def _event_collaboration_metadata(
node_map: Dict[str, Any],
node_id: str,
@@ -842,6 +930,7 @@ class AIWorkflowService:
elapsed_time=elapsed,
tokens_used=result.tokens_used,
metadata=_node_result_metadata(result),
events=_node_result_events(result),
inputs=copy.deepcopy(node_inputs),
target_node_id=target_node_id,
))
@@ -1024,6 +1113,7 @@ class AIWorkflowService:
elapsed_time=elapsed,
tokens_used=result.tokens_used,
metadata=_node_result_metadata(result),
events=_node_result_events(result),
inputs=copy.deepcopy(node_inputs),
branch=branch_id,
branch_label=branch_labels.get(branch_id, branch_id),
@@ -1252,6 +1342,7 @@ class AIWorkflowService:
elapsed_time=elapsed,
tokens_used=result.tokens_used,
metadata=_node_result_metadata(result),
events=_node_result_events(result),
branch=branch_id,
branch_label=branch_labels.get(branch_id, branch_id),
))
@@ -1545,6 +1636,7 @@ class AIWorkflowService:
elapsed_time=elapsed,
tokens_used=result.tokens_used,
metadata=_node_result_metadata(result),
events=_node_result_events(result),
loop_iteration=iteration,
))
@@ -1593,7 +1685,10 @@ class AIWorkflowService:
current_id,
node_type,
status='waiting',
output=result.output,
elapsed_time=elapsed,
metadata=_node_result_metadata(result),
events=_node_result_events(result),
loop_iteration=iteration,
))
@@ -1733,15 +1828,20 @@ class AIWorkflowService:
)
self._db.add(run)
await self._db.flush()
event_logs: List[dict] = []
try:
# 使用异步生成器执行工作流
async for event in self._execute_workflow_stream(workflow, inputs, run, conversation_history, use_draft):
yield event
event_log = _make_stream_event_log_entry(event)
if event_log:
event_logs.append(event_log)
run.execution_log = _merge_execution_logs(run.execution_log or [], event_logs)
# 在关键事件后提交数据库状态
event_type = event.get('type')
if event_type in ('node_complete', 'waiting_input', 'complete', 'error'):
if event_log or event_type in ('node_complete', 'waiting_input', 'complete', 'error'):
await self._db.commit()
# 更新运行记录
@@ -1771,13 +1871,19 @@ class AIWorkflowService:
await self._db.commit()
# 发送错误事件给前端
yield {
error_event = {
'type': 'error',
'message': str(e),
'error_message': str(e),
'run_id': run.id,
'workflow_id': workflow.id,
}
error_log = _make_stream_event_log_entry(error_event)
if error_log:
event_logs.append(error_log)
run.execution_log = _merge_execution_logs(run.execution_log or [], event_logs)
await self._db.commit()
yield error_event
async def _execute_workflow_stream(
self,
@@ -2026,6 +2132,7 @@ class AIWorkflowService:
output=result.output,
elapsed_time=elapsed,
metadata=_node_result_metadata(result),
events=_node_result_events(result),
inputs=copy.deepcopy(node_inputs),
target_node_id=target_node_id,
)
@@ -2067,6 +2174,7 @@ class AIWorkflowService:
elapsed_time=elapsed,
tokens_used=result.tokens_used,
metadata=_node_result_metadata(result),
events=_node_result_events(result),
inputs=copy.deepcopy(node_inputs),
target_node_id=target_node_id,
)
@@ -2557,6 +2665,7 @@ class AIWorkflowService:
output=result.output,
elapsed_time=elapsed,
metadata=_node_result_metadata(result),
events=_node_result_events(result),
target_node_id=target_node_id,
)
logs.append(log_entry)
@@ -2595,6 +2704,7 @@ class AIWorkflowService:
elapsed_time=elapsed,
tokens_used=result.tokens_used,
metadata=_node_result_metadata(result),
events=_node_result_events(result),
target_node_id=target_node_id,
)
logs.append(log_entry)