Agentic Cognitive Governance Protocol - Add AI agent oversight with real-time feedback in 5 lines of code
Project description
ACGP
Agentic Cognitive Governance Protocol - Lightweight Python SDK for AI agent oversight integration.
Developed and maintained by MeaningStack
Overview
ACGP provides a minimal-footprint way to integrate Steward Agent oversight into any AI agent pipeline. It works with any Python agent framework (CrewAI, LangChain, custom agents) without blocking your agent's execution.
Client Code -> ACGP (lightweight wrapper) -> Steward API -> Backend Infrastructure
|
v
SSE Stream (real-time feedback)
Features
- Framework-agnostic: Works with any Python agent (CrewAI, LangChain, AutoGen, custom)
- Non-blocking: Async queue ensures traces never slow down your agent
- Real-time feedback: SSE streaming delivers oversight feedback via callbacks
- Auto-batching: Accumulates events, flushes at configurable thresholds
- Retry & circuit breaker: Resilient transport with exponential backoff
- Graceful degradation: Works offline, queues events locally
- Tiny footprint: ~1.5MB core dependencies
Installation
pip install acgp
# With optional framework support
pip install acgp[crewai]
pip install acgp[langchain]
pip install acgp[all]
Quick Start
Real-time Feedback (Recommended)
The recommended pattern uses a callback for real-time feedback. Feedback arrives automatically via SSE streaming as traces are processed.
from acgp import ACGPClient
# Define callback to handle feedback
def handle_feedback(trace_id: str, feedback_data: dict):
print(f"Received feedback for trace {trace_id}")
ctq = feedback_data.get('ctq_scores', {})
print(f" Composite Score: {ctq.get('composite_score')}")
print(f" Feedback: {feedback_data.get('feedback')}")
# React to low quality scores
if ctq.get('composite_score', 1.0) < 0.5:
print(" Warning: Low quality score - consider adjustments")
# Check for human-in-the-loop requirement
if feedback_data.get('requires_hitl'):
print(" ALERT: Human review required!")
# Initialize with callback
client = ACGPClient(
api_key="sk_live_...",
endpoint="https://api.steward-agent.com",
on_feedback=handle_feedback, # Enables real-time SSE streaming
)
# Trace your agents - use unique agent_tag for each agent
with client.trace(
agent_type="planner",
agent_tag="myapp-query-planner", # Unique identifier for this agent
goal="Analyze user query"
) as trace:
result = my_planner_agent.run(query)
trace.set_output(result)
trace.set_reasoning(my_planner_agent.reasoning_log)
with client.trace(
agent_type="executor",
agent_tag="myapp-task-executor", # Unique identifier for this agent
goal="Execute the plan"
) as trace:
result = my_executor_agent.run(plan)
trace.set_output(result)
# Wait for all feedback before shutdown (auto-flushes pending traces)
client.wait_for_feedback(timeout=30.0)
client.close()
Manual Polling (Alternative)
For simpler use cases or when callbacks are not suitable:
from acgp import ACGPClient
client = ACGPClient(
api_key="sk_live_...",
endpoint="https://api.steward-agent.com",
)
# Capture trace
with client.trace(agent_type="planner", goal="Analyze query") as trace:
result = my_agent.run(query)
trace.set_output(result)
# Flush and poll for results
client.flush()
score = client.get_score(trace.trace_id) # Auto-retries until processed
feedback = client.get_feedback(trace.trace_id)
print(f"Score: {score}")
print(f"Feedback: {feedback}")
client.close()
API Reference
ACGPClient
The main client for interacting with the Steward Agent backend.
client = ACGPClient(
api_key="sk_live_...", # Required: Your API key
endpoint="https://...", # Required: Backend URL
on_feedback=callback_fn, # Optional: Real-time feedback callback
on_flush_complete=flush_callback, # Optional: Called after each batch flush
sse_config=SSEConfig(...), # Optional: SSE configuration
debug=False, # Optional: Enable debug logging
)
Methods
| Method | Description |
|---|---|
trace(agent_type, goal, ...) |
Context manager for tracing agent execution |
capture_trace(...) |
Manually capture a trace without context manager |
flush() |
Force-send all pending traces immediately |
wait_for_feedback(timeout, auto_flush) |
Wait for SSE feedback (auto-flushes by default) |
get_score(trace_id) |
Get CTQ score for a trace (with retry) |
get_feedback(trace_id) |
Get oversight feedback for a trace (with retry) |
get_metrics() |
Get SDK metrics (queue size, SSE status, etc.) |
close(timeout) |
Gracefully shutdown client |
Properties
| Property | Description |
|---|---|
sse_connected |
Boolean indicating if SSE stream is connected |
pending_feedback_count |
Number of traces waiting for feedback |
TraceContext
Returned by client.trace() context manager.
| Parameter | Required | Description |
|---|---|---|
agent_type |
Yes | Type of agent: "planner", "executor", or custom |
goal |
Yes | What the agent is trying to accomplish |
agent_tag |
Recommended | Unique identifier for this agent (defaults to "A"/"B" if not provided) |
with client.trace(
agent_type="planner",
agent_tag="myapp-planner", # Recommended: unique tag
goal="Plan the task"
) as trace:
# Set output (required)
trace.set_output("The plan is...")
# Set reasoning (recommended)
trace.set_reasoning("Step 1: Analyzed input. Step 2: ...")
# Add context information
trace.add_context("model", "gpt-4")
trace.add_context("temperature", 0.7)
# Set documents used (for RAG agents)
trace.set_documents(
documents=["doc1.pdf", "doc2.pdf"],
similarity_scores=[0.95, 0.87]
)
# Access trace_id after context exits
print(f"Trace ID: {trace.trace_id}")
Configuration
Environment Variables
# Required
ACGP_API_KEY=sk_live_...
ACGP_ENDPOINT=https://api.steward-agent.com
# Optional
ACGP_ENABLED=true # Master switch (default: true)
ACGP_DEBUG=false # Debug logging (default: false)
ACGP_FLUSH_INTERVAL=5.0 # Seconds between auto-flushes (default: 5.0)
ACGP_BATCH_SIZE=100 # Max traces per batch (default: 100)
ACGP_MAX_QUEUE_SIZE=10000 # Max queued traces (default: 10000)
ACGP_MAX_RETRIES=3 # Retry attempts (default: 3)
ACGP_FAIL_SILENTLY=true # Don't raise on errors (default: true)
SSEConfig
Customize SSE streaming behavior:
from acgp import ACGPClient, SSEConfig
sse_config = SSEConfig(
reconnect_delay=1.0, # Initial reconnect delay (seconds)
max_reconnect_delay=30.0, # Maximum reconnect delay
connection_timeout=10.0, # Connection timeout
read_timeout=300.0, # Read timeout (5 minutes)
)
client = ACGPClient(
api_key="sk_...",
on_feedback=handle_feedback,
sse_config=sse_config,
)
Decorator Pattern
For simple function-based agents:
from acgp import ACGPClient, acgp_trace
client = ACGPClient(api_key="sk_...")
@acgp_trace(client, agent_type="planner", agent_tag="myapp-planner")
def my_planning_agent(query: str) -> str:
# Your agent logic
return plan
@acgp_trace(client, agent_type="executor", agent_tag="myapp-executor")
def my_execution_agent(plan: str) -> str:
# Your agent logic
return result
# Traces are automatically captured
plan = my_planning_agent("What should I do?")
result = my_execution_agent(plan)
Agent Tagging Best Practices
The agent_tag parameter identifies individual agents in your dashboard. While it defaults to "A" for planners and "B" for executors, you should provide unique, descriptive tags to avoid confusion when:
- Running multiple applications with the same API key
- Having multiple agents of the same type
- Distinguishing agents across different environments
Recommended Naming Conventions
# Format: {app-name}-{agent-role} or {app-name}/{agent-role}
# Good: Unique, descriptive tags
with client.trace(
agent_type="planner",
agent_tag="inventory-system-planner",
goal="Plan inventory update"
) as trace:
...
with client.trace(
agent_type="executor",
agent_tag="inventory-system-executor",
goal="Execute inventory changes"
) as trace:
...
# For multi-agent systems
with client.trace(agent_type="planner", agent_tag="crm-lead-scorer", goal="...") as trace: ...
with client.trace(agent_type="executor", agent_tag="crm-email-sender", goal="...") as trace: ...
with client.trace(agent_type="executor", agent_tag="crm-data-enricher", goal="...") as trace: ...
Why Unique Tags Matter
Without unique tags, agents from different apps get mixed in the dashboard:
# App 1: E-commerce
with client.trace(agent_type="planner", goal="...") as trace: # Gets tag "A"
...
# App 2: Customer Support (same API key)
with client.trace(agent_type="planner", goal="...") as trace: # Also gets tag "A"
...
# Dashboard shows both under "Agent A" - impossible to distinguish!
With unique tags:
# App 1: E-commerce
with client.trace(agent_type="planner", agent_tag="ecommerce-order-planner", goal="...") as trace:
...
# App 2: Customer Support
with client.trace(agent_type="planner", agent_tag="support-ticket-planner", goal="...") as trace:
...
# Dashboard shows each agent separately with clear identification
Framework Integration Examples
CrewAI
from crewai import Agent, Task, Crew
from acgp import ACGPClient
client = ACGPClient(api_key="sk_...", on_feedback=handle_feedback)
agent = Agent(role="Researcher", goal="Research topics", ...)
task = Task(description="Research AI trends", agent=agent)
crew = Crew(agents=[agent], tasks=[task])
with client.trace(
agent_type="unified",
agent_tag="crewai-research-agent", # Unique tag for this CrewAI agent
goal="Research AI trends"
) as trace:
result = crew.kickoff()
trace.set_output(str(result.raw))
trace.set_reasoning("CrewAI research workflow")
client.wait_for_feedback(timeout=30.0)
client.close()
LangChain
from langchain.agents import initialize_agent
from acgp import ACGPClient
client = ACGPClient(api_key="sk_...", on_feedback=handle_feedback)
agent = initialize_agent(tools, llm, agent="zero-shot-react-description")
with client.trace(
agent_type="executor",
agent_tag="langchain-react-agent", # Unique tag
goal=user_query
) as trace:
result = agent.run(user_query)
trace.set_output(result)
client.wait_for_feedback(timeout=30.0)
client.close()
GPT-Researcher
from gpt_researcher import GPTResearcher
from acgp import ACGPClient
client = ACGPClient(api_key="sk_...", on_feedback=handle_feedback)
researcher = GPTResearcher(query=query)
# Trace planner phase
with client.trace(
agent_type="planner",
agent_tag="gpt-researcher-planner", # Unique tag
goal=f"Plan research: {query}"
) as trace:
await researcher.conduct_research()
trace.set_output(f"Agent: {researcher.agent}")
# Trace executor phase
with client.trace(
agent_type="executor",
agent_tag="gpt-researcher-writer", # Unique tag
goal=f"Generate report: {query}"
) as trace:
report = await researcher.write_report()
trace.set_output(report)
client.wait_for_feedback(timeout=30.0)
client.close()
Metrics & Monitoring
metrics = client.get_metrics()
print(f"SDK Enabled: {metrics['enabled']}")
print(f"Configured: {metrics['configured']}")
print(f"Queue Started: {metrics['queue_started']}")
print(f"SSE Enabled: {metrics['sse_enabled']}")
print(f"SSE Connected: {metrics['sse_connected']}")
print(f"Pending Feedback: {metrics['pending_feedback']}")
print(f"Events Sent: {metrics.get('events_sent', 0)}")
print(f"Events Dropped: {metrics.get('events_dropped', 0)}")
Error Handling
By default, the SDK fails silently to avoid disrupting your agent's execution:
# Default: fail_silently=True
client = ACGPClient(api_key="sk_...") # Errors logged but not raised
# Strict mode: raise exceptions
client = ACGPClient(api_key="sk_...", fail_silently=False)
try:
with client.trace(...) as trace:
# ...
except Exception as e:
print(f"Trace failed: {e}")
Architecture
acgp/
├── __init__.py # Public exports
├── client.py # ACGPClient - main entry point
├── config.py # Configuration management
├── transport/ # Network layer
│ ├── async_queue.py # Non-blocking event queue
│ ├── batch_sender.py # Batched HTTP delivery
│ ├── sse_client.py # Real-time SSE streaming
│ └── retry.py # Retry with circuit breaker
├── models/ # Data models
└── adapters/ # Framework-specific adapters
SDK Trace Capture Flow
How the SDK captures and transmits traces without blocking agent execution.
sequenceDiagram
autonumber
participant Agent as Client Agent
participant Wrapper as StewardAgentWrapper
participant Trace as TraceContext
participant Queue as AsyncEventQueue
participant Sender as BatchSender
participant API as Backend API
Agent->>Wrapper: wrapped_agent.run(input)
Wrapper->>Trace: Create trace context
Trace->>Trace: Record start timestamp
Trace->>Trace: Capture goal
Wrapper->>Agent: Execute original agent
Agent-->>Wrapper: Return result
Wrapper->>Trace: Set output
Wrapper->>Trace: Set reasoning (if available)
Wrapper->>Trace: Record end timestamp
Trace->>Queue: Enqueue trace event
Note over Queue: Non-blocking enqueue
Queue-->>Trace: Queued (immediate return)
Wrapper-->>Agent: Return result (no delay)
Note over Queue,Sender: Background processing
loop Every flush_interval (5s default)
Queue->>Sender: Flush batch
Sender->>API: POST /ingest/batch
API-->>Sender: 202 Accepted
end
Key Points:
- Wrapper intercepts agent execution transparently
- Trace context captures timing, goal, reasoning, output
- Queue is non-blocking - agent never waits for governance
- BatchSender groups traces for efficient transmission
Core Dependencies
httpx>=0.25.0- Modern async HTTP clientpydantic>=2.0.0- Data validationpython-dotenv>=1.0.0- Environment variable loading
Development
# Clone repository
git clone https://github.com/meaningstack/acgp.git
cd acgp
# Install in development mode
pip install -e ".[dev]"
# Run tests
pytest
# Format code
black acgp/
ruff check acgp/
# Type checking
mypy acgp/
Contributing
Contributions are welcome! Please read our Contributing Guide for details.
License
MIT License - see LICENSE for details.
Copyright (c) 2025 MeaningStack
Support
- Website: https://meaningstack.com
- Documentation: TBD
- Issues: TBD
- Email: TBD
Project details
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
Built Distribution
Filter files by name, interpreter, ABI, and platform.
If you're not sure about the file name format, learn more about wheel file names.
Copy a direct link to the current filters
File details
Details for the file acgp-0.1.1.tar.gz.
File metadata
- Download URL: acgp-0.1.1.tar.gz
- Upload date:
- Size: 62.3 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.2.0 CPython/3.10.15
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
36a9f61b0eb477077925a2edd41f4e075974bb1eae91e2ac8a59860ed16a4a4e
|
|
| MD5 |
7c028ff72af5727c573bd00a59b83b6c
|
|
| BLAKE2b-256 |
bda6fb06777c68d35cf2e823e051952581d0c3ca6c0e230647d8f80636ab2b4e
|
File details
Details for the file acgp-0.1.1-py3-none-any.whl.
File metadata
- Download URL: acgp-0.1.1-py3-none-any.whl
- Upload date:
- Size: 55.1 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.2.0 CPython/3.10.15
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
3f7ce35c633bcce5bcc3b44c00e5d9877ddce0972e6e65082d57fe0386e15f5b
|
|
| MD5 |
ea1e05eff617a1a3d6b3001c517359c1
|
|
| BLAKE2b-256 |
576863772fcdabeb7ee3b922d842192518293bd2794cc3c226b59b36ee22b9ad
|