[FFmpeg-devel] [PATCH] avformat: Add AMQP version 0-9-1 protocol support
Andriy Gelman
andriy.gelman at gmail.com
Sat Feb 1 21:02:14 EET 2020
From: Andriy Gelman <andriy.gelman at gmail.com>
Supports connecting to a RabbitMQ broker via AMQP version 0-9-1. The
broker can redistribute content to other clients based on "exchange" and
"routing_key" fields.
---
Compilation notes:
- Requires librabbitmq-dev package (on ubuntu).
- The pkg-config libprabbitmq.pc has a corrupt entry.
The line "Libs.private: rt; -lpthread" should be changed to
"Libs.private: -lrt -lpthread". I have made a bug report.
- Compile FFmpeg with --enable-librabbitmq
To run an example:
#
# Start the RabbitMQ broker (I use docker)
# The following starts the broker on localhost:5672. A webui is available on
# localhost:15672 (User/password is "guest" by default)
#
$ docker run -it --rm --name rabbitmq -p 127.0.0.1:5672:5672 -p 127.0.0.1:15672:15672 rabbitmq:3-management
#
# Stream to the RabbitMQ broker:
#
$ ./ffmpeg -re -f lavfi -i yuvtestsrc -codec:v libx264 -f mpegts -routing_key "amqp" -exchange "amq.direct" amqp://localhost:5672
#
# Connect any number of clients to fetch data from the broker:
# The clients are filtered by the routing_key and exchange.
#
$ ./ffplay -routing_key "amqp" -exchange "amq.direct" amqp://localhost:5672
Changelog | 1 +
configure | 5 +
doc/general.texi | 1 +
doc/protocols.texi | 53 ++++++++
libavformat/Makefile | 1 +
libavformat/libamqp.c | 272 ++++++++++++++++++++++++++++++++++++++++
libavformat/protocols.c | 1 +
7 files changed, 334 insertions(+)
create mode 100644 libavformat/libamqp.c
diff --git a/Changelog b/Changelog
index a4d20a94310..0d2c1dcc2d9 100644
--- a/Changelog
+++ b/Changelog
@@ -33,6 +33,7 @@ version <next>:
- Argonaut Games ADPCM decoder
- Argonaut Games ASF demuxer
- xfade video filter
+- AMQP protocol (RabbitMQ)
version 4.2:
diff --git a/configure b/configure
index c02dbcc8b23..e421ecb5004 100755
--- a/configure
+++ b/configure
@@ -254,6 +254,7 @@ External library support:
--enable-libopenmpt enable decoding tracked files via libopenmpt [no]
--enable-libopus enable Opus de/encoding via libopus [no]
--enable-libpulse enable Pulseaudio input via libpulse [no]
+ --enable-librabbitmq enable rabbitmq library [no]
--enable-librav1e enable AV1 encoding via rav1e [no]
--enable-librsvg enable SVG rasterization via librsvg [no]
--enable-librubberband enable rubberband needed for rubberband filter [no]
@@ -1786,6 +1787,7 @@ EXTERNAL_LIBRARY_LIST="
libopenmpt
libopus
libpulse
+ librabbitmq
librav1e
librsvg
librtmp
@@ -3430,6 +3432,8 @@ unix_protocol_deps="sys_un_h"
unix_protocol_select="network"
# external library protocols
+libamqp_protocol_deps="librabbitmq"
+libamqp_protocol_select="network"
librtmp_protocol_deps="librtmp"
librtmpe_protocol_deps="librtmp"
librtmps_protocol_deps="librtmp"
@@ -6305,6 +6309,7 @@ enabled libopus && {
}
}
enabled libpulse && require_pkg_config libpulse libpulse pulse/pulseaudio.h pa_context_new
+enabled librabbitmq && require_pkg_config librabbitmq "librabbitmq >= 0.7.1" amqp.h amqp_new_connection
enabled librav1e && require_pkg_config librav1e "rav1e >= 0.1.0" rav1e.h rav1e_context_new
enabled librsvg && require_pkg_config librsvg librsvg-2.0 librsvg-2.0/librsvg/rsvg.h rsvg_handle_render_cairo
enabled librtmp && require_pkg_config librtmp librtmp librtmp/rtmp.h RTMP_Socket
diff --git a/doc/general.texi b/doc/general.texi
index 85db50462c2..4057a07632d 100644
--- a/doc/general.texi
+++ b/doc/general.texi
@@ -1326,6 +1326,7 @@ performance on systems without hardware floating point support).
@multitable @columnfractions .4 .1
@item Name @tab Support
+ at item AMQP @tab X
@item file @tab X
@item FTP @tab X
@item Gopher @tab X
diff --git a/doc/protocols.texi b/doc/protocols.texi
index 5e8c97d1649..3d236291e77 100644
--- a/doc/protocols.texi
+++ b/doc/protocols.texi
@@ -51,6 +51,59 @@ in microseconds.
A description of the currently available protocols follows.
+ at section amqp
+
+Advanced Message Queueing Protocol (AMQP) version 0-9-1 is a broker based
+publish-subscribe communication protocol.
+
+FFmpeg must be compiled with --enable-librabbitmq to support AMQP. A separate
+AMQP broker must also be run. An example open-source AMQP broker is RabbitMQ.
+
+When connecting to the broker, a client sets an "exchange" and a "routing key".
+These keys are used to filter connections: A streaming client will only receive
+the data that matches their "exchange" and "routing key".
+
+After starting the broker, an FFmpeg client may stream data to the broker using
+the command:
+
+ at example
+ffmpeg -re -i input -f mpegts amqp://[user:password@@]hostname:port
+ at end example
+
+Where hostname and port is the location of the broker. The client may also
+set a user/password for authentication. The defaults for both fields are
+"guest".
+
+A separate instance can stream from the broker using the command:
+ at example
+ffplay amqp://[user:password@@]hostname:port
+ at end example
+
+The protocol supports the following options:
+
+ at table @option
+
+ at item routing_key
+Sets the routing key. The default value is "amqp". Clients can
+only stream data which has the same key. Multiple clients may stream data to the
+broker with different keys.
+
+ at item exchange
+Sets the exchange to use on the broker. The default value is "amqp.direct". A
+broker may have multiple exchanges which are configured on the broker side.
+
+ at item pkt_size
+Maximum size of each packet sent/received to the broker. Default is 131072.
+Minimum is 4096 and max is any large value (representable by an int). When
+receiving packets, this sets an internal buffer size in FFmpeg. It should be
+equal to or greater than the size of the sent packets to the broker. Otherwise
+the received message may be truncated causing decoding errors.
+
+ at item connection_timeout
+The timeout in milliseconds during the initial connection to the broker.
+
+ at end table
+
@section async
Asynchronous data filling wrapper for input stream.
diff --git a/libavformat/Makefile b/libavformat/Makefile
index ba6ea8c4a62..8889f60cb92 100644
--- a/libavformat/Makefile
+++ b/libavformat/Makefile
@@ -627,6 +627,7 @@ OBJS-$(CONFIG_UDPLITE_PROTOCOL) += udp.o ip.o
OBJS-$(CONFIG_UNIX_PROTOCOL) += unix.o
# external library protocols
+OBJS-$(CONFIG_LIBAMQP_PROTOCOL) += libamqp.o
OBJS-$(CONFIG_LIBRTMP_PROTOCOL) += librtmp.o
OBJS-$(CONFIG_LIBRTMPE_PROTOCOL) += librtmp.o
OBJS-$(CONFIG_LIBRTMPS_PROTOCOL) += librtmp.o
diff --git a/libavformat/libamqp.c b/libavformat/libamqp.c
new file mode 100644
index 00000000000..2b45fdaf193
--- /dev/null
+++ b/libavformat/libamqp.c
@@ -0,0 +1,272 @@
+/*
+ * AMQP Protocol
+ * Copyright (c) 2020 Andriy Gelman
+ *
+ * This file is part of FFmpeg.
+ *
+ * FFmpeg is free software; you can redistribute it and/or
+ * modify it under the terms of the GNU Lesser General Public
+ * License as published by the Free Software Foundation; either
+ * version 2.1 of the License, or (at your option) any later version.
+ *
+ * FFmpeg is distributed in the hope that it will be useful,
+ * but WITHOUT ANY WARRANTY; without even the implied warranty of
+ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU
+ * Lesser General Public License for more details.
+ *
+ * You should have received a copy of the GNU Lesser General Public
+ * License along with FFmpeg; if not, write to the Free Software
+ * Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
+ */
+
+#include <amqp.h>
+#include <amqp_tcp_socket.h>
+#include <sys/time.h>
+#include "avformat.h"
+#include "libavutil/avstring.h"
+#include "libavutil/opt.h"
+#include "libavutil/time.h"
+#include "network.h"
+#include "url.h"
+
+typedef struct AMQPContext {
+ const AVClass *class;
+ amqp_connection_state_t conn;
+ amqp_socket_t *socket;
+ const char *routing_key;
+ const char *exchange;
+ int pkt_size;
+ int connection_timeout;
+ int pkt_size_overflow;
+} AMQPContext;
+
+#define ARRAY_LEN 1024
+#define DEFAULT_CHANNEL 1
+
+#define OFFSET(x) offsetof(AMQPContext, x)
+#define D AV_OPT_FLAG_DECODING_PARAM
+#define E AV_OPT_FLAG_ENCODING_PARAM
+static const AVOption options[] = {
+ { "pkt_size", "Maximum send/read packet size", OFFSET(pkt_size), AV_OPT_TYPE_INT, { .i64 = 131072 }, 4096, INT_MAX, .flags = D | E },
+ { "routing_key", "Key to filter streams", OFFSET(routing_key), AV_OPT_TYPE_STRING, { .str = "amqp" }, 0, 0, .flags = D | E },
+ { "exchange", "Exchange to send/read packets", OFFSET(exchange), AV_OPT_TYPE_STRING, { .str = "amq.direct" }, 0, 0, .flags = D | E },
+ { "connection_timeout", "Initial connection timeout", OFFSET(connection_timeout), AV_OPT_TYPE_INT, { .i64 = -1 }, -1, INT_MAX, .flags = D | E},
+ { NULL }
+};
+
+static int amqp_proto_open(URLContext *h, const char *uri, int flags)
+{
+ int ret, server_msg;
+ char hostname[ARRAY_LEN], credentials[ARRAY_LEN];
+ int port;
+ const char *user, *password;
+ char *end;
+ struct timeval tval = { 0 };
+
+ amqp_rpc_reply_t broker_reply;
+
+ AMQPContext *s = h->priv_data;
+
+ h->is_streamed = 1;
+ h->max_packet_size = s->pkt_size;
+
+ av_url_split(NULL, 0, credentials, sizeof(credentials),
+ hostname, sizeof(hostname), &port, NULL, 0, uri);
+
+ if (hostname[0] == '\0' || port < 0 || port > 65535 ) {
+ av_log(h, AV_LOG_ERROR, "Invalid hostname/port\n");
+ return AVERROR(EIO);
+ }
+
+ user = av_strtok(credentials, ":", &end);
+ if (!user)
+ user = "guest";
+
+ password = av_strtok(NULL, ":", &end);
+ if (!password)
+ password = "guest";
+
+ s->conn = amqp_new_connection();
+ if (!s->conn) {
+ av_log(h, AV_LOG_ERROR, "Error creating connection\n");
+ return AVERROR_EXTERNAL;
+ }
+
+ s->socket = amqp_tcp_socket_new(s->conn);
+ if (!s->socket) {
+ av_log(h, AV_LOG_ERROR, "Error creating socket\n");
+ goto destroy_connection;
+ }
+
+ if (s->connection_timeout > 0) {
+ tval.tv_sec = s->connection_timeout / 1000000;
+ tval.tv_usec = s->connection_timeout % 1000000;
+ ret = amqp_socket_open_noblock(s->socket, hostname, port, &tval);
+ }
+ else
+ ret = amqp_socket_open_noblock(s->socket, hostname, port, NULL);
+
+ if (ret) {
+ av_log(h, AV_LOG_ERROR, "Error connecting to server\n");
+ goto destroy_connection;
+ }
+
+ broker_reply = amqp_login(s->conn, "/", 0, s->pkt_size, 0,
+ AMQP_SASL_METHOD_PLAIN, user, password);
+
+ if (broker_reply.reply_type != AMQP_RESPONSE_NORMAL) {
+ av_log(h, AV_LOG_ERROR, "Error login\n");
+ server_msg = AMQP_ACCESS_REFUSED;
+ goto close_connection;
+ }
+
+ amqp_channel_open(s->conn, DEFAULT_CHANNEL);
+ broker_reply = amqp_get_rpc_reply(s->conn);
+
+ if (broker_reply.reply_type != AMQP_RESPONSE_NORMAL) {
+ av_log(h, AV_LOG_ERROR, "Error set channel\n");
+ server_msg = AMQP_CHANNEL_ERROR;
+ goto close_connection;
+ }
+
+ if (h->flags & AVIO_FLAG_READ) {
+ amqp_bytes_t queuename;
+ char queuename_buff[ARRAY_LEN];
+ amqp_queue_declare_ok_t *r;
+
+ r = amqp_queue_declare(s->conn, DEFAULT_CHANNEL, amqp_empty_bytes,
+ 0, 0, 0, 1, amqp_empty_table);
+ broker_reply = amqp_get_rpc_reply(s->conn);
+ if (!r || broker_reply.reply_type != AMQP_RESPONSE_NORMAL) {
+ av_log(h, AV_LOG_ERROR, "Error declare queue\n");
+ server_msg = AMQP_RESOURCE_ERROR;
+ goto close_channel;
+ }
+
+ /* backup queuename */
+ queuename.bytes = queuename_buff;
+ queuename.len = FFMIN(r->queue.len, ARRAY_LEN);
+ memcpy(queuename.bytes, r->queue.bytes, queuename.len);
+
+ amqp_queue_bind(s->conn, DEFAULT_CHANNEL, queuename, amqp_cstring_bytes(s->exchange),
+ amqp_cstring_bytes(s->routing_key), amqp_empty_table);
+
+ broker_reply = amqp_get_rpc_reply(s->conn);
+ if (broker_reply.reply_type != AMQP_RESPONSE_NORMAL) {
+ av_log(h, AV_LOG_ERROR, "Queue bind error\n");
+ server_msg = AMQP_INTERNAL_ERROR;
+ goto close_channel;
+ }
+
+ amqp_basic_consume(s->conn, DEFAULT_CHANNEL, queuename, amqp_empty_bytes,
+ 0, 1, 0, amqp_empty_table);
+
+ broker_reply = amqp_get_rpc_reply(s->conn);
+ if (broker_reply.reply_type != AMQP_RESPONSE_NORMAL) {
+ av_log(h, AV_LOG_ERROR, "Set consume error\n");
+ server_msg = AMQP_INTERNAL_ERROR;
+ goto close_channel;
+ }
+ }
+
+ return 0;
+
+close_channel:
+ amqp_channel_close(s->conn, DEFAULT_CHANNEL, server_msg);
+close_connection:
+ amqp_connection_close(s->conn, server_msg);
+destroy_connection:
+ amqp_destroy_connection(s->conn);
+
+ return AVERROR_EXTERNAL;
+}
+
+static int amqp_proto_write(URLContext *h, const unsigned char *buf, int size)
+{
+ int ret;
+ AMQPContext *s = h->priv_data;
+ int fd = amqp_socket_get_sockfd(s->socket);
+
+ amqp_bytes_t message = { size, (void *)buf };
+ amqp_basic_properties_t props;
+
+ ret = ff_network_wait_fd_timeout(fd, 1, h->rw_timeout, &h->interrupt_callback);
+ if (ret)
+ return ret;
+
+ props._flags = AMQP_BASIC_CONTENT_TYPE_FLAG | AMQP_BASIC_DELIVERY_MODE_FLAG;
+ props.content_type = amqp_cstring_bytes("octet/stream");
+ props.delivery_mode = 2; /* persistent delivery mode */
+
+ ret = amqp_basic_publish(s->conn, DEFAULT_CHANNEL, amqp_cstring_bytes(s->exchange),
+ amqp_cstring_bytes(s->routing_key), 0, 0,
+ &props, message);
+
+ if (ret) {
+ av_log(h, AV_LOG_ERROR, "Error publish\n");
+ return AVERROR_EXTERNAL;
+ }
+
+ return size;
+}
+
+static int amqp_proto_read(URLContext *h, unsigned char *buf, int size)
+{
+ AMQPContext *s = h->priv_data;
+ int fd = amqp_socket_get_sockfd(s->socket);
+ int ret;
+
+ amqp_rpc_reply_t broker_reply;
+ amqp_envelope_t envelope;
+
+ ret = ff_network_wait_fd_timeout(fd, 0, h->rw_timeout, &h->interrupt_callback);
+ if (ret)
+ return ret;
+
+ amqp_maybe_release_buffers(s->conn);
+ broker_reply = amqp_consume_message(s->conn, &envelope, NULL, 0);
+
+ if (broker_reply.reply_type != AMQP_RESPONSE_NORMAL)
+ return AVERROR_EXTERNAL;
+
+ if (envelope.message.body.len > size) {
+ s->pkt_size_overflow = FFMAX(s->pkt_size_overflow, envelope.message.body.len);
+ av_log(h, AV_LOG_WARNING, "Message exceeds space in the buffer. "
+ "Message will be truncated. Setting -pkt_size %d "
+ "may resolve the issue.\n", s->pkt_size_overflow);
+ envelope.message.body.len = size;
+ }
+
+ memcpy(buf, envelope.message.body.bytes, envelope.message.body.len);
+ amqp_destroy_envelope(&envelope);
+
+ return envelope.message.body.len;
+}
+
+static int amqp_proto_close(URLContext *h)
+{
+ AMQPContext *s = h->priv_data;
+ amqp_channel_close(s->conn, DEFAULT_CHANNEL, AMQP_REPLY_SUCCESS);
+ amqp_connection_close(s->conn, AMQP_REPLY_SUCCESS);
+ amqp_destroy_connection(s->conn);
+
+ return 0;
+}
+
+static const AVClass amqp_context_class = {
+ .class_name = "amqp",
+ .item_name = av_default_item_name,
+ .option = options,
+ .version = LIBAVUTIL_VERSION_INT,
+};
+
+const URLProtocol ff_libamqp_protocol = {
+ .name = "amqp",
+ .url_close = amqp_proto_close,
+ .url_open = amqp_proto_open,
+ .url_read = amqp_proto_read,
+ .url_write = amqp_proto_write,
+ .priv_data_size = sizeof(AMQPContext),
+ .priv_data_class = &amqp_context_class,
+ .flags = URL_PROTOCOL_FLAG_NETWORK,
+};
diff --git a/libavformat/protocols.c b/libavformat/protocols.c
index 29fb99e7fa3..f1b8eab0fd6 100644
--- a/libavformat/protocols.c
+++ b/libavformat/protocols.c
@@ -60,6 +60,7 @@ extern const URLProtocol ff_tls_protocol;
extern const URLProtocol ff_udp_protocol;
extern const URLProtocol ff_udplite_protocol;
extern const URLProtocol ff_unix_protocol;
+extern const URLProtocol ff_libamqp_protocol;
extern const URLProtocol ff_librtmp_protocol;
extern const URLProtocol ff_librtmpe_protocol;
extern const URLProtocol ff_librtmps_protocol;
--
2.25.0
More information about the ffmpeg-devel
mailing list