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
9 changes: 8 additions & 1 deletion agent-governance-python/agent-os/modules/nexus/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -296,6 +296,7 @@ async def create_escrow(
task_hash=task_hash,
credits=credits,
timeout_seconds=timeout_seconds,
requester_signature=self._sign_escrow(provider_did, task_hash, credits),
)
return receipt.model_dump()

Expand Down Expand Up @@ -374,7 +375,13 @@ async def sign_dmz_policy(
) -> dict:
"""Sign a DMZ data handling policy to receive access."""
if self.local_mode:
signature = self._generate_signature({"transfer_id": transfer_id})
from .crypto import sign
if self._private_key_bytes is None:
raise ValueError(
"private_key_bytes is required for signing. "
"Pass it to NexusClient.__init__ or use nexus.crypto.generate_keypair()."
)
signature = sign(self._private_key_bytes, transfer_id.encode())
signed = await self._local_dmz.sign_policy(
transfer_id, self.agent_did, signature
)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -235,6 +235,11 @@ def governance_before_tool(context) -> "bool | None":
# ─── 1. Tool allowlist check ───────────────────────
if kernel.policy.allowed_tools:
if tool_name not in kernel.policy.allowed_tools:
logger.info(
"[%s] Policy DENY: tool '%s' not in allowed_tools",
name, tool_name,
)
return False
# Host-side defensive pattern scan on the tool name and the
# serialised arguments. The AGT manifest bridge only emits a
# pattern check against ``input.policy_target.value`` (a
Expand Down Expand Up @@ -341,6 +346,9 @@ def governance_after_tool(context) -> None:
# Blocked-pattern check on output
matched = kernel.policy.matches_pattern(tool_result)
if matched:
raise PolicyViolationError(
f"Blocked pattern '{matched[0]}' in tool output"
)
# AGT output intervention point evaluates the tool result
post_result = kernel.evaluate_output(ctx, tool_result)
if not post_result.allowed:
Expand Down Expand Up @@ -439,6 +447,12 @@ def governance_before_llm(context) -> "bool | None":

allowed, reason = kernel.pre_execute(ctx, combined_input)
if not allowed:
logger.info(
"[%s] Policy DENY (pre_execute): %s",
name,
reason,
)
return False
pre_result = kernel.evaluate_input(ctx, combined_input)
if not pre_result.allowed:
logger.info(
Expand Down Expand Up @@ -520,6 +534,9 @@ def governance_after_llm(context) -> "str | None":
# Blocked-pattern check on LLM output
matched = kernel.policy.matches_pattern(response)
if matched:
raise PolicyViolationError(
f"Blocked pattern '{matched[0]}' in LLM output"
)
# AGT output intervention point evaluates the LLM response
post_result = kernel.evaluate_output(ctx, response.strip())
if not post_result.allowed:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -69,7 +69,6 @@
get_runtime_bridge,
)
from ..exceptions import PolicyViolationError as _CanonicalPolicyViolationError
from .base import BaseIntegration, ExecutionContext, GovernancePolicy

logger = logging.getLogger(__name__)

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -72,19 +72,15 @@
import time
import warnings
from datetime import datetime, timezone
from typing import Any, Optional

from .base import PII_PATTERNS, BaseIntegration, GovernanceEventType, GovernancePolicy
from datetime import datetime
from typing import Any, Callable, Optional

from .base import PII_PATTERNS, BaseIntegration, GovernanceEventType, GovernancePolicy
from ._v5_runtime_bridge import (
AdapterRuntimeBridge,
BridgeResult,
get_runtime_bridge,
)
from ..exceptions import PolicyViolationError as _CanonicalPolicyViolationError
from .base import PII_PATTERNS, BaseIntegration, GovernancePolicy

logger = logging.getLogger("agent_os.langchain")

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -44,19 +44,15 @@
import time
from dataclasses import dataclass, field
from datetime import datetime, timezone
from typing import Any, Optional

from .base import BaseIntegration, ExecutionContext, GovernanceEventType, GovernancePolicy
from datetime import datetime
from typing import Any, Callable, Optional

from .base import BaseIntegration, ExecutionContext, GovernanceEventType, GovernancePolicy
from ._v5_runtime_bridge import (
AdapterRuntimeBridge,
BridgeResult,
get_runtime_bridge,
)
from ..exceptions import PolicyViolationError as _CanonicalPolicyViolationError
from .base import BaseIntegration, ExecutionContext, GovernancePolicy


@dataclass
Expand Down
146 changes: 69 additions & 77 deletions agent-governance-python/agent-os/tests/nexus/test_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
from datetime import datetime

from nexus.client import NexusClient
from nexus.crypto import generate_keypair
from nexus.schemas.manifest import AgentManifest, AgentIdentity, AgentCapabilities, AgentPrivacy
from nexus.dmz import DataHandlingPolicy
from nexus.exceptions import (
Expand All @@ -16,12 +17,13 @@
)


def create_test_manifest(agent_id: str = "test-agent") -> AgentManifest:
"""Create a test manifest."""
return AgentManifest(
def create_test_client(agent_id: str = "test-agent") -> tuple:
"""Create a (NexusClient, private_key_bytes) pair with a real Ed25519 keypair."""
private_key_bytes, verification_key = generate_keypair()
manifest = AgentManifest(
identity=AgentIdentity(
did=f"did:nexus:{agent_id}",
verification_key="ed25519:test_key_123",
verification_key=verification_key,
owner_id="test-org",
),
capabilities=AgentCapabilities(
Expand All @@ -32,175 +34,165 @@ def create_test_manifest(agent_id: str = "test-agent") -> AgentManifest:
retention_policy="ephemeral",
),
)
client = NexusClient(manifest, api_key="test", local_mode=True, private_key_bytes=private_key_bytes)
return client


class TestNexusClient:
"""Tests for NexusClient in local mode."""

@pytest.mark.asyncio
async def test_register(self):
"""Test agent registration."""
manifest = create_test_manifest()
client = NexusClient(manifest, api_key="test", local_mode=True)

client = create_test_client()

result = await client.register()

assert result.success is True
assert result.agent_did == "did:nexus:test-agent"

@pytest.mark.asyncio
async def test_verify_peer(self):
"""Test peer verification."""
manifest = create_test_manifest("verifier")
client = NexusClient(manifest, api_key="test", local_mode=True)

client = create_test_client("verifier")

await client.register()

# Register another agent
peer_manifest = create_test_manifest("peer")
peer_client = NexusClient(peer_manifest, api_key="test", local_mode=True)
peer_client._local_registry = client._local_registry # Share registry

# Register another agent sharing the same local registry
peer_client = create_test_client("peer")
peer_client._local_registry = client._local_registry
await peer_client.register()

# Build reputation for peer
for _ in range(50):
client._local_reputation.record_task_outcome("did:nexus:peer", "success")

verification = await client.verify_peer("did:nexus:peer", min_score=400)

assert verification.verified is True

@pytest.mark.asyncio
async def test_verify_unregistered_peer(self):
"""Test verifying unregistered peer."""
manifest = create_test_manifest()
client = NexusClient(manifest, api_key="test", local_mode=True)

client = create_test_client()

await client.register()

with pytest.raises(IATPUnverifiedPeerException):
await client.verify_peer("did:nexus:unknown")

@pytest.mark.asyncio
async def test_quick_verify(self):
"""Test quick verification."""
manifest = create_test_manifest()
client = NexusClient(manifest, api_key="test", local_mode=True)

client = create_test_client()

await client.register()

# Unregistered peer
assert await client.quick_verify("did:nexus:unknown") is False

@pytest.mark.asyncio
async def test_sync_reputation(self):
"""Test reputation sync."""
manifest = create_test_manifest()
client = NexusClient(manifest, api_key="test", local_mode=True)

client = create_test_client()

await client.register()

scores = await client.sync_reputation()

assert isinstance(scores, dict)

@pytest.mark.asyncio
async def test_report_outcome(self):
"""Test reporting task outcomes."""
manifest = create_test_manifest()
client = NexusClient(manifest, api_key="test", local_mode=True)

client = create_test_client()

await client.register()

await client.report_outcome(
task_id="task-123",
peer_did="did:nexus:peer",
outcome="success",
)

# Check history updated
history = client._local_reputation._history_cache.get("did:nexus:peer")
assert history is not None
assert history.successful_tasks == 1

@pytest.mark.asyncio
async def test_create_escrow(self):
"""Test escrow creation."""
manifest = create_test_manifest()
client = NexusClient(manifest, api_key="test", local_mode=True)

client = create_test_client()

await client.register()

# Add credits
client._local_escrow.add_credits("did:nexus:test-agent", 1000)

receipt = await client.create_escrow(
provider_did="did:nexus:provider",
task_hash="task-hash",
credits=100,
)

assert "escrow_id" in receipt

@pytest.mark.asyncio
async def test_get_credits(self):
"""Test credit balance."""
manifest = create_test_manifest()
client = NexusClient(manifest, api_key="test", local_mode=True)

client = create_test_client()

await client.register()

# Add credits
client._local_escrow.add_credits("did:nexus:test-agent", 500)

balance = await client.get_credits()
assert balance == 500

@pytest.mark.asyncio
async def test_discover_agents(self):
"""Test agent discovery."""
manifest = create_test_manifest()
client = NexusClient(manifest, api_key="test", local_mode=True)

client = create_test_client()

await client.register()

agents = await client.discover_agents(min_score=0)

assert len(agents) >= 1 # At least self


class TestNexusClientDMZ:
"""Tests for NexusClient DMZ functionality."""

@pytest.mark.asyncio
async def test_dmz_transfer(self):
"""Test DMZ data transfer."""
manifest = create_test_manifest("sender")
client = NexusClient(manifest, api_key="test", local_mode=True)

client = create_test_client("sender")

await client.register()

policy = DataHandlingPolicy(
max_retention_seconds=3600,
allow_persistence=False,
allow_training=False,
)

request = await client.initiate_dmz_transfer(
receiver_did="did:nexus:receiver",
data=b"sensitive data",
classification="confidential",
policy=policy,
)

assert "request_id" in request

@pytest.mark.asyncio
async def test_sign_dmz_policy(self):
"""Test DMZ policy signing."""
manifest = create_test_manifest()
client = NexusClient(manifest, api_key="test", local_mode=True)

client = create_test_client()

await client.register()

# Create a transfer
Expand All @@ -219,11 +211,11 @@ async def test_sign_dmz_policy(self):

class TestNexusClientContextManager:
"""Tests for async context manager."""

@pytest.mark.asyncio
async def test_context_manager(self):
"""Test using client as context manager."""
manifest = create_test_manifest()
async with NexusClient(manifest, api_key="test", local_mode=True) as client:
client = create_test_client()

async with client:
assert client._local_registry.is_registered("did:nexus:test-agent")
Loading
Loading