Two agents, two companies, one protocol: building a visible agent-to-agent negotiation
How consortium works: Microsoft Agent Framework 1.x, A2A for the negotiation, MCP for the work, and a data plane that says no. There is a particular kind of agent demo I find unsatisfying. Five agents sit in a group c
How consortium works: Microsoft Agent Framework 1.x, A2A for the negotiation, MCP for the work,
and a data plane that says no.
There is a particular kind of agent demo I find unsatisfying. Five agents sit in a group chat, a framework orchestrates their turn-taking, and the "negotiation" is a string of tokens inside one process's memory. You can print it, but you cannot point at a protocol. Nothing crossed a wire.
So I built the version I wanted to look at: two agents that belong to two different companies, talking over a real protocol, where the thing that actually enforces the agreement is not a language model.
BuyerCorp's agent wants one number: average order value by region for Q3. It does not have the data. VendorCorp's agent has the data and will not release anything until they agree which columns, for what purpose, for how long. They negotiate that contract over A2A. When the deal is struck, the query runs through MCP tools β and the MCP server refuses calls that step outside the contract, in code, before reading a single row.
Everything is visible in a Streamlit observer, one message at a time.
It is ~2,900 lines of Python across four processes and a shared core/ package. This is how it is built, and the eight things that cost me real time.
The shape of it
| Process | Port | What it owns |
|---|---|---|
mcp_server.py |
8123 | the CSV, the audit log, the scope check |
vendor_agent.py |
8101 | VendorCorp's agent, the contract grant |
buyer_agent.py |
8102 | BuyerCorp's agent, the negotiation driver |
observer.py |
8501 | rendering. Writes nothing but config |
Two design choices do most of the work:
The control plane and the data plane are different protocols. A2A carries the argument. MCP carries the data. The bridge between them is a file: when VendorCorp accepts, it publishes the agreed scope to data/contract.json, which the MCP server re-reads on every call.
The buyer talks to the vendor's MCP server directly. That is what makes the refusal interesting. If the vendor's own agent refused its own tool call, the enforcement would be a courtesy. Instead an independent process β one the buyer cannot negotiate with β evaluates the request against the contract and says no.
Step 1: the negotiation is a schema, not a vibe
Every message is one JSON object:
{
"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"
}
The rule I set myself: reject anything that does not parse, retry once, then fail the run. No "best effort" path, because a best-effort negotiation is how you get a contract nobody agreed to.
The validator is deliberately picky:
# 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)}")
Three details worth calling out:
-
expect_fromstops an agent answering on the other side's behalf. If VendorCorp's reply comes back with"from": "buyer", that is a protocol violation, not a cosmetic slip. -
Unknown columns are rejected against the peer's published schema. Where does the buyer learn the column names? From the vendor's agent card. A model that invents
avg_order_valuenever reaches the data plane. -
Extra keys are rejected. Tempting to add
"urgency": "high", and exactly the kind of drift a schema exists to stop.
The retry lives in the model call, not in the 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
Step 2: serving an agent as an A2A endpoint
Microsoft Agent Framework ships the bridge in agent-framework-a2a: A2AExecutor adapts an agent to the SDK's AgentExecutor, and a2a-sdk supplies the routes.
# 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 only needs something that satisfies SupportsAgentRun β run() plus create_session(). That is the seam I used to put the protocol logic in one place: a VendorService class that validates the inbound message, decides whether the turn needs the model at all, and logs
both sides of the wire.
The agent card is also how the vendor publishes its negotiable schema and its data plane:
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
)
At startup the vendor connects to its own MCP server, calls tools/list, and embeds the real input schemas in every grant it issues. So the buyer learns the tool signatures from the same place it learns the columns.
Step 3: calling a peer, and the context_id question
A negotiation has memory. VendorCorp must remember what it offered, or an acceptance means nothing. The natural key is the A2A context_id β but does a client-supplied context_id survive across separate message/send requests?
I tested it before building on it, and it does:
session = AgentSession(service_session_id=A2AServiceSessionId(context_id=conversation))
final = await peer.run(body, session=session)
The vendor's A2AExecutor calls create_session(session_id=task.context_id), and task.context_id
came back as exactly the string the buyer sent (conv-A, conv-A, conv-B in the probe). That is the whole basis for the vendor's per-conversation state:
self._offers: dict[str, dict[str, Any]] = {} # context_id β last scope VendorCorp offered
self._transcripts: dict[str, list[dict[str, Any]]] = {}
Fetching the card is one line with a caveat:
from a2a.client.card_resolver import parse_agent_card
card = parse_agent_card(response.json())
parse_agent_card exists because the served card carries legacy fields (preferredTransport,
url, protocolVersion) that the protobuf AgentCard does not have. json_format.ParseDict
without ignore_unknown_fields=True throws on them. Use the SDK's helper.
Step 4: the data plane is an MCP server, not a function
mcp_server.py is a standalone process on streamable HTTP, so that "kill the MCP server mid-run" is a thing you can actually do to it.
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)
...
On the caller side, ToolError becomes isError on the wire, which Agent Framework's MCP client raises as ToolExecutionException with the server's text intact:
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)
That distinction is the demo. A ToolExecutionException is a policy decision from the data owner. Anything else β ConnectionError, MCP server failed to initialize β is a transport failure, and the observer renders them differently.
Step 5: where enforcement lives
core/contract.py holds the one function that decides anything real:
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
Four refusal codes, deliberately distinct, because "the column does not exist" and "you were not granted the column" are different incidents.
And two rules that are code, not prompt β because a prompt is a request and code is a guarantee:
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
If the model over-grants, the grant is narrower than the model's reply, and the log says so.
The second rule closes the hole I thought about hardest. What stops a buyer sending an accept whose scope quietly includes customer_email? Nothing in the message format. So VendorCorp honours an acceptance against its own last offer, never against what the buyer claims to accept:
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."}
The grant step also skips the model entirely. Once both sides agree, publishing a contract is a deterministic act; routing it through an LLM adds a chance of hallucinating terms into a document that is about to become an access-control policy.
Step 6: the observer is a log reader
Four processes append to one JSONL file. Cross-process ordering is the only hard part: take an exclusive fcntl lock, read the tail to find the last seq, append, 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())
The UI rule I cared about: never render a message that did not happen. No seeded timeline, no placeholder exchange. An empty log draws an empty state that tells you what to start.
State is derived, not stored:
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"
A stored flag would lie the moment a process dies mid-run. Derived state cannot.
Same instinct for the health dots in the sidebar: a 200 from a port proves nothing about which process answered, so each check requires an identity field:
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"β¦
The MCP row goes further: it is reported up only when the vendor successfully discovered tools from it at startup, which is evidence the data plane answered a real tools/list.
Eight things that cost time
-
A CSV reader that silently produced garbage.
[{field: _coerce(raw) for field, raw in zip(fields, record)}]whererecordis aDictReaderrow β zipping a list against a dict iterates its keys, so every cell became its own column name and every mean wasnull. No exception anywhere. Found by printing one result. Correct form:{field: _coerce(record[field]) for field in fields}. -
httpx awaits response hooks. The shim for the provider quirk below was written as a plain function, which produced
TypeError: object NoneType can't be used in 'await' expressionwrapped in a framework exception that did not mention hooks. It must beasync defwithawait response.aread(). -
api.particle.aireturnsmetadataas an object on structured replies while the OpenAI SDK types itstr | None, so valid responses fail client-side validation. Fixed with the hook from (2): null the field, rewrite the body, never touch a non-JSON (SSE) response. -
deepseek-v4.1-flashspends the whole completion budget on hidden reasoning and returnscontent: "".reasoning_effort: "none"in the request body fixes it; without it, every agent turn looks like an empty model. -
Streaming A2A returned nothing.
A2AExecutor(stream=False)puts its reply in a task artifact; askingA2AAgent.run(stream=True)for it yielded zero updates and an empty final response on these versions. Non-streaming works, and one request/response per negotiation message is a more honest trace anyway. Card advertisesstreaming=Falseto match. -
API surface drift in Agent Framework 1.x.
MCPClientis gone β it isMCPStdioTool/MCPStreamableHTTPToolnow,tool.functionsis a list not a dict,call_tool()returnslist[Content]rather than a string, andToolExecutionExceptionlives inagent_framework.exceptions, not the package root.A2AContinuationTokenis not re-exported through theagent_framework.a2alazy shim either. -
st.fragmentin Streamlit 1.65 returns a plain function β no.run(), so dynamicrun_everyis not available the way the docs' examples imply. Fixed the interval instead. Alsouse_container_widthis deprecated in favour ofwidth='stretch'. -
My own
run.shlied. It reported "observer up" because the port was open β held by a stale instance it did not own, while its own process had exited with "Port 8501 is not available". Now it refuses to start a service whose port is held by someone else and checks the child it spawned is still alive.
What it actually produces
Negotiation, from the captured event log:
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
The vendor dropped order_id; the buyer accepted the narrower scope; the contract carries a 30-day expiry and the A2A context_id it was negotiated under. The answer is 202 of 624 rows:
| region | n_rows | mean(order_value) |
|---|---|---|
| APAC | 22 | 138.3195 |
| EMEA | 44 | 173.9941 |
| LATAM | 27 | 162.1626 |
| NA-East | 50 | 188.0616 |
| NA-West | 59 | 224.1207 |
Same question, but naming the email column, and the data plane answers:
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.
Kill the MCP server mid-run and the run ends instead of hanging:
{"plane": "mcp", "actor": "buyercorp", "kind": "connection-error",
"error": "ToolException: MCP server failed to initialize: [Errno 61] Connection refused"}
One honest caveat. Those runs were driven through scripts/mock_provider.py, a scripted OpenAI-compatible endpoint, because no provider key was configured when I captured them. The aggregates and every refusal string are model-independent β the MCP server computes them from the committed CSV, and tests/test_core.py recomputes the means with statistics.fmean and asserts equality. What a real model changes is the wording and column choices inside the negotiation messages. The harness is a plumbing test, and I have labelled it as one rather than letting it pass for a demo.
Try it, or read the interesting files
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
If you want to read three files rather than eight:
-
core/contract.pyβ the entire policy model in ~120 lines. -
mcp_server.pyβ where a refusal is born. -
vendor_agent.pyβ the part where an agent decides not to call its model.
Code & more: https://www.dailybuild.xyz/project/279-consortium
Originally published by Dev.to AI. Aggregated on AIWithGhost for educational purposes β full credit and traffic to the original publisher.