From d728c4ab7d62d698fdf0ab4a0428df964bb1b067 Mon Sep 17 00:00:00 2001 From: reinierlakhan Date: Mon, 7 Sep 2026 01:54:35 -0400 Subject: [PATCH] python-codecs: arm zygote fork self-destruct at call time, not first frame (#2139) CodecFork.generateVideoFramesLibav and generateVideoFramesGstreamer are async generator functions, so their bodies -- including the call_later(10, timeoutExit) that lets the worker kill itself -- do not run until the caller requests the first frame. The PluginFork handle is discarded at handout in LibavGenerator and GstreamerGenerator, so a worker that is handed a session but never iterated has no watchdog and no reachable terminate(), and idles for the life of the plugin host. aclose() on a never-started async generator returns without executing the body, so even a consumer that closes the RPC proxy cannot trigger the self-destruct. Make the two methods coroutines that return the inner generator. rpc.py does await maybe_await(invoke(*args)) before serializing, so the async iterator proxy seen by the caller is unchanged, but the watchdog is armed as soon as the fork is committed to a session. Also keep the PluginFork handle at the call sites and terminate it if the session cannot be handed to the caller. Fixes #2138 Claude-Session: https://claude.ai/code/session_018kjDzffi8VZKA1mn6nrGfo Co-authored-by: Claude Opus 5 (1M context) --- plugins/python-codecs/src/main.py | 58 ++++++++++++++++++++++--------- 1 file changed, 41 insertions(+), 17 deletions(-) diff --git a/plugins/python-codecs/src/main.py b/plugins/python-codecs/src/main.py index c411fa3bd..8aac0574c 100644 --- a/plugins/python-codecs/src/main.py +++ b/plugins/python-codecs/src/main.py @@ -26,6 +26,16 @@ except: pass +def terminate_zygote(fork: scrypted_sdk.PluginFork) -> None: + # a fork popped from the zygote pool is only reachable through this handle. + # kill it if the decode session could not be handed to the caller, otherwise + # the worker idles forever with no one able to stop it. + try: + fork.terminate() + except: + pass + + class LibavGenerator(scrypted_sdk.ScryptedDeviceBase, scrypted_sdk.VideoFrameGenerator): def __init__(self, nativeId: Union[str, None], z): super().__init__(nativeId) @@ -38,8 +48,13 @@ class LibavGenerator(scrypted_sdk.ScryptedDeviceBase, scrypted_sdk.VideoFrameGen # todo remove filter: Any = None, ) -> scrypted_sdk.VideoFrame: - forked: CodecFork = await self.zygote().result - return await forked.generateVideoFramesLibav(mediaObject, options) + fork = self.zygote() + try: + forked: CodecFork = await fork.result + return await forked.generateVideoFramesLibav(mediaObject, options) + except: + terminate_zygote(fork) + raise class GstreamerGenerator( @@ -59,15 +74,20 @@ class GstreamerGenerator( filter: Any = None, ) -> scrypted_sdk.VideoFrame: start = time.time() - forked: CodecFork = await self.zygote().result - print("fork", time.time() - start) - return await forked.generateVideoFramesGstreamer( - mediaObject, - options, - self.storage.getItem("h264Decoder"), - self.storage.getItem("h265Decoder"), - self.storage.getItem("postProcessPipeline"), - ) + fork = self.zygote() + try: + forked: CodecFork = await fork.result + print("fork", time.time() - start) + return await forked.generateVideoFramesGstreamer( + mediaObject, + options, + self.storage.getItem("h264Decoder"), + self.storage.getItem("h265Decoder"), + self.storage.getItem("postProcessPipeline"), + ) + except: + terminate_zygote(fork) + raise async def getSettings(self) -> List[Setting]: return [ @@ -237,32 +257,36 @@ class CodecFork: h265Decoder: str, postProcessPipeline: str, ) -> AsyncGenerator[scrypted_sdk.VideoFrame, Any]: + # armed here rather than inside the generator body: the body does not run + # until the caller requests a frame, so a session that is never iterated + # would otherwise have no watchdog at all. loop = asyncio.get_event_loop() self.timeout = loop.call_later(10, self.timeoutExit) - async for data in self.generateVideoFrames( + return self.generateVideoFrames( gstreamer.generateVideoFramesGstreamer( mediaObject, options, h264Decoder, h265Decoder, postProcessPipeline ), "gstreamer", options and options.get("firstFrameOnly"), - ): - yield data + ) async def generateVideoFramesLibav( self, mediaObject: scrypted_sdk.MediaObject, options: scrypted_sdk.VideoFrameGeneratorOptions = None, ) -> AsyncGenerator[scrypted_sdk.VideoFrame, Any]: + # armed here rather than inside the generator body: the body does not run + # until the caller requests a frame, so a session that is never iterated + # would otherwise have no watchdog at all. loop = asyncio.get_event_loop() self.timeout = loop.call_later(10, self.timeoutExit) - async for data in self.generateVideoFrames( + return self.generateVideoFrames( libav.generateVideoFramesLibav(mediaObject, options), "libav", options and options.get("firstFrameOnly"), - ): - yield data + ) async def fork():