0
0

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?

A2Aエージェント同士をつなぐと何が起きるのか? 【A2A深掘りシリーズ 3】

0
Last updated at Posted at 2026-09-06

はじめに

AIエージェント同士をつなぐオープン標準 A2A(Agent2Agent)プロトコル を、SDKを使わない生実装と実測ログで調べていくシリーズの第3回です。

ここまでの2回では、エージェントに話しかけるクライアントは人間(curl)でした。しかしA2Aは「Agent2Agent」、つまりエージェント同士の通信のためのプロトコルです。今回はついにエージェントを2体つなぎ、「エージェントがエージェントに依頼する」を実際に動かします。

本記事のゴールは、次の3点を実感を持って言えるようになることです。

  • エージェント同士の連携とは、片方のエージェントが、もう片方へのクライアントになることである
  • それはA2Aの基本機能(Agent Cardによる発見と、SendMessageによる依頼)だけで実現できる
  • 2つの会話(利用者⇔受付、受付⇔ワーカー)をつなぐのはプロトコルではなく、間に立つエージェントの実装である

検証環境: macOS / Node.js v25.6.1 / A2Aプロトコル v1.0(仕様。本文とコード中の§番号はこの仕様の節番号です)

前提: 第2回で作成した task-server.mjs を手元に用意してください(本記事では改変せずそのまま使います)。

登場人物: 受付エージェントとワーカー

今回の構成はこうです。

  • ワーカー: 第2回で作った Slow Calc Agent(足し算に約3秒かかる)を、一切改変せずそのまま使います
  • 受付エージェント: 今回新しく作ります。計算依頼を受け付けますが、自分では計算せず、ワーカーへA2Aで取り次ぎます

A2Aでエージェント同士を連携させるとき、間に立つエージェントは2つの役割を同時に担います。利用者に対してはA2Aサーバとして振る舞い、ワーカーに対してはA2Aクライアントとして振る舞います。「エージェント同士の連携」とは、このサーバ役とクライアント役を兼ねるエージェントが間に入ることであって、プロトコルに連携専用のメソッドが増えるわけではありません。

受付エージェントを作る

受付エージェントの実装は、次の4つの部品でできています。

  • ワーカーの発見: getWorkerEndpoint() がワーカーのAgent Cardを取得し、エンドポイントを知る。第1回で人間がcurlでたどった「発見 → 接続」の2段階を、受付が自動でなぞります
  • ワーカーの呼び出し: callWorker() がNode標準の fetch でJSON-RPCを送る。クライアント側も依存ゼロで書けます
  • ワーカーの状態の反映: settleFromWorker() が、ワーカーから返ってきたタスクの状態(COMPLETED / INPUT_REQUIRED / REJECTED)を読み、受付自身のタスクをどの状態に進めるかを決める。ワーカーの状態をそのままコピーするのではなく、「ワーカーがこうなったから、受付としてはこうする」と受付側で決め直す処理で、その対応関係(REJECTEDを受付でもREJECTEDにするか、別のワーカーに再依頼するか等)は受付の実装判断。実運用のオーケストレーターのように配下に複数のワーカーを持つ構成では、この決め直しが連携の中核になる
  • 自分のタスク管理: 利用者からの依頼ごとに受付自身のTaskを作り(taskIdの採番、状態の遷移、履歴と成果物の保持)、利用者からのGetTaskに答え、ブロッキング(blocking)のSendMessageには自タスクが終端か中断に達するまで応答を待たせる。利用者に対するA2Aサーバとしての責務はここが担う。実装は第2回のサーバの状態管理部分をそのまま流用し、ストリーミング(SSE)だけ省いている(実装がSSEに対応しないので、受付のAgent Cardの capabilities.streamingfalse と宣言している)

受付エージェントのソースコード全文(依存ゼロ・約300行)は、本文に載せるには長いので巻末の付録に掲載しました。付録のソースコードを、第2回task-server.mjs と同じディレクトリに reception-server.mjs として保存してください。

通信ログ(wire log)は受付側が reception-wire.log に「利用者⇔受付」と「受付⇔ワーカー」の両方向を記録します。ワーカー側は従来どおり wire.log に記録するので、同じ1往復が、立場の違う2つのログに現れることになります。

受付に依頼を取り次がせる

ワーカー(第2回task-server.mjs)と受付(reception-server.mjs)の2つのサーバを、それぞれ別のターミナルで起動します。

node task-server.mjs        # ワーカー(第2回のまま。ポート4102)
node reception-server.mjs   # 受付(ポート4103)

curlで通信の中身を確認していきます。手順は次の6ステップです。

  1. 2枚のAgent Cardを見比べる
  2. 利用者から受付に依頼する(ブロッキング)
  3. 通信ログで「受付が送ったMessage」を見る
  4. taskIdとcontextIdの独立を確認する
  5. 中断(INPUT_REQUIRED)を伝搬させる
  6. ワーカー側の履歴を確認する

後続のcurlコマンドで繰り返し使う値(受付とワーカーのURL、共通のHTTPヘッダ)を、あらかじめシェル変数に入れておきます。

R=http://localhost:4103
W=http://localhost:4102
H1='Content-Type: application/json'
H2='A2A-Version: 1.0'

ステップ1: 2枚のAgent Cardを見比べる

curl -s $R/.well-known/agent-card.json
curl -s $W/.well-known/agent-card.json

capabilitiesの部分だけ抜粋して並べます。

Calc Reception Agent: { "streaming": false, "pushNotifications": false, ... }
Slow Calc Agent:      { "streaming": true,  "pushNotifications": false, ... }

エージェントごとに対応能力が違う、というだけの当たり前の光景ですが、クライアントは接続前にこの差を知って振る舞いを変えられます(受付はSSE非対応なので、利用者が結果を待たずに戻る呼び方、つまり returnImmediately を選ぶ場合は、SSEではなくポーリングで確認することになります)。

ステップ2: 利用者から受付に依頼する(ブロッキング)

[利用者 → 受付] に、第2回と全く同じ形の依頼を送ります。裏で何が起きるかは、次の図のとおりです。

time curl -s -X POST $R/a2a/v1 -H "$H1" -H "$H2" \
  -d '{"jsonrpc":"2.0","id":1,"method":"SendMessage","params":{"message":{"messageId":"u-1","role":"ROLE_USER","parts":[{"text":"{\"a\": 19, \"b\": 23}"}]}}}'

約3.9秒後(実測3.862秒)、COMPLETEDのTaskが返ります。レスポンスの注目部分を抜粋します。

{
  "result": {
    "task": {
      "id": "d165d4d3-f9a1-4610-9983-ea26fd7e9a87",
      "contextId": "c646dec3-b63a-474c-b0d7-8662a2d4ceaa",
      "status": { "state": "TASK_STATE_COMPLETED" },
      "artifacts": [
        {
          "name": "計算結果",
          "parts": [{ "text": "19 + 23 = 42" }],
          "metadata": {
            "source": "Slow Calc Agent",
            "workerTaskId": "3833ec71-1b0f-4397-956b-7cc0fcdbb404"
          }
        }
      ],
      "metadata": { "workerTaskId": "3833ec71-1b0f-4397-956b-7cc0fcdbb404" },
      "history": [ (省略: 利用者の依頼と、受付の実況Messageが並ぶ) ]
    }
  }
}

利用者から見えるのは受付のTaskだけです。ワーカーの存在は、受付が metadata に書いた workerTaskId を覗かない限り見えません(このmetadataへの記載はプロトコルの決まりではなく、本実装が学習用に「あえて公開する」と決めた実装判断です)。

ステップ3: 通信ログで「受付が送ったMessage」を見る

ここが本記事のハイライトです。受付がワーカーに何を送ったのか、reception-wire.log の [受付 → ワーカー] エントリを見てみます。

{
  "jsonrpc": "2.0",
  "id": 1,
  "method": "SendMessage",
  "params": {
    "message": {
      "messageId": "ce55994e-622e-435c-bcec-924afcc94e41",
      "role": "ROLE_USER",
      "parts": [
        { "text": "{\"a\": 19, \"b\": 23}" }
      ]
    }
  }
}

注目は2点です。

  1. roleROLE_USER。送り主は人間ではなくエージェント(受付)なのに、です。第1回では「依頼する側がROLE_USER、エージェント側がROLE_AGENT」と説明しましたが、その実物がこれです。依頼する側に立った者が、エージェントであってもROLE_USERを名乗ります
  2. messageId が、利用者がステップ2のリクエストで指定した u-1 という文字列ではなく、受付が採番したUUIDになっています。利用者のMessageを転送しているのではなく、受付が新しい依頼を自分の名義で作成しているのです。messageIdはそのMessageの作成者が付ける(A2A仕様 §4.1.4 Messageの「created by the message creator」)、という決まりがここでも貫かれています。なお、付録コードのコメントや状況Messageで使っている「転送」という語も、この「受付が新しい依頼を作って送る」という意味です

ステップ4: taskIdとcontextIdの独立を確認する

受付のTaskとワーカーのTaskのIDを並べます。ワーカー側のTaskは、ステップ2の応答に含まれていた workerTaskId を使い、[利用者 → ワーカー] に直接GetTaskすると確認できます。

curl -s -X POST $W/a2a/v1 -H "$H1" -H "$H2" \
  -d '{"jsonrpc":"2.0","id":2,"method":"GetTask","params":{"id":"<workerTaskIdの値>","historyLength":0}}'
受付のTask ワーカーのTask
taskId d165d4d3-... 3833ec71-...
contextId c646dec3-... 000b7a81-...

完全に別物です。1つの依頼が2区間(ホップ)を流れても、会話を横断する共通のIDは存在しません。taskIdもcontextIdもホップごとに独立していて、2つのタスクの対応付けは受付が自分で管理します(本実装ではmetadataで公開していますが、それは受付の実装判断です)。分散トレーシングのような「端から端まで追えるID」はA2Aの仕様には無い、というのは運用を考えるうえで押さえておきたい事実です。

ステップ5: 中断(INPUT_REQUIRED)を伝搬させる

第2回で見た「足し算 a + b のbが無いと聞き返す」動きは、間に受付が挟まるとどうなるでしょうか。[利用者 → 受付] にaだけの依頼を送ります。

curl -s -X POST $R/a2a/v1 -H "$H1" -H "$H2" \
  -d '{"jsonrpc":"2.0","id":3,"method":"SendMessage","params":{"message":{"messageId":"u-2","role":"ROLE_USER","parts":[{"text":"{\"a\": 19}"}]}}}'

受付からの応答は INPUT_REQUIRED で、質問文はこうなっています。

"text": "計算担当からの質問です: a=19 を受け取りました。b はいくつですか?"

ワーカーの質問が、受付によって一段包まれて届きました。受付のtaskIdを添えて「23」と答えます。

curl -s -X POST $R/a2a/v1 -H "$H1" -H "$H2" \
  -d '{"jsonrpc":"2.0","id":4,"method":"SendMessage","params":{"message":{"messageId":"u-3","taskId":"<受付のtaskIdの値>","role":"ROLE_USER","parts":[{"text":"23"}]}}}'

受付はこの回答を、ワーカー側タスクの再開のMessage(受付名義)として送り直し、2つのタスクが連動してCOMPLETEDになります。

大事なのは、この連動はプロトコルの機能ではないことです。プロトコル上は「利用者⇔受付の会話」と「受付⇔ワーカーの会話」という独立した2つの会話があるだけで、両者をつないでいるのは受付の実装(settleFromWorker() によるワーカーの状態の反映と、taskIdの対応表)です。中断の伝搬が成立するのは、どちらのホップも「INPUT_REQUIREDという状態」と「taskId付きSendMessageによる再開」という同じ決まりで表しているからで、受付はワーカーの状態を見て自分の状態を決め直すだけで済んでいます。

ステップ6: ワーカー側の履歴を確認する

完了後、ワーカーに直接GetTaskして history(履歴)を見ます。

curl -s -X POST $W/a2a/v1 -H "$H1" -H "$H2" \
  -d '{"jsonrpc":"2.0","id":5,"method":"GetTask","params":{"id":"<workerTaskIdの値>"}}'

履歴の要点を抜き出すとこうなっています。

ROLE_USER  messageId=7e29ec85-... : {"a": 19}
ROLE_AGENT messageId=fbeb27a3-... : a=19 を受け取りました。b はいくつですか?
ROLE_USER  messageId=a7b8039b-... : 23
ROLE_AGENT messageId=83fc75be-... : 19 + 23 の計算を始めます

ワーカーから見た「会話の相手」は受付です。ROLE_USERの2つのMessageはどちらも受付が採番・作成したもので、利用者の存在はワーカーの履歴のどこにも現れません。第1回からの主張「エージェントの中身は相手から見えなくてよい(opaque)」は、多段になっても保たれています。ワーカーは相手が人間のcurlか受付エージェントかを区別できず、区別する必要もないのです。

実測して初めて分かったこと

1. エージェント連携はクライアント実装の再利用でしかない

受付がワーカーに送っていたのは、私たちが第1回第2回でcurlから送っていたのと同じリクエストでした。ワーカーは無改変のまま、相手が人間からエージェントに変わったことに気づきもしません。「A2Aエージェント同士をつなぐ」とは、新しいプロトコル機能を使うことではなく、エージェントの中にクライアントを1つ持たせることでした。

2. roleは多段でも「立場」を表し続ける

受付が送るMessageはROLE_USERでした。roleは会話上の立場であって、送り主が人間かエージェントかを表しません。依頼する側に立った者がROLE_USERを名乗ります。ワーカーの履歴には、受付の発言がROLE_USERとして残ります。

3. ホップを横断するIDは無い。追跡の手がかりは自分で設計する

事実としては2つです。taskId・contextIdはホップごとに独立していて、それぞれのサーバが自分の会話のために採番します。2つの会話をつなぐのは受付の実装(taskIdの対応表と、ワーカーの状態を自タスクに反映する処理)だけです。中断の伝搬も、両ホップが状態を同じ列挙値(TASK_STATE_*)で表しているから状態の対応付けだけで成立しているのであって、プロトコルが面倒を見てくれるわけではありません。

なぜこうなっているかを考えると、IDが独立なのは不備ではなく、第1回で触れた「相手の中身は見えなくてよい(opaque)」という設計思想の帰結です。利用者から見えるのは受付のTaskだけで、その裏でワーカーが何体いるか、どのタスクに対応しているかは受付の内部事情です。この独立のおかげで、受付はワーカーの障害時に別のワーカーへ依頼し直したり、依頼を複数のワーカーに分けたりしても、利用者に渡したtaskIdを変えずに済みます。

代償として、端から端まで追える共通IDは仕様に無く、受付が持つ対応表だけが2つの会話をつなぐ唯一の手がかりになります。受付がこの対応表を失えば、利用者のTaskはワーカー側のどの作業と対応するのか分からなくなります。本実装は対応をmetadataで公開して観察しやすくしましたが、実運用では「対応表をどう永続化するか」「追跡IDをどこで運ぶか(Task.metadataに載せるか、HTTP層のトレースヘッダに任せるか)」を、間に立つエージェントの設計で決めなければなりません。多段構成のトレーサビリティは設計課題として残ります。

MCPとの対比

比較の前提として、本シリーズが対象にしてきたMCPは2025-06-18版です。

観点 A2A(v1.0) MCP(2025-06-18版)
多段連携の作り方 受付エージェントがワーカーのA2Aクライアントになる(本記事の構成) MCPサーバが別のMCPサーバのMCPクライアントになる
下流の作業状態を上流へ伝える手段 Taskの状態(TASK_STATE_* の列挙値)が標準で決まっている 作業状態を表す標準の値は無い。progressToken による進捗数値の通知のみ

多段連携の作り方そのものは、どちらも「間に立つサーバが、下流に対してはクライアントになる」という同じ形です。違いは、下流の作業状態を上流にどう伝えるかにあります。A2Aでは作業状態を表す値が TASK_STATE_WORKINGTASK_STATE_INPUT_REQUIRED のように標準で決まっているため、受付は「ワーカーの状態を見て自分の状態を決める」だけで進捗の中継が成立します。今回の構成を作って初めて実感した差です。MCPには progressToken で進捗を数値として通知する仕組みはありますが、作業状態を表す標準の値は無いため、下流の状態を上流の状態に機械的に対応付けることはできません。

なお、MCP 2025-11-25版では実験的機能(experimental)としてTasksが追加され、input_required などの作業状態を表す値が入り始めています。

まとめ

  • エージェント同士の連携に専用機能は無い。片方のエージェントがもう片方のA2Aクライアントになるだけで、呼ばれる側は相手が人間かエージェントかを区別できない(する必要もない)
  • roleは会話上の立場。依頼する側のエージェントはROLE_USERを名乗る
  • taskId・contextIdはホップごとに独立。多段の連動(中断の伝搬を含む)はプロトコルではなく、間に立つエージェントの実装がつくる。端から端まで追えるIDは無く、トレーサビリティは設計課題として残る
  • 受付が返すTaskの状態は、ワーカーの状態のコピーではなく「ワーカーがこうなったから、受付としてはこうする」と決め直した結果
  • messageIdは受付が採番し直す。利用者が付けたIDはワーカーには届かず、受付は利用者のMessageをそのまま送るのではなく、新しい依頼を自分の名義で作っている

次回(第4回: 今回の取り次ぎはブロッキング、つまり受付がワーカーの応答を待ち続ける方式でした。利用者の待ち時間3.9秒の内訳はほぼ丸ごと「受付がワーカーを待って接続を保持している時間」で、ホップが増えるほど待ち時間とタイムアウト設計が積み上がります。次回はこれを、受領書だけもらって接続を切り、結果はワーカーからの逆向きHTTP POSTで受け取る**プッシュ通知(webhook)**に置き換えます。「通知を受け取る側もHTTPサーバになる」という、エージェント2体ならではの構成が見えてきます。

付録: reception-server.mjs の全文

本文の体験で使用した受付エージェントの全文です(ブロッキングで取り次ぐ方式。第4回でプッシュ通知に対応した版へ拡張します)。

// A2A v1.0 エージェント間連携の学習用「受付エージェント」(依存ゼロNode, JSON-RPCバインディング)
//
// 役割: 計算依頼の受付窓口。自分では計算せず、Slow Calc Agent(第2回の task-server.mjs。同じディレクトリに置く)に
//       A2Aで依頼を転送する。つまり「A2Aサーバでありながら、別のエージェントに対してはA2Aクライアント」。
//
// 本記事(第3回)の観察ポイント:
//   - エージェントが送る側に回っても、Messageの role は ROLE_USER(Roleは会話の立場であって人間/AIの区別ではない)
//   - taskId / contextId はホップごとに独立(利用者⇔受付のTaskと、受付⇔ワーカーのTaskは別物)
//   - INPUT_REQUIRED は伝搬する: ワーカーの質問を受付が自タスクの中断として利用者へ中継し、
//     利用者の回答を受付がワーカーのタスク再開として転送する
//   - 受付が送るMessageのmessageIdは受付が発番する(利用者のmessageIdを使い回さない。発番者=作成者)
//
// 対応メソッド: SendMessage / GetTask(ストリーミング非対応 = Cardの capabilities.streaming: false)
// 取り次ぎ方式: ワーカーの応答(終端または中断)までHTTP接続を保持して待つ(blocking)

import http from 'node:http';
import crypto from 'node:crypto';
import fs from 'node:fs';
import { fileURLToPath } from 'node:url';

const PORT = 4103;
const BASE = `http://localhost:${PORT}`;
const WORKER_BASE = process.env.WORKER_BASE ?? 'http://localhost:4102';
const WIRE_LOG = fileURLToPath(new URL('./reception-wire.log', import.meta.url));

const AGENT_CARD = {
  name: 'Calc Reception Agent',
  description: '計算依頼の受付窓口。依頼内容を確認し、計算はSlow Calc AgentにA2Aで取り次ぐ(エージェント間連携の学習用)',
  version: '0.1.0',
  supportedInterfaces: [
    { url: `${BASE}/a2a/v1`, protocolBinding: 'JSONRPC', protocolVersion: '1.0' },
  ],
  capabilities: { streaming: false, pushNotifications: false, extendedAgentCard: false },
  defaultInputModes: ['application/json', 'text/plain'],
  defaultOutputModes: ['text/plain'],
  skills: [
    {
      id: 'delegate-add',
      name: '足し算の取り次ぎ',
      description: '足し算の依頼を受け付け、計算担当のエージェントに取り次いで結果を返す',
      tags: ['calculator', 'delegation'],
      examples: ['{"a": 19, "b": 23}', '19 + 23'],
      inputModes: ['application/json', 'text/plain'],
      outputModes: ['text/plain'],
    },
  ],
};
const CARD_ETAG = `"${AGENT_CARD.version}"`;

const PARSE_ERROR = -32700;
const INVALID_REQUEST = -32600;
const INVALID_PARAMS = -32602;
const METHOD_NOT_FOUND = -32601;
const TASK_NOT_FOUND = -32001;          // A2A仕様 §5.4 TaskNotFoundError
const UNSUPPORTED_OPERATION = -32004;   // A2A仕様 §5.4 UnsupportedOperationError
const VERSION_NOT_SUPPORTED = -32009;   // A2A仕様 §5.4 VersionNotSupportedError

function wire(direction, text) {
  fs.appendFileSync(WIRE_LOG, `[${new Date().toISOString()}] [${direction}]\n${text}\n\n`);
}

// ---- ワーカーの発見と呼び出し(ここが「A2Aクライアント」側の実装) ----

let workerEndpoint; // Agent Cardから得たRPCエンドポイント(発見は初回のみ、以後キャッシュ)

async function getWorkerEndpoint() {
  if (workerEndpoint) return workerEndpoint;
  const url = `${WORKER_BASE}/.well-known/agent-card.json`;
  wire('受付 → ワーカー', `GET ${url}`);
  const res = await fetch(url);
  const card = await res.json();
  wire('ワーカー → 受付', `HTTP ${res.status}\n\n${JSON.stringify(card, null, 2)}`);
  workerEndpoint = card.supportedInterfaces[0].url; // 発見 → 接続の2段階(第1回と同じ流儀)
  return workerEndpoint;
}

let outboundRpcId = 1; // 受付→ワーカーのJSON-RPC idは受付が発番する(利用者側のidとは別の名前空間)

async function callWorker(method, params) {
  const endpoint = await getWorkerEndpoint();
  const payload = { jsonrpc: '2.0', id: outboundRpcId++, method, params };
  const body = JSON.stringify(payload, null, 2);
  wire('受付 → ワーカー', `POST ${endpoint}\ncontent-type: application/json\na2a-version: 1.0\n\n${body}`);
  const res = await fetch(endpoint, {
    method: 'POST',
    headers: { 'Content-Type': 'application/json', 'A2A-Version': '1.0' },
    body,
  });
  const text = await res.text();
  wire('ワーカー → 受付', `HTTP ${res.status}\n\n${text}`);
  return JSON.parse(text);
}

// ---- 受付自身のTask管理(task-server.mjsの簡略版。ストリーミングなし) ----
// task = { id, contextId, status, artifacts, history, workerTaskId, waiters }
const tasks = new Map();

const TERMINAL = new Set(['TASK_STATE_COMPLETED', 'TASK_STATE_FAILED', 'TASK_STATE_CANCELED', 'TASK_STATE_REJECTED']);
const INTERRUPTED = new Set(['TASK_STATE_INPUT_REQUIRED', 'TASK_STATE_AUTH_REQUIRED']);

function agentMessage(text, taskId, contextId) {
  return { messageId: crypto.randomUUID(), taskId, contextId, role: 'ROLE_AGENT', parts: [{ text }] };
}

function setState(task, state, noteText) {
  const message = noteText ? agentMessage(noteText, task.id, task.contextId) : undefined;
  task.status = { state, timestamp: new Date().toISOString(), ...(message ? { message } : {}) };
  if (message) task.history.push(message);
  if (TERMINAL.has(state) || INTERRUPTED.has(state)) {
    task.waiters.splice(0).forEach((resolve) => resolve());
  }
}

function taskView(task, historyLength) {
  const view = {
    id: task.id,
    contextId: task.contextId,
    status: task.status,
    artifacts: task.artifacts,
    // 取り次ぎ先の対応付けは(プロトコル規定外の情報なので)metadataで公開する
    ...(task.workerTaskId ? { metadata: { workerTaskId: task.workerTaskId } } : {}),
  };
  if (historyLength === 0) return view;
  const h = historyLength > 0 ? task.history.slice(-historyLength) : task.history;
  return { ...view, history: h };
}

function waitForSettled(task) {
  if (TERMINAL.has(task.status.state) || INTERRUPTED.has(task.status.state)) return Promise.resolve();
  return new Promise((resolve) => task.waiters.push(resolve));
}

// ---- 取り次ぎ本体 ----

// ワーカーから返ってきたタスクの状態を読み、受付自身のタスクをどの状態に進めるかを決める
function settleFromWorker(task, workerTask) {
  task.workerTaskId = workerTask.id;
  const state = workerTask.status.state;
  if (state === 'TASK_STATE_COMPLETED') {
    for (const a of workerTask.artifacts ?? []) {
      task.artifacts.push({
        artifactId: crypto.randomUUID(), // 受付の成果物として発番し直す(出所はmetadataに残す)
        name: a.name,
        parts: a.parts,
        metadata: { source: 'Slow Calc Agent', workerTaskId: workerTask.id, workerArtifactId: a.artifactId },
      });
    }
    setState(task, 'TASK_STATE_COMPLETED', '計算担当から結果を受け取りました');
  } else if (state === 'TASK_STATE_INPUT_REQUIRED') {
    const question = workerTask.status.message?.parts?.map((p) => p.text ?? '').join('') ?? '(質問文なし)';
    setState(task, 'TASK_STATE_INPUT_REQUIRED', `計算担当からの質問です: ${question}`);
  } else if (state === 'TASK_STATE_REJECTED') {
    const reason = workerTask.status.message?.parts?.map((p) => p.text ?? '').join('') ?? '';
    setState(task, 'TASK_STATE_REJECTED', `計算担当に断られました: ${reason}`);
  } else {
    setState(task, 'TASK_STATE_FAILED', `計算担当のタスクが ${state} になりました`);
  }
}

// 新規依頼をワーカーへ転送する。Messageは受付が新規に作る(messageIdも受付が発番)。
// ワーカーの応答(終端または中断)までHTTP接続を保持して待つ
async function delegate(task, text) {
  try {
    setState(task, 'TASK_STATE_WORKING', '依頼を計算担当(Slow Calc Agent)に取り次ぎます');
    const out = await callWorker('SendMessage', {
      message: { messageId: crypto.randomUUID(), role: 'ROLE_USER', parts: [{ text }] },
    });
    if (out.error) return setState(task, 'TASK_STATE_FAILED', `計算担当がエラーを返しました: ${out.error.code} ${out.error.message}`);
    settleFromWorker(task, out.result.task);
  } catch (e) {
    setState(task, 'TASK_STATE_FAILED', `計算担当に接続できません: ${e.message}`);
  }
}

// 利用者の回答を、ワーカー側タスクの再開(taskId付きSendMessage)として転送する
async function forwardAnswer(task, text) {
  try {
    const out = await callWorker('SendMessage', {
      message: { messageId: crypto.randomUUID(), taskId: task.workerTaskId, role: 'ROLE_USER', parts: [{ text }] },
    });
    if (out.error) return setState(task, 'TASK_STATE_FAILED', `計算担当がエラーを返しました: ${out.error.code} ${out.error.message}`);
    settleFromWorker(task, out.result.task);
  } catch (e) {
    setState(task, 'TASK_STATE_FAILED', `計算担当に接続できません: ${e.message}`);
  }
}

// ---- JSON-RPCメソッド ----
async function handleSendMessage(params) {
  const msg = params?.message;
  if (!Array.isArray(msg?.parts) || msg.parts.length === 0) {
    return { error: { code: INVALID_PARAMS, message: 'Invalid parameters' } };
  }
  const text = msg.parts.map((p) => p.text ?? '').join('');

  let task;
  let work;
  if (msg.taskId) {
    task = tasks.get(msg.taskId);
    if (!task) return { error: { code: TASK_NOT_FOUND, message: 'Task not found' } };
    if (msg.contextId && msg.contextId !== task.contextId) {
      return { error: { code: INVALID_PARAMS, message: 'contextId does not match the referenced task' } };
    }
    if (task.status.state !== 'TASK_STATE_INPUT_REQUIRED') {
      return { error: { code: UNSUPPORTED_OPERATION, message: 'Task is not waiting for input' } };
    }
    task.history.push({ ...msg, contextId: task.contextId });
    // 再開は同期的にWORKINGへ(INPUT_REQUIREDのままだとブロッキング待機が即返ってしまう)
    setState(task, 'TASK_STATE_WORKING', '回答を計算担当へ転送します');
    work = forwardAnswer(task, text);
  } else {
    task = {
      id: crypto.randomUUID(),
      contextId: msg.contextId ?? crypto.randomUUID(),
      status: {},
      artifacts: [],
      history: [],
      waiters: [],
    };
    // 履歴のMessageは単体で読めるように所属を補完してから保存する
    task.history.push({ ...msg, taskId: task.id, contextId: task.contextId });
    tasks.set(task.id, task);
    setState(task, 'TASK_STATE_SUBMITTED');
    work = delegate(task, text);
  }

  // A2A仕様 §3.2.2: 既定はブロッキング。取り次ぎ(work)の完了 = 自タスクが終端または中断に達した時点
  if (params?.configuration?.returnImmediately) {
    void work; // 受領書を先に返し、取り次ぎは裏で続ける
  } else {
    await work;
    await waitForSettled(task);
  }
  return { result: { task: taskView(task) } };
}

const server = http.createServer((req, res) => {
  let body = '';
  req.on('data', (c) => (body += c));
  req.on('end', async () => {
    const headerDump = ['host', 'content-type', 'a2a-version', 'if-none-match']
      .filter((h) => req.headers[h] !== undefined)
      .map((h) => `${h}: ${req.headers[h]}`)
      .join('\n');
    wire('利用者 → 受付', `${req.method} ${req.url}\n${headerDump}${body ? '\n\n' + body : ''}`);

    const reply = (status, headers, payload) => {
      wire('受付 → 利用者', `HTTP ${status}\n${Object.entries(headers).map(([k, v]) => `${k}: ${v}`).join('\n')}${payload ? '\n\n' + payload : ''}`);
      res.writeHead(status, headers);
      res.end(payload);
    };

    if (req.method === 'GET' && req.url === '/.well-known/agent-card.json') {
      if (req.headers['if-none-match'] === CARD_ETAG) {
        return reply(304, { ETag: CARD_ETAG, 'Cache-Control': 'max-age=3600' }, '');
      }
      return reply(200, { 'Content-Type': 'application/json', ETag: CARD_ETAG, 'Cache-Control': 'max-age=3600' }, JSON.stringify(AGENT_CARD, null, 2));
    }

    if (req.method === 'POST' && req.url === '/a2a/v1') {
      const respond = (obj) => reply(200, { 'Content-Type': 'application/json' }, JSON.stringify(obj, null, 2));
      const err = (id, code, message) => respond({ jsonrpc: '2.0', id, error: { code, message } });

      let rpc;
      try {
        rpc = JSON.parse(body);
      } catch {
        return err(null, PARSE_ERROR, 'Invalid JSON payload');
      }
      if (rpc.jsonrpc !== '2.0' || typeof rpc.method !== 'string') {
        return err(rpc.id ?? null, INVALID_REQUEST, 'Request payload validation error');
      }
      if ((req.headers['a2a-version'] ?? '0.3') !== '1.0') {
        return err(rpc.id, VERSION_NOT_SUPPORTED, 'Version not supported');
      }

      switch (rpc.method) {
        case 'SendMessage': {
          const out = await handleSendMessage(rpc.params);
          return respond({ jsonrpc: '2.0', id: rpc.id, ...(out.error ? { error: out.error } : out) });
        }
        case 'GetTask': {
          const task = tasks.get(rpc.params?.id);
          if (!task) return err(rpc.id, TASK_NOT_FOUND, 'Task not found');
          // rpc GetTask returns (Task): resultはTaskそのもの(result.taskに包まない)
          return respond({ jsonrpc: '2.0', id: rpc.id, result: taskView(task, rpc.params?.historyLength) });
        }
        default:
          // ストリーミング等は非対応(Cardの capabilities.streaming: false と整合)
          return err(rpc.id, METHOD_NOT_FOUND, 'Method not found');
      }
    }

    reply(404, { 'Content-Type': 'application/json' }, JSON.stringify({ error: 'not found' }));
  });
});

server.listen(PORT, () => {
  console.log(`Calc Reception Agent (A2A v1.0, JSON-RPC binding) : ${BASE}`);
  console.log(`  Agent Card : ${BASE}/.well-known/agent-card.json`);
  console.log(`  JSON-RPC   : POST ${BASE}/a2a/v1`);
  console.log(`  worker     : ${WORKER_BASE}(Slow Calc Agent)`);
  console.log(`  wire log   : ${WIRE_LOG}`);
});
0
0
0

Register as a new user and use Qiita more conveniently

  1. You get articles that match your needs
  2. You can efficiently read back useful information
  3. You can use dark theme
What you can do with signing up
0
0

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?