fix(agent): preserve workflow node execution logs

This commit is contained in:
Stephen Zhou 2026-08-06 11:01:54 +08:00
parent a7283bccb3
commit 8d2eac55d7
No known key found for this signature in database
2 changed files with 405 additions and 68 deletions

View File

@ -224,8 +224,7 @@ class AgentObservabilityService:
)
if source_filter.kind in {"all", "workflow"}:
rows.extend(
self.serialize_log_message(message)
for message in self._list_workflow_messages(
self._list_workflow_messages(
app=app,
agent_id=agent_id,
conversation_id=conversation_id,
@ -308,22 +307,20 @@ class AgentObservabilityService:
workflow_app = aliased(App)
stmt = (
select(
Conversation,
WorkflowNodeExecutionModel.id.label("node_execution_id"),
WorkflowNodeExecutionModel.title.label("node_title"),
WorkflowNodeExecutionModel.status.label("node_status"),
WorkflowNodeExecutionModel.created_by_role.label("node_created_by_role"),
WorkflowNodeExecutionModel.created_by.label("node_created_by"),
WorkflowNodeExecutionModel.created_at.label("node_created_at"),
WorkflowNodeExecutionModel.finished_at.label("node_finished_at"),
workflow_app,
WorkflowAgentNodeBinding.workflow_id,
WorkflowAgentNodeBinding.workflow_version,
WorkflowAgentNodeBinding.node_id,
func.count(sa.distinct(Message.id)).label("message_count"),
func.max(Message.created_at).label("created_at"),
func.max(Message.updated_at).label("updated_at"),
func.sum(sa.case((Message.status == MessageStatus.PAUSED, 1), else_=0)).label("paused_count"),
func.sum(
sa.case((or_(Message.error.is_not(None), Message.status == MessageStatus.ERROR), 1), else_=0)
).label("failed_count"),
)
.select_from(Message)
.join(Conversation, Conversation.id == Message.conversation_id)
.join(WorkflowRun, WorkflowRun.id == Message.workflow_run_id)
.select_from(WorkflowNodeExecutionModel)
.join(WorkflowRun, WorkflowRun.id == WorkflowNodeExecutionModel.workflow_run_id)
.join(
WorkflowAgentNodeBinding,
and_(
@ -334,40 +331,32 @@ class AgentObservabilityService:
WorkflowAgentNodeBinding.workflow_version == WorkflowRun.version,
),
)
.join(
WorkflowNodeExecutionModel,
and_(
WorkflowNodeExecutionModel.workflow_run_id == WorkflowRun.id,
WorkflowNodeExecutionModel.node_id == WorkflowAgentNodeBinding.node_id,
),
)
.join(workflow_app, workflow_app.id == WorkflowAgentNodeBinding.app_id)
.where(Message.workflow_run_id.is_not(None), Conversation.app_id == WorkflowAgentNodeBinding.app_id)
.group_by(
Conversation.id,
workflow_app.id,
WorkflowAgentNodeBinding.workflow_id,
WorkflowAgentNodeBinding.workflow_version,
WorkflowAgentNodeBinding.node_id,
.where(
WorkflowNodeExecutionModel.tenant_id == app.tenant_id,
WorkflowNodeExecutionModel.app_id == WorkflowAgentNodeBinding.app_id,
WorkflowNodeExecutionModel.workflow_id == WorkflowAgentNodeBinding.workflow_id,
WorkflowNodeExecutionModel.node_id == WorkflowAgentNodeBinding.node_id,
)
)
stmt = self._apply_observability_filters(stmt, params=params, source_filter=source_filter)
stmt = self._apply_workflow_node_filters(stmt, params=params, workflow_app=workflow_app)
stmt = self._apply_workflow_source_filter(stmt, source_filter)
rows = list(self._session.execute(stmt).all())
return [
self._serialize_conversation_log(
conversation=row[0],
message_count=row.message_count,
paused_count=row.paused_count,
failed_count=row.failed_count,
self._serialize_workflow_execution_log(
node_execution_id=row.node_execution_id,
title=row.node_title,
status=row.node_status,
created_by_role=row.node_created_by_role,
created_by=row.node_created_by,
created_at=row.node_created_at,
finished_at=row.node_finished_at,
source=self._serialize_workflow_source(
app=row[1],
app=row[7],
workflow_id=row.workflow_id,
workflow_version=row.workflow_version,
node_id=row.node_id,
),
created_at=row.created_at,
updated_at=row.updated_at,
)
for row in rows
]
@ -443,10 +432,12 @@ class AgentObservabilityService:
conversation_id: str,
params: AgentLogQueryParams,
source_filter: AgentSourceFilter,
) -> list[Message]:
) -> list[dict[str, Any]]:
workflow_app = aliased(App)
stmt = (
select(Message)
.join(WorkflowRun, WorkflowRun.id == Message.workflow_run_id)
select(WorkflowNodeExecutionModel)
.select_from(WorkflowNodeExecutionModel)
.join(WorkflowRun, WorkflowRun.id == WorkflowNodeExecutionModel.workflow_run_id)
.join(
WorkflowAgentNodeBinding,
and_(
@ -457,18 +448,23 @@ class AgentObservabilityService:
WorkflowAgentNodeBinding.workflow_version == WorkflowRun.version,
),
)
.join(
WorkflowNodeExecutionModel,
and_(
WorkflowNodeExecutionModel.workflow_run_id == WorkflowRun.id,
WorkflowNodeExecutionModel.node_id == WorkflowAgentNodeBinding.node_id,
),
.join(workflow_app, workflow_app.id == WorkflowAgentNodeBinding.app_id)
.where(
WorkflowNodeExecutionModel.id == conversation_id,
WorkflowNodeExecutionModel.tenant_id == app.tenant_id,
WorkflowNodeExecutionModel.app_id == WorkflowAgentNodeBinding.app_id,
WorkflowNodeExecutionModel.workflow_id == WorkflowAgentNodeBinding.workflow_id,
WorkflowNodeExecutionModel.node_id == WorkflowAgentNodeBinding.node_id,
)
.where(Message.conversation_id == conversation_id)
)
stmt = self._apply_message_filters(stmt, params=params, source_filter=source_filter)
stmt = self._apply_workflow_node_filters(stmt, params=params, workflow_app=workflow_app)
stmt = self._apply_workflow_source_filter(stmt, source_filter)
return list(self._session.scalars(stmt.order_by(Message.created_at.desc(), Message.id.desc())).all())
executions = list(
self._session.scalars(
stmt.order_by(WorkflowNodeExecutionModel.created_at.desc(), WorkflowNodeExecutionModel.id.desc())
).all()
)
return [self.serialize_workflow_node_message(execution) for execution in executions]
def _list_workflow_sources(self, *, app: App, agent_id: str) -> list[dict[str, Any]]:
workflow_app = aliased(App)
@ -523,6 +519,62 @@ class AgentObservabilityService:
stmt = cls._apply_status_filter(stmt, params.statuses)
return stmt
@classmethod
def _apply_workflow_node_filters(cls, stmt, *, params: AgentLogQueryParams, workflow_app):
if params.start:
stmt = stmt.where(WorkflowNodeExecutionModel.created_at >= params.start)
if params.end:
stmt = stmt.where(WorkflowNodeExecutionModel.created_at < params.end)
if params.keyword:
escaped_keyword = escape_like_pattern(params.keyword)
pattern = f"%{escaped_keyword}%"
stmt = stmt.where(
or_(
WorkflowNodeExecutionModel.inputs.ilike(pattern, escape="\\"),
WorkflowNodeExecutionModel.outputs.ilike(pattern, escape="\\"),
WorkflowNodeExecutionModel.error.ilike(pattern, escape="\\"),
WorkflowNodeExecutionModel.title.ilike(pattern, escape="\\"),
workflow_app.name.ilike(pattern, escape="\\"),
)
)
if params.statuses:
stmt = cls._apply_workflow_node_status_filter(stmt, params.statuses)
return stmt
@staticmethod
def _apply_workflow_node_status_filter(stmt, statuses: tuple[str, ...]):
conditions = []
for status in statuses:
normalized = status.strip().lower()
if normalized in {"success", "normal"}:
conditions.append(WorkflowNodeExecutionModel.status == WorkflowNodeExecutionStatus.SUCCEEDED)
elif normalized in {"failed", "error"}:
conditions.append(
WorkflowNodeExecutionModel.status.in_(
(
WorkflowNodeExecutionStatus.FAILED,
WorkflowNodeExecutionStatus.EXCEPTION,
WorkflowNodeExecutionStatus.STOPPED,
)
)
)
elif normalized == "paused":
conditions.append(
WorkflowNodeExecutionModel.status.in_(
(
WorkflowNodeExecutionStatus.PAUSED,
WorkflowNodeExecutionStatus.PENDING,
WorkflowNodeExecutionStatus.RUNNING,
WorkflowNodeExecutionStatus.RETRY,
)
)
)
else:
raise ValueError(f"Unsupported status: {status}")
if not conditions:
return stmt
return stmt.where(or_(*conditions))
@staticmethod
def _apply_workflow_source_filter(stmt, source_filter: AgentSourceFilter):
if source_filter.app_id:
@ -569,8 +621,8 @@ class AgentObservabilityService:
source: dict[str, Any],
created_at: datetime | None,
updated_at: datetime | None,
user_rate: float | None = None,
operation_rate: float | None = None,
user_rate: float | None,
operation_rate: float | None,
) -> dict[str, Any]:
return {
"id": conversation.id,

View File

@ -6,8 +6,10 @@ from typing import Protocol
import pytest
from core.app.entities.app_invoke_entities import InvokeFrom
from graphon.enums import WorkflowNodeExecutionStatus
from models.enums import (
ConversationFromSource,
CreatorUserRole,
FeedbackFromSource,
FeedbackRating,
MessageStatus,
@ -276,23 +278,16 @@ def test_list_logs_sorts_by_requested_field(monkeypatch: pytest.MonkeyPatch) ->
def test_list_log_messages_merges_deduplicates_and_sorts_sources(monkeypatch: pytest.MonkeyPatch) -> None:
service = AgentObservabilityService(session=None)
webapp_message = SimpleNamespace(id="shared", created_at=10, updated_at=30)
workflow_messages = [
SimpleNamespace(id="shared", created_at=10, updated_at=20),
SimpleNamespace(id="workflow-only", created_at=20, updated_at=10),
webapp_message = SimpleNamespace(id="shared")
webapp_row = {"id": "shared", "created_at": 10, "updated_at": 30}
workflow_rows = [
{"id": "shared", "created_at": 10, "updated_at": 20},
{"id": "workflow-only", "created_at": 20, "updated_at": 10},
]
monkeypatch.setattr(service, "_list_webapp_messages", lambda **kwargs: [webapp_message])
monkeypatch.setattr(service, "_list_message_feedbacks", lambda **kwargs: {})
monkeypatch.setattr(
service,
"serialize_log_message",
lambda message, feedbacks=(): {
"id": message.id,
"created_at": message.created_at,
"updated_at": message.updated_at,
},
)
monkeypatch.setattr(service, "_list_workflow_messages", lambda **kwargs: workflow_messages)
monkeypatch.setattr(service, "serialize_log_message", lambda message, feedbacks=(): webapp_row)
monkeypatch.setattr(service, "_list_workflow_messages", lambda **kwargs: workflow_rows)
payload = service.list_log_messages(
app=SimpleNamespace(id="agent-app"), # type: ignore[arg-type]
@ -306,10 +301,7 @@ def test_list_log_messages_merges_deduplicates_and_sorts_sources(monkeypatch: py
)
assert payload == {
"data": [
{"id": "shared", "created_at": 10, "updated_at": 20},
{"id": "workflow-only", "created_at": 20, "updated_at": 10},
],
"data": [workflow_rows[0], workflow_rows[1]],
"page": 1,
"limit": 20,
"total": 2,
@ -371,6 +363,151 @@ def test_list_webapp_conversation_logs_includes_feedback_rates(monkeypatch: pyte
assert rows[0]["operation_rate"] == 1.0
def test_list_workflow_logs_uses_node_executions_without_messages() -> None:
created_at = datetime(2026, 7, 23, 7, 0, 19, tzinfo=UTC)
node_execution = SimpleNamespace(
id="node-execution-1",
title="Agent",
status=WorkflowNodeExecutionStatus.SUCCEEDED,
created_by_role=CreatorUserRole.END_USER,
created_by="end-user-1",
created_at=created_at,
finished_at=None,
)
workflow_app = SimpleNamespace(
id="workflow-app-1",
name="Marketing Department",
icon_type=None,
icon=None,
icon_background=None,
)
class FakeRow:
node_execution_id = node_execution.id
node_title = node_execution.title
node_status = node_execution.status
node_created_by_role = node_execution.created_by_role
node_created_by = node_execution.created_by
node_created_at = node_execution.created_at
node_finished_at = node_execution.finished_at
workflow_id = "workflow-1"
workflow_version = "v1"
node_id = "node-1"
def __getitem__(self, index: int):
return (None, None, None, None, None, None, None, workflow_app)[index]
class FakeResult:
def all(self):
return [FakeRow()]
class FakeSession:
def __init__(self):
self.query = ""
def execute(self, stmt):
self.query = str(stmt)
return FakeResult()
session = FakeSession()
service = AgentObservabilityService(session)
rows = service._list_workflow_conversation_logs(
app=SimpleNamespace(tenant_id="tenant-1"), # type: ignore[arg-type]
agent_id="agent-1",
params=AgentLogQueryParams(),
source_filter=AgentObservabilityService.resolve_source_filter("workflow:workflow-app-1"),
)
assert "FROM workflow_node_executions" in session.query
assert "JOIN messages" not in session.query
assert rows[0]["id"] == "node-execution-1"
assert rows[0]["source"]["app_name"] == "Marketing Department"
def test_list_workflow_messages_uses_node_execution_identity(monkeypatch: pytest.MonkeyPatch) -> None:
node_execution = SimpleNamespace(id="node-execution-1")
class FakeScalarResult:
def all(self):
return [node_execution]
class FakeSession:
def __init__(self):
self.query = ""
def scalars(self, stmt):
self.query = str(stmt)
return FakeScalarResult()
session = FakeSession()
service = AgentObservabilityService(session)
serialized = {"id": "node-execution-1", "conversation_id": "node-execution-1"}
monkeypatch.setattr(service, "serialize_workflow_node_message", lambda execution: serialized)
rows = service._list_workflow_messages(
app=SimpleNamespace(tenant_id="tenant-1"), # type: ignore[arg-type]
agent_id="agent-1",
conversation_id="node-execution-1",
params=AgentLogQueryParams(),
source_filter=AgentObservabilityService.resolve_source_filter("workflow:workflow-app-1"),
)
assert rows == [serialized]
assert "FROM workflow_node_executions" in session.query
assert "workflow_node_executions.id =" in session.query
assert "JOIN messages" not in session.query
def test_apply_workflow_node_filters_supports_time_keyword_and_status() -> None:
class FakeStmt:
def __init__(self):
self.conditions = []
def where(self, *conditions):
self.conditions.extend(conditions)
return self
stmt = FakeStmt()
params = AgentLogQueryParams(
start=datetime(2026, 7, 1, tzinfo=UTC),
end=datetime(2026, 8, 1, tzinfo=UTC),
keyword="meeting_100%",
statuses=("success",),
)
result = AgentObservabilityService._apply_workflow_node_filters(
stmt,
params=params,
workflow_app=SimpleNamespace(name=observability_service_module.App.name),
)
assert result is stmt
assert len(stmt.conditions) == 4
def test_apply_workflow_node_status_filter_supports_all_status_groups() -> None:
class FakeStmt:
def __init__(self):
self.conditions = []
def where(self, *conditions):
self.conditions.extend(conditions)
return self
stmt = FakeStmt()
result = AgentObservabilityService._apply_workflow_node_status_filter(stmt, ("normal", "error", "paused"))
assert result is stmt
assert len(stmt.conditions) == 1
empty_stmt = FakeStmt()
assert AgentObservabilityService._apply_workflow_node_status_filter(empty_stmt, ()) is empty_stmt
assert empty_stmt.conditions == []
with pytest.raises(ValueError, match="Unsupported status"):
AgentObservabilityService._apply_workflow_node_status_filter(FakeStmt(), ("unknown",))
def test_source_serializers_return_structured_frontend_shape() -> None:
app = SimpleNamespace(
id="app-1",
@ -524,6 +661,106 @@ def test_serialize_log_message_returns_frontend_log_shape() -> None:
}
def test_serialize_workflow_node_message_returns_frontend_log_shape() -> None:
created_at = datetime(2026, 7, 23, 7, 0, 19, tzinfo=UTC)
finished_at = datetime(2026, 7, 23, 7, 0, 28, tzinfo=UTC)
node_execution = SimpleNamespace(
id="node-execution-1",
title="Agent",
inputs=(
'{"agent_backend_request":{"composition":{"layers":['
'{"name":"workflow_node_job_prompt","config":{"user":"Summarize the meeting"}},'
'{"name":"workflow_user_prompt","config":{"user":"Focus on action items"}}]}}}'
),
outputs='{"output":"Alice owns the follow-up."}',
execution_metadata=(
'{"agent_log":{"agent_backend":{"usage":{"prompt_tokens":10,"completion_tokens":5,'
'"total_tokens":15,"total_price":"0.0015","currency":"USD","latency":1.25}}}}'
),
status=WorkflowNodeExecutionStatus.SUCCEEDED,
error=None,
elapsed_time=1.5,
created_by_role=CreatorUserRole.END_USER,
created_by="end-user-1",
created_at=created_at,
finished_at=finished_at,
)
payload = AgentObservabilityService.serialize_workflow_node_message(node_execution) # type: ignore[arg-type]
assert payload == {
"id": "node-execution-1",
"message_id": "node-execution-1",
"conversation_id": "node-execution-1",
"query": "Summarize the meeting\n\nFocus on action items",
"answer": "Alice owns the follow-up.",
"status": "success",
"error": None,
"from_end_user_id": "end-user-1",
"from_account_id": None,
"feedback_enabled": False,
"feedbacks": [],
"message_tokens": 10,
"answer_tokens": 5,
"total_tokens": 15,
"total_price": "0.0015",
"currency": "USD",
"latency": 1.25,
"created_at": int(created_at.timestamp()),
"updated_at": int(finished_at.timestamp()),
}
def test_serialize_workflow_node_message_handles_sparse_runtime_data() -> None:
created_at = datetime(2026, 7, 23, 7, 0, 19, tzinfo=UTC)
node_execution = SimpleNamespace(
id="node-execution-2",
title="Fallback prompt",
inputs={
"agent_backend_request": {
"composition": {
"layers": [
None,
{"name": "unrelated", "config": {"user": "ignored"}},
{"name": "workflow_user_prompt", "config": {"user": " "}},
]
}
}
},
outputs={"output": {"structured": True}},
execution_metadata={
"agent_log": {
"agent_backend": {
"usage": {
"prompt_tokens": "2",
"completion_tokens": 3,
}
}
}
},
status=WorkflowNodeExecutionStatus.PAUSED,
error=None,
elapsed_time=2,
created_by_role=CreatorUserRole.ACCOUNT.value,
created_by="account-1",
created_at=created_at,
finished_at=None,
)
payload = AgentObservabilityService.serialize_workflow_node_message(node_execution) # type: ignore[arg-type]
assert payload["query"] == "Fallback prompt"
assert payload["answer"] == '{"output": {"structured": true}}'
assert payload["status"] == "paused"
assert payload["from_end_user_id"] is None
assert payload["from_account_id"] == "account-1"
assert payload["total_tokens"] == 5
assert payload["total_price"] == "0"
assert payload["currency"] == ""
assert payload["latency"] == 2.0
assert payload["updated_at"] == int(created_at.timestamp())
def test_positive_feedback_rate_uses_rated_messages_as_denominator() -> None:
assert AgentObservabilityService._positive_feedback_rate(like_count=2, total_count=4) == 0.5
assert AgentObservabilityService._positive_feedback_rate(like_count=0, total_count=1) == 0
@ -615,6 +852,54 @@ def test_list_conversation_feedback_rates_maps_user_and_admin_sources() -> None:
)
def test_workflow_node_serialization_helpers_handle_invalid_values() -> None:
assert AgentObservabilityService._json_mapping(None) == {}
assert AgentObservabilityService._json_mapping("not-json") == {}
assert AgentObservabilityService._json_mapping("[]") == {}
assert AgentObservabilityService._mapping_value({"value": []}, "value") == {}
assert AgentObservabilityService._int_value(None) == 0
assert AgentObservabilityService._int_value("not-a-number") == 0
assert (
AgentObservabilityService._workflow_node_query(
{"agent_backend_request": {"composition": {"layers": "invalid"}}}, fallback="fallback"
)
== "fallback"
)
assert AgentObservabilityService._workflow_node_answer({"output": 1, "text": "fallback text"}) == "fallback text"
assert AgentObservabilityService._workflow_node_answer({}) == ""
def test_serialize_workflow_execution_log_uses_node_execution_identity() -> None:
created_at = datetime(2026, 7, 23, 7, 0, 19, tzinfo=UTC)
node_execution = SimpleNamespace(
id="node-execution-1",
title="Agent",
status=WorkflowNodeExecutionStatus.FAILED,
created_by_role=CreatorUserRole.ACCOUNT,
created_by="account-1",
created_at=created_at,
finished_at=None,
)
payload = AgentObservabilityService._serialize_workflow_execution_log(
node_execution_id=node_execution.id,
title=node_execution.title,
status=node_execution.status,
created_by_role=node_execution.created_by_role,
created_by=node_execution.created_by,
created_at=node_execution.created_at,
finished_at=node_execution.finished_at,
source={"id": "workflow:app-1:workflow-1:v1:node-1"},
)
assert payload["id"] == "node-execution-1"
assert payload["conversation_id"] == "node-execution-1"
assert payload["message_count"] == 1
assert payload["end_user_id"] is None
assert payload["status"] == "failed"
assert payload["unread"] is False
def test_build_charts_and_summary_match_monitoring_metrics() -> None:
rows = [
{