安装、配置并运行 Janus,然后完成你的第一个 Agent 到 Agent 任务交接。
Janus 仓库自带生产可用的 Dockerfile 和 docker-compose.yml。克隆即可获得所需的一切:
git clone https://github.com/agentium-lab/Janus.git
cd Janus
克隆后你会得到:
Dockerfile — 多阶段构建(Go 编译 → 精简 Alpine 运行时),产出 janus-api 二进制docker-compose.yml — 一条命令启动 Janus + PostgreSQL + NATS + Redisconfigs/janus.example.yaml — 参考配置文件migrations/ — PostgreSQL 数据库迁移(启动时自动执行)在仓库根目录下,一条命令即可构建 Janus 镜像并启动全部四个服务(janus-api、postgres:16、nats:2 含 JetStream、redis:7):
docker compose up -d --build
首次构建需要几分钟(编译 Go 二进制)。后续启动是即时的。验证所有服务健康:
curl -s http://localhost:8080/healthz
# → {"status":"ok"}
curl -s http://localhost:8080/readyz
# → {"status":"ready"} (检查 PG + NATS + Redis 连通性)
Janus 暴露两个端口:
8080 — HTTP REST API 与协议网关9090 — gRPC API若需在 Docker 之外运行 Janus(开发、自定义构建、调试):
# 启动三个依赖服务
docker compose -f deployments/smoke-deps.compose.yaml up -d postgres nats redis
# 编译服务端二进制
cd server && go build -o janus-api ./cmd/janus-api/
# 通过环境变量运行
JANUS_PG_HOST=localhost JANUS_PG_USER=janus JANUS_PG_DATABASE=janus \
JANUS_NATS_URL=nats://localhost:4222 JANUS_REDIS_ADDR=localhost:6379 \
JANUS_AUTH_ENABLED=false ./janus-api
Janus 通过环境变量(优先级最高)和/或 YAML 配置文件读取配置。Compose 文件已为本地开发配好合理默认值。
最常调整的变量(均以 JANUS_ 为前缀):
| 变量 | 默认值 | 说明 |
|---|---|---|
JANUS_CONFIG_FILE | — | YAML 配置文件路径(覆盖以下所有项) |
JANUS_HTTP_PORT | 8080 | HTTP REST API 端口 |
JANUS_GRPC_PORT | 9090 | gRPC API 端口 |
JANUS_PG_HOST | localhost | PostgreSQL 主机 |
JANUS_PG_PORT | 5432 | PostgreSQL 端口 |
JANUS_PG_USER / PASSWORD / DATABASE | janus | PostgreSQL 凭据 |
JANUS_NATS_URL | nats://localhost:4222 | NATS 连接地址 |
JANUS_REDIS_ADDR | localhost:6379 | Redis 地址(限流与心跳) |
JANUS_AUTH_ENABLED | true | 开启 API 密钥认证 |
JANUS_MIGRATION_AUTO | false | 启动时自动执行数据库迁移 |
JANUS_TLS_ENABLED | false | 开启 TLS(配合 CERT_FILE / KEY_FILE / CLIENT_CA_FILE) |
复杂部署时,通过 JANUS_CONFIG_FILE 指定 YAML 文件。可从参考文件复制起步:
cp configs/janus.example.yaml configs/janus.local.yaml
# 然后运行:
JANUS_CONFIG_FILE=configs/janus.local.yaml ./janus-api
YAML 与环境变量对应,并额外提供 outbox、心跳与可观测性的结构化分段:
server:
http_port: 8080
grpc_port: 9090
postgres:
host: localhost
port: 5432
user: janus
password: ${JANUS_PG_PASSWORD} # 支持环境变量插值
database: janus
max_conns: 20
nats:
url: "nats://localhost:4222"
cache:
driver: redis
addr: "localhost:6379"
log:
level: info # debug | info | warn | error
format: json # json | text
migration:
auto: true
path: migrations/
outbox:
enable: true
batch_size: 100
max_retries: 5 # 重试耗尽 → 死信队列
observability:
metrics:
enabled: true
path: /metrics
tracing:
enabled: false
endpoint: "" # OTLP collector 地址
第二个 YAML 文件 —— 项目配置 —— 让你以代码形式声明 Agent、预算与治理策略。它与上面的服务端配置互补(后者管连接与运行时):项目文件定义你的租户拓扑。通过 Janus CLI 应用:
# 1. 在当前目录生成初始 janus.project.yaml
janus project init
# 2. 编辑它,声明 Agent、预算与策略
# (完整 schema:configs/janus.project.example.yaml)
# 3. 预览变更(不实际应用)
janus project diff
# 4. 应用(幂等 —— 可安全在 CI/CD 中重复执行)
janus project apply
CLI 如何查找项目文件 —— 按以下顺序解析(先命中者优先):
| 优先级 | 来源 | 示例 |
|---|---|---|
| 1 | --file 参数 | janus project apply --file configs/janus.project.yaml |
| 2 | JANUS_PROJECT_FILE 环境变量 | export JANUS_PROJECT_FILE=configs/janus.project.yaml |
| 3 | 当前目录下的 janus.project.yaml(默认) | 运行 janus project init 生成它 |
JANUS_PROJECT_FILE 告诉 CLI 去哪里找文件 —— API 服务端不直接读取它。端到端示例:
# 生成初始文件
janus project init
# 在 YAML 中声明 Agent、预算与策略(见下文)
# 显式指定路径应用文件
export JANUS_PROJECT_FILE=configs/janus.project.yaml
janus project diff # 审查
janus project apply # 应用(幂等)
项目文件声明完整的租户拓扑 —— Agent 及其能力、租户级与团队级预算、allow/deny/approve 策略:
version: v1
default_tenant: acme
defaults:
protocol: custom-sdk
mailbox:
ack_wait_seconds: 300
max_deliver: 5
policy:
priority: 100
tenants:
acme:
name: Acme Engineering
agents:
code-review:
team: engineering
capabilities:
- id: code_review
data_classifications: [public, internal, confidential]
concurrency: 4
budgets:
tenant:
tpm: 2000000# 每分钟 token 数
teams:
engineering:
concurrency: 20
tpm: 600000
policies:
approve:
capabilities: [prod_deploy]# 需要人工审批
allow:
- agent: code-review
capability: code_review
deny:
- agent: intern-agent
tool: deploy.prod
data_classification:
deny:
- team: interns
classifications: [confidential, restricted]
完整 schema 见仓库中的 configs/janus.project.example.yaml。
每个字段都映射到 Janus 核心资源。下表说明文件每一层级的含义。
顶层
| 字段 | 类型 | 说明 |
|---|---|---|
version | string | Schema 版本。当前为 v1。必填。 |
default_tenant | string | 未传 --tenant 时默认操作的租户。必须存在于 tenants 下。 |
defaults | object | 所有租户与 Agent 继承的默认值,可在本地覆盖。 |
tenants | map | 租户声明,以租户 ID 为键。必填。 |
defaults — 被所有租户继承:
| 字段 | 默认值 | 说明 |
|---|---|---|
protocol | custom-sdk | Agent 通信协议。当 Agent 未单独设置时生效。 |
classification | — | 默认数据分类。取值:public、internal、confidential、restricted。 |
mailbox.ack_wait_seconds | — | 任务未被确认前,多少秒后被重新投递到邮箱。 |
mailbox.max_deliver | — | 投递尝试达到上限后,任务进入死信队列。 |
mailbox.retention_seconds | — | 已完成的邮箱状态保留多久。 |
capacity.max_concurrency | 1 | 每个 Agent 默认的最大并发进行中任务数。 |
policy.priority | — | 自动生成的策略规则的优先级(越大越先评估)。 |
租户(tenants.<id>):
| 字段 | 类型 | 说明 |
|---|---|---|
name | string | 人类可读的显示名称。 |
agents | map | Agent 声明,以 Agent ID 为键。 |
budgets | object | 租户、团队、Agent、模型与任务作用域的成本与速率限制。 |
policies | object | allow、deny、审批与数据分类规则,编译为策略模板。 |
Agent(tenants.<id>.agents.<id>):
| 字段 | 类型 | 说明 |
|---|---|---|
name | string | 显示名称。默认取 Agent ID。 |
team | string | 团队 ID。启用团队级预算与策略定位。 |
protocol | string | 为该 Agent 覆盖 defaults.protocol。 |
endpoint | string | Agent 接收任务通知的回调地址。 |
description | string | 该 Agent 职责的自由文本描述。 |
capabilities | list | 该 Agent 能处理的能力。每项是字符串或含 id、description、data_classifications 的对象。必填 — 至少一个。 |
concurrency | int | 最大并发进行中任务数。覆盖 capacity.max_concurrency。 |
rpm / tpm | int | 每分钟请求数与每分钟 token 数,用于预算执行。 |
mailbox | object | Agent 级邮箱覆盖。若省略,自动创建默认邮箱 <agent-id>.default。 |
预算限额 — 每个预算作用域接受相同的字段:
| 字段 | 说明 |
|---|---|
rpm | 每分钟最大请求数。 |
tpm | 每分钟最大 token 数。 |
concurrency | 最大并发操作数。 |
daily_usd | 规划中。暂未强制执行——待可信 token 计量落地。 |
monthly_usd | 规划中。暂未强制执行——待可信 token 计量落地。 |
预算作用域:tenant(整个租户)、teams.<id>、agents.<id>、models.<id>、model_providers.<id>、tasks.<id>。
策略 — 通过模板编译为策略规则:
| 字段 | 说明 |
|---|---|
approve.capabilities | 执行前需要人工审批的能力。 |
approve.tools | 调用前需要人工审批的工具。 |
allow / deny | 将 agent 或 team 绑定到 capability 或 tool。每条绑定必须且只能设置 agent/team 之一、capability/tool 之一。 |
data_classification.allow/deny | 控制 Agent 或团队可以接收哪些数据分类。 |
一个真实项目文件:三个 Agent、两个团队,含团队级预算与治理策略:
version: v1
default_tenant: acme
defaults:
protocol: custom-sdk
mailbox:
ack_wait_seconds: 300
max_deliver: 5
policy:
priority: 100
tenants:
acme:
name: Acme Engineering
agents:
code-review: # Agent 1:审核 Pull Request
team: engineering
capabilities:
- id: code_review
data_classifications: [public, internal, confidential]
concurrency: 4
test-runner: # Agent 2:运行测试套件
team: engineering
capabilities: [test_run]
concurrency: 8
deploy-bot: # Agent 3:生产部署
team: sre
endpoint: https://deploy.internal.acme.com
capabilities:
- id: prod_deploy
description: Deploy to production
concurrency: 1
mailbox:
ack_wait_seconds: 600 # 部署耗时更长
budgets:
tenant:
tpm: 2000000
teams:
engineering:
concurrency: 20
tpm: 600000
sre:
concurrency: 4
policies:
approve:
capabilities: [prod_deploy] # 需要人工审批
tools: [deploy.prod]
allow:
- agent: code-review
capability: code_review
- agent: test-runner
capability: test_run
deny:
- agent: deploy-bot
tool: deploy.prod # 未经审批则被阻止
data_classification:
deny:
- team: interns
classifications: [confidential, restricted]
janus project apply 创建缺失的资源、更新已有资源,但不会删除任何内容。先运行 janus project diff 预览变更。项目文件并不是唯一入口 —— 你可以随时用 janus agent add 向已存在的租户添加 Agent。该命令会注册 Agent、创建默认邮箱,并把变更持久化回 janus.project.yaml:
janus agent add summarizer --tenant acme --team engineering --capability summarize
# → Agent summarizer added to tenant acme and saved to janus.project.yaml
重复传入 --capability 可授予多个能力,用 --classification 限制每个能力的数据分类;不传 --mailbox 时自动创建 <agent-id>.default 邮箱:
janus agent add indexer --tenant acme --endpoint https://idx.internal.acme.com --capability index --capability search --concurrency 4
janus agent status indexer --tenant acme
janus agent heartbeat indexer --tenant acme
如果只想做一次性的运行时注册(不写入项目文件),使用 janus agent register;--protocol 可选 a2a、acp 或 custom-sdk(默认 a2a):
janus agent register --id adhoc-worker --name "Adhoc Worker" --protocol custom-sdk
Janus 对每一次 API 调用都使用按租户作用域的 API 密钥鉴权(默认开启)。可选地,可叠加 mTLS 用于更严格的机器到机器(M2M)部署。
当 auth.enabled 开启(默认)时,每个请求都必须通过 X-API-Key 或 Authorization: Bearer <key> 请求头携带有效密钥。密钥绑定单一租户 —— 为租户 A 签发的密钥无法访问租户 B(服务端以 403 拒绝该调用)。
用 CLI 创建密钥。原始密钥 —— janus_<64 位十六进制> —— 只显示一次,请立即保存:
janus --tenant acme api-key create --name ci-bot
# → {"tenant_id":"acme","name":"ci-bot","prefix":"janus_a1","key":"janus_a1b2c3...","created_at":"..."} (请保存 —— 只显示一次)
服务端只保存密钥的 SHA-256 哈希以及用于查询的 8 位前缀,因此原始密钥永远无法被找回。用同样的方式列出与撤销密钥:
janus --tenant acme api-key list
janus --tenant acme api-key revoke <key-id>
用 --api-key 或 JANUS_API_KEY 环境变量将密钥交给 CLI 或 SDK;直接发 HTTP 请求时以请求头携带:
export JANUS_API_KEY=janus_a1b2c3d4e5…
curl -H "Authorization: Bearer $JANUS_API_KEY" http://localhost:8080/v1/tenants/acme/tasks
对于更严格的机器到机器安全,可启用 TLS 并要求客户端证书。当设置了 client_ca_file 时,Janus 会校验客户端证书(双向 TLS,最低 TLS 1.2):
# configs/janus.local.yaml
auth:
enabled: true
tls:
enabled: true
cert_file: /etc/janus/server.crt
key_file: /etc/janus/server.key
client_ca_file: /etc/janus/ca.crt # 设置后要求客户端证书(mTLS)
或使用环境变量(均以 JANUS_ 开头):
export JANUS_AUTH_ENABLED=true
export JANUS_TLS_ENABLED=true
export JANUS_TLS_CERT_FILE=/etc/janus/server.crt
export JANUS_TLS_KEY_FILE=/etc/janus/server.key
export JANUS_TLS_CLIENT_CA_FILE=/etc/janus/ca.crt
生产环境请保持 auth.enabled 开启 —— 关闭后所有 API 端点都不再鉴权。
启用 target_type: "intent"(自然语言 → 能力路由)需要配置 OpenAI 兼容的 LLM:
# 环境变量
export JANUS_LLM_ENABLED=true
export JANUS_LLM_API_KEY=sk-xxx
export JANUS_LLM_MODEL=gpt-4o-mini
export JANUS_LLM_BASE_URL=https://api.openai.com/v1 # 或 Ollama/vLLM/Azure
未配置 LLM 时,target_type: "intent" 降级为关键词匹配(无外部依赖)。配置后,Janus 通过 LLM 将自然语言映射到最匹配的能力,经在线 agent catalog 校验后确定性路由。
pip install janus-broker
创建文件 publish.py:
from janus_broker import JanusClient
client = JanusClient("http://localhost:8080", tenant_id="acme")
# 为审核 Agent 创建邮箱
client.create_mailbox("review-mb", agent_id="reviewer")
# 向邮箱发布任务
resp = client.publish_task({
"id": "task-001",
"source_agent": "product",
"target_type": "mailbox",
"target_value": "review-mb",
"envelope": {
"type": "code_review",
"payload": {"pr_url": "https://github.com/org/repo/pull/42"},
"priority": "high"
}
})
print("Published:", resp.task_id)
python publish.py
Janus 支持五种任务投递方式:
agent-2 解析为其活跃邮箱。team: "sre",然后以 target_type: "group", target_value: "sre" 发布。JANUS_LLM_ENABLED=true)。Janus 将你的描述映射到最匹配的能力。# 示例:投递给特定 Agent
client.publish_task({
"id": "task-002", "source_agent": "product",
"target_type": "agent", "target_value": "reviewer",
"envelope": {"type": "code_review", "payload": {"pr_url": "https://github.com/org/repo/pull/42"} }
})
# 示例:按能力投递
client.publish_task({
"id": "task-003", "source_agent": "product",
"target_type": "capability", "target_value": "go-code-review",
"envelope": {"type": "go_pr_review", "payload": {"repo": "myapp"} }
})
# 示例:投递给团队
client.publish_task({
"id": "task-004", "source_agent": "product",
"target_type": "group", "target_value": "sre",
"envelope": {"type": "deploy", "payload": {"version": "2.0.0"} }
})
创建文件 worker.py:
from janus_broker import JanusClient
client = JanusClient("http://localhost:8080", tenant_id="acme")
# 从邮箱拉取任务
result = client.pull_task("review-mb", "reviewer")
print(f"Got task: {result.task.id}")
# 开始处理(获取租约)
client.start_task(result.task.id, result.lease.lease_id)
# 处理中……(实际代码里在这里做真正的处理)
# 确认完成
client.ack_task(result.task.id, {
"lease_id": result.lease.lease_id,
"result_ref": "s3://review-result.json"
})
print("Task completed")
python worker.py
使用 Janus CLI 查看任务状态:
janus task list --tenant acme
# → task-001 completed 2026-07-31T12:00:00Z
janus mailbox list --tenant acme
# → review-mb tasks: 1 pending: 0
janus CLI 是管理 Agent、任务、邮箱、API 密钥、策略与项目配置的主要工具。它通过 HTTP 与 Janus API 服务端通信 —— 用 --server 指定服务端地址。
所有命令都接受以下持久参数。可通过环境变量一次性设置,避免重复输入:
| 参数 | 默认值 | 说明 |
|---|---|---|
--server | http://localhost:8080 | Janus API 服务端地址。 |
--tenant | default | 要操作的租户 ID。 |
--api-key | $JANUS_API_KEY | 认证用 API 密钥。回退到环境变量。 |
--file | — | 项目配置文件路径。覆盖 $JANUS_PROJECT_FILE 与默认的 janus.project.yaml。 |
# 为你的 shell 会话做一次性设置
export JANUS_API_KEY="jak_..."
export JANUS_PROJECT_FILE=configs/janus.project.yaml
# 之后所有 janus 命令都会自动读取这些值
janus task list --tenant acme
| 命令 | 子命令 | 用途 |
|---|---|---|
project | init、validate、diff、apply、sync | 声明式租户拓扑管理。 |
tenant | add | 创建租户并持久化到项目文件。 |
agent | register、add、status、heartbeat | 注册 Agent 并检查存活。 |
task | publish、status、cancel、replay、events | 发布与查看任务。 |
mailbox | create、status、pause、resume、pull、ack、nack | 管理邮箱并处理任务。 |
api-key | create、list、revoke | 租户作用域 API 密钥生命周期。 |
policy | allow-agent、deny-agent、allow-team、deny-team、require-approval、allow-classification、deny-classification、allow-tool、deny-tool | 通过模板生成治理策略规则。 |
dlq | query、replay、discard | 死信队列检视与恢复。 |
dashboard | — | 交互式 TUI,监控任务、邮箱与 Agent。 |
项目命令以代码形式管理你的租户拓扑(详见配置):
janus project init # 生成 janus.project.yaml
janus project validate # 检查文件是否有错误
janus project diff # 预览与线上 API 的差异
janus project apply --all-tenants # 应用所有租户
janus project sync --overwrite # 将线上资源拉回文件
API 密钥按租户作用域划分。每个环境或服务创建一把,轮换时撤销旧密钥。原始密钥只显示一次:
janus --tenant acme api-key create --name ci-bot
# → {"tenant_id":"acme","name":"ci-bot","prefix":"janus_a1","key":"janus_a1b2c3...","created_at":"..."} (请保存 —— 只显示一次)
janus --tenant acme api-key list
janus --tenant acme api-key revoke <key-id>
发布、拉取并确认任务 —— 与 SDK 封装的同一套生命周期,只是从命令行操作:
# 向邮箱发布任务
janus task publish --tenant acme --source product --mailbox review-mb
# 查看任务状态
janus task status task-001 --tenant acme
# 从邮箱拉取下一个任务
janus mailbox pull review-mb --agent reviewer
# 检视死信队列并重放
janus dlq query --tenant acme
janus dlq replay task-001 --tenant acme
启动交互式终端仪表盘,实时监控任务、邮箱与 Agent 健康状态:
janus dashboard --server http://localhost:8080 --tenant acme
对于 Kubernetes,Janus 提供 Helm chart,内置存活/就绪探针、迁移 Job 钩子、HPA、PDB 与 Prometheus 采集注解:
helm install janus deployments/helm/janus-core/
完整的 chart values 与运维手册见 GitHub 仓库。
Janus Core 是无状态计算节点 —— 所有持久状态都外部化到 PostgreSQL、NATS JetStream 与 Redis。你可以横向运行多个副本,任意副本都能处理任意请求;单个副本崩溃不会丢失业务状态,流量在前端负载均衡后面均摊。
多副本并行消费 outbox 事件时,靠数据库级租约互斥,做到既不重复投递也不丢事件。迁移 000011_outbox_worker_lease 给 outbox_events 表加了租约字段,000012_outbox_dedupe_key 加了幂等去重键:
locked_by text -- 抢占该行的 worker
locked_at timestamptz
lease_expires_at timestamptz -- 到期后被其他 worker 回收
dedupe_key text -- (tenant_id, dedupe_key) 唯一,防重复投递
某个副本崩溃后,它持有的行在租约到期时由存活副本接管;dedupe_key 保证重试与并发调度不会重复插入同一次投递。
Helm chart 默认启用 HorizontalPodAutoscaler,按 CPU / 内存自动扩缩:
autoscaling:
enabled: true
minReplicas: 2
maxReplicas: 10
targetCPUUtilizationPercentage: 70
targetMemoryUtilizationPercentage: 80
同时配了 PodDisruptionBudget(minAvailable: 1),滚动更新与节点维护期间始终至少一个副本可用。
计算层水平扩展时,状态层各自按自己的方式扩展:
readOnlyRootFilesystem 运行,唯一持久卷只用于产物落盘(artifacts)—— 不含任何运行态,因此副本是无状态、可随时替换的。本节通过一个真实的多 Agent 协作场景——处理客户发错货投诉——展示 Janus 的路由、可靠性和治理功能如何在一个类生产工作流中协同工作。
target_type: intent 进行 LLM 驱动的自然语言路由target_type: capability)客户投诉:"我订购了红色手机,却收到了蓝色的。需要换货。"
Janus 协调四个 Agent来处理:
wrong_item_investigate,团队:logistics)reship,团队:warehouse)retrieve_wrong,团队:warehouse)notify,团队:support)
客户:"订的红色收到了蓝色"
│
▼ intent LLM → "wrong_item_investigate"
│
物流员 ─── 拉取订单 ─── 💀 崩溃(不 ACK) ─── 租约过期 ─── 恢复 ─── ACK
│
│ 数据:{order: "ORD-789", ordered: "red", shipped: "blue"}
▼
├── ► 发货员 ─── 发正确的红色 ─── ACK {tracking: "SF789"}
├── ► 回收员 ─── 取回错误的蓝色 ─── ACK {pickup: "PU123"}
└── ► 通知员 ─── "换货已启动" ─── ACK
│
▼
通知员 ─── "红色已发SF789,蓝色取件PU123" ─── ACK
核心设计:物流员的调查结果(哪个商品正确、哪个错误)驱动下游 Agent 的行为——发货员读取调查数据知道发什么;回收员读取调查数据知道回收什么。Agent 之间不是各干各的,而是通过数据流协作。
每个 Agent 声明自己能做什么(capability)和所属团队(team)。Janus 用能力做匹配,用团队做分组路由。
# 注册四个 Agent
for agent in [
{"id": "logistics-bot", "display_name": "物流员", "team": "logistics"},
{"id": "shipping-bot", "display_name": "发货员", "team": "warehouse"},
{"id": "return-bot", "display_name": "回收员", "team": "warehouse"},
{"id": "notify-bot", "display_name": "通知员", "team": "support"}
]:
client.create_agent(agent["id"], agent["display_name"], agent["team"])
然后为每个 Agent 分配能力和创建邮箱(任务队列):
# 能力:每个 Agent 能做什么
caps = {
"logistics-bot": "wrong_item_investigate",
"shipping-bot": "reship",
"return-bot": "retrieve_wrong",
"notify-bot": "notify",
}
# 创建邮箱(物流邮箱 ACK 超时设为 5 秒,用于租约过期测试)
client.create_mailbox("logistics-mb", "logistics-bot", ack_wait=5)
client.create_mailbox("shipping-mb", "shipping-bot")
client.create_mailbox("return-mb", "return-bot")
client.create_mailbox("notify-mb", "notify-bot")
不指定具体的 agent 或能力,而是用 target_type: intent 配合自然语言描述。Janus 的 LLM 解析器将请求映射到最匹配的能力:
# 客户投诉 → LLM 解析为 "wrong_item_investigate"
task = client.publish_task({
"id": "task-001",
"source_agent": "customer",
"target_type": "intent",
"target_value": "收到错误商品,订购红色手机收到蓝色,需要换货",
"envelope": {"type": "complaint", "payload": {"order": "ORD-789"}}
})
# Janus 自动解析:
# target_type: "capability"
# target_value: "wrong_item_investigate"
# → 路由到 logistics-mb(物流员的邮箱)
Janus 保证至少一次投递。当 Agent 拉取任务后崩溃而未 ACK,租约过期后任务自动返回队列:
# 物流员拉取任务……
delivery = client.pull_task("logistics-mb", "logistics-bot")
# 💀 模拟崩溃 —— 不 ACK,直接停止
# 租约过期后(配置的 5 秒),Janus 自动将任务重新投递到邮箱
# Agent 恢复后重新拉取:
delivery = client.pull_task("logistics-mb", "logistics-bot")
# ← 同一个任务!Janus 在租约过期后重新投递了它
# 正常 ACK,附带调查结果(传给下游 Agent):
client.ack_task(delivery.task_id, delivery.lease_id,
result_ref='{ "confirmed": true, "order": "ORD-789", "ordered": "red", "shipped": "blue" }')
确认发错货后,Janus 协调三个并行工作流。每个分支用不同的路由方式:
# 分支 A:发出正确商品(能力路由)
client.publish_task({
"id": "ship-001", "source_agent": "logistics-bot",
"target_type": "capability", "target_value": "reship",
"envelope": {"payload": {"order": "ORD-789", "correct": "red"}}
})
# → Janus 找到发货员(capability: reship),路由到 shipping-mb
# 分支 B:回收错误商品(能力路由)
client.publish_task({
"id": "ret-001", "source_agent": "logistics-bot",
"target_type": "capability", "target_value": "retrieve_wrong",
"envelope": {"payload": {"order": "ORD-789", "wrong": "blue"}}
})
# 分支 C:通知客户(能力路由到 support 团队)
client.publish_task({
"id": "notify-001", "source_agent": "logistics-bot",
"target_type": "capability", "target_value": "notify",
"envelope": {"payload": {"message": "换货已启动"}}
})
每个 Agent 独立拉取、处理、ACK 自己的任务,产出结果数据流向下游:
# 发货员:拉取 → 处理 → ACK(带快递单号)
delivery = client.pull_task("shipping-mb", "shipping-bot")
client.ack_task(delivery.task_id, delivery.lease_id,
result_ref='{ "tracking": "SF789", "reshipped": true }')
# 回收员:拉取 → 处理 → ACK(带取件信息)
delivery = client.pull_task("return-mb", "return-bot")
client.ack_task(delivery.task_id, delivery.lease_id,
result_ref='{ "pickup": "PU123", "scheduled": true }')
# 通知员:拉取 → ACK(初始通知)
delivery = client.pull_task("notify-mb", "notify-bot")
client.ack_task(delivery.task_id, delivery.lease_id)
发货和回收都完成后,发布最终通知,合并两路结果——展示了 Agent 之间如何通过 Janus 传递和组合数据:
# 最终通知合并发货 + 回收结果
client.publish_task({
"id": "notify-final", "source_agent": "shipping-bot",
"target_type": "capability", "target_value": "notify",
"envelope": {"payload": {
"shipping": {"tracking": "SF789"},
"return": {"pickup": "PU123"},
"message": "红色手机已发出(SF789),蓝色手机已安排取件(PU123)"
}})
JANUS_AUTH_ENABLED=false——在将 Janus 暴露到网络前,请重新开启(并通过 janus api-key create 创建 API 密钥)。