Files
2026-07-13 12:43:34 +08:00

200 lines
9.4 KiB
Python

from typing import List, Dict, Union, Any
from sqlalchemy import text, func, and_
from sqlalchemy.orm import Session
from superagi.models.events import Event
class AnalyticsHelper:
def __init__(self, session: Session, organisation_id: int):
self.session = session
self.organisation_id = organisation_id
def calculate_run_completed_metrics(self) -> Dict[str, Dict[str, Union[int, List[Dict[str, int]]]]]:
agent_model_query = self.session.query(
Event.event_property['model'].label('model'),
Event.agent_id
).filter_by(event_name="agent_created", org_id=self.organisation_id).subquery()
agent_runs_query = self.session.query(
agent_model_query.c.model,
func.count(Event.id).label('runs')
).join(Event, and_(Event.agent_id == agent_model_query.c.agent_id, Event.org_id == self.organisation_id)).filter(Event.event_name.in_(['run_completed', 'run_iteration_limit_crossed'])).group_by(agent_model_query.c.model).subquery()
agent_tokens_query = self.session.query(
agent_model_query.c.model,
func.sum(text("(event_property->>'tokens_consumed')::int")).label('tokens')
).join(Event, and_(Event.agent_id == agent_model_query.c.agent_id, Event.org_id == self.organisation_id)).filter(Event.event_name.in_(['run_completed', 'run_iteration_limit_crossed'])).group_by(agent_model_query.c.model).subquery()
agent_count_query = self.session.query(
agent_model_query.c.model,
func.count(agent_model_query.c.agent_id).label('agents')
).group_by(agent_model_query.c.model).subquery()
agents = self.session.query(agent_count_query).all()
runs = self.session.query(agent_runs_query).all()
tokens = self.session.query(agent_tokens_query).all()
metrics = {
'agent_details': {
'total_agents': sum([item.agents for item in agents]),
'model_metrics': [{'name': item.model, 'value': item.agents} for item in agents]
},
'run_details': {
'total_runs': sum([item.runs for item in runs]),
'model_metrics': [{'name': item.model, 'value': item.runs} for item in runs]
},
'tokens_details': {
'total_tokens': sum([item.tokens for item in tokens]),
'model_metrics': [{'name': item.model, 'value': item.tokens} for item in tokens]
},
}
return metrics
def fetch_agent_data(self) -> Dict[str, List[Dict[str, Any]]]:
agent_subquery = self.session.query(
Event.agent_id,
Event.event_property['agent_name'].label('agent_name'),
Event.event_property['model'].label('model')
).filter_by(event_name="agent_created", org_id=self.organisation_id).subquery()
run_subquery = self.session.query(
Event.agent_id,
func.sum(text("(event_property->>'tokens_consumed')::int")).label('total_tokens'),
func.sum(text("(event_property->>'calls')::int")).label('total_calls'),
func.count(Event.id).label('runs_completed'),
).filter(and_(Event.event_name.in_(['run_completed', 'run_iteration_limit_crossed']), Event.org_id == self.organisation_id)).group_by(Event.agent_id).subquery()
tool_subquery = self.session.query(
Event.agent_id,
func.array_agg(Event.event_property['tool_name'].distinct()).label('tools_used'),
).filter_by(event_name="tool_used", org_id=self.organisation_id).group_by(Event.agent_id).subquery()
start_time_subquery = self.session.query(
Event.agent_id,
Event.event_property['agent_execution_id'].label('agent_execution_id'),
func.min(func.extract('epoch', Event.created_at)).label('start_time')
).filter_by(event_name="run_created", org_id=self.organisation_id).group_by(Event.agent_id, Event.event_property['agent_execution_id']).subquery()
end_time_subquery = self.session.query(
Event.agent_id,
Event.event_property['agent_execution_id'].label('agent_execution_id'),
func.max(func.extract('epoch', Event.created_at)).label('end_time')
).filter(and_(Event.event_name.in_(['run_completed', 'run_iteration_limit_crossed']), Event.org_id == self.organisation_id)).group_by(Event.agent_id, Event.event_property['agent_execution_id']).subquery()
time_diff_subquery = self.session.query(
start_time_subquery.c.agent_id,
(func.avg(end_time_subquery.c.end_time - start_time_subquery.c.start_time)).label('avg_run_time')
).join(end_time_subquery, start_time_subquery.c.agent_execution_id == end_time_subquery.c.agent_execution_id). \
group_by(start_time_subquery.c.agent_id).subquery()
query = self.session.query(
agent_subquery.c.agent_id,
agent_subquery.c.agent_name,
agent_subquery.c.model,
run_subquery.c.total_tokens,
run_subquery.c.total_calls,
run_subquery.c.runs_completed,
tool_subquery.c.tools_used,
time_diff_subquery.c.avg_run_time
).outerjoin(run_subquery, run_subquery.c.agent_id == agent_subquery.c.agent_id) \
.outerjoin(tool_subquery, tool_subquery.c.agent_id == agent_subquery.c.agent_id) \
.outerjoin(time_diff_subquery, time_diff_subquery.c.agent_id == agent_subquery.c.agent_id)
result = query.all()
agent_details = [{
"name": row.agent_name,
"agent_id": row.agent_id,
"runs_completed": row.runs_completed if row.runs_completed else 0,
"total_calls": row.total_calls if row.total_calls else 0,
"total_tokens": row.total_tokens if row.total_tokens else 0,
"tools_used": row.tools_used,
"model_name": row.model,
"avg_run_time": row.avg_run_time if row.avg_run_time else 0,
} for row in result]
return {'agent_details': agent_details}
def fetch_agent_runs(self, agent_id: int) -> List[Dict[str, int]]:
agent_runs = []
completed_subquery = self.session.query(
Event.event_property['agent_execution_id'].label('completed_agent_execution_id'),
Event.event_property['tokens_consumed'].label('tokens_consumed'),
Event.event_property['calls'].label('calls'),
Event.updated_at
).filter(Event.event_name.in_(['run_completed','run_iteration_limit_crossed']), Event.agent_id == agent_id, Event.org_id == self.organisation_id).subquery()
created_subquery = self.session.query(
Event.event_property['agent_execution_id'].label('created_agent_execution_id'),
Event.event_property['agent_execution_name'].label('agent_execution_name'),
Event.created_at
).filter(Event.event_name == "run_created", Event.agent_id == agent_id, Event.org_id == self.organisation_id).subquery()
query = self.session.query(
created_subquery.c.agent_execution_name,
completed_subquery.c.tokens_consumed,
completed_subquery.c.calls,
created_subquery.c.created_at,
completed_subquery.c.updated_at
).join(completed_subquery, completed_subquery.c.completed_agent_execution_id == created_subquery.c.created_agent_execution_id)
result = query.all()
agent_runs = [{
'name': row.agent_execution_name,
'tokens_consumed': int(row.tokens_consumed) if row.tokens_consumed else 0,
'calls': int(row.calls) if row.calls else 0,
'created_at': row.created_at,
'updated_at': row.updated_at
} for row in result]
return agent_runs
def get_active_runs(self) -> List[Dict[str, str]]:
running_executions = []
end_event_subquery = self.session.query(
Event.event_property['agent_execution_id'].label('agent_execution_id'),
).filter(
Event.event_name.in_(['run_completed', 'run_iteration_limit_crossed']),
Event.org_id == self.organisation_id
).subquery()
start_subquery = self.session.query(
Event.event_property['agent_execution_id'].label('agent_execution_id'),
Event.event_property['agent_execution_name'].label('agent_execution_name'),
Event.created_at,
Event.agent_id
).filter_by(event_name="run_created", org_id = self.organisation_id).subquery()
agent_created_subquery = self.session.query(
Event.event_property['agent_name'].label('agent_name'),
Event.agent_id
).filter_by(event_name="agent_created", org_id = self.organisation_id).subquery()
query = self.session.query(
start_subquery.c.agent_execution_name,
start_subquery.c.created_at,
agent_created_subquery.c.agent_name
).select_from(start_subquery)
query = query.outerjoin(end_event_subquery, start_subquery.c.agent_execution_id == end_event_subquery.c.agent_execution_id).filter(end_event_subquery.c.agent_execution_id == None)
query = query.join(agent_created_subquery, start_subquery.c.agent_id == agent_created_subquery.c.agent_id)
result = query.all()
running_executions = [{
'name': row.agent_execution_name,
'created_at': row.created_at,
'agent_name': row.agent_name or 'Unknown',
} for row in result]
return running_executions