mirror of
https://github.com/langgenius/dify.git
synced 2026-09-29 17:07:38 +08:00
fix(agent): attribute nested tool activity to its source
Resolve Agent logs, details and usage statistics from the executed source workflow and binding while preserving the outer workflow run relationship.
This commit is contained in:
@@ -14,12 +14,19 @@ from sqlalchemy.orm import aliased
|
||||
|
||||
from configs import dify_config
|
||||
from core.app.entities.app_invoke_entities import InvokeFrom
|
||||
from core.workflow.node_execution_process_data import WORKFLOW_AGENT_BINDING_ID_KEY
|
||||
from graphon.enums import WorkflowNodeExecutionStatus
|
||||
from libs.helper import convert_datetime_to_date, escape_like_pattern, to_timestamp
|
||||
from models.agent import WorkflowAgentNodeBinding
|
||||
from models.enums import CreatorUserRole, FeedbackFromSource, FeedbackRating, MessageStatus
|
||||
from models.model import App, Conversation, Message, MessageFeedback
|
||||
from models.workflow import WorkflowNodeExecutionModel, WorkflowRun, WorkflowType
|
||||
from models.workflow import (
|
||||
Workflow,
|
||||
WorkflowNodeExecutionModel,
|
||||
WorkflowNodeExecutionTriggeredFrom,
|
||||
WorkflowRun,
|
||||
WorkflowType,
|
||||
)
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
@@ -326,17 +333,15 @@ class AgentObservabilityService:
|
||||
and_(
|
||||
WorkflowAgentNodeBinding.tenant_id == app.tenant_id,
|
||||
WorkflowAgentNodeBinding.agent_id == agent_id,
|
||||
WorkflowAgentNodeBinding.app_id == WorkflowRun.app_id,
|
||||
WorkflowAgentNodeBinding.workflow_id == WorkflowRun.workflow_id,
|
||||
WorkflowAgentNodeBinding.workflow_version == WorkflowRun.version,
|
||||
self._workflow_node_binding_condition(),
|
||||
),
|
||||
)
|
||||
.join(workflow_app, workflow_app.id == WorkflowAgentNodeBinding.app_id)
|
||||
.join(
|
||||
workflow_app,
|
||||
and_(workflow_app.id == WorkflowAgentNodeBinding.app_id, workflow_app.tenant_id == app.tenant_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_workflow_node_filters(stmt, params=params, workflow_app=workflow_app)
|
||||
@@ -443,18 +448,16 @@ class AgentObservabilityService:
|
||||
and_(
|
||||
WorkflowAgentNodeBinding.tenant_id == app.tenant_id,
|
||||
WorkflowAgentNodeBinding.agent_id == agent_id,
|
||||
WorkflowAgentNodeBinding.app_id == WorkflowRun.app_id,
|
||||
WorkflowAgentNodeBinding.workflow_id == WorkflowRun.workflow_id,
|
||||
WorkflowAgentNodeBinding.workflow_version == WorkflowRun.version,
|
||||
self._workflow_node_binding_condition(),
|
||||
),
|
||||
)
|
||||
.join(workflow_app, workflow_app.id == WorkflowAgentNodeBinding.app_id)
|
||||
.join(
|
||||
workflow_app,
|
||||
and_(workflow_app.id == WorkflowAgentNodeBinding.app_id, workflow_app.tenant_id == app.tenant_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,
|
||||
)
|
||||
)
|
||||
stmt = self._apply_workflow_node_filters(stmt, params=params, workflow_app=workflow_app)
|
||||
@@ -466,6 +469,65 @@ class AgentObservabilityService:
|
||||
)
|
||||
return [self.serialize_workflow_node_message(execution) for execution in executions]
|
||||
|
||||
def _workflow_node_binding_condition(self, *, statistics: bool = False) -> sa.ColumnElement[bool]:
|
||||
"""Resolve source ownership independently of a Tool's outer lifecycle run."""
|
||||
node = aliased(WorkflowNodeExecutionModel, name="wne") if statistics else WorkflowNodeExecutionModel
|
||||
binding = aliased(WorkflowAgentNodeBinding, name="wanb") if statistics else WorkflowAgentNodeBinding
|
||||
run = aliased(WorkflowRun, name="wr") if statistics else WorkflowRun
|
||||
process_data_json = func.nullif(node.process_data, "")
|
||||
if self._session.get_bind().dialect.name == "sqlite":
|
||||
binding_id = func.json_extract(process_data_json, f"$.{WORKFLOW_AGENT_BINDING_ID_KEY}")
|
||||
else:
|
||||
binding_id = sa.cast(process_data_json, sa.JSON)[WORKFLOW_AGENT_BINDING_ID_KEY].as_string()
|
||||
|
||||
node_table = "wne" if statistics else WorkflowNodeExecutionModel.__tablename__
|
||||
source_version = (
|
||||
select(Workflow.version)
|
||||
.where(
|
||||
Workflow.id == sa.literal_column(f"{node_table}.workflow_id"),
|
||||
Workflow.app_id == sa.literal_column(f"{node_table}.app_id"),
|
||||
Workflow.tenant_id == sa.literal_column(f"{node_table}.tenant_id"),
|
||||
)
|
||||
.scalar_subquery()
|
||||
)
|
||||
old_binding = aliased(WorkflowAgentNodeBinding, name="old_binding")
|
||||
single_binding_version = (
|
||||
select(func.min(old_binding.workflow_version))
|
||||
.where(
|
||||
old_binding.tenant_id == sa.literal_column(f"{node_table}.tenant_id"),
|
||||
old_binding.app_id == sa.literal_column(f"{node_table}.app_id"),
|
||||
old_binding.workflow_id == sa.literal_column(f"{node_table}.workflow_id"),
|
||||
old_binding.node_id == sa.literal_column(f"{node_table}.node_id"),
|
||||
)
|
||||
.having(func.count() == 1)
|
||||
.scalar_subquery()
|
||||
)
|
||||
is_tool = node.triggered_from == WorkflowNodeExecutionTriggeredFrom.WORKFLOW_TOOL
|
||||
# Old Tool rows may outlive their source Workflow. A sole source binding
|
||||
# is still unambiguous; never choose among versions using the caller run.
|
||||
old_version = sa.case((is_tool, func.coalesce(source_version, single_binding_version)), else_=run.version)
|
||||
return and_(
|
||||
node.tenant_id == run.tenant_id,
|
||||
binding.tenant_id == node.tenant_id,
|
||||
binding.app_id == node.app_id,
|
||||
binding.workflow_id == node.workflow_id,
|
||||
binding.node_id == node.node_id,
|
||||
or_(is_tool, and_(node.app_id == run.app_id, node.workflow_id == run.workflow_id)),
|
||||
or_(
|
||||
sa.cast(binding.id, sa.String) == binding_id,
|
||||
and_(binding_id.is_(None), binding.workflow_version == old_version),
|
||||
),
|
||||
)
|
||||
|
||||
def _statistics_workflow_node_binding_join_sql(self) -> str:
|
||||
# Reuse the log/detail owner predicate in both raw-SQL statistics queries.
|
||||
return str(
|
||||
self._workflow_node_binding_condition(statistics=True).compile(
|
||||
dialect=self._session.get_bind().dialect,
|
||||
compile_kwargs={"literal_binds": True},
|
||||
)
|
||||
)
|
||||
|
||||
def _list_workflow_sources(self, *, app: App, agent_id: str) -> list[dict[str, Any]]:
|
||||
workflow_app = aliased(App)
|
||||
stmt = (
|
||||
@@ -947,6 +1009,7 @@ WHERE
|
||||
("agent_log", "agent_backend", "usage", "completion_tokens"), "BIGINT"
|
||||
)
|
||||
binding_filters = self._statistics_workflow_binding_filters_sql(source_filter)
|
||||
binding_join = self._statistics_workflow_node_binding_join_sql()
|
||||
run_date_filters = ""
|
||||
args: dict[str, Any] = {
|
||||
"tz": params.timezone,
|
||||
@@ -981,17 +1044,14 @@ WHERE
|
||||
COALESCE(SUM(COALESCE(wne.elapsed_time, 0)), 0) AS latency,
|
||||
COALESCE(SUM(COALESCE({completion_tokens}, 0)), 0) AS answer_tokens
|
||||
FROM workflow_runs wr
|
||||
JOIN workflow_agent_node_bindings wanb
|
||||
ON wanb.tenant_id = :tenant_id
|
||||
AND wanb.agent_id = :agent_id
|
||||
AND wanb.app_id = wr.app_id
|
||||
AND wanb.workflow_id = wr.workflow_id
|
||||
AND wanb.workflow_version = wr.version
|
||||
{binding_filters}
|
||||
JOIN workflow_node_executions wne
|
||||
ON wne.workflow_run_id = wr.id
|
||||
AND wne.node_id = wanb.node_id
|
||||
WHERE wr.type != :chat_workflow_type{run_date_filters}
|
||||
JOIN workflow_agent_node_bindings wanb
|
||||
ON {binding_join}
|
||||
AND wanb.agent_id = :agent_id
|
||||
{binding_filters}
|
||||
JOIN apps source_app ON source_app.id = wanb.app_id AND source_app.tenant_id = wr.tenant_id
|
||||
WHERE wr.tenant_id = :tenant_id AND wr.type != :chat_workflow_type{run_date_filters}
|
||||
GROUP BY wr.id, wr.created_by_role, wr.created_by, wr.created_at
|
||||
)
|
||||
SELECT
|
||||
@@ -1085,24 +1145,23 @@ WHERE
|
||||
workflow_binding_filters.append("wanb.node_id = :node_id")
|
||||
return f"AND {' AND '.join(workflow_binding_filters)}" if workflow_binding_filters else ""
|
||||
|
||||
@classmethod
|
||||
def _statistics_workflow_message_scope_sql(cls, source_filter: AgentSourceFilter) -> str:
|
||||
binding_filters = cls._statistics_workflow_binding_filters_sql(source_filter)
|
||||
def _statistics_workflow_message_scope_sql(self, source_filter: AgentSourceFilter) -> str:
|
||||
binding_filters = self._statistics_workflow_binding_filters_sql(source_filter)
|
||||
binding_join = self._statistics_workflow_node_binding_join_sql()
|
||||
return f"""m.workflow_run_id IS NOT NULL
|
||||
AND EXISTS (
|
||||
SELECT 1
|
||||
FROM workflow_runs wr
|
||||
JOIN workflow_agent_node_bindings wanb
|
||||
ON wanb.tenant_id = :tenant_id
|
||||
AND wanb.agent_id = :agent_id
|
||||
AND wanb.app_id = wr.app_id
|
||||
AND wanb.workflow_id = wr.workflow_id
|
||||
AND wanb.workflow_version = wr.version
|
||||
{binding_filters}
|
||||
JOIN workflow_node_executions wne
|
||||
ON wne.workflow_run_id = wr.id
|
||||
AND wne.node_id = wanb.node_id
|
||||
JOIN workflow_agent_node_bindings wanb
|
||||
ON {binding_join}
|
||||
AND wanb.agent_id = :agent_id
|
||||
{binding_filters}
|
||||
JOIN apps source_app ON source_app.id = wanb.app_id AND source_app.tenant_id = wr.tenant_id
|
||||
WHERE wr.id = m.workflow_run_id
|
||||
AND wr.app_id = m.app_id
|
||||
AND wr.tenant_id = :tenant_id
|
||||
AND wr.type = :chat_workflow_type
|
||||
)"""
|
||||
|
||||
|
||||
@@ -20,6 +20,7 @@ from models.enums import (
|
||||
)
|
||||
from models.model import App, AppMode, Conversation, Message, MessageFeedback
|
||||
from models.workflow import (
|
||||
Workflow,
|
||||
WorkflowExecutionStatus,
|
||||
WorkflowNodeExecutionModel,
|
||||
WorkflowNodeExecutionTriggeredFrom,
|
||||
@@ -272,15 +273,6 @@ def test_statistics_workflow_app_source_covers_all_versions_and_nodes() -> None:
|
||||
assert "wanb.node_id = :node_id" not in scope_sql
|
||||
|
||||
|
||||
def test_statistics_workflow_chat_context_only_uses_chat_runs() -> None:
|
||||
source_filter = AgentObservabilityService.resolve_source_filter("workflow:app-2")
|
||||
|
||||
scope_sql = AgentObservabilityService._statistics_workflow_message_scope_sql(source_filter)
|
||||
|
||||
assert "wr.id = m.workflow_run_id" in scope_sql
|
||||
assert "wr.type = :chat_workflow_type" in scope_sql
|
||||
|
||||
|
||||
def test_workflow_metadata_numeric_sql_supports_postgresql_and_mysql(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
apply_config_overrides(monkeypatch, DB_TYPE="postgresql")
|
||||
|
||||
@@ -489,6 +481,97 @@ def test_list_workflow_logs_uses_node_executions_without_messages(sqlite_session
|
||||
assert rows[0]["source"]["app_name"] == "Marketing Department"
|
||||
|
||||
|
||||
@pytest.mark.parametrize("source_identity", ["persisted-binding", "source-version", "unique-legacy-binding"])
|
||||
def test_nested_agent_logs_and_messages_keep_source_identity(sqlite_session: Session, source_identity: str) -> None:
|
||||
source_app = _app(app_id="workflow-app-1", name="Source workflow", mode=AppMode.WORKFLOW)
|
||||
source = Workflow(
|
||||
id="workflow-1",
|
||||
tenant_id="tenant-1",
|
||||
app_id=source_app.id,
|
||||
type=WorkflowType.WORKFLOW,
|
||||
version="v1",
|
||||
graph="{}",
|
||||
_features="{}",
|
||||
created_by="account-1",
|
||||
)
|
||||
run = _workflow_run()
|
||||
run.app_id, run.workflow_id, run.version = "caller-app", "caller-workflow", "caller-version"
|
||||
execution = _node_execution()
|
||||
execution.triggered_from = WorkflowNodeExecutionTriggeredFrom.WORKFLOW_TOOL
|
||||
binding = _workflow_binding()
|
||||
if source_identity == "persisted-binding":
|
||||
execution.process_data = json.dumps({"workflow_agent_binding_id": binding.id})
|
||||
unused_binding = _workflow_binding(binding_id="unused-binding")
|
||||
unused_binding.workflow_version = "v2"
|
||||
sqlite_session.add_all([source_app, run, execution, binding])
|
||||
if source_identity != "unique-legacy-binding":
|
||||
sqlite_session.add_all([source, unused_binding])
|
||||
sqlite_session.commit()
|
||||
service = AgentObservabilityService(sqlite_session)
|
||||
params = AgentLogQueryParams(sources=("workflow:workflow-app-1",))
|
||||
|
||||
logs = service.list_logs(app=_app(app_id="agent-app"), agent_id="agent-1", params=params)
|
||||
messages = service.list_log_messages(
|
||||
app=_app(app_id="agent-app"), agent_id="agent-1", conversation_id=execution.id, params=params
|
||||
)
|
||||
|
||||
assert logs["total"] == messages["total"] == 1
|
||||
assert logs["data"][0]["id"] == messages["data"][0]["id"] == execution.id
|
||||
assert logs["data"][0]["source"]["id"] == "workflow:workflow-app-1:workflow-1:v1:node-1"
|
||||
assert messages["data"][0]["total_tokens"] == 454_064
|
||||
for source_filter in ("workflow:caller-app", "workflow:workflow-app-1:workflow-1:v2:node-1"):
|
||||
params = AgentLogQueryParams(sources=(source_filter,))
|
||||
assert service.list_logs(app=_app(app_id="agent-app"), agent_id="agent-1", params=params)["total"] == 0
|
||||
assert (
|
||||
service.list_log_messages(
|
||||
app=_app(app_id="agent-app"), agent_id="agent-1", conversation_id=execution.id, params=params
|
||||
)["total"]
|
||||
== 0
|
||||
)
|
||||
|
||||
if source_identity == "unique-legacy-binding":
|
||||
# Without a captured binding or source Workflow, another Agent/version
|
||||
# makes attribution ambiguous, even when filtering for only agent-1.
|
||||
unused_binding.agent_id = "another-agent"
|
||||
sqlite_session.add(unused_binding)
|
||||
sqlite_session.commit()
|
||||
params = AgentLogQueryParams(sources=("workflow",))
|
||||
assert service.list_logs(app=_app(app_id="agent-app"), agent_id="agent-1", params=params)["total"] == 0
|
||||
assert (
|
||||
service.list_log_messages(
|
||||
app=_app(app_id="agent-app"), agent_id="agent-1", conversation_id=execution.id, params=params
|
||||
)["total"]
|
||||
== 0
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.parametrize("mismatch", ["run-tenant", "source-tenant", "binding-id", "unrelated-run"])
|
||||
def test_workflow_agent_observability_rejects_mismatched_ownership(sqlite_session: Session, mismatch: str) -> None:
|
||||
source_app = _app(app_id="workflow-app-1", mode=AppMode.WORKFLOW)
|
||||
run, execution, binding = _workflow_run(), _node_execution(), _workflow_binding()
|
||||
execution.process_data = json.dumps({"workflow_agent_binding_id": binding.id})
|
||||
if mismatch == "run-tenant":
|
||||
run.tenant_id = "other-tenant"
|
||||
elif mismatch == "source-tenant":
|
||||
source_app.tenant_id = "other-tenant"
|
||||
elif mismatch == "binding-id":
|
||||
execution.process_data = json.dumps({"workflow_agent_binding_id": "missing-binding"})
|
||||
else:
|
||||
run.app_id, run.workflow_id = "unrelated-app", "unrelated-workflow"
|
||||
sqlite_session.add_all([source_app, run, execution, binding])
|
||||
sqlite_session.commit()
|
||||
service = AgentObservabilityService(sqlite_session)
|
||||
params = AgentLogQueryParams(sources=("workflow",))
|
||||
|
||||
assert service.list_logs(app=_app(app_id="agent-app"), agent_id="agent-1", params=params)["total"] == 0
|
||||
assert (
|
||||
service.list_log_messages(
|
||||
app=_app(app_id="agent-app"), agent_id="agent-1", conversation_id=execution.id, params=params
|
||||
)["total"]
|
||||
== 0
|
||||
)
|
||||
|
||||
|
||||
def test_list_workflow_messages_uses_node_execution_identity(
|
||||
monkeypatch: pytest.MonkeyPatch, sqlite_session: Session
|
||||
) -> None:
|
||||
|
||||
@@ -0,0 +1,215 @@
|
||||
"""Real SQL coverage for Agent statistics inside a Workflow Tool invocation."""
|
||||
|
||||
import json
|
||||
from datetime import datetime
|
||||
from decimal import Decimal
|
||||
from sqlite3 import Connection
|
||||
|
||||
import pytest
|
||||
from sqlalchemy import select
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from core.workflow.node_execution_process_data import WORKFLOW_AGENT_BINDING_ID_KEY
|
||||
from models.enums import ConversationFromSource, InvokeFrom
|
||||
from models.model import App, AppMode, Message
|
||||
from models.workflow import WorkflowNodeExecutionModel, WorkflowNodeExecutionTriggeredFrom, WorkflowRun, WorkflowType
|
||||
from services.agent import observability_service as observability_service_module
|
||||
from services.agent.observability_service import AgentObservabilityService, AgentStatisticsQueryParams
|
||||
from tests.unit_tests.config_override import apply_config_overrides
|
||||
from tests.unit_tests.services.agent.test_agent_observability_service import (
|
||||
_app,
|
||||
_conversation,
|
||||
_message,
|
||||
_node_execution,
|
||||
_workflow_binding,
|
||||
_workflow_run,
|
||||
)
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def statistics_service(monkeypatch: pytest.MonkeyPatch, sqlite_session: Session) -> AgentObservabilityService:
|
||||
# Keep the production joins/aggregation; adapt only MySQL functions to SQLite.
|
||||
apply_config_overrides(monkeypatch, DB_TYPE="mysql")
|
||||
monkeypatch.setattr(observability_service_module, "convert_datetime_to_date", lambda field: f"DATE({field})")
|
||||
connection = sqlite_session.connection().connection.dbapi_connection
|
||||
assert isinstance(connection, Connection)
|
||||
connection.create_function("JSON_UNQUOTE", 1, lambda value: value)
|
||||
return AgentObservabilityService(sqlite_session)
|
||||
|
||||
|
||||
def _seed_nested_agent(sqlite_session: Session, *, workflow_type: WorkflowType) -> App:
|
||||
agent_app = _app(app_id="agent-app")
|
||||
outer_app = _app(
|
||||
app_id="outer-app", mode=AppMode.ADVANCED_CHAT if workflow_type == WorkflowType.CHAT else AppMode.WORKFLOW
|
||||
)
|
||||
source_app = _app(app_id="source-app", mode=AppMode.WORKFLOW)
|
||||
run = _workflow_run(workflow_type=workflow_type)
|
||||
run.app_id = outer_app.id
|
||||
run.workflow_id = "outer-workflow"
|
||||
run.version = "outer-v1"
|
||||
source_binding = _workflow_binding(app_id=source_app.id, binding_id="source-binding-v1", node_id="agent-node")
|
||||
source_binding.workflow_id = "source-workflow"
|
||||
source_binding.workflow_version = "source-v1"
|
||||
unused_binding = _workflow_binding(app_id=source_app.id, binding_id="source-binding-v2", node_id="agent-node")
|
||||
unused_binding.workflow_id = source_binding.workflow_id
|
||||
unused_binding.workflow_version = "source-v2"
|
||||
outer_binding = _workflow_binding(app_id=outer_app.id, node_id="agent-node")
|
||||
outer_binding.workflow_id = run.workflow_id
|
||||
outer_binding.workflow_version = run.version
|
||||
other_agent_binding = _workflow_binding(app_id=source_app.id, node_id="other-agent-node")
|
||||
other_agent_binding.workflow_id = source_binding.workflow_id
|
||||
other_agent_binding.workflow_version = source_binding.workflow_version
|
||||
other_agent_binding.agent_id = "other-agent"
|
||||
sqlite_session.add_all(
|
||||
[agent_app, outer_app, source_app, run, source_binding, unused_binding, outer_binding, other_agent_binding]
|
||||
)
|
||||
|
||||
for index, (binding, tokens, completion_tokens, price, latency) in enumerate(
|
||||
[
|
||||
(source_binding, 7, 2, "0.01", 2.0),
|
||||
(source_binding, 11, 4, "0.02", 3.0),
|
||||
(other_agent_binding, 900, 800, "9.00", 90.0),
|
||||
]
|
||||
):
|
||||
execution = _node_execution(execution_id=f"tool-agent-execution-{index}")
|
||||
execution.app_id = source_app.id
|
||||
execution.workflow_id = source_binding.workflow_id
|
||||
execution.node_id = binding.node_id
|
||||
execution.triggered_from = WorkflowNodeExecutionTriggeredFrom.WORKFLOW_TOOL
|
||||
execution.process_data = json.dumps({WORKFLOW_AGENT_BINDING_ID_KEY: binding.id})
|
||||
execution.elapsed_time = latency
|
||||
execution.execution_metadata = json.dumps(
|
||||
{
|
||||
"agent_log": {
|
||||
"agent_backend": {
|
||||
"usage": {"total_tokens": tokens, "completion_tokens": completion_tokens, "total_price": price}
|
||||
}
|
||||
}
|
||||
}
|
||||
)
|
||||
sqlite_session.add(execution)
|
||||
|
||||
if workflow_type == WorkflowType.CHAT:
|
||||
conversation = _conversation(app_id=outer_app.id)
|
||||
conversation.mode = AppMode.ADVANCED_CHAT
|
||||
conversation.from_source = ConversationFromSource.API
|
||||
message = _message(app_id=outer_app.id, created_at=datetime(2026, 7, 22, 8, 0))
|
||||
message.workflow_run_id = run.id
|
||||
message.from_source = ConversationFromSource.API
|
||||
message.invoke_from = InvokeFrom.WEB_APP
|
||||
message.from_account_id = None
|
||||
message.from_end_user_id = "end-user-1"
|
||||
sqlite_session.add_all([conversation, message])
|
||||
sqlite_session.commit()
|
||||
return agent_app
|
||||
|
||||
|
||||
@pytest.mark.parametrize("workflow_type", [WorkflowType.WORKFLOW, WorkflowType.CHAT])
|
||||
def test_nested_agent_daily_statistics_preserve_source_usage_and_date(
|
||||
sqlite_session: Session, statistics_service: AgentObservabilityService, workflow_type: WorkflowType
|
||||
) -> None:
|
||||
app = _seed_nested_agent(sqlite_session, workflow_type=workflow_type)
|
||||
# Workflow metrics remain grouped by the outer run's date; Chatflow metrics
|
||||
# retain the original message's date and usage. Child nodes executed July 23.
|
||||
day = 22 if workflow_type == WorkflowType.CHAT else 21
|
||||
date = f"2026-07-{day}"
|
||||
expected_tokens = 7 if workflow_type == WorkflowType.CHAT else 18
|
||||
expected_price = Decimal("0.0001") if workflow_type == WorkflowType.CHAT else Decimal("0.03")
|
||||
expected_latency_ms = 1250.0 if workflow_type == WorkflowType.CHAT else 5000.0
|
||||
expected_tps = 3.2 if workflow_type == WorkflowType.CHAT else 1.2
|
||||
|
||||
for source in ("workflow:source-app", "workflow:source-app:source-workflow:source-v1:agent-node", "workflow"):
|
||||
payload = statistics_service.get_statistics_summary(
|
||||
app=app,
|
||||
agent_id="agent-1",
|
||||
params=AgentStatisticsQueryParams(
|
||||
source=source,
|
||||
start=datetime(2026, 7, day),
|
||||
end=datetime(2026, 7, day + 1),
|
||||
),
|
||||
)
|
||||
|
||||
summary = payload["summary"]
|
||||
assert summary["total_messages"] == 1
|
||||
assert summary["total_conversations"] == 1
|
||||
assert summary["total_end_users"] == 1
|
||||
assert summary["total_tokens"] == expected_tokens
|
||||
assert Decimal(summary["total_price"]) == expected_price
|
||||
assert summary["average_response_time"] == expected_latency_ms
|
||||
assert summary["tokens_per_second"] == expected_tps
|
||||
assert payload["charts"]["daily_messages"] == [{"date": date, "message_count": 1}]
|
||||
assert payload["charts"]["average_response_time"] == [{"date": date, "latency": expected_latency_ms}]
|
||||
|
||||
for unrelated_source in (
|
||||
"workflow:outer-app",
|
||||
"workflow:source-app:source-workflow:source-v2:agent-node",
|
||||
"workflow:source-app:source-workflow:source-v1:other-agent-node",
|
||||
):
|
||||
payload = statistics_service.get_statistics_summary(
|
||||
app=app, agent_id="agent-1", params=AgentStatisticsQueryParams(source=unrelated_source)
|
||||
)
|
||||
assert payload["summary"]["total_messages"] == 0
|
||||
assert payload["summary"]["total_tokens"] == 0
|
||||
assert payload["charts"]["daily_messages"] == []
|
||||
|
||||
|
||||
@pytest.mark.parametrize("workflow_type", [WorkflowType.WORKFLOW, WorkflowType.CHAT])
|
||||
@pytest.mark.parametrize("invalid_owner", ["outer_tenant", "source_tenant", "node_app", "node_workflow", "binding_id"])
|
||||
def test_nested_agent_statistics_reject_mismatched_ownership(
|
||||
sqlite_session: Session,
|
||||
statistics_service: AgentObservabilityService,
|
||||
workflow_type: WorkflowType,
|
||||
invalid_owner: str,
|
||||
) -> None:
|
||||
app = _seed_nested_agent(sqlite_session, workflow_type=workflow_type)
|
||||
if invalid_owner == "outer_tenant":
|
||||
run = sqlite_session.get(WorkflowRun, "workflow-run-1")
|
||||
assert run is not None
|
||||
run.tenant_id = "other-tenant"
|
||||
elif invalid_owner == "source_tenant":
|
||||
source_app = sqlite_session.get(App, "source-app")
|
||||
assert source_app is not None
|
||||
source_app.tenant_id = "other-tenant"
|
||||
else:
|
||||
executions = sqlite_session.scalars(
|
||||
select(WorkflowNodeExecutionModel).where(
|
||||
WorkflowNodeExecutionModel.id.in_(["tool-agent-execution-0", "tool-agent-execution-1"])
|
||||
)
|
||||
).all()
|
||||
assert len(executions) == 2
|
||||
for execution in executions:
|
||||
if invalid_owner == "node_app":
|
||||
execution.app_id = "other-app"
|
||||
elif invalid_owner == "node_workflow":
|
||||
execution.workflow_id = "other-workflow"
|
||||
else:
|
||||
execution.process_data = json.dumps(
|
||||
{WORKFLOW_AGENT_BINDING_ID_KEY: "binding-source-app-other-agent-node"}
|
||||
)
|
||||
sqlite_session.commit()
|
||||
|
||||
payload = statistics_service.get_statistics_summary(
|
||||
app=app, agent_id="agent-1", params=AgentStatisticsQueryParams(source="workflow:source-app")
|
||||
)
|
||||
|
||||
assert payload["summary"]["total_messages"] == 0
|
||||
assert payload["summary"]["total_tokens"] == 0
|
||||
assert payload["charts"]["daily_messages"] == []
|
||||
|
||||
|
||||
def test_nested_agent_chat_statistics_reject_message_from_another_app(
|
||||
sqlite_session: Session, statistics_service: AgentObservabilityService
|
||||
) -> None:
|
||||
app = _seed_nested_agent(sqlite_session, workflow_type=WorkflowType.CHAT)
|
||||
message = sqlite_session.get(Message, "message-1")
|
||||
assert message is not None
|
||||
message.app_id = "other-app"
|
||||
sqlite_session.commit()
|
||||
|
||||
payload = statistics_service.get_statistics_summary(
|
||||
app=app, agent_id="agent-1", params=AgentStatisticsQueryParams(source="workflow:source-app")
|
||||
)
|
||||
|
||||
assert payload["summary"]["total_messages"] == 0
|
||||
assert payload["summary"]["total_tokens"] == 0
|
||||
assert payload["charts"]["daily_messages"] == []
|
||||
Reference in New Issue
Block a user