fix: dedupe workflow run event logs
This commit is contained in:
@@ -286,6 +286,8 @@ def _make_stream_event_log_entry(event: Dict[str, Any]) -> Optional[dict]:
|
|||||||
timestamp = event_data.get('timestamp') or datetime.now().isoformat()
|
timestamp = event_data.get('timestamp') or datetime.now().isoformat()
|
||||||
event_data.setdefault('timestamp', timestamp)
|
event_data.setdefault('timestamp', timestamp)
|
||||||
communication = dict(event_data.get('communication') or {})
|
communication = dict(event_data.get('communication') or {})
|
||||||
|
collaboration = dict(event_data.get('collaboration') or {})
|
||||||
|
actor = dict(communication.get('actor') or {})
|
||||||
communication.setdefault('event', event_type)
|
communication.setdefault('event', event_type)
|
||||||
communication.setdefault('timestamp', timestamp)
|
communication.setdefault('timestamp', timestamp)
|
||||||
|
|
||||||
@@ -294,6 +296,8 @@ def _make_stream_event_log_entry(event: Dict[str, Any]) -> Optional[dict]:
|
|||||||
node_label = (
|
node_label = (
|
||||||
event_data.get('node_label')
|
event_data.get('node_label')
|
||||||
or event_data.get('workflow_name')
|
or event_data.get('workflow_name')
|
||||||
|
or collaboration.get('node_label')
|
||||||
|
or actor.get('node_label')
|
||||||
or communication.get('message')
|
or communication.get('message')
|
||||||
or event_type
|
or event_type
|
||||||
)
|
)
|
||||||
@@ -317,7 +321,7 @@ def _make_stream_event_log_entry(event: Dict[str, Any]) -> Optional[dict]:
|
|||||||
'events': [event_data],
|
'events': [event_data],
|
||||||
'metadata': {
|
'metadata': {
|
||||||
'stream_event': True,
|
'stream_event': True,
|
||||||
'collaboration': event_data.get('collaboration') or {},
|
'collaboration': collaboration,
|
||||||
'communication': communication,
|
'communication': communication,
|
||||||
},
|
},
|
||||||
}
|
}
|
||||||
@@ -334,7 +338,19 @@ def _merge_execution_logs(node_logs: List[dict], event_logs: List[dict]) -> List
|
|||||||
log for log in (node_logs or [])
|
log for log in (node_logs or [])
|
||||||
if not (log.get('metadata') or {}).get('stream_event')
|
if not (log.get('metadata') or {}).get('stream_event')
|
||||||
]
|
]
|
||||||
merged = base_logs + list(event_logs or [])
|
completed_base_nodes = {
|
||||||
|
log.get('node_id')
|
||||||
|
for log in base_logs
|
||||||
|
if log.get('node_id') and log.get('status') in {'completed', 'failed', 'waiting'}
|
||||||
|
}
|
||||||
|
filtered_event_logs = [
|
||||||
|
log for log in (event_logs or [])
|
||||||
|
if not (
|
||||||
|
log.get('type') == 'node_complete'
|
||||||
|
and log.get('node_id') in completed_base_nodes
|
||||||
|
)
|
||||||
|
]
|
||||||
|
merged = base_logs + filtered_event_logs
|
||||||
return sorted(
|
return sorted(
|
||||||
merged,
|
merged,
|
||||||
key=lambda item: (_log_timestamp(item), 0 if item.get('metadata', {}).get('stream_event') else 1),
|
key=lambda item: (_log_timestamp(item), 0 if item.get('metadata', {}).get('stream_event') else 1),
|
||||||
|
|||||||
Reference in New Issue
Block a user