@@ -1061,7 +1061,7 @@ static int process_input(int file_index, AVPacket *pkt)
InputStream *ist;
int ret, i;
- ret = ifile_get_packet(ifile, pkt);
+ ret = 0;
if (ret == 1) {
/* the input file is looped: flush the decoders */
@@ -860,18 +860,6 @@ int64_t of_filesize(OutputFile *of);
int ifile_open(const OptionsContext *o, const char *filename, Scheduler *sch);
void ifile_close(InputFile **f);
-/**
- * Get next input packet from the demuxer.
- *
- * @param pkt the packet is written here when this function returns 0
- * @return
- * - 0 when a packet has been read successfully
- * - 1 when stream end was reached, but the stream is looped;
- * caller should flush decoders and read from this demuxer again
- * - a negative error code on failure
- */
-int ifile_get_packet(InputFile *f, AVPacket *pkt);
-
int ist_output_add(InputStream *ist, OutputStream *ost);
int ist_filter_add(InputStream *ist, InputFilter *ifilter, int is_simple);
@@ -21,8 +21,6 @@
#include "ffmpeg.h"
#include "ffmpeg_sched.h"
-#include "objpool.h"
-#include "thread_queue.h"
#include "libavutil/avassert.h"
#include "libavutil/avstring.h"
@@ -34,7 +32,6 @@
#include "libavutil/pixdesc.h"
#include "libavutil/time.h"
#include "libavutil/timestamp.h"
-#include "libavutil/thread.h"
#include "libavcodec/packet.h"
@@ -65,6 +62,9 @@ typedef struct DemuxStream {
double ts_scale;
+ // scheduler returned EOF for this stream
+ int finished;
+
int streamcopy_needed;
int wrap_correction_done;
@@ -115,11 +115,10 @@ typedef struct Demuxer {
double readrate_initial_burst;
Scheduler *sch;
- ThreadQueue *thread_queue;
- int thread_queue_size;
- pthread_t thread;
int read_started;
+ int nb_streams_used;
+ int nb_streams_finished;
} Demuxer;
static DemuxStream *ds_from_ist(InputStream *ist)
@@ -503,6 +502,8 @@ static int input_packet_process(Demuxer *d, AVPacket *pkt)
av_ts2timestr(input_files[ist->file_index]->ts_offset, &AV_TIME_BASE_Q));
}
+ pkt->stream_index = ds->sch_idx_stream;
+
return 0;
}
@@ -565,9 +566,14 @@ static void *input_thread(void *arg)
discard_unused_programs(f);
+ // XXX
+ d->read_started = 1;
+
d->wallclock_start = av_gettime_relative();
while (1) {
+ DemuxStream *ds;
+
ret = av_read_frame(f->ctx, pkt);
if (ret == AVERROR(EAGAIN)) {
@@ -575,25 +581,32 @@ static void *input_thread(void *arg)
continue;
}
if (ret < 0) {
-#if 0
if (d->loop) {
- /* signal looping to the consumer thread */
- pkt->opaque = (void*)(intptr_t)PKT_OPAQUE_SEEK;
- ret = tq_send(d->thread_queue, 0, pkt);
- if (ret >= 0)
- ret = seek_to_start(d);
+ /* signal looping to our consumers */
+ for (int i = 0; i < f->nb_streams; i++) {
+ pkt->opaque = (void*)(intptr_t)PKT_OPAQUE_SEEK;
+ pkt->stream_index = i;
+
+ ret = sch_demux_send(d->sch, f->index, pkt);
+ // XXX
+ if (ret >= 0)
+ //ret = seek_to_start(d);
+ if (ret < 0)
+ break;
+ }
if (ret >= 0)
continue;
/* fallthrough to the error path */
}
-#endif
if (ret == AVERROR_EOF)
av_log(d, AV_LOG_VERBOSE, "EOF while reading input\n");
- else
+ else {
av_log(d, AV_LOG_ERROR, "Error during demuxing: %s\n",
av_err2str(ret));
+ ret = exit_on_error ? ret : 0;
+ }
break;
}
@@ -605,8 +618,9 @@ static void *input_thread(void *arg)
/* the following test is needed in case new streams appear
dynamically in stream : we ignore them */
- if (pkt->stream_index >= f->nb_streams ||
- f->streams[pkt->stream_index]->discard) {
+ ds = pkt->stream_index < f->nb_streams ?
+ ds_from_ist(f->streams[pkt->stream_index]) : NULL;
+ if (!ds || ds->ist.discard || ds->finished) {
report_new_stream(d, pkt);
av_packet_unref(pkt);
continue;
@@ -630,40 +644,47 @@ static void *input_thread(void *arg)
if (f->readrate)
readrate_sleep(d);
- ret = tq_send(d->thread_queue, 0, pkt);
- if (ret < 0) {
- if (ret != AVERROR_EOF)
+ ret = sch_demux_send(d->sch, f->index, pkt);
+ if (ret == AVERROR_EOF) {
+ av_packet_unref(pkt);
+
+ av_log(ds, AV_LOG_VERBOSE, "All consumers done\n");
+ ds->finished = 1;
+
+ if (++d->nb_streams_finished == d->nb_streams_used) {
+ av_log(f, AV_LOG_VERBOSE, "All streams' consumers done\n");
+ break;
+ }
+ continue;
+ } else if (ret < 0) {
+ if (ret != AVERROR_EXIT)
av_log(f, AV_LOG_ERROR,
- "Unable to send packet to main thread: %s\n",
+ "Unable to send demuxed packet to consumers: %s\n",
av_err2str(ret));
break;
}
}
+ // EOF/EXIT is normal termination
+ if (ret == AVERROR_EOF || ret == AVERROR_EXIT)
+ ret = 0;
+
finish:
- av_assert0(ret < 0);
- tq_send_finish(d->thread_queue, 0);
+ sch_demux_send(d->sch, f->index, NULL);
av_packet_free(&pkt);
av_log(d, AV_LOG_VERBOSE, "Terminating demuxer thread\n");
- return NULL;
+ return (void*)(intptr_t)ret;
}
+// XXX
+#if 0
static void thread_stop(Demuxer *d)
{
InputFile *f = &d->f;
- if (!d->thread_queue)
- return;
-
- tq_receive_finish(d->thread_queue, 0);
-
- pthread_join(d->thread, NULL);
-
- tq_free(&d->thread_queue);
-
//av_thread_message_queue_free(&f->audio_duration_queue);
}
@@ -671,22 +692,7 @@ static int thread_start(Demuxer *d)
{
int ret;
InputFile *f = &d->f;
- ObjPool *op;
- if (d->thread_queue_size <= 0)
- d->thread_queue_size = (nb_input_files > 1 ? 8 : 1);
-
- op = objpool_alloc_packets();
- if (!op)
- return AVERROR(ENOMEM);
-
- d->thread_queue = tq_alloc(1, d->thread_queue_size, op, pkt_move);
- if (!d->thread_queue) {
- objpool_free(&op);
- return AVERROR(ENOMEM);
- }
-
-#if 0
if (d->loop) {
int nb_audio_dec = 0;
@@ -704,20 +710,6 @@ static int thread_start(Demuxer *d)
f->audio_duration_queue_size = nb_audio_dec;
}
}
-#endif
-
- if ((ret = pthread_create(&d->thread, NULL, input_thread, d))) {
- av_log(d, AV_LOG_ERROR, "pthread_create failed: %s. Try to increase `ulimit -v` or decrease `ulimit -s`.\n", strerror(ret));
- ret = AVERROR(ret);
- goto fail;
- }
-
- d->read_started = 1;
-
- return 0;
-fail:
- tq_free(&d->thread_queue);
- return ret;
}
int ifile_get_packet(InputFile *f, AVPacket *pkt)
@@ -725,12 +717,6 @@ int ifile_get_packet(InputFile *f, AVPacket *pkt)
Demuxer *d = demuxer_from_ifile(f);
int ret, dummy;
- if (!d->thread_queue) {
- ret = thread_start(d);
- if (ret < 0)
- return ret;
- }
-
ret = tq_receive(d->thread_queue, &dummy, pkt);
if (ret < 0)
return ret;
@@ -744,6 +730,7 @@ int ifile_get_packet(InputFile *f, AVPacket *pkt)
return 0;
}
+#endif
static void demux_final_stats(Demuxer *d)
{
@@ -813,8 +800,6 @@ void ifile_close(InputFile **pf)
if (!f)
return;
- thread_stop(d);
-
if (d->read_started)
demux_final_stats(d);
@@ -846,7 +831,11 @@ static int ist_use(InputStream *ist, int decoding_needed)
ds->sch_idx_stream = ret;
}
- ist->discard = 0;
+ if (ist->discard) {
+ ist->discard = 0;
+ d->nb_streams_used++;
+ }
+
ist->st->discard = ist->user_set_discard;
ist->decoding_needed |= decoding_needed;
ds->streamcopy_needed |= !decoding_needed;
@@ -1647,8 +1636,6 @@ int ifile_open(const OptionsContext *o, const char *filename, Scheduler *sch)
"since neither -readrate nor -re were given\n");
}
- d->thread_queue_size = o->thread_queue_size;
-
/* Add all the streams from the given input file to the demuxer */
for (int i = 0; i < ic->nb_streams; i++) {
ret = ist_add(o, d, ic->streams[i]);