项目让买卖双方 Agent 通过 A2A 协商数据使用范围,再通过 MCP 工具执行查询。MCP 服务在读取数据前以代码拒绝越权调用,Streamlit 界面逐条展示协商过程。
consortium 的工作原理:Microsoft Agent Framework 1.x、用于协商的 A2A、用于执行工作的 MCP,以及一个会明确拒绝请求的数据平面。
有一类智能体演示总让我觉得不够满意。五个智能体坐在一个群聊里,由框架安排它们轮流发言,所谓的“协商”不过是某个进程内存中的一串 token。你可以把它打印出来,却无法指出它使用了什么协议。没有任何消息真正经过网络传输。
于是,我构建了自己想看到的版本:两个分属不同公司的智能体,通过真实的协议交流,而真正负责执行协议约定的并不是语言模型。
BuyerCorp 的智能体只想得到一项指标:第三季度各地区的平均订单金额。它没有数据。VendorCorp 的智能体拥有数据,但在双方就使用哪些列、出于什么目的、保留多长时间达成一致之前,它不会放出任何数据。双方通过 A2A 协商这份合约。达成协议后,查询通过 MCP 工具执行——而 MCP 服务器会在读取任何一行数据之前,通过代码拒绝超出合约范围的调用。
整个过程都可以在 Streamlit 观察界面中逐条查看。
这个项目约有 2,900 行 Python 代码,分布在四个进程和一个共享的 core/ 包中。下面介绍它的构建方式,以及八个真正耗费了我时间的问题。

两个设计选择承担了大部分工作:
控制平面和数据平面使用不同的协议。A2A 负责传递协商内容,MCP 负责传递数据。连接两者的是一个文件:VendorCorp 接受提案时,会将双方约定的范围写入 data/contract.json,MCP 服务器在每次调用时都会重新读取这个文件。
买方直接与卖方的 MCP 服务器通信。这正是拒绝机制有意思的地方。如果只是卖方自己的智能体拒绝调用自己的工具,那么这种约束不过是出于自觉。这里则由一个独立进程——买方无法与之协商的进程——根据合约评估请求,并明确拒绝。
每条消息都是一个 JSON 对象:
{
"from": "buyer|vendor",
"intent": "propose|counter|accept|reject|request",
"scope": {
"columns": ["region", "order_value", "order_date"],
"purpose": "average order value by region for Q3 2025",
"retention_days": 30
},
"rationale": "one sentence"
}
我给自己定下的规则是:拒绝任何无法解析的内容,重试一次,再失败就终止本次运行。不提供“尽力而为”的处理路径,因为尽力而为的协商,最终可能产生一份谁都没有同意过的合约。
校验器刻意设计得很严格:
# core/negotiation.py
TOP_LEVEL = {"from", "intent", "scope", "rationale"}
SCOPE_KEYS = {"columns", "purpose", "retention_days"}
extra = sorted(set(obj) - TOP_LEVEL)
if extra:
raise SchemaError(f"unexpected top-level key(s) {extra}; allowed: {sorted(TOP_LEVEL)}")
if expect_from and side != expect_from:
raise SchemaError(f'expected "from": "{expect_from}", got "{side}" '
f"(a reply may not speak for the other party)")
if known_columns is not None:
unknown = sorted(set(columns) - known_columns)
if unknown:
raise SchemaError(f"column(s) {unknown} are not in the dataset schema {sorted(known_columns)}")
有三个细节值得指出:
expect_from 防止智能体代替另一方作答。如果 VendorCorp 的回复中出现了 "from": "buyer",那就是违反协议,而不是无关紧要的小失误。
校验器会根据对方公布的 schema 拒绝未知列。买方从哪里得知列名?从卖方的智能体卡片中。模型即使凭空编造出 avg_order_value,也根本到不了数据平面。
额外的键也会被拒绝。添加 "urgency": "high" 可能很诱人,但 schema 的存在,正是为了阻止这类偏离。
重试逻辑位于模型调用中,而不是 UI 中:
# core/model.py — one round trip that must yield a parseable JSON object
try:
parsed = extract_json(raw)
except ValueError as exc:
if attempts > 1:
emit("model", side, "parse-failed", turn=turn, attempts=attempts, error=str(exc), raw=raw[:800])
raise ProtocolError(f"{side} reply was not valid JSON after one retry: {exc}")
text = prompt + "\n\nYour previous reply could not be parsed as a single JSON object" \
f" ({exc}). Reply again with ONLY the JSON object: no prose, no code fence."
continue
Microsoft Agent Framework 在 agent-framework-a2a 中提供了这层桥接:A2AExecutor 将智能体适配为 SDK 的 AgentExecutor,a2a-sdk 则提供路由。
# vendor_agent.py
from a2a.server.request_handlers import DefaultRequestHandler
from a2a.server.routes import add_a2a_routes_to_fastapi, create_agent_card_routes, create_jsonrpc_routes
from a2a.server.tasks import InMemoryTaskStore
from agent_framework.a2a import A2AExecutor
card = AgentCard(
name="VendorCorp Data Steward",
version="1.0.0",
provider={"url": url, "organization": "VendorCorp"},
default_input_modes=["text"],
default_output_modes=["text"],
capabilities=AgentCapabilities(streaming=False),
supported_interfaces=[AgentInterface(url=url, protocol_binding="JSONRPC")],
skills=[...],
)
handler = DefaultRequestHandler(
agent_executor=A2AExecutor(service, stream=False),
task_store=InMemoryTaskStore(),
agent_card=card,
)
app = FastAPI()
add_a2a_routes_to_fastapi(
app,
agent_card_routes=create_agent_card_routes(card),
jsonrpc_routes=create_jsonrpc_routes(handler, rpc_url="/"),
)
A2AExecutor 只需要一个满足 SupportsAgentRun 的对象,也就是实现 run() 和 create_session()。我利用这个接入点,将协议逻辑集中在一个地方:VendorService 类负责校验传入的消息,判断这一轮是否真的需要调用模型,并记录通信双方的消息。
卖方也通过智能体卡片公布可协商的 schema 及其数据平面:
AgentSkill(
id="vendor-dataset-schema",
name="Negotiable dataset schema",
description="Column names VendorCorp will negotiate access to. Names only: no rows and no "
"values are exposed before a contract is granted.",
tags=columns(), # ← the schema the buyer validates proposals against
)
启动时,卖方会连接自己的 MCP 服务器,调用 tools/list,并在每次授予访问权限时附上真实的输入 schema。因此,买方从同一个地方获取列名和工具签名。
协商需要记忆。VendorCorp 必须记住自己提出过什么条件,否则接受提案就毫无意义。最自然的键是 A2A 的 context_id——但客户端提供的 context_id,能否在多次独立的 message/send 请求之间保持不变?
在以此为基础继续构建之前,我先做了测试,答案是可以:
session = AgentSession(service_session_id=A2AServiceSessionId(context_id=conversation))
final = await peer.run(body, session=session)
卖方的 A2AExecutor 会调用 create_session(session_id=task.context_id),而返回的 task.context_id 与买方发送的字符串完全一致(测试中依次为 conv-A、conv-A、conv-B)。这就是卖方按会话保存状态的全部基础:
self._offers: dict[str, dict[str, Any]] = {} # context_id → last scope VendorCorp offered
self._transcripts: dict[str, list[dict[str, Any]]] = {}
获取卡片只需一行代码,但有一个注意事项:
from a2a.client.card_resolver import parse_agent_card
card = parse_agent_card(response.json())
之所以提供 parse_agent_card,是因为服务端返回的卡片包含一些旧版字段(preferredTransport、url、protocolVersion),而 protobuf 的 AgentCard 中没有这些字段。调用 json_format.ParseDict 时,如果没有设置 ignore_unknown_fields=True,就会因这些字段抛出异常。请使用 SDK 提供的辅助函数。
mcp_server.py 是一个通过 Streamable HTTP 提供服务的独立进程,因此你确实可以在运行过程中终止 MCP 服务器。
from mcp.server.fastmcp import FastMCP, ToolError
mcp = FastMCP("vendorcorp-data", host="127.0.0.1", port=8123,
streamable_http_path="/mcp", stateless_http=True, json_response=True)
@mcp.tool(name="query_records", description="Aggregate VendorCorp's order dataset. …")
async def query_records(filters: dict[str, Any], group_by: str | None,
aggregate: dict[str, Any]) -> dict[str, Any]:
dataset_columns = set(columns())
referenced = _referenced_columns(filters, group_by, aggregate)
allowed, reason, contract = evaluate(referenced, dataset_columns)
if not allowed:
_record_call("query_records", arguments, decision="refusal", reason=reason, contract=contract)
raise ToolError(reason) # ← the refusal originates HERE
result = run_query(filters, group_by or None, aggregate)
...
在调用方,ToolError 会在协议传输中转换为 isError,随后 Agent Framework 的 MCP 客户端会抛出 ToolExecutionException,并完整保留服务器返回的文本:
try:
out = await tool.call_tool("query_records", **args)
except ToolExecutionException as exc:
reason = str(exc) # the vendor's words, not ours
emit("mcp", "buyercorp", "refusal", tool="query_records", arguments=args, reason=reason)
这个演示要展示的,正是这种区别。ToolExecutionException 表示数据所有者作出的策略决定。其他异常——ConnectionError、MCP server failed to initialize——则属于传输故障,观察界面会以不同方式呈现它们。
core/contract.py 中只有一个函数真正作出决定:
def evaluate(referenced, dataset_columns, *, now=None) -> tuple[bool, str, dict | None]:
unknown = sorted(referenced - dataset_columns)
if unknown:
return False, f"UNKNOWN_COLUMN: {unknown} is not a column of the VendorCorp dataset…", contract
if contract is None:
return False, "NO_CONTRACT: no negotiated contract has been published…", contract
if _expired(contract, now):
return False, f"CONTRACT_EXPIRED: {contract['contract_id']} expired at {contract['expires_at']}…", contract
denied = sorted(referenced - granted)
if denied:
return False, (f"SCOPE_DENIED: {denied} was not granted by contract {contract['contract_id']}. "
f"Granted columns: {sorted(granted)}. Purpose of record: "
f"{contract['scope']['purpose']!r}. This call is refused at the MCP server "
f"and was not read from the dataset."), contract
return True, "", contract
这里刻意区分了四种拒绝码,因为“这一列不存在”和“你没有获得这一列的访问权限”是两种不同的情况。
另外还有两条规则,写在代码里,而不是提示词里——因为提示词只是请求,代码才提供保证:
PII_COLUMNS = {"customer_name", "customer_email"}
def _apply_policy(self, message):
kept = [c for c in message["scope"]["columns"] if c not in PII_COLUMNS]
removed = sorted(set(message["scope"]["columns"]) - set(kept))
if removed:
emit("contract", "vendorcorp", "policy-filter", removed=removed,
reason="PII_COLUMNS is enforced in vendor_agent.py after the model replies")
message["intent"] = "counter"
message["scope"]["retention_days"] = min(retention, MAX_RETENTION_DAYS) # 90-day cap
如果模型授予的权限过多,实际授权范围就会比模型回复中的范围更窄,日志也会明确记录这一点。
第二条规则堵住了我花最多心思考虑的漏洞。如果买方发送一条 accept 消息,却在范围里悄悄加入 customer_email,有什么能阻止它?消息格式本身没有任何限制。因此,VendorCorp 只依据自己最后一次报价来确认接受,而不会依据买方声称接受的内容:
buyer_cols, offered_cols = set(accepted["scope"]["columns"]), set(offered["columns"])
widening = sorted(buyer_cols - offered_cols)
if widening:
return {"from": "vendor", "intent": "reject", "rationale":
f"An accept may only confirm the scope VendorCorp offered; {widening} was never offered, "
"so this reads as a new proposal rather than an acceptance."}
授权步骤也完全绕过模型。一旦双方达成一致,发布合约就是一个确定性的操作;让它经过 LLM,只会增加模型往文档里凭空添加条款的风险,而这份文档马上就要成为访问控制策略。
四个进程向同一个 JSONL 文件追加内容。唯一棘手的部分是跨进程排序:获取 fcntl 排他锁,读取文件尾部以找到最后一个 seq,然后追加内容并执行 fsync。
with _lock, open(LOG_PATH, "a+", encoding="utf-8") as handle:
fcntl.flock(handle.fileno(), fcntl.LOCK_EX)
event["seq"] = _next_seq(handle) # parse the last line of the tail
handle.write(json.dumps(event) + "\n")
handle.flush(); os.fsync(handle.fileno())
我最在意的一条 UI 规则是:绝不显示没有发生过的消息。不预填时间线,也不放占位对话。日志为空时,界面就显示空状态,告诉你应该启动什么。
状态是推导出来的,而不是存储的:
def derive_state(events):
if not events: return "idle", "No event has been written yet…"
starts = [i for i, e in enumerate(events) if e["plane"] == "run" and e.get("kind") == "start"]
if not starts: return "idle", "…no negotiation has been requested."
tail = events[starts[-1]:]
if any(e.get("kind") in ("error", "transport-error", "connection-error") for e in tail):
return "error", …
if any(e["plane"] == "mcp" and e.get("kind") in ("call", "connect", "result", "refusal") for e in tail):
return "executing", "Calling the VendorCorp MCP data plane"
return "negotiating", "A2A messages are exchanging between the two agents"
如果把状态存成标志位,进程在运行中途退出的那一刻,这个标志位就会失真。推导出来的状态不会。
侧边栏的健康状态指示灯也遵循同样的思路:端口返回 200,并不能证明究竟是哪个进程作出了响应,所以每次检查都必须验证一个身份字段:
if isinstance(body, dict) and body.get("service") == "consortium-vendor-agent":
…report up…
elif body and body.get("service"):
…report "port answered by X — not this build"…
MCP 那一行还更进一步:只有卖方在启动时成功从 MCP 服务发现了工具,才会显示服务正常,因为这能证明数据平面确实响应过一次真实的 tools/list 请求。
一个悄无声息地产生垃圾数据的 CSV 读取器。[{field: _coerce(raw) for field, raw in zip(fields, record)}] 中,record 是 DictReader 返回的一行——将列表与字典传给 zip 时,遍历的是字典的键,因此每个单元格都变成了自己的列名,所有均值都成了 null。整个过程没有抛出任何异常。打印一条结果后才发现问题。正确写法是:{field: _coerce(record[field]) for field in fields}。
httpx 会 await 响应钩子。为下文提到的服务商特殊行为编写的适配函数,最初用了普通函数,结果触发了 TypeError: object NoneType can't be used in 'await' expression,外层又被框架异常包装,完全没有提到钩子。它必须写成 async def,并调用 await response.aread()。
api.particle.ai 在结构化回复中将 metadata 返回为对象,而 OpenAI SDK 将其类型定义为 str | None,因此有效响应也会在客户端校验时失败。用第(2)项中的钩子解决:将该字段设为 null,重写响应体,绝不处理非 JSON(SSE)响应。
deepseek-v4.1-flash 把整个补全预算都花在隐藏推理上,返回 content: ""。在请求体中加入 reasoning_effort: "none" 就能解决;如果不加,智能体每一轮都会像是收到了模型的空回复。
流式 A2A 没有返回任何内容。A2AExecutor(stream=False) 会将回复放进任务产物中;在这些版本里,通过 A2AAgent.run(stream=True) 获取它,收到的更新数量为零,最终响应也是空的。非流式调用可以正常工作,而且每条协商消息对应一次请求和响应,本来就能更如实地呈现交互轨迹。Card 也声明 streaming=False,与实际行为保持一致。
Agent Framework 1.x 的 API 接口发生了变化。MCPClient 已经移除——现在用的是 MCPStdioTool / MCPStreamableHTTPTool;tool.functions 是列表而不是字典;call_tool() 返回 list[Content] 而不是字符串;ToolExecutionException 位于 agent_framework.exceptions,而不是包的根模块。agent_framework.a2a 的延迟加载兼容层也没有重新导出 A2AContinuationToken。
Streamlit 1.65 中的 st.fragment 返回的是普通函数——没有 .run(),因此无法按文档示例所暗示的方式动态设置 run_every。于是改用了固定间隔。另外,use_container_width 已弃用,应改用 width='stretch'。
我自己的 run.sh 也报了假状态。它因为端口已打开就报告“observer up”——实际上,占用端口的是一个不归它管理的旧实例,而它自己启动的进程已经报出 "Port 8501 is not available" 并退出。现在,如果服务端口被其他进程占用,它会拒绝启动该服务,并检查自己创建的子进程是否仍然存活。
以下是从捕获的事件日志中提取的协商过程:
7 a2a buyercorp send buyer:propose ['region','order_date','order_value','order_id']
8 a2a vendorcorp inbound buyer:propose
10 a2a vendorcorp outbound vendor:counter ['order_date','order_value','region']
14 a2a buyercorp send buyer:accept ['order_date','order_value','region']
17 a2a vendorcorp outbound GRANT c-020682
22 mcp buyercorp call query_records
23 mcp vendorcorp result query_records
27 run buyercorp answer
卖方去掉了 order_id;买方接受了更窄的范围;合约的有效期为 30 天,并记录了协商时使用的 A2A context_id。回答涉及 624 行中的 202 行:
还是同一个问题,但这次指明了邮箱列,数据平面给出的响应是:
SCOPE_DENIED: ['customer_email'] was not granted by contract c-61cc5c. Granted columns:
['order_date', 'order_value', 'region']. Purpose of record: 'aggregate order value by region for
Q3 2025'. This call is refused at the MCP server and was not read from the dataset.
在运行中途终止 MCP 服务,运行会随之结束,而不会一直挂起:
{"plane": "mcp", "actor": "buyercorp", "kind": "connection-error",
"error": "ToolException: MCP server failed to initialize: [Errno 61] Connection refused"}
有一点需要坦诚说明:这些运行是通过 scripts/mock_provider.py 驱动的,这是一个按预设脚本响应、兼容 OpenAI 的端点,因为我记录这些运行时并未配置模型提供商的 API 密钥。聚合结果和每一条拒绝响应字符串都与模型无关——MCP 服务器根据仓库中已提交的 CSV 计算它们,tests/test_core.py 则使用 statistics.fmean 重新计算均值,并断言两者相等。真实模型会改变的是协商消息中的措辞和列的选择。这套测试框架用于验证各组件之间的连接是否正常,我也明确将其标注为此类测试,而没有把它包装成演示。
试着运行,或阅读值得关注的文件
git clone https://github.com/harishkotra/consortium
cd consortium && uv venv --python 3.12 .venv
uv pip install -p .venv/bin/python -r requirements.txt
cp .env.example .env # put a key in, or point at scripts/mock_provider.py
./run.sh start # http://127.0.0.1:8501
如果你只想读三个文件,而不是八个:
core/contract.py —— 用约 120 行代码实现的完整策略模型。
mcp_server.py —— 拒绝响应产生的地方。
vendor_agent.py —— 智能体决定不调用其模型的那部分代码。
代码及更多内容:https://www.dailybuild.xyz/project/279-consortium
部分评论可能仅对已登录的访客可见。登录后即可查看所有评论。
如需采取进一步措施,你可以考虑屏蔽此人和/或举报滥用行为。