diff --git a/agent.py b/agent.py index 22c099d..6589a02 100644 --- a/agent.py +++ b/agent.py @@ -1,189 +1,839 @@ -"""agent-compliance-auditor agent. +"""Read-only A2A agent compliance auditor. -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 +This agent is deterministic by design: it does not call an LLM, does not execute +arbitrary commands, does not mutate production systems, and treats all audited +content as untrusted evidence. External integrations are optional and +caller-configured; when they are absent the corresponding evidence is marked as +``unknown`` / ``not_tested`` rather than guessed. """ from __future__ import annotations +import asyncio +import hashlib import json -from pathlib import Path -from typing import Any +import re +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 from a2a_pack import ( A2AAgent, - LLMProvisioning, - {{ auth_type }}, + ConsumerSetup, + ConsumerSetupField, + EgressPolicy, + NoAuth, Pricing, + Resources, RunContext, WorkspaceAccess, WorkspaceMode, ) -from a2a_pack.context import LLMCreds +from a2a_pack.context import AgentEvent +from a2a_pack.workspace import FileType -class AgentComplianceAuditorConfig(BaseModel): - pass +MAX_TEXT_BYTES = 256_000 +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 = """\ -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 = "agent-compliance-auditor/.deepagents/skills/" -DEEPAGENTS_RECURSION_LIMIT = 500 +class StrictModel(BaseModel): + model_config = ConfigDict(extra="forbid") -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" - 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" config_model = AgentComplianceAuditorConfig - auth_model = {{ auth_type }} - - # 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 + auth_model = NoAuth 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.", + caller_pays_llm=False, + 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( - max_files=64, + max_files=MAX_FILES_PER_AUDIT + 10, allowed_modes=(WorkspaceMode.READ_ONLY, WorkspaceMode.READ_WRITE_OVERLAY), require_reason=False, + deny_patterns=DENIED_REPO_GLOBS, + max_total_size_bytes=12 * 1024 * 1024, + ) + egress = EgressPolicy( + allow_hosts=( + "gitea.a2a.svc.cluster.local", + "argocd-server.argocd.svc.cluster.local", + "kubernetes.default.svc", + "registry-1.docker.io", + "ghcr.io", + ), + deny_internet_by_default=True, + ) + tools_used = ("httpx", "pydantic", "workspace") + consumer_setup = ConsumerSetup.from_fields( + ConsumerSetupField.config("GITEA_BASE_URL", label="Gitea base URL", required=False, input_type="url"), + 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), ) - 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}, - ) - await ctx.emit_progress("deepagent finished") - return _last_message_text(state) - - def _build_deep_agent( - 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 - - @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, + @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") + async def inventory_agent(self, ctx: RunContext[NoAuth], request: InventoryRequest) -> InventoryResult: + audit_id = _audit_id(request.target, request.options.audit_id) + await _event(ctx, "audit_started", audit_id=audit_id, skill="inventory_agent", target=request.target.agent_name) + inventory, evidence, warnings, status = await _collect_inventory(self, ctx, request.target, request.options, audit_id) + output_dir = _output_dir(audit_id) + files = await _write_json_bundle(ctx, output_dir, {"inventory.json": inventory}, warnings) + await _event(ctx, "audit_completed", audit_id=audit_id, skill="inventory_agent", findings=0) + return InventoryResult(audit_id=audit_id, target=request.target, status=status, output_dir=output_dir, generated_files=files, inventory=inventory, evidence=evidence, warnings=warnings) + + @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") + async def audit_configuration(self, ctx: RunContext[NoAuth], request: AuditConfigurationRequest) -> FindingResult: + audit_id = _audit_id(request.target, request.options.audit_id) + await _event(ctx, "audit_started", audit_id=audit_id, skill="audit_configuration", target=request.target.agent_name) + inventory = request.inventory or (await _collect_inventory(self, ctx, request.target, request.options, audit_id))[0] + findings, evidence, warnings = _configuration_findings(request.target, 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_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 + warnings.append(f"Repository workspace unavailable: {type(exc).__name__}") + return {} + files: dict[str, str] = {} + for match in view.files[:max_files]: + path = match.path.replace("\\", "/") + if not _repo_path_allowed(path): + continue + try: + raw = await view.read(path) + except Exception as exc: # noqa: BLE001 + warnings.append(f"Could not read {path}: {type(exc).__name__}") + continue + files[path] = raw[:MAX_TEXT_BYTES].decode("utf-8", errors="replace") + return files -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}" +async def _write_json_bundle(ctx: RunContext[NoAuth], output_dir: str, files: dict[str, Any], warnings: list[str]) -> list[OutputFile]: + out: list[OutputFile] = [] + for name, payload in files.items(): + rel = f"{output_dir}/{name}" + if name.endswith(".md") and isinstance(payload, str): + data = payload.encode("utf-8") + mime = "text/markdown" + 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]: - """Copy packaged DeepAgents skills into the invocation workspace. +def _configuration_findings(target: AuditTarget, inventory: dict[str, Any]) -> tuple[list[Finding], list[Evidence], list[str]]: + 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 - ``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 [] + if not card: + 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.")) + if len(skills) == 0: + 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: + schema = skill.get("input_schema") or skill.get("inputSchema") or {} + out_schema = skill.get("output_schema") or skill.get("outputSchema") or {} + if not _schema_strict(schema) or not out_schema: + findings.append(_finding(target, "A2A-SCHEMA-002", f"Skill {skill.get('name', '')} 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.")) + if exposure.get("public") is True and not exposure.get("documented_public_reason"): + 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 resources.get("max_runtime_seconds") in (None, "unknown"): + 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.")) + if secrets: + 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.")) + if repo.get("suspicious_instruction_count", 0) > 0: + 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.")) + 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: - messages = state.get("messages") or [] - if not messages: - return json.dumps(state, default=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]]: + findings: list[Finding] = [] + evidence: list[Evidence] = [] + 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): - 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 _finding(target: AuditTarget, control_id: str, title: str, severity: Severity, confidence: Confidence, evidence_ref: str, impact: str, remediation: str, verification: str) -> Finding: + base = f"{target.agent_name}|{control_id}|{title}" + suffix = hashlib.sha256(base.encode()).hexdigest()[:10] + 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) + + +def _dedupe_findings(findings: list[Finding]) -> list[Finding]: + seen: dict[str, Finding] = {} + rank = {Severity.CRITICAL: 5, Severity.HIGH: 4, Severity.MEDIUM: 3, Severity.LOW: 2, Severity.INFO: 1} + for f in findings: + prev = seen.get(f.id) + if prev is None or rank[f.severity] > rank[prev.severity]: + 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}")