diff --git a/maestro_worker_python/request_logging.py b/maestro_worker_python/request_logging.py new file mode 100644 index 0000000..cabcf92 --- /dev/null +++ b/maestro_worker_python/request_logging.py @@ -0,0 +1,36 @@ +"""Access-log instrumentation hardened against a missing peer address. + +json_logging 1.5.1 dereferences Starlette's optional ``Request.client`` +unguarded, so a request whose peer socket is already gone (kubelet resetting a +timed-out readiness probe, typically) loses its access-log line to a formatter +traceback on stderr. Delete this once bobbui/json-logging-python#116 reaches a +release. +""" + +import json_logging +from json_logging.framework.fastapi import ( + FastAPIAppRequestInstrumentationConfigurator, + FastAPIRequestInfoExtractor, + FastAPIResponseInfoExtractor, +) +from json_logging.frameworks import register_framework_support +from starlette.requests import Request + + +class ClientSafeRequestInfoExtractor(FastAPIRequestInfoExtractor): + def get_remote_ip(self, request: Request) -> str: + return request.client.host if request.client else json_logging.EMPTY_VALUE + + def get_remote_port(self, request: Request) -> int | str: + return request.client.port if request.client else json_logging.EMPTY_VALUE + + +def register_client_safe_request_extractor() -> None: + """Must run before ``init_fastapi``, which freezes the extractor in a singleton.""" + register_framework_support( + "fastapi", + app_configurator=None, + app_request_instrumentation_configurator=FastAPIAppRequestInstrumentationConfigurator, + request_info_extractor_class=ClientSafeRequestInfoExtractor, + response_info_extractor_class=FastAPIResponseInfoExtractor, + ) diff --git a/maestro_worker_python/serve.py b/maestro_worker_python/serve.py index 7a262eb..5a1964b 100644 --- a/maestro_worker_python/serve.py +++ b/maestro_worker_python/serve.py @@ -21,6 +21,7 @@ from .health import get_health_metadata from .kill_process import kill_child_processes, terminate_current_process from .load_worker import load_worker +from .request_logging import register_client_safe_request_extractor from .response import ValidationError, WorkerResponse @@ -87,6 +88,7 @@ async def lifespan(_app: FastAPI) -> AsyncIterator[None]: logging.basicConfig(level=settings.log_level.upper()) if settings.enable_json_logging: + register_client_safe_request_extractor() json_logging.init_fastapi(enable_json=True) json_logging.init_request_instrument(app) json_logging.config_root_logger() diff --git a/pyproject.toml b/pyproject.toml index 010a333..c835820 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -1,6 +1,6 @@ [project] name = "maestro-worker-python" -version = "5.1.1" +version = "5.1.2" description = "Utility to run workers on Moises/Maestro" readme = "README.md" requires-python = ">=3.10,<3.14" diff --git a/tests/test_request_logging.py b/tests/test_request_logging.py new file mode 100644 index 0000000..cec09b3 --- /dev/null +++ b/tests/test_request_logging.py @@ -0,0 +1,87 @@ +import subprocess +import sys +import textwrap + +import json_logging +import pytest +from starlette.requests import Request + +from maestro_worker_python.request_logging import ( + ClientSafeRequestInfoExtractor, + register_client_safe_request_extractor, +) + + +def _request(client: tuple[str, int] | None) -> Request: + return Request({"type": "http", "method": "GET", "path": "/health", "headers": [], "client": client}) + + +@pytest.fixture +def extractor() -> ClientSafeRequestInfoExtractor: + return ClientSafeRequestInfoExtractor() + + +def test_reports_peer_address_when_present(extractor): + request = _request(("10.0.0.1", 54321)) + assert extractor.get_remote_ip(request) == "10.0.0.1" + assert extractor.get_remote_port(request) == 54321 + + +def test_reports_empty_value_when_peer_is_gone(extractor): + request = _request(None) + assert extractor.get_remote_ip(request) == json_logging.EMPTY_VALUE + assert extractor.get_remote_port(request) == json_logging.EMPTY_VALUE + + +def test_registration_replaces_only_the_request_extractor(monkeypatch: pytest.MonkeyPatch): + before = json_logging._framework_support_map["fastapi"] + monkeypatch.setitem(json_logging._framework_support_map, "fastapi", dict(before)) + + register_client_safe_request_extractor() + + after = json_logging._framework_support_map["fastapi"] + assert after["request_info_extractor_class"] is ClientSafeRequestInfoExtractor + assert {k: v for k, v in after.items() if k != "request_info_extractor_class"} == { + k: v for k, v in before.items() if k != "request_info_extractor_class" + } + + +# json_logging's init mutates process-wide logging state and caches the extractor in a +# singleton, so the serving path can only be exercised honestly in a fresh interpreter. +SERVE_A_REQUEST_WITHOUT_A_PEER = """ +import asyncio, sys +from maestro_worker_python import serve + +scope = { + "type": "http", "asgi": {"version": "3.0"}, "http_version": "1.1", "scheme": "http", + "method": "GET", "path": "/health", "raw_path": b"/health", "query_string": b"", + "root_path": "", "headers": [], "client": None, "server": ("testserver", 80), +} + +async def receive(): + return {"type": "http.request", "body": b"", "more_body": False} + +messages = [] + +async def send(message): + messages.append(message) + +asyncio.run(serve.app(scope, receive, send)) +print([m["status"] for m in messages if m["type"] == "http.response.start"], file=sys.stderr) +""" + + +def test_serving_a_request_without_a_peer_logs_cleanly(tmp_path): + worker_path = tmp_path / "worker.py" + worker_path.write_text("class MoisesWorker:\n pass\n") + + result = subprocess.run( + [sys.executable, "-c", textwrap.dedent(SERVE_A_REQUEST_WITHOUT_A_PEER)], + capture_output=True, + text=True, + env={"ENABLE_JSON_LOGGING": "true", "MODEL_PATH": str(worker_path)}, + ) + + assert result.returncode == 0, result.stderr + assert "[200]" in result.stderr + assert "Logging error" not in result.stderr, result.stderr diff --git a/uv.lock b/uv.lock index d30db63..6258f4d 100644 --- a/uv.lock +++ b/uv.lock @@ -220,7 +220,7 @@ wheels = [ [[package]] name = "maestro-worker-python" -version = "5.1.1" +version = "5.1.2" source = { editable = "." } dependencies = [ { name = "fastapi" },