a2a-source-edit: write agent.py

This commit is contained in:
a2a-cloud
2026-07-13 03:01:58 +00:00
parent b63e45d185
commit c937c99f65

936
agent.py
View File

@@ -1,189 +1,839 @@
"""agent-compliance-auditor agent. """Read-only A2A agent compliance auditor.
Starter stack: This agent is deterministic by design: it does not call an LLM, does not execute
- DeepAgents for tool-calling orchestration arbitrary commands, does not mutate production systems, and treats all audited
- Caller-provided LLM credentials via ctx.llm content as untrusted evidence. External integrations are optional and
- A tiny model-call middleware hook you can replace with tracing, caller-configured; when they are absent the corresponding evidence is marked as
routing, rate limits, or policy checks ``unknown`` / ``not_tested`` rather than guessed.
""" """
from __future__ import annotations from __future__ import annotations
import asyncio
import hashlib
import json import json
from pathlib import Path import re
from typing import Any import time
from datetime import UTC, datetime
from enum import Enum
from fnmatch import fnmatch
from typing import Any, Literal
from urllib.parse import urlparse
from pydantic import BaseModel import httpx
from pydantic import AnyHttpUrl, BaseModel, ConfigDict, Field, field_validator, model_validator
import a2a_pack as a2a import a2a_pack as a2a
from a2a_pack import ( from a2a_pack import (
A2AAgent, A2AAgent,
LLMProvisioning, ConsumerSetup,
{{ auth_type }}, ConsumerSetupField,
EgressPolicy,
NoAuth,
Pricing, Pricing,
Resources,
RunContext, RunContext,
WorkspaceAccess, WorkspaceAccess,
WorkspaceMode, WorkspaceMode,
) )
from a2a_pack.context import LLMCreds from a2a_pack.context import AgentEvent
from a2a_pack.workspace import FileType
class AgentComplianceAuditorConfig(BaseModel): MAX_TEXT_BYTES = 256_000
pass MAX_EVIDENCE_BYTES = 8_000
MAX_FILES_PER_AUDIT = 80
MAX_FINDINGS = 200
HTTP_TIMEOUT_SECONDS = 8.0
ALLOWED_REPO_GLOBS = (
"agent.py",
"a2a.yaml",
"requirements.txt",
"README.md",
"pyproject.toml",
"Dockerfile",
"*.py",
"*.yaml",
"*.yml",
"*.json",
"*.md",
"tests/*.py",
"tests/**/*.py",
"skills/**/SKILL.md",
)
DENIED_REPO_GLOBS = (
".git/**",
"**/.env",
"**/.env.*",
"**/*secret*",
"**/*credential*",
"**/id_rsa",
"**/*.pem",
"**/*.key",
)
SECRET_PATTERNS = (
re.compile(r"(?i)(api[_-]?key|token|secret|password)\s*[:=]\s*['\"]?([A-Za-z0-9_./+=-]{12,})"),
re.compile(r"sk-[A-Za-z0-9]{20,}"),
re.compile(r"(?i)bearer\s+[A-Za-z0-9_./+=-]{16,}"),
)
WRITE_INDICATORS = (
"kubectl apply",
"kubectl delete",
"kubectl patch",
"kubectl scale",
"helm upgrade",
"git push",
"write_file(",
"ctx.secret(",
"OPENAI_API_KEY",
"A2A_LITELLM_KEY",
)
SYSTEM_PROMPT = """\ class StrictModel(BaseModel):
You are a compact tool-calling agent. model_config = ConfigDict(extra="forbid")
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 = "agent-compliance-auditor/.deepagents/skills/"
DEEPAGENTS_RECURSION_LIMIT = 500
class AgentComplianceAuditor(A2AAgent[AgentComplianceAuditorConfig, {{ auth_type }}]): class Severity(str, Enum):
CRITICAL = "critical"
HIGH = "high"
MEDIUM = "medium"
LOW = "low"
INFO = "info"
class Confidence(str, Enum):
HIGH = "high"
MEDIUM = "medium"
LOW = "low"
class ProbeMode(str, Enum):
NONE = "none"
HEALTH_ONLY = "health_only"
SYNTHETIC = "synthetic"
class AuditScope(str, Enum):
CARD = "card"
CONFIGURATION = "configuration"
RUNTIME = "runtime"
REPOSITORY = "repository"
RISK = "risk"
REMEDIATION = "remediation"
class Evidence(StrictModel):
id: str
source: str
path: str | None = None
status: Literal["observed", "missing", "unknown", "not_tested", "redacted"]
excerpt: str = Field(default="", max_length=MAX_EVIDENCE_BYTES)
sha256: str | None = None
class Finding(StrictModel):
id: str
control_id: str
title: str
severity: Severity
confidence: Confidence
status: Literal["open", "passed", "unknown", "not_tested"] = "open"
evidence_refs: list[str] = Field(default_factory=list, max_length=20)
impact: str
remediation: str
verification_procedure: str
class OutputFile(StrictModel):
path: str
artifact_uri: str | None = None
mime_type: str
size_bytes: int = Field(ge=0)
class AuditTarget(StrictModel):
agent_name: str = Field(min_length=1, max_length=128, pattern=r"^[a-z0-9][a-z0-9-]{0,127}$")
agent_url: AnyHttpUrl | None = None
repository: str | None = Field(default=None, max_length=256)
namespace: str | None = Field(default=None, max_length=128)
deployment: str | None = Field(default=None, max_length=128)
@field_validator("repository", "namespace", "deployment")
@classmethod
def _clean_optional(cls, value: str | None) -> str | None:
if value is None:
return None
clean = value.strip()
return clean or None
class AuditOptions(StrictModel):
audit_id: str | None = Field(default=None, max_length=80, pattern=r"^[A-Za-z0-9_.-]+$")
max_files: int = Field(default=40, ge=0, le=MAX_FILES_PER_AUDIT)
probe_mode: ProbeMode = ProbeMode.HEALTH_ONLY
include_repository: bool = True
include_runtime: bool = True
include_remediation: bool = True
allowed_hosts: list[str] = Field(default_factory=list, max_length=20)
allowed_namespaces: list[str] = Field(default_factory=list, max_length=20)
class InventoryRequest(StrictModel):
target: AuditTarget
options: AuditOptions = Field(default_factory=AuditOptions)
class AuditConfigurationRequest(StrictModel):
target: AuditTarget
options: AuditOptions = Field(default_factory=AuditOptions)
inventory: dict[str, Any] | None = None
class AuditRuntimeRequest(StrictModel):
target: AuditTarget
options: AuditOptions = Field(default_factory=AuditOptions)
inventory: dict[str, Any] | None = None
class AssessRiskRequest(StrictModel):
target: AuditTarget
options: AuditOptions = Field(default_factory=AuditOptions)
inventory: dict[str, Any] | None = None
configuration_findings: list[Finding] = Field(default_factory=list, max_length=MAX_FINDINGS)
runtime_findings: list[Finding] = Field(default_factory=list, max_length=MAX_FINDINGS)
class GenerateRemediationRequest(StrictModel):
target: AuditTarget
options: AuditOptions = Field(default_factory=AuditOptions)
findings: list[Finding] = Field(default_factory=list, max_length=MAX_FINDINGS)
risk_summary: dict[str, Any] | None = None
class InventoryResult(StrictModel):
audit_id: str
target: AuditTarget
status: Literal["ok", "partial", "unreachable", "setup_required"]
output_dir: str
generated_files: list[OutputFile] = Field(default_factory=list)
inventory: dict[str, Any]
evidence: list[Evidence] = Field(default_factory=list)
warnings: list[str] = Field(default_factory=list)
mutations_performed: bool = False
class FindingResult(StrictModel):
audit_id: str
target: AuditTarget
status: Literal["ok", "partial", "setup_required"]
output_dir: str
generated_files: list[OutputFile] = Field(default_factory=list)
findings: list[Finding]
evidence: list[Evidence] = Field(default_factory=list)
warnings: list[str] = Field(default_factory=list)
mutations_performed: bool = False
class RiskSummaryResult(StrictModel):
audit_id: str
target: AuditTarget
status: Literal["ok", "partial"]
output_dir: str
generated_files: list[OutputFile] = Field(default_factory=list)
risk_summary: dict[str, Any]
findings: list[Finding]
warnings: list[str] = Field(default_factory=list)
mutations_performed: bool = False
class RemediationResult(StrictModel):
audit_id: str
target: AuditTarget
status: Literal["ok", "partial"]
output_dir: str
generated_files: list[OutputFile] = Field(default_factory=list)
remediation_markdown: str
sarif: dict[str, Any]
warnings: list[str] = Field(default_factory=list)
mutations_performed: bool = False
class AgentComplianceAuditorConfig(StrictModel):
default_allowed_hosts: list[str] = Field(default_factory=list, max_length=50)
default_allowed_namespaces: list[str] = Field(default_factory=list, max_length=50)
class AgentComplianceAuditor(A2AAgent[AgentComplianceAuditorConfig, NoAuth]):
name = "agent-compliance-auditor" name = "agent-compliance-auditor"
description = "Read-only compliance auditor for deployed A2A agents and repositories." description = (
"Read-only compliance auditor for deployed A2A agents and repositories. "
"Produces bounded evidence, stable findings, risk summaries, and remediation artifacts."
)
version = "0.1.0" version = "0.1.0"
config_model = AgentComplianceAuditorConfig config_model = AgentComplianceAuditorConfig
auth_model = {{ auth_type }} auth_model = NoAuth
# 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.
llm_provisioning = LLMProvisioning.PLATFORM
pricing = Pricing( pricing = Pricing(
price_per_call_usd=0.0, price_per_call_usd=0.0,
caller_pays_llm=True, caller_pays_llm=False,
notes="Starter agent uses the caller's saved LLM credential via ctx.llm.", notes="Deterministic read-only auditing; no LLM calls and no production mutations.",
) )
resources = Resources(cpu="500m", memory="512Mi", max_runtime_seconds=900)
workspace_access = WorkspaceAccess.dynamic( workspace_access = WorkspaceAccess.dynamic(
max_files=64, max_files=MAX_FILES_PER_AUDIT + 10,
allowed_modes=(WorkspaceMode.READ_ONLY, WorkspaceMode.READ_WRITE_OVERLAY), allowed_modes=(WorkspaceMode.READ_ONLY, WorkspaceMode.READ_WRITE_OVERLAY),
require_reason=False, require_reason=False,
deny_patterns=DENIED_REPO_GLOBS,
max_total_size_bytes=12 * 1024 * 1024,
) )
tools_used = ("deepagents", "langchain") egress = EgressPolicy(
allow_hosts=(
@a2a.tool(description="Ask the starter DeepAgent to answer with tool calls when useful") "gitea.a2a.svc.cluster.local",
async def ask(self, ctx: RunContext[{{ auth_type }}], prompt: str) -> str: "argocd-server.argocd.svc.cluster.local",
creds = ctx.llm "kubernetes.default.svc",
await ctx.emit_progress(f"llm: {creds.model} via {creds.source}") "registry-1.docker.io",
if not creds.api_key: "ghcr.io",
return ( ),
"LLM key required. Add an LLM credential in Settings > LLM " deny_internet_by_default=True,
"credentials before running this agent; for local --invoke "
"runs set AGENT_LLM_KEY."
) )
graph = self._build_deep_agent(ctx=ctx, creds=creds) tools_used = ("httpx", "pydantic", "workspace")
state = await graph.ainvoke( consumer_setup = ConsumerSetup.from_fields(
{"messages": [{"role": "user", "content": prompt}]}, ConsumerSetupField.config("GITEA_BASE_URL", label="Gitea base URL", required=False, input_type="url"),
config={"recursion_limit": DEEPAGENTS_RECURSION_LIMIT}, ConsumerSetupField.secret("GITEA_TOKEN", label="Gitea read-only token", required=False),
ConsumerSetupField.config("KUBERNETES_BASE_URL", label="Kubernetes API base URL", required=False, input_type="url"),
ConsumerSetupField.secret("KUBERNETES_TOKEN", label="Kubernetes read-only bearer token", required=False),
ConsumerSetupField.config("ARGOCD_BASE_URL", label="Argo CD base URL", required=False, input_type="url"),
ConsumerSetupField.secret("ARGOCD_TOKEN", label="Argo CD read-only token", required=False),
ConsumerSetupField.config("REGISTRY_BASE_URL", label="Registry metadata base URL", required=False, input_type="url"),
ConsumerSetupField.secret("REGISTRY_TOKEN", label="Registry read-only token", required=False),
) )
await ctx.emit_progress("deepagent finished")
return _last_message_text(state)
def _build_deep_agent( @a2a.tool(description="Inventory a deployed A2A agent card, repository evidence, manifests, exposure, provenance, resources, and health", timeout_seconds=180, idempotent=True, cost_class="read-only")
self, async def inventory_agent(self, ctx: RunContext[NoAuth], request: InventoryRequest) -> InventoryResult:
*, audit_id = _audit_id(request.target, request.options.audit_id)
ctx: RunContext[{{ auth_type }}], await _event(ctx, "audit_started", audit_id=audit_id, skill="inventory_agent", target=request.target.agent_name)
creds: LLMCreds, inventory, evidence, warnings, status = await _collect_inventory(self, ctx, request.target, request.options, audit_id)
) -> Any: output_dir = _output_dir(audit_id)
# Lazy imports keep `a2a card` usable before local dependencies are files = await _write_json_bundle(ctx, output_dir, {"inventory.json": inventory}, warnings)
# installed. `a2a deploy` installs requirements.txt during the build. await _event(ctx, "audit_completed", audit_id=audit_id, skill="inventory_agent", findings=0)
from a2a_pack.deepagents import create_a2a_deep_agent return InventoryResult(audit_id=audit_id, target=request.target, status=status, output_dir=output_dir, generated_files=files, inventory=inventory, evidence=evidence, warnings=warnings)
from langchain.agents.middleware import wrap_model_call
from langchain_core.tools import tool
@tool @a2a.tool(description="Audit static configuration and repository files against A2A security, privacy, and platform controls", timeout_seconds=180, idempotent=True, cost_class="read-only")
def text_stats(text: str) -> str: async def audit_configuration(self, ctx: RunContext[NoAuth], request: AuditConfigurationRequest) -> FindingResult:
"""Return exact word, character, and line counts for text.""" audit_id = _audit_id(request.target, request.options.audit_id)
words = [part for part in text.split() if part.strip()] await _event(ctx, "audit_started", audit_id=audit_id, skill="audit_configuration", target=request.target.agent_name)
return json.dumps( inventory = request.inventory or (await _collect_inventory(self, ctx, request.target, request.options, audit_id))[0]
{ findings, evidence, warnings = _configuration_findings(request.target, inventory)
"characters": len(text), findings = _dedupe_findings(findings)
"words": len(words), output_dir = _output_dir(audit_id)
"lines": len(text.splitlines()) or 1, files = await _write_json_bundle(ctx, output_dir, {"findings.json": {"findings": [_dump(f) for f in findings]}}, warnings)
await _event(ctx, "audit_completed", audit_id=audit_id, skill="audit_configuration", findings=len(findings))
return FindingResult(audit_id=audit_id, target=request.target, status="ok", output_dir=output_dir, generated_files=files, findings=findings, evidence=evidence, warnings=warnings)
@a2a.tool(description="Run bounded non-destructive runtime probes and audit deployment health without mutations", timeout_seconds=180, idempotent=True, cost_class="read-only")
async def audit_runtime(self, ctx: RunContext[NoAuth], request: AuditRuntimeRequest) -> FindingResult:
audit_id = _audit_id(request.target, request.options.audit_id)
await _event(ctx, "audit_started", audit_id=audit_id, skill="audit_runtime", target=request.target.agent_name)
inventory = request.inventory or (await _collect_inventory(self, ctx, request.target, request.options, audit_id))[0]
findings, evidence, warnings = await _runtime_findings(self, ctx, request.target, request.options, inventory)
findings = _dedupe_findings(findings)
output_dir = _output_dir(audit_id)
files = await _write_json_bundle(ctx, output_dir, {"findings.json": {"findings": [_dump(f) for f in findings]}}, warnings)
await _event(ctx, "audit_completed", audit_id=audit_id, skill="audit_runtime", findings=len(findings))
return FindingResult(audit_id=audit_id, target=request.target, status="ok", output_dir=output_dir, generated_files=files, findings=findings, evidence=evidence, warnings=warnings)
@a2a.tool(description="Deduplicate findings, map evidence to controls, and score severity, confidence, and aggregate risk", timeout_seconds=120, idempotent=True, cost_class="read-only")
async def assess_risk(self, ctx: RunContext[NoAuth], request: AssessRiskRequest) -> RiskSummaryResult:
audit_id = _audit_id(request.target, request.options.audit_id)
await _event(ctx, "audit_started", audit_id=audit_id, skill="assess_risk", target=request.target.agent_name)
all_findings = _dedupe_findings([*request.configuration_findings, *request.runtime_findings])
summary = _risk_summary(all_findings)
output_dir = _output_dir(audit_id)
files = await _write_json_bundle(ctx, output_dir, {"risk-summary.json": summary, "findings.json": {"findings": [_dump(f) for f in all_findings]}}, [])
await _event(ctx, "audit_completed", audit_id=audit_id, skill="assess_risk", findings=len(all_findings), score=summary["risk_score"])
return RiskSummaryResult(audit_id=audit_id, target=request.target, status="ok", output_dir=output_dir, generated_files=files, risk_summary=summary, findings=all_findings)
@a2a.tool(description="Generate actionable remediation and a machine-readable SARIF report for compliance findings", timeout_seconds=120, idempotent=True, cost_class="read-only")
async def generate_remediation(self, ctx: RunContext[NoAuth], request: GenerateRemediationRequest) -> RemediationResult:
audit_id = _audit_id(request.target, request.options.audit_id)
await _event(ctx, "audit_started", audit_id=audit_id, skill="generate_remediation", target=request.target.agent_name)
findings = _dedupe_findings(request.findings)
risk_summary = request.risk_summary or _risk_summary(findings)
markdown = _remediation_markdown(request.target, audit_id, findings, risk_summary)
sarif = _sarif(request.target, findings)
output_dir = _output_dir(audit_id)
files = await _write_json_bundle(ctx, output_dir, {"remediation.md": markdown, "sarif.json": sarif, "risk-summary.json": risk_summary}, [])
await _event(ctx, "audit_completed", audit_id=audit_id, skill="generate_remediation", findings=len(findings))
return RemediationResult(audit_id=audit_id, target=request.target, status="ok", output_dir=output_dir, generated_files=files, remediation_markdown=markdown, sarif=sarif)
async def _collect_inventory(agent: AgentComplianceAuditor, ctx: RunContext[NoAuth], target: AuditTarget, options: AuditOptions, audit_id: str) -> tuple[dict[str, Any], list[Evidence], list[str], Literal["ok", "partial", "unreachable", "setup_required"]]:
warnings: list[str] = []
evidence: list[Evidence] = []
status: Literal["ok", "partial", "unreachable", "setup_required"] = "ok"
allowed_hosts = _allowed_hosts(agent, options)
allowed_namespaces = _allowed_namespaces(agent, options)
await _event(ctx, "audit_progress", audit_id=audit_id, phase="resolve_identity")
identity = {
"agent_name": target.agent_name,
"agent_url": str(target.agent_url) if target.agent_url else None,
"repository": target.repository,
"namespace": target.namespace,
"deployment": target.deployment,
"authorized": _namespace_allowed(target.namespace, allowed_namespaces),
} }
if not identity["authorized"]:
status = "partial"
warnings.append("Target namespace is outside the configured allowlist; runtime probes skipped.")
await _event(ctx, "audit_progress", audit_id=audit_id, phase="agent_card")
card = await _fetch_agent_card(ctx, str(target.agent_url) if target.agent_url else None, allowed_hosts, warnings)
if card.get("status") == "unreachable":
status = "unreachable" if status == "ok" else status
evidence.append(_evidence("card", "agent_card", "observed" if card.get("card") else "unknown", card))
await _event(ctx, "audit_progress", audit_id=audit_id, phase="repository_static_inspection")
repo_files = await _read_workspace_repo(ctx, options.max_files if options.include_repository else 0, warnings)
repo_summary = _summarize_repo(repo_files)
evidence.extend(_repo_evidence(repo_files))
manifests = _extract_manifests(repo_files)
skills = _extract_skills(card.get("card"), repo_files)
secrets = _detect_secret_references(repo_files)
exposure = _extract_exposure(card.get("card"), manifests)
resources = _extract_resources(card.get("card"), manifests)
provenance = _extract_provenance(manifests, repo_files)
rbac = _extract_rbac(manifests)
health = {"status": "not_tested", "reason": "probe_mode is none"}
if options.probe_mode is not ProbeMode.NONE and identity["authorized"]:
await _event(ctx, "audit_progress", audit_id=audit_id, phase="bounded_runtime_probe")
health = await _health_probe(str(target.agent_url) if target.agent_url else None, allowed_hosts, warnings)
inventory = {
"schema_version": "2026-07-13",
"audit_id": audit_id,
"generated_at": datetime.now(UTC).isoformat(),
"identity": identity,
"agent_card": card,
"skills": skills,
"schemas": _extract_schemas(card.get("card")),
"manifests": manifests,
"rbac": rbac,
"secrets_references": secrets,
"network_exposure": exposure,
"image_provenance": provenance,
"resource_limits": resources,
"deployment_health": health,
"repository": repo_summary,
"controls": _controls_catalog(),
"unknowns": _unknowns(card, manifests, health, repo_files),
"mutations_performed": False,
}
return inventory, evidence, warnings, status
async def _fetch_agent_card(ctx: RunContext[NoAuth], agent_url: str | None, allowed_hosts: set[str], warnings: list[str]) -> dict[str, Any]:
if not agent_url:
return {"status": "unknown", "reason": "agent_url not provided"}
parsed = urlparse(agent_url)
if not _safe_http_url(parsed, allowed_hosts):
warnings.append("Agent URL host is not allowlisted or is unsafe; card fetch skipped.")
return {"status": "not_tested", "reason": "host_not_allowlisted"}
url = agent_url.rstrip("/") + "/.well-known/agent-card"
try:
async with httpx.AsyncClient(timeout=httpx.Timeout(HTTP_TIMEOUT_SECONDS), follow_redirects=False) as client:
resp = await client.get(url, headers={"accept": "application/json"})
if resp.status_code >= 400:
return {"status": "unreachable", "http_status": resp.status_code, "body_excerpt": _redact(resp.text[:500])}
data = _bounded_json(resp.text, warnings)
return {"status": "ok", "url": url, "card": data}
except Exception as exc: # noqa: BLE001
warnings.append(f"Agent card fetch failed: {type(exc).__name__}")
return {"status": "unreachable", "error_type": type(exc).__name__}
async def _health_probe(agent_url: str | None, allowed_hosts: set[str], warnings: list[str]) -> dict[str, Any]:
if not agent_url:
return {"status": "unknown", "reason": "agent_url not provided"}
parsed = urlparse(agent_url)
if not _safe_http_url(parsed, allowed_hosts):
return {"status": "not_tested", "reason": "host_not_allowlisted"}
url = agent_url.rstrip("/") + "/healthz"
started = time.perf_counter()
try:
async with httpx.AsyncClient(timeout=httpx.Timeout(HTTP_TIMEOUT_SECONDS), follow_redirects=False) as client:
resp = await client.get(url)
return {"status": "ok" if resp.status_code < 500 else "degraded", "http_status": resp.status_code, "latency_ms": round((time.perf_counter() - started) * 1000, 2)}
except Exception as exc: # noqa: BLE001
warnings.append(f"Health probe failed: {type(exc).__name__}")
return {"status": "unreachable", "error_type": type(exc).__name__}
async def _read_workspace_repo(ctx: RunContext[NoAuth], max_files: int, warnings: list[str]) -> dict[str, str]:
if max_files <= 0:
return {}
try:
view = await ctx.workspace.open_view(
purpose="Read-only static inspection of A2A agent source files",
hints=("agent.py a2a.yaml requirements README tests skills",),
file_types=(FileType.PYTHON, FileType.YAML, FileType.JSON, FileType.MARKDOWN, FileType.TOML),
max_files=max_files,
mode=WorkspaceMode.READ_ONLY,
reason="Compliance audit reads bounded source files and never writes repository paths.",
) )
except Exception as exc: # noqa: BLE001
@wrap_model_call warnings.append(f"Repository workspace unavailable: {type(exc).__name__}")
async def log_model_call(request: Any, handler: Any) -> Any: return {}
messages = request.state.get("messages", []) files: dict[str, str] = {}
print( for match in view.files[:max_files]:
"[middleware] model_call " path = match.path.replace("\\", "/")
f"model={creds.model} source={creds.source} messages={len(messages)}" if not _repo_path_allowed(path):
) continue
return await handler(request) try:
raw = await view.read(path)
backend = ctx.workspace_backend() except Exception as exc: # noqa: BLE001
skill_sources = _seed_runtime_skills(backend, ctx) warnings.append(f"Could not read {path}: {type(exc).__name__}")
# create_a2a_deep_agent resolves provider:model strings with continue
# langchain.init_chat_model from ctx.llm, preserving LiteLLM routing, files[path] = raw[:MAX_TEXT_BYTES].decode("utf-8", errors="replace")
# provider-specific extra body, and runtime model overrides. return files
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,
)
def _runtime_skills_root(ctx: RunContext[Any]) -> str: async def _write_json_bundle(ctx: RunContext[NoAuth], output_dir: str, files: dict[str, Any], warnings: list[str]) -> list[OutputFile]:
workspace = getattr(ctx, "_workspace", None) out: list[OutputFile] = []
prefixes = tuple(getattr(workspace, "write_prefixes", ()) or ()) for name, payload in files.items():
if not prefixes: rel = f"{output_dir}/{name}"
outputs_prefix = getattr(workspace, "outputs_prefix", None) if name.endswith(".md") and isinstance(payload, str):
prefixes = (outputs_prefix or "outputs/",) data = payload.encode("utf-8")
prefix = str(prefixes[0]).strip("/") mime = "text/markdown"
return f"/{prefix}/{RUNTIME_SKILLS_DIR}" if prefix else f"/{RUNTIME_SKILLS_DIR}" else:
data = json.dumps(payload, indent=2, sort_keys=True, default=str).encode("utf-8")
mime = "application/json"
artifact_uri: str | None = None
try:
ref = await ctx.write_artifact(rel, data, mime)
await ctx.emit_artifact(ref)
artifact_uri = ref.uri
except Exception as exc: # noqa: BLE001
warnings.append(f"Artifact write failed for {rel}: {type(exc).__name__}")
try:
ws = ctx.workspace
if getattr(ws, "is_writable_output", lambda _p: False)(rel) and hasattr(ws, "write_bytes"):
ws.write_bytes(rel, data)
except Exception:
pass
out.append(OutputFile(path=rel, artifact_uri=artifact_uri, mime_type=mime, size_bytes=len(data)))
return out
def _seed_runtime_skills(backend: Any, ctx: RunContext[Any]) -> list[str]: def _configuration_findings(target: AuditTarget, inventory: dict[str, Any]) -> tuple[list[Finding], list[Evidence], list[str]]:
"""Copy packaged DeepAgents skills into the invocation workspace. findings: list[Finding] = []
evidence: list[Evidence] = []
warnings: list[str] = []
card = (inventory.get("agent_card") or {}).get("card") or {}
repo = inventory.get("repository") or {}
resources = inventory.get("resource_limits") or {}
secrets = inventory.get("secrets_references") or []
exposure = inventory.get("network_exposure") or {}
skills = inventory.get("skills") or []
DeepAgents loads skills from its backend, while source-controlled if not card:
``skills/`` folders live in the image. This bridge lets generated agents findings.append(_finding(target, "A2A-CARD-001", "Agent Card could not be verified", Severity.HIGH, Confidence.MEDIUM, "card", "Consumers cannot validate skills, schemas, runtime, or billing metadata.", "Ensure /.well-known/agent-card is reachable over an approved host and returns a valid A2A Card.", "Fetch /.well-known/agent-card and validate that all public skills and runtime metadata are present."))
ship reusable SKILL.md bundles without giving up durable A2A workspace if len(skills) == 0:
files. findings.append(_finding(target, "A2A-SCHEMA-001", "No public skills discovered", Severity.HIGH, Confidence.MEDIUM, "skills", "Callers cannot rely on a stable typed API contract.", "Expose typed @a2a.tool methods with Pydantic-compatible input and output schemas.", "Reload the live Agent Card and confirm skills[].input_schema and output_schema are non-null."))
""" for skill in skills:
root = Path(__file__).parent / "skills" schema = skill.get("input_schema") or skill.get("inputSchema") or {}
if not root.exists(): out_schema = skill.get("output_schema") or skill.get("outputSchema") or {}
return [] if not _schema_strict(schema) or not out_schema:
runtime_skills_root = _runtime_skills_root(ctx) findings.append(_finding(target, "A2A-SCHEMA-002", f"Skill {skill.get('name', '<unknown>')} has weak schema", Severity.MEDIUM, Confidence.HIGH, "schemas", "Weak or nullable schemas increase integration errors and prompt-injection attack surface.", "Use strict Pydantic models or explicit JSON Schema with type=object and additionalProperties=false.", "Inspect the Agent Card skill schemas and run negative validation tests for extra fields."))
uploads: list[tuple[str, bytes]] = [] if exposure.get("public") is True and not exposure.get("documented_public_reason"):
for path in root.rglob("*"): findings.append(_finding(target, "A2A-NET-001", "Agent is publicly exposed without documented reason", Severity.MEDIUM, Confidence.MEDIUM, "network_exposure", "Unexpected public exposure can broaden abuse and data-exfiltration paths.", "Set expose.public=false unless public access is required; document the reason and authentication boundary.", "Review a2a.yaml and ingress policy; confirm exposure matches intended audience."))
if path.is_file(): if resources.get("max_runtime_seconds") in (None, "unknown"):
rel = path.relative_to(root).as_posix() findings.append(_finding(target, "A2A-OPS-001", "Runtime limit is unknown", Severity.LOW, Confidence.MEDIUM, "resource_limits", "Unbounded or unknown runtime budgets can cause reliability and cost issues.", "Declare Resources(... max_runtime_seconds=...) in agent.py and mirror runtime.resources in a2a.yaml for larger pods.", "Load Agent Card runtime.resources and deployment manifest to confirm limits."))
uploads.append((runtime_skills_root + rel, path.read_bytes())) if secrets:
if uploads: findings.append(_finding(target, "A2A-SECRET-001", "Potential secret or provider-key reference found in source", Severity.HIGH, Confidence.HIGH, "secrets_references", "Secrets in source or logs can be exfiltrated by callers or dependencies.", "Remove literal secrets, provider-key environment reads, and broad secret access. Use ctx.consumer_secret or declared runtime secrets only.", "Run secret scanning and verify no returned artifacts contain secret values."))
backend.upload_files(uploads) if repo.get("suspicious_instruction_count", 0) > 0:
return [runtime_skills_root] findings.append(_finding(target, "A2A-UNTRUSTED-001", "Repository contains malicious instruction-like text", Severity.MEDIUM, Confidence.HIGH, "repository", "Auditors or LLM-backed tools could be manipulated if repo text is treated as instructions.", "Keep repository content in evidence fields only; never concatenate it into privileged prompts without quoting and policy guards.", "Re-run audit and confirm malicious text appears only as redacted evidence, not as executed instructions."))
return [] evidence.append(_evidence("configuration", "configuration_rules", "observed", {"finding_count": len(findings)}))
return findings[:MAX_FINDINGS], evidence, warnings
def _last_message_text(state: dict[str, Any]) -> str: async def _runtime_findings(agent: AgentComplianceAuditor, ctx: RunContext[NoAuth], target: AuditTarget, options: AuditOptions, inventory: dict[str, Any]) -> tuple[list[Finding], list[Evidence], list[str]]:
messages = state.get("messages") or [] findings: list[Finding] = []
if not messages: evidence: list[Evidence] = []
return json.dumps(state, default=str) warnings: list[str] = []
health = inventory.get("deployment_health") or {"status": "unknown"}
identity = inventory.get("identity") or {}
if identity.get("authorized") is False:
findings.append(_finding(target, "A2A-TENANT-001", "Runtime audit skipped by namespace allowlist", Severity.HIGH, Confidence.HIGH, "identity", "Auditing outside the authorized tenant boundary could violate isolation policy.", "Configure allowed_namespaces for the target tenant or run the audit from the owning tenant.", "Confirm namespace allowlist contains only tenant-owned namespaces and retry."))
if health.get("status") in {"unreachable", "degraded"}:
findings.append(_finding(target, "A2A-HEALTH-001", "Agent health probe failed or degraded", Severity.MEDIUM, Confidence.MEDIUM, "deployment_health", "Consumers may experience failed calls or stale Agent Cards.", "Inspect deployment status, readiness probes, service routing, and recent rollout events.", "Repeat health-only probe and verify HTTP 2xx/3xx or documented readiness semantics."))
if health.get("status") in {"unknown", "not_tested"}:
findings.append(_finding(target, "A2A-HEALTH-002", "Runtime health state is unknown", Severity.LOW, Confidence.MEDIUM, "deployment_health", "Unknown runtime state leaves operational risk unmeasured.", "Provide an allowlisted agent_url and enable health_only probes; configure read-only Kubernetes/Argo metadata if needed.", "Run audit_runtime with probe_mode=health_only and verify deployment_health.status."))
evidence.append(_evidence("runtime", "runtime_rules", "observed", {"finding_count": len(findings), "health": health}))
return findings[:MAX_FINDINGS], evidence, warnings
content = getattr(messages[-1], "content", None)
if isinstance(content, str): def _finding(target: AuditTarget, control_id: str, title: str, severity: Severity, confidence: Confidence, evidence_ref: str, impact: str, remediation: str, verification: str) -> Finding:
return content base = f"{target.agent_name}|{control_id}|{title}"
if isinstance(content, list): suffix = hashlib.sha256(base.encode()).hexdigest()[:10]
parts: list[str] = [] return Finding(id=f"{control_id.lower()}-{suffix}", control_id=control_id, title=title, severity=severity, confidence=confidence, evidence_refs=[evidence_ref], impact=impact, remediation=remediation, verification_procedure=verification)
for item in content:
if isinstance(item, dict):
text = item.get("text") or item.get("content") def _dedupe_findings(findings: list[Finding]) -> list[Finding]:
if text: seen: dict[str, Finding] = {}
parts.append(str(text)) rank = {Severity.CRITICAL: 5, Severity.HIGH: 4, Severity.MEDIUM: 3, Severity.LOW: 2, Severity.INFO: 1}
elif item: for f in findings:
parts.append(str(item)) prev = seen.get(f.id)
return "\n".join(parts) if parts else json.dumps(content, default=str) if prev is None or rank[f.severity] > rank[prev.severity]:
return str(content or messages[-1]) seen[f.id] = f
return sorted(seen.values(), key=lambda f: (-rank[f.severity], f.control_id, f.id))[:MAX_FINDINGS]
def _risk_summary(findings: list[Finding]) -> dict[str, Any]:
weights = {Severity.CRITICAL: 100, Severity.HIGH: 40, Severity.MEDIUM: 15, Severity.LOW: 5, Severity.INFO: 1}
counts = {s.value: 0 for s in Severity}
score = 0
for f in findings:
counts[f.severity.value] += 1
score += weights[f.severity]
score = min(score, 100)
if counts["critical"]:
rating = "critical"
elif score >= 70:
rating = "high"
elif score >= 30:
rating = "medium"
elif score > 0:
rating = "low"
else:
rating = "clean"
return {"schema_version": "2026-07-13", "risk_score": score, "risk_rating": rating, "finding_counts": counts, "top_controls": sorted({f.control_id for f in findings})[:20], "unknown_or_not_tested": [f.control_id for f in findings if f.status in {"unknown", "not_tested"}], "mutations_performed": False}
def _remediation_markdown(target: AuditTarget, audit_id: str, findings: list[Finding], risk_summary: dict[str, Any]) -> str:
lines = [f"# Remediation plan for {target.agent_name}", "", f"Audit ID: `{audit_id}`", f"Risk rating: **{risk_summary.get('risk_rating', 'unknown')}** ({risk_summary.get('risk_score', 0)}/100)", "", "This plan is read-only output. No production systems or repositories were changed.", ""]
if not findings:
lines.extend(["## No open findings", "", "The synthetic audit found no policy violations in the supplied evidence. Unknown/not-tested areas should still be reviewed."])
for f in findings:
lines.extend([f"## {f.severity.value.upper()}{f.title}", "", f"- Finding ID: `{f.id}`", f"- Control: `{f.control_id}`", f"- Confidence: `{f.confidence.value}`", f"- Evidence references: {', '.join(f.evidence_refs) or 'none'}", f"- Impact: {f.impact}", f"- Remediation: {f.remediation}", f"- Verification: {f.verification_procedure}", ""])
return "\n".join(lines)
def _sarif(target: AuditTarget, findings: list[Finding]) -> dict[str, Any]:
rules = []
results = []
for f in findings:
rules.append({"id": f.control_id, "name": f.title, "shortDescription": {"text": f.title}, "help": {"text": f.remediation}})
results.append({"ruleId": f.control_id, "level": _sarif_level(f.severity), "message": {"text": f"{f.title}: {f.impact}"}, "partialFingerprints": {"findingId": f.id}, "properties": {"severity": f.severity.value, "confidence": f.confidence.value, "verification": f.verification_procedure}})
return {"version": "2.1.0", "$schema": "https://json.schemastore.org/sarif-2.1.0.json", "runs": [{"tool": {"driver": {"name": "agent-compliance-auditor", "informationUri": "https://docs.a2acloud.io/", "rules": rules}}, "invocations": [{"executionSuccessful": True, "properties": {"target": target.agent_name, "mutations_performed": False}}], "results": results}]}
def _sarif_level(severity: Severity) -> str:
return {Severity.CRITICAL: "error", Severity.HIGH: "error", Severity.MEDIUM: "warning", Severity.LOW: "note", Severity.INFO: "none"}[severity]
def _summarize_repo(files: dict[str, str]) -> dict[str, Any]:
suspicious = 0
indicators: list[str] = []
for path, text in files.items():
lower = text.lower()
if any(token in lower for token in ("ignore previous instructions", "exfiltrate", "send secrets", "system prompt")):
suspicious += 1
indicators.append(path)
return {"files_inspected": len(files), "paths": sorted(files), "suspicious_instruction_count": suspicious, "suspicious_instruction_paths": indicators[:20], "bounded_bytes_per_file": MAX_TEXT_BYTES}
def _extract_manifests(files: dict[str, str]) -> dict[str, Any]:
manifests: dict[str, Any] = {}
for path in ("a2a.yaml", "agent.py", "requirements.txt", "Dockerfile"):
if path in files:
manifests[path] = {"present": True, "sha256": _sha(files[path]), "excerpt": _redact(files[path][:1200])}
else:
manifests[path] = {"present": False}
return manifests
def _extract_skills(card: Any, files: dict[str, str]) -> list[dict[str, Any]]:
if isinstance(card, dict):
skills = card.get("skills") or []
if isinstance(skills, list):
return [s for s in skills if isinstance(s, dict)]
text = files.get("agent.py", "")
names = re.findall(r"async\s+def\s+([a-zA-Z_][a-zA-Z0-9_]*)\s*\(", text)
return [{"name": n, "source": "agent.py"} for n in names]
def _extract_schemas(card: Any) -> dict[str, Any]:
schemas: dict[str, Any] = {}
if isinstance(card, dict):
for skill in card.get("skills") or []:
if isinstance(skill, dict):
name = str(skill.get("name") or skill.get("id") or "unknown")
schemas[name] = {"input_schema": skill.get("input_schema") or skill.get("inputSchema"), "output_schema": skill.get("output_schema") or skill.get("outputSchema")}
return schemas
def _detect_secret_references(files: dict[str, str]) -> list[dict[str, str]]:
out: list[dict[str, str]] = []
for path, text in files.items():
for pattern in SECRET_PATTERNS:
for match in pattern.finditer(text):
out.append({"path": path, "match": _redact(match.group(0))})
for indicator in WRITE_INDICATORS:
if indicator in text:
out.append({"path": path, "match": _redact(indicator)})
return out[:50]
def _extract_exposure(card: Any, manifests: dict[str, Any]) -> dict[str, Any]:
public = None
text = json.dumps(manifests, default=str).lower()
if "public: true" in text or "public=true" in text:
public = True
if "public: false" in text or "public=false" in text:
public = False
return {"public": public, "documented_public_reason": "public_reason" in text or "exposure_reason" in text, "unknown_when_null": public is None}
def _extract_resources(card: Any, manifests: dict[str, Any]) -> dict[str, Any]:
runtime = card.get("runtime") if isinstance(card, dict) else {}
resources = runtime.get("resources") if isinstance(runtime, dict) else {}
out = dict(resources or {}) if isinstance(resources, dict) else {}
if "max_runtime_seconds" not in out:
text = json.dumps(manifests, default=str)
match = re.search(r"max_runtime_seconds[\s:=]+([0-9]+)", text)
out["max_runtime_seconds"] = int(match.group(1)) if match else "unknown"
return out
def _extract_provenance(manifests: dict[str, Any], files: dict[str, str]) -> dict[str, Any]:
dockerfile = files.get("Dockerfile", "")
digest_pinned = "@sha256:" in dockerfile
return {"dockerfile_present": bool(dockerfile), "digest_pinned_base_image": digest_pinned if dockerfile else "unknown", "source_hashes": {k: v.get("sha256") for k, v in manifests.items() if isinstance(v, dict) and v.get("sha256")}}
def _extract_rbac(manifests: dict[str, Any]) -> dict[str, Any]:
text = json.dumps(manifests, default=str).lower()
return {"service_account_mentions": text.count("serviceaccount"), "cluster_role_mentions": text.count("clusterrole"), "status": "observed" if "role" in text else "unknown"}
def _unknowns(card: dict[str, Any], manifests: dict[str, Any], health: dict[str, Any], files: dict[str, str]) -> list[str]:
out: list[str] = []
if not card.get("card"):
out.append("agent_card")
if not files:
out.append("repository_files")
if health.get("status") in {"unknown", "not_tested"}:
out.append("deployment_health")
if not any(v.get("present") for v in manifests.values() if isinstance(v, dict)):
out.append("source_manifests")
return out
def _controls_catalog() -> list[dict[str, str]]:
return [
{"control_id": "A2A-CARD-001", "title": "Agent Card availability and integrity"},
{"control_id": "A2A-SCHEMA-001", "title": "Public skill inventory"},
{"control_id": "A2A-SCHEMA-002", "title": "Strict non-null schemas"},
{"control_id": "A2A-NET-001", "title": "Network exposure minimization"},
{"control_id": "A2A-OPS-001", "title": "Runtime resource limits"},
{"control_id": "A2A-SECRET-001", "title": "Secret handling and redaction"},
{"control_id": "A2A-UNTRUSTED-001", "title": "Untrusted content isolation"},
{"control_id": "A2A-TENANT-001", "title": "Tenant and namespace isolation"},
{"control_id": "A2A-HEALTH-001", "title": "Deployment health"},
]
def _schema_strict(schema: Any) -> bool:
return isinstance(schema, dict) and schema.get("type") == "object" and schema.get("additionalProperties") is False and bool(schema.get("properties"))
def _repo_evidence(files: dict[str, str]) -> list[Evidence]:
return [_evidence(f"repo:{path}", "repository", "observed", text[:1000], path=path) for path, text in list(files.items())[:MAX_FILES_PER_AUDIT]]
def _evidence(eid: str, source: str, status: Literal["observed", "missing", "unknown", "not_tested", "redacted"], payload: Any, path: str | None = None) -> Evidence:
text = payload if isinstance(payload, str) else json.dumps(payload, sort_keys=True, default=str)
redacted = _redact(text)[:MAX_EVIDENCE_BYTES]
return Evidence(id=eid, source=source, path=path, status=status, excerpt=redacted, sha256=_sha(text))
def _audit_id(target: AuditTarget, supplied: str | None) -> str:
if supplied:
return supplied
digest = hashlib.sha256(f"{target.agent_name}|{target.repository or ''}|{int(time.time())}".encode()).hexdigest()[:12]
return f"audit-{target.agent_name}-{digest}"
def _output_dir(audit_id: str) -> str:
clean = re.sub(r"[^A-Za-z0-9_.-]", "-", audit_id)[:80]
return f"outputs/compliance/{clean}"
def _repo_path_allowed(path: str) -> bool:
clean = path.replace("\\", "/").lstrip("/")
if ".." in clean.split("/"):
return False
if any(fnmatch(clean, pat) for pat in DENIED_REPO_GLOBS):
return False
return any(fnmatch(clean, pat) for pat in ALLOWED_REPO_GLOBS)
def _allowed_hosts(agent: AgentComplianceAuditor, options: AuditOptions) -> set[str]:
return {h.lower().strip() for h in [*agent.config.default_allowed_hosts, *options.allowed_hosts] if h.strip()}
def _allowed_namespaces(agent: AgentComplianceAuditor, options: AuditOptions) -> set[str]:
return {n.strip() for n in [*agent.config.default_allowed_namespaces, *options.allowed_namespaces] if n.strip()}
def _namespace_allowed(namespace: str | None, allowed_namespaces: set[str]) -> bool:
if not allowed_namespaces:
return True
return bool(namespace and namespace in allowed_namespaces)
def _safe_http_url(parsed: Any, allowed_hosts: set[str]) -> bool:
if parsed.scheme not in {"https", "http"}:
return False
host = (parsed.hostname or "").lower()
if not host or host in {"localhost", "127.0.0.1", "0.0.0.0"} or host.startswith("169.254.") or host.endswith(".local"):
return False
return host in allowed_hosts
def _bounded_json(text: str, warnings: list[str]) -> Any:
if len(text.encode("utf-8")) > MAX_TEXT_BYTES:
warnings.append("HTTP JSON response exceeded evidence bound and was truncated before parsing.")
text = text[:MAX_TEXT_BYTES]
try:
return json.loads(text)
except json.JSONDecodeError:
warnings.append("HTTP response was not valid JSON.")
return {"invalid_json_excerpt": _redact(text[:1000])}
def _redact(text: str) -> str:
out = text
for pattern in SECRET_PATTERNS:
out = pattern.sub(lambda m: m.group(0)[:8] + "…REDACTED", out)
return out
def _sha(text: str) -> str:
return hashlib.sha256(text.encode("utf-8", errors="replace")).hexdigest()
def _dump(model_or_value: Any) -> Any:
if isinstance(model_or_value, BaseModel):
return model_or_value.model_dump(mode="json")
return model_or_value
async def _event(ctx: RunContext[NoAuth], kind: str, **payload: Any) -> None:
safe_payload = json.loads(json.dumps(payload, default=str))
await ctx.emit_event(AgentEvent(kind=kind, payload=safe_payload))
if kind in {"audit_started", "audit_progress", "audit_completed"}:
message = safe_payload.get("phase") or safe_payload.get("skill") or kind
await ctx.emit_progress(f"{kind}: {message}")