This commit is contained in:
a2a-platform
2026-06-12 18:15:55 +00:00
parent 81ea321f3c
commit 88f343d5d6
3 changed files with 4878 additions and 1 deletions

130
agent.py
View File

@@ -1,6 +1,7 @@
from __future__ import annotations
import base64
import builtins
import json
import re
from typing import Any
@@ -1625,6 +1626,135 @@ class Openpannel(A2AAgent):
body=body,
)
@skill(
name="get_project_geo_from_events",
description=(
"Summarize OpenPanel visitor geography from raw exported events. "
"Use this instead of the stale /insights/{projectId}/traffic/geo route."
),
tags=("Insights", "Geo"),
timeout_seconds=120,
)
async def get_project_geo_from_events(
self,
ctx: RunContext,
projectId: str,
start: str | None = None,
end: str | None = None,
range: str = "today",
limit: int = 250,
max_pages: int = 4,
) -> dict[str, Any]:
per_page = max(1, min(int(limit), 1000))
pages = max(1, min(int(max_pages), 10))
aggregate: dict[str, dict[str, Any]] = {}
samples: list[dict[str, Any]] = []
total_events = 0
for page in builtins.range(1, pages + 1):
params: dict[str, Any] = {
"projectId": projectId,
"page": page,
"limit": per_page,
}
if start or end:
if start:
params["start"] = start
if end:
params["end"] = end
else:
params["range"] = range
response = await self._request(
ctx,
"get_export_events",
parameters=params,
body=None,
)
if not response.get("ok"):
return response
result = response.get("result") if isinstance(response, dict) else None
data = result.get("data", []) if isinstance(result, dict) else []
meta = result.get("meta", {}) if isinstance(result, dict) else {}
if not isinstance(data, list) or not data:
break
total_events += len(data)
for event in data:
if not isinstance(event, dict):
continue
country = self._clean_geo_value(event.get("country"))
region = self._clean_geo_value(event.get("region"))
city = self._clean_geo_value(event.get("city"))
key = "|".join([country, region, city])
bucket = aggregate.setdefault(
key,
{
"country": country or None,
"region": region or None,
"city": city or None,
"events": 0,
"sessions": set(),
"latest": None,
},
)
bucket["events"] += 1
session_id = event.get("sessionId") or event.get("session_id")
if session_id:
bucket["sessions"].add(str(session_id))
created_at = event.get("createdAt") or event.get("created_at")
if created_at and (bucket["latest"] is None or str(created_at) > str(bucket["latest"])):
bucket["latest"] = created_at
if len(samples) < 10:
samples.append(
{
"createdAt": created_at,
"name": event.get("name"),
"path": event.get("path"),
"country": country or None,
"region": region or None,
"city": city or None,
"browser": event.get("browser"),
"os": event.get("os"),
}
)
if isinstance(meta, dict):
current = int(meta.get("current") or page)
page_count = int(meta.get("pages") or page)
if current >= page_count:
break
rows = []
for bucket in aggregate.values():
rows.append(
{
"country": bucket["country"],
"region": bucket["region"],
"city": bucket["city"],
"events": bucket["events"],
"sessions": len(bucket["sessions"]),
"latest": bucket["latest"],
}
)
rows.sort(key=lambda row: (row["latest"] or "", row["events"]), reverse=True)
return {
"ok": True,
"operation_id": "get_project_geo_from_events",
"projectId": projectId,
"range": {"start": start, "end": end, "range": range},
"events_scanned": total_events,
"geo": rows,
"sample_events": samples,
}
@staticmethod
def _clean_geo_value(value: Any) -> str:
text = str(value or "").replace("\x00", "").strip()
return text
def _operation_tools(self, ctx: RunContext) -> list[StructuredTool]:
tools: list[StructuredTool] = []
for operation_id, operation in OPERATIONS.items():