FFmpeg
Loading...
Searching...
No Matches
libzmq.c
Go to the documentation of this file.
1/*
2 * ZeroMQ Protocol
3 * Copyright (c) 2019 Andriy Gelman
4 *
5 * This file is part of FFmpeg.
6 *
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.
11 *
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.
16 *
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
20 */
21
22#include <zmq.h>
23#include "url.h"
24#include "network.h"
25#include "libavutil/avstring.h"
26#include "libavutil/opt.h"
27#include "libavutil/time.h"
28
29#define ZMQ_STRERROR zmq_strerror(zmq_errno())
30
31typedef struct ZMQContext {
32 const AVClass *class;
33 void *context;
34 void *socket;
36 int pkt_size_overflow; /*keep track of the largest packet during overflow*/
38
39#define OFFSET(x) offsetof(ZMQContext, x)
40#define D AV_OPT_FLAG_DECODING_PARAM
41#define E AV_OPT_FLAG_ENCODING_PARAM
42static const AVOption options[] = {
43 { "pkt_size", "Maximum send/read packet size", OFFSET(pkt_size), AV_OPT_TYPE_INT, { .i64 = 131072 }, -1, INT_MAX, .flags = D | E },
44 { NULL }
45};
46
47static int zmq_proto_wait(URLContext *h, void *socket, int write)
48{
49 int ret;
50 int ev = write ? ZMQ_POLLOUT : ZMQ_POLLIN;
51 zmq_pollitem_t items = { .socket = socket, .fd = 0, .events = ev, .revents = 0 };
52 ret = zmq_poll(&items, 1, POLLING_TIME);
53 if (ret == -1) {
54 av_log(h, AV_LOG_ERROR, "Error occurred during zmq_poll(): %s\n", ZMQ_STRERROR);
55 return AVERROR_EXTERNAL;
56 }
57 return items.revents & ev ? 0 : AVERROR(EAGAIN);
58}
59
60static int zmq_proto_wait_timeout(URLContext *h, void *socket, int write, int64_t timeout, AVIOInterruptCB *int_cb)
61{
62 int ret;
63 int64_t wait_start = 0;
64
65 while (1) {
67 return AVERROR_EXIT;
68 ret = zmq_proto_wait(h, socket, write);
69 if (ret != AVERROR(EAGAIN))
70 return ret;
71 if (timeout > 0) {
72 if (!wait_start)
73 wait_start = av_gettime_relative();
74 else if (av_gettime_relative() - wait_start > timeout)
75 return AVERROR(ETIMEDOUT);
76 }
77 }
78}
79
80static int zmq_proto_open(URLContext *h, const char *uri, int flags)
81{
82 int ret;
83 ZMQContext *s = h->priv_data;
84 s->pkt_size_overflow = 0;
85 h->is_streamed = 1;
86
87 if (s->pkt_size > 0)
88 h->max_packet_size = s->pkt_size;
89
90 s->context = zmq_ctx_new();
91 if (!s->context) {
92 /*errno not set on failure during zmq_ctx_new()*/
93 av_log(h, AV_LOG_ERROR, "Error occurred during zmq_ctx_new()\n");
94 return AVERROR_EXTERNAL;
95 }
96
97 if (!av_strstart(uri, "zmq:", &uri)) {
98 av_log(h, AV_LOG_ERROR, "URL %s lacks prefix\n", uri);
99 return AVERROR(EINVAL);
100 }
101
102 /*publish during write*/
103 if (h->flags & AVIO_FLAG_WRITE) {
104 s->socket = zmq_socket(s->context, ZMQ_PUB);
105 if (!s->socket) {
106 av_log(h, AV_LOG_ERROR, "Error occurred during zmq_socket(): %s\n", ZMQ_STRERROR);
107 goto fail_term;
108 }
109
110 ret = zmq_bind(s->socket, uri);
111 if (ret == -1) {
112 av_log(h, AV_LOG_ERROR, "Error occurred during zmq_bind(): %s\n", ZMQ_STRERROR);
113 goto fail_close;
114 }
115 }
116
117 /*subscribe for read*/
118 if (h->flags & AVIO_FLAG_READ) {
119 s->socket = zmq_socket(s->context, ZMQ_SUB);
120 if (!s->socket) {
121 av_log(h, AV_LOG_ERROR, "Error occurred during zmq_socket(): %s\n", ZMQ_STRERROR);
122 goto fail_term;
123 }
124
125 ret = zmq_setsockopt(s->socket, ZMQ_SUBSCRIBE, "", 0);
126 if (ret == -1) {
127 av_log(h, AV_LOG_ERROR, "Error occurred during zmq_setsockopt(): %s\n", ZMQ_STRERROR);
128 goto fail_close;
129 }
130
131 ret = zmq_connect(s->socket, uri);
132 if (ret == -1) {
133 av_log(h, AV_LOG_ERROR, "Error occurred during zmq_connect(): %s\n", ZMQ_STRERROR);
134 goto fail_close;
135 }
136 }
137 return 0;
138
139fail_close:
140 zmq_close(s->socket);
141fail_term:
142 zmq_ctx_term(s->context);
143 return AVERROR_EXTERNAL;
144}
145
146static int zmq_proto_write(URLContext *h, const unsigned char *buf, int size)
147{
148 int ret;
149 ZMQContext *s = h->priv_data;
150
151 ret = zmq_proto_wait_timeout(h, s->socket, 1, h->rw_timeout, &h->interrupt_callback);
152 if (ret)
153 return ret;
154 ret = zmq_send(s->socket, buf, size, 0);
155 if (ret == -1) {
156 av_log(h, AV_LOG_ERROR, "Error occurred during zmq_send(): %s\n", ZMQ_STRERROR);
157 return AVERROR_EXTERNAL;
158 }
159 return ret; /*number of bytes sent*/
160}
161
162static int zmq_proto_read(URLContext *h, unsigned char *buf, int size)
163{
164 int ret;
165 ZMQContext *s = h->priv_data;
166
167 ret = zmq_proto_wait_timeout(h, s->socket, 0, h->rw_timeout, &h->interrupt_callback);
168 if (ret)
169 return ret;
170 ret = zmq_recv(s->socket, buf, size, 0);
171 if (ret == -1) {
172 av_log(h, AV_LOG_ERROR, "Error occurred during zmq_recv(): %s\n", ZMQ_STRERROR);
173 return AVERROR_EXTERNAL;
174 }
175 if (ret > size) {
176 s->pkt_size_overflow = FFMAX(s->pkt_size_overflow, ret);
177 av_log(h, AV_LOG_WARNING, "Message exceeds available space in the buffer. Message will be truncated. Setting -pkt_size %d may resolve the issue.\n", s->pkt_size_overflow);
178 ret = size;
179 }
180 return ret; /*number of bytes read*/
181}
182
184{
185 ZMQContext *s = h->priv_data;
186 zmq_close(s->socket);
187 zmq_ctx_term(s->context);
188 return 0;
189}
190
192 .class_name = "zmq",
193 .item_name = av_default_item_name,
194 .option = options,
195 .version = LIBAVUTIL_VERSION_INT,
196};
197
199 .name = "zmq",
200 .url_close = zmq_proto_close,
201 .url_open = zmq_proto_open,
202 .url_read = zmq_proto_read,
203 .url_write = zmq_proto_write,
204 .priv_data_size = sizeof(ZMQContext),
205 .priv_data_class = &zmq_context_class,
207};
#define E
Definition avdct.c:34
#define D
Definition avdct.c:35
int ff_check_interrupt(AVIOInterruptCB *cb)
Check if the user has requested to interrupt a blocking function associated with cb.
Definition avio.c:922
#define AVIO_FLAG_READ
read-only
Definition avio.h:617
#define AVIO_FLAG_WRITE
write-only
Definition avio.h:618
#define flags(name, subs,...)
Definition cbs_h264.c:74
#define s(width, name)
Definition cbs_vp9.c:198
#define NULL
Definition coverity.c:32
long long int64_t
Definition coverity.c:34
const AVIOInterruptCB int_cb
Definition ffmpeg.c:322
@ AV_OPT_TYPE_INT
Underlying C type is int.
Definition opt.h:258
#define AVERROR_EXIT
Immediate exit was requested; the called function should not be restarted.
Definition error.h:58
#define AVERROR_EXTERNAL
Generic error in an external library.
Definition error.h:59
#define AVERROR(e)
Definition error.h:45
#define AV_LOG_WARNING
Something somehow does not look correct.
Definition log.h:216
#define AV_LOG_ERROR
Something went wrong and cannot losslessly be recovered.
Definition log.h:210
const char * av_default_item_name(void *ptr)
Return the context name.
Definition log.c:241
int av_strstart(const char *str, const char *pfx, const char **ptr)
Return non-zero if pfx is a prefix of str.
Definition avstring.c:36
#define LIBAVUTIL_VERSION_INT
Definition version.h:85
static int zmq_proto_read(URLContext *h, unsigned char *buf, int size)
Definition libzmq.c:162
const URLProtocol ff_libzmq_protocol
Definition libzmq.c:198
static const AVClass zmq_context_class
Definition libzmq.c:191
static int zmq_proto_wait(URLContext *h, void *socket, int write)
Definition libzmq.c:47
static int zmq_proto_close(URLContext *h)
Definition libzmq.c:183
static int zmq_proto_wait_timeout(URLContext *h, void *socket, int write, int64_t timeout, AVIOInterruptCB *int_cb)
Definition libzmq.c:60
#define OFFSET(x)
Definition libzmq.c:39
static int zmq_proto_write(URLContext *h, const unsigned char *buf, int size)
Definition libzmq.c:146
static int zmq_proto_open(URLContext *h, const char *uri, int flags)
Definition libzmq.c:80
#define ZMQ_STRERROR
Definition libzmq.c:29
#define FFMAX(a, b)
Definition macros.h:47
#define POLLING_TIME
Definition network.h:249
AVOptions.
Describe the class of an AVClass context structure.
Definition log.h:76
Callback for checking whether to abort blocking functions.
Definition avio.h:59
AVOption.
Definition opt.h:428
int pkt_size
Definition libzmq.c:35
void * socket
Definition libzmq.c:34
int pkt_size_overflow
Definition libzmq.c:36
void * context
Definition libzmq.c:33
#define av_log(a,...)
int64_t av_gettime_relative(void)
Get the current time in microseconds since some unspecified starting point.
Definition time.c:57
int size
unbuffered private I/O API
#define URL_PROTOCOL_FLAG_NETWORK
Definition url.h:33