KubeMQ
AiwayMCPTools

Queue Tools

MCP tools for durable point-to-point queue messaging — queue_send, queue_receive, and queue_peek over the tools/call method.

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

KubeMQ's queue messaging 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 (port 9090). Channels are created automatically on first use.

ToolOperationIdempotent
queue_sendEnqueue a message for durable, single-consumer deliveryNo
queue_receiveReceive and consume messages (destructive read)No
queue_peekInspect 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 for details.

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.

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

A channel that starts with the reserved prefix _AGENTS_. is rejected with a -32602 Invalid Params error. See Channel resolution.

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

ArgumentTypeRequiredDefaultDescription
channelstringYesQueue channel name. Must not start with the reserved prefix _AGENTS_.
bodystringYesMessage body (string)
metadatastringNo""Optional metadata
tagsobjectNo{}Optional key-value tags (string values)
delay_secondsintegerNo0Delay delivery by N seconds (0 = immediately visible)
expiration_secondsintegerNo0Message expiry in seconds (0 = no expiration)
max_receive_countintegerNo0Max receive attempts before dead-letter (0 = unlimited)
dead_letter_queuestringNo""Dead-letter queue channel name

Usage

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" }
      }
    }
  }'
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}");
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)
}
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();
    }
}
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()
}
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())
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
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(())
}
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)")
    }
}
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);

Response

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

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

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

ArgumentTypeRequiredDefaultDescription
channelstringYesQueue channel name
max_messagesintegerYes1Max messages to receive (clamped to 1–100)
wait_timeout_secondsintegerNo5Long-poll wait timeout in seconds (1–60)

Usage

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
      }
    }
  }'
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}");
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)
}
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();
    }
}
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()
}
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())
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
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(())
}
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)")
    }
}
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);

Response

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:

{
  "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

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

ArgumentTypeRequiredDefaultDescription
channelstringYesQueue channel name
max_messagesintegerYes1Max messages to peek (clamped to 1–100)

Usage

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
      }
    }
  }'
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}");
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)
}
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();
    }
}
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()
}
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())
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
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(())
}
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)")
    }
}
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);

Response

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

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

Receive vs peek

Behaviorqueue_receivequeue_peek
Removes messages from the queueYesNo
IdempotentNoYes
Long-poll waitwait_timeout_seconds (default 5s)Fixed 1s internally
Typical useProcessing work itemsMonitoring, debugging

Was this page helpful?

On this page