/* transcode.c -- see transcode.h
 *
 * THREADING MODEL (important):
 * libavformat's custom-IO read callback is expected to BLOCK until data
 * is available or real EOF -- it is not designed for an EAGAIN/"try me
 * later" protocol (avformat_find_stream_info() and the steady-state
 * av_read_frame() path will busy-spin on EAGAIN instead of backing off,
 * which pegs a CPU core and never makes progress).
 *
 * So: tc_write() (called from the network receive thread) only pushes
 * bytes into a mutex+cond-guarded ring buffer and returns immediately --
 * it never touches libavformat/libavcodec itself, and never blocks the
 * caller. All demux/decode/encode/mux work happens on a dedicated
 * per-channel pump thread started by tc_start(), whose blocking read
 * callback (ring_read_cb) sleeps on the condvar until enough bytes
 * exist or tc_stop() requests shutdown. This keeps the NIC receive
 * thread (other channels' siblings) never stalled by encode work,
 * exactly like the original ffmpeg-subprocess design intended. */
#define _GNU_SOURCE
#include "transcode.h"
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <pthread.h>
#include <errno.h>
#include <time.h>

#include <libavformat/avformat.h>
#include <libavformat/avio.h>
#include <libavcodec/avcodec.h>
#include <libswresample/swresample.h>
#include <libavutil/opt.h>
#include <libavutil/audio_fifo.h>
#include <libavutil/channel_layout.h>
#include <libavutil/error.h>

#define TS188        188
#define IN_RINGSZ    (8<<20)      /* 8MB ring for incoming raw TS bytes */
#define CUSTOM_IOBUF (64*1024)    /* libav custom-IO buffer chunk size */
#define MAX_ASTREAMS 16

/* -------- one decode/encode chain per audio input stream -------- */
typedef struct {
    int           in_idx;          /* index into AVFormatContext->streams (input) */
    AVCodecContext *dec_ctx;
    AVCodecContext *enc_ctx;
    SwrContext    *swr;
    AVAudioFifo   *fifo;           /* resampled samples awaiting full encoder frames */
    int           out_stream_idx;  /* index into output AVFormatContext->streams */
    int64_t       next_pts;        /* in encoder time_base units */
    int           active;
} AStream;

struct Transcoder {
    TcOpts   opts;
    char     name[64];
    tc_pkt_cb cb;
    void     *user;

    /* ---- custom input IO: ring buffer fed by tc_write() ---- */
    uint8_t  *ring;
    size_t    ring_cap;
    size_t    ring_head, ring_tail; /* byte offsets, wrap via %cap */
    size_t    ring_fill;
    int       eof_requested;
    pthread_mutex_t ring_mtx;
    pthread_cond_t  ring_cv;        /* signaled whenever ring_fill grows, eof_requested set,
                                        stop_requested set, OR probe_done transitions */
    pthread_t       pump_tid;
    int             pump_running;
    int             stop_requested;
    int             probe_done;     /* 1 once pump thread has finished avformat_open_input +
                                        avformat_find_stream_info (success OR failure) */

    AVIOContext     *in_avio;
    AVFormatContext *ifmt;

    /* ---- video chain ---- */
    int             v_in_idx;
    AVCodecContext *v_dec_ctx;
    AVCodecContext *v_enc_ctx;
    int             v_out_idx;
    int             have_video;
    int64_t         v_next_pts;

    /* ---- audio chains (one per mapped input audio stream) ---- */
    AStream  astreams[MAX_ASTREAMS];
    int      n_astreams;

    /* ---- custom output IO: emits 188B TS packets via cb ---- */
    AVIOContext      *out_avio;
    AVFormatContext  *ofmt;
    uint8_t           out_carry[TS188];
    int               out_carry_n;

    int      started;     /* header written */
    int      failed;      /* unrecoverable error -- tc_write becomes no-op */
    uint64_t drop_bytes;  /* total bytes dropped due to ring overflow, mutex-protected */
};

/* ================================================================
   RING BUFFER (input) -- mutex+condvar guarded, real blocking reads
   ================================================================ */

/* Called from the network thread via tc_write(). Never blocks the
 * caller: drops oldest bytes on overflow rather than waiting. */
static void ring_push(Transcoder *t, const uint8_t *data, int len){
    pthread_mutex_lock(&t->ring_mtx);
    if((size_t)len > t->ring_cap - t->ring_fill) {
        size_t need = (size_t)len - (t->ring_cap - t->ring_fill);
        t->ring_head = (t->ring_head + need) % t->ring_cap;
        t->ring_fill -= need;
        t->drop_bytes += need;
    }
    size_t first = t->ring_cap - t->ring_tail;
    if(first >= (size_t)len){
        memcpy(t->ring + t->ring_tail, data, len);
    } else {
        memcpy(t->ring + t->ring_tail, data, first);
        memcpy(t->ring, data + first, (size_t)len - first);
    }
    t->ring_tail = (t->ring_tail + (size_t)len) % t->ring_cap;
    t->ring_fill += (size_t)len;
    pthread_cond_signal(&t->ring_cv);
    pthread_mutex_unlock(&t->ring_mtx);
}

/* Called only from the pump thread (inside libavformat). BLOCKS until
 * data is available, EOF is requested, or shutdown is requested -- this
 * is the behavior avformat_find_stream_info()/av_read_frame() actually
 * expect from a custom-IO source; returning EAGAIN here instead causes
 * libavformat to busy-spin forever rather than back off. */
static int ring_read_cb(void *opaque, uint8_t *buf, int wanted){
    Transcoder *t = (Transcoder*)opaque;
    pthread_mutex_lock(&t->ring_mtx);
    while(t->ring_fill == 0 && !t->eof_requested && !t->stop_requested)
        pthread_cond_wait(&t->ring_cv, &t->ring_mtx);
    if(t->stop_requested){ pthread_mutex_unlock(&t->ring_mtx); return AVERROR_EOF; }
    if(t->ring_fill == 0 && t->eof_requested){ pthread_mutex_unlock(&t->ring_mtx); return AVERROR_EOF; }

    int n = wanted;
    if((size_t)n > t->ring_fill) n = (int)t->ring_fill;
    size_t first = t->ring_cap - t->ring_head;
    if(first >= (size_t)n){
        memcpy(buf, t->ring + t->ring_head, n);
    } else {
        memcpy(buf, t->ring + t->ring_head, first);
        memcpy(buf + first, t->ring, (size_t)n - first);
    }
    t->ring_head = (t->ring_head + (size_t)n) % t->ring_cap;
    t->ring_fill -= (size_t)n;
    pthread_mutex_unlock(&t->ring_mtx);
    return n;
}

/* ================================================================
   OUTPUT: custom-IO write callback -> slice into 188B packets -> cb
   ================================================================ */
static int out_write_cb(void *opaque, uint8_t *buf, int buf_size){
    Transcoder *t = (Transcoder*)opaque;
    int off = 0;
    /* top up any carried partial packet first */
    if(t->out_carry_n > 0){
        int need = TS188 - t->out_carry_n;
        int take = need < buf_size ? need : buf_size;
        memcpy(t->out_carry + t->out_carry_n, buf, take);
        t->out_carry_n += take; off += take;
        if(t->out_carry_n == TS188){
            t->cb(t->user, t->out_carry);
            t->out_carry_n = 0;
        }
    }
    while(buf_size - off >= TS188){
        t->cb(t->user, buf + off);
        off += TS188;
    }
    if(off < buf_size){
        int rem = buf_size - off;
        memcpy(t->out_carry, buf + off, rem);
        t->out_carry_n = rem;
    }
    return buf_size;
}

/* ================================================================
   ENCODER SETUP HELPERS
   ================================================================ */
static AVCodecContext *open_x264_encoder(const TcOpts *o, AVCodecContext *dec_ctx,
                                          AVRational frame_rate, const char *chan){
    const AVCodec *enc = avcodec_find_encoder_by_name("libx264");
    if(!enc){
        fprintf(stderr,"[%s] libx264 encoder not found in this libavcodec build\n", chan);
        return NULL;
    }
    AVCodecContext *c = avcodec_alloc_context3(enc);
    if(!c) return NULL;

    c->width  = dec_ctx->width;
    c->height = dec_ctx->height;
    c->pix_fmt = AV_PIX_FMT_YUV420P;
    c->time_base = av_inv_q(frame_rate.num ? frame_rate : (AVRational){25,1});
    c->framerate = frame_rate.num ? frame_rate : (AVRational){25,1};

    int gop = 50;
    if(o->segment_secs > 0){
        int fr = frame_rate.num && frame_rate.den ? (frame_rate.num/frame_rate.den) : 25;
        gop = fr * o->segment_secs;
        if(gop < 2) gop = 2;
    }
    c->gop_size     = gop;
    c->keyint_min   = gop/2 > 0 ? gop/2 : 1;
    c->max_b_frames = 0; /* keep simple GOP structure -> reliable forced IDR at boundaries */

    if(o->video_crf > 0){
        av_opt_set_int(c->priv_data, "crf", o->video_crf, 0);
    } else {
        long br = 2000000;
        if(o->video_bitrate[0]){
            br = atol(o->video_bitrate) * (strchr(o->video_bitrate,'k')||strchr(o->video_bitrate,'K') ? 1000 : 1);
        }
        c->bit_rate = br;
        c->rc_max_rate = br;
        c->rc_buffer_size = (int)(br/2);
    }
    av_opt_set(c->priv_data, "preset", o->video_preset[0] ? o->video_preset : "ultrafast", 0);
    if(o->video_tune[0])
        av_opt_set(c->priv_data, "tune", o->video_tune, 0);
    av_opt_set(c->priv_data, "profile", "main", 0);
    av_opt_set(c->priv_data, "level", "4.0", 0);
    /* sc_threshold=0 so our forced keyframes (via AV_PKT_FLAG_KEY / gop) are
     * the only keyframe boundaries -- matches the original ffmpeg recipe's
     * -sc_threshold 0, keeping segment cuts predictable. */
    av_opt_set_int(c->priv_data, "sc_threshold", 0, 0);
    if(o->video_threads > 0) c->thread_count = o->video_threads;

    c->flags |= AV_CODEC_FLAG_GLOBAL_HEADER;
    return c;
}

/* Open an AAC encoder, preferring libfdk_aac when requested/available,
 * falling back to native ffmpeg "aac" encoder transparently. */
static AVCodecContext *open_aac_encoder(const TcOpts *o, int sample_rate, int channels,
                                         const char *chan){
    const AVCodec *enc = NULL;
    int using_fdk = 0;
    if(o->prefer_fdk){
        enc = avcodec_find_encoder_by_name("libfdk_aac");
        if(enc) using_fdk = 1;
    }
    if(!enc) enc = avcodec_find_encoder_by_name("aac");
    if(!enc){
        fprintf(stderr,"[%s] no AAC encoder available\n", chan);
        return NULL;
    }
    if(o->prefer_fdk && !using_fdk && o->verbose)
        fprintf(stderr,"[%s] libfdk_aac not available, falling back to native aac encoder\n", chan);

    AVCodecContext *c = avcodec_alloc_context3(enc);
    if(!c) return NULL;

    c->sample_rate = sample_rate > 0 ? sample_rate : 48000;
    av_channel_layout_default(&c->ch_layout, channels > 0 ? channels : 2);
    c->sample_fmt = enc->sample_fmts ? enc->sample_fmts[0] : AV_SAMPLE_FMT_FLTP;
    c->time_base = (AVRational){1, c->sample_rate};

    int cbr = (strcmp(o->aac_mode,"cbr")==0) || !o->aac_mode[0];
    if(cbr){
        long br = 128000;
        if(o->audio_bitrate[0])
            br = atol(o->audio_bitrate) * (strchr(o->audio_bitrate,'k')||strchr(o->audio_bitrate,'K') ? 1000 : 1);
        c->bit_rate = br;
    } else if(using_fdk){
        int vbr = 3;
        if(strncmp(o->aac_mode,"vbr",3)==0){
            int v = o->aac_mode[3]-'0';
            if(v>=1 && v<=5) vbr = v;
        }
        av_opt_set_int(c->priv_data, "vbr", vbr, 0);
    } else {
        /* native aac encoder: no real VBR knob like fdk's 1-5; approximate
         * with a reasonable bitrate so behavior is still sane without fdk. */
        c->bit_rate = 160000;
    }
    c->flags |= AV_CODEC_FLAG_GLOBAL_HEADER;
    return c;
}

/* ================================================================
   STREAM SETUP: walk input streams, decide video/audio mapping,
   open decoders + encoders + output streams per TcOpts cases 1-5.
   ================================================================ */
static int setup_streams(Transcoder *t){
    AVFormatContext *ifmt = t->ifmt;
    const TcOpts *o = &t->opts;

    int audio_seen = 0; /* running count of audio streams encountered, for audio_map matching */

    for(unsigned i=0; i<ifmt->nb_streams; i++){
        AVStream *st = ifmt->streams[i];
        enum AVMediaType type = st->codecpar->codec_type;

        if(type == AVMEDIA_TYPE_VIDEO && !t->have_video){
            t->v_in_idx = (int)i;
            if(o->video_h264){
                const AVCodec *dec = avcodec_find_decoder(st->codecpar->codec_id);
                if(!dec){ fprintf(stderr,"[%s] no decoder for input video codec\n", t->name); continue; }
                t->v_dec_ctx = avcodec_alloc_context3(dec);
                avcodec_parameters_to_context(t->v_dec_ctx, st->codecpar);
                t->v_dec_ctx->pkt_timebase = st->time_base;
                if(avcodec_open2(t->v_dec_ctx, dec, NULL) < 0){
                    fprintf(stderr,"[%s] failed to open video decoder\n", t->name);
                    avcodec_free_context(&t->v_dec_ctx);
                    continue;
                }
                AVRational fr = st->avg_frame_rate.num ? st->avg_frame_rate : (AVRational){25,1};
                t->v_enc_ctx = open_x264_encoder(o, t->v_dec_ctx, fr, t->name);
                if(!t->v_enc_ctx) continue;
                if(avcodec_open2(t->v_enc_ctx, t->v_enc_ctx->codec, NULL) < 0){
                    fprintf(stderr,"[%s] failed to open libx264 encoder\n", t->name);
                    avcodec_free_context(&t->v_enc_ctx);
                    continue;
                }
                AVStream *ost = avformat_new_stream(t->ofmt, NULL);
                avcodec_parameters_from_context(ost->codecpar, t->v_enc_ctx);
                ost->time_base = t->v_enc_ctx->time_base;
                t->v_out_idx = ost->index;
                t->have_video = 1;
            } else {
                /* video copy/passthrough: just mirror the stream into output */
                AVStream *ost = avformat_new_stream(t->ofmt, NULL);
                avcodec_parameters_copy(ost->codecpar, st->codecpar);
                ost->time_base = st->time_base;
                t->v_out_idx = ost->index;
                t->have_video = 1; /* marks "video stream present"; v_enc_ctx stays NULL -> copy path */
            }
            continue;
        }

        if(type == AVMEDIA_TYPE_AUDIO){
            int this_audio_idx = audio_seen++;
            int want = 0;
            if(!o->audio_aac){
                want = 0; /* not transcoding audio at all in this build path */
            } else if(o->audio_map < 0){
                want = 1; /* case 2: ALL audio streams */
            } else {
                want = (o->audio_map == this_audio_idx); /* case 4: one specific stream */
            }
            if(!want) continue;
            if(t->n_astreams >= MAX_ASTREAMS) continue;

            AStream *as = &t->astreams[t->n_astreams];
            memset(as, 0, sizeof *as);
            as->in_idx = (int)i;

            const AVCodec *dec = avcodec_find_decoder(st->codecpar->codec_id);
            if(!dec){
                fprintf(stderr,"[%s] no decoder for input audio stream %d (codec_id=%d)\n",
                        t->name, this_audio_idx, st->codecpar->codec_id);
                continue;
            }
            as->dec_ctx = avcodec_alloc_context3(dec);
            avcodec_parameters_to_context(as->dec_ctx, st->codecpar);
            as->dec_ctx->pkt_timebase = st->time_base;
            if(avcodec_open2(as->dec_ctx, dec, NULL) < 0){
                fprintf(stderr,"[%s] failed to open audio decoder for stream %d\n", t->name, this_audio_idx);
                avcodec_free_context(&as->dec_ctx);
                continue;
            }

            int in_ch = as->dec_ctx->ch_layout.nb_channels > 0 ? as->dec_ctx->ch_layout.nb_channels : 2;
            int in_rate = as->dec_ctx->sample_rate > 0 ? as->dec_ctx->sample_rate : 48000;

            as->enc_ctx = open_aac_encoder(o, in_rate, in_ch, t->name);
            if(!as->enc_ctx){ avcodec_free_context(&as->dec_ctx); continue; }
            if(avcodec_open2(as->enc_ctx, as->enc_ctx->codec, NULL) < 0){
                fprintf(stderr,"[%s] failed to open AAC encoder for stream %d\n", t->name, this_audio_idx);
                avcodec_free_context(&as->dec_ctx);
                avcodec_free_context(&as->enc_ctx);
                continue;
            }

            int ret = swr_alloc_set_opts2(&as->swr,
                    &as->enc_ctx->ch_layout, as->enc_ctx->sample_fmt, as->enc_ctx->sample_rate,
                    &as->dec_ctx->ch_layout, as->dec_ctx->sample_fmt, as->dec_ctx->sample_rate,
                    0, NULL);
            if(ret < 0 || !as->swr || swr_init(as->swr) < 0){
                fprintf(stderr,"[%s] failed to init resampler for audio stream %d\n", t->name, this_audio_idx);
                avcodec_free_context(&as->dec_ctx);
                avcodec_free_context(&as->enc_ctx);
                continue;
            }

            as->fifo = av_audio_fifo_alloc(as->enc_ctx->sample_fmt,
                                            as->enc_ctx->ch_layout.nb_channels, 1);

            AVStream *ost = avformat_new_stream(t->ofmt, NULL);
            avcodec_parameters_from_context(ost->codecpar, as->enc_ctx);
            ost->time_base = as->enc_ctx->time_base;
            as->out_stream_idx = ost->index;
            as->active = 1;
            t->n_astreams++;
            continue;
        }
        /* other stream types (subtitles, data) are dropped -- not requested */
    }

    if(!t->have_video && t->n_astreams == 0){
        fprintf(stderr,"[%s] no streams selected for output (check options)\n", t->name);
        return -1;
    }
    return 0;
}

/* ================================================================
   tc_start
   ================================================================ */
/* forward decls -- defined below, used by pump_thread_fn above their definition */
static void process_video_packet(Transcoder *t, AVPacket *pkt);
static void process_audio_packet(Transcoder *t, AStream *as, AVPacket *pkt);

/* ================================================================
   PUMP THREAD -- owns all libavformat/libavcodec calls. Runs the
   blocking open/probe once, then loops av_read_frame() (which itself
   blocks inside ring_read_cb when starved) until tc_stop() requests
   shutdown or real EOF is signaled.
   ================================================================ */
/* Mark probing finished (success or failure) and wake anyone in tc_stop()
 * waiting on probe_done. Must be called exactly once per pump thread,
 * whether the outcome was success or failure. */
static void mark_probe_done(Transcoder *t, int failed){
    pthread_mutex_lock(&t->ring_mtx);
    if(failed) t->failed = 1;
    t->probe_done = 1;
    pthread_cond_broadcast(&t->ring_cv);
    pthread_mutex_unlock(&t->ring_mtx);
}

static void *pump_thread_fn(void *arg){
    Transcoder *t = (Transcoder*)arg;

    const AVInputFormat *fmt = av_find_input_format("mpegts");
    AVDictionary *d = NULL;
    av_dict_set(&d, "analyzeduration", "1000000", 0);
    av_dict_set(&d, "probesize", "5000000", 0);
    int ret = avformat_open_input(&t->ifmt, NULL, fmt, &d);
    av_dict_free(&d);
    if(ret < 0){
        char errbuf[128]; av_strerror(ret, errbuf, sizeof errbuf);
        fprintf(stderr,"[%s] avformat_open_input failed: %s\n", t->name, errbuf);
        mark_probe_done(t, 1); return NULL;
    }
    if(avformat_find_stream_info(t->ifmt, NULL) < 0){
        fprintf(stderr,"[%s] avformat_find_stream_info failed\n", t->name);
        mark_probe_done(t, 1); return NULL;
    }

    if(avformat_alloc_output_context2(&t->ofmt, NULL, "mpegts", NULL) < 0){
        fprintf(stderr,"[%s] avformat_alloc_output_context2 failed\n", t->name);
        mark_probe_done(t, 1); return NULL;
    }
    uint8_t *outbuf = av_malloc(CUSTOM_IOBUF);
    t->out_avio = avio_alloc_context(outbuf, CUSTOM_IOBUF, 1, t, NULL, out_write_cb, NULL);
    if(!t->out_avio){
        fprintf(stderr,"[%s] avio_alloc_context (out) failed\n", t->name);
        mark_probe_done(t, 1); return NULL;
    }
    t->ofmt->pb = t->out_avio;
    t->ofmt->flags |= AVFMT_FLAG_CUSTOM_IO;

    if(setup_streams(t) < 0){ mark_probe_done(t, 1); return NULL; }

    AVDictionary *muxopts = NULL;
    av_dict_set(&muxopts, "mpegts_flags", "resend_headers", 0);
    if(avformat_write_header(t->ofmt, &muxopts) < 0){
        fprintf(stderr,"[%s] avformat_write_header failed\n", t->name);
        av_dict_free(&muxopts);
        mark_probe_done(t, 1); return NULL;
    }
    av_dict_free(&muxopts);
    t->started = 1;
    if(t->opts.verbose)
        fprintf(stderr,"[%s] transcoder live: video=%s audio_streams=%d\n",
                t->name, t->have_video ? (t->v_enc_ctx?"libx264":"copy") : "none", t->n_astreams);
    mark_probe_done(t, 0);


    AVPacket *pkt = av_packet_alloc();
    for(;;){
        pthread_mutex_lock(&t->ring_mtx);
        int stop = t->stop_requested;
        pthread_mutex_unlock(&t->ring_mtx);
        if(stop) break;

        int rret = av_read_frame(t->ifmt, pkt);
        if(rret < 0) break; /* real EOF or unrecoverable demux error */

        if(pkt->stream_index == t->v_in_idx){
            process_video_packet(t, pkt);
        } else {
            for(int i=0;i<t->n_astreams;i++){
                if(t->astreams[i].active && t->astreams[i].in_idx == pkt->stream_index){
                    process_audio_packet(t, &t->astreams[i], pkt);
                    break;
                }
            }
        }
        av_packet_unref(pkt);
    }
    av_packet_free(&pkt);
    return NULL;
}

Transcoder *tc_start(const TcOpts *opts, tc_pkt_cb cb, void *user, const char *chan_name){
    Transcoder *t = calloc(1, sizeof *t);
    if(!t) return NULL;
    t->opts = *opts;
    t->cb = cb;
    t->user = user;
    snprintf(t->name, sizeof t->name, "%s", chan_name ? chan_name : "ch");
    t->v_out_idx = -1;

    t->ring_cap = IN_RINGSZ;
    t->ring = malloc(t->ring_cap);
    if(!t->ring){ free(t); return NULL; }
    pthread_mutex_init(&t->ring_mtx, NULL);
    pthread_cond_init(&t->ring_cv, NULL);

    uint8_t *inbuf = av_malloc(CUSTOM_IOBUF);
    /* read_packet given, write_packet NULL: this AVIOContext is read-only (input side) */
    t->in_avio = avio_alloc_context(inbuf, CUSTOM_IOBUF, 0, t, ring_read_cb, NULL, NULL);
    if(!t->in_avio){
        fprintf(stderr,"[%s] avio_alloc_context (in) failed\n", t->name);
        free(t->ring); free(t);
        return NULL;
    }

    t->ifmt = avformat_alloc_context();
    t->ifmt->pb = t->in_avio;
    t->ifmt->flags |= AVFMT_FLAG_CUSTOM_IO;

    t->failed = 0;
    t->pump_running = 1;
    if(pthread_create(&t->pump_tid, NULL, pump_thread_fn, t) != 0){
        fprintf(stderr,"[%s] failed to start pump thread\n", t->name);
        t->pump_running = 0;
        av_freep(&t->in_avio->buffer);
        avio_context_free(&t->in_avio);
        avformat_free_context(t->ifmt); t->ifmt = NULL;
        free(t->ring); free(t);
        return NULL;
    }
    return t;
}

/* ================================================================
   ENCODE/MUX HELPERS
   ================================================================ */
static void mux_encoded_packets(Transcoder *t, AVCodecContext *enc_ctx, int out_idx,
                                 AVRational enc_tb){
    AVPacket *pkt = av_packet_alloc();
    int rc;
    while((rc = avcodec_receive_packet(enc_ctx, pkt)) == 0){
        pkt->stream_index = out_idx;
        av_packet_rescale_ts(pkt, enc_tb, t->ofmt->streams[out_idx]->time_base);
        int wret = av_interleaved_write_frame(t->ofmt, pkt);
        if(wret < 0 && t->opts.verbose){
            char eb[128]; av_strerror(wret, eb, sizeof eb);
            fprintf(stderr,"[%s] mux write failed (stream %d): %s\n", t->name, out_idx, eb);
        }
        av_packet_unref(pkt);
    }
    if(rc != AVERROR(EAGAIN) && rc != AVERROR_EOF && t->opts.verbose){
        char eb[128]; av_strerror(rc, eb, sizeof eb);
        fprintf(stderr,"[%s] encoder error: %s\n", t->name, eb);
    }
    av_packet_free(&pkt);
}

static void process_video_packet(Transcoder *t, AVPacket *pkt){
    if(!t->have_video) return;
    if(!t->v_enc_ctx){
        /* copy mode: just rewrite stream index + timebase and pass through */
        AVStream *ist = t->ifmt->streams[t->v_in_idx];
        AVPacket *op = av_packet_clone(pkt);
        op->stream_index = t->v_out_idx;
        av_packet_rescale_ts(op, ist->time_base, t->ofmt->streams[t->v_out_idx]->time_base);
        av_interleaved_write_frame(t->ofmt, op);
        av_packet_free(&op);
        return;
    }
    if(avcodec_send_packet(t->v_dec_ctx, pkt) < 0) return;
    AVFrame *frm = av_frame_alloc();
    while(avcodec_receive_frame(t->v_dec_ctx, frm) == 0){
        /* best_effort_timestamp is in the DECODER's time_base (typically
         * the input stream's 1/90000 MPEG-TS clock). The encoder has its
         * own time_base (e.g. 1/25 for 25fps). Feeding the raw decoder
         * PTS straight into the encoder without rescaling corrupts every
         * downstream duration/timestamp by the time_base ratio -- this
         * previously inflated reported file duration by ~3600x. */
        int64_t ts = frm->best_effort_timestamp;
        if(ts == AV_NOPTS_VALUE) ts = frm->pts;
        frm->pts = av_rescale_q(ts, t->v_dec_ctx->pkt_timebase, t->v_enc_ctx->time_base);
        if(avcodec_send_frame(t->v_enc_ctx, frm) == 0)
            mux_encoded_packets(t, t->v_enc_ctx, t->v_out_idx, t->v_enc_ctx->time_base);
        av_frame_unref(frm);
    }
    av_frame_free(&frm);
}

static void process_audio_packet(Transcoder *t, AStream *as, AVPacket *pkt){
    if(avcodec_send_packet(as->dec_ctx, pkt) < 0) return;
    AVFrame *frm = av_frame_alloc();
    while(avcodec_receive_frame(as->dec_ctx, frm) == 0){
        /* resample into encoder's expected format/rate/layout */
        uint8_t **converted = NULL;
        int out_samples = (int)av_rescale_rnd(swr_get_delay(as->swr, as->dec_ctx->sample_rate) + frm->nb_samples,
                                               as->enc_ctx->sample_rate, as->dec_ctx->sample_rate, AV_ROUND_UP);
        av_samples_alloc_array_and_samples(&converted, NULL, as->enc_ctx->ch_layout.nb_channels,
                                            out_samples, as->enc_ctx->sample_fmt, 0);
        int got = swr_convert(as->swr, converted, out_samples,
                               (const uint8_t**)frm->extended_data, frm->nb_samples);
        if(got > 0 && av_audio_fifo_realloc(as->fifo, av_audio_fifo_size(as->fifo) + got) == 0)
            av_audio_fifo_write(as->fifo, (void**)converted, got);
        av_freep(&converted[0]); av_freep(&converted);
        av_frame_unref(frm);

        /* drain fixed-size frames to the encoder as they become available */
        int frame_sz = as->enc_ctx->frame_size > 0 ? as->enc_ctx->frame_size : 1024;
        while(av_audio_fifo_size(as->fifo) >= frame_sz){
            AVFrame *ef = av_frame_alloc();
            ef->nb_samples = frame_sz;
            ef->format = as->enc_ctx->sample_fmt;
            ef->sample_rate = as->enc_ctx->sample_rate;
            av_channel_layout_copy(&ef->ch_layout, &as->enc_ctx->ch_layout);
            av_frame_get_buffer(ef, 0);
            av_audio_fifo_read(as->fifo, (void**)ef->data, frame_sz);
            ef->pts = as->next_pts;
            as->next_pts += frame_sz;
            int sret = avcodec_send_frame(as->enc_ctx, ef);
            if(sret == 0)
                mux_encoded_packets(t, as->enc_ctx, as->out_stream_idx, as->enc_ctx->time_base);
            else if(t->opts.verbose) {
                char eb[128]; av_strerror(sret, eb, sizeof eb);
                fprintf(stderr,"[%s] audio encoder send_frame failed: %s\n", t->name, eb);
            }
            av_frame_free(&ef);
        }
    }
    av_frame_free(&frm);
}

/* ================================================================
   tc_write -- feed raw bytes into the ring; returns immediately.
   All demux/decode/encode/mux work happens on the pump thread.
   ================================================================ */
void tc_write(Transcoder *t, const uint8_t *data, int len){
    if(!t || t->failed || len <= 0) return;
    ring_push(t, data, len);
}

/* ================================================================
   tc_stop -- signal shutdown, join the pump thread, then flush
   encoders / write trailer / free everything. Safe because once the
   pump thread has exited, nothing else touches t concurrently.
   ================================================================ */
void tc_stop(Transcoder *t){
    if(!t) return;

    if(t->pump_running){
        /* Wait for the pump thread to finish its open/probe phase (success
         * or failure) before requesting shutdown. Without this, a stop
         * issued immediately after start can race the probe and the
         * channel would end up with zero streams selected. Bounded to 5s
         * so a genuinely stuck/garbage source can't hang shutdown. */
        struct timespec deadline;
        clock_gettime(CLOCK_REALTIME, &deadline);
        deadline.tv_sec += 5;
        pthread_mutex_lock(&t->ring_mtx);
        while(!t->probe_done){
            if(pthread_cond_timedwait(&t->ring_cv, &t->ring_mtx, &deadline) == ETIMEDOUT)
                break;
        }

        /* Signal "no more bytes are coming" (graceful EOF), NOT an abrupt
         * stop -- the ring may still hold data tc_write() already
         * accepted (e.g. the whole rest of a short test file, or
         * whatever arrived just before the caller decided to stop) and
         * the demuxer may also be holding already-parsed packets
         * internally from the probe phase. Both must still be drained
         * through process_*_packet()/mux or that audio/video is silently
         * lost. eof_requested makes ring_read_cb return real EOF only
         * once the ring is actually empty, which lets av_read_frame's
         * loop in pump_thread_fn finish naturally. */
        t->eof_requested = 1;
        pthread_cond_signal(&t->ring_cv);
        pthread_mutex_unlock(&t->ring_mtx);

        /* Bounded wait for the pump thread to drain the ring naturally
         * and exit its read loop via real EOF. If it's still running
         * after the grace period (e.g. genuinely stuck), force it via
         * stop_requested so shutdown can't hang forever. */
        struct timespec join_deadline;
        clock_gettime(CLOCK_REALTIME, &join_deadline);
        join_deadline.tv_sec += 5;
#if defined(__USE_GNU) || defined(_GNU_SOURCE)
        int jret = pthread_timedjoin_np(t->pump_tid, NULL, &join_deadline);
#else
        int jret = pthread_join(t->pump_tid, NULL); /* portable fallback: unbounded */
#endif
        if(jret != 0){
            pthread_mutex_lock(&t->ring_mtx);
            t->stop_requested = 1;
            pthread_cond_signal(&t->ring_cv);
            pthread_mutex_unlock(&t->ring_mtx);
            pthread_join(t->pump_tid, NULL);
        }
        t->pump_running = 0;
    }

    if(t->started){
        if(t->v_enc_ctx){
            avcodec_send_frame(t->v_enc_ctx, NULL);
            mux_encoded_packets(t, t->v_enc_ctx, t->v_out_idx, t->v_enc_ctx->time_base);
        }
        for(int i=0;i<t->n_astreams;i++){
            AStream *as = &t->astreams[i];
            if(!as->active) continue;
            avcodec_send_frame(as->enc_ctx, NULL);
            mux_encoded_packets(t, as->enc_ctx, as->out_stream_idx, as->enc_ctx->time_base);
        }
        av_write_trailer(t->ofmt);
        if(t->out_carry_n > 0){
            /* pad final partial TS packet with stuffing so we don't emit a
             * short packet to pkt_write(), which assumes 188-byte units */
            memset(t->out_carry + t->out_carry_n, 0xFF, TS188 - t->out_carry_n);
            t->cb(t->user, t->out_carry);
            t->out_carry_n = 0;
        }
    }
    if(t->v_dec_ctx) avcodec_free_context(&t->v_dec_ctx);
    if(t->v_enc_ctx) avcodec_free_context(&t->v_enc_ctx);
    for(int i=0;i<t->n_astreams;i++){
        AStream *as = &t->astreams[i];
        if(as->dec_ctx) avcodec_free_context(&as->dec_ctx);
        if(as->enc_ctx) avcodec_free_context(&as->enc_ctx);
        if(as->swr) swr_free(&as->swr);
        if(as->fifo) av_audio_fifo_free(as->fifo);
    }
    if(t->ofmt){
        if(t->out_avio) av_freep(&t->out_avio->buffer);
        avio_context_free(&t->out_avio);
        avformat_free_context(t->ofmt);
    }
    if(t->ifmt){
        if(t->ifmt->iformat) avformat_close_input(&t->ifmt);
        else avformat_free_context(t->ifmt);
    }
    if(t->in_avio){
        av_freep(&t->in_avio->buffer);
        avio_context_free(&t->in_avio);
    }
    pthread_mutex_destroy(&t->ring_mtx);
    pthread_cond_destroy(&t->ring_cv);
    free(t->ring);
    free(t);
}

uint64_t tc_drop_count(const Transcoder *t){
    if(!t) return 0;
    /* drop_bytes is written under ring_mtx by the writer thread(s); a
     * relaxed read here is fine for a stats counter (matches the old
     * pipe_drops counter's consistency guarantees, which had none either). */
    return t->drop_bytes;
}
