amq-rpc 0.1.0

RabbitMQ RPC library
Documentation
use serde_json::Value;
use async_trait::async_trait;
use amq_rpc::{AmqpConnection, RpcServer, CommandHandler, RpcResult, RpcError};

struct MathHandler;

#[async_trait]
impl CommandHandler for MathHandler {
    async fn handle(&self, command: &str, args: Vec<Value>) -> RpcResult<Value> {
        match command {
            "add" => {
                if args.len() != 2 {
                    return Err(RpcError::InvalidCommand {
                        message: "add requires exactly 2 arguments".to_string(),
                    });
                }

                let a = args[0].as_f64().unwrap_or(0.0);
                let b = args[1].as_f64().unwrap_or(0.0);
                Ok(Value::Number(serde_json::Number::from_f64(a + b).unwrap()))
            }
            "multiply" => {
                if args.len() != 2 {
                    return Err(RpcError::InvalidCommand {
                        message: "multiply requires exactly 2 arguments".to_string(),
                    });
                }

                let a = args[0].as_f64().unwrap_or(0.0);
                let b = args[1].as_f64().unwrap_or(0.0);
                Ok(Value::Number(serde_json::Number::from_f64(a * b).unwrap()))
            }
            _ => Err(RpcError::InvalidCommand {
                message: format!("Unknown math command: {}", command),
            }),
        }
    }
}

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    let connection = AmqpConnection::new("amqp://guest:guest@localhost:5672");
    let server = RpcServer::new(connection, "math_queue".to_string());

    server.register_handler("add", MathHandler).await;
    server.register_handler("multiply", MathHandler).await;

    println!("Starting RPC server on queue: math_queue");
    server.start().await?;

    tokio::signal::ctrl_c().await?;
    println!("Shutting down server...");
    server.close().await?;

    Ok(())
}