FFmpeg coverage


Directory: ../../../ffmpeg/
File: src/fftools/thread_queue.c
Date: 2026-09-28 04:46:15
Exec Total Coverage
Lines: 119 139 85.6%
Functions: 9 9 100.0%
Branches: 56 74 75.7%

Line Branch Exec Source
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"
23 #include "libavutil/container_fifo.h"
24 #include "libavutil/error.h"
25 #include "libavutil/fifo.h"
26 #include "libavutil/frame.h"
27 #include "libavutil/intreadwrite.h"
28 #include "libavutil/mem.h"
29 #include "libavutil/thread.h"
30
31 #include "libavcodec/packet.h"
32
33 #include "thread_queue.h"
34
35 enum {
36 FINISHED_SEND = (1 << 0),
37 FINISHED_RECV = (1 << 1),
38 };
39
40 struct ThreadQueue {
41 int choked;
42 int *finished;
43 unsigned int nb_streams;
44 size_t queue_size;
45
46 enum ThreadQueueType type;
47
48 AVContainerFifo *fifo;
49 AVFifo *fifo_stream_index;
50 size_t *stream_count;
51
52 pthread_mutex_t lock;
53 pthread_cond_t cond_read;
54 pthread_cond_t *cond_write;
55 unsigned int nb_cond_write;
56 };
57
58 32668 void tq_free(ThreadQueue **ptq)
59 {
60 32668 ThreadQueue *tq = *ptq;
61
62
2/2
✓ Branch 0 taken 7 times.
✓ Branch 1 taken 32661 times.
32668 if (!tq)
63 7 return;
64
65 32661 av_container_fifo_free(&tq->fifo);
66 32661 av_fifo_freep2(&tq->fifo_stream_index);
67 32661 av_freep(&tq->stream_count);
68
69 32661 av_freep(&tq->finished);
70
71
2/2
✓ Branch 0 taken 40380 times.
✓ Branch 1 taken 32661 times.
73041 for (unsigned i = 0; i < tq->nb_cond_write; i++)
72 40380 pthread_cond_destroy(&tq->cond_write[i]);
73 32661 av_freep(&tq->cond_write);
74
75 32661 pthread_cond_destroy(&tq->cond_read);
76 32661 pthread_mutex_destroy(&tq->lock);
77
78 32661 av_freep(ptq);
79 }
80
81 32661 ThreadQueue *tq_alloc(unsigned int nb_streams, size_t queue_size,
82 enum ThreadQueueType type)
83 {
84 ThreadQueue *tq;
85 int ret;
86
87 32661 tq = av_mallocz(sizeof(*tq));
88
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 32661 times.
32661 if (!tq)
89 ✗ return NULL;
90
91 32661 ret = pthread_cond_init(&tq->cond_read, NULL);
92
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 32661 times.
32661 if (ret) {
93 ✗ av_freep(&tq);
94 ✗ return NULL;
95 }
96
97 32661 ret = pthread_mutex_init(&tq->lock, NULL);
98
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 32661 times.
32661 if (ret) {
99 ✗ pthread_cond_destroy(&tq->cond_read);
100 ✗ av_freep(&tq);
101 ✗ return NULL;
102 }
103
104 32661 tq->cond_write = av_calloc(nb_streams, sizeof(*tq->cond_write));
105
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 32661 times.
32661 if (!tq->cond_write)
106 ✗ goto fail;
107
2/2
✓ Branch 0 taken 40380 times.
✓ Branch 1 taken 32661 times.
73041 for (tq->nb_cond_write = 0; tq->nb_cond_write < nb_streams; tq->nb_cond_write++) {
108 40380 ret = pthread_cond_init(&tq->cond_write[tq->nb_cond_write], NULL);
109
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 40380 times.
40380 if (ret)
110 ✗ goto fail;
111 }
112
113 32661 tq->finished = av_calloc(nb_streams, sizeof(*tq->finished));
114
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 32661 times.
32661 if (!tq->finished)
115 ✗ goto fail;
116 32661 tq->nb_streams = nb_streams;
117 32661 tq->queue_size = queue_size;
118 32661 tq->type = type;
119
120 32661 tq->fifo = (type == THREAD_QUEUE_FRAMES) ?
121
2/2
✓ Branch 0 taken 16692 times.
✓ Branch 1 taken 15969 times.
32661 av_container_fifo_alloc_avframe(0) : av_container_fifo_alloc_avpacket(0);
122
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 32661 times.
32661 if (!tq->fifo)
123 ✗ goto fail;
124
125
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 32661 times.
32661 av_assert0(queue_size);
126
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 32661 times.
32661 if (nb_streams > SIZE_MAX / queue_size)
127 ✗ goto fail; // treat like OOM
128
129 32661 tq->fifo_stream_index = av_fifo_alloc2(queue_size * nb_streams, sizeof(unsigned), 0);
130
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 32661 times.
32661 if (!tq->fifo_stream_index)
131 ✗ goto fail;
132
133 32661 tq->stream_count = av_calloc(nb_streams, sizeof(*tq->stream_count));
134
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 32661 times.
32661 if (!tq->stream_count)
135 ✗ goto fail;
136
137 32661 return tq;
138 ✗ fail:
139 ✗ tq_free(&tq);
140 ✗ return NULL;
141 }
142
143 2666068 static int can_write(ThreadQueue *tq, unsigned int stream_idx)
144 {
145 2666068 return tq->stream_count[stream_idx] < tq->queue_size;
146 }
147
148 2009233 int tq_send(ThreadQueue *tq, unsigned int stream_idx, void *data)
149 {
150 int *finished;
151 int ret;
152
153
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 2009233 times.
2009233 av_assert0(stream_idx < tq->nb_streams);
154 2009233 finished = &tq->finished[stream_idx];
155
156 2009233 pthread_mutex_lock(&tq->lock);
157
158
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 2009233 times.
2009233 if (*finished & FINISHED_SEND) {
159 ✗ ret = AVERROR(EINVAL);
160 ✗ goto finish;
161 }
162
163
4/4
✓ Branch 0 taken 2666068 times.
✓ Branch 1 taken 6970 times.
✓ Branch 3 taken 663805 times.
✓ Branch 4 taken 2002263 times.
2673038 while (!(*finished & FINISHED_RECV) && !can_write(tq, stream_idx))
164 663805 pthread_cond_wait(&tq->cond_write[stream_idx], &tq->lock);
165
166
2/2
✓ Branch 0 taken 6970 times.
✓ Branch 1 taken 2002263 times.
2009233 if (*finished & FINISHED_RECV) {
167 6970 ret = AVERROR_EOF;
168 6970 *finished |= FINISHED_SEND;
169 } else {
170 2002263 ret = av_fifo_write(tq->fifo_stream_index, &stream_idx, 1);
171
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 2002263 times.
2002263 if (ret < 0)
172 ✗ goto finish;
173
174 2002263 ret = av_container_fifo_write(tq->fifo, data, 0);
175
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 2002263 times.
2002263 if (ret < 0)
176 ✗ goto finish;
177
178 2002263 tq->stream_count[stream_idx]++;
179 2002263 pthread_cond_broadcast(&tq->cond_read); // signal downstream
180 }
181
182 2009233 finish:
183 2009233 pthread_mutex_unlock(&tq->lock);
184
185 2009233 return ret;
186 }
187
188 3161526 static int receive_locked(ThreadQueue *tq, int *stream_idx,
189 void *data)
190 {
191 3161526 unsigned int nb_finished = 0;
192
193
2/2
✓ Branch 0 taken 7487 times.
✓ Branch 1 taken 3154039 times.
3161526 if (tq->choked)
194 7487 return AVERROR(EAGAIN);
195
196
2/2
✓ Branch 1 taken 1988986 times.
✓ Branch 2 taken 1165721 times.
3154707 while (av_container_fifo_read(tq->fifo, data, 0) >= 0) {
197 unsigned idx;
198 int ret;
199
200 1988986 ret = av_fifo_read(tq->fifo_stream_index, &idx, 1);
201
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 1988986 times.
1988986 av_assert0(ret >= 0);
202
203 // signal upstream if the fifo is no longer full
204
2/2
✓ Branch 0 taken 795699 times.
✓ Branch 1 taken 1193287 times.
1988986 if (tq->stream_count[idx]-- == tq->queue_size)
205 795699 pthread_cond_broadcast(&tq->cond_write[idx]);
206
207
2/2
✓ Branch 0 taken 668 times.
✓ Branch 1 taken 1988318 times.
1988986 if (tq->finished[idx] & FINISHED_RECV) {
208 668 (tq->type == THREAD_QUEUE_FRAMES) ?
209
2/2
✓ Branch 0 taken 561 times.
✓ Branch 1 taken 107 times.
668 av_frame_unref(data) : av_packet_unref(data);
210 668 continue;
211 }
212
213 1988318 *stream_idx = idx;
214 1988318 return 0;
215 }
216
217
2/2
✓ Branch 0 taken 1452839 times.
✓ Branch 1 taken 1144476 times.
2597315 for (unsigned int i = 0; i < tq->nb_streams; i++) {
218
2/2
✓ Branch 0 taken 1413986 times.
✓ Branch 1 taken 38853 times.
1452839 if (!tq->finished[i])
219 1413986 continue;
220
221 /* return EOF to the consumer at most once for each stream */
222
2/2
✓ Branch 0 taken 21245 times.
✓ Branch 1 taken 17608 times.
38853 if (!(tq->finished[i] & FINISHED_RECV)) {
223 21245 tq->finished[i] |= FINISHED_RECV;
224 21245 *stream_idx = i;
225 21245 return AVERROR_EOF;
226 }
227
228 17608 nb_finished++;
229 }
230
231
2/2
✓ Branch 0 taken 9166 times.
✓ Branch 1 taken 1135310 times.
1144476 return nb_finished == tq->nb_streams ? AVERROR_EOF : AVERROR(EAGAIN);
232 }
233
234 2050893 int tq_receive(ThreadQueue *tq, int *stream_idx, void *data, int flags)
235 {
236 int ret;
237
238 2050893 *stream_idx = -1;
239
240 2050893 pthread_mutex_lock(&tq->lock);
241
242 while (1) {
243 3161526 ret = receive_locked(tq, stream_idx, data);
244
245
4/4
✓ Branch 0 taken 1142797 times.
✓ Branch 1 taken 2018729 times.
✓ Branch 2 taken 1110633 times.
✓ Branch 3 taken 32164 times.
3161526 if (ret == AVERROR(EAGAIN) && !(flags & THREAD_QUEUE_FLAG_NO_BLOCK)) {
246 1110633 pthread_cond_wait(&tq->cond_read, &tq->lock);
247 1110633 continue;
248 }
249
250 2050893 break;
251 }
252
253 2050893 pthread_mutex_unlock(&tq->lock);
254
255 2050893 return ret;
256 }
257
258 51207 void tq_send_finish(ThreadQueue *tq, unsigned int stream_idx)
259 {
260
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 51207 times.
51207 av_assert0(stream_idx < tq->nb_streams);
261
262 51207 pthread_mutex_lock(&tq->lock);
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 51207 tq->finished[stream_idx] |= FINISHED_SEND;
268 51207 tq->choked = 0;
269 51207 pthread_cond_broadcast(&tq->cond_read);
270
271 51207 pthread_mutex_unlock(&tq->lock);
272 51207 }
273
274 41402 void tq_receive_finish(ThreadQueue *tq, unsigned int stream_idx)
275 {
276
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 41402 times.
41402 av_assert0(stream_idx < tq->nb_streams);
277
278 41402 pthread_mutex_lock(&tq->lock);
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 41402 tq->finished[stream_idx] |= FINISHED_RECV;
284 41402 pthread_cond_broadcast(&tq->cond_write[stream_idx]);
285
286 41402 pthread_mutex_unlock(&tq->lock);
287 41402 }
288
289 27436 void tq_choke(ThreadQueue *tq, int choked)
290 {
291 27436 pthread_mutex_lock(&tq->lock);
292
293 27436 int prev_choked = tq->choked;
294 27436 tq->choked = choked;
295
2/2
✓ Branch 0 taken 24468 times.
✓ Branch 1 taken 2968 times.
27436 if (choked != prev_choked)
296 24468 pthread_cond_broadcast(&tq->cond_read);
297
298 27436 pthread_mutex_unlock(&tq->lock);
299 27436 }
300