From patchwork Sat Feb 1 19:02:14 2020 Content-Type: text/plain; charset="utf-8" MIME-Version: 1.0 Content-Transfer-Encoding: 7bit X-Patchwork-Submitter: Andriy Gelman X-Patchwork-Id: 17647 Return-Path: X-Original-To: patchwork@ffaux-bg.ffmpeg.org Delivered-To: patchwork@ffaux-bg.ffmpeg.org Received: from ffbox0-bg.mplayerhq.hu (ffbox0-bg.ffmpeg.org [79.124.17.100]) by ffaux.localdomain (Postfix) with ESMTP id 7B8A9446DA8 for ; Sat, 1 Feb 2020 21:02:32 +0200 (EET) Received: from [127.0.1.1] (localhost [127.0.0.1]) by ffbox0-bg.mplayerhq.hu (Postfix) with ESMTP id 56E3868AEF8; Sat, 1 Feb 2020 21:02:32 +0200 (EET) X-Original-To: ffmpeg-devel@ffmpeg.org Delivered-To: ffmpeg-devel@ffmpeg.org Received: from mail-qt1-f177.google.com (mail-qt1-f177.google.com [209.85.160.177]) by ffbox0-bg.mplayerhq.hu (Postfix) with ESMTPS id 541C368AEEC for ; Sat, 1 Feb 2020 21:02:25 +0200 (EET) Received: by mail-qt1-f177.google.com with SMTP id w8so8175563qts.11 for ; Sat, 01 Feb 2020 11:02:25 -0800 (PST) DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=gmail.com; s=20161025; h=from:to:cc:subject:date:message-id:mime-version :content-transfer-encoding; bh=SaTHNj0p7C4tykByvxWmTtamBOmU0p+Shme3oPWXflw=; b=hQDok9zAtlFOZcrjg4dDVg/fdwJcgpxw2n4+1c4syaXVUAgXaPMWjgewN6e2IxpS00 JQeLN1gj5/+DcEkCY6u1C30GG8GXkoxXHv2Be2XcmaxcudtI2/ajO7v0lVtuWXqK33fQ uh4quTHKu4UmbyZVLnByFWJ8b5Y02aVdCl+fNHj03+0S+3iURx5uzvNdZQ62XkTXR/I4 A0HmNn4J0ZhwxGJDWLjTZZam+xTnjqU2K3k1Xre6G69KUDysHWLIy0AYpx5+wxED8YhJ lBf48tM9mE636/NrXfBAUTDDtUirFauxjhI54LmhdAXY6FV+c7+2hDbg/RnszeyXnsL2 R18w== X-Google-DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=1e100.net; s=20161025; h=x-gm-message-state:from:to:cc:subject:date:message-id:mime-version :content-transfer-encoding; bh=SaTHNj0p7C4tykByvxWmTtamBOmU0p+Shme3oPWXflw=; b=CDcplWnCrJ1mTVxWMg3NxNf/KSQ7u93Q35qB+lpAi41fVYxlByLpUtsVKHi6i5+O8r tF2m9E+UU4MyQ1sccSbB5MmFgEy3YTb/9TWIEbOfwndUsb5iqDa/cJhzdmhap93XFe+g Yp82yA2uHzMueGzumdAHI0/xVc52zb46GbHR97H6EBpWSnPWsI4714kBe+RBVec74ANg A2cupWgSXG3bSkwfDwNSrFoUTnVBCQrI7Opps7Lfq93km6FNi6EAINPI8kNamM/QgYvy BmsOKPqQjFN87Bgcfbp3Rk12nwYtFZtTV1LlYZ88Gf+XLkN6qRuU1wLUTSaoX9DJ8aLn A4Kg== X-Gm-Message-State: APjAAAUKdGbrsuvYJVwhaabdvAsban9TO2u7FmlGaYvhQWnChnqr4dhi nw7st8AdlRRA6hnww4RNXzEC4jhu X-Google-Smtp-Source: APXvYqxA8Ux7rwfaCu9COJ7rUUbgtMHdmQmGBFBk6iC5E5m8/3XTwwd6YMy9tm/dr7/soEUKaM9JSA== X-Received: by 2002:ac8:6644:: with SMTP id j4mr15950298qtp.90.1580583743256; Sat, 01 Feb 2020 11:02:23 -0800 (PST) Received: from localhost.localdomain (c-71-232-27-28.hsd1.ma.comcast.net. [71.232.27.28]) by smtp.gmail.com with ESMTPSA id 3sm6784491qte.59.2020.02.01.11.02.22 (version=TLS1_3 cipher=TLS_AES_256_GCM_SHA384 bits=256/256); Sat, 01 Feb 2020 11:02:22 -0800 (PST) From: Andriy Gelman X-Google-Original-From: Andriy Gelman To: ffmpeg-devel@ffmpeg.org Date: Sat, 1 Feb 2020 14:02:14 -0500 Message-Id: <20200201190214.10457-1-andriy.gelman@gmail.com> X-Mailer: git-send-email 2.25.0 MIME-Version: 1.0 Subject: [FFmpeg-devel] [PATCH] avformat: Add AMQP version 0-9-1 protocol support X-BeenThere: ffmpeg-devel@ffmpeg.org X-Mailman-Version: 2.1.20 Precedence: list List-Id: FFmpeg development discussions and patches List-Unsubscribe: , List-Archive: List-Post: List-Help: List-Subscribe: , Reply-To: FFmpeg development discussions and patches Cc: Andriy Gelman Errors-To: ffmpeg-devel-bounces@ffmpeg.org Sender: "ffmpeg-devel" From: Andriy Gelman 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 : - 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 +@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. +@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: + +@example +ffmpeg -re -i input -f mpegts amqp://[user:password@@]hostname:port +@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: +@example +ffplay amqp://[user:password@@]hostname:port +@end example + +The protocol supports the following options: + +@table @option + +@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. + +@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. + +@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. + +@item connection_timeout +The timeout in milliseconds during the initial connection to the broker. + +@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 +#include +#include +#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;