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 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/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..9731ffb54 --- /dev/null +++ b/gunicorn/workers/ginterpreter.py @@ -0,0 +1,312 @@ +# +# 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 ssl +import sys +import time +from dataclasses import dataclass + +from gunicorn.config import Config, NewSSLContext, PreRequest, PostRequest +from gunicorn.glogging import Logger +from gunicorn.workers.base import Worker + + +def _check_interpreter_pool_available(): + """Check if InterpreterPoolExecutor is available (Python 3.14+).""" + try: + from concurrent.futures import InterpreterPoolExecutor # noqa: F401 # pylint: disable=unused-import + return True + except ImportError: + raise RuntimeError( + "InterpreterPoolExecutor requires Python 3.14+. " + f"Current version: {sys.version_info.major}.{sys.version_info.minor}" + ) + + +@dataclass +class InterpreterState: + cfg: Config | None = None + log: Logger | None = None + ssl_context: ssl.SSLContext | None = None + wsgi_app: object = None + + +_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 + + +def _init_interpreter(cfg_dict, app_uri) -> None: + """Initialize the interpreter with WSGI app and config.""" + from gunicorn.util import import_app + + cfg = _config_from_dict(cfg_dict) + _interpreter_state.cfg = cfg + _interpreter_state.wsgi_app = import_app(app_uri) + + if cfg.is_ssl: + import ssl # pylint: disable=reimported + context = ssl.create_default_context( + ssl.Purpose.CLIENT_AUTH, cafile=cfg.ca_certs + ) + context.load_cert_chain( + certfile=cfg.certfile, keyfile=cfg.keyfile + ) + context.verify_mode = cfg.cert_reqs + if cfg.ciphers: + context.set_ciphers(cfg.ciphers) + _interpreter_state.ssl_context = context + + _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 # pylint: disable=reimported + from datetime import datetime + + 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 + + if cfg 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 + if ssl_context is not None: + sock = ssl_context.wrap_socket( + sock, + server_side=True, + suppress_ragged_eofs=cfg.suppress_ragged_eofs, + do_handshake_on_connect=cfg.do_handshake_on_connect, + ) + if not cfg.do_handshake_on_connect: + sock.do_handshake() + + sock.settimeout(cfg.timeout) + + 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 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 + 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: + pass + + +class InterpreterWorker(Worker): + """Worker using InterpreterPoolExecutor for true parallelism.""" + + 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() + + 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: + raise NotImplementedError( + "ssl_context hook is not supported with ginterpreter worker " + "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() + + 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), + ) + self.max_conns = min(self.cfg.threads, self.cfg.worker_connections) + + super().init_process() + + def _extract_config(self) -> dict: + 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), + '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.settings['sendfile'].value, + 'workers': cfg.workers, + '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): + client, addr = listener.accept() + fd = client.fileno() + family = client.family + server = listener.getsockname() + 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: + self.log.exception("Request failed in sub-interpreter") + + 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: + 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) + 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 + + 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) + super().handle_quit(sig, frame) 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 diff --git a/tests/test_ginterpreter.py b/tests/test_ginterpreter.py new file mode 100644 index 000000000..d6095d1bb --- /dev/null +++ b/tests/test_ginterpreter.py @@ -0,0 +1,146 @@ +# +# This file is part of gunicorn released under the MIT license. +# See the NOTICE for more information. + +"""Tests for the ginterpreter worker.""" + +import os +import ssl +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_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 TestHandleRequest: + + 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() + + + @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, + )