This is an automated email from the git hooks/post-receive script.

Git pushed a commit to branch master
in repository ffmpeg.

commit 15e07fbcd1739acffbcd9be21df29d2901b9df1e
Author:     Kacper Michajłow <[email protected]>
AuthorDate: Mon Jun 15 02:52:11 2026 +0200
Commit:     Kacper Michajłow <[email protected]>
CommitDate: Mon Jul 27 17:05:20 2026 +0000

    avformat/libcurl: implement event loop and read path
    
    Add the per-AVFormatContext curl_multi worker thread, the command queue,
    write/header/xferinfo callbacks and a FIFO with pause/unpause
    backpressure. url_open2 probes the response on a dedicated thread.
    url_read drains the FIFO and yields to the avio layer so the interrupt
    callback is honoured. An explicit "libcurl:" URL prefix forces the
    protocol.
    
    Signed-off-by: Kacper Michajłow <[email protected]>
---
 configure              |   2 +-
 libavformat/avformat.c |   5 +
 libavformat/internal.h |  15 +-
 libavformat/libcurl.c  | 570 ++++++++++++++++++++++++++++++++++++++++++++++++-
 4 files changed, 584 insertions(+), 8 deletions(-)

diff --git a/configure b/configure
index c9d8736dac..61a948d14d 100755
--- a/configure
+++ b/configure
@@ -4117,7 +4117,7 @@ ipns_gateway_protocol_select="https_protocol"
 libamqp_protocol_deps="librabbitmq"
 libamqp_protocol_select="network"
 librist_protocol_deps="librist"
-libcurl_protocol_deps="libcurl"
+libcurl_protocol_deps="libcurl threads"
 libcurl_protocol_select="network"
 librist_protocol_select="network"
 librtmp_protocol_deps="librtmp"
diff --git a/libavformat/avformat.c b/libavformat/avformat.c
index db9fad2f2d..ae12f00975 100644
--- a/libavformat/avformat.c
+++ b/libavformat/avformat.c
@@ -19,6 +19,8 @@
  * Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
  */
 
+#include "config_components.h"
+
 #include <math.h>
 #include "libavutil/avassert.h"
 #include "libavutil/avstring.h"
@@ -187,6 +189,9 @@ void avformat_free_context(AVFormatContext *s)
     av_freep(&s->chapters);
     av_dict_free(&s->metadata);
     av_dict_free(&si->id3v2_meta);
+#if CONFIG_LIBCURL_PROTOCOL
+    ff_curl_loop_free(&si->curl_loop);
+#endif
     av_packet_free(&si->pkt);
     av_packet_free(&si->parse_pkt);
     ff_packet_list_free(&si->packet_buffer);
diff --git a/libavformat/internal.h b/libavformat/internal.h
index 759f1a9492..b77292a532 100644
--- a/libavformat/internal.h
+++ b/libavformat/internal.h
@@ -23,8 +23,6 @@
 
 #include <stdint.h>
 
-#include "config_components.h"
-
 #include "packet_internal.h"
 
 #include "avformat.h"
@@ -120,6 +118,13 @@ typedef struct FFFormatContext {
     AVDictionary *id3v2_meta;
 
     int missing_streams;
+
+    /**
+     * Shared libcurl event loop, created on demand on the first use. Freed on
+     * context free. This allows to share libcurl state across URLContexts,
+     * scoped to this context.
+     */
+    struct CurlLoop *curl_loop;
 } FFFormatContext;
 
 static av_always_inline FFFormatContext *ffformatcontext(AVFormatContext *s)
@@ -607,6 +612,12 @@ int ff_copy_whiteblacklists(AVFormatContext *dst, const 
AVFormatContext *src);
  */
 int ff_format_io_close(AVFormatContext *s, AVIOContext **pb);
 
+/**
+ * Release a libcurl event loop and set *loop to NULL.
+ * No-op when @p loop or *loop is NULL.
+ */
+void ff_curl_loop_free(struct CurlLoop **loop);
+
 /**
  * Utility function to check if the file uses http or https protocol
  *
diff --git a/libavformat/libcurl.c b/libavformat/libcurl.c
index b4a4995873..ebb00ede13 100644
--- a/libavformat/libcurl.c
+++ b/libavformat/libcurl.c
@@ -23,29 +23,589 @@
 
 #include <curl/curl.h>
 
+#include "libavutil/avstring.h"
+#include "libavutil/error.h"
+#include "libavutil/fifo.h"
+#include "libavutil/macros.h"
+#include "libavutil/mem.h"
 #include "libavutil/opt.h"
+#include "libavutil/thread.h"
+#include "libavutil/time.h"
 
 #include "avformat.h"
+#include "http.h"
 #include "internal.h"
 #include "url.h"
 
-typedef struct CurlContext {
-    const AVClass *class;
-} CurlContext;
+#define CURL_DEFAULT_BUFFER_SIZE (4 << 20)
+/* Blocking waits wake up this often so url_read()/open can poll the interrupt
+ * callback. */
+#define CURL_WAIT_US 100000
+
+typedef struct CurlContext CurlContext;
+
+enum cmd_kind {
+    CMD_ADD,     /* add the easy handle to the multi and start the transfer */
+    CMD_REMOVE,  /* remove the easy handle from the multi */
+    CMD_UNPAUSE, /* resume a transfer paused because the FIFO was full */
+};
+
+typedef struct CurlCmd {
+    enum cmd_kind   kind;
+    CurlContext    *ctx;
+    int             sync;   /* caller waits for completion, flips by done */
+    int             done;
+    struct CurlCmd *next;
+} CurlCmd;
+
+typedef struct CurlLoop {
+    pthread_t       thread;
+    CURLM          *multi;
+
+    pthread_mutex_t mutex;   /* guards the command queue, exit and cmd->done */
+    pthread_cond_t  cond;    /* signaled when a sync command completes */
+    CurlCmd        *cmd_head, *cmd_tail;
+    int             exit;
+} CurlLoop;
+
+struct CurlContext {
+    const AVClass  *class;
+    URLContext     *h;
+
+    CurlLoop       *loop;
+    int             private_loop;  /* loop is owned by this context (not 
shared) */
+    CURL           *easy;
+
+    int64_t         buffer_size;
+
+    /* Producer bookkeeping, touched only by the loop thread. */
+    int             active;     /* currently added to the multi */
+
+    /* Probe result. Set by the loop thread, read by url_open() once probed. */
+    int             probed;
+    int             stream_ok;
+    int             seekable;
+    int64_t         content_size;
+
+    /* Shared transfer state, guarded by mutex. */
+    pthread_mutex_t mutex;
+    pthread_cond_t  cond;
+    AVFifo         *fifo;
+    int             paused;      /* write callback paused, FIFO was full */
+    int             eof;         /* producer delivered all data */
+    int             error;       /* AVERROR for an unrecoverable failure, or 0 
*/
+    int             aborted;     /* transfer should stop (open was 
interrupted) */
+};
+
+/* Guards lazy creation of a format context's shared loop. */
+static AVMutex curl_loop_lock = AV_MUTEX_INITIALIZER;
+
+static int curlcode_to_averror(CURLcode code)
+{
+    switch (code) {
+    case CURLE_OK:                       return 0;
+    case CURLE_URL_MALFORMAT:
+    case CURLE_UNSUPPORTED_PROTOCOL:     return AVERROR(EINVAL);
+    case CURLE_COULDNT_RESOLVE_PROXY:
+    case CURLE_COULDNT_RESOLVE_HOST:     return AVERROR(EHOSTUNREACH);
+    case CURLE_COULDNT_CONNECT:          return AVERROR(ECONNREFUSED);
+    case CURLE_OPERATION_TIMEDOUT:       return AVERROR(ETIMEDOUT);
+    case CURLE_LOGIN_DENIED:
+    case CURLE_REMOTE_ACCESS_DENIED:     return AVERROR(EACCES);
+    case CURLE_OUT_OF_MEMORY:            return AVERROR(ENOMEM);
+    case CURLE_PEER_FAILED_VERIFICATION:
+    case CURLE_SSL_CACERT_BADFILE:       return AVERROR_INVALIDDATA;
+    default:                             return AVERROR(EIO);
+    }
+}
+
+/* ------------------------------------------------------------------------- */
+/* curl callbacks (run on the loop thread)                                   */
+/* ------------------------------------------------------------------------- */
+
+static size_t write_callback(char *ptr, size_t size, size_t nmemb, void 
*userdata)
+{
+    CurlContext *c = userdata;
+    size_t bytes = size * nmemb;
+    size_t space;
+
+    pthread_mutex_lock(&c->mutex);
+
+    if (c->aborted || !c->stream_ok) {
+        pthread_mutex_unlock(&c->mutex);
+        return CURL_WRITEFUNC_ERROR;
+    }
+
+    space = av_fifo_can_write(c->fifo);
+    if (space < bytes) {
+        /* pause the transfer and wait for the consumer to drain. */
+        c->paused = 1;
+        pthread_mutex_unlock(&c->mutex);
+        return CURL_WRITEFUNC_PAUSE;
+    }
+
+    av_fifo_write(c->fifo, ptr, bytes);
+    c->paused = 0;
+    pthread_cond_broadcast(&c->cond);
+    pthread_mutex_unlock(&c->mutex);
+
+    return bytes;
+}
+
+static size_t header_callback(char *ptr, size_t size, size_t nitems, void 
*userdata)
+{
+    CurlContext *c = userdata;
+    size_t len = size * nitems;
+    size_t n = len;
+    long status = 0;
+
+    /* Act only on the blank line that terminates a header block. */
+    while (n && (ptr[n - 1] == '\r' || ptr[n - 1] == '\n'))
+        n--;
+    if (n)
+        return len;
+
+    curl_easy_getinfo(c->easy, CURLINFO_RESPONSE_CODE, &status);
+
+    /* Interim (1xx) and redirect (3xx) responses produce an intermediate 
header
+     * block, wait for the final one. */
+    if (status < 200 || (status >= 300 && status < 400))
+        return len;
+
+    pthread_mutex_lock(&c->mutex);
+    if (status >= 200 && status < 300) {
+        c->stream_ok = 1;
+    } else {
+        c->stream_ok = 0;
+        if (!c->error)
+            c->error = ff_http_averror(status, AVERROR(EIO));
+    }
+    c->probed = 1;
+    pthread_cond_broadcast(&c->cond);
+    pthread_mutex_unlock(&c->mutex);
+
+    return len;
+}
+
+static int xferinfo_callback(void *userdata, curl_off_t dltotal, curl_off_t 
dlnow,
+                             curl_off_t ultotal, curl_off_t ulnow)
+{
+    CurlContext *c = userdata;
+    int aborted;
+    pthread_mutex_lock(&c->mutex);
+    aborted = c->aborted;
+    pthread_mutex_unlock(&c->mutex);
+    return aborted; /* non-zero aborts the transfer */
+}
+
+/* Transfer finished (or failed) */
+static void on_done(CurlContext *c, CURLcode code)
+{
+    pthread_mutex_lock(&c->mutex);
+    if (!c->probed) {
+        /* Connection died before any usable header arrived. */
+        c->probed = 1;
+        c->stream_ok = 0;
+        if (!c->error)
+            c->error = curlcode_to_averror(code);
+    } else if (code == CURLE_OK && !c->aborted) {
+        c->eof = 1;
+    } else if (!c->aborted && !c->error) {
+        c->error = curlcode_to_averror(code);
+    }
+    pthread_cond_broadcast(&c->cond);
+    pthread_mutex_unlock(&c->mutex);
+}
+
+/* ------------------------------------------------------------------------- */
+/* event loop thread + command queue                                         */
+/* ------------------------------------------------------------------------- */
+
+static void execute_command(CurlLoop *loop, CurlCmd *cmd)
+{
+    CurlContext *c = cmd->ctx;
+
+    switch (cmd->kind) {
+    case CMD_ADD:
+        c->active = 1;
+        curl_multi_add_handle(loop->multi, c->easy);
+        break;
+    case CMD_REMOVE:
+        if (c->active) {
+            curl_multi_remove_handle(loop->multi, c->easy);
+            c->active = 0;
+        }
+        break;
+    case CMD_UNPAUSE:
+        curl_easy_pause(c->easy, CURLPAUSE_CONT);
+        break;
+    }
+}
+
+static void *curl_worker(void *arg)
+{
+    CurlLoop *loop = arg;
+
+    ff_thread_setname("curl");
+
+    while (1) {
+        CurlCmd *cmd;
+        CURLMsg *msg;
+        int running = 0, left = 0, do_exit;
+
+        pthread_mutex_lock(&loop->mutex);
+        cmd = loop->cmd_head;
+        if (cmd) {
+            loop->cmd_head = cmd->next;
+            if (!loop->cmd_head)
+                loop->cmd_tail = NULL;
+        }
+        do_exit = loop->exit;
+        pthread_mutex_unlock(&loop->mutex);
+
+        if (cmd) {
+            execute_command(loop, cmd);
+            if (cmd->sync) {
+                pthread_mutex_lock(&loop->mutex);
+                cmd->done = 1;
+                pthread_cond_broadcast(&loop->cond);
+                pthread_mutex_unlock(&loop->mutex);
+            } else {
+                av_free(cmd);
+            }
+            continue; /* drain the whole queue before pumping curl */
+        }
+
+        if (do_exit)
+            break;
+
+        curl_multi_perform(loop->multi, &running);
+
+        while ((msg = curl_multi_info_read(loop->multi, &left))) {
+            CurlContext *c = NULL;
+            if (msg->msg != CURLMSG_DONE)
+                continue;
+            curl_easy_getinfo(msg->easy_handle, CURLINFO_PRIVATE, &c);
+            curl_multi_remove_handle(loop->multi, msg->easy_handle);
+            if (c) {
+                c->active = 0;
+                on_done(c, msg->data.result);
+            }
+        }
+
+        curl_multi_poll(loop->multi, NULL, 0, 1000, NULL);
+    }
+
+    return NULL;
+}
+
+/* Dispatch a command to the loop. For sync commands the caller blocks until 
the
+ * loop thread has executed it. Returns 0 or a negative AVERROR. */
+static int curl_dispatch(CurlLoop *loop, enum cmd_kind kind, CurlContext *c, 
int sync)
+{
+    CurlCmd stackcmd = {0};
+    CurlCmd *cmd = sync ? &stackcmd : av_mallocz(sizeof(*cmd));
+
+    if (!cmd)
+        return AVERROR(ENOMEM);
+
+    cmd->kind = kind;
+    cmd->ctx  = c;
+    cmd->sync = sync;
+
+    pthread_mutex_lock(&loop->mutex);
+    if (loop->cmd_tail)
+        loop->cmd_tail->next = cmd;
+    else
+        loop->cmd_head = cmd;
+    loop->cmd_tail = cmd;
+    curl_multi_wakeup(loop->multi);
+    if (sync) {
+        while (!cmd->done)
+            pthread_cond_wait(&loop->cond, &loop->mutex);
+    }
+    pthread_mutex_unlock(&loop->mutex);
+
+    return 0;
+}
+
+static CurlLoop *curl_loop_create(void)
+{
+    CurlLoop *loop = av_mallocz(sizeof(*loop));
+    if (!loop)
+        return NULL;
+
+    if (pthread_mutex_init(&loop->mutex, NULL))
+        goto fail;
+    if (pthread_cond_init(&loop->cond, NULL)) {
+        pthread_mutex_destroy(&loop->mutex);
+        goto fail;
+    }
+
+    if (curl_global_init(CURL_GLOBAL_DEFAULT) != CURLE_OK)
+        goto fail2;
+
+    loop->multi = curl_multi_init();
+    if (!loop->multi)
+        goto fail3;
+    curl_multi_setopt(loop->multi, CURLMOPT_PIPELINING, CURLPIPE_MULTIPLEX);
+
+    if (pthread_create(&loop->thread, NULL, curl_worker, loop)) {
+        curl_multi_cleanup(loop->multi);
+        goto fail3;
+    }
+
+    return loop;
+
+fail3:
+    curl_global_cleanup();
+fail2:
+    pthread_cond_destroy(&loop->cond);
+    pthread_mutex_destroy(&loop->mutex);
+fail:
+    av_free(loop);
+    return NULL;
+}
+
+static void curl_loop_destroy(CurlLoop *loop)
+{
+    pthread_mutex_lock(&loop->mutex);
+    loop->exit = 1;
+    curl_multi_wakeup(loop->multi);
+    pthread_mutex_unlock(&loop->mutex);
+
+    pthread_join(loop->thread, NULL);
+
+    curl_multi_cleanup(loop->multi);
+    pthread_cond_destroy(&loop->cond);
+    pthread_mutex_destroy(&loop->mutex);
+    av_free(loop);
+
+    /* Released after the thread is joined and the multi handle is gone. */
+    curl_global_cleanup();
+}
+
+/* Attach a context to its event loop. With an owning AVFormatContext the loop 
is
+ * created lazily, cached on it, and shared across the demuxer's transfers so 
curl
+ * reuses connections; it is freed at format teardown. Without one the context
+ * gets a private loop freed on close. */
+static int curl_loop_attach(CurlContext *c, AVFormatContext *avfc)
+{
+    if (!avfc) {
+        c->loop = curl_loop_create();
+        c->private_loop = 1;
+        return c->loop ? 0 : AVERROR(ENOMEM);
+    }
+
+    pthread_mutex_lock(&curl_loop_lock);
+    c->loop = ffformatcontext(avfc)->curl_loop;
+    if (!c->loop) {
+        c->loop = curl_loop_create();
+        ffformatcontext(avfc)->curl_loop = c->loop;
+    }
+    pthread_mutex_unlock(&curl_loop_lock);
+
+    return c->loop ? 0 : AVERROR(ENOMEM);
+}
+
+void ff_curl_loop_free(struct CurlLoop **loop)
+{
+    if (loop && *loop) {
+        curl_loop_destroy(*loop);
+        *loop = NULL;
+    }
+}
+
+/* ------------------------------------------------------------------------- */
+/* URLProtocol callbacks                                                     */
+/* ------------------------------------------------------------------------- */
+
+static int libcurl_close(URLContext *h);
+
+static void setup_curl(CurlContext *c)
+{
+    CURL *e = c->easy;
+    const char *url = c->h->filename;
+
+    /* Drop an optional "libcurl:" prefix that forces this protocol. */
+    av_strstart(url, "libcurl:", &url);
+
+    curl_easy_setopt(e, CURLOPT_URL, url);
+    curl_easy_setopt(e, CURLOPT_PRIVATE, c);
+    curl_easy_setopt(e, CURLOPT_NOSIGNAL, 1L);
+
+    curl_easy_setopt(e, CURLOPT_WRITEFUNCTION, write_callback);
+    curl_easy_setopt(e, CURLOPT_WRITEDATA, c);
+    curl_easy_setopt(e, CURLOPT_HEADERFUNCTION, header_callback);
+    curl_easy_setopt(e, CURLOPT_HEADERDATA, c);
+
+    curl_easy_setopt(e, CURLOPT_NOPROGRESS, 0L);
+    curl_easy_setopt(e, CURLOPT_XFERINFOFUNCTION, xferinfo_callback);
+    curl_easy_setopt(e, CURLOPT_XFERINFODATA, c);
+
+    curl_easy_setopt(e, CURLOPT_FOLLOWLOCATION, 1L);
+    curl_easy_setopt(e, CURLOPT_TCP_KEEPALIVE, 1L);
+    curl_easy_setopt(e, CURLOPT_ACCEPT_ENCODING, "");
+}
+
+static void curl_cond_wait(CurlContext *c)
+{
+    int64_t t = av_gettime() + CURL_WAIT_US;
+    struct timespec ts = { .tv_sec  = t / 1000000,
+                           .tv_nsec = (t % 1000000) * 1000 };
+    pthread_cond_timedwait(&c->cond, &c->mutex, &ts);
+}
+
+/* Block until the transfer has been probed, the stream errored, or the open 
was
+ * interrupted. Returns 0, or a negative AVERROR. */
+static int wait_for_probe(CurlContext *c)
+{
+    URLContext *h = c->h;
+    int ret = 0;
+
+    pthread_mutex_lock(&c->mutex);
+    while (!c->probed && !c->error) {
+        if (ff_check_interrupt(&h->interrupt_callback)) {
+            c->aborted = 1;
+            ret = AVERROR_EXIT;
+            break;
+        }
+        curl_cond_wait(c);
+    }
+    if (!ret) {
+        if (!c->stream_ok)
+            ret = c->error ? c->error : AVERROR(EIO);
+    }
+    pthread_mutex_unlock(&c->mutex);
+
+    return ret;
+}
 
 static int libcurl_open(URLContext *h, const char *url, int flags,
                         AVDictionary **options)
 {
-    return AVERROR(ENOSYS);
+    /* Guard against non-thread-safe libcurl builds. This should never happen,
+     * since libcurl is used only on platforms with thread support, and thread
+     * safety is enabled unconditionally in libcurl when the platform supports
+     * threads or atomics. */
+    curl_version_info_data *info = curl_version_info(CURLVERSION_NOW);
+    if (!(info->features & CURL_VERSION_THREADSAFE))
+        return AVERROR(ENOSYS);
+
+    CurlContext *c = h->priv_data;
+    int ret;
+
+    c->h = h;
+    c->content_size = -1;
+    if (c->buffer_size <= 0)
+        c->buffer_size = CURL_DEFAULT_BUFFER_SIZE;
+
+    ret = pthread_mutex_init(&c->mutex, NULL);
+    if (ret)
+        return AVERROR(ret);
+    ret = pthread_cond_init(&c->cond, NULL);
+    if (ret) {
+        pthread_mutex_destroy(&c->mutex);
+        return AVERROR(ret);
+    }
+
+    c->fifo = av_fifo_alloc2(c->buffer_size, 1, 0);
+    if (!c->fifo) {
+        ret = AVERROR(ENOMEM);
+        goto fail;
+    }
+
+    ret = curl_loop_attach(c, h->avfc);
+    if (ret < 0)
+        goto fail;
+
+    c->easy = curl_easy_init();
+    if (!c->easy) {
+        ret = AVERROR(ENOMEM);
+        goto fail;
+    }
+    setup_curl(c);
+
+    ret = curl_dispatch(c->loop, CMD_ADD, c, 0);
+    if (ret < 0)
+        goto fail;
+
+    ret = wait_for_probe(c);
+    if (ret < 0)
+        goto fail;
+
+    h->is_streamed = !c->seekable;
+
+    return 0;
+
+fail:
+    libcurl_close(h);
+    return ret;
 }
 
 static int libcurl_read(URLContext *h, unsigned char *buf, int size)
 {
-    return AVERROR(ENOSYS);
+    CurlContext *c = h->priv_data;
+    int nonblock = h->flags & AVIO_FLAG_NONBLOCK;
+    int ret;
+
+    pthread_mutex_lock(&c->mutex);
+    while (1) {
+        size_t avail = av_fifo_can_read(c->fifo);
+
+        if (avail) {
+            int n = FFMIN(avail, (size_t)size);
+            int unpause;
+            av_fifo_read(c->fifo, buf, n);
+            /* Resume a paused transfer once the FIFO is at least half empty. 
*/
+            unpause = c->paused && av_fifo_can_write(c->fifo) * 2 >= 
c->buffer_size;
+            pthread_mutex_unlock(&c->mutex);
+            if (unpause)
+                curl_dispatch(c->loop, CMD_UNPAUSE, c, 0);
+            return n;
+        }
+        if (c->error) {
+            ret = c->error;
+            break;
+        }
+        if (c->eof) {
+            ret = AVERROR_EOF;
+            break;
+        }
+        if (nonblock) {
+            ret = AVERROR(EAGAIN);
+            break;
+        }
+        curl_cond_wait(c);
+        /* Return to the avio layer so it can poll the interrupt callback. */
+        nonblock = 1;
+    }
+    pthread_mutex_unlock(&c->mutex);
+
+    return ret;
 }
 
 static int libcurl_close(URLContext *h)
 {
+    CurlContext *c = h->priv_data;
+
+    if (c->loop) {
+        if (c->easy) {
+            /* Ensure the handle is out of the multi before we free it. */
+            curl_dispatch(c->loop, CMD_REMOVE, c, 1);
+            curl_easy_cleanup(c->easy);
+            c->easy = NULL;
+        }
+        /* A shared loop outlives the transfer for connection reuse. */
+        if (c->private_loop)
+            curl_loop_destroy(c->loop);
+        c->loop = NULL;
+    }
+
+    av_fifo_freep2(&c->fifo);
+    pthread_cond_destroy(&c->cond);
+    pthread_mutex_destroy(&c->mutex);
+
     return 0;
 }
 

_______________________________________________
ffmpeg-cvslog mailing list -- [email protected]
To unsubscribe send an email to [email protected]

Reply via email to