Skip to content

Commit 7a7d4a9

Browse files
committed
Honour timeout when muxing and closing an output container
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. Arm the read timeout around both. Each mux gets the full timeout, as each demux already does. A close shares one deadline across writing the trailer and flushing, so it cannot outlast the timeout it was given. The close side has no test of its own. A timed-out write leaves its error on the AVIO context, so the trailer then fails at once rather than blocking, and I could not get a peer to stall only between the last mux and the close: loopback buffers are far larger than the 32KB AVIO buffer, and AF_UNIX ones are small enough that the header blocks first.
1 parent ac35ccc commit 7a7d4a9

4 files changed

Lines changed: 83 additions & 6 deletions

File tree

‎CHANGELOG.rst‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -65,7 +65,8 @@ Fixes:
6565
- 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.
6666
- Frames returned by flushing a codec context directly (``CodecContext.decode()`` with no packet) now carry the stream's ``time_base`` instead of ``None``.
6767
- ``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.
68-
- 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. Muxing and closing still ignore ``timeout``, and Muxing and closing still ignore ``timeout``, and :meth:`.OutputContainer.close` now raises rather than freeing a context another thread is still muxing or closing. By :gh-user:`adrianrfreedman` in (:pr:`2412`).
68+
- 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`).
69+
- ``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`).
6970

7071

7172
18.X and Below

‎av/container/core.py‎

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -500,10 +500,10 @@ def open(
500500
Honored only when ``file`` is a file-like object. Defaults to 32768 (32k).
501501
:param timeout: How many seconds to wait for data before giving up, as a float, or a
502502
``(open timeout, read timeout)`` tuple. The open timeout covers both connecting
503-
and reading or writing the header. Writing honours it only while opening, so it
504-
is the supported way to give up on an output that never connects. Muxing and
505-
closing still block indefinitely, and calling :meth:`.OutputContainer.close`
506-
from another thread to break out of them raises instead.
503+
and reading or writing the header. The read timeout covers each subsequent
504+
demux, mux, or close, so a stalled peer gives up rather than blocking forever.
505+
Each demux and mux gets the full timeout; a close shares one across writing
506+
the trailer and flushing, so it cannot outlast the timeout either.
507507
:param callable io_open: Custom I/O callable for opening files/streams.
508508
This option is intended for formats that need to open additional
509509
file-like objects to ``file`` using custom I/O.

‎av/container/output.py‎

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -62,16 +62,23 @@ def close_output(self: OutputContainer) -> cython.void:
6262
# segmentation fault. Therefore no matter whether it succeeds or not
6363
# we must absolutely set enum.done.
6464
ret: cython.int
65+
self.set_timeout(self.read_timeout)
6566
try:
67+
self.start_timeout()
6668
with cython.nogil:
6769
ret = lib.av_write_trailer(self.ptr)
6870
self.err_check(ret)
6971
finally:
7072
if self.file is None and not (
7173
self.ptr.oformat.flags & lib.AVFMT_NOFILE
7274
):
75+
# No fresh deadline: the trailer and this flush share one,
76+
# so closing cannot outlast the timeout. The point here is
77+
# to stop the flush hanging, not to report on it, so its
78+
# return goes unchecked as it always has.
7379
with cython.nogil:
7480
lib.avio_closep(cython.address(self.ptr.pb))
81+
self.set_timeout(None)
7582
self._myflag |= 8 # enum.done = True
7683
finally:
7784
# Drop the context so a closed output reports itself as closed:
@@ -755,12 +762,15 @@ def _mux_one(self, packet: Packet) -> cython.void:
755762
self.err_check(lib.av_packet_ref(self.packet_ptr, packet.ptr))
756763

757764
ret: cython.int
765+
self.set_timeout(self.read_timeout)
766+
self.start_timeout()
758767
self._blocking_depth += 1
759768
try:
760769
with cython.nogil:
761770
ret = lib.av_interleaved_write_frame(self.ptr, self.packet_ptr)
762771
finally:
763772
self._blocking_depth -= 1
773+
self.set_timeout(None)
764774
self.err_check(ret)
765775

766776
@cython.cfunc

‎tests/test_output_blocking.py‎

Lines changed: 67 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@
22
import threading
33
import time
44

5+
import numpy as np
56
import pytest
67

78
import av
@@ -15,12 +16,22 @@
1516
MIN_TICKS = 100
1617

1718

19+
# A writer that never blocks has nothing to time out, so give up once the
20+
# socket has swallowed more than any plausible buffer.
21+
MAX_FRAMES = 500
22+
23+
1824
class SilentServer:
19-
"""Accepts connections and then says nothing, so the handshake never ends."""
25+
"""Accepts connections and then neither reads nor writes.
26+
27+
An RTMP handshake never completes against it, and a socket written to it
28+
fills up and stays full.
29+
"""
2030

2131
def __init__(self) -> None:
2232
self.sock = socket.socket()
2333
self.sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
34+
self.sock.setsockopt(socket.SOL_SOCKET, socket.SO_RCVBUF, 1024)
2435
self.sock.bind(("127.0.0.1", 0))
2536
self.sock.listen(4)
2637
self.port: int = self.sock.getsockname()[1]
@@ -152,3 +163,58 @@ def test_a_failed_header_write_closes_the_connection(self) -> None:
152163
pass # Drain whatever the muxer wrote before it gave up.
153164
except TimeoutError:
154165
raise AssertionError("the connection was left open") from None
166+
167+
168+
class TestOutputWriteTimeout(TestCase):
169+
"""Muxing and closing over a peer that has stopped reading."""
170+
171+
def setUp(self) -> None:
172+
self.server = SilentServer()
173+
174+
def tearDown(self) -> None:
175+
self.server.close()
176+
177+
def _open(self) -> av.container.OutputContainer:
178+
container = av.open(
179+
f"tcp://127.0.0.1:{self.server.port}",
180+
"w",
181+
format="mpegts",
182+
timeout=WINDOW,
183+
)
184+
stream = container.add_stream("mpeg4", rate=30)
185+
stream.width = 640
186+
stream.height = 480
187+
stream.pix_fmt = "yuv420p"
188+
return container
189+
190+
def _fill(self, container: av.container.OutputContainer) -> None:
191+
"""Mux noise, which compresses badly, until the socket blocks."""
192+
stream = container.streams.video[0]
193+
rgb = np.random.randint(0, 256, (480, 640, 3), dtype=np.uint8)
194+
frame = av.VideoFrame.from_ndarray(rgb, format="rgb24")
195+
for i in range(MAX_FRAMES):
196+
frame.pts = i
197+
for packet in stream.encode(frame):
198+
container.mux(packet)
199+
raise AssertionError("the socket swallowed everything without blocking")
200+
201+
def test_mux_honours_the_timeout(self) -> None:
202+
container = self._open()
203+
start = time.monotonic()
204+
with pytest.raises(av.error.ExitError):
205+
self._fill(container)
206+
assert time.monotonic() - start >= WINDOW, "the write gave up early"
207+
208+
def test_close_frees_the_container_after_a_stalled_write(self) -> None:
209+
container = self._open()
210+
with pytest.raises(av.error.ExitError):
211+
self._fill(container)
212+
213+
# A timed-out write leaves its error on the AVIO context, so the
214+
# trailer fails straight away rather than blocking. That makes this a
215+
# test of the teardown, not of the close timeout: close() must report
216+
# the failure and still free the context.
217+
with pytest.raises(av.error.ExitError):
218+
container.close()
219+
with pytest.raises(AssertionError, match="not open"):
220+
container.add_stream("mpeg4", rate=30)

0 commit comments

Comments
 (0)