mirror of
https://github.com/elisspace/core.git
synced 2026-09-30 06:19:55 +00:00
Remove incomplete segment on stream restart (#59532)
This commit is contained in:
committed by
Paulus Schoutsen
parent
04e1dc3a10
commit
0f0ca36aa8
@@ -22,8 +22,8 @@ from aiohttp import web
|
||||
import async_timeout
|
||||
import pytest
|
||||
|
||||
from homeassistant.components.stream import Stream
|
||||
from homeassistant.components.stream.core import Segment, StreamOutput
|
||||
from homeassistant.components.stream.worker import SegmentBuffer
|
||||
|
||||
TEST_TIMEOUT = 7.0 # Lower than 9s home assistant timeout
|
||||
|
||||
@@ -34,7 +34,7 @@ class WorkerSync:
|
||||
def __init__(self):
|
||||
"""Initialize WorkerSync."""
|
||||
self._event = None
|
||||
self._original = Stream._worker_finished
|
||||
self._original = SegmentBuffer.discontinuity
|
||||
|
||||
def pause(self):
|
||||
"""Pause the worker before it finalizes the stream."""
|
||||
@@ -45,7 +45,7 @@ class WorkerSync:
|
||||
logging.debug("waking blocked worker")
|
||||
self._event.set()
|
||||
|
||||
def blocking_finish(self, stream: Stream):
|
||||
def blocking_discontinuity(self, stream: SegmentBuffer):
|
||||
"""Intercept call to pause stream worker."""
|
||||
# Worker is ending the stream, which clears all output buffers.
|
||||
# Block the worker thread until the test has a chance to verify
|
||||
@@ -63,8 +63,8 @@ def stream_worker_sync(hass):
|
||||
"""Patch StreamOutput to allow test to synchronize worker stream end."""
|
||||
sync = WorkerSync()
|
||||
with patch(
|
||||
"homeassistant.components.stream.Stream._worker_finished",
|
||||
side_effect=sync.blocking_finish,
|
||||
"homeassistant.components.stream.worker.SegmentBuffer.discontinuity",
|
||||
side_effect=sync.blocking_discontinuity,
|
||||
autospec=True,
|
||||
):
|
||||
yield sync
|
||||
|
||||
@@ -448,3 +448,33 @@ async def test_hls_max_segments_discontinuity(hass, hls_stream, stream_worker_sy
|
||||
|
||||
stream_worker_sync.resume()
|
||||
stream.stop()
|
||||
|
||||
|
||||
async def test_remove_incomplete_segment_on_exit(hass, stream_worker_sync):
|
||||
"""Test that the incomplete segment gets removed when the worker thread quits."""
|
||||
await async_setup_component(hass, "stream", {"stream": {}})
|
||||
|
||||
stream = create_stream(hass, STREAM_SOURCE, {})
|
||||
stream_worker_sync.pause()
|
||||
stream.start()
|
||||
hls = stream.add_provider(HLS_PROVIDER)
|
||||
|
||||
segment = Segment(sequence=0, stream_id=0, duration=SEGMENT_DURATION)
|
||||
hls.put(segment)
|
||||
segment = Segment(sequence=1, stream_id=0, duration=SEGMENT_DURATION)
|
||||
hls.put(segment)
|
||||
segment = Segment(sequence=2, stream_id=0, duration=0)
|
||||
hls.put(segment)
|
||||
await hass.async_block_till_done()
|
||||
|
||||
segments = hls._segments
|
||||
assert len(segments) == 3
|
||||
assert not segments[-1].complete
|
||||
stream_worker_sync.resume()
|
||||
stream._thread_quit.set()
|
||||
stream._thread.join()
|
||||
stream._thread = None
|
||||
await hass.async_block_till_done()
|
||||
assert segments[-1].complete
|
||||
assert len(segments) == 2
|
||||
stream.stop()
|
||||
|
||||
Reference in New Issue
Block a user