2 * Advanced Message Queuing Protocol (AMQP) 0-9-1
3 * Copyright (c) 2020 Andriy Gelman
5 * This file is part of FFmpeg.
7 * FFmpeg is free software; you can redistribute it and/or
8 * modify it under the terms of the GNU Lesser General Public
9 * License as published by the Free Software Foundation; either
10 * version 2.1 of the License, or (at your option) any later version.
12 * FFmpeg is distributed in the hope that it will be useful,
13 * but WITHOUT ANY WARRANTY; without even the implied warranty of
14 * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU
15 * Lesser General Public License for more details.
17 * You should have received a copy of the GNU Lesser General Public
18 * License along with FFmpeg; if not, write to the Free Software
19 * Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
23 #include <amqp_tcp_socket.h>
26 #include "libavutil/avstring.h"
27 #include "libavutil/opt.h"
28 #include "libavutil/time.h"
31 #include "urldecode.h"
33 typedef struct AMQPContext {
35 amqp_connection_state_t conn;
36 amqp_socket_t *socket;
38 const char *routing_key;
40 int64_t connection_timeout;
41 int pkt_size_overflow;
45 #define DEFAULT_CHANNEL 1
47 #define OFFSET(x) offsetof(AMQPContext, x)
48 #define D AV_OPT_FLAG_DECODING_PARAM
49 #define E AV_OPT_FLAG_ENCODING_PARAM
50 static const AVOption options[] = {
51 { "pkt_size", "Maximum send/read packet size", OFFSET(pkt_size), AV_OPT_TYPE_INT, { .i64 = 131072 }, 4096, INT_MAX, .flags = D | E },
52 { "exchange", "Exchange to send/read packets", OFFSET(exchange), AV_OPT_TYPE_STRING, { .str = "amq.direct" }, 0, 0, .flags = D | E },
53 { "routing_key", "Key to filter streams", OFFSET(routing_key), AV_OPT_TYPE_STRING, { .str = "amqp" }, 0, 0, .flags = D | E },
54 { "connection_timeout", "Initial connection timeout", OFFSET(connection_timeout), AV_OPT_TYPE_DURATION, { .i64 = -1 }, -1, INT64_MAX, .flags = D | E},
58 static int amqp_proto_open(URLContext *h, const char *uri, int flags)
61 char hostname[STR_LEN], credentials[STR_LEN];
63 const char *user, *password = NULL;
64 const char *user_decoded, *password_decoded;
66 amqp_rpc_reply_t broker_reply;
67 struct timeval tval = { 0 };
69 AMQPContext *s = h->priv_data;
72 h->max_packet_size = s->pkt_size;
74 av_url_split(NULL, 0, credentials, sizeof(credentials),
75 hostname, sizeof(hostname), &port, NULL, 0, uri);
80 if (hostname[0] == '\0' || port <= 0 || port > 65535 ) {
81 av_log(h, AV_LOG_ERROR, "Invalid hostname/port\n");
82 return AVERROR(EINVAL);
85 p = strchr(credentials, ':');
91 if (!password || *password == '\0')
94 password_decoded = ff_urldecode(password, 0);
95 if (!password_decoded)
96 return AVERROR(ENOMEM);
102 user_decoded = ff_urldecode(user, 0);
104 av_freep(&password_decoded);
105 return AVERROR(ENOMEM);
108 s->conn = amqp_new_connection();
110 av_freep(&user_decoded);
111 av_freep(&password_decoded);
112 av_log(h, AV_LOG_ERROR, "Error creating connection\n");
113 return AVERROR_EXTERNAL;
116 s->socket = amqp_tcp_socket_new(s->conn);
118 av_log(h, AV_LOG_ERROR, "Error creating socket\n");
119 goto destroy_connection;
122 if (s->connection_timeout < 0)
123 s->connection_timeout = (h->rw_timeout > 0 ? h->rw_timeout : 5000000);
125 tval.tv_sec = s->connection_timeout / 1000000;
126 tval.tv_usec = s->connection_timeout % 1000000;
127 ret = amqp_socket_open_noblock(s->socket, hostname, port, &tval);
130 av_log(h, AV_LOG_ERROR, "Error connecting to server: %s\n",
131 amqp_error_string2(ret));
132 goto destroy_connection;
135 broker_reply = amqp_login(s->conn, "/", 0, s->pkt_size, 0,
136 AMQP_SASL_METHOD_PLAIN, user_decoded, password_decoded);
138 if (broker_reply.reply_type != AMQP_RESPONSE_NORMAL) {
139 av_log(h, AV_LOG_ERROR, "Error login\n");
140 server_msg = AMQP_ACCESS_REFUSED;
141 goto close_connection;
144 amqp_channel_open(s->conn, DEFAULT_CHANNEL);
145 broker_reply = amqp_get_rpc_reply(s->conn);
147 if (broker_reply.reply_type != AMQP_RESPONSE_NORMAL) {
148 av_log(h, AV_LOG_ERROR, "Error set channel\n");
149 server_msg = AMQP_CHANNEL_ERROR;
150 goto close_connection;
153 if (h->flags & AVIO_FLAG_READ) {
154 amqp_bytes_t queuename;
155 char queuename_buff[STR_LEN];
156 amqp_queue_declare_ok_t *r;
158 r = amqp_queue_declare(s->conn, DEFAULT_CHANNEL, amqp_empty_bytes,
159 0, 0, 0, 1, amqp_empty_table);
160 broker_reply = amqp_get_rpc_reply(s->conn);
161 if (!r || broker_reply.reply_type != AMQP_RESPONSE_NORMAL) {
162 av_log(h, AV_LOG_ERROR, "Error declare queue\n");
163 server_msg = AMQP_RESOURCE_ERROR;
167 /* store queuename */
168 queuename.bytes = queuename_buff;
169 queuename.len = FFMIN(r->queue.len, STR_LEN);
170 memcpy(queuename.bytes, r->queue.bytes, queuename.len);
172 amqp_queue_bind(s->conn, DEFAULT_CHANNEL, queuename,
173 amqp_cstring_bytes(s->exchange),
174 amqp_cstring_bytes(s->routing_key), amqp_empty_table);
176 broker_reply = amqp_get_rpc_reply(s->conn);
177 if (broker_reply.reply_type != AMQP_RESPONSE_NORMAL) {
178 av_log(h, AV_LOG_ERROR, "Queue bind error\n");
179 server_msg = AMQP_INTERNAL_ERROR;
183 amqp_basic_consume(s->conn, DEFAULT_CHANNEL, queuename, amqp_empty_bytes,
184 0, 1, 0, amqp_empty_table);
186 broker_reply = amqp_get_rpc_reply(s->conn);
187 if (broker_reply.reply_type != AMQP_RESPONSE_NORMAL) {
188 av_log(h, AV_LOG_ERROR, "Set consume error\n");
189 server_msg = AMQP_INTERNAL_ERROR;
194 av_freep(&user_decoded);
195 av_freep(&password_decoded);
199 amqp_channel_close(s->conn, DEFAULT_CHANNEL, server_msg);
201 amqp_connection_close(s->conn, server_msg);
203 amqp_destroy_connection(s->conn);
205 av_freep(&user_decoded);
206 av_freep(&password_decoded);
207 return AVERROR_EXTERNAL;
210 static int amqp_proto_write(URLContext *h, const unsigned char *buf, int size)
213 AMQPContext *s = h->priv_data;
214 int fd = amqp_socket_get_sockfd(s->socket);
216 amqp_bytes_t message = { size, (void *)buf };
217 amqp_basic_properties_t props;
219 ret = ff_network_wait_fd_timeout(fd, 1, h->rw_timeout, &h->interrupt_callback);
223 props._flags = AMQP_BASIC_CONTENT_TYPE_FLAG | AMQP_BASIC_DELIVERY_MODE_FLAG;
224 props.content_type = amqp_cstring_bytes("octet/stream");
225 props.delivery_mode = 2; /* persistent delivery mode */
227 ret = amqp_basic_publish(s->conn, DEFAULT_CHANNEL, amqp_cstring_bytes(s->exchange),
228 amqp_cstring_bytes(s->routing_key), 0, 0,
232 av_log(h, AV_LOG_ERROR, "Error publish: %s\n", amqp_error_string2(ret));
233 return AVERROR_EXTERNAL;
239 static int amqp_proto_read(URLContext *h, unsigned char *buf, int size)
241 AMQPContext *s = h->priv_data;
242 int fd = amqp_socket_get_sockfd(s->socket);
245 amqp_rpc_reply_t broker_reply;
246 amqp_envelope_t envelope;
248 ret = ff_network_wait_fd_timeout(fd, 0, h->rw_timeout, &h->interrupt_callback);
252 amqp_maybe_release_buffers(s->conn);
253 broker_reply = amqp_consume_message(s->conn, &envelope, NULL, 0);
255 if (broker_reply.reply_type != AMQP_RESPONSE_NORMAL)
256 return AVERROR_EXTERNAL;
258 if (envelope.message.body.len > size) {
259 s->pkt_size_overflow = FFMAX(s->pkt_size_overflow, envelope.message.body.len);
260 av_log(h, AV_LOG_WARNING, "Message exceeds space in the buffer. "
261 "Message will be truncated. Setting -pkt_size %d "
262 "may resolve this issue.\n", s->pkt_size_overflow);
264 size = FFMIN(size, envelope.message.body.len);
266 memcpy(buf, envelope.message.body.bytes, size);
267 amqp_destroy_envelope(&envelope);
272 static int amqp_proto_close(URLContext *h)
274 AMQPContext *s = h->priv_data;
275 amqp_channel_close(s->conn, DEFAULT_CHANNEL, AMQP_REPLY_SUCCESS);
276 amqp_connection_close(s->conn, AMQP_REPLY_SUCCESS);
277 amqp_destroy_connection(s->conn);
282 static const AVClass amqp_context_class = {
283 .class_name = "amqp",
284 .item_name = av_default_item_name,
286 .version = LIBAVUTIL_VERSION_INT,
289 const URLProtocol ff_libamqp_protocol = {
291 .url_close = amqp_proto_close,
292 .url_open = amqp_proto_open,
293 .url_read = amqp_proto_read,
294 .url_write = amqp_proto_write,
295 .priv_data_size = sizeof(AMQPContext),
296 .priv_data_class = &amqp_context_class,
297 .flags = URL_PROTOCOL_FLAG_NETWORK,