Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
36 changes: 36 additions & 0 deletions maestro_worker_python/request_logging.py
Original file line number Diff line number Diff line change
@@ -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,
)
2 changes: 2 additions & 0 deletions maestro_worker_python/serve.py
Original file line number Diff line number Diff line change
Expand Up @@ -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


Expand Down Expand Up @@ -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()
Expand Down
2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
@@ -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"
Expand Down
87 changes: 87 additions & 0 deletions tests/test_request_logging.py
Original file line number Diff line number Diff line change
@@ -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
2 changes: 1 addition & 1 deletion uv.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Loading