123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211 |
- // Copyright 2007 - 2021, Alan Antonuk and the rabbitmq-c contributors.
- // SPDX-License-Identifier: mit
- #include <stdint.h>
- #include <stdio.h>
- #include <stdlib.h>
- #include <string.h>
- #include <rabbitmq-c/amqp.h>
- #include <rabbitmq-c/tcp_socket.h>
- #include <assert.h>
- #include "utils.h"
- int main(int argc, char *argv[]) {
- char const *hostname;
- int port, status;
- char const *exchange;
- char const *routingkey;
- char const *messagebody;
- amqp_socket_t *socket = NULL;
- amqp_connection_state_t conn;
- amqp_bytes_t reply_to_queue;
- if (argc < 6) { /* minimum number of mandatory arguments */
- fprintf(stderr,
- "usage:\namqp_rpc_sendstring_client host port exchange routingkey "
- "messagebody\n");
- return 1;
- }
- hostname = argv[1];
- port = atoi(argv[2]);
- exchange = argv[3];
- routingkey = argv[4];
- messagebody = argv[5];
- /*
- establish a channel that is used to connect RabbitMQ server
- */
- conn = amqp_new_connection();
- socket = amqp_tcp_socket_new(conn);
- if (!socket) {
- die("creating TCP socket");
- }
- status = amqp_socket_open(socket, hostname, port);
- if (status) {
- die("opening TCP socket");
- }
- die_on_amqp_error(amqp_login(conn, "/", 0, 131072, 0, AMQP_SASL_METHOD_PLAIN,
- "guest", "guest"),
- "Logging in");
- amqp_channel_open(conn, 1);
- die_on_amqp_error(amqp_get_rpc_reply(conn), "Opening channel");
- /*
- create private reply_to queue
- */
- {
- amqp_queue_declare_ok_t *r = amqp_queue_declare(
- conn, 1, amqp_empty_bytes, 0, 0, 0, 1, amqp_empty_table);
- die_on_amqp_error(amqp_get_rpc_reply(conn), "Declaring queue");
- reply_to_queue = amqp_bytes_malloc_dup(r->queue);
- if (reply_to_queue.bytes == NULL) {
- fprintf(stderr, "Out of memory while copying queue name");
- return 1;
- }
- }
- /*
- send the message
- */
- {
- /*
- set properties
- */
- amqp_basic_properties_t props;
- props._flags = AMQP_BASIC_CONTENT_TYPE_FLAG |
- AMQP_BASIC_DELIVERY_MODE_FLAG | AMQP_BASIC_REPLY_TO_FLAG |
- AMQP_BASIC_CORRELATION_ID_FLAG;
- props.content_type = amqp_cstring_bytes("text/plain");
- props.delivery_mode = 2; /* persistent delivery mode */
- props.reply_to = amqp_bytes_malloc_dup(reply_to_queue);
- if (props.reply_to.bytes == NULL) {
- fprintf(stderr, "Out of memory while copying queue name");
- return 1;
- }
- props.correlation_id = amqp_cstring_bytes("1");
- /*
- publish
- */
- die_on_error(amqp_basic_publish(conn, 1, amqp_cstring_bytes(exchange),
- amqp_cstring_bytes(routingkey), 0, 0,
- &props, amqp_cstring_bytes(messagebody)),
- "Publishing");
- amqp_bytes_free(props.reply_to);
- }
- /*
- wait an answer
- */
- {
- amqp_basic_consume(conn, 1, reply_to_queue, amqp_empty_bytes, 0, 1, 0,
- amqp_empty_table);
- die_on_amqp_error(amqp_get_rpc_reply(conn), "Consuming");
- amqp_bytes_free(reply_to_queue);
- {
- amqp_frame_t frame;
- int result;
- amqp_basic_deliver_t *d;
- amqp_basic_properties_t *p;
- size_t body_target;
- size_t body_received;
- for (;;) {
- amqp_maybe_release_buffers(conn);
- result = amqp_simple_wait_frame(conn, &frame);
- printf("Result: %d\n", result);
- if (result < 0) {
- break;
- }
- printf("Frame type: %u channel: %u\n", frame.frame_type, frame.channel);
- if (frame.frame_type != AMQP_FRAME_METHOD) {
- continue;
- }
- printf("Method: %s\n", amqp_method_name(frame.payload.method.id));
- if (frame.payload.method.id != AMQP_BASIC_DELIVER_METHOD) {
- continue;
- }
- d = (amqp_basic_deliver_t *)frame.payload.method.decoded;
- printf("Delivery: %u exchange: %.*s routingkey: %.*s\n",
- (unsigned)d->delivery_tag, (int)d->exchange.len,
- (char *)d->exchange.bytes, (int)d->routing_key.len,
- (char *)d->routing_key.bytes);
- result = amqp_simple_wait_frame(conn, &frame);
- if (result < 0) {
- break;
- }
- if (frame.frame_type != AMQP_FRAME_HEADER) {
- fprintf(stderr, "Expected header!");
- abort();
- }
- p = (amqp_basic_properties_t *)frame.payload.properties.decoded;
- if (p->_flags & AMQP_BASIC_CONTENT_TYPE_FLAG) {
- printf("Content-type: %.*s\n", (int)p->content_type.len,
- (char *)p->content_type.bytes);
- }
- printf("----\n");
- body_target = (size_t)frame.payload.properties.body_size;
- body_received = 0;
- while (body_received < body_target) {
- result = amqp_simple_wait_frame(conn, &frame);
- if (result < 0) {
- break;
- }
- if (frame.frame_type != AMQP_FRAME_BODY) {
- fprintf(stderr, "Expected body!");
- abort();
- }
- body_received += frame.payload.body_fragment.len;
- assert(body_received <= body_target);
- amqp_dump(frame.payload.body_fragment.bytes,
- frame.payload.body_fragment.len);
- }
- if (body_received != body_target) {
- /* Can only happen when amqp_simple_wait_frame returns <= 0 */
- /* We break here to close the connection */
- break;
- }
- /* everything was fine, we can quit now because we received the reply */
- break;
- }
- }
- }
- /*
- closing
- */
- die_on_amqp_error(amqp_channel_close(conn, 1, AMQP_REPLY_SUCCESS),
- "Closing channel");
- die_on_amqp_error(amqp_connection_close(conn, AMQP_REPLY_SUCCESS),
- "Closing connection");
- die_on_error(amqp_destroy_connection(conn), "Ending connection");
- return 0;
- }
|