Oracle监视器🔮
多智能体人工智能系统的实时可观察性平台
    
______________________________________________________________________
🎯 什么是Oracle Monitor?
想象一下,调试一个分布式系统,其中多个AI代理正在跨Kubernetes Pod、消息队列和第三方API编排复杂的工作流。传统监控工具为您展示 *什么* 发生了。Oracle监视器显示 *一切* --实时。
Oracle监视器 是专门为多智能体AI系统构建的综合可观测性平台。它提供:
- 统一状态快照:随时查看整个系统的状态——从Kubernetes pod健康状况到代理推理链再到LLM令牌消费
- 实时智能:WebSocket驱动的仪表板会在基础设施中的任何地方发生变化时立即更新
- 历史时间旅行:浏览过去的系统状态,调试几小时或几天前发生的问题
- 代理人反思:在代理推理任务和做出决策时,窥探代理的“想法”
- 成本情报:跟踪在LLM提供商中花费的每一个API调用、令牌和美元
这不仅仅是监控。它对你的人工智能基础设施来说是无所不知的。
______________________________________________________________________
🌟 Oracle Monitor存在的原因
在RVCE从事自主代理系统工作期间,我一直遇到同样令人沮丧的问题: 没有全面的上下文,分布式调试是不可能的.
当代理发生故障时,故障通常会从一系列事件中级联:
- Kafka消息延迟200毫秒
- 代理在等待该消息时超时
- 它重试,达到OpenAI速率限制
- pod被OOMKilled,因为它在重试逻辑期间泄漏了内存
- Kubernetes重新安排了它,但任务现在是孤立的
传统工具会单独向我展示每一件作品。普罗米修斯将展示OOMKill。Kafka UI将显示消息延迟。OpenAI的仪表板将显示速率限制。但连接这些点需要脑力体操和数小时的原木考古。
Oracle Monitor通过将所有内容聚合到一个带有时间戳的快照中来解决这个问题。 现在,当出现问题时,我可以倒退到那个确切的时刻,在一个统一的视图中看到完整的画面——每个代理的状态、每个队列深度、每个API调用、每个吊舱重启。
______________________________________________________________________
🏗️ 系统架构
该平台围绕三个核心原则进行设计:
- 非侵入性观察:收藏家在不干扰系统的情况下监视您的系统
- 集中式聚合:所有遥测数据流都通过一个聚合点
- 实时同步:状态更改通过Supabase实时即时传播到UI
高级数据流
graph TB
subgraph "Agent Runtime Environment"
A1[Research Agent]
A2[Writer Agent]
A3[Analysis Agent]
end
subgraph "Infrastructure Layer"
K8S[Kubernetes API]
KAFKA[Redpanda/Kafka]
LLM[OpenAI/Anthropic APIs]
end
subgraph "Collection Layer"
C1[K8s Collector
Pod Metrics]
C2[Kafka Collector
Queue Depth]
C3[LLM Collector
Token Usage]
end
subgraph "Intelligence Layer"
AGG[State Aggregator
Schema Validator]
DB[(Supabase
Time-Series Storage)]
end
subgraph "Presentation Layer"
UI[React Dashboard
Glassmorphism UI]
end
A1 --> KAFKA
A2 --> KAFKA
A3 --> KAFKA
A1 -.->|API Calls| LLM
A2 -.->|API Calls| LLM
K8S -.->|Watch API| C1
KAFKA -.->|Consumer Groups| C2
LLM -.->|Usage Logs| C3
C1 --> AGG
C2 --> AGG
C3 --> AGG
AGG -->|Validated Snapshots| DB
DB -->|WebSocket Events| UI
style AGG fill:#ff6b6b,stroke:#c92a2a,stroke-width:3px
style DB fill:#4dabf7,stroke:#1971c2,stroke-width:2px
style UI fill:#51cf66,stroke:#2f9e44,stroke-width:2px组件分解
🤖 自治代理 (/agents)
你系统中的实际工人。这些Python微服务执行任务、做出决策并与外部API交互。
当前实施:
- 研究代理:使用Playwright进行网络抓取和信息收集
- 编剧代理:使用GPT-4进行内容合成和报告生成
- 分析代理:数据处理和模式识别
遥测发射:
# Agents emit structured events to Kafka
{
"agent_id": "research-agent-7f8d9",
"timestamp": "2024-02-06T14:32:11.482Z",
"event_type": "task_started",
"task_id": "research-tech-trends-2024",
"internal_state": {
"reasoning": "Breaking down query into 3 sub-searches",
"tools_planned": ["web_search", "scrape_page", "summarize"]
}
}📡 收集器 (/collectors)
专业观察员,在不干扰运营的情况下监控基础设施的不同层。
K8s收集器 (k8s_collector/)
- 手表吊舱生命周期事件(创建、运行、失败、OOMKilled)
- 报废资源使用指标(CPU、内存、网络I/O)
- 跟踪副本计数和运行状况检查
- 检测重启循环和碰撞模式
卡夫卡收藏家 (kafka_collector/)
- 监控所有主题分区的消费者延迟
- 计算端到端消息延迟
- 识别死信队列堆积
- 跟踪吞吐量和背压指示器
LLM收集器 (llm_collector/)
- 通过代理中间件拦截API请求/响应
- 聚合代币消费(提示+完成)
- 使用当前定价计算每个请求的成本
- 检测速率限制命中和配额耗尽
⚡ 状态聚合器 (/aggregator)
Oracle Monitor的大脑。此服务:
- 消耗 来自所有收集器的遥测流
- 合并 将部分状态转换为统一快照
- 验证 反对
oracle_state.schema.json确保数据完整性 - 持续 每2秒进行一次Supabase快照
- 触发器 实时更新连接的仪表板
聚合逻辑:
# Pseudo-code for aggregation cycle
while True:
snapshot = {
"timestamp": now(),
"kubernetes": k8s_collector.get_state(),
"kafka": kafka_collector.get_state(),
"llm_usage": llm_collector.get_state(),
"agents": [
research_agent.get_state(),
writer_agent.get_state(),
]
}
# Validate against JSON Schema
validate(snapshot, oracle_state_schema)
# Persist to time-series database
supabase.table("system_snapshots").insert(snapshot)
sleep(2)💾 Supabase后端
同时充当持久层和实时事件总线。
数据库架构:
-- Main snapshots table
CREATE TABLE system_snapshots (
id UUID PRIMARY KEY DEFAULT uuid_generate_v4(),
timestamp TIMESTAMPTZ NOT NULL,
snapshot JSONB NOT NULL,
-- Indexed for fast time-range queries
CONSTRAINT snapshots_timestamp_idx UNIQUE (timestamp)
);
-- Enable real-time subscriptions
ALTER TABLE system_snapshots REPLICA IDENTITY FULL;为什么选择Supabase?
- 内置WebSocket支持(无需自定义Socket.io服务器)
- JSONB列类型,用于复杂查询的快速索引
- 多租户部署的行级安全
- 丰富的免费开发层
🎨 Glassmorphism仪表板 (/frontend)
受苹果设计语言启发的现代React应用程序。UI是使用以下方式构建的:
- 维特:开发过程中闪电般快速的HMR
- 尾风CSS:使用自定义glassmorphism组件进行实用性优先的造型
- Recharts:用于时间序列可视化的可组合图表库
- Supabase JS:快照更新的实时订阅
主要特点:
- 代理网格:显示每个代理当前任务和推理的实时卡
- 基础设施健康:每个Kubernetes pod的CPU/内存仪表
- 消息流:实时Kafka主题吞吐量图
- 成本跟踪器:按年龄细分的API LLM支出总额
- 提醒票务员:滚动横幅显示系统异常(速率限制、OOMKill等)
______________________________________________________________________
📂 存储库结构
oracle-monitor/
│
├── 🤖 agents/ # Agent implementations
│ ├── base/
│ │ ├── agent.py # Abstract BaseAgent class
│ │ ├── kafka_client.py # Kafka producer wrapper
│ │ └── telemetry.py # Structured logging utilities
│ │
│ ├── research_agent/
│ │ ├── agent.py # Research-specific logic
│ │ ├── tools/ # Web scraping, search tools
│ │ └── Dockerfile
│ │
│ └── writer_agent/
│ ├── agent.py # Writing-specific logic
│ ├── templates/ # Report templates
│ └── Dockerfile
│
├── 🔌 api/ # FastAPI backend
│ ├── main.py # API routes and startup
│ ├── models.py # Pydantic schemas
│ └── config.py # Environment configuration
│
├── 📡 collectors/ # Monitoring services
│ ├── k8s_collector/
│ │ ├── collector.py # Pod watcher implementation
│ │ ├── metrics.py # Resource usage calculations
│ │ └── Dockerfile
│ │
│ ├── kafka_collector/
│ │ ├── collector.py # Consumer lag monitoring
│ │ └── Dockerfile
│ │
│ └── llm_collector/
│ ├── collector.py # API call interceptor
│ ├── pricing.py # Token cost calculations
│ └── Dockerfile
│
├── ⚡ aggregator/ # State aggregation service
│ ├── aggregator.py # Main aggregation loop
│ ├── validator.py # JSON Schema validation
│ └── Dockerfile
│
├── 🎨 frontend/ # React dashboard
│ ├── public/
│ ├── src/
│ │ ├── components/
│ │ │ ├── AgentCard.jsx # Individual agent status
│ │ │ ├── InfrastructureGrid.jsx
│ │ │ ├── CostTracker.jsx
│ │ │ └── AlertTicker.jsx
│ │ │
│ │ ├── services/
│ │ │ └── supabase.js # Real-time subscriptions
│ │ │
│ │ ├── App.jsx
│ │ └── main.jsx
│ │
│ ├── package.json
│ └── vite.config.js
│
├── 🏗️ infrastructure/ # Deployment configs
│ ├── kubernetes/
│ │ ├── secrets/
│ │ │ └── api-keys-secret.yaml
│ │ │
│ │ ├── deployments/
│ │ │ ├── research-agent.yaml
│ │ │ ├── writer-agent.yaml
│ │ │ ├── kafka-collector.yaml
│ │ │ └── aggregator.yaml
│ │ │
│ │ └── services/
│ │ └── kafka-service.yaml
│ │
│ └── kafka/
│ └── redpanda-config.yaml # Message broker setup
│
├── 📜 schema/ # Data contracts
│ └── oracle_state.schema.json # System snapshot schema
│
├── 📖 docs/ # Documentation
│ ├── ARCHITECTURE.md
│ ├── API_REFERENCE.md
│ └── TROUBLESHOOTING.md
│
├── RUN_INSTRUCTIONS.md # Step-by-step setup guide
├── requirements.txt # Python dependencies
└── README.md # You are here______________________________________________________________________
🚀 入门指南
先决条件
在开始之前,请确保已安装以下内容:
- Docker桌面 (v20.10+)已启用Kubernetes
- 迷你 (v1.25+)适用于本地Kubernetes集群
- Python 3.11+ 使用pip
- Node.js 18+ 使用npm或yarn
- kubectl 的 CLI工具
- 辅助数据库账户 (免费版运行良好)
步骤1:克隆和设置
# Clone the repository
git clone https://github.com/pranavjambur/oracle-monitor.git
cd oracle-monitor
# Install Python dependencies
pip install -r requirements.txt
# Install frontend dependencies
cd frontend
npm install
cd ..步骤2:配置辅助数据库
- 在以下位置创建新项目 网站 supabase.com
- 运行架构迁移:
-- In Supabase SQL Editor
CREATE TABLE system_snapshots (
id UUID PRIMARY KEY DEFAULT uuid_generate_v4(),
timestamp TIMESTAMPTZ NOT NULL,
snapshot JSONB NOT NULL
);
CREATE INDEX idx_snapshots_timestamp ON system_snapshots(timestamp DESC);- 创建
.env项目根目录中的文件:
SUPABASE_URL=https://your-project.supabase.co
SUPABASE_ANON_KEY=your-anon-key
OPENAI_API_KEY=sk-...
ANTHROPIC_API_KEY=sk-ant-...步骤3:启动本地Kubernetes
# Start Minikube cluster
minikube start --driver=docker --cpus=4 --memory=8192
# Verify cluster is running
kubectl cluster-info
# Create namespace
kubectl create namespace oracle-monitor
kubectl config set-context --current --namespace=oracle-monitor步骤4:部署基础设施
# Build agent Docker images
minikube image build -t oracle-research-agent:latest -f agents/research_agent/Dockerfile .
minikube image build -t oracle-writer-agent:latest -f agents/writer_agent/Dockerfile .
# Create Kubernetes secrets for API keys
kubectl create secret generic api-keys \
--from-literal=openai-key=$OPENAI_API_KEY \
--from-literal=anthropic-key=$ANTHROPIC_API_KEY \
--from-literal=supabase-url=$SUPABASE_URL \
--from-literal=supabase-key=$SUPABASE_ANON_KEY
# Deploy Kafka/Redpanda
kubectl apply -f infrastructure/kafka/redpanda-config.yaml
# Deploy collectors and aggregator
kubectl apply -f infrastructure/kubernetes/deployments/
# Wait for pods to be ready
kubectl wait --for=condition=ready pod --all --timeout=300s步骤5:启动仪表板
cd frontend
npm run dev打开浏览器 http://localhost:5173 您应该看到Oracle Monitor仪表板栩栩如生!
______________________________________________________________________
💡 使用示例
场景1:调试失败的研究任务
一个研究代理在处理任务时崩溃。以下是Oracle Monitor如何帮助您进行调试:
- 导航到时间线:向后拖动到崩溃时间戳
- 检查代理状态:看到代理人正在推理:
"Awaiting web_search tool response" - 查看Kafka指标:注意5秒的延迟峰值
tool-responses话题 - 检查LLM日志:请参阅OpenAI API的429速率限制错误
- 连点成线:工具响应延迟,代理超时并重试,达到速率限制
根本原因:工具响应主题上的消费者并行性不足。
修复:增加Kafka分区数量和消费者副本。
场景2:优化代币成本
你注意到你的法学硕士成本本周翻了一番。Oracle Monitor显示:
- 成本跟踪显示:Writer Agent消耗的令牌比上周多3倍
- 代理人反思:最近的任务显示了过于冗长的推理链
- 即时分析:系统提示已更新,以包含更多示例
行动:系统提示中的示例计数减少,每个任务的令牌减少了60%。
______________________________________________________________________
🔧 高级配置
调整快照频率
编辑 aggregator/aggregator.py:
SNAPSHOT_INTERVAL_SECONDS = 2 # Change to 5 for less granular data启用其他LLM提供程序
增添 llm_collector/pricing.py:
PRICING = {
"openai": {"gpt-4": {"input": 0.03, "output": 0.06}},
"anthropic": {"claude-3-opus": {"input": 0.015, "output": 0.075}},
"cohere": {"command": {"input": 0.001, "output": 0.002}} # Add new
}自定义代理实现
通过扩展创建新代理 BaseAgent:
from agents.base.agent import BaseAgent
class MyCustomAgent(BaseAgent):
def process_task(self, task: dict):
self.emit_telemetry("task_started", {
"reasoning": "Analyzing customer sentiment",
"tools_planned": ["sentiment_analysis", "report_generator"]
})
# Your custom logic here
result = self.analyze_sentiment(task["data"])
self.emit_telemetry("task_completed", {
"result_summary": result
})______________________________________________________________________
🤝 贡献
欢迎投稿!以下是您可以提供帮助的方式:
报告Bug
发现bug了吗?请通过以下方式打开问题:
- 重现步骤
- 预期行为与实际行为
- 相关日志来自 `kubectl logs
`
功能请求
有主意吗?打开一个标记为的问题 enhancement 并描述:
- 它解决的用例
- 它如何适应当前的架构
- 任何实施想法
拉取请求
- 分叉回购
- 创建要素分支(
git checkout -b feature/amazing-feature) - 进行更改
- 如果适用,添加测试
- 以明确的信息提交(
git commit -m 'Add sentiment analysis agent') - 推你的叉子(
git push origin feature/amazing-feature) - 打开拉取请求
______________________________________________________________________
📊 性能基准
在配备Minikube的MacBook Pro(M1,16GB RAM)上测试:
| 度量 | 值 |
|---|---|
| 快照延迟 | ~150ms(收集到UI渲染) |
| 仪表板FPS | 60fps,含10个活性剂 |
| 内存占用 | 总计约2GB(所有服务) |
| Supabase写入/分钟 | 30(以2秒的快照间隔) |
| 历史查询速度 | 24小时范围内\<100ms |
______________________________________________________________________
🛠️ 故障排除
Pod持续重启
# Check pod logs
kubectl logs
--previous
# Common issue: Missing API keys
kubectl get secret api-keys -o yaml仪表板未更新
- 验证是否启用了Supabase实时(项目设置→ API → 实时)
- 检查浏览器控制台是否存在WebSocket错误
- 确认聚合器正在写入快照:
kubectl logs deployment/aggregator高卡夫卡滞后
# Check consumer group status
kubectl exec -it kafka-0 -- kafka-consumer-groups.sh \
--bootstrap-server localhost:9092 \
--describe --group oracle-collectors有关更多帮助,请参阅 故障排除.md.
______________________________________________________________________
📄 许可证
此项目根据MIT许可证获得许可-请参阅 许可证 文件以获取详细信息。
______________________________________________________________________
🙏 致谢
这个项目是我工作的一部分 班加罗尔R.V.工程学院特别感谢:
- 我的分布式系统架构指导导师
- Kubernetes、React和FastAPI等令人难以置信的工具的开源社区
- 使每个人都可以访问实时数据库的Supabase
______________________________________________________________________
📬 联系
普拉纳夫·詹布尔\ 计算机科学与工程\ 班加罗尔R.V.工程学院
📧 pranavvjambur.cs23@rvce.edu.in\ 🔗 | 领英
______________________________________________________________________
使用Python、FastAPI、React、Kubernetes和Supabase构建
*如果这个项目对你有帮助,考虑给它一个!*
