From 29a46b541dfb2c1327463ca0d678e77efed8d33b Mon Sep 17 00:00:00 2001 From: Loren Eteval Date: Wed, 9 Sep 2026 12:17:47 +0800 Subject: [PATCH] Clean up partial external core startup Signed-off-by: Loren Eteval --- Furious/Backends/ExternalCore/AGENTS.md | 4 +- Furious/Backends/ExternalCore/Process.py | 42 +++++++++--- tests/test_external_core.py | 85 +++++++++++++++++++++++- 3 files changed, 119 insertions(+), 12 deletions(-) diff --git a/Furious/Backends/ExternalCore/AGENTS.md b/Furious/Backends/ExternalCore/AGENTS.md index ca4aa79..28e5df1 100644 --- a/Furious/Backends/ExternalCore/AGENTS.md +++ b/Furious/Backends/ExternalCore/AGENTS.md @@ -17,7 +17,9 @@ the intentionally different direct-subprocess scope for user-selected executable - One runtime owns its exact `Popen`, stdout/stderr pipes and readers, watcher, partial-line buffer, exit callback, and reaping path. Shutdown terminates that process, uses only platform-specific escalation for its PID when necessary, - kills as a last resort, joins readers, and remains bounded and idempotent after partial startup. + kills as a last resort, joins readers, and remains bounded and idempotent after partial startup. Register each + successfully started reader immediately; a later reader/watcher startup failure unwinds the child and every pipe, + including pipes with no reader. Disposal is terminal and must reject new process acquisition. - Keep execution liveness, configured proxy endpoints, and semantic readiness distinct. An immediate or later exit is interpreted once at this runtime boundary and retains actionable code/reason context for the shared startup workflow. - Application tun2socks is an explicit profile capability. It requires a usable SOCKS endpoint and a separate remote diff --git a/Furious/Backends/ExternalCore/Process.py b/Furious/Backends/ExternalCore/Process.py index da16bd3..97e5a74 100644 --- a/Furious/Backends/ExternalCore/Process.py +++ b/Furious/Backends/ExternalCore/Process.py @@ -181,8 +181,6 @@ class ExternalCoreProcess(CoreRuntime): def _startReaders(self, process: subprocess.Popen): """Start one bounded-lifetime reader for each captured output pipe.""" - readers = [] - for stream, label in ( (process.stdout, 'stdout'), (process.stderr, 'stderr'), @@ -195,12 +193,11 @@ class ExternalCoreProcess(CoreRuntime): args=(stream, label), daemon=True, ) + thread.start() - readers.append(thread) - - with self._lock: - self._readerThreads = readers + with self._lock: + self._readerThreads.append(thread) def _joinReaders(self, process: Optional[subprocess.Popen] = None): """Finish pipe readers, closing inherited pipes if descendants retain them.""" @@ -217,7 +214,7 @@ class ExternalCoreProcess(CoreRuntime): pending = tuple(thread for thread in readers if thread.is_alive()) - if pending and process is not None: + if process is not None: for stream in (process.stdout, process.stderr): if stream is None: continue @@ -262,7 +259,15 @@ class ExternalCoreProcess(CoreRuntime): with self._lock: self._watcherThread = watcher - watcher.start() + try: + watcher.start() + except Exception: + # Any non-exit exceptions + + with self._lock: + self._watcherThread = None + + raise @staticmethod def _creationOptions() -> dict: @@ -285,6 +290,9 @@ class ExternalCoreProcess(CoreRuntime): def start(self): """Validate and launch execution or raise ``RuntimeStartError``.""" + if self.state is RuntimeState.Disposed: + raise RuntimeStartError('Runtime has already been disposed') + config = self._configuration if not isinstance(config, ConfigExternalCore): @@ -366,8 +374,22 @@ class ExternalCoreProcess(CoreRuntime): self._process = process self.setState(RuntimeState.Alive) - self._startReaders(process) - self._startWatcher(process) + try: + self._startReaders(process) + self._startWatcher(process) + except Exception as ex: + # Any non-exit exceptions + + try: + self.stop() + except Exception: + # Any non-exit exceptions + + logger.exception('failed to clean up partial external core startup') + + self.setState(RuntimeState.Failed) + + raise RuntimeStartError('Failed to start core', details=str(ex)) from ex logger.info(f'external core process started with PID {process.pid}') diff --git a/tests/test_external_core.py b/tests/test_external_core.py index f55bdf2..28161b2 100644 --- a/tests/test_external_core.py +++ b/tests/test_external_core.py @@ -24,7 +24,7 @@ from Furious.Backends.ExternalCore.Plugin import ( ExternalCorePlugin, ExternalCoreRuntimeFactory, ) -from Furious.Interface import CoreRuntime, RuntimeExit, RuntimeStartError +from Furious.Interface import CoreRuntime, RuntimeExit, RuntimeStartError, RuntimeState from Furious.Plugins.API import ( CoreRuntimeRequest, SubscriptionItem, @@ -100,6 +100,89 @@ class ExternalCoreProcessTest(unittest.TestCase): } ) + def testDisposedRuntimeCannotAcquireAnotherProcess(self): + """Disposal is terminal even when the stored launch specification is valid.""" + runtime = ExternalCoreProcess( + self.configuration(['-c', 'pass'], str(Path.cwd())) + ) + runtime.dispose() + + try: + with mock.patch( + 'Furious.Backends.ExternalCore.Process.subprocess.Popen' + ) as spawn: + with self.assertRaises(RuntimeStartError): + runtime.start() + + spawn.assert_not_called() + finally: + runtime.dispose() + + def testPartialThreadStartupReapsChildAndClosesEveryPipe(self): + """A failed reader or watcher releases all earlier acquisitions.""" + originalStart = threading.Thread.start + originalPopen = subprocess.Popen + + for failAt in (1, 2, 3): + with self.subTest(failAt=failAt): + threads = [] + children = [] + runtime = ExternalCoreProcess( + self.configuration( + ['-u', '-c', 'import time; time.sleep(60)'], str(Path.cwd()) + ) + ) + + def startThread(thread): + threads.append(thread) + + if len(threads) == failAt: + raise RuntimeError('thread startup fixture failure') + + originalStart(thread) + + def spawn(*args, **kwargs): + child = originalPopen(*args, **kwargs) + children.append(child) + return child + + try: + with mock.patch('threading.Thread.start', startThread), mock.patch( + 'Furious.Backends.ExternalCore.Process.subprocess.Popen', + side_effect=spawn, + ): + with self.assertRaises(RuntimeStartError): + runtime.start() + + child = children[0] + + self.assertIsNotNone(child.poll()) + self.assertTrue(child.stdout.closed) + self.assertTrue(child.stderr.closed) + self.assertTrue(all(not thread.is_alive() for thread in threads)) + + self.assertIsNone(runtime.process) + self.assertIsNone(runtime._watcherThread) + self.assertFalse(runtime._readerThreads) + self.assertIs(runtime.state, RuntimeState.Failed) + finally: + for child in children: + if child.poll() is None: + child.kill() + child.wait(timeout=5) + + for thread in threads: + if thread.ident is not None: + thread.join(5) + + for child in children: + for stream in (child.stdout, child.stderr): + if stream is not None: + stream.close() + + runtime._watcherThread = None + runtime.dispose() + def testWindowsTaskkillHasBoundedWait(self): """Never allow the fully mocked host shutdown command to wait forever.""" process = mock.Mock(pid=1234)