diff --git a/CHANGELOG.rst b/CHANGELOG.rst index e68481872..2bbdd74ed 100644 --- a/CHANGELOG.rst +++ b/CHANGELOG.rst @@ -65,6 +65,8 @@ Fixes: - Fix crashes from indexes that were turned into C pointer arithmetic without being range checked. ``MotionVectors[i]`` only checked the upper bound, so a negative index read off the front of the buffer (``mvs[-1]`` now returns the last vector, as with any sequence); ``VideoFormatComponent`` and ``AudioPlane`` accepted any index at all; and ``BitmapSubtitlePlane`` and ``VideoBlockParams`` were missing their lower bounds. - Frames returned by flushing a codec context directly (``CodecContext.decode()`` with no packet) now carry the stream's ``time_base`` instead of ``None``. - ``VideoFrame.reformat()`` (and so ``to_ndarray(format=...)``, ``to_rgb()``, ``to_image()``) now shares one ``SwsContext`` per thread instead of allocating one per frame. FFmpeg 8's swscale retains megabytes of graph state per context, which showed up as large RSS growth when many frames were alive at once. +- Writing to a network URL no longer blocks every other Python thread, and ``timeout`` now applies to opening an output container. ``avio_open()``, ``avformat_write_header()``, ``av_write_trailer()``, and ``avio_closep()`` held the GIL, so an unreachable RTMP server froze the whole process, and the interrupt callback was only installed for demuxing, so nothing could end the wait. :meth:`.OutputContainer.close` now raises rather than freeing a context another thread is still muxing or closing. By :gh-user:`adrianrfreedman` in (:pr:`2412`). +- ``timeout`` now applies to muxing and closing an output container, not just to opening it. Only opening armed the interrupt callback, so a peer that accepted the connection and then stopped reading left ``av_interleaved_write_frame()`` and ``av_write_trailer()`` blocked forever. Each mux gets the full timeout, and a close shares one across writing the trailer and flushing, so neither can outlast it. By :gh-user:`adrianrfreedman` in (:pr:`2414`). 18.X and Below diff --git a/av/container/core.py b/av/container/core.py index fe756a159..62c1dd74f 100755 --- a/av/container/core.py +++ b/av/container/core.py @@ -292,16 +292,20 @@ def __cinit__( # We need the context before we open the input AND setup Python IO. self.ptr = lib.avformat_alloc_context() - # Setup interrupt callback - if self.open_timeout is not None or self.read_timeout is not None: - self.ptr.interrupt_callback.callback = interrupt_cb - self.ptr.interrupt_callback.opaque = cython.address( - self.interrupt_callback_info - ) - if acodec is not None: self.ptr.audio_codec_id = getattr(AudioCodec, acodec) + # Setup interrupt callback. Muxing needs it as much as demuxing does, + # since writing the header to a network URL can block indefinitely. + if self.open_timeout is not None or self.read_timeout is not None: + # Start disarmed, so nothing between here and the first + # start_timeout() can be interrupted by a zeroed deadline. + self.set_timeout(None) + self.ptr.interrupt_callback.callback = interrupt_cb + self.ptr.interrupt_callback.opaque = cython.address( + self.interrupt_callback_info + ) + self.ptr.flags |= lib.AVFMT_FLAG_GENPTS self.ptr.opaque = cython.cast(cython.p_void, self) @@ -495,7 +499,11 @@ def open( :param int buffer_size: Size of buffer for Python input/output operations in bytes. Honored only when ``file`` is a file-like object. Defaults to 32768 (32k). :param timeout: How many seconds to wait for data before giving up, as a float, or a - ``(open timeout, read timeout)`` tuple. + ``(open timeout, read timeout)`` tuple. The open timeout covers both connecting + and reading or writing the header. The read timeout covers each subsequent + demux, mux, or close, so a stalled peer gives up rather than blocking forever. + Each demux and mux gets the full timeout; a close shares one across writing + the trailer and flushing, so it cannot outlast the timeout either. :param callable io_open: Custom I/O callable for opening files/streams. This option is intended for formats that need to open additional file-like objects to ``file`` using custom I/O. diff --git a/av/container/output.pxd b/av/container/output.pxd index fc98c8828..41150f27b 100644 --- a/av/container/output.pxd +++ b/av/container/output.pxd @@ -7,6 +7,9 @@ from av.stream cimport Stream cdef class OutputContainer(Container): cdef lib.AVPacket *packet_ptr + # How many nogil libav calls are in flight, so close() can refuse + # to free the context while another thread is still inside one. + cdef int _blocking_depth cdef dict _extradata_bsfs cdef list[Packet] _buffered_packets cdef _buffer_for_extradata(self, Packet packet) diff --git a/av/container/output.py b/av/container/output.py index 44f4990eb..b7b733f31 100644 --- a/av/container/output.py +++ b/av/container/output.py @@ -49,6 +49,7 @@ def close_output(self: OutputContainer) -> cython.void: self._mux_one(packet) self.streams = StreamContainer() + self._blocking_depth += 1 try: if self._myflag & 12 == 4: # enum.started and not enum.done # If the underlying Python IO file was already closed (e.g. during @@ -60,13 +61,24 @@ def close_output(self: OutputContainer) -> cython.void: # We must only ever call av_write_trailer *once*, otherwise we get a # segmentation fault. Therefore no matter whether it succeeds or not # we must absolutely set enum.done. + ret: cython.int + self.set_timeout(self.read_timeout) try: - self.err_check(lib.av_write_trailer(self.ptr)) + self.start_timeout() + with cython.nogil: + ret = lib.av_write_trailer(self.ptr) + self.err_check(ret) finally: if self.file is None and not ( self.ptr.oformat.flags & lib.AVFMT_NOFILE ): - lib.avio_closep(cython.address(self.ptr.pb)) + # No fresh deadline: the trailer and this flush share one, + # so closing cannot outlast the timeout. The point here is + # to stop the flush hanging, not to report on it, so its + # return goes unchecked as it always has. + with cython.nogil: + lib.avio_closep(cython.address(self.ptr.pb)) + self.set_timeout(None) self._myflag |= 8 # enum.done = True finally: # Drop the context so a closed output reports itself as closed: @@ -76,6 +88,7 @@ def close_output(self: OutputContainer) -> cython.void: with cython.nogil: lib.avformat_free_context(self.ptr) self.ptr = cython.NULL + self._blocking_depth -= 1 @cython.final @@ -591,17 +604,54 @@ def start_encoding(self): # Open the output file, if needed. name_obj: bytes = os.fsencode(self.name if self.file is None else "") name: cython.p_char = name_obj - if self.ptr.pb == cython.NULL and not self.ptr.oformat.flags & lib.AVFMT_NOFILE: - err_check( - lib.avio_open(cython.address(self.ptr.pb), name, lib.AVIO_FLAG_WRITE) - ) + ret: cython.int + opened_pb: cython.bint = False + all_options: Dictionary + options: Dictionary + options_ptr: cython.pointer[cython.pointer[lib.AVDictionary]] + + self.set_timeout(self.open_timeout) + self.start_timeout() + self._blocking_depth += 1 + try: + if ( + self.ptr.pb == cython.NULL + and not self.ptr.oformat.flags & lib.AVFMT_NOFILE + ): + # avio_open() would pass the protocol a NULL interrupt + # callback, so a stalled connect could never be timed out. + with cython.nogil: + ret = lib.avio_open2( + cython.address(self.ptr.pb), + name, + lib.AVIO_FLAG_WRITE, + cython.address(self.ptr.interrupt_callback), + cython.NULL, + ) + err_check(ret) + opened_pb = True - # Copy the metadata dict. - dict_to_avdict(cython.address(self.ptr.metadata), self.metadata) + # Copy the metadata dict. + dict_to_avdict(cython.address(self.ptr.metadata), self.metadata) - all_options: Dictionary = Dictionary(self.options, self.container_options) - options: Dictionary = all_options.copy() - self.err_check(lib.avformat_write_header(self.ptr, cython.address(options.ptr))) + all_options = Dictionary(self.options, self.container_options) + options = all_options.copy() + options_ptr = cython.address(options.ptr) + with cython.nogil: + ret = lib.avformat_write_header(self.ptr, options_ptr) + try: + self.err_check(ret) + except Exception: + # started is never set, so close_output() will not close pb. + # Nothing else will either, and a stalled header write is an + # expected path now that it can time out. + if opened_pb: + with cython.nogil: + lib.avio_closep(cython.address(self.ptr.pb)) + raise + finally: + self._blocking_depth -= 1 + self.set_timeout(None) # Track option usage... for k in all_options: @@ -668,6 +718,13 @@ def default_subtitle_codec(self): return lib.avcodec_get_name(self.format.optr.subtitle_codec) def close(self): + if self._blocking_depth: + # Another thread is inside libav without the GIL, so freeing the + # context here would be a use-after-free. Pass ``timeout`` to + # :func:`av.open` to give up on an open that never connects. + raise RuntimeError( + "Cannot close an OutputContainer while another thread is writing to it" + ) close_output(self) def mux(self, packets): @@ -704,8 +761,16 @@ def _mux_one(self, packet: Packet) -> cython.void: # takes ownership of the reference. self.err_check(lib.av_packet_ref(self.packet_ptr, packet.ptr)) - with cython.nogil: - ret: cython.int = lib.av_interleaved_write_frame(self.ptr, self.packet_ptr) + ret: cython.int + self.set_timeout(self.read_timeout) + self.start_timeout() + self._blocking_depth += 1 + try: + with cython.nogil: + ret = lib.av_interleaved_write_frame(self.ptr, self.packet_ptr) + finally: + self._blocking_depth -= 1 + self.set_timeout(None) self.err_check(ret) @cython.cfunc diff --git a/av/video/codeccontext.py b/av/video/codeccontext.py index c01ee95c8..33c4840f4 100644 --- a/av/video/codeccontext.py +++ b/av/video/codeccontext.py @@ -14,6 +14,7 @@ @cython.cfunc +@cython.nogil @cython.exceptval(check=False) def _get_hw_format( ctx: cython.pointer[lib.AVCodecContext], @@ -114,8 +115,15 @@ def _encode_upload_frame(self, vframe: VideoFrame) -> VideoFrame: ) hwframe: VideoFrame = alloc_video_frame() - err_check(lib.av_hwframe_get_buffer(self.ptr.hw_frames_ctx, hwframe.ptr, 0)) - err_check(lib.av_hwframe_transfer_data(hwframe.ptr, vframe.ptr, 0)) + + res: cython.int + transfer_res: cython.int = 0 + with cython.nogil: + res = lib.av_hwframe_get_buffer(self.ptr.hw_frames_ctx, hwframe.ptr, 0) + if res == 0: + transfer_res = lib.av_hwframe_transfer_data(hwframe.ptr, vframe.ptr, 0) + err_check(res) + err_check(transfer_res) hwframe._copy_internal_attributes(vframe, data_layout=False) hwframe._init_user_attributes() @@ -180,7 +188,10 @@ def _transfer_hwframe(self, frame: Frame): return frame frame_sw: Frame = self._alloc_next_frame() - err_check(lib.av_hwframe_transfer_data(frame_sw.ptr, frame.ptr, 0)) + res: cython.int + with cython.nogil: + res = lib.av_hwframe_transfer_data(frame_sw.ptr, frame.ptr, 0) + err_check(res) frame_sw._copy_internal_attributes(frame, data_layout=False) return frame_sw diff --git a/av/video/reformatter.py b/av/video/reformatter.py index 12b78bc23..492437c57 100644 --- a/av/video/reformatter.py +++ b/av/video/reformatter.py @@ -287,9 +287,12 @@ def _reformat( dst_color_primaries: cython.int, threads: cython.int, ): + res: cython.int if frame.ptr.hw_frames_ctx: frame_sw = alloc_video_frame() - err_check(lib.av_hwframe_transfer_data(frame_sw.ptr, frame.ptr, 0)) + with cython.nogil: + res = lib.av_hwframe_transfer_data(frame_sw.ptr, frame.ptr, 0) + err_check(res) frame_sw._copy_internal_attributes(frame, data_layout=False) frame_sw._init_user_attributes() frame = frame_sw diff --git a/include/avformat.pxd b/include/avformat.pxd index 4a4181b4d..318205433 100644 --- a/include/avformat.pxd +++ b/include/avformat.pxd @@ -179,6 +179,10 @@ cdef extern from "libavformat/avformat.h" nogil: cdef int av_interleaved_write_frame(AVFormatContext *ctx, AVPacket *pkt) cdef int av_write_frame(AVFormatContext *ctx, AVPacket *pkt) cdef int avio_open(AVIOContext **s, const char *url, int flags) + cdef int avio_open2( + AVIOContext **s, const char *url, int flags, + const AVIOInterruptCB *int_cb, AVDictionary **options + ) cdef int64_t avio_size(AVIOContext *s) cdef const AVOutputFormat* av_guess_format( const char *short_name, const char *filename, const char *mime_type diff --git a/tests/test_output_blocking.py b/tests/test_output_blocking.py new file mode 100644 index 000000000..ff96dc9e7 --- /dev/null +++ b/tests/test_output_blocking.py @@ -0,0 +1,220 @@ +import socket +import threading +import time + +import numpy as np +import pytest + +import av + +from .common import TestCase + +# The main thread sleeps in 1 ms slices while another thread is stuck in the +# RTMP handshake. It wakes roughly a thousand times if the GIL is free and a +# handful of times if it is not, so this threshold sits well clear of both. +WINDOW = 1.0 +MIN_TICKS = 100 + + +# A writer that never blocks has nothing to time out, so give up once the +# socket has swallowed more than any plausible buffer. +MAX_FRAMES = 500 + + +class SilentServer: + """Accepts connections and then neither reads nor writes. + + An RTMP handshake never completes against it, and a socket written to it + fills up and stays full. + """ + + def __init__(self) -> None: + self.sock = socket.socket() + self.sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) + self.sock.setsockopt(socket.SOL_SOCKET, socket.SO_RCVBUF, 1024) + self.sock.bind(("127.0.0.1", 0)) + self.sock.listen(4) + self.port: int = self.sock.getsockname()[1] + self.accepted: list[socket.socket] = [] + self.thread = threading.Thread(target=self._accept, daemon=True) + self.thread.start() + + def _accept(self) -> None: + while True: + try: + conn, _ = self.sock.accept() + except OSError: + return + self.accepted.append(conn) + + def close(self) -> None: + self.sock.close() + for conn in self.accepted: + conn.close() + + +def has_rtmp() -> bool: + """Whether FFmpeg was built with the RTMP protocol. + + Port 1 on loopback refuses at once, so the probe either fails looking the + protocol up, before any connect, or fails connecting. + """ + try: + with av.open("rtmp://127.0.0.1:1/x", "w", format="flv", timeout=1) as container: + container.start_encoding() + except av.error.ProtocolNotFoundError: + return False + except Exception: + pass + return True + + +@pytest.mark.skipif(not has_rtmp(), reason="FFmpeg was built without RTMP") +class TestOutputBlocking(TestCase): + def setUp(self) -> None: + self.server = SilentServer() + + def tearDown(self) -> None: + self.server.close() + + def _push( + self, timeout: float, containers: list | None = None + ) -> tuple[threading.Thread, list[BaseException]]: + raised: list[BaseException] = [] + + def run() -> None: + try: + container = av.open( + f"rtmp://127.0.0.1:{self.server.port}/live/x", + "w", + format="flv", + timeout=timeout, + ) + if containers is not None: + containers.append(container) + stream = container.add_stream("h264", rate=30) + stream.width = 320 + stream.height = 240 + stream.pix_fmt = "yuv420p" + container.start_encoding() + except BaseException as e: + raised.append(e) + + thread = threading.Thread(target=run, daemon=True) + thread.start() + return thread, raised + + def test_start_encoding_releases_the_gil(self) -> None: + thread, _ = self._push(WINDOW * 3) + + ticks = 0 + deadline = time.monotonic() + WINDOW + while time.monotonic() < deadline: + ticks += 1 + time.sleep(0.001) + + assert thread.is_alive(), "the handshake completed, so nothing was blocking" + assert ticks > MIN_TICKS, f"main thread only ran {ticks} times" + thread.join(WINDOW * 8) + + def test_start_encoding_honours_the_timeout(self) -> None: + thread, raised = self._push(WINDOW) + thread.join(WINDOW * 8) + assert not thread.is_alive(), "timeout did not interrupt the handshake" + assert raised, "the handshake returned instead of timing out" + + def test_close_refuses_to_free_a_container_in_use(self) -> None: + containers: list[av.container.OutputContainer] = [] + thread, _ = self._push(WINDOW * 2, containers) + + # The server only accepts once the writing thread is inside the + # connect, which is where the context stops being ours to free. + deadline = time.monotonic() + WINDOW + while not self.server.accepted and time.monotonic() < deadline: + time.sleep(0.001) + assert self.server.accepted, "the writing thread never connected" + + with pytest.raises(RuntimeError, match="another thread"): + containers[0].close() + thread.join(WINDOW * 8) + + +class TestFailedHeaderWrite(TestCase): + def test_a_failed_header_write_closes_the_connection(self) -> None: + """mp4 cannot carry PCM, so the muxer rejects it after the connect.""" + server = SilentServer() + self.addCleanup(server.close) + + container = av.open(f"tcp://127.0.0.1:{server.port}", "w", format="mp4") + container.add_stream("pcm_s16le") + with pytest.raises(av.error.ArgumentError): + container.start_encoding() + + deadline = time.monotonic() + WINDOW + while not server.accepted and time.monotonic() < deadline: + time.sleep(0.001) + assert server.accepted, "the writer never connected" + + # started was never set, so nothing downstream would close pb. + conn = server.accepted[0] + conn.settimeout(WINDOW) + try: + while conn.recv(4096): + pass # Drain whatever the muxer wrote before it gave up. + except TimeoutError: + raise AssertionError("the connection was left open") from None + + +class TestOutputWriteTimeout(TestCase): + """Muxing and closing over a peer that has stopped reading.""" + + def setUp(self) -> None: + self.server = SilentServer() + + def tearDown(self) -> None: + self.server.close() + + def _open(self) -> av.container.OutputContainer: + container = av.open( + f"tcp://127.0.0.1:{self.server.port}", + "w", + format="mpegts", + timeout=WINDOW, + ) + stream = container.add_stream("mpeg4", rate=30) + stream.width = 640 + stream.height = 480 + stream.pix_fmt = "yuv420p" + return container + + def _fill(self, container: av.container.OutputContainer) -> None: + """Mux noise, which compresses badly, until the socket blocks.""" + stream = container.streams.video[0] + rgb = np.random.randint(0, 256, (480, 640, 3), dtype=np.uint8) + frame = av.VideoFrame.from_ndarray(rgb, format="rgb24") + for i in range(MAX_FRAMES): + frame.pts = i + for packet in stream.encode(frame): + container.mux(packet) + raise AssertionError("the socket swallowed everything without blocking") + + def test_mux_honours_the_timeout(self) -> None: + container = self._open() + start = time.monotonic() + with pytest.raises(av.error.ExitError): + self._fill(container) + assert time.monotonic() - start >= WINDOW, "the write gave up early" + + def test_close_frees_the_container_after_a_stalled_write(self) -> None: + container = self._open() + with pytest.raises(av.error.ExitError): + self._fill(container) + + # A timed-out write leaves its error on the AVIO context, so the + # trailer fails straight away rather than blocking. That makes this a + # test of the teardown, not of the close timeout: close() must report + # the failure and still free the context. + with pytest.raises(av.error.ExitError): + container.close() + with pytest.raises(AssertionError, match="not open"): + container.add_stream("mpeg4", rate=30)