# Queue Tools (/aiway/mcp/tools/queues)



The queue tools expose KubeMQ's durable, point-to-point queue messaging as MCP tools, so an AI model can enqueue work, consume it, and inspect a backlog without removing it.

## Overview [#overview]

KubeMQ's [queue messaging](/learn/queues) is a durable, at-least-once, single-consumer channel: a message is held until exactly one consumer receives it. The MCP connector surfaces three queue operations as tools, all invoked through the standard `tools/call` method against the single `POST /mcp` endpoint on the [shared HTTP server](/connectors/concepts/shared-http-server) (port `9090`). Channels are created automatically on first use.

| Tool            | Operation                                                     | Idempotent |
| --------------- | ------------------------------------------------------------- | ---------- |
| `queue_send`    | Enqueue a message for durable, single-consumer delivery       | No         |
| `queue_receive` | Receive and consume messages (destructive read)               | No         |
| `queue_peek`    | Inspect messages without removing them (non-destructive read) | Yes        |

The MCP connector is **enabled by default** — start kubemq-server and `/mcp` is live. To disable it, set `CONNECTORSMCP_ENABLE=false`. See [Configuration](/aiway/mcp/configuration) for details.

## How it works [#how-it-works]

A `tools/call` request flows through the MCP connector, which translates the tool arguments into a native KubeMQ queue operation against the broker.

<Mermaid
  chart="`
graph LR
AI[&#x22;AI model / agent&#x22;]
MCP[&#x22;MCP connector<br/>/mcp&#x22;]
Q{{&#x22;Queue channel<br/>example-queue&#x22;}}
BROKER[&#x22;KubeMQ broker&#x22;]

AI -- &#x22;tools/call queue_send&#x22; --> MCP
MCP -- enqueue --> Q
Q --> BROKER
BROKER -. &#x22;tools/call queue_receive&#x22; .-> MCP
MCP -. messages .-> AI

class AI external
class MCP aiway
class Q queue
class BROKER broker
`"
/>

*An AI model sends and receives queue messages through the MCP connector, which bridges to the KubeMQ broker.*

<Callout type="warn">
  A `channel` that starts with the reserved prefix `_AGENTS_.` is rejected with a `-32602` Invalid Params error. See [Channel resolution](/aiway/mcp/guides/channel-resolution).
</Callout>

## queue\_send [#queue_send]

> Send a message to a KubeMQ queue channel. Use for reliable, persistent messaging with at-least-once delivery. Each call enqueues a new message (not idempotent).

The channel is created automatically if it does not exist. Optional policy fields delay visibility (`delay_seconds`), set a time-to-live (`expiration_seconds`), and route poison messages to a dead-letter queue (`max_receive_count` plus `dead_letter_queue`). Both `channel` and `body` are required — omitting either returns a `-32602` Invalid Params error.

### Input schema [#input-schema]

| Argument             | Type    | Required | Default | Description                                                             |
| -------------------- | ------- | -------- | ------- | ----------------------------------------------------------------------- |
| `channel`            | string  | Yes      | —       | Queue channel name. Must not start with the reserved prefix `_AGENTS_.` |
| `body`               | string  | Yes      | —       | Message body (string)                                                   |
| `metadata`           | string  | No       | `""`    | Optional metadata                                                       |
| `tags`               | object  | No       | `{}`    | Optional key-value tags (string values)                                 |
| `delay_seconds`      | integer | No       | `0`     | Delay delivery by N seconds (`0` = immediately visible)                 |
| `expiration_seconds` | integer | No       | `0`     | Message expiry in seconds (`0` = no expiration)                         |
| `max_receive_count`  | integer | No       | `0`     | Max receive attempts before dead-letter (`0` = unlimited)               |
| `dead_letter_queue`  | string  | No       | `""`    | Dead-letter queue channel name                                          |

### Usage [#usage]

<Tabs groupId="language" items="['curl','C#','Go','Java','Kotlin','Python','Ruby','Rust','Swift','TypeScript']">
  <Tab value="curl">
    ```bash
    curl -X POST http://localhost:9090/mcp \
      -H 'Content-Type: application/json' \
      -H 'Accept: application/json, text/event-stream' \
      -d '{
        "jsonrpc": "2.0",
        "id": 2,
        "method": "tools/call",
        "params": {
          "name": "queue_send",
          "arguments": {
            "channel": "example-queue",
            "body": "Hello from MCP",
            "metadata": "example-metadata",
            "tags": { "env": "dev", "source": "mcp-example" }
          }
        }
      }'
    ```
  </Tab>

  <Tab value="C#">
    ```csharp
    using ModelContextProtocol.Client;

    var url = Environment.GetEnvironmentVariable("KUBEMQ_MCP_URL") ?? "http://localhost:9090";
    var transport = new HttpClientTransport(new HttpClientTransportOptions { Endpoint = new Uri($"{url}/mcp") });
    await using var client = await McpClientFactory.CreateAsync(transport);

    var result = await client.CallToolAsync("queue_send", new Dictionary<string, object>
    {
        ["channel"] = "example-queue",
        ["body"] = "Hello from C# MCP",
        ["metadata"] = "example-metadata",
        ["tags"] = new Dictionary<string, string> { ["env"] = "dev", ["source"] = "mcp-example" },
    });

    Console.WriteLine($"Result: {result}");
    ```
  </Tab>

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

    import (
        "context"
        "fmt"
        "log"
        "os"

        "github.com/mark3labs/mcp-go/client"
        "github.com/mark3labs/mcp-go/mcp"
    )

    func main() {
        url := os.Getenv("KUBEMQ_MCP_URL")
        if url == "" {
            url = "http://localhost:9090"
        }

        c, err := client.NewStreamableHttpClient(url + "/mcp")
        if err != nil {
            log.Fatal(err)
        }
        defer c.Close()

        ctx := context.Background()
        if err := c.Start(ctx); err != nil {
            log.Fatal(err)
        }

        result, err := c.CallTool(ctx, mcp.CallToolRequest{
            Params: mcp.CallToolParams{
                Name: "queue_send",
                Arguments: map[string]any{
                    "channel":  "example-queue",
                    "body":     "Hello from Go MCP",
                    "metadata": "example-metadata",
                    "tags":     map[string]any{"env": "dev", "source": "mcp-example"},
                },
            },
        })
        if err != nil {
            log.Fatal(err)
        }

        fmt.Printf("Result: %+v\n", result)
    }
    ```
  </Tab>

  <Tab value="Java">
    ```java
    import io.modelcontextprotocol.sdk.McpClient;
    import io.modelcontextprotocol.sdk.client.transport.HttpClientStreamableHttpTransport;
    import io.modelcontextprotocol.spec.McpSchema.CallToolRequest;

    import java.util.Map;

    public class QueueSend {
        public static void main(String[] args) {
            String url = System.getenv().getOrDefault("KUBEMQ_MCP_URL", "http://localhost:9090");
            var transport = HttpClientStreamableHttpTransport.builder(url).endpoint("/mcp").build();
            var client = McpClient.sync(transport).build();
            client.initialize();

            var result = client.callTool(new CallToolRequest(
                "queue_send",
                Map.of(
                    "channel", "example-queue",
                    "body", "Hello from Java MCP",
                    "metadata", "example-metadata",
                    "tags", Map.of("env", "dev", "source", "mcp-example")
                )
            ));
            System.out.println(result);

            client.closeGracefully();
        }
    }
    ```
  </Tab>

  <Tab value="Kotlin">
    ```kotlin
    import io.modelcontextprotocol.kotlin.sdk.Implementation
    import io.modelcontextprotocol.kotlin.sdk.client.Client
    import io.modelcontextprotocol.kotlin.sdk.client.StreamableHttpClientTransport
    import io.ktor.client.*
    import io.ktor.client.plugins.sse.*
    import kotlinx.coroutines.runBlocking

    fun main() = runBlocking {
        val url = System.getenv("KUBEMQ_MCP_URL") ?: "http://localhost:9090"

        val httpClient = HttpClient { install(SSE) }
        val transport = StreamableHttpClientTransport(client = httpClient, url = "$url/mcp")
        val client = Client(clientInfo = Implementation(name = "kubemq-mcp-kotlin-example", version = "1.0.0"))
        client.connect(transport)

        val result = client.callTool("queue_send", mapOf(
            "channel" to "example-queue",
            "body" to "Hello from Kotlin MCP",
            "metadata" to "example-metadata",
            "tags" to mapOf("env" to "dev", "source" to "mcp-example")
        ))

        println("Result: $result")

        client.close()
        httpClient.close()
    }
    ```
  </Tab>

  <Tab value="Python">
    ```python
    import asyncio
    import os

    from mcp.client.streamable_http import streamablehttp_client
    from mcp import ClientSession

    KUBEMQ_MCP_URL = os.environ.get("KUBEMQ_MCP_URL", "http://localhost:9090")


    async def main():
        async with streamablehttp_client(f"{KUBEMQ_MCP_URL}/mcp") as (read, write, _):
            async with ClientSession(read, write) as session:
                await session.initialize()

                result = await session.call_tool("queue_send", {
                    "channel": "example-queue",
                    "body": "Hello from Python MCP",
                    "metadata": "example-metadata",
                    "tags": {"env": "dev", "source": "mcp-example"},
                })

                print(f"IsError: {result.isError}")
                for content in result.content:
                    print(f"Result: {content.text}")


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

  <Tab value="Ruby">
    ```ruby
    require "mcp"

    url = ENV.fetch("KUBEMQ_MCP_URL", "http://localhost:9090")

    client = MCP::Client.new(
      transport: MCP::Transport::StreamableHTTP.new("#{url}/mcp"),
      name: "kubemq-mcp-ruby-example",
      version: "1.0.0"
    )
    client.initialize_handshake

    result = client.call_tool("queue_send", {
      "channel" => "example-queue",
      "body" => "Hello from Ruby MCP",
      "metadata" => "example-metadata",
      "tags" => { "env" => "dev", "source" => "mcp-example" },
    })

    puts "Result: #{result}"

    client.close
    ```
  </Tab>

  <Tab value="Rust">
    ```rust
    use rmcp::transport::streamable_http::StreamableHttpClientTransport;
    use rmcp::service::RunService;
    use serde_json::json;

    #[tokio::main]
    async fn main() -> anyhow::Result<()> {
        let url = std::env::var("KUBEMQ_MCP_URL")
            .unwrap_or_else(|_| "http://localhost:9090".to_string());

        let transport = StreamableHttpClientTransport::from_uri(format!("{url}/mcp"))?;
        let client = ().serve(transport).await?;

        let result = client.call_tool("queue_send", json!({
            "channel": "example-queue",
            "body": "Hello from Rust MCP",
            "metadata": "example-metadata",
            "tags": {"env": "dev", "source": "mcp-example"}
        })).await?;

        println!("Result: {result:#?}");
        Ok(())
    }
    ```
  </Tab>

  <Tab value="Swift">
    ```swift
    import Foundation
    import MCP

    @main
    struct QueueSend {
        static func main() async throws {
            let url = ProcessInfo.processInfo.environment["KUBEMQ_MCP_URL"] ?? "http://localhost:9090"

            let transport = HTTPClientTransport(endpoint: URL(string: "\(url)/mcp")!, streaming: true)
            let client = Client(name: "kubemq-mcp-swift-example", version: "1.0.0")
            try await client.connect(transport: transport)

            let result = try await client.callTool("queue_send", arguments: [
                "channel": "example-queue",
                "body": "Hello from Swift MCP",
                "metadata": "example-metadata",
                "tags": ["env": "dev", "source": "mcp-example"],
            ])

            print("Result: \(result)")
        }
    }
    ```
  </Tab>

  <Tab value="TypeScript">
    ```typescript
    import { Client } from "@modelcontextprotocol/sdk/client/index.js";
    import { StreamableHTTPClientTransport } from "@modelcontextprotocol/sdk/client/streamableHttp.js";

    const KUBEMQ_MCP_URL = process.env.KUBEMQ_MCP_URL || "http://localhost:9090";

    async function main() {
      const transport = new StreamableHTTPClientTransport(
        new URL(`${KUBEMQ_MCP_URL}/mcp`)
      );
      const client = new Client({ name: "kubemq-mcp-ts-example", version: "1.0.0" });
      await client.connect(transport);

      const result = await client.callTool({
        name: "queue_send",
        arguments: {
          channel: "example-queue",
          body: "Hello from TypeScript MCP",
          metadata: "example-metadata",
          tags: { env: "dev", source: "mcp-example" },
        },
      });
      console.log(JSON.stringify(result, null, 2));

      await client.close();
    }

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

### Response [#response]

The tool result wraps a JSON string in the standard MCP `content[]` / `isError` envelope. The text payload reports the enqueue result:

```json
{
  "content": [
    {
      "type": "text",
      "text": "{\"message_id\":\"a1b2c3...\",\"sent_at\":\"2026-06-08T10:00:00Z\",\"is_error\":false}"
    }
  ],
  "isError": false
}
```

## queue\_receive [#queue_receive]

> Receive messages from a KubeMQ queue channel. Messages are auto-acknowledged on receipt (destructive read). Not idempotent — failed processing requires re-enqueue.

This is a destructive read: returned messages are removed from the queue, so each message is delivered to exactly one consumer. Set `max_messages` to drain a batch (clamped to the `1–100` range), and `wait_timeout_seconds` to long-poll for messages that have not yet arrived. Receiving from a non-existent channel is not an error — it returns an empty `messages` array. Because the read is not idempotent, do not blindly retry: a retry may consume *additional* messages rather than re-fetch the same ones.

### Input schema [#input-schema-1]

| Argument               | Type    | Required | Default | Description                                  |
| ---------------------- | ------- | -------- | ------- | -------------------------------------------- |
| `channel`              | string  | Yes      | —       | Queue channel name                           |
| `max_messages`         | integer | Yes      | `1`     | Max messages to receive (clamped to `1–100`) |
| `wait_timeout_seconds` | integer | No       | `5`     | Long-poll wait timeout in seconds (`1–60`)   |

### Usage [#usage-1]

<Tabs groupId="language" items="['curl','C#','Go','Java','Kotlin','Python','Ruby','Rust','Swift','TypeScript']">
  <Tab value="curl">
    ```bash
    curl -X POST http://localhost:9090/mcp \
      -H 'Content-Type: application/json' \
      -H 'Accept: application/json, text/event-stream' \
      -d '{
        "jsonrpc": "2.0",
        "id": 3,
        "method": "tools/call",
        "params": {
          "name": "queue_receive",
          "arguments": {
            "channel": "example-queue",
            "max_messages": 5
          }
        }
      }'
    ```
  </Tab>

  <Tab value="C#">
    ```csharp
    using ModelContextProtocol.Client;

    var url = Environment.GetEnvironmentVariable("KUBEMQ_MCP_URL") ?? "http://localhost:9090";
    var transport = new HttpClientTransport(new HttpClientTransportOptions { Endpoint = new Uri($"{url}/mcp") });
    await using var client = await McpClientFactory.CreateAsync(transport);

    var result = await client.CallToolAsync("queue_receive", new Dictionary<string, object>
    {
        ["channel"] = "example-queue",
        ["max_messages"] = 5,
    });

    Console.WriteLine($"Result: {result}");
    ```
  </Tab>

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

    import (
        "context"
        "fmt"
        "log"
        "os"

        "github.com/mark3labs/mcp-go/client"
        "github.com/mark3labs/mcp-go/mcp"
    )

    func main() {
        url := os.Getenv("KUBEMQ_MCP_URL")
        if url == "" {
            url = "http://localhost:9090"
        }

        c, err := client.NewStreamableHttpClient(url + "/mcp")
        if err != nil {
            log.Fatal(err)
        }
        defer c.Close()

        ctx := context.Background()
        if err := c.Start(ctx); err != nil {
            log.Fatal(err)
        }

        result, err := c.CallTool(ctx, mcp.CallToolRequest{
            Params: mcp.CallToolParams{
                Name: "queue_receive",
                Arguments: map[string]any{
                    "channel":      "example-queue",
                    "max_messages": 5,
                },
            },
        })
        if err != nil {
            log.Fatal(err)
        }

        fmt.Printf("Result: %+v\n", result)
    }
    ```
  </Tab>

  <Tab value="Java">
    ```java
    import io.modelcontextprotocol.sdk.McpClient;
    import io.modelcontextprotocol.sdk.client.transport.HttpClientStreamableHttpTransport;
    import io.modelcontextprotocol.spec.McpSchema.CallToolRequest;

    import java.util.Map;

    public class QueueReceive {
        public static void main(String[] args) {
            String url = System.getenv().getOrDefault("KUBEMQ_MCP_URL", "http://localhost:9090");
            var transport = HttpClientStreamableHttpTransport.builder(url).endpoint("/mcp").build();
            var client = McpClient.sync(transport).build();
            client.initialize();

            var result = client.callTool(new CallToolRequest(
                "queue_receive",
                Map.of(
                    "channel", "example-queue",
                    "max_messages", 5
                )
            ));
            System.out.println(result);

            client.closeGracefully();
        }
    }
    ```
  </Tab>

  <Tab value="Kotlin">
    ```kotlin
    import io.modelcontextprotocol.kotlin.sdk.Implementation
    import io.modelcontextprotocol.kotlin.sdk.client.Client
    import io.modelcontextprotocol.kotlin.sdk.client.StreamableHttpClientTransport
    import io.ktor.client.*
    import io.ktor.client.plugins.sse.*
    import kotlinx.coroutines.runBlocking

    fun main() = runBlocking {
        val url = System.getenv("KUBEMQ_MCP_URL") ?: "http://localhost:9090"

        val httpClient = HttpClient { install(SSE) }
        val transport = StreamableHttpClientTransport(client = httpClient, url = "$url/mcp")
        val client = Client(clientInfo = Implementation(name = "kubemq-mcp-kotlin-example", version = "1.0.0"))
        client.connect(transport)

        val result = client.callTool("queue_receive", mapOf(
            "channel" to "example-queue",
            "max_messages" to 5
        ))

        println("Result: $result")

        client.close()
        httpClient.close()
    }
    ```
  </Tab>

  <Tab value="Python">
    ```python
    import asyncio
    import os

    from mcp.client.streamable_http import streamablehttp_client
    from mcp import ClientSession

    KUBEMQ_MCP_URL = os.environ.get("KUBEMQ_MCP_URL", "http://localhost:9090")


    async def main():
        async with streamablehttp_client(f"{KUBEMQ_MCP_URL}/mcp") as (read, write, _):
            async with ClientSession(read, write) as session:
                await session.initialize()

                result = await session.call_tool("queue_receive", {
                    "channel": "example-queue",
                    "max_messages": 5,
                })

                print(f"IsError: {result.isError}")
                for content in result.content:
                    print(f"Result: {content.text}")


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

  <Tab value="Ruby">
    ```ruby
    require "mcp"

    url = ENV.fetch("KUBEMQ_MCP_URL", "http://localhost:9090")

    client = MCP::Client.new(
      transport: MCP::Transport::StreamableHTTP.new("#{url}/mcp"),
      name: "kubemq-mcp-ruby-example",
      version: "1.0.0"
    )
    client.initialize_handshake

    result = client.call_tool("queue_receive", {
      "channel" => "example-queue",
      "max_messages" => 5,
    })

    puts "Result: #{result}"

    client.close
    ```
  </Tab>

  <Tab value="Rust">
    ```rust
    use rmcp::transport::streamable_http::StreamableHttpClientTransport;
    use rmcp::service::RunService;
    use serde_json::json;

    #[tokio::main]
    async fn main() -> anyhow::Result<()> {
        let url = std::env::var("KUBEMQ_MCP_URL")
            .unwrap_or_else(|_| "http://localhost:9090".to_string());

        let transport = StreamableHttpClientTransport::from_uri(format!("{url}/mcp"))?;
        let client = ().serve(transport).await?;

        let result = client.call_tool("queue_receive", json!({
            "channel": "example-queue",
            "max_messages": 5
        })).await?;

        println!("Result: {result:#?}");
        Ok(())
    }
    ```
  </Tab>

  <Tab value="Swift">
    ```swift
    import Foundation
    import MCP

    @main
    struct QueueReceive {
        static func main() async throws {
            let url = ProcessInfo.processInfo.environment["KUBEMQ_MCP_URL"] ?? "http://localhost:9090"

            let transport = HTTPClientTransport(endpoint: URL(string: "\(url)/mcp")!, streaming: true)
            let client = Client(name: "kubemq-mcp-swift-example", version: "1.0.0")
            try await client.connect(transport: transport)

            let result = try await client.callTool("queue_receive", arguments: [
                "channel": "example-queue",
                "max_messages": 5,
            ])

            print("Result: \(result)")
        }
    }
    ```
  </Tab>

  <Tab value="TypeScript">
    ```typescript
    import { Client } from "@modelcontextprotocol/sdk/client/index.js";
    import { StreamableHTTPClientTransport } from "@modelcontextprotocol/sdk/client/streamableHttp.js";

    const KUBEMQ_MCP_URL = process.env.KUBEMQ_MCP_URL || "http://localhost:9090";

    async function main() {
      const transport = new StreamableHTTPClientTransport(
        new URL(`${KUBEMQ_MCP_URL}/mcp`)
      );
      const client = new Client({ name: "kubemq-mcp-ts-example", version: "1.0.0" });
      await client.connect(transport);

      const result = await client.callTool({
        name: "queue_receive",
        arguments: {
          channel: "example-queue",
          max_messages: 5,
        },
      });
      console.log(JSON.stringify(result, null, 2));

      await client.close();
    }

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

### Response [#response-1]

The text payload is a JSON object holding the received messages and a count. Each message carries its `id`, `channel`, `body`, `metadata`, and (when present) `timestamp`, `sequence`, and `tags`:

```json
{
  "content": [
    {
      "type": "text",
      "text": "{\"messages\":[{\"id\":\"...\",\"channel\":\"example-queue\",\"body\":\"Hello from MCP\",\"metadata\":\"example-metadata\",\"tags\":{\"env\":\"dev\",\"source\":\"mcp-example\"}}],\"messages_count\":1,\"is_error\":false}"
    }
  ],
  "isError": false
}
```

## queue\_peek [#queue_peek]

> Peek at messages in a KubeMQ queue without removing them. Idempotent — does not modify queue state.

Peek is a non-destructive, idempotent read: inspected messages remain in the queue and can still be consumed later by `queue_receive`. Use it for monitoring, debugging, or letting an agent reason about pending work before deciding to consume it. Set `max_messages` to control how many messages to inspect (clamped to `1–100`). As with receive, peeking a non-existent channel returns an empty result rather than an error.

### Input schema [#input-schema-2]

| Argument       | Type    | Required | Default | Description                               |
| -------------- | ------- | -------- | ------- | ----------------------------------------- |
| `channel`      | string  | Yes      | —       | Queue channel name                        |
| `max_messages` | integer | Yes      | `1`     | Max messages to peek (clamped to `1–100`) |

### Usage [#usage-2]

<Tabs groupId="language" items="['curl','C#','Go','Java','Kotlin','Python','Ruby','Rust','Swift','TypeScript']">
  <Tab value="curl">
    ```bash
    curl -X POST http://localhost:9090/mcp \
      -H 'Content-Type: application/json' \
      -H 'Accept: application/json, text/event-stream' \
      -d '{
        "jsonrpc": "2.0",
        "id": 4,
        "method": "tools/call",
        "params": {
          "name": "queue_peek",
          "arguments": {
            "channel": "example-queue",
            "max_messages": 5
          }
        }
      }'
    ```
  </Tab>

  <Tab value="C#">
    ```csharp
    using ModelContextProtocol.Client;

    var url = Environment.GetEnvironmentVariable("KUBEMQ_MCP_URL") ?? "http://localhost:9090";
    var transport = new HttpClientTransport(new HttpClientTransportOptions { Endpoint = new Uri($"{url}/mcp") });
    await using var client = await McpClientFactory.CreateAsync(transport);

    var result = await client.CallToolAsync("queue_peek", new Dictionary<string, object>
    {
        ["channel"] = "example-queue",
        ["max_messages"] = 5,
    });

    Console.WriteLine($"Result: {result}");
    ```
  </Tab>

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

    import (
        "context"
        "fmt"
        "log"
        "os"

        "github.com/mark3labs/mcp-go/client"
        "github.com/mark3labs/mcp-go/mcp"
    )

    func main() {
        url := os.Getenv("KUBEMQ_MCP_URL")
        if url == "" {
            url = "http://localhost:9090"
        }

        c, err := client.NewStreamableHttpClient(url + "/mcp")
        if err != nil {
            log.Fatal(err)
        }
        defer c.Close()

        ctx := context.Background()
        if err := c.Start(ctx); err != nil {
            log.Fatal(err)
        }

        result, err := c.CallTool(ctx, mcp.CallToolRequest{
            Params: mcp.CallToolParams{
                Name: "queue_peek",
                Arguments: map[string]any{
                    "channel":      "example-queue",
                    "max_messages": 5,
                },
            },
        })
        if err != nil {
            log.Fatal(err)
        }

        fmt.Printf("Result: %+v\n", result)
    }
    ```
  </Tab>

  <Tab value="Java">
    ```java
    import io.modelcontextprotocol.sdk.McpClient;
    import io.modelcontextprotocol.sdk.client.transport.HttpClientStreamableHttpTransport;
    import io.modelcontextprotocol.spec.McpSchema.CallToolRequest;

    import java.util.Map;

    public class QueuePeek {
        public static void main(String[] args) {
            String url = System.getenv().getOrDefault("KUBEMQ_MCP_URL", "http://localhost:9090");
            var transport = HttpClientStreamableHttpTransport.builder(url).endpoint("/mcp").build();
            var client = McpClient.sync(transport).build();
            client.initialize();

            var result = client.callTool(new CallToolRequest(
                "queue_peek",
                Map.of(
                    "channel", "example-queue",
                    "max_messages", 5
                )
            ));
            System.out.println(result);

            client.closeGracefully();
        }
    }
    ```
  </Tab>

  <Tab value="Kotlin">
    ```kotlin
    import io.modelcontextprotocol.kotlin.sdk.Implementation
    import io.modelcontextprotocol.kotlin.sdk.client.Client
    import io.modelcontextprotocol.kotlin.sdk.client.StreamableHttpClientTransport
    import io.ktor.client.*
    import io.ktor.client.plugins.sse.*
    import kotlinx.coroutines.runBlocking

    fun main() = runBlocking {
        val url = System.getenv("KUBEMQ_MCP_URL") ?: "http://localhost:9090"

        val httpClient = HttpClient { install(SSE) }
        val transport = StreamableHttpClientTransport(client = httpClient, url = "$url/mcp")
        val client = Client(clientInfo = Implementation(name = "kubemq-mcp-kotlin-example", version = "1.0.0"))
        client.connect(transport)

        val result = client.callTool("queue_peek", mapOf(
            "channel" to "example-queue",
            "max_messages" to 5
        ))

        println("Result: $result")

        client.close()
        httpClient.close()
    }
    ```
  </Tab>

  <Tab value="Python">
    ```python
    import asyncio
    import os

    from mcp.client.streamable_http import streamablehttp_client
    from mcp import ClientSession

    KUBEMQ_MCP_URL = os.environ.get("KUBEMQ_MCP_URL", "http://localhost:9090")


    async def main():
        async with streamablehttp_client(f"{KUBEMQ_MCP_URL}/mcp") as (read, write, _):
            async with ClientSession(read, write) as session:
                await session.initialize()

                result = await session.call_tool("queue_peek", {
                    "channel": "example-queue",
                    "max_messages": 5,
                })

                print(f"IsError: {result.isError}")
                for content in result.content:
                    print(f"Result: {content.text}")


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

  <Tab value="Ruby">
    ```ruby
    require "mcp"

    url = ENV.fetch("KUBEMQ_MCP_URL", "http://localhost:9090")

    client = MCP::Client.new(
      transport: MCP::Transport::StreamableHTTP.new("#{url}/mcp"),
      name: "kubemq-mcp-ruby-example",
      version: "1.0.0"
    )
    client.initialize_handshake

    result = client.call_tool("queue_peek", {
      "channel" => "example-queue",
      "max_messages" => 5,
    })

    puts "Result: #{result}"

    client.close
    ```
  </Tab>

  <Tab value="Rust">
    ```rust
    use rmcp::transport::streamable_http::StreamableHttpClientTransport;
    use rmcp::service::RunService;
    use serde_json::json;

    #[tokio::main]
    async fn main() -> anyhow::Result<()> {
        let url = std::env::var("KUBEMQ_MCP_URL")
            .unwrap_or_else(|_| "http://localhost:9090".to_string());

        let transport = StreamableHttpClientTransport::from_uri(format!("{url}/mcp"))?;
        let client = ().serve(transport).await?;

        let result = client.call_tool("queue_peek", json!({
            "channel": "example-queue",
            "max_messages": 5
        })).await?;

        println!("Result: {result:#?}");
        Ok(())
    }
    ```
  </Tab>

  <Tab value="Swift">
    ```swift
    import Foundation
    import MCP

    @main
    struct QueuePeek {
        static func main() async throws {
            let url = ProcessInfo.processInfo.environment["KUBEMQ_MCP_URL"] ?? "http://localhost:9090"

            let transport = HTTPClientTransport(endpoint: URL(string: "\(url)/mcp")!, streaming: true)
            let client = Client(name: "kubemq-mcp-swift-example", version: "1.0.0")
            try await client.connect(transport: transport)

            let result = try await client.callTool("queue_peek", arguments: [
                "channel": "example-queue",
                "max_messages": 5,
            ])

            print("Result: \(result)")
        }
    }
    ```
  </Tab>

  <Tab value="TypeScript">
    ```typescript
    import { Client } from "@modelcontextprotocol/sdk/client/index.js";
    import { StreamableHTTPClientTransport } from "@modelcontextprotocol/sdk/client/streamableHttp.js";

    const KUBEMQ_MCP_URL = process.env.KUBEMQ_MCP_URL || "http://localhost:9090";

    async function main() {
      const transport = new StreamableHTTPClientTransport(
        new URL(`${KUBEMQ_MCP_URL}/mcp`)
      );
      const client = new Client({ name: "kubemq-mcp-ts-example", version: "1.0.0" });
      await client.connect(transport);

      const result = await client.callTool({
        name: "queue_peek",
        arguments: {
          channel: "example-queue",
          max_messages: 5,
        },
      });
      console.log(JSON.stringify(result, null, 2));

      await client.close();
    }

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

### Response [#response-2]

Same envelope as `queue_receive`, but the messages stay in the queue. Peeked messages omit `tags` and report `messages_count`:

```json
{
  "content": [
    {
      "type": "text",
      "text": "{\"messages\":[{\"id\":\"...\",\"channel\":\"example-queue\",\"body\":\"Hello from MCP\",\"metadata\":\"example-metadata\"}],\"messages_count\":1}"
    }
  ],
  "isError": false
}
```

## Receive vs peek [#receive-vs-peek]

| Behavior                        | `queue_receive`                     | `queue_peek`          |
| ------------------------------- | ----------------------------------- | --------------------- |
| Removes messages from the queue | Yes                                 | No                    |
| Idempotent                      | No                                  | Yes                   |
| Long-poll wait                  | `wait_timeout_seconds` (default 5s) | Fixed 1s internally   |
| Typical use                     | Processing work items               | Monitoring, debugging |

## Related [#related]

<Cards>
  <Card title="Queues concept" href="/learn/queues" description="The underlying durable, point-to-point messaging model." />

  <Card title="Event tools" href="/aiway/mcp/tools/events" description="Pub/sub and events-store tools — events_publish and more." />

  <Card title="Tools reference" href="/aiway/mcp/reference/tools-reference" description="Full 15-tool catalog with parameters and response shapes." />

  <Card title="Error codes" href="/aiway/mcp/reference/error-codes" description="JSON-RPC error codes and isError semantics for tool results." />
</Cards>
