FFmpeg
Loading...
Searching...
No Matches
thread_queue.c
Go to the documentation of this file.
1/*
2 * This file is part of FFmpeg.
3 *
4 * FFmpeg is free software; you can redistribute it and/or
5 * modify it under the terms of the GNU Lesser General Public
6 * License as published by the Free Software Foundation; either
7 * version 2.1 of the License, or (at your option) any later version.
8 *
9 * FFmpeg is distributed in the hope that it will be useful,
10 * but WITHOUT ANY WARRANTY; without even the implied warranty of
11 * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU
12 * Lesser General Public License for more details.
13 *
14 * You should have received a copy of the GNU Lesser General Public
15 * License along with FFmpeg; if not, write to the Free Software
16 * Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
17 */
18
19#include <stdint.h>
20#include <string.h>
21
22#include "libavutil/avassert.h"
24#include "libavutil/error.h"
25#include "libavutil/fifo.h"
26#include "libavutil/frame.h"
28#include "libavutil/mem.h"
29#include "libavutil/thread.h"
30
31#include "libavcodec/packet.h"
32
33#include "thread_queue.h"
34
35enum {
36 FINISHED_SEND = (1 << 0),
37 FINISHED_RECV = (1 << 1),
38};
39
57
59{
60 ThreadQueue *tq = *ptq;
61
62 if (!tq)
63 return;
64
68
69 av_freep(&tq->finished);
70
71 for (unsigned i = 0; i < tq->nb_cond_write; i++)
73 av_freep(&tq->cond_write);
74
77
78 av_freep(ptq);
79}
80
81ThreadQueue *tq_alloc(unsigned int nb_streams, size_t queue_size,
83{
84 ThreadQueue *tq;
85 int ret;
86
87 tq = av_mallocz(sizeof(*tq));
88 if (!tq)
89 return NULL;
90
91 ret = pthread_cond_init(&tq->cond_read, NULL);
92 if (ret) {
93 av_freep(&tq);
94 return NULL;
95 }
96
97 ret = pthread_mutex_init(&tq->lock, NULL);
98 if (ret) {
100 av_freep(&tq);
101 return NULL;
102 }
103
104 tq->cond_write = av_calloc(nb_streams, sizeof(*tq->cond_write));
105 if (!tq->cond_write)
106 goto fail;
107 for (tq->nb_cond_write = 0; tq->nb_cond_write < nb_streams; tq->nb_cond_write++) {
109 if (ret)
110 goto fail;
111 }
112
113 tq->finished = av_calloc(nb_streams, sizeof(*tq->finished));
114 if (!tq->finished)
115 goto fail;
117 tq->queue_size = queue_size;
118 tq->type = type;
119
120 tq->fifo = (type == THREAD_QUEUE_FRAMES) ?
122 if (!tq->fifo)
123 goto fail;
124
125 av_assert0(queue_size);
126 if (nb_streams > SIZE_MAX / queue_size)
127 goto fail; // treat like OOM
128
129 tq->fifo_stream_index = av_fifo_alloc2(queue_size * nb_streams, sizeof(unsigned), 0);
130 if (!tq->fifo_stream_index)
131 goto fail;
132
133 tq->stream_count = av_calloc(nb_streams, sizeof(*tq->stream_count));
134 if (!tq->stream_count)
135 goto fail;
136
137 return tq;
138fail:
139 tq_free(&tq);
140 return NULL;
141}
142
143static int can_write(ThreadQueue *tq, unsigned int stream_idx)
144{
145 return tq->stream_count[stream_idx] < tq->queue_size;
146}
147
148int tq_send(ThreadQueue *tq, unsigned int stream_idx, void *data)
149{
150 int *finished;
151 int ret;
152
153 av_assert0(stream_idx < tq->nb_streams);
154 finished = &tq->finished[stream_idx];
155
157
158 if (*finished & FINISHED_SEND) {
159 ret = AVERROR(EINVAL);
160 goto finish;
161 }
162
163 while (!(*finished & FINISHED_RECV) && !can_write(tq, stream_idx))
164 pthread_cond_wait(&tq->cond_write[stream_idx], &tq->lock);
165
166 if (*finished & FINISHED_RECV) {
167 ret = AVERROR_EOF;
168 *finished |= FINISHED_SEND;
169 } else {
170 ret = av_fifo_write(tq->fifo_stream_index, &stream_idx, 1);
171 if (ret < 0)
172 goto finish;
173
174 ret = av_container_fifo_write(tq->fifo, data, 0);
175 if (ret < 0)
176 goto finish;
177
178 tq->stream_count[stream_idx]++;
179 pthread_cond_broadcast(&tq->cond_read); // signal downstream
180 }
181
182finish:
184
185 return ret;
186}
187
188static int receive_locked(ThreadQueue *tq, int *stream_idx,
189 void *data)
190{
191 unsigned int nb_finished = 0;
192
193 if (tq->choked)
194 return AVERROR(EAGAIN);
195
196 while (av_container_fifo_read(tq->fifo, data, 0) >= 0) {
197 unsigned idx;
198 int ret;
199
200 ret = av_fifo_read(tq->fifo_stream_index, &idx, 1);
201 av_assert0(ret >= 0);
202
203 // signal upstream if the fifo is no longer full
204 if (tq->stream_count[idx]-- == tq->queue_size)
206
207 if (tq->finished[idx] & FINISHED_RECV) {
208 (tq->type == THREAD_QUEUE_FRAMES) ?
210 continue;
211 }
212
213 *stream_idx = idx;
214 return 0;
215 }
216
217 for (unsigned int i = 0; i < tq->nb_streams; i++) {
218 if (!tq->finished[i])
219 continue;
220
221 /* return EOF to the consumer at most once for each stream */
222 if (!(tq->finished[i] & FINISHED_RECV)) {
223 tq->finished[i] |= FINISHED_RECV;
224 *stream_idx = i;
225 return AVERROR_EOF;
226 }
227
228 nb_finished++;
229 }
230
231 return nb_finished == tq->nb_streams ? AVERROR_EOF : AVERROR(EAGAIN);
232}
233
234int tq_receive(ThreadQueue *tq, int *stream_idx, void *data, int flags)
235{
236 int ret;
237
238 *stream_idx = -1;
239
241
242 while (1) {
243 ret = receive_locked(tq, stream_idx, data);
244
245 if (ret == AVERROR(EAGAIN) && !(flags & THREAD_QUEUE_FLAG_NO_BLOCK)) {
246 pthread_cond_wait(&tq->cond_read, &tq->lock);
247 continue;
248 }
249
250 break;
251 }
252
254
255 return ret;
256}
257
258void tq_send_finish(ThreadQueue *tq, unsigned int stream_idx)
259{
260 av_assert0(stream_idx < tq->nb_streams);
261
263
264 /* mark the stream as send-finished;
265 * next time the consumer thread tries to read this stream it will get
266 * an EOF and recv-finished flag will be set */
267 tq->finished[stream_idx] |= FINISHED_SEND;
268 tq->choked = 0;
270
272}
273
274void tq_receive_finish(ThreadQueue *tq, unsigned int stream_idx)
275{
276 av_assert0(stream_idx < tq->nb_streams);
277
279
280 /* mark the stream as recv-finished;
281 * next time the producer thread tries to send for this stream, it will
282 * get an EOF and send-finished flag will be set */
283 tq->finished[stream_idx] |= FINISHED_RECV;
284 pthread_cond_broadcast(&tq->cond_write[stream_idx]);
285
287}
288
289void tq_choke(ThreadQueue *tq, int choked)
290{
292
293 int prev_choked = tq->choked;
294 tq->choked = choked;
295 if (choked != prev_choked)
297
299}
static void finish(void)
simple assert() macros that are a bit more flexible than ISO C assert().
#define av_assert0(cond)
assert() equivalent, that is always enabled.
Definition avassert.h:42
#define flags(name, subs,...)
Definition cbs_h264.c:74
#define i(width, name, range_min, range_max)
Definition cbs_h264.c:63
int av_container_fifo_write(AVContainerFifo *cf, void *obj, unsigned flags)
Write the contents of obj to the FIFO.
void av_container_fifo_free(AVContainerFifo **pcf)
Free a AVContainerFifo and everything in it.
int av_container_fifo_read(AVContainerFifo *cf, void *obj, unsigned flags)
Read the next available object from the FIFO into obj.
AVContainerFifo * av_container_fifo_alloc_avframe(unsigned flags)
Allocate an AVContainerFifo instance for AVFrames.
#define NULL
Definition coverity.c:32
error code definitions
static unsigned int nb_streams
Definition ffprobe.c:353
A generic FIFO API.
reference-counted frame API
#define fail
Definition test.h:479
void av_packet_unref(AVPacket *pkt)
Wipe the packet.
Definition packet.c:434
AVContainerFifo * av_container_fifo_alloc_avpacket(unsigned flags)
Allocate an AVContainerFifo instance for AVPacket.
Definition packet.c:695
#define AVERROR_EOF
End of file.
Definition error.h:57
#define AVERROR(e)
Definition error.h:45
AVFifo * av_fifo_alloc2(size_t nb_elems, size_t elem_size, unsigned int flags)
Allocate and initialize an AVFifo with a given element size.
Definition fifo.c:47
void av_fifo_freep2(AVFifo **f)
Free an AVFifo and reset pointer to NULL.
Definition fifo.c:286
int av_fifo_write(AVFifo *f, const void *buf, size_t nb_elems)
Write data into a FIFO.
Definition fifo.c:188
int av_fifo_read(AVFifo *f, void *buf, size_t nb_elems)
Read data from a FIFO.
Definition fifo.c:240
void av_frame_unref(AVFrame *frame)
Unreference all the buffers referenced by frame and reset the frame fields.
Definition frame.c:496
uint32_t type
Definition jpegmpfenc.c:80
void * av_calloc(size_t nmemb, size_t size)
Definition mem.c:370
Memory handling functions.
const char data[16]
Definition mxf.c:149
static av_always_inline int pthread_cond_broadcast(pthread_cond_t *cond)
Definition os2threads.h:168
static av_always_inline int pthread_mutex_lock(pthread_mutex_t *mutex)
Definition os2threads.h:119
static av_always_inline int pthread_cond_destroy(pthread_cond_t *cond)
Definition os2threads.h:150
static av_always_inline int pthread_mutex_init(pthread_mutex_t *mutex, const pthread_mutexattr_t *attr)
Definition os2threads.h:104
static av_always_inline int pthread_cond_init(pthread_cond_t *cond, const pthread_condattr_t *attr)
Definition os2threads.h:139
_fmutex pthread_mutex_t
Definition os2threads.h:53
static av_always_inline int pthread_mutex_unlock(pthread_mutex_t *mutex)
Definition os2threads.h:132
static av_always_inline int pthread_cond_wait(pthread_cond_t *cond, pthread_mutex_t *mutex)
Definition os2threads.h:198
static av_always_inline int pthread_mutex_destroy(pthread_mutex_t *mutex)
Definition os2threads.h:112
AVContainerFifo is a FIFO for "containers" - dynamically allocated reusable structs (e....
Definition fifo.c:35
pthread_mutex_t lock
unsigned int nb_cond_write
unsigned int nb_streams
size_t queue_size
enum ThreadQueueType type
pthread_cond_t cond_read
size_t * stream_count
AVContainerFifo * fifo
AVFifo * fifo_stream_index
pthread_cond_t * cond_write
#define av_mallocz(s)
#define av_freep(p)
ThreadQueue * tq_alloc(unsigned int nb_streams, size_t queue_size, enum ThreadQueueType type)
Allocate a queue for sending data between threads.
void tq_send_finish(ThreadQueue *tq, unsigned int stream_idx)
Mark the given stream finished from the sending side.
int tq_send(ThreadQueue *tq, unsigned int stream_idx, void *data)
Send an item for the given stream to the queue.
static int receive_locked(ThreadQueue *tq, int *stream_idx, void *data)
void tq_choke(ThreadQueue *tq, int choked)
Prevent further reads from the thread queue until it is unchoked.
static int can_write(ThreadQueue *tq, unsigned int stream_idx)
int tq_receive(ThreadQueue *tq, int *stream_idx, void *data, int flags)
Read the next item from the queue.
void tq_free(ThreadQueue **ptq)
@ FINISHED_SEND
@ FINISHED_RECV
void tq_receive_finish(ThreadQueue *tq, unsigned int stream_idx)
Mark the given stream finished from the receiving side.
ThreadQueueType
@ THREAD_QUEUE_FRAMES
@ THREAD_QUEUE_FLAG_NO_BLOCK