From a2af007893dd6240f36d13f8deeed8427c002243 Mon Sep 17 00:00:00 2001 From: AN Long Date: Sat, 7 Feb 2026 01:21:48 +0900 Subject: [PATCH 01/14] feat: introduce ginterpreter worker --- gunicorn/workers/__init__.py | 1 + gunicorn/workers/ginterpreter.py | 188 +++++++++++++++++++++++++++++++ tests/test_ginterpreter.py | 77 +++++++++++++ 3 files changed, 266 insertions(+) create mode 100644 gunicorn/workers/ginterpreter.py create mode 100644 tests/test_ginterpreter.py diff --git a/gunicorn/workers/__init__.py b/gunicorn/workers/__init__.py index ad77416d1..ce9b55d48 100644 --- a/gunicorn/workers/__init__.py +++ b/gunicorn/workers/__init__.py @@ -12,4 +12,5 @@ "tornado": "gunicorn.workers.gtornado.TornadoWorker", "gthread": "gunicorn.workers.gthread.ThreadWorker", "asgi": "gunicorn.workers.gasgi.ASGIWorker", + "ginterpreter": "gunicorn.workers.ginterpreter.InterpreterWorker", } diff --git a/gunicorn/workers/ginterpreter.py b/gunicorn/workers/ginterpreter.py new file mode 100644 index 000000000..d855478c1 --- /dev/null +++ b/gunicorn/workers/ginterpreter.py @@ -0,0 +1,188 @@ +# +# This file is part of gunicorn released under the MIT license. +# See the NOTICE for more information. + +"""InterpreterPoolExecutor-based worker using Python 3.14+ sub-interpreters.""" + +import errno +import os +import select +import sys + +from . import base + + +def _check_interpreter_pool_available(): + """Check if InterpreterPoolExecutor is available (Python 3.14+).""" + try: + from concurrent.futures import InterpreterPoolExecutor # noqa: F401 + return True + except ImportError: + return False + + +_interpreter_state = { + 'wsgi_app': None, + 'cfg_dict': None, +} + + +def _init_interpreter(cfg_dict, app_uri): + """Initialize the interpreter with WSGI app and config.""" + from gunicorn.util import import_app + + _interpreter_state['cfg_dict'] = cfg_dict + _interpreter_state['wsgi_app'] = import_app(app_uri) + + +def _handle_request_in_interpreter(fd, client_addr, server_addr, family): + """Handle a single HTTP request in a sub-interpreter.""" + import socket + import types + + from gunicorn.http.parser import RequestParser + from gunicorn.http.wsgi import create + + cfg_dict = _interpreter_state['cfg_dict'] + wsgi_app = _interpreter_state['wsgi_app'] + + if cfg_dict is None or wsgi_app is None: + os.close(fd) + return + + sock = socket.fromfd(fd, family, socket.SOCK_STREAM) + try: + os.close(fd) + sock.settimeout(cfg_dict.get('timeout', 30)) + + cfg = types.SimpleNamespace(**cfg_dict) + cfg.forwarded_allow_networks = lambda: [] + cfg.proxy_allow_networks = lambda: [] + + parser = RequestParser(cfg, sock, client_addr) + try: + req = next(parser) + except StopIteration: + return + + if not req: + return + + resp, environ = create(req, sock, client_addr, server_addr, cfg) + environ['wsgi.multithread'] = True + environ['wsgi.multiprocess'] = True + + respiter = wsgi_app(environ, resp.start_response) + try: + for item in respiter: + resp.write(item) + resp.close() + finally: + if hasattr(respiter, 'close'): + respiter.close() + + except socket.timeout: + pass + except OSError as e: + if e.errno not in (errno.EPIPE, errno.ECONNRESET, errno.ENOTCONN): + raise + finally: + try: + sock.close() + except Exception: + pass + + +class InterpreterWorker(base.Worker): + """Worker using InterpreterPoolExecutor for true parallelism.""" + + def init_process(self): + if not _check_interpreter_pool_available(): + raise RuntimeError( + "InterpreterPoolExecutor requires Python 3.14+. " + f"Current version: {sys.version_info.major}.{sys.version_info.minor}" + ) + + from concurrent.futures import InterpreterPoolExecutor + + self.cfg_dict = self._extract_config() + + self.app_uri = getattr(self.app, 'app_uri', None) or self.app.cfg.wsgi_app + if not self.app_uri: + raise RuntimeError( + "ginterpreter worker requires wsgi_app config to be set. " + "Use 'gunicorn myapp:app' or set wsgi_app in config." + ) + + self.executor = InterpreterPoolExecutor( + max_workers=self.cfg.threads, + initializer=_init_interpreter, + initargs=(self.cfg_dict, self.app_uri), + ) + + super().init_process() + + def _extract_config(self): + cfg = self.cfg + return { + 'limit_request_line': cfg.limit_request_line, + 'limit_request_fields': cfg.limit_request_fields, + 'limit_request_field_size': cfg.limit_request_field_size, + 'strip_header_spaces': cfg.strip_header_spaces, + 'permit_obsolete_folding': cfg.permit_obsolete_folding, + 'header_map': cfg.header_map, + 'casefold_http_method': cfg.casefold_http_method, + 'permit_unconventional_http_method': cfg.permit_unconventional_http_method, + 'permit_unconventional_http_version': cfg.permit_unconventional_http_version, + 'forwarded_allow_ips': list(cfg.forwarded_allow_ips), + 'forwarder_headers': list(cfg.forwarder_headers), + 'secure_scheme_headers': dict(cfg.secure_scheme_headers), + 'proxy_protocol': cfg.proxy_protocol, + 'proxy_allow_ips': list(cfg.proxy_allow_ips), + 'is_ssl': cfg.is_ssl, + 'sendfile': cfg.sendfile, + 'workers': cfg.workers, + 'errorlog': cfg.errorlog, + 'timeout': cfg.timeout, + } + + def accept(self, listener): + client, addr = listener.accept() + fd = client.fileno() + family = client.family + server = listener.getsockname() + self.executor.submit( + _handle_request_in_interpreter, + fd, addr, server, family, + ) + client.detach() + + def run(self): + for listener in self.sockets: + listener.setblocking(False) + + while self.alive: + self.notify() + + if self.ppid != os.getppid(): + self.log.info("Parent changed, shutting down: %s", self) + break + + try: + ready = select.select(self.sockets, [], [], 1.0) + for listener in ready[0]: + try: + self.accept(listener) + except OSError as e: + if e.errno not in (errno.EAGAIN, errno.ECONNABORTED, + errno.EWOULDBLOCK): + raise + except OSError as e: + if e.errno != errno.EINTR: + raise + + self.executor.shutdown(wait=True) + + def handle_quit(self, sig, frame): + self.executor.shutdown(wait=False) + super().handle_quit(sig, frame) diff --git a/tests/test_ginterpreter.py b/tests/test_ginterpreter.py new file mode 100644 index 000000000..d634bfb5f --- /dev/null +++ b/tests/test_ginterpreter.py @@ -0,0 +1,77 @@ +# +# This file is part of gunicorn released under the MIT license. +# See the NOTICE for more information. + +"""Tests for the ginterpreter worker.""" + +import os +from unittest import mock + +import pytest + +from gunicorn.config import Config +from gunicorn.workers import ginterpreter + + +def _create_worker(cfg=None): + """Create a worker instance for testing.""" + if cfg is None: + cfg = Config() + cfg.set('workers', 1) + cfg.set('threads', 4) + + mock_app = mock.Mock() + mock_app.app_uri = 'myapp:application' + mock_app.cfg = mock.Mock() + mock_app.cfg.wsgi_app = None + + return ginterpreter.InterpreterWorker( + age=1, + ppid=os.getpid(), + sockets=[], + app=mock_app, + timeout=30, + cfg=cfg, + log=mock.Mock(), + ) + + +class TestInterpreterWorker: + + def test_extract_config(self): + worker = _create_worker() + cfg_dict = worker._extract_config() + assert isinstance(cfg_dict, dict) + assert 'limit_request_line' in cfg_dict + assert 'timeout' in cfg_dict + for value in cfg_dict.values(): + assert isinstance(value, (int, bool, str, list, dict, type(None))) + + def test_handle_quit(self): + worker = _create_worker() + worker.executor = mock.Mock() + with pytest.raises(SystemExit): + worker.handle_quit(None, None) + worker.executor.shutdown.assert_called_once_with(wait=False) + + +class TestAccept: + + def test_submit_to_executor(self): + worker = _create_worker() + worker.executor = mock.Mock() + + mock_client = mock.Mock() + mock_client.fileno.return_value = 7 + mock_client.family = 2 + mock_listener = mock.Mock() + mock_listener.accept.return_value = (mock_client, ('127.0.0.1', 8000)) + mock_listener.getsockname.return_value = ('0.0.0.0', 9000) + + worker.accept(mock_listener) + + worker.executor.submit.assert_called_once_with( + ginterpreter._handle_request_in_interpreter, + 7, ('127.0.0.1', 8000), ('0.0.0.0', 9000), 2, + ) + mock_client.detach.assert_called_once() \ No newline at end of file From 7727acac2ca7c31b0b80c46367246ada054c21fb Mon Sep 17 00:00:00 2001 From: AN Long Date: Sat, 7 Feb 2026 01:55:28 +0900 Subject: [PATCH 02/14] chore: remove the unnessesary dup and close on sock fd --- gunicorn/workers/ginterpreter.py | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/gunicorn/workers/ginterpreter.py b/gunicorn/workers/ginterpreter.py index d855478c1..3fa0f27bd 100644 --- a/gunicorn/workers/ginterpreter.py +++ b/gunicorn/workers/ginterpreter.py @@ -50,9 +50,8 @@ def _handle_request_in_interpreter(fd, client_addr, server_addr, family): os.close(fd) return - sock = socket.fromfd(fd, family, socket.SOCK_STREAM) + sock = socket.socket(family, socket.SOCK_STREAM, fileno=fd) try: - os.close(fd) sock.settimeout(cfg_dict.get('timeout', 30)) cfg = types.SimpleNamespace(**cfg_dict) From 04e9d9db38cacebf72bf1966b301dc30ab083cb2 Mon Sep 17 00:00:00 2001 From: AN Long Date: Sun, 8 Feb 2026 00:35:30 +0900 Subject: [PATCH 03/14] Mute lint error --- gunicorn/workers/ginterpreter.py | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/gunicorn/workers/ginterpreter.py b/gunicorn/workers/ginterpreter.py index 3fa0f27bd..a922d085a 100644 --- a/gunicorn/workers/ginterpreter.py +++ b/gunicorn/workers/ginterpreter.py @@ -15,7 +15,7 @@ def _check_interpreter_pool_available(): """Check if InterpreterPoolExecutor is available (Python 3.14+).""" try: - from concurrent.futures import InterpreterPoolExecutor # noqa: F401 + from concurrent.futures import InterpreterPoolExecutor # noqa: F401 # pylint: disable=unused-import return True except ImportError: return False @@ -54,7 +54,7 @@ def _handle_request_in_interpreter(fd, client_addr, server_addr, family): try: sock.settimeout(cfg_dict.get('timeout', 30)) - cfg = types.SimpleNamespace(**cfg_dict) + cfg = types.SimpleNamespace(**cfg_dict) # pylint: disable=not-a-mapping cfg.forwarded_allow_networks = lambda: [] cfg.proxy_allow_networks = lambda: [] @@ -102,7 +102,7 @@ def init_process(self): f"Current version: {sys.version_info.major}.{sys.version_info.minor}" ) - from concurrent.futures import InterpreterPoolExecutor + from concurrent.futures import InterpreterPoolExecutor # pylint: disable=no-name-in-module self.cfg_dict = self._extract_config() From 72b8e6009676af320290215f7cfa3a7a822e9080 Mon Sep 17 00:00:00 2001 From: AN Long Date: Sun, 8 Feb 2026 01:02:51 +0900 Subject: [PATCH 04/14] Track connections --- gunicorn/workers/ginterpreter.py | 18 +++++++++++++++++- 1 file changed, 17 insertions(+), 1 deletion(-) diff --git a/gunicorn/workers/ginterpreter.py b/gunicorn/workers/ginterpreter.py index a922d085a..a18341727 100644 --- a/gunicorn/workers/ginterpreter.py +++ b/gunicorn/workers/ginterpreter.py @@ -95,6 +95,11 @@ def _handle_request_in_interpreter(fd, client_addr, server_addr, family): class InterpreterWorker(base.Worker): """Worker using InterpreterPoolExecutor for true parallelism.""" + def __init__(self, *args, **kwargs): + super().__init__(*args, **kwargs) + self.nr_conns = 0 + self.pending_futures = set() + def init_process(self): if not _check_interpreter_pool_available(): raise RuntimeError( @@ -150,12 +155,23 @@ def accept(self, listener): fd = client.fileno() family = client.family server = listener.getsockname() - self.executor.submit( + self.nr_conns += 1 + future = self.executor.submit( _handle_request_in_interpreter, fd, addr, server, family, ) + future.add_done_callback(self._on_request_complete) + self.pending_futures.add(future) client.detach() + def _on_request_complete(self, future): + self.pending_futures.discard(future) + self.nr_conns -= 1 + try: + future.result() + except Exception as e: + self.log.exception("Request failed in sub-interpreter") + def run(self): for listener in self.sockets: listener.setblocking(False) From 30c6015e49b5ab0c9a998101869367bbea0a1265 Mon Sep 17 00:00:00 2001 From: AN Long Date: Sun, 8 Feb 2026 01:13:34 +0900 Subject: [PATCH 05/14] Handle graceful shutdown --- gunicorn/workers/ginterpreter.py | 14 +++++++++++++- 1 file changed, 13 insertions(+), 1 deletion(-) diff --git a/gunicorn/workers/ginterpreter.py b/gunicorn/workers/ginterpreter.py index a18341727..eac56a1d1 100644 --- a/gunicorn/workers/ginterpreter.py +++ b/gunicorn/workers/ginterpreter.py @@ -7,6 +7,7 @@ import errno import os import select +import time import sys from . import base @@ -196,7 +197,18 @@ def run(self): if e.errno != errno.EINTR: raise - self.executor.shutdown(wait=True) + for listener in self.sockets: + listener.close() + + graceful_timeout = time.monotonic() + self.cfg.graceful_timeout + while self.nr_conns > 0: + self.notify() + time_remaining = graceful_timeout - time.monotonic() + if time_remaining <= 0: + break + time.sleep(min(time_remaining, 1.0)) + + self.executor.shutdown(wait=False) def handle_quit(self, sig, frame): self.executor.shutdown(wait=False) From b6eece73dd11d8a7f62a78196f75fcfac4deb64c Mon Sep 17 00:00:00 2001 From: AN Long Date: Sun, 8 Feb 2026 01:35:30 +0900 Subject: [PATCH 06/14] Support ssl --- gunicorn/workers/ginterpreter.py | 44 ++++++++++++++++++++++++++++++++ 1 file changed, 44 insertions(+) diff --git a/gunicorn/workers/ginterpreter.py b/gunicorn/workers/ginterpreter.py index eac56a1d1..e3023fa30 100644 --- a/gunicorn/workers/ginterpreter.py +++ b/gunicorn/workers/ginterpreter.py @@ -10,6 +10,7 @@ import time import sys +from gunicorn.config import NewSSLContext from . import base @@ -25,6 +26,7 @@ def _check_interpreter_pool_available(): _interpreter_state = { 'wsgi_app': None, 'cfg_dict': None, + 'ssl_context': None, } @@ -35,10 +37,24 @@ def _init_interpreter(cfg_dict, app_uri): _interpreter_state['cfg_dict'] = cfg_dict _interpreter_state['wsgi_app'] = import_app(app_uri) + if cfg_dict.get('is_ssl'): + import ssl + context = ssl.create_default_context( + ssl.Purpose.CLIENT_AUTH, cafile=cfg_dict.get('ca_certs') + ) + context.load_cert_chain( + certfile=cfg_dict['certfile'], keyfile=cfg_dict.get('keyfile') + ) + context.verify_mode = cfg_dict.get('cert_reqs', ssl.CERT_NONE) + if cfg_dict.get('ciphers'): + context.set_ciphers(cfg_dict['ciphers']) + _interpreter_state['ssl_context'] = context + def _handle_request_in_interpreter(fd, client_addr, server_addr, family): """Handle a single HTTP request in a sub-interpreter.""" import socket + import ssl import types from gunicorn.http.parser import RequestParser @@ -53,6 +69,17 @@ def _handle_request_in_interpreter(fd, client_addr, server_addr, family): sock = socket.socket(family, socket.SOCK_STREAM, fileno=fd) try: + ssl_context = _interpreter_state['ssl_context'] + if ssl_context is not None: + sock = ssl_context.wrap_socket( + sock, + server_side=True, + suppress_ragged_eofs=cfg_dict.get('suppress_ragged_eofs', True), + do_handshake_on_connect=cfg_dict.get('do_handshake_on_connect', True), + ) + if not cfg_dict.get('do_handshake_on_connect', True): + sock.do_handshake() + sock.settimeout(cfg_dict.get('timeout', 30)) cfg = types.SimpleNamespace(**cfg_dict) # pylint: disable=not-a-mapping @@ -83,6 +110,9 @@ def _handle_request_in_interpreter(fd, client_addr, server_addr, family): except socket.timeout: pass + except ssl.SSLError as e: + if e.args[0] != ssl.SSL_ERROR_EOF: + raise except OSError as e: if e.errno not in (errno.EPIPE, errno.ECONNRESET, errno.ENOTCONN): raise @@ -110,6 +140,13 @@ def init_process(self): from concurrent.futures import InterpreterPoolExecutor # pylint: disable=no-name-in-module + if self.cfg.is_ssl and self.cfg.ssl_context is not NewSSLContext.ssl_context: + self.log.warning( + "ssl_context hook is not supported with ginterpreter worker " + "because callables cannot be shared across sub-interpreters. " + "The hook will be ignored; SSL context is created from config values only." + ) + self.cfg_dict = self._extract_config() self.app_uri = getattr(self.app, 'app_uri', None) or self.app.cfg.wsgi_app @@ -145,6 +182,13 @@ def _extract_config(self): 'proxy_protocol': cfg.proxy_protocol, 'proxy_allow_ips': list(cfg.proxy_allow_ips), 'is_ssl': cfg.is_ssl, + 'certfile': cfg.certfile, + 'keyfile': cfg.keyfile, + 'ca_certs': cfg.ca_certs, + 'cert_reqs': cfg.cert_reqs, + 'ciphers': cfg.ciphers, + 'suppress_ragged_eofs': cfg.suppress_ragged_eofs, + 'do_handshake_on_connect': cfg.do_handshake_on_connect, 'sendfile': cfg.sendfile, 'workers': cfg.workers, 'errorlog': cfg.errorlog, From 0855d6bf472784af0daae655d3f2fe92a79b1bb3 Mon Sep 17 00:00:00 2001 From: AN Long Date: Sun, 8 Feb 2026 12:21:42 +0900 Subject: [PATCH 07/14] Add access log support --- gunicorn/workers/ginterpreter.py | 36 +++++++++++++++++++++++++++++++- 1 file changed, 35 insertions(+), 1 deletion(-) diff --git a/gunicorn/workers/ginterpreter.py b/gunicorn/workers/ginterpreter.py index e3023fa30..bdf307680 100644 --- a/gunicorn/workers/ginterpreter.py +++ b/gunicorn/workers/ginterpreter.py @@ -32,6 +32,9 @@ def _check_interpreter_pool_available(): def _init_interpreter(cfg_dict, app_uri): """Initialize the interpreter with WSGI app and config.""" + import types + + from gunicorn.glogging import Logger from gunicorn.util import import_app _interpreter_state['cfg_dict'] = cfg_dict @@ -50,23 +53,32 @@ def _init_interpreter(cfg_dict, app_uri): context.set_ciphers(cfg_dict['ciphers']) _interpreter_state['ssl_context'] = context + cfg_ns = types.SimpleNamespace(**cfg_dict) + _interpreter_state['log'] = Logger(cfg_ns) + def _handle_request_in_interpreter(fd, client_addr, server_addr, family): """Handle a single HTTP request in a sub-interpreter.""" import socket import ssl import types + from datetime import datetime from gunicorn.http.parser import RequestParser from gunicorn.http.wsgi import create cfg_dict = _interpreter_state['cfg_dict'] wsgi_app = _interpreter_state['wsgi_app'] + log = _interpreter_state['log'] if cfg_dict is None or wsgi_app is None: os.close(fd) return + request_start = datetime.now() + resp = None + environ = None + sock = socket.socket(family, socket.SOCK_STREAM, fileno=fd) try: ssl_context = _interpreter_state['ssl_context'] @@ -117,6 +129,12 @@ def _handle_request_in_interpreter(fd, client_addr, server_addr, family): if e.errno not in (errno.EPIPE, errno.ECONNRESET, errno.ENOTCONN): raise finally: + try: + if resp is not None and environ is not None: + request_time = datetime.now() - request_start + log.access(resp, req, environ, request_time) + except Exception: + pass try: sock.close() except Exception: @@ -191,8 +209,24 @@ def _extract_config(self): 'do_handshake_on_connect': cfg.do_handshake_on_connect, 'sendfile': cfg.sendfile, 'workers': cfg.workers, - 'errorlog': cfg.errorlog, 'timeout': cfg.timeout, + # logging + 'accesslog': cfg.accesslog, + 'access_log_format': cfg.access_log_format, + 'errorlog': cfg.errorlog, + 'loglevel': cfg.loglevel, + 'capture_output': False, + 'syslog': cfg.syslog, + 'syslog_addr': cfg.syslog_addr, + 'syslog_prefix': cfg.syslog_prefix, + 'syslog_facility': cfg.syslog_facility, + 'disable_redirect_access_to_syslog': cfg.disable_redirect_access_to_syslog, + 'logconfig': cfg.logconfig, + 'logconfig_dict': cfg.logconfig_dict, + 'logconfig_json': cfg.logconfig_json, + 'user': cfg.user, + 'group': cfg.group, + 'proc_name': cfg.proc_name, } def accept(self, listener): From 0539df74df80ac4967c6b1bfa74ada67670d0011 Mon Sep 17 00:00:00 2001 From: AN Long Date: Sun, 8 Feb 2026 13:42:50 +0900 Subject: [PATCH 08/14] Raise NotImplmentedError while hooks are set --- gunicorn/workers/ginterpreter.py | 19 +++++++++++++++---- 1 file changed, 15 insertions(+), 4 deletions(-) diff --git a/gunicorn/workers/ginterpreter.py b/gunicorn/workers/ginterpreter.py index bdf307680..4191ef42a 100644 --- a/gunicorn/workers/ginterpreter.py +++ b/gunicorn/workers/ginterpreter.py @@ -10,7 +10,7 @@ import time import sys -from gunicorn.config import NewSSLContext +from gunicorn.config import NewSSLContext, PreRequest, PostRequest from . import base @@ -159,10 +159,21 @@ def init_process(self): from concurrent.futures import InterpreterPoolExecutor # pylint: disable=no-name-in-module if self.cfg.is_ssl and self.cfg.ssl_context is not NewSSLContext.ssl_context: - self.log.warning( + raise NotImplementedError( "ssl_context hook is not supported with ginterpreter worker " - "because callables cannot be shared across sub-interpreters. " - "The hook will be ignored; SSL context is created from config values only." + "because callables cannot be shared across sub-interpreters." + ) + + if self.cfg.pre_request is not PreRequest.pre_request: + raise NotImplementedError( + "pre_request hook is not supported with ginterpreter worker " + "because callables cannot be shared across sub-interpreters." + ) + + if self.cfg.post_request is not PostRequest.post_request: + raise NotImplementedError( + "post_request hook is not supported with ginterpreter worker " + "because callables cannot be shared across sub-interpreters." ) self.cfg_dict = self._extract_config() From c0e54aa95d05143d80844267387c79ec45eccc3e Mon Sep 17 00:00:00 2001 From: AN Long Date: Sun, 8 Feb 2026 20:51:33 +0900 Subject: [PATCH 09/14] Using Config instead of SimpleNamespace and some type refactoring --- gunicorn/workers/ginterpreter.py | 97 +++++++++++++++++--------------- tests/test_ginterpreter.py | 11 +--- 2 files changed, 52 insertions(+), 56 deletions(-) diff --git a/gunicorn/workers/ginterpreter.py b/gunicorn/workers/ginterpreter.py index 4191ef42a..8a06ae682 100644 --- a/gunicorn/workers/ginterpreter.py +++ b/gunicorn/workers/ginterpreter.py @@ -7,11 +7,14 @@ import errno import os import select -import time +import ssl import sys +import time +from dataclasses import dataclass -from gunicorn.config import NewSSLContext, PreRequest, PostRequest -from . import base +from gunicorn.config import Config, NewSSLContext, PreRequest, PostRequest +from gunicorn.glogging import Logger +from gunicorn.workers.base import Worker def _check_interpreter_pool_available(): @@ -20,58 +23,69 @@ def _check_interpreter_pool_available(): from concurrent.futures import InterpreterPoolExecutor # noqa: F401 # pylint: disable=unused-import return True except ImportError: - return False + raise RuntimeError( + "InterpreterPoolExecutor requires Python 3.14+. " + f"Current version: {sys.version_info.major}.{sys.version_info.minor}" + ) -_interpreter_state = { - 'wsgi_app': None, - 'cfg_dict': None, - 'ssl_context': None, -} +@dataclass +class InterpreterState: + cfg: Config | None = None + log: Logger | None = None + ssl_context: ssl.SSLContext | None = None + wsgi_app: object = None -def _init_interpreter(cfg_dict, app_uri): - """Initialize the interpreter with WSGI app and config.""" - import types +_interpreter_state = InterpreterState() + + +def _config_from_dict(cfg_dict: dict) -> Config: + cfg = Config() + for key, value in cfg_dict.items(): + if key in cfg.settings: + cfg.settings[key].value = value + return cfg - from gunicorn.glogging import Logger + +def _init_interpreter(cfg_dict, app_uri) -> None: + """Initialize the interpreter with WSGI app and config.""" from gunicorn.util import import_app - _interpreter_state['cfg_dict'] = cfg_dict - _interpreter_state['wsgi_app'] = import_app(app_uri) + cfg = _config_from_dict(cfg_dict) + _interpreter_state.cfg = cfg + _interpreter_state.wsgi_app = import_app(app_uri) - if cfg_dict.get('is_ssl'): + if cfg.is_ssl: import ssl context = ssl.create_default_context( - ssl.Purpose.CLIENT_AUTH, cafile=cfg_dict.get('ca_certs') + ssl.Purpose.CLIENT_AUTH, cafile=cfg.ca_certs ) context.load_cert_chain( - certfile=cfg_dict['certfile'], keyfile=cfg_dict.get('keyfile') + certfile=cfg.certfile, keyfile=cfg.keyfile ) - context.verify_mode = cfg_dict.get('cert_reqs', ssl.CERT_NONE) - if cfg_dict.get('ciphers'): - context.set_ciphers(cfg_dict['ciphers']) - _interpreter_state['ssl_context'] = context + context.verify_mode = cfg.cert_reqs + if cfg.ciphers: + context.set_ciphers(cfg.ciphers) + _interpreter_state.ssl_context = context - cfg_ns = types.SimpleNamespace(**cfg_dict) - _interpreter_state['log'] = Logger(cfg_ns) + _interpreter_state.log = Logger(cfg) def _handle_request_in_interpreter(fd, client_addr, server_addr, family): """Handle a single HTTP request in a sub-interpreter.""" import socket import ssl - import types from datetime import datetime from gunicorn.http.parser import RequestParser from gunicorn.http.wsgi import create - cfg_dict = _interpreter_state['cfg_dict'] + cfg = _interpreter_state['cfg'] wsgi_app = _interpreter_state['wsgi_app'] log = _interpreter_state['log'] - if cfg_dict is None or wsgi_app is None: + if cfg is None or wsgi_app is None: os.close(fd) return @@ -81,22 +95,18 @@ def _handle_request_in_interpreter(fd, client_addr, server_addr, family): sock = socket.socket(family, socket.SOCK_STREAM, fileno=fd) try: - ssl_context = _interpreter_state['ssl_context'] + ssl_context = _interpreter_state.get('ssl_context') if ssl_context is not None: sock = ssl_context.wrap_socket( sock, server_side=True, - suppress_ragged_eofs=cfg_dict.get('suppress_ragged_eofs', True), - do_handshake_on_connect=cfg_dict.get('do_handshake_on_connect', True), + suppress_ragged_eofs=cfg.suppress_ragged_eofs, + do_handshake_on_connect=cfg.do_handshake_on_connect, ) - if not cfg_dict.get('do_handshake_on_connect', True): + if not cfg.do_handshake_on_connect: sock.do_handshake() - sock.settimeout(cfg_dict.get('timeout', 30)) - - cfg = types.SimpleNamespace(**cfg_dict) # pylint: disable=not-a-mapping - cfg.forwarded_allow_networks = lambda: [] - cfg.proxy_allow_networks = lambda: [] + sock.settimeout(cfg.timeout) parser = RequestParser(cfg, sock, client_addr) try: @@ -141,7 +151,7 @@ def _handle_request_in_interpreter(fd, client_addr, server_addr, family): pass -class InterpreterWorker(base.Worker): +class InterpreterWorker(Worker): """Worker using InterpreterPoolExecutor for true parallelism.""" def __init__(self, *args, **kwargs): @@ -150,11 +160,7 @@ def __init__(self, *args, **kwargs): self.pending_futures = set() def init_process(self): - if not _check_interpreter_pool_available(): - raise RuntimeError( - "InterpreterPoolExecutor requires Python 3.14+. " - f"Current version: {sys.version_info.major}.{sys.version_info.minor}" - ) + _check_interpreter_pool_available() from concurrent.futures import InterpreterPoolExecutor # pylint: disable=no-name-in-module @@ -193,7 +199,7 @@ def init_process(self): super().init_process() - def _extract_config(self): + def _extract_config(self) -> dict: cfg = self.cfg return { 'limit_request_line': cfg.limit_request_line, @@ -210,7 +216,6 @@ def _extract_config(self): 'secure_scheme_headers': dict(cfg.secure_scheme_headers), 'proxy_protocol': cfg.proxy_protocol, 'proxy_allow_ips': list(cfg.proxy_allow_ips), - 'is_ssl': cfg.is_ssl, 'certfile': cfg.certfile, 'keyfile': cfg.keyfile, 'ca_certs': cfg.ca_certs, @@ -218,7 +223,7 @@ def _extract_config(self): 'ciphers': cfg.ciphers, 'suppress_ragged_eofs': cfg.suppress_ragged_eofs, 'do_handshake_on_connect': cfg.do_handshake_on_connect, - 'sendfile': cfg.sendfile, + 'sendfile': cfg.settings['sendfile'].value, 'workers': cfg.workers, 'timeout': cfg.timeout, # logging @@ -259,7 +264,7 @@ def _on_request_complete(self, future): self.nr_conns -= 1 try: future.result() - except Exception as e: + except Exception: self.log.exception("Request failed in sub-interpreter") def run(self): diff --git a/tests/test_ginterpreter.py b/tests/test_ginterpreter.py index d634bfb5f..08dff1c4b 100644 --- a/tests/test_ginterpreter.py +++ b/tests/test_ginterpreter.py @@ -38,15 +38,6 @@ def _create_worker(cfg=None): class TestInterpreterWorker: - def test_extract_config(self): - worker = _create_worker() - cfg_dict = worker._extract_config() - assert isinstance(cfg_dict, dict) - assert 'limit_request_line' in cfg_dict - assert 'timeout' in cfg_dict - for value in cfg_dict.values(): - assert isinstance(value, (int, bool, str, list, dict, type(None))) - def test_handle_quit(self): worker = _create_worker() worker.executor = mock.Mock() @@ -74,4 +65,4 @@ def test_submit_to_executor(self): ginterpreter._handle_request_in_interpreter, 7, ('127.0.0.1', 8000), ('0.0.0.0', 9000), 2, ) - mock_client.detach.assert_called_once() \ No newline at end of file + mock_client.detach.assert_called_once() From 8ad4062356807124662a77a7aff9216218da9fc0 Mon Sep 17 00:00:00 2001 From: AN Long Date: Sun, 8 Feb 2026 21:11:18 +0900 Subject: [PATCH 10/14] Add some tests --- gunicorn/workers/ginterpreter.py | 8 ++-- tests/test_ginterpreter.py | 80 +++++++++++++++++++++++++++++++- 2 files changed, 83 insertions(+), 5 deletions(-) diff --git a/gunicorn/workers/ginterpreter.py b/gunicorn/workers/ginterpreter.py index 8a06ae682..50ef78bcc 100644 --- a/gunicorn/workers/ginterpreter.py +++ b/gunicorn/workers/ginterpreter.py @@ -81,9 +81,9 @@ def _handle_request_in_interpreter(fd, client_addr, server_addr, family): from gunicorn.http.parser import RequestParser from gunicorn.http.wsgi import create - cfg = _interpreter_state['cfg'] - wsgi_app = _interpreter_state['wsgi_app'] - log = _interpreter_state['log'] + cfg = _interpreter_state.cfg + wsgi_app = _interpreter_state.wsgi_app + log = _interpreter_state.log if cfg is None or wsgi_app is None: os.close(fd) @@ -95,7 +95,7 @@ def _handle_request_in_interpreter(fd, client_addr, server_addr, family): sock = socket.socket(family, socket.SOCK_STREAM, fileno=fd) try: - ssl_context = _interpreter_state.get('ssl_context') + ssl_context = _interpreter_state.ssl_context if ssl_context is not None: sock = ssl_context.wrap_socket( sock, diff --git a/tests/test_ginterpreter.py b/tests/test_ginterpreter.py index 08dff1c4b..d6095d1bb 100644 --- a/tests/test_ginterpreter.py +++ b/tests/test_ginterpreter.py @@ -5,6 +5,7 @@ """Tests for the ginterpreter worker.""" import os +import ssl from unittest import mock import pytest @@ -46,7 +47,7 @@ def test_handle_quit(self): worker.executor.shutdown.assert_called_once_with(wait=False) -class TestAccept: +class TestHandleRequest: def test_submit_to_executor(self): worker = _create_worker() @@ -66,3 +67,80 @@ def test_submit_to_executor(self): 7, ('127.0.0.1', 8000), ('0.0.0.0', 9000), 2, ) mock_client.detach.assert_called_once() + + + @mock.patch('gunicorn.http.wsgi.create') + @mock.patch('gunicorn.http.parser.RequestParser') + @mock.patch('socket.socket') + def test_handle_request(self, mock_socket_cls, mock_parser_cls, mock_create): + mock_sock = mock.Mock() + mock_socket_cls.return_value = mock_sock + + mock_req = mock.Mock() + mock_parser_cls.return_value = iter([mock_req]) + + mock_resp = mock.Mock() + mock_environ = {'wsgi.multithread': False, 'wsgi.multiprocess': False} + mock_create.return_value = (mock_resp, mock_environ) + + mock_wsgi_app = mock.Mock(return_value=[b'response']) + + cfg = ginterpreter._config_from_dict({ + 'timeout': 30, + 'forwarded_allow_ips': [], + 'proxy_allow_ips': [], + }) + ginterpreter._interpreter_state = ginterpreter.InterpreterState( + cfg=cfg, + ssl_context=None, + wsgi_app=mock_wsgi_app, + log=mock.Mock(), + ) + + ginterpreter._handle_request_in_interpreter( + 7, ('127.0.0.1', 8000), ('0.0.0.0', 9000), 2 + ) + + mock_socket_cls.assert_called_once_with(2, mock.ANY, fileno=7) + mock_sock.settimeout.assert_called_once_with(30) + mock_wsgi_app.assert_called_once() + + +class TestSSL: + + @mock.patch('socket.socket') + def test_handle_request_wraps_socket_with_ssl(self, mock_socket_cls): + mock_sock = mock.Mock() + mock_socket_cls.return_value = mock_sock + + mock_ssl_ctx = mock.Mock() + mock_wrapped = mock.Mock() + mock_ssl_ctx.wrap_socket.return_value = mock_wrapped + mock_wrapped.settimeout = mock.Mock() + + cfg = ginterpreter._config_from_dict({ + 'suppress_ragged_eofs': True, + 'do_handshake_on_connect': True, + 'timeout': 30, + 'forwarded_allow_ips': [], + 'proxy_allow_ips': [], + }) + ginterpreter._interpreter_state = ginterpreter.InterpreterState( + cfg=cfg, + ssl_context=mock_ssl_ctx, + wsgi_app=mock.Mock(), + log=mock.Mock(), + ) + + with mock.patch('gunicorn.http.parser.RequestParser') as mock_parser: + mock_parser.return_value = iter([]) + ginterpreter._handle_request_in_interpreter( + 7, ('127.0.0.1', 8000), ('0.0.0.0', 9000), 2 + ) + + mock_ssl_ctx.wrap_socket.assert_called_once_with( + mock_sock, + server_side=True, + suppress_ragged_eofs=True, + do_handshake_on_connect=True, + ) From 8b0156276134225e19b878cb56a1751662e01575 Mon Sep 17 00:00:00 2001 From: AN Long Date: Sat, 14 Mar 2026 17:09:53 +0900 Subject: [PATCH 11/14] Fix tox-lint --- gunicorn/workers/ginterpreter.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/gunicorn/workers/ginterpreter.py b/gunicorn/workers/ginterpreter.py index 50ef78bcc..259c00435 100644 --- a/gunicorn/workers/ginterpreter.py +++ b/gunicorn/workers/ginterpreter.py @@ -57,7 +57,7 @@ def _init_interpreter(cfg_dict, app_uri) -> None: _interpreter_state.wsgi_app = import_app(app_uri) if cfg.is_ssl: - import ssl + import ssl # pylint: disable=reimported context = ssl.create_default_context( ssl.Purpose.CLIENT_AUTH, cafile=cfg.ca_certs ) @@ -75,7 +75,7 @@ def _init_interpreter(cfg_dict, app_uri) -> None: def _handle_request_in_interpreter(fd, client_addr, server_addr, family): """Handle a single HTTP request in a sub-interpreter.""" import socket - import ssl + import ssl # pylint: disable=reimported from datetime import datetime from gunicorn.http.parser import RequestParser From dbf72b5fc0a319301e5d4e53e19791207618113a Mon Sep 17 00:00:00 2001 From: AN Long Date: Sat, 14 Mar 2026 17:45:17 +0900 Subject: [PATCH 12/14] Add ginterpreter.md --- docs/content/ginterpreter.md | 51 ++++++++++++++++++++++++++++++++++++ mkdocs.yml | 1 + 2 files changed, 52 insertions(+) create mode 100644 docs/content/ginterpreter.md diff --git a/docs/content/ginterpreter.md b/docs/content/ginterpreter.md new file mode 100644 index 000000000..7f0d7bcdb --- /dev/null +++ b/docs/content/ginterpreter.md @@ -0,0 +1,51 @@ +# Interpreter Worker + +!!! warning "Experimental" + The `ginterpreter` worker is experimental and requires Python 3.14+. The API + and behavior may change in future releases. + +The interpreter worker uses Python's `InterpreterPoolExecutor` to handle each +request in a separate sub-interpreter. Each sub-interpreter runs in its own +thread with an independent GIL, enabling true CPU parallelism without multiple +processes. + +## Quick Start + +```bash +gunicorn myapp:app --worker-class ginterpreter --threads 4 +``` + +Or in a configuration file: + +```python +# gunicorn.conf.py +worker_class = "ginterpreter" +threads = 4 +``` + +## Configuration + +The interpreter worker uses the standard gunicorn settings. The most relevant ones: + +| Setting | Default | Description | +|---------|---------|-------------| +| `threads` | `1` | Number of sub-interpreters (i.e. concurrent requests per worker) | +| `workers` | `1` | Number of worker processes | +| `timeout` | `30` | Request timeout in seconds | +| `graceful_timeout` | `30` | Time to wait for in-flight requests on shutdown | + +## Known Limitations + +The following features are **not supported**: + +- **`ssl_context` hook** — SSL contexts cannot be shared across sub-interpreters. Built-in SSL via `certfile`/`keyfile` works normally. +- **`pre_request` / `post_request` hooks** — Callables cannot be passed to sub-interpreters. +- **Keepalive connections** — Each connection is closed after the response. +- **HTTP/2** +- **Sendfile** +- **`max_requests` / `max_requests_jitter`** + +## See Also + +- [Settings Reference](reference/settings.md) - All available settings +- [Design](design.md) - Worker architecture overview diff --git a/mkdocs.yml b/mkdocs.yml index f01f0da79..704f3e26f 100644 --- a/mkdocs.yml +++ b/mkdocs.yml @@ -17,6 +17,7 @@ nav: - Docker: guides/docker.md - HTTP/2: guides/http2.md - ASGI Worker: asgi.md + - Interpreter Worker: ginterpreter.md - Dirty Arbiters: dirty.md - uWSGI Protocol: uwsgi.md - Signals: signals.md From 506f949c1c1a5a2fe1bb58ecbad3138c80ad00e9 Mon Sep 17 00:00:00 2001 From: AN Long Date: Sat, 14 Mar 2026 17:49:38 +0900 Subject: [PATCH 13/14] Update ginterpreter worker in design.md --- docs/content/design.md | 19 +++++++++++++++++++ 1 file changed, 19 insertions(+) diff --git a/docs/content/design.md b/docs/content/design.md index c06c1cbd5..19c85c4da 100644 --- a/docs/content/design.md +++ b/docs/content/design.md @@ -124,6 +124,23 @@ Choose a worker type based on your application's needs. gunicorn myapp:app -k tornado ``` +=== "Interpreter (Experimental)" + + !!! warning "Experimental — requires Python 3.14+" + + **Sub-interpreter** worker using `InterpreterPoolExecutor`. Each request runs + in its own sub-interpreter with an independent GIL, enabling true CPU + parallelism without forking extra processes. + + - True parallelism for CPU-bound workloads + - No keep-alive support + - Hooks (`ssl_context`, `pre_request`, `post_request`) are not supported + - See the [Interpreter Worker](ginterpreter.md) guide for details + + ```bash + gunicorn myapp:app -k ginterpreter --threads 4 + ``` + ## Comparison | Worker | Concurrency Model | Keep-Alive | Best For | @@ -134,6 +151,7 @@ Choose a worker type based on your application's needs. | `gevent` | Greenlets | ✅ | I/O-bound, WebSockets, streaming | | `eventlet` | Greenlets | ✅ | **Deprecated** - use `gevent` instead | | `tornado` | Tornado IOLoop | ✅ | Native Tornado applications | +| `ginterpreter` *(experimental)* | Sub-interpreters | ❌ | CPU-bound apps on Python 3.14+ | !!! tip "Quick Decision Guide" @@ -141,6 +159,7 @@ Choose a worker type based on your application's needs. - **Need keep-alive or moderate concurrency?** → `gthread` - **WebSockets, streaming, long-polling?** → `gevent` or ASGI worker - **FastAPI, Starlette, or async framework?** → ASGI worker + - **CPU-bound and on Python 3.14+?** → `ginterpreter` *(experimental)* ## When to Use Async Workers From 87c1a783cbf1bc203d1fb2069d332b2a60a35854 Mon Sep 17 00:00:00 2001 From: AN Long Date: Sat, 14 Mar 2026 18:36:59 +0900 Subject: [PATCH 14/14] Using max(max_conns, threads) to select the fds from kernel --- gunicorn/workers/ginterpreter.py | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/gunicorn/workers/ginterpreter.py b/gunicorn/workers/ginterpreter.py index 259c00435..9731ffb54 100644 --- a/gunicorn/workers/ginterpreter.py +++ b/gunicorn/workers/ginterpreter.py @@ -158,6 +158,7 @@ def __init__(self, *args, **kwargs): super().__init__(*args, **kwargs) self.nr_conns = 0 self.pending_futures = set() + self.max_conns = 0 def init_process(self): _check_interpreter_pool_available() @@ -196,6 +197,7 @@ def init_process(self): initializer=_init_interpreter, initargs=(self.cfg_dict, self.app_uri), ) + self.max_conns = min(self.cfg.threads, self.cfg.worker_connections) super().init_process() @@ -279,7 +281,8 @@ def run(self): break try: - ready = select.select(self.sockets, [], [], 1.0) + listeners = self.sockets if self.nr_conns < self.max_conns else [] + ready = select.select(listeners, [], [], 1.0) for listener in ready[0]: try: self.accept(listener)