mirror of
https://github.com/LorenEteval/Furious.git
synced 2026-10-02 03:48:03 +03:00
Clean up partial external core startup
Signed-off-by: Loren Eteval <loren.eteval@proton.me>
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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}')
|
||||
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user