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;