a2a-source-edit: write agent.py

This commit is contained in:
a2a-cloud
2026-07-19 19:51:27 +00:00
parent 91a048f592
commit 4eaa384ef6

810
agent.py
View File

@@ -1,189 +1,717 @@
"""production-proof-v116-high-utili-7255-1 agent.
"""Production Proof v116 dashboard agent.
Starter stack:
- DeepAgents for tool-calling orchestration
- Caller-provided LLM credentials via ctx.llm
- A tiny model-call middleware hook you can replace with tracing,
routing, rate limits, or policy checks
A deterministic full-stack A2A product for software agencies: paste recurring
ops updates, get filters, trend data, exceptions, durable run history, a CSV
artifact, MCP-callable tools, and persisted execution receipts.
"""
from __future__ import annotations
import csv
import hashlib
import io
import json
from pathlib import Path
import os
import re
import uuid
from contextlib import contextmanager
from datetime import UTC, datetime
from typing import Any
from pydantic import BaseModel
from pydantic import BaseModel, Field
import a2a_pack as a2a
from a2a_pack import (
A2AAgent,
AccountAccess,
AgentDatabase,
AgentDatabaseEnv,
AgentDatabaseMigrations,
AgentPlatformResources,
LLMProvisioning,
{{ auth_type }},
PlatformUserAuth,
Pricing,
Resources,
RunContext,
WorkspaceAccess,
WorkspaceMode,
)
from a2a_pack.context import LLMCreds
DEFAULT_UPDATES = """2026-07-01 | Acme Portal | green | Sprint 18 shipped client dashboard, 3 blockers cleared, utilization 81%, margin 34%.
2026-07-02 | Beacon CRM | yellow | API integration waiting on customer credentials; utilization 74%, margin 27%, risk: delayed approval.
2026-07-03 | Cedar Mobile | red | Escalated crash fix is overdue; utilization 92%, margin 18%, blocker: App Store review.
2026-07-04 | Delta RevOps | green | Retainer reporting automated; utilization 68%, margin 39%, next: expansion proposal."""
MAX_UPDATE_CHARS = 6000
MAX_ROWS = 50
DATABASE_NAME = "production-proof-v116-high-utili-7255-1-data"
class ProductionProofV116HighUtili72551Config(BaseModel):
pass
default_window_label: str = "This week"
SYSTEM_PROMPT = """\
You are a compact tool-calling agent.
Use the text_stats tool when the user asks about text, counts, summaries,
or anything where exact length/word numbers would help. Mention tool results
briefly instead of dumping raw JSON.
"""
RUNTIME_SKILLS_DIR = "production-proof-v116-high-utili-7255-1/.deepagents/skills/"
DEEPAGENTS_RECURSION_LIMIT = 500
class DashboardRow(BaseModel):
date: str
client: str
status: str
update: str
utilization: int | None = None
margin: int | None = None
exception: bool = False
class ProductionProofV116HighUtili72551(A2AAgent[ProductionProofV116HighUtili72551Config, {{ auth_type }}]):
class ArtifactDescriptor(BaseModel):
name: str
media_type: str
uri: str
size_bytes: int
content: str
class DashboardResult(BaseModel):
status: str
run_id: str
receipt_id: str
created_at: str
agency_name: str
window_label: str
filters: dict[str, Any]
trends: dict[str, Any]
exceptions: list[dict[str, Any]]
dashboard_data: list[DashboardRow]
history: list[dict[str, Any]]
artifact: ArtifactDescriptor
summary: str
class HistoryResult(BaseModel):
status: str
history: list[dict[str, Any]]
class ScheduledRollupResult(BaseModel):
status: str
run_key: str
run_id: str
receipt_id: str
already_processed: bool
summary: str
dashboard_data: list[DashboardRow]
class ProductionProofV116HighUtili72551(
A2AAgent[ProductionProofV116HighUtili72551Config, PlatformUserAuth]
):
name = "production-proof-v116-high-utili-7255-1"
description = "One-page dashboard for software agencies to turn recurring operational updates into filters, trends, exceptions, durable history, downloads, receipts, and MCP tools."
description = (
"One-page dashboard for software agencies to turn recurring operational "
"updates into filters, trends, exceptions, durable history, downloads, "
"receipts, and MCP tools."
)
version = "0.1.0"
config_model = ProductionProofV116HighUtili72551Config
auth_model = {{ auth_type }}
auth_model = PlatformUserAuth
# Hosted generated agents read the caller's saved LLM credential through
# ctx.llm. The platform may proxy that credential through LiteLLM, but agent
# code never reads provider keys, LiteLLM master keys, or OPENAI_API_KEY
# directly.
# Required by the account-funded trial contract: the platform funds exactly
# three skill calls for signed-in accounts, then requires the caller's saved
# BYOK LLM credential. This product is deterministic and does not read
# ctx.llm or provider keys directly.
llm_provisioning = LLMProvisioning.PLATFORM
account_access = AccountAccess(required=True, platform_skill_calls=3, after_trial="byok")
pricing = Pricing(
price_per_call_usd=0.0,
caller_pays_llm=True,
notes="Starter agent uses the caller's saved LLM credential via ctx.llm.",
notes="Account required. First 3 platform skill calls are funded; after that BYOK is required by the platform.",
)
workspace_access = WorkspaceAccess.dynamic(
max_files=64,
allowed_modes=(WorkspaceMode.READ_ONLY, WorkspaceMode.READ_WRITE_OVERLAY),
require_reason=False,
)
tools_used = ("deepagents", "langchain")
@a2a.tool(description="Ask the starter DeepAgent to answer with tool calls when useful")
async def ask(self, ctx: RunContext[{{ auth_type }}], prompt: str) -> str:
creds = ctx.llm
await ctx.emit_progress(f"llm: {creds.model} via {creds.source}")
if not creds.api_key:
return (
"LLM key required. Add an LLM credential in Settings > LLM "
"credentials before running this agent; for local --invoke "
"runs set AGENT_LLM_KEY."
)
graph = self._build_deep_agent(ctx=ctx, creds=creds)
state = await graph.ainvoke(
{"messages": [{"role": "user", "content": prompt}]},
config={"recursion_limit": DEEPAGENTS_RECURSION_LIMIT},
resources = Resources(cpu="500m", memory="512Mi", max_runtime_seconds=120)
platform_resources = AgentPlatformResources(
databases=(
AgentDatabase(
name=DATABASE_NAME,
scope="user",
access_mode="read_write",
env=AgentDatabaseEnv(url="DATABASE_URL"),
migrations=AgentDatabaseMigrations(path="db/migrations"),
),
)
await ctx.emit_progress("deepagent finished")
return _last_message_text(state)
)
tools_used = ("postgres", "mcp", "artifacts", "scheduled-rollup")
capabilities = {
"mcp": {"tools": ["build_dashboard", "list_history", "scheduled_rollup"]},
"schedules": {
"tool": "scheduled_rollup",
"idempotency_key": "run_key",
"note": "Safe to call from a recurring platform schedule or MCP client.",
},
"artifacts": {"outputs": ["agency-ops-dashboard.csv"]},
}
def _build_deep_agent(
@a2a.tool(
description="Build the agency ops dashboard, persist history/receipt, and emit a CSV artifact.",
timeout_seconds=60,
cost_class="standard",
grant_mode="read_write_overlay",
grant_allow_patterns=("outputs/dashboard/**",),
grant_outputs_prefix="outputs/dashboard/",
grant_write_prefixes=("outputs/dashboard/",),
)
async def build_dashboard(
self,
*,
ctx: RunContext[{{ auth_type }}],
creds: LLMCreds,
) -> Any:
# Lazy imports keep `a2a card` usable before local dependencies are
# installed. `a2a deploy` installs requirements.txt during the build.
from a2a_pack.deepagents import create_a2a_deep_agent
from langchain.agents.middleware import wrap_model_call
from langchain_core.tools import tool
ctx: RunContext[PlatformUserAuth],
updates_text: str = Field(default=DEFAULT_UPDATES, max_length=MAX_UPDATE_CHARS),
agency_name: str = Field(default="Proofline Software Agency", max_length=80),
window_label: str = Field(default="This week", max_length=60),
min_severity: str = Field(default="all", pattern="^(all|green|yellow|red)$"),
) -> DashboardResult:
tenant = _tenant_key(ctx)
rows = _analyze_updates(updates_text)
filtered = _filter_rows(rows, min_severity)
trends = _build_trends(rows)
exceptions = _build_exceptions(rows)
filters = _build_filters(rows, min_severity)
summary = _summary(rows, exceptions)
csv_text = _rows_to_csv(filtered)
artifact_ref = await ctx.write_artifact(
"agency-ops-dashboard.csv",
csv_text.encode("utf-8"),
"text/csv",
)
await ctx.emit_artifact(artifact_ref)
@tool
def text_stats(text: str) -> str:
"""Return exact word, character, and line counts for text."""
words = [part for part in text.split() if part.strip()]
return json.dumps(
{
"characters": len(text),
"words": len(words),
"lines": len(text.splitlines()) or 1,
}
created_at = _now_iso()
run_id = str(uuid.uuid4())
receipt_id = str(uuid.uuid4())
input_hash = _sha256_json(
{
"updates_text": updates_text,
"agency_name": agency_name,
"window_label": window_label,
"min_severity": min_severity,
}
)
history: list[dict[str, Any]] = []
try:
history = _persist_dashboard_run(
tenant_key=tenant,
run_id=run_id,
receipt_id=receipt_id,
created_at=created_at,
skill_name="build_dashboard",
input_hash=input_hash,
agency_name=agency_name,
window_label=window_label,
min_severity=min_severity,
rows=rows,
filtered=filtered,
filters=filters,
trends=trends,
exceptions=exceptions,
summary=summary,
artifact_name=artifact_ref.name,
artifact_media_type=artifact_ref.mime_type,
artifact_size=artifact_ref.size_bytes,
)
except RuntimeError as exc:
await ctx.emit_error(str(exc), code="database_unavailable")
return DashboardResult(
status="setup_required",
run_id=run_id,
receipt_id=receipt_id,
created_at=created_at,
agency_name=agency_name,
window_label=window_label,
filters=filters,
trends=trends,
exceptions=exceptions,
dashboard_data=filtered,
history=[],
artifact=ArtifactDescriptor(
name=artifact_ref.name,
media_type=artifact_ref.mime_type,
uri=artifact_ref.uri,
size_bytes=artifact_ref.size_bytes,
content=csv_text,
),
summary=f"{summary} Database persistence needs setup: {exc}",
)
@wrap_model_call
async def log_model_call(request: Any, handler: Any) -> Any:
messages = request.state.get("messages", [])
print(
"[middleware] model_call "
f"model={creds.model} source={creds.source} messages={len(messages)}"
await ctx.emit_event(
a2a.AgentEvent(
kind="production_receipt_persisted",
payload={"receipt_id": receipt_id, "run_id": run_id},
)
return await handler(request)
)
return DashboardResult(
status="ok",
run_id=run_id,
receipt_id=receipt_id,
created_at=created_at,
agency_name=agency_name,
window_label=window_label,
filters=filters,
trends=trends,
exceptions=exceptions,
dashboard_data=filtered,
history=history,
artifact=ArtifactDescriptor(
name=artifact_ref.name,
media_type=artifact_ref.mime_type,
uri=artifact_ref.uri,
size_bytes=artifact_ref.size_bytes,
content=csv_text,
),
summary=summary,
)
backend = ctx.workspace_backend()
skill_sources = _seed_runtime_skills(backend, ctx)
# create_a2a_deep_agent resolves provider:model strings with
# langchain.init_chat_model from ctx.llm, preserving LiteLLM routing,
# provider-specific extra body, and runtime model overrides.
return create_a2a_deep_agent(
ctx,
creds=creds,
backend=backend,
skills=skill_sources or None,
tools=[text_stats],
middleware=[log_model_call],
system_prompt=SYSTEM_PROMPT,
@a2a.tool(
description="List durable dashboard history for the signed-in platform user.",
timeout_seconds=30,
cost_class="standard",
idempotent=True,
)
async def list_history(
self,
ctx: RunContext[PlatformUserAuth],
limit: int = Field(default=8, ge=1, le=25),
) -> HistoryResult:
tenant = _tenant_key(ctx)
try:
history = _load_history(tenant, limit=limit)
except RuntimeError as exc:
await ctx.emit_error(str(exc), code="database_unavailable")
return HistoryResult(status="setup_required", history=[])
return HistoryResult(status="ok", history=history)
@a2a.tool(
description="Idempotent scheduled/MCP rollup that stores at most one dashboard run per run_key.",
timeout_seconds=60,
cost_class="standard",
idempotent=True,
grant_mode="read_write_overlay",
grant_allow_patterns=("outputs/dashboard/**",),
grant_outputs_prefix="outputs/dashboard/",
grant_write_prefixes=("outputs/dashboard/",),
)
async def scheduled_rollup(
self,
ctx: RunContext[PlatformUserAuth],
run_key: str = Field(default="demo-weekly-rollup", max_length=80, pattern="^[A-Za-z0-9_.:-]+$"),
updates_text: str = Field(default=DEFAULT_UPDATES, max_length=MAX_UPDATE_CHARS),
) -> ScheduledRollupResult:
tenant = _tenant_key(ctx)
rows = _analyze_updates(updates_text)
exceptions = _build_exceptions(rows)
summary = _summary(rows, exceptions)
csv_text = _rows_to_csv(rows)
artifact_ref = await ctx.write_artifact(
"agency-ops-dashboard.csv",
csv_text.encode("utf-8"),
"text/csv",
)
await ctx.emit_artifact(artifact_ref)
created_at = _now_iso()
run_id = str(uuid.uuid4())
receipt_id = str(uuid.uuid4())
input_hash = _sha256_json({"run_key": run_key, "updates_text": updates_text})
try:
stored = _persist_scheduled_rollup(
tenant_key=tenant,
run_key=run_key,
run_id=run_id,
receipt_id=receipt_id,
created_at=created_at,
input_hash=input_hash,
rows=rows,
summary=summary,
)
except RuntimeError as exc:
await ctx.emit_error(str(exc), code="database_unavailable")
return ScheduledRollupResult(
status="setup_required",
run_key=run_key,
run_id=run_id,
receipt_id=receipt_id,
already_processed=False,
summary=f"{summary} Database persistence needs setup: {exc}",
dashboard_data=rows,
)
return ScheduledRollupResult(
status="ok",
run_key=run_key,
run_id=stored["run_id"],
receipt_id=stored["receipt_id"],
already_processed=stored["already_processed"],
summary=summary,
dashboard_data=rows,
)
def _runtime_skills_root(ctx: RunContext[Any]) -> str:
workspace = getattr(ctx, "_workspace", None)
prefixes = tuple(getattr(workspace, "write_prefixes", ()) or ())
if not prefixes:
outputs_prefix = getattr(workspace, "outputs_prefix", None)
prefixes = (outputs_prefix or "outputs/",)
prefix = str(prefixes[0]).strip("/")
return f"/{prefix}/{RUNTIME_SKILLS_DIR}" if prefix else f"/{RUNTIME_SKILLS_DIR}"
def _tenant_key(ctx: RunContext[PlatformUserAuth]) -> str:
auth = ctx.auth
stable_id = auth.user_id if auth.user_id is not None else auth.sub
if not stable_id:
raise PermissionError("stable platform identity required")
return f"user:{stable_id}"
def _seed_runtime_skills(backend: Any, ctx: RunContext[Any]) -> list[str]:
"""Copy packaged DeepAgents skills into the invocation workspace.
DeepAgents loads skills from its backend, while source-controlled
``skills/`` folders live in the image. This bridge lets generated agents
ship reusable SKILL.md bundles without giving up durable A2A workspace
files.
"""
root = Path(__file__).parent / "skills"
if not root.exists():
return []
runtime_skills_root = _runtime_skills_root(ctx)
uploads: list[tuple[str, bytes]] = []
for path in root.rglob("*"):
if path.is_file():
rel = path.relative_to(root).as_posix()
uploads.append((runtime_skills_root + rel, path.read_bytes()))
if uploads:
backend.upload_files(uploads)
return [runtime_skills_root]
return []
@contextmanager
def _db_conn():
database_url = os.environ.get("DATABASE_URL")
if not database_url:
raise RuntimeError("DATABASE_URL is not configured")
try:
import psycopg
except Exception as exc: # pragma: no cover - dependency installed in deploy image
raise RuntimeError("psycopg is not installed") from exc
options = "-c statement_timeout=5000 -c lock_timeout=3000 -c idle_in_transaction_session_timeout=10000"
try:
with psycopg.connect(database_url, autocommit=False, options=options) as conn:
yield conn
except Exception as exc: # noqa: BLE001
raise RuntimeError("managed database is unavailable") from exc
def _last_message_text(state: dict[str, Any]) -> str:
messages = state.get("messages") or []
if not messages:
return json.dumps(state, default=str)
def _persist_dashboard_run(
*,
tenant_key: str,
run_id: str,
receipt_id: str,
created_at: str,
skill_name: str,
input_hash: str,
agency_name: str,
window_label: str,
min_severity: str,
rows: list[DashboardRow],
filtered: list[DashboardRow],
filters: dict[str, Any],
trends: dict[str, Any],
exceptions: list[dict[str, Any]],
summary: str,
artifact_name: str,
artifact_media_type: str,
artifact_size: int,
) -> list[dict[str, Any]]:
payload = {
"rows": [row.model_dump() for row in rows],
"filtered_rows": [row.model_dump() for row in filtered],
"filters": filters,
"trends": trends,
"exceptions": exceptions,
"artifact": {
"name": artifact_name,
"media_type": artifact_media_type,
"size_bytes": artifact_size,
},
}
receipt_payload = {
"receipt_id": receipt_id,
"run_id": run_id,
"status": "ok",
"rows": len(filtered),
"exceptions": len(exceptions),
"artifact_name": artifact_name,
}
with _db_conn() as conn:
with conn.transaction():
conn.execute(
"""
INSERT INTO dashboard_runs (
id, tenant_key, agency_name, window_label, min_severity,
input_hash, summary, filters, trends, exceptions,
dashboard_rows, artifact_name, artifact_media_type,
artifact_size_bytes, created_at
) VALUES (%s, %s, %s, %s, %s, %s, %s, %s::jsonb, %s::jsonb, %s::jsonb, %s::jsonb, %s, %s, %s, %s)
""",
(
run_id,
tenant_key,
agency_name,
window_label,
min_severity,
input_hash,
summary,
json.dumps(filters),
json.dumps(trends),
json.dumps(exceptions),
json.dumps(payload["filtered_rows"]),
artifact_name,
artifact_media_type,
artifact_size,
created_at,
),
)
conn.execute(
"""
INSERT INTO execution_receipts (
id, tenant_key, run_id, skill_name, status, input_hash,
result_overview, created_at
) VALUES (%s, %s, %s, %s, %s, %s, %s::jsonb, %s)
""",
(
receipt_id,
tenant_key,
run_id,
skill_name,
"ok",
input_hash,
json.dumps(receipt_payload),
created_at,
),
)
conn.commit()
return _load_history(tenant_key, limit=8)
content = getattr(messages[-1], "content", None)
if isinstance(content, str):
return content
if isinstance(content, list):
parts: list[str] = []
for item in content:
if isinstance(item, dict):
text = item.get("text") or item.get("content")
if text:
parts.append(str(text))
elif item:
parts.append(str(item))
return "\n".join(parts) if parts else json.dumps(content, default=str)
return str(content or messages[-1])
def _persist_scheduled_rollup(
*,
tenant_key: str,
run_key: str,
run_id: str,
receipt_id: str,
created_at: str,
input_hash: str,
rows: list[DashboardRow],
summary: str,
) -> dict[str, Any]:
with _db_conn() as conn:
with conn.transaction():
existing = conn.execute(
"SELECT run_id, receipt_id FROM scheduled_rollups WHERE tenant_key = %s AND run_key = %s",
(tenant_key, run_key),
).fetchone()
if existing:
return {"run_id": str(existing[0]), "receipt_id": str(existing[1]), "already_processed": True}
conn.execute(
"""
INSERT INTO scheduled_rollups (tenant_key, run_key, run_id, receipt_id, created_at)
VALUES (%s, %s, %s, %s, %s)
""",
(tenant_key, run_key, run_id, receipt_id, created_at),
)
conn.execute(
"""
INSERT INTO dashboard_runs (
id, tenant_key, agency_name, window_label, min_severity,
input_hash, summary, filters, trends, exceptions,
dashboard_rows, artifact_name, artifact_media_type,
artifact_size_bytes, created_at
) VALUES (%s, %s, %s, %s, %s, %s, %s, %s::jsonb, %s::jsonb, %s::jsonb, %s::jsonb, %s, %s, %s, %s)
""",
(
run_id,
tenant_key,
"Scheduled agency ops",
run_key,
"all",
input_hash,
summary,
json.dumps(_build_filters(rows, "all")),
json.dumps(_build_trends(rows)),
json.dumps(_build_exceptions(rows)),
json.dumps([row.model_dump() for row in rows]),
"agency-ops-dashboard.csv",
"text/csv",
len(_rows_to_csv(rows).encode("utf-8")),
created_at,
),
)
conn.execute(
"""
INSERT INTO execution_receipts (
id, tenant_key, run_id, skill_name, status, input_hash,
result_overview, created_at
) VALUES (%s, %s, %s, %s, %s, %s, %s::jsonb, %s)
""",
(
receipt_id,
tenant_key,
run_id,
"scheduled_rollup",
"ok",
input_hash,
json.dumps({"run_key": run_key, "rows": len(rows), "status": "ok"}),
created_at,
),
)
conn.commit()
return {"run_id": run_id, "receipt_id": receipt_id, "already_processed": False}
def _load_history(tenant_key: str, *, limit: int) -> list[dict[str, Any]]:
with _db_conn() as conn:
rows = conn.execute(
"""
SELECT id, agency_name, window_label, summary, jsonb_array_length(dashboard_rows),
jsonb_array_length(exceptions), artifact_name, created_at
FROM dashboard_runs
WHERE tenant_key = %s
ORDER BY created_at DESC
LIMIT %s
""",
(tenant_key, limit),
).fetchall()
return [
{
"run_id": str(row[0]),
"agency_name": row[1],
"window_label": row[2],
"summary": row[3],
"row_count": int(row[4] or 0),
"exception_count": int(row[5] or 0),
"artifact_name": row[6],
"created_at": row[7].isoformat() if hasattr(row[7], "isoformat") else str(row[7]),
}
for row in rows
]
def _analyze_updates(updates_text: str) -> list[DashboardRow]:
text = (updates_text or "").strip()[:MAX_UPDATE_CHARS]
if not text:
text = DEFAULT_UPDATES
parsed: list[DashboardRow] = []
for raw_line in text.splitlines()[:MAX_ROWS]:
line = raw_line.strip(" -\t")
if not line:
continue
parts = [part.strip() for part in re.split(r"\s+\|\s+", line, maxsplit=3)]
if len(parts) >= 4:
date, client, status_raw, update = parts[0], parts[1], parts[2], parts[3]
else:
date = _extract_date(line) or "unspecified"
client = _extract_client(line)
status_raw = _infer_status(line)
update = line
status = _normalize_status(status_raw, update)
utilization = _extract_percent(update, "utilization")
margin = _extract_percent(update, "margin")
exception = status == "red" or bool(re.search(r"\b(blocker|blocked|overdue|escalat|risk|delay|waiting)\b", update, re.I))
parsed.append(
DashboardRow(
date=date[:24],
client=client[:80] or "Unassigned client",
status=status,
update=update[:500],
utilization=utilization,
margin=margin,
exception=exception,
)
)
if not parsed:
parsed = _analyze_updates(DEFAULT_UPDATES)
return parsed
def _filter_rows(rows: list[DashboardRow], min_severity: str) -> list[DashboardRow]:
if min_severity == "all":
return rows
severity = {"green": 1, "yellow": 2, "red": 3}
floor = severity.get(min_severity, 1)
return [row for row in rows if severity.get(row.status, 1) >= floor]
def _build_filters(rows: list[DashboardRow], min_severity: str) -> dict[str, Any]:
return {
"status": sorted({row.status for row in rows}),
"clients": sorted({row.client for row in rows}),
"active_min_severity": min_severity,
"exception_only_available": any(row.exception for row in rows),
}
def _build_trends(rows: list[DashboardRow]) -> dict[str, Any]:
total = max(len(rows), 1)
by_status = {status: sum(1 for row in rows if row.status == status) for status in ("green", "yellow", "red")}
utilizations = [row.utilization for row in rows if row.utilization is not None]
margins = [row.margin for row in rows if row.margin is not None]
return {
"total_updates": len(rows),
"status_counts": by_status,
"exception_rate": round(sum(1 for row in rows if row.exception) / total, 2),
"average_utilization": round(sum(utilizations) / len(utilizations), 1) if utilizations else None,
"average_margin": round(sum(margins) / len(margins), 1) if margins else None,
"trend_points": [
{"label": row.date, "client": row.client, "status": row.status, "utilization": row.utilization or 0}
for row in rows
],
}
def _build_exceptions(rows: list[DashboardRow]) -> list[dict[str, Any]]:
return [
{
"client": row.client,
"date": row.date,
"status": row.status,
"reason": _exception_reason(row.update, row.status),
"update": row.update,
}
for row in rows
if row.exception
][:20]
def _summary(rows: list[DashboardRow], exceptions: list[dict[str, Any]]) -> str:
red = sum(1 for row in rows if row.status == "red")
yellow = sum(1 for row in rows if row.status == "yellow")
return f"Built dashboard from {len(rows)} updates: {red} red, {yellow} yellow, {len(exceptions)} exception(s) needing follow-up."
def _rows_to_csv(rows: list[DashboardRow]) -> str:
buf = io.StringIO()
writer = csv.DictWriter(
buf,
fieldnames=["date", "client", "status", "utilization", "margin", "exception", "update"],
)
writer.writeheader()
for row in rows:
writer.writerow(row.model_dump())
return buf.getvalue()
def _normalize_status(status: str, update: str) -> str:
combined = f"{status} {update}".lower()
if "red" in combined or re.search(r"\b(overdue|blocked|blocker|escalat|critical)\b", combined):
return "red"
if "yellow" in combined or re.search(r"\b(risk|delay|waiting|watch)\b", combined):
return "yellow"
return "green"
def _extract_percent(text: str, label: str) -> int | None:
match = re.search(label + r"\D{0,12}(\d{1,3})\s*%", text, re.I)
if not match:
return None
value = int(match.group(1))
return value if 0 <= value <= 100 else None
def _extract_date(text: str) -> str | None:
match = re.search(r"\b(20\d{2}-\d{2}-\d{2})\b", text)
return match.group(1) if match else None
def _extract_client(text: str) -> str:
match = re.search(r"(?:client|for)\s+([A-Z][A-Za-z0-9 &-]{2,40})", text)
return match.group(1).strip() if match else "Operations"
def _infer_status(text: str) -> str:
return _normalize_status("", text)
def _exception_reason(update: str, status: str) -> str:
lower = update.lower()
for keyword in ("blocker", "blocked", "overdue", "escalated", "risk", "delayed", "waiting"):
if keyword in lower:
return keyword
return "red status" if status == "red" else "attention needed"
def _sha256_json(payload: dict[str, Any]) -> str:
return hashlib.sha256(json.dumps(payload, sort_keys=True).encode("utf-8")).hexdigest()
def _now_iso() -> str:
return datetime.now(UTC).replace(microsecond=0).isoformat()