From a3c851e3b9c481aa85cab9110daff6cc98149247 Mon Sep 17 00:00:00 2001 From: a2a-cloud Date: Sun, 19 Jul 2026 18:57:52 +0000 Subject: [PATCH] a2a-source-edit: write agent.py --- agent.py | 625 +++++++++++++++++++++++++++++++++++++++++++------------ 1 file changed, 487 insertions(+), 138 deletions(-) diff --git a/agent.py b/agent.py index 5077f6c..e3137d6 100644 --- a/agent.py +++ b/agent.py @@ -1,189 +1,538 @@ -"""high-utility-one-page-startups-w-e3c7-3 agent. +"""Property manager one-page document generator. -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 that turns a short property-management +intake into a polished one-page client document, emits a markdown artifact, and +persists an execution receipt when the managed database is available. """ from __future__ import annotations +import base64 +import binascii +import hashlib import json -from pathlib import Path -from typing import Any +import os +import re +import uuid +from datetime import datetime, timezone +from typing import Annotated, Any, Literal -from pydantic import BaseModel +from pydantic import BaseModel, Field, field_validator 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 +from a2a_pack.workspace import FileUpload, UploadedFile + +OUTPUT_FILENAME = "high-utility-one-page-startups-w-e3c7-3-output.md" +MAX_TEXT_CHARS = 5_000 +MAX_UPLOAD_BYTES = 256_000 +MAX_BROWSER_UPLOADS = 2 +ALLOWED_MEDIA_TYPES = { + "text/plain", + "text/markdown", + "application/json", + "text/csv", +} class HighUtilityOnePageStartupsWE3c73Config(BaseModel): - pass + default_city: str = "Austin" -SYSTEM_PROMPT = """\ -You are a compact tool-calling agent. +class BrowserDocument(BaseModel): + filename: str = Field(..., min_length=1, max_length=120) + media_type: str = Field(..., min_length=3, max_length=80) + data_base64: str = Field(..., min_length=1, max_length=((MAX_UPLOAD_BYTES * 4) // 3) + 16) -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. -""" + @field_validator("filename") + @classmethod + def clean_filename(cls, value: str) -> str: + cleaned = value.replace("\\", "/").split("/")[-1].strip() + if not cleaned or cleaned in {".", ".."}: + raise ValueError("filename is required") + if any(ch in cleaned for ch in "\x00\r\n"): + raise ValueError("filename contains invalid characters") + return cleaned[:120] -RUNTIME_SKILLS_DIR = "high-utility-one-page-startups-w-e3c7-3/.deepagents/skills/" -DEEPAGENTS_RECURSION_LIMIT = 500 + @field_validator("media_type") + @classmethod + def allowed_media_type(cls, value: str) -> str: + cleaned = value.strip().lower() + if cleaned not in ALLOWED_MEDIA_TYPES: + raise ValueError( + "unsupported media type; use text/plain, text/markdown, application/json, or text/csv" + ) + return cleaned -class HighUtilityOnePageStartupsWE3c73(A2AAgent[HighUtilityOnePageStartupsWE3c73Config, {{ auth_type }}]): +class DocumentRequest(BaseModel): + client_name: str = Field(..., min_length=1, max_length=120) + property_name: str = Field(..., min_length=1, max_length=160) + property_type: Literal["multifamily", "single-family", "commercial", "mixed-use", "hoa"] = "multifamily" + city: str = Field(..., min_length=1, max_length=80) + goal: Literal["owner update", "leasing plan", "maintenance brief", "management proposal"] = "owner update" + intake: str = Field(..., min_length=20, max_length=MAX_TEXT_CHARS) + tone: Literal["professional", "warm", "executive", "urgent"] = "professional" + include_next_steps: bool = True + uploaded_documents: list[BrowserDocument] = Field(default_factory=list, max_length=MAX_BROWSER_UPLOADS) + + +class UploadedDocumentRequest(BaseModel): + client_name: str = Field(..., min_length=1, max_length=120) + property_name: str = Field(..., min_length=1, max_length=160) + city: str = Field(..., min_length=1, max_length=80) + goal: Literal["owner update", "leasing plan", "maintenance brief", "management proposal"] = "owner update" + intake: str = Field(..., min_length=20, max_length=MAX_TEXT_CHARS) + + +class DocumentDescriptor(BaseModel): + name: str + filename: str + media_type: str + size_bytes: int + artifact_uri: str + content_base64: str + + +class DocumentResult(BaseModel): + status: Literal["ok", "validation_error", "persistence_warning"] + document_preview: str + document: DocumentDescriptor | None + receipt_id: str + persisted_receipt: bool + warnings: list[str] = Field(default_factory=list) + highlights: list[str] = Field(default_factory=list) + + +class UploadValidationResult(BaseModel): + status: Literal["ok", "validation_error"] + extracted_text: str = "" + filename: str = "" + media_type: str = "" + size_bytes: int = 0 + message: str = "" + + +class HighUtilityOnePageStartupsWE3c73(A2AAgent[HighUtilityOnePageStartupsWE3c73Config, PlatformUserAuth]): name = "high-utility-one-page-startups-w-e3c7-3" - description = "One-page document generator for property managers that turns a short intake into a polished client document with preview and artifact download." + description = ( + "One-page document generator for property managers that turns a short intake " + "into a polished client document with preview and markdown artifact download." + ) version = "0.1.0" config_model = HighUtilityOnePageStartupsWE3c73Config - 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 trial/BYOK launch contract. The implementation is + # deterministic and does not read ctx.llm or provider secrets. 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="Includes 3 platform-funded skill calls per account, then requires BYOK per launch contract.", ) + resources = Resources(cpu="500m", memory="512Mi", max_runtime_seconds=120) workspace_access = WorkspaceAccess.dynamic( - max_files=64, + max_files=16, allowed_modes=(WorkspaceMode.READ_ONLY, WorkspaceMode.READ_WRITE_OVERLAY), require_reason=False, + max_total_size_bytes=2 * MAX_UPLOAD_BYTES, ) - 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}, + tools_used = ("artifacts", "mcp", "postgres") + platform_resources = AgentPlatformResources( + databases=( + AgentDatabase( + name="high-utility-one-page-startups-w-e3c7-3-data", + 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) + ) - def _build_deep_agent( + @a2a.tool( + name="generate_document", + description="Generate a polished one-page property-management client document, preview, receipt, and markdown artifact.", + timeout_seconds=120, + idempotent=False, + cost_class="standard", + grant_mode="read_write_overlay", + grant_allow_patterns=("outputs/documents/**",), + grant_outputs_prefix="outputs/documents/", + grant_write_prefixes=("outputs/documents/",), + ) + async def generate_document( 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], + request: DocumentRequest, + ) -> DocumentResult: + tenant = _tenant_key(ctx) + upload_texts: list[str] = [] + warnings: list[str] = [] + for item in request.uploaded_documents: + parsed = _decode_browser_document(item) + if parsed.status != "ok": + return DocumentResult( + status="validation_error", + document_preview="", + document=None, + receipt_id=_receipt_id(tenant, request.model_dump(mode="json")), + persisted_receipt=False, + warnings=[parsed.message], + ) + upload_texts.append(f"Attachment {parsed.filename}: {parsed.extracted_text}") - @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, - } - ) - - @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)}" - ) - return await handler(request) - - 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, + document_text, highlights = _build_document(request, upload_texts) + data = document_text.encode("utf-8") + ref = await ctx.write_artifact(OUTPUT_FILENAME, data, "text/markdown") + await ctx.emit_artifact(ref) + receipt_id = _receipt_id(tenant, {**request.model_dump(mode="json"), "artifact": ref.uri}) + persisted = await _persist_receipt( + tenant_key=tenant, + receipt_id=receipt_id, + skill_name="generate_document", + inputs=request.model_dump(mode="json"), + result={ + "filename": OUTPUT_FILENAME, + "artifact_uri": ref.uri, + "size_bytes": len(data), + "preview_hash": hashlib.sha256(data).hexdigest(), + }, + ) + if not persisted: + warnings.append("Receipt persistence is unavailable in this local/runtime context.") + return DocumentResult( + status="ok" if persisted else "persistence_warning", + document_preview=document_text, + document=DocumentDescriptor( + name="document", + filename=OUTPUT_FILENAME, + media_type="text/markdown", + size_bytes=len(data), + artifact_uri=ref.uri, + content_base64=base64.b64encode(data).decode("ascii"), + ), + receipt_id=receipt_id, + persisted_receipt=persisted, + warnings=warnings, + highlights=highlights, ) + @a2a.tool( + name="validate_browser_upload", + description="Validate a bounded browser base64 upload and extract safe text for the document generator.", + timeout_seconds=30, + idempotent=True, + ) + async def validate_browser_upload( + self, + ctx: RunContext[PlatformUserAuth], + document: BrowserDocument, + ) -> UploadValidationResult: + _tenant_key(ctx) # proves platform auth and tenant boundary for MCP/API calls + return _decode_browser_document(document) -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}" + @a2a.tool( + name="generate_document_from_upload", + description="External client upload path: generate a document using a typed FileUpload plus form fields.", + timeout_seconds=120, + idempotent=False, + cost_class="standard", + grant_mode="read_write_overlay", + grant_allow_patterns=("uploads/**", "outputs/documents/**"), + grant_outputs_prefix="outputs/documents/", + grant_write_prefixes=("outputs/documents/",), + ) + async def generate_document_from_upload( + self, + ctx: RunContext[PlatformUserAuth], + uploaded_file: Annotated[ + UploadedFile, + FileUpload( + accept=tuple(sorted(ALLOWED_MEDIA_TYPES)), + max_bytes=MAX_UPLOAD_BYTES, + description="Optional intake notes or source text for the generated document.", + ), + ], + request: UploadedDocumentRequest, + ) -> DocumentResult: + try: + raw = ctx.workspace.read_bytes(uploaded_file.path) + except Exception as exc: # noqa: BLE001 + return DocumentResult( + status="validation_error", + document_preview="", + document=None, + receipt_id="", + persisted_receipt=False, + warnings=[f"Unable to read uploaded file: {type(exc).__name__}"], + ) + browser_doc = BrowserDocument( + filename=uploaded_file.filename, + media_type=uploaded_file.media_type, + data_base64=base64.b64encode(raw[: MAX_UPLOAD_BYTES + 1]).decode("ascii"), + ) + parsed = _decode_browser_document(browser_doc) + if parsed.status != "ok": + return DocumentResult( + status="validation_error", + document_preview="", + document=None, + receipt_id="", + persisted_receipt=False, + warnings=[parsed.message], + ) + full_request = DocumentRequest( + client_name=request.client_name, + property_name=request.property_name, + city=request.city, + goal=request.goal, + intake=request.intake, + uploaded_documents=[browser_doc], + ) + return await self.generate_document(ctx, full_request) -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 [] +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 stable_id is None or str(stable_id).strip() == "": + raise PermissionError("stable platform identity required") + return f"user:{stable_id}" -def _last_message_text(state: dict[str, Any]) -> str: - messages = state.get("messages") or [] - if not messages: - return json.dumps(state, default=str) +def _decode_browser_document(document: BrowserDocument) -> UploadValidationResult: + try: + raw = base64.b64decode(document.data_base64, validate=True) + except (binascii.Error, ValueError): + return UploadValidationResult( + status="validation_error", + filename=document.filename, + media_type=document.media_type, + message="Upload must be valid base64.", + ) + if len(raw) > MAX_UPLOAD_BYTES: + return UploadValidationResult( + status="validation_error", + filename=document.filename, + media_type=document.media_type, + size_bytes=len(raw), + message=f"Upload is too large; maximum is {MAX_UPLOAD_BYTES} bytes.", + ) + text = raw.decode("utf-8", errors="replace") + text = _compact_whitespace(text)[:2_000] + if not text.strip(): + return UploadValidationResult( + status="validation_error", + filename=document.filename, + media_type=document.media_type, + size_bytes=len(raw), + message="Upload did not contain readable UTF-8 text.", + ) + return UploadValidationResult( + status="ok", + extracted_text=text, + filename=document.filename, + media_type=document.media_type, + size_bytes=len(raw), + message="Upload accepted.", + ) - 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 _build_document(request: DocumentRequest, upload_texts: list[str]) -> tuple[str, list[str]]: + intake = _compact_whitespace(request.intake) + sentences = _split_sentences(intake) + bullets = _select_bullets(sentences, upload_texts) + goal_label = request.goal.title() + tone_line = { + "professional": "clear, practical, and client-ready", + "warm": "approachable while staying specific and action-oriented", + "executive": "concise, decision-focused, and outcome-led", + "urgent": "direct, timely, and focused on immediate resolution", + }[request.tone] + next_steps = _next_steps(request.goal, request.property_type) if request.include_next_steps else [] + highlights = bullets[:3] + + lines = [ + f"# {goal_label}: {request.property_name}", + "", + f"**Prepared for:** {request.client_name} ", + f"**Asset:** {request.property_name} ({request.property_type.replace('-', ' ')}) ", + f"**Market:** {request.city} ", + f"**Date:** {datetime.now(timezone.utc).date().isoformat()}", + "", + "## Executive Summary", + f"This one-page brief turns the current intake into a {tone_line} client document for {request.client_name}. " + f"The immediate objective is to support a polished {request.goal} for {request.property_name} while keeping the action plan simple enough to start this week.", + "", + "## Situation Snapshot", + ] + lines.extend(f"- {bullet}" for bullet in bullets) + lines.extend([ + "", + "## Recommended Client Message", + _client_message(request, bullets), + "", + "## Operating Priorities", + ]) + priorities = _priorities(request.goal, request.property_type) + lines.extend(f"- **{title}:** {detail}" for title, detail in priorities) + if next_steps: + lines.extend(["", "## Next Steps"]) + lines.extend(f"{idx}. {step}" for idx, step in enumerate(next_steps, start=1)) + if upload_texts: + lines.extend(["", "## Source Notes Used"]) + for item in upload_texts[:MAX_BROWSER_UPLOADS]: + lines.append(f"- {_truncate(item, 220)}") + lines.extend([ + "", + "---", + "Generated by high-utility-one-page-startups-w-e3c7-3 for fast property-management client communication.", + ]) + return "\n".join(lines) + "\n", highlights + + +def _client_message(request: DocumentRequest, bullets: list[str]) -> str: + opener = { + "owner update": "Here is the concise owner-facing update I recommend sending now:", + "leasing plan": "Here is the leasing-plan narrative I recommend sharing with stakeholders:", + "maintenance brief": "Here is the maintenance brief language I recommend using:", + "management proposal": "Here is the proposal-ready positioning I recommend leading with:", + }[request.goal] + return ( + f"{opener} {request.property_name} is being managed with focus on " + f"{bullets[0].rstrip('.').lower()}. The plan is to keep communication tight, " + f"make the next operational move visible, and give {request.client_name} a clear basis for approval or follow-up." + ) + + +def _priorities(goal: str, property_type: str) -> list[tuple[str, str]]: + common = [ + ("Owner confidence", "Lead with what changed, what is controlled, and what decision is needed."), + ("Resident or tenant experience", "Convert intake details into visible service improvements and clear timelines."), + ] + by_goal = { + "owner update": ("Reporting cadence", "Send a short weekly status note until the current issue or initiative is closed."), + "leasing plan": ("Demand capture", "Refresh pricing, response speed, and showing follow-up before adding new spend."), + "maintenance brief": ("Resolution path", "Separate urgent safety items from cosmetic work and assign owners to each."), + "management proposal": ("Value proof", "Tie management actions to NOI protection, retention, and reduced owner friction."), + } + type_priority = ("Asset fit", f"Tune the plan for a {property_type.replace('-', ' ')} property rather than a generic template.") + return [common[0], by_goal[goal], type_priority, common[1]] + + +def _next_steps(goal: str, property_type: str) -> list[str]: + if goal == "leasing plan": + return [ + "Confirm current vacancy, pricing, concessions, and lead-response time within 24 hours.", + "Publish one refreshed leasing message and track qualified inquiries for seven days.", + "Send the owner a short results note with the next pricing or marketing adjustment.", + ] + if goal == "maintenance brief": + return [ + "Classify open items by safety, revenue risk, and resident impact.", + "Assign each item an owner, budget range, and target completion date.", + "Share a completion photo or vendor note for every closed item.", + ] + if goal == "management proposal": + return [ + f"Lead with the two highest-friction problems typical for this {property_type.replace('-', ' ')} asset.", + "Attach a 30-day transition checklist with owner approvals clearly marked.", + "Close with the first measurable operating win the manager will deliver.", + ] + return [ + "Send the one-page update to the client and ask for one approval or correction.", + "Turn the operating priorities into assigned tasks with dates.", + "Follow up with a short progress note after the first completed action.", + ] + + +def _split_sentences(text: str) -> list[str]: + parts = re.split(r"(?<=[.!?])\s+|\n+", text) + return [_truncate(part.strip(" -•\t"), 180) for part in parts if part.strip()] + + +def _select_bullets(sentences: list[str], upload_texts: list[str]) -> list[str]: + bullets = [] + for sentence in sentences[:5]: + clean = sentence.rstrip(".") + if clean: + bullets.append(clean[0].upper() + clean[1:] + ".") + for text in upload_texts: + if len(bullets) >= 6: + break + bullets.append(_truncate(_compact_whitespace(text), 180).rstrip(".") + ".") + while len(bullets) < 4: + fallback = [ + "Client communication should be concise, specific, and tied to the next operating decision.", + "The manager should turn the intake into a visible action plan rather than a long narrative.", + "Follow-up should include an owner, due date, and evidence of completion.", + "The document should be usable immediately in an email or meeting recap.", + ][len(bullets)] + bullets.append(fallback) + return bullets[:6] + + +def _compact_whitespace(text: str) -> str: + return re.sub(r"\s+", " ", text).strip() + + +def _truncate(text: str, limit: int) -> str: + return text if len(text) <= limit else text[: limit - 1].rstrip() + "…" + + +def _receipt_id(tenant_key: str, payload: dict[str, Any]) -> str: + material = json.dumps(payload, sort_keys=True, separators=(",", ":")) + digest = hashlib.sha256(f"{tenant_key}:{material}:{uuid.uuid4().hex}".encode("utf-8")).hexdigest() + return f"rct_{digest[:24]}" + + +async def _persist_receipt( + *, + tenant_key: str, + receipt_id: str, + skill_name: str, + inputs: dict[str, Any], + result: dict[str, Any], +) -> bool: + database_url = os.environ.get("DATABASE_URL") + if not database_url: + return False + try: + import psycopg + from psycopg.types.json import Jsonb + + options = "-c statement_timeout=5000 -c lock_timeout=3000 -c idle_in_transaction_session_timeout=5000" + with psycopg.connect(database_url, options=options) as conn: + with conn.transaction(): + with conn.cursor() as cur: + cur.execute( + """ + INSERT INTO execution_receipts + (tenant_key, receipt_id, skill_name, input_json, result_json) + VALUES (%s, %s, %s, %s, %s) + ON CONFLICT (tenant_key, receipt_id) DO UPDATE SET + result_json = EXCLUDED.result_json, + updated_at = NOW() + """, + (tenant_key, receipt_id, skill_name, Jsonb(inputs), Jsonb(result)), + ) + return True + except Exception: + return False