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.
| 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 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
| 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
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.closeuse 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
| 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
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.closeuse 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
| 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
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.closeuse 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
| 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
Was this page helpful?
Tools Overview
Map of all 15 KubeMQ MCP tools — 11 core messaging tools plus 4 agent-bridge tools — and the shared tools/call response shape.
Events Tools
Publish ephemeral and persistent events and read the events store as MCP tools — events_publish, events_store_publish, and events_store_read on KubeMQ.