# Streaming Task Pipeline (/aiway/a2a/scenarios/streaming-task-pipeline)



This scenario streams a **long-running agent task** end to end. A single client opens a
`message/stream&#x60; request, the gateway relays the agent's progress as Server-Sent Events,
and the client consumes a sequence of **`task.*` envelopes** — interim `task.status`
updates, one or more `task.artifact` results, and a terminal `task.done`. It builds on
[streaming (SSE)](/aiway/a2a/streaming) and the per-agent
[virtual subscriber](/aiway/a2a/architecture).

## The setup [#the-setup]

The agent is a plain HTTP server registered by URL — no KubeMQ SDK runs on it. When the
client `POST`s a `message/stream` request to `/a2a/{agent_id}`, the gateway opens an SSE
relay through the agent's virtual subscriber: it sends a Query carrying the stream
channel, the agent streams SSE events back, and the gateway relays each one to the caller
as a `task.*` event. The client reads the stream until it sees the terminal `task.done`
(or `task.error`), and **cancelling is just disconnecting** — when the caller's connection
drops, the gateway sends a `stream_cancel` query to the agent and tears the relay down.

<Mermaid
  chart="`
graph LR
CLIENT[&#x22;Client<br/>(task consumer)&#x22;]
GW[&#x22;A2A Gateway<br/>:9090&#x22;]
VSUB[&#x22;Virtual Subscriber&#x22;]
BROKER[&#x22;Message Broker&#x22;]
AGENT[&#x22;task-events-agent-01<br/>(HTTP SSE server)&#x22;]

CLIENT -- &#x22;POST /a2a/...<br/>message/stream&#x22; --> GW
GW --> VSUB
VSUB --> BROKER
BROKER -. &#x22;_AGENTS_.agents/<id>&#x22; .-> AGENT
AGENT -. &#x22;task.status / task.artifact / task.done&#x22; .-> VSUB
VSUB -. &#x22;SSE relay&#x22; .-> CLIENT

class CLIENT client
class GW,VSUB aiway
class BROKER broker
class AGENT external
`"
/>

*The gateway relays the agent's SSE task events back to the caller and cancels the agent when the caller disconnects.*

## The task envelopes [#the-task-envelopes]

Each SSE frame names an event type and carries a JSON envelope. The agent emits four
envelope kinds; the gateway maps them to these SSE event names:

| SSE event       | Envelope `type` | Meaning                                              |
| --------------- | --------------- | ---------------------------------------------------- |
| `task.status`   | `status_update` | Interim progress (`status`, `progress`, `total`)     |
| `task.artifact` | `artifact`      | A produced result (`name`, `data`)                   |
| `task.done`     | `done`          | Terminal success (`final_result`) — stream closes    |
| `task.error`    | `error`         | Terminal failure (`code`, `message`) — stream closes |

The wire frame for a status update looks like this:

```text
event: task.status
data: {"stream_id":"<uuid>","type":"status_update","payload":{"status":"working","progress":1,"total":3}}

```

A comment keepalive (`: keepalive`) arrives every 30 seconds on an otherwise idle stream,
and the stream closes on the first `task.done` or `task.error`.

## Step 1 — Open the stream [#step-1--open-the-stream]

Send a `message/stream` JSON-RPC request to the agent. The gateway responds with
`Content-Type: text/event-stream` and begins relaying the agent's `task.*` events. Pass
`Accept: text/event-stream` and read the response body line by line.

<Tabs groupId="language" items="['curl','C#','Go','Java','Python','TypeScript']">
  <Tab value="curl">
    ```bash
    # -N disables curl buffering so events print as they arrive.
    curl -N -X POST http://localhost:9090/a2a/task-events-agent-01 \
      -H 'Content-Type: application/json' \
      -H 'Accept: text/event-stream' \
      -d '{
        "jsonrpc": "2.0",
        "id": 1,
        "method": "message/stream",
        "params": {"message": {"parts": [{"text": "Show me all event types"}]}}
      }'
    ```
  </Tab>

  <Tab value="C#">
    ```csharp
    using System.Text;
    using System.Text.Json.Nodes;

    const string KubeMqUrl = "http://localhost:9090";
    const string AgentId = "task-events-agent-01";

    var payload = new JsonObject
    {
        ["jsonrpc"] = "2.0",
        ["id"] = 1,
        ["method"] = "message/stream",
        ["params"] = new JsonObject
        {
            ["message"] = new JsonObject
            {
                ["parts"] = new JsonArray(new JsonObject { ["text"] = "Show me all event types" })
            }
        }
    };

    using var client = new HttpClient { Timeout = TimeSpan.FromSeconds(60) };
    var request = new HttpRequestMessage(HttpMethod.Post, $"{KubeMqUrl}/a2a/{AgentId}")
    {
        Content = new StringContent(payload.ToJsonString(), Encoding.UTF8, "application/json")
    };

    Console.WriteLine("Connecting to SSE stream...");
    using var resp = await client.SendAsync(request, HttpCompletionOption.ResponseHeadersRead);
    using var stream = await resp.Content.ReadAsStreamAsync();
    using var reader = new StreamReader(stream);
    // ...continue reading events in Step 2
    ```
  </Tab>

  <Tab value="Go">
    ```go
    package main

    import (
        "bytes"
        "encoding/json"
        "fmt"
        "net/http"
        "os"
    )

    const (
        kubemqURL = "http://localhost:9090"
        agentID   = "task-events-agent-01"
    )

    func main() {
        payload := map[string]interface{}{
            "jsonrpc": "2.0",
            "id":      1,
            "method":  "message/stream",
            "params": map[string]interface{}{
                "message": map[string]interface{}{
                    "parts": []map[string]interface{}{{"text": "Show me all event types"}},
                },
            },
        }

        data, _ := json.Marshal(payload)
        resp, err := http.Post(kubemqURL+"/a2a/"+agentID, "application/json", bytes.NewReader(data))
        if err != nil {
            fmt.Fprintf(os.Stderr, "Request failed: %v\n", err)
            os.Exit(1)
        }
        defer resp.Body.Close()
        fmt.Println("Connecting to SSE stream...")
        // ...continue reading events in Step 2
    }
    ```
  </Tab>

  <Tab value="Java">
    ```java
    import com.fasterxml.jackson.databind.ObjectMapper;

    import java.net.URI;
    import java.net.http.HttpClient;
    import java.net.http.HttpRequest;
    import java.net.http.HttpResponse;
    import java.time.Duration;
    import java.util.List;
    import java.util.Map;

    public class Client {

        static final String KUBEMQ_URL = "http://localhost:9090";
        static final String AGENT_ID = "task-events-agent-01";
        static final ObjectMapper MAPPER = new ObjectMapper();

        public static void main(String[] args) throws Exception {
            var payload = Map.of(
                "jsonrpc", "2.0",
                "id", 1,
                "method", "message/stream",
                "params", Map.of(
                    "message", Map.of("parts", List.of(Map.of("text", "Show me all event types")))
                )
            );

            var client = HttpClient.newBuilder().connectTimeout(Duration.ofSeconds(60)).build();
            var req = HttpRequest.newBuilder()
                .uri(URI.create(KUBEMQ_URL + "/a2a/" + AGENT_ID))
                .header("Content-Type", "application/json")
                .header("Accept", "text/event-stream")
                .timeout(Duration.ofSeconds(60))
                .POST(HttpRequest.BodyPublishers.ofString(MAPPER.writeValueAsString(payload)))
                .build();

            System.out.println("Connecting to SSE stream...");
            var resp = client.send(req, HttpResponse.BodyHandlers.ofLines());
            // ...continue reading events in Step 2
        }
    }
    ```
  </Tab>

  <Tab value="Python">
    ```python
    """Stream a long-running task and consume its task.* envelopes."""

    import asyncio
    import json
    from collections import Counter

    import httpx
    from httpx_sse import aconnect_sse

    KUBEMQ_URL = "http://localhost:9090"
    AGENT_ID = "task-events-agent-01"


    async def main() -> None:
        payload = {
            "jsonrpc": "2.0",
            "id": 1,
            "method": "message/stream",
            "params": {"message": {"parts": [{"text": "Show me all event types"}]}},
        }

        async with httpx.AsyncClient(timeout=60) as client:
            print("Connecting to SSE stream...")
            async with aconnect_sse(
                client, "POST", f"{KUBEMQ_URL}/a2a/{AGENT_ID}", json=payload
            ) as event_source:
                ...  # consume events in Step 2


    if __name__ == "__main__":
        asyncio.run(main())
    ```
  </Tab>

  <Tab value="TypeScript">
    ```typescript
    const KUBEMQ_URL = "http://localhost:9090";
    const AGENT_ID = "task-events-agent-01";

    async function main() {
      const request = {
        jsonrpc: "2.0",
        id: 1,
        method: "message/stream",
        params: { message: { parts: [{ text: "Send me task events" }] } },
      };

      const resp = await fetch(`${KUBEMQ_URL}/a2a/${AGENT_ID}`, {
        method: "POST",
        headers: {
          "Content-Type": "application/json",
          Accept: "text/event-stream",
        },
        body: JSON.stringify(request),
      });

      const reader = resp.body!.getReader();
      // ...continue reading events in Step 2
    }

    main().catch(console.error);
    ```
  </Tab>
</Tabs>

## Step 2 — Consume the task.\* envelopes [#step-2--consume-the-task-envelopes]

Read the stream frame by frame, dispatch on the SSE event name, and pull the result out of
`payload.payload`. Keep a count of each event type for a summary, and **break on the
terminal envelope** (`task.done` or `task.error`) — that frame closes the stream.

<Tabs groupId="language" items="['curl','C#','Go','Java','Python','TypeScript']">
  <Tab value="curl">
    ```text
    event: task.status
    data: {"stream_id":"...","type":"status_update","payload":{"status":"working","progress":1,"total":3}}

    event: task.status
    data: {"stream_id":"...","type":"status_update","payload":{"status":"working","progress":2,"total":3}}

    event: task.artifact
    data: {"stream_id":"...","type":"artifact","payload":{"name":"result.json","data":{"key":"value"}}}

    event: task.done
    data: {"stream_id":"...","type":"done","payload":{"final_result":"completed"}}
    ```
  </Tab>

  <Tab value="C#">
    ```csharp
    var eventTypes = new Dictionary<string, int>();
    string? eventType = null;

    while (!reader.EndOfStream)
    {
        var line = await reader.ReadLineAsync();
        if (line == null) break;

        if (line.StartsWith("event: "))
            eventType = line[7..];
        else if (line.StartsWith("data: ") && eventType != null)
        {
            eventTypes[eventType] = eventTypes.GetValueOrDefault(eventType) + 1;
            var data = JsonNode.Parse(line[6..])!;

            if (eventType == "task.status")
            {
                var p = data["payload"]!["payload"]!;
                Console.WriteLine($"  [STATUS]   progress={p["progress"]}/{p["total"]}");
            }
            else if (eventType == "task.artifact")
                Console.WriteLine($"  [ARTIFACT] name={data["payload"]!["payload"]!["name"]}");
            else if (eventType == "task.done")
                Console.WriteLine($"  [DONE]     result={data["payload"]!["payload"]!["final_result"]}");
            else if (eventType == "task.error")
                Console.WriteLine($"  [ERROR]    {data["payload"]!["payload"]}");

            if (eventType is "task.done" or "task.error")
                break;
        }
        else if (line.Length == 0)
            eventType = null;
    }

    Console.WriteLine($"\n--- Event Summary ---");
    foreach (var (et, count) in eventTypes.OrderBy(kv => kv.Key))
        Console.WriteLine($"  {et}: {count}");
    Console.WriteLine($"  Total: {eventTypes.Values.Sum()}");
    ```
  </Tab>

  <Tab value="Go">
    ```go
    scanner := bufio.NewScanner(resp.Body)
    eventType := ""
    counts := map[string]int{}

    for scanner.Scan() {
        line := scanner.Text()
        if strings.HasPrefix(line, "event: ") {
            eventType = strings.TrimPrefix(line, "event: ")
        } else if strings.HasPrefix(line, "data: ") {
            counts[eventType]++
            dataStr := strings.TrimPrefix(line, "data: ")

            var d map[string]interface{}
            json.Unmarshal([]byte(dataStr), &d)
            payload, _ := d["payload"].(map[string]interface{})

            inner, _ := payload["payload"].(map[string]interface{})
            if inner == nil {
                inner = payload
            }
            switch eventType {
            case "task.status":
                fmt.Printf("  [STATUS]   progress=%.0f/%.0f\n", inner["progress"], inner["total"])
            case "task.artifact":
                fmt.Printf("  [ARTIFACT] name=%v\n", inner["name"])
            case "task.done":
                fmt.Printf("  [DONE]     result=%v\n", inner["final_result"])
            case "task.error":
                fmt.Printf("  [ERROR]    %v\n", inner)
            }

            if eventType == "task.done" || eventType == "task.error" {
                break
            }
        }
    }
    ```
  </Tab>

  <Tab value="Java">
    ```java
    Map<String, Integer> eventTypes = new LinkedHashMap<>();
    String currentEvent = null;

    for (var it = resp.body().iterator(); it.hasNext(); ) {
        String line = it.next();
        if (line.startsWith("event: ")) {
            currentEvent = line.substring(7).trim();
        } else if (line.startsWith("data: ") && currentEvent != null) {
            var data = MAPPER.readTree(line.substring(6));
            eventTypes.merge(currentEvent, 1, Integer::sum);

            var inner = data.path("payload").path("payload");
            switch (currentEvent) {
                case "task.status" ->
                    System.out.println("  [STATUS]   progress=" + inner.path("progress") + "/" + inner.path("total"));
                case "task.artifact" ->
                    System.out.println("  [ARTIFACT] name=" + inner.path("name").asText());
                case "task.done" ->
                    System.out.println("  [DONE]     result=" + inner.path("final_result").asText());
                case "task.error" ->
                    System.out.println("  [ERROR]    " + inner);
            }

            if ("task.done".equals(currentEvent) || "task.error".equals(currentEvent)) break;
        }
    }
    ```
  </Tab>

  <Tab value="Python">
    ```python
    event_types: Counter[str] = Counter()

    async for event in event_source.aiter_sse():
        data = json.loads(event.data)
        event_types[event.event] += 1

        if event.event == "task.status":
            inner = data["payload"]["payload"]
            print(f"  [STATUS]   progress={inner['progress']}/{inner['total']}")
        elif event.event == "task.artifact":
            inner = data["payload"]["payload"]
            print(f"  [ARTIFACT] name={inner['name']}")
        elif event.event == "task.done":
            inner = data["payload"]["payload"]
            print(f"  [DONE]     result={inner['final_result']}")
        elif event.event == "task.error":
            print(f"  [ERROR]    {data['payload']['payload']}")

        if event.event in ("task.done", "task.error"):
            break

    print(f"\n--- Event Summary ---")
    for event_type, count in sorted(event_types.items()):
        print(f"  {event_type}: {count}")
    print(f"  Total: {sum(event_types.values())}")
    ```
  </Tab>

  <Tab value="TypeScript">
    ```typescript
    function parseSSE(chunk: string): Array<{ event: string; data: string }> {
      const events: Array<{ event: string; data: string }> = [];
      let currentEvent = "";
      let currentData = "";

      for (const line of chunk.split("\n")) {
        if (line.startsWith("event: ")) {
          currentEvent = line.slice(7).trim();
        } else if (line.startsWith("data: ")) {
          currentData = line.slice(6).trim();
        } else if (line === "" && currentEvent) {
          events.push({ event: currentEvent, data: currentData });
          currentEvent = "";
          currentData = "";
        }
      }
      return events;
    }

    const decoder = new TextDecoder();
    const counts: Record<string, number> = {};

    while (true) {
      const { done, value } = await reader.read();
      if (done) break;

      const events = parseSSE(decoder.decode(value, { stream: true }));
      for (const evt of events) {
        counts[evt.event] = (counts[evt.event] || 0) + 1;
        const payload = JSON.parse(evt.data);
        const inner = payload.payload?.payload ?? payload.payload;

        switch (evt.event) {
          case "task.status":
            console.log(`[STATUS]   progress=${inner.progress}/${inner.total} status=${inner.status}`);
            break;
          case "task.artifact":
            console.log(`[ARTIFACT] name=${inner.name} data=${JSON.stringify(inner.data)}`);
            break;
          case "task.done":
            console.log(`[DONE]     result=${inner.final_result}`);
            break;
          case "task.error":
            console.log(`[ERROR]    ${inner.message}`);
            break;
        }

        if (evt.event === "task.done" || evt.event === "task.error") {
          console.log("\n=== Event Summary ===");
          for (const [type, count] of Object.entries(counts)) {
            console.log(`  ${type}: ${count}`);
          }
          reader.cancel();
          return;
        }
      }
    }
    ```
  </Tab>
</Tabs>

A complete run prints the interim status, the artifact, the terminal result, and a summary:

```text
Connecting to SSE stream...
  [STATUS]   progress=1/3
  [STATUS]   progress=2/3
  [ARTIFACT] name=result.json
  [DONE]     result=completed

--- Event Summary ---
  task.artifact: 1
  task.done: 1
  task.status: 2
  Total: 4
```

## Step 3 — Cancel by disconnecting [#step-3--cancel-by-disconnecting]

There is no explicit "cancel" RPC for the caller — **cancelling a streaming task is just
closing the connection**. When the client stops reading and disconnects before the terminal
envelope, the gateway detects the dropped connection and sends a `stream_cancel` query to
the agent's virtual subscriber (on `_AGENTS_.agents/{agent_id}`, with the `stream_id` and a
10-second timeout). The virtual subscriber closes its HTTP SSE connection to the agent and
tears down the relay, so the agent stops doing work.

<Tabs groupId="language" items="['curl','C#','Go','Java','Python','TypeScript']">
  <Tab value="curl">
    ```bash
    # Read only the first events, then Ctrl-C (or pipe through head) to disconnect.
    # The gateway sends stream_cancel to the agent on disconnect.
    curl -N -X POST http://localhost:9090/a2a/slow-stream-agent-01 \
      -H 'Content-Type: application/json' \
      -H 'Accept: text/event-stream' \
      -d '{"jsonrpc":"2.0","id":1,"method":"message/stream",
           "params":{"message":{"parts":[{"text":"I will disconnect early"}]}}}' \
      | head -n 6
    ```
  </Tab>

  <Tab value="Go">
    ```go
    // Break out of the read loop early and close the body — the gateway
    // detects the disconnect and cancels the agent's stream.
    resp, _ := http.Post(kubemqURL+"/a2a/"+agentID, "application/json", bytes.NewReader(data))

    count := 0
    eventType := ""
    scanner := bufio.NewScanner(resp.Body)
    for scanner.Scan() {
        line := scanner.Text()
        if strings.HasPrefix(line, "event: ") {
            eventType = strings.TrimPrefix(line, "event: ")
        } else if strings.HasPrefix(line, "data: ") {
            count++
            var d map[string]interface{}
            json.Unmarshal([]byte(strings.TrimPrefix(line, "data: ")), &d)
            p, _ := d["payload"].(map[string]interface{})
            fmt.Printf("  Event %d: [%s] progress=%.0f\n", count, eventType, p["progress"])

            if count >= maxEvents {
                fmt.Printf("\nDisconnecting after %d events...\n", maxEvents)
                break
            }
        }
    }
    resp.Body.Close() // closing the connection triggers stream_cancel
    fmt.Println("KubeMQ will detect the disconnect and clean up the stream.")
    ```
  </Tab>

  <Tab value="Python">
    ```python
    """Disconnect from the SSE stream after a couple of events to cancel the task."""
    MAX_EVENTS = 2

    async with httpx.AsyncClient(timeout=30) as client:
        async with aconnect_sse(
            client, "POST", f"{KUBEMQ_URL}/a2a/{AGENT_ID}", json=payload
        ) as event_source:
            count = 0
            async for event in event_source.aiter_sse():
                count += 1
                data = json.loads(event.data)
                print(f"  Event {count}: [{event.event}] progress={data['payload'].get('progress')}")
                if count >= MAX_EVENTS:
                    print(f"\nDisconnecting after {MAX_EVENTS} events...")
                    break  # leaving the context closes the connection -> stream_cancel

    print("Client disconnected.")
    print("KubeMQ will detect the disconnect and clean up the stream.")
    ```
  </Tab>
</Tabs>

<Callout type="info">
  The gateway also closes the stream on its own when the idle timer fires
  (`MaxSSEIdleSeconds`, default 300s): it emits a `task.error` with code `-32001` and message
  `"stream idle timeout"`, then best-effort cancels the agent. See
  [SSE behavior](/aiway/a2a/guides/sse-behavior) for the full wire-level rules.
</Callout>

## Related [#related]

<Cards>
  <Card title="Streaming (SSE)" href="/aiway/a2a/streaming" description="The message/stream operation, envelope types, and SSE relay model." />

  <Card title="SSE behavior" href="/aiway/a2a/guides/sse-behavior" description="Keepalive, idle timeout, reconnection, and cancel-on-disconnect details." />

  <Card title="How it works" href="/aiway/a2a/architecture" description="The virtual-subscriber bridge that relays agent SSE events over the broker." />

  <Card title="Error handling" href="/aiway/a2a/error-handling" description="Transport vs application errors and the A2A error-code catalog." />
</Cards>
