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
,