Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .cspell-repo-terms.txt
Original file line number Diff line number Diff line change
Expand Up @@ -892,6 +892,7 @@ pytest
pytestmark
pythonhosted
pythonstartup
readyz
pyupgrade
pyyaml
qualname
Expand Down
5 changes: 5 additions & 0 deletions agent-governance-python/agent-mesh/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -117,6 +117,11 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

- `AuditLog.export()` and `AuditLog.export_cloudevents()` now return all matching
records instead of silently capping exports at 10,000 records.
- **Empty policy sets are visible to readiness probes.** The policy server and
governance sidecar now return `503 Not Ready` when no enabled policy rules
are loaded. Policy and generation status responses expose `effective_rules`
and `load_warnings` so operators can distinguish an empty policy set from a
healthy deployment.
- **Pending-message batch isolation.** A single malformed entry in a relay-supplied
`pending_messages` batch no longer aborts the drain; the failure is surfaced
through the error handler and the remaining queued messages are still delivered.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -161,7 +161,7 @@ def from_yaml(cls, path: str | Path) -> "TrustPolicy":
A fully-constructed ``TrustPolicy`` instance.
"""
path = Path(path)
with open(path, "r") as f:
with open(path, "r", encoding="utf-8") as f:
data = yaml.safe_load(f)
return cls(**data)

Expand All @@ -173,7 +173,7 @@ def to_yaml(self, path: str | Path) -> None:
"""
path = Path(path)
data = self.model_dump(mode="json")
with open(path, "w") as f:
with open(path, "w", encoding="utf-8") as f:
yaml.dump(data, f, default_flow_style=False, sort_keys=False)


Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,12 @@
_start_time: float = 0.0


def create_base_app(component: str, description: str) -> FastAPI:
def create_base_app(
component: str,
description: str,
*,
include_readyz: bool = True,
) -> FastAPI:
"""Create a FastAPI app with standard health/metrics endpoints."""
global _start_time
_start_time = time.monotonic()
Expand All @@ -40,9 +45,11 @@ def create_base_app(component: str, description: str) -> FastAPI:
async def healthz() -> dict[str, str]:
return {"status": "ok", "component": component}

@app.get("/readyz", tags=["health"])
async def readyz() -> dict[str, str]:
return {"status": "ready", "component": component}
if include_readyz:

@app.get("/readyz", tags=["health"])
async def readyz() -> dict[str, str]:
return {"status": "ready", "component": component}

@app.get("/metrics", tags=["observability"])
async def metrics() -> PlainTextResponse:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@

import yaml
from fastapi import HTTPException
from fastapi.responses import JSONResponse
from pydantic import BaseModel, Field

from agentmesh.governance.policy import PolicyDecision, PolicyEngine
Expand All @@ -29,6 +30,7 @@
app = create_base_app(
"policy-server",
"Evaluates governance policies against agent actions.",
include_readyz=False,
)

POLICY_DIR = os.getenv("AGENTMESH_POLICY_DIR", "/etc/agentmesh/policies")
Expand All @@ -39,11 +41,42 @@
_trust_evaluator: PolicyEvaluator | None = None
_loaded_count: int = 0
_GOVERNANCE_ONLY_KEYS = frozenset({"agent", "agents", "default_action", "extends", "scope"})
_effective_rule_count: int = 0
_load_warnings: list[str] = []
Comment thread
Ricky-G marked this conversation as resolved.


@app.get("/readyz", tags=["health"], response_model=None)
async def readyz() -> JSONResponse:
payload = {
"status": "ready" if _effective_rule_count > 0 else "not-ready",
"component": "policy-server",
"total_loaded": _loaded_count,
"effective_rules": _effective_rule_count,
"policy_dir": POLICY_DIR,
"load_warnings": list(_load_warnings),
}
if _effective_rule_count == 0:
return JSONResponse(status_code=503, content=payload)
return JSONResponse(content=payload)


Comment thread
Ricky-G marked this conversation as resolved.
def _validate_load_warnings() -> None:
"""Record a warning when a load completes without effective rules."""
global _load_warnings

_load_warnings = []
if _effective_rule_count == 0:
warning = (
f"Policy load validation: no effective rules loaded from {POLICY_DIR}; "
"readiness remains blocked until an enabled policy rule is loaded."
)
logger.warning(warning)
_load_warnings.append(warning)


def _load_policies() -> None:
"""Load all YAML/JSON policy files from POLICY_DIR."""
global _engine, _trust_policies, _trust_evaluator, _loaded_count
global _engine, _trust_policies, _trust_evaluator, _loaded_count, _effective_rule_count

policy_path = Path(POLICY_DIR)
Comment thread
Ricky-G marked this conversation as resolved.
if not policy_path.is_dir():
Expand All @@ -58,6 +91,7 @@ def _load_policies() -> None:
local_engine = PolicyEngine()
local_trust: list[TrustPolicy] = []
governance_count = 0
effective_rule_count = 0
errors: list[tuple[str, Exception]] = []

try:
Expand All @@ -80,8 +114,9 @@ def _load_policies() -> None:
continue

try:
local_engine.load_yaml(content)
policy = local_engine.load_yaml(content)
governance_count += 1
effective_rule_count += sum(rule.enabled for rule in policy.rules)
logger.info("Loaded governance policy: %s", f.name)
except Exception as ge:
try:
Expand Down Expand Up @@ -117,8 +152,9 @@ def _load_policies() -> None:

for f in (path for path in discovered if path.suffix == ".json"):
try:
local_engine.load_json(f.read_text(encoding="utf-8"))
policy = local_engine.load_json(f.read_text(encoding="utf-8"))
governance_count += 1
effective_rule_count += sum(rule.enabled for rule in policy.rules)
except Exception as exc:
errors.append((f.name, exc))

Expand All @@ -138,11 +174,15 @@ def _load_policies() -> None:
_trust_evaluator = PolicyEvaluator(_trust_policies) if _trust_policies else None

_loaded_count = governance_count + len(_trust_policies)
_effective_rule_count = effective_rule_count + sum(
len(policy.rules) for policy in _trust_policies
)
logger.info(
"Loaded %d governance + %d trust policies",
governance_count,
len(_trust_policies),
)
_validate_load_warnings()


@app.on_event("startup")
Expand Down Expand Up @@ -232,8 +272,10 @@ async def list_policies() -> dict[str, Any]:
"""List all loaded policies."""
return {
"total_loaded": _loaded_count,
"effective_rules": _effective_rule_count,
"trust_policies": len(_trust_policies),
"policy_dir": POLICY_DIR,
"load_warnings": list(_load_warnings),
}


Expand All @@ -251,7 +293,9 @@ async def reload_policies() -> dict[str, Any]:
return {
"status": "reloaded",
"total_loaded": _loaded_count,
"effective_rules": _effective_rule_count,
"trust_policies": len(_trust_policies),
"load_warnings": list(_load_warnings),
}


Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@
from typing import Any, Literal

from fastapi import FastAPI
from fastapi.responses import PlainTextResponse
from fastapi.responses import JSONResponse, PlainTextResponse
from pydantic import BaseModel, ConfigDict, Field

from agentmesh.governance.policy import PolicyEngine as _PolicyEngine
Expand Down Expand Up @@ -59,24 +59,30 @@ async def health() -> dict[str, str]:
"""Liveness probe."""
return {"status": "ok", "component": "governance-sidecar"}

@app.get("/ready", tags=["health"])
async def ready() -> dict[str, Any]:
"""Readiness probe. Reports loaded policy count."""
return {
"status": "ready",
def _readiness_response() -> JSONResponse:
generation = _policy_state[1]
payload: dict[str, Any] = {
"status": "ready" if generation.effective_rules > 0 else "not-ready",
"component": "governance-sidecar",
**_policy_state[1].model_dump(exclude={"files"}),
**generation.model_dump(exclude={"files"}),
}
status_code = 200 if generation.effective_rules > 0 else 503
return JSONResponse(status_code=status_code, content=payload)

@app.get("/ready", tags=["health"], response_model=None)
async def ready() -> JSONResponse:
"""Readiness probe. Reports loaded policy count."""
return _readiness_response()

@app.get("/healthz", tags=["health"])
async def healthz() -> dict[str, str]:
"""Kubernetes-style liveness probe."""
return {"status": "ok", "component": "governance-sidecar"}

@app.get("/readyz", tags=["health"])
async def readyz() -> dict[str, str]:
@app.get("/readyz", tags=["health"], response_model=None)
async def readyz() -> JSONResponse:
"""Kubernetes-style readiness probe."""
return {"status": "ready", "component": "governance-sidecar"}
return _readiness_response()

# ── Metrics endpoint ─────────────────────────────────────────────

Expand Down Expand Up @@ -209,16 +215,18 @@ class PolicyFileLoad(BaseModel):


class PolicyLoadGeneration(BaseModel):
"""Immutable manifest of a completed load; counts refer to files, not unique names."""
"""Immutable manifest of a completed load."""

model_config = ConfigDict(frozen=True)
policy_set_id: str
policy_set_status: Literal["complete", "degraded", "rejected", "not_loaded"]
policies_discovered: int
policies_loaded: int
effective_rules: int = 0
policies_failed: int
directory_status: Literal["available", "unavailable", "not_loaded"]
files: tuple[PolicyFileLoad, ...]
load_warnings: tuple[str, ...] = ()


# ── Internal state ───────────────────────────────────────────────────
Expand All @@ -231,6 +239,7 @@ class PolicyLoadGeneration(BaseModel):
policy_set_status="not_loaded",
policies_discovered=0,
policies_loaded=0,
effective_rules=0,
policies_failed=0,
directory_status="not_loaded",
files=(),
Expand All @@ -246,6 +255,7 @@ def _load_policies() -> PolicyLoadGeneration:

_policy_dir = os.getenv("AGT_POLICY_DIR", "/etc/agt/policies")
engine = PolicyEngine()
loaded_policies: dict[str, Any] = {}

policy_path = Path(_policy_dir)
files = []
Expand All @@ -266,7 +276,8 @@ def _load_policies() -> PolicyLoadGeneration:
try:
content = f.read_bytes()
digest = hashlib.sha256(content).hexdigest()
loader(content.decode("utf-8"))
policy = loader(content.decode("utf-8"))
loaded_policies[policy.name] = policy
except Exception as exc:
error_type = type(exc).__name__
logger.warning("Skipped policy %r: %s", f.name, error_type)
Expand All @@ -285,6 +296,10 @@ def _load_policies() -> PolicyLoadGeneration:
}
canonical = json.dumps(manifest, sort_keys=True, separators=(",", ":"))
failed = sum(entry.status == "failed" for entry in files)
effective_rules = sum(
sum(rule.enabled for rule in policy.rules)
for policy in loaded_policies.values()
)

# Fail-closed (#3536 review): when files fail, publish the generation
# as 'degraded' so evaluate_policy can deny based on policies_failed.
Expand All @@ -302,14 +317,25 @@ def _load_policies() -> PolicyLoadGeneration:
if failed or directory_status == "unavailable":
policy_set_status = "degraded"

load_warnings: tuple[str, ...] = ()
if effective_rules == 0:
warning = (
f"Policy load validation: no effective rules loaded from {_policy_dir}; "
"readiness remains blocked until an enabled policy rule is loaded."
)
logger.warning(warning)
load_warnings = (warning,)

generation = PolicyLoadGeneration(
policy_set_id="sha256:" + hashlib.sha256(canonical.encode("utf-8")).hexdigest(),
policy_set_status=policy_set_status,
policies_discovered=len(files),
policies_loaded=len(files) - failed,
effective_rules=effective_rules,
policies_failed=failed,
directory_status=directory_status,
files=tuple(files),
load_warnings=load_warnings,
)
serialized = generation.model_dump_json()
_policy_state = (engine, generation)
Expand Down
Loading
Loading