Skip to main content
Glama
ivanmartinezmorales

RabbitMQ MCP Server

MessageConsumer.ts4.08 kB
import { MCPTool, MCPInput } from "mcp-framework"; import { z } from "zod"; import RabbitMQConnection from "./RabbitMQConnection.js"; const ConsumeSchema = z.object({ action: z.enum(["consume", "get", "stop"]).describe("Action: 'consume' for continuous, 'get' for single message, 'stop' to cancel consumer"), queue_name: z.string().describe("Name of the queue to consume from"), consumer_tag: z.string().optional().describe("Consumer tag for identification (required for 'stop' action)"), count: z.number().optional().describe("Number of messages to get (only for 'get' action, default: 1)"), no_ack: z.boolean().optional().describe("Whether to auto-acknowledge messages (default: false)"), exclusive: z.boolean().optional().describe("Whether consumer should be exclusive (default: false)"), }); class MessageConsumer extends MCPTool { name = "message_consumer"; description = "Consume messages from RabbitMQ queues"; schema = ConsumeSchema; private static consumers: Map<string, any> = new Map(); async execute(input: MCPInput<this>) { const channel = RabbitMQConnection.getChannel(); if (!channel) { return "Error: No active RabbitMQ connection. Please connect first."; } try { switch (input.action) { case "get": const count = input.count ?? 1; const messages = []; for (let i = 0; i < count; i++) { const message = await channel.get(input.queue_name, { noAck: input.no_ack ?? false }); if (message) { messages.push({ content: message.content.toString(), properties: message.properties, fields: message.fields }); if (!input.no_ack) { channel.ack(message); } } else { break; } } return { message: `Retrieved ${messages.length} message(s) from queue '${input.queue_name}'`, messages, queue: input.queue_name }; case "consume": const consumerTag = input.consumer_tag ?? `consumer_${Date.now()}`; const consumer = await channel.consume( input.queue_name, (message) => { if (message) { console.log(`Received message: ${message.content.toString()}`); if (!input.no_ack) { channel.ack(message); } } }, { noAck: input.no_ack ?? false, exclusive: input.exclusive ?? false, consumerTag } ); MessageConsumer.consumers.set(consumerTag, consumer); return { message: `Started consuming from queue '${input.queue_name}'`, consumer_tag: consumer.consumerTag, queue: input.queue_name, note: "Messages will be logged to console. Use 'stop' action to cancel consumer." }; case "stop": if (!input.consumer_tag) { return "Error: consumer_tag is required for 'stop' action"; } if (!MessageConsumer.consumers.has(input.consumer_tag)) { return `Error: No consumer found with tag '${input.consumer_tag}'`; } await channel.cancel(input.consumer_tag); MessageConsumer.consumers.delete(input.consumer_tag); return { message: `Consumer '${input.consumer_tag}' stopped`, consumer_tag: input.consumer_tag }; default: return "Invalid action. Use 'consume', 'get', or 'stop'"; } } catch (error) { return `Error: ${error instanceof Error ? error.message : 'Unknown error'}`; } } static getActiveConsumers(): string[] { return Array.from(MessageConsumer.consumers.keys()); } } export default MessageConsumer;

Latest Blog Posts

MCP directory API

We provide all the information about MCP servers via our MCP API.

curl -X GET 'https://glama.ai/api/mcp/v1/servers/ivanmartinezmorales/rabbitmq-mcp'

If you have feedback or need assistance with the MCP directory API, please join our Discord server