fix(a2a): complete tasks under message-tool-only reply policy (#154317)

Related: #144534

## What Problem This Solves

Fixes A2A tasks remaining `TASK_STATE_WORKING` after a completed agent turn when `messages.visibleReplies="message_tool"` suppresses the adapter's correlated final reply.

## User Impact

User impact: authenticated A2A callers receive the final artifact through `SendMessage`/`GetTask` without needing an outbound peer URL. Peer authentication, allowlisting, task isolation, command rejection, retention, and permissions are unchanged.

## Why This Change Was Made

A2A source replies complete an existing local task; the generic message tool starts a separate outbound message. Request automatic source delivery through A2A's existing task-store callback. The resolver's explicit `strictMessageToolOnly=true` precedence remains unchanged.

The regression exercises the registered HTTP endpoint through a full Gateway and a deterministic loopback model, rather than asserting a forwarded option.

## Evidence

- On main `10c229159c`, authenticated `SendMessage` returned a working task; after the model returned its answer and the Gateway recorded the turn as completed, `GetTask` still returned the same working task without an artifact. Changing only the global policy to automatic completed the control task.
- With this fix and the original message-tool-only policy, the full Gateway returned the exact final artifact for the same authenticated task ID. Wrong bearer tokens returned HTTP 401; another authenticated peer could not read that task. Distinct peers reusing a context ID and one peer using distinct contexts received their own artifacts. Peer slash commands remained rejected.
- 59 focused tests passed: 36 source-delivery policy tests (including strict precedence), 12 A2A inbound tests, 10 task-store tests (including FIFO and peer isolation), and the new real-Gateway A2A regression. The new E2E case took 9.027 seconds (31.78 seconds including its Vitest setup); hosted CI timing is not yet available.
- Candidate runtime build, scoped lint/format, both line-count ratchets, and staged whitespace checks passed. All task-owned Gateway/model processes stopped; their loopback ports no longer accepted connections.

**Known limitation:** a held-first, concurrent same-peer/same-context control still failed before the second model request with `Async work scope is closed`. The identical control failed on unchanged main with automatic delivery, so this is a separate pre-existing queued-turn lifecycle failure, not a passing concurrent FIFO proof or a fix claimed here. The existing task-store FIFO tests pass. This PR does not claim to resolve every delivery-confirmation path in #144534.

Runtime proof used the OpenClaw agent runtime with a deterministic loopback provider, not native Codex or a paid model.

Co-authored-by: Ayaan Zaidi <hi@obviy.us>
This commit is contained in:
Operations Team 2026-09-21 06:02:35 -07:00 • committed by GitHub
parent 33ecdbb346
commit 5769409f2b
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
2 changed files with 127 additions and 0 deletions

View file

@ -122,6 +122,9 @@ export async function dispatchA2aInbound(params: A2aInboundDispatchParams): Prom
params.store.fail(params.taskId, error);
},
},
// Source replies complete the correlated task; the generic message tool
// starts a separate outbound message without that task correlation.
replyOptions: { sourceReplyDeliveryMode: "automatic" },
replyPipeline: {},
});
if (dispatch.admission.kind !== "dispatch") {

View file

@ -0,0 +1,124 @@
import { randomUUID } from "node:crypto";
import { expect, test } from "vitest";
import { startQaMockOpenAiServer } from "../extensions/qa-lab/api.js";
import {
createOpenClawTestInstance,
type OpenClawTestInstance,
} from "./helpers/openclaw-test-instance.js";
import { runQaGatewayFixture } from "./helpers/qa-gateway-cleanup.js";
test("A2A completes correlated tasks under message-tool-only source policy", async () => {
const model = await startQaMockOpenAiServer({ modelRefs: ["a2a-proof/a2a-proof"] });
let instance: OpenClawTestInstance | undefined;
await runQaGatewayFixture(
async () => {
instance = await createOpenClawTestInstance({
name: "a2a-correlated-replies",
config: {
plugins: { allow: ["a2a", "openai"], slots: { memory: "none" } },
agents: {
defaults: {
heartbeat: { every: "0m" },
model: { primary: "a2a-proof/a2a-proof" },
models: { "a2a-proof/a2a-proof": { agentRuntime: { id: "openclaw" } } },
skipBootstrap: true,
skills: [],
},
},
tools: { profile: "minimal", alsoAllow: ["message"] },
messages: { visibleReplies: "message_tool" },
channels: {
a2a: {
enabled: true,
rateLimitPerMinute: 0,
peers: {
hermes: { token: "a2a-hermes-test-token" },
crew: { token: "a2a-crew-test-token" },
},
},
},
models: {
mode: "replace",
providers: {
"a2a-proof": {
baseUrl: `${model.baseUrl}/v1`,
apiKey: "a2a-model-test-token",
api: "openai-responses",
request: { allowPrivateNetwork: true },
models: [
{
id: "a2a-proof",
name: "A2A proof",
api: "openai-responses",
reasoning: false,
input: ["text"],
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 },
contextWindow: 128_000,
maxTokens: 4096,
},
],
},
},
},
},
env: {
OPENCLAW_SKIP_CHANNELS: undefined,
OPENCLAW_SKIP_PROVIDERS: undefined,
OPENCLAW_TEST_MINIMAL_GATEWAY: undefined,
},
});
await instance.startGateway();
const endpoint = `http://127.0.0.1:${instance.port}/a2a/v1`;
const rpc = async (method: string, params: unknown, token: string) =>
fetch(endpoint, {
method: "POST",
headers: { "content-type": "application/json", authorization: `Bearer ${token}` },
body: JSON.stringify({ jsonrpc: "2.0", id: randomUUID(), method, params }),
});
expect((await rpc("GetTask", { id: "unknown" }, "wrong-token")).status).toBe(401);
const contextId = "shared-peer-context";
const completeTask = async (peer: string) => {
const token = `a2a-${peer}-test-token`;
const text = `A2A-${peer.toUpperCase()}-COMPLETE`;
const response = await rpc(
"SendMessage",
{
message: {
messageId: randomUUID(),
contextId,
role: "ROLE_USER",
parts: [{ text: `Reply exactly \`${text}\`` }],
},
configuration: { returnImmediately: true },
},
token,
);
expect(response.status).toBe(200);
const accepted = await response.json();
expect(accepted.result.task.status.state).toBe("TASK_STATE_WORKING");
const taskId = accepted.result.task.id;
await expect
.poll(async () => (await (await rpc("GetTask", { id: taskId }, token)).json()).result, {
timeout: 30_000,
})
.toMatchObject({
id: taskId,
contextId,
status: { state: "TASK_STATE_COMPLETED" },
artifacts: [{ parts: [{ text }] }],
});
return taskId;
};
const [hermesTask, crewTask] = await Promise.all([
completeTask("hermes"),
completeTask("crew"),
]);
expect(hermesTask).not.toBe(crewTask);
const denied = await rpc("GetTask", { id: hermesTask }, "a2a-crew-test-token");
expect(await denied.json()).toMatchObject({ error: { code: -32001 } });
},
() => instance?.cleanup(),
() => model.stop(),
);
}, 120_000);