diff --git a/nextcloud_mcp_server/cli.py b/nextcloud_mcp_server/cli.py index aa714930..a787fdeb 100644 --- a/nextcloud_mcp_server/cli.py +++ b/nextcloud_mcp_server/cli.py @@ -6,6 +6,7 @@ import click import uvicorn from nextcloud_mcp_server.config import ( + Settings, get_database_url, get_settings, is_ephemeral_token_db, @@ -17,7 +18,12 @@ from nextcloud_mcp_server.migrations import ( show_migration_history, upgrade_database, ) -from nextcloud_mcp_server.observability import get_uvicorn_logging_config +from nextcloud_mcp_server.observability import ( + get_uvicorn_logging_config, + setup_logging, + setup_metrics, + setup_tracing, +) from nextcloud_mcp_server.server import AVAILABLE_APPS from .app import get_app @@ -284,6 +290,44 @@ def run( ) +def _init_worker_observability(settings: Settings) -> None: + """Configure logging, metrics, and tracing for the ingest worker. + + Mirrors the observability bootstrap the API pod performs in its lifespan + (``app.py``), but for the standalone ``worker`` entrypoint which never runs + uvicorn. Without this the worker emits plain-text logs and serves no + ``/metrics`` endpoint, so the astrolabe_* document-pipeline metrics and the + ``document_processor.parse`` spans (recorded in the shared registry/processor + code the worker executes) stay invisible in external split-worker mode + (Deck #310 / #175). + """ + # Structured logging first, so every subsequent startup line is JSON like + # the API's — the worker entrypoint never went through uvicorn's log_config. + setup_logging( + log_format=settings.log_format, + log_level=settings.log_level, + include_trace_context=settings.log_include_trace_context, + ) + + if settings.metrics_enabled: + setup_metrics(port=settings.metrics_port) + logger.info( + "Prometheus metrics enabled on dedicated port %s", settings.metrics_port + ) + + if settings.otel_exporter_otlp_endpoint: + setup_tracing( + service_name=settings.otel_service_name, + otlp_endpoint=settings.otel_exporter_otlp_endpoint, + otlp_verify_ssl=settings.otel_exporter_verify_ssl, + sampling_rate=settings.otel_traces_sampler_arg, + ) + logger.info( + "OpenTelemetry tracing enabled (endpoint: %s)", + settings.otel_exporter_otlp_endpoint, + ) + + @click.command() @click.option( "--concurrency", @@ -319,6 +363,15 @@ def worker(concurrency: int | None): f"resolved INGEST_QUEUE={settings.ingest_queue!r}" ) + # Initialize observability once the config is known to be runnable. The + # always-on API pod does this in its lifespan (app.py); the worker has its + # own entrypoint, so without this it emits plain-text logs and exposes no + # /metrics — leaving the ingest workload (which does the real + # parse/embed/upsert work, and where the astrolabe_* pipeline metrics + + # document_processor.parse spans are recorded) invisible in external + # split-worker mode (Deck #310, unblocks #175). + _init_worker_observability(settings) + from nextcloud_mcp_server.vector.queue.procrastinate import ( # noqa: PLC0415 INGEST_QUEUE_NAME, apply_ingest_queue_schema, diff --git a/tests/test_cli.py b/tests/test_cli.py index 9d25dae9..c54fcc60 100644 --- a/tests/test_cli.py +++ b/tests/test_cli.py @@ -1,11 +1,12 @@ """Tests for CLI options using Click's testing utilities.""" import os +from types import SimpleNamespace import pytest from click.testing import CliRunner -from nextcloud_mcp_server.cli import run +from nextcloud_mcp_server.cli import _init_worker_observability, run, worker @pytest.fixture @@ -324,3 +325,146 @@ def test_stdio_calls_get_stdio_mcp(runner, clean_env, monkeypatch): assert result.exit_code == 0, result.output assert called_with.get("transport") == "stdio" assert called_with.get("enabled_apps") is None + + +# --------------------------------------------------------------------------- +# Ingest worker observability bootstrap (Deck #310 / #175) +# --------------------------------------------------------------------------- + + +def _fake_settings(**overrides): + """A lightweight settings stand-in for the worker observability helper. + + The helper only reads attributes, so a SimpleNamespace avoids running the + real Settings.__post_init__ validation/derivation. + """ + base = dict( + ingest_queue="postgres", + log_format="json", + log_level="INFO", + log_include_trace_context=True, + metrics_enabled=True, + metrics_port=9090, + otel_exporter_otlp_endpoint=None, + otel_service_name="nextcloud-mcp-server", + otel_exporter_verify_ssl=False, + otel_traces_sampler_arg=1.0, + ) + base.update(overrides) + return SimpleNamespace(**base) + + +@pytest.fixture +def patched_observability(monkeypatch): + """Patch the worker's observability entrypoints and record their kwargs.""" + calls: dict[str, dict] = {} + monkeypatch.setattr( + "nextcloud_mcp_server.cli.setup_logging", + lambda **kw: calls.__setitem__("logging", kw), + ) + monkeypatch.setattr( + "nextcloud_mcp_server.cli.setup_metrics", + lambda **kw: calls.__setitem__("metrics", kw), + ) + monkeypatch.setattr( + "nextcloud_mcp_server.cli.setup_tracing", + lambda **kw: calls.__setitem__("tracing", kw), + ) + return calls + + +def test_init_worker_observability_configures_logging(patched_observability): + """Worker initializes structured logging from settings (AC: JSON logs).""" + _init_worker_observability(_fake_settings()) + + assert patched_observability["logging"] == { + "log_format": "json", + "log_level": "INFO", + "include_trace_context": True, + } + + +def test_init_worker_observability_starts_metrics_when_enabled(patched_observability): + """Worker starts the Prometheus server on the configured port (AC: /metrics).""" + _init_worker_observability(_fake_settings(metrics_port=9123)) + + assert patched_observability["metrics"] == {"port": 9123} + + +def test_init_worker_observability_skips_metrics_when_disabled(patched_observability): + """METRICS_ENABLED=false leaves the worker without a metrics server.""" + _init_worker_observability(_fake_settings(metrics_enabled=False)) + + assert "metrics" not in patched_observability + # Logging is still configured regardless of the metrics toggle. + assert "logging" in patched_observability + + +def test_init_worker_observability_sets_up_tracing_when_endpoint( + patched_observability, +): + """An OTLP endpoint enables tracing so worker spans (parse/embed) export.""" + _init_worker_observability( + _fake_settings( + otel_exporter_otlp_endpoint="http://otel:4317", + otel_traces_sampler_arg=0.5, + ) + ) + + assert patched_observability["tracing"] == { + "service_name": "nextcloud-mcp-server", + "otlp_endpoint": "http://otel:4317", + "otlp_verify_ssl": False, + "sampling_rate": 0.5, + } + + +def test_init_worker_observability_skips_tracing_without_endpoint( + patched_observability, +): + """No OTLP endpoint → tracing stays disabled (matches API pod behavior).""" + _init_worker_observability(_fake_settings(otel_exporter_otlp_endpoint=None)) + + assert "tracing" not in patched_observability + + +def test_worker_initializes_observability_on_postgres_queue(runner, monkeypatch): + """The worker command wires up observability once config is runnable.""" + monkeypatch.setattr( + "nextcloud_mcp_server.cli.get_settings", + lambda: _fake_settings(ingest_queue="postgres"), + ) + + called = {} + + def fake_init(settings): + called["settings"] = settings + # Stop before the procrastinate/worker machinery. + raise SystemExit(0) + + monkeypatch.setattr( + "nextcloud_mcp_server.cli._init_worker_observability", fake_init + ) + + result = runner.invoke(worker, []) + assert result.exit_code == 0, result.output + assert called.get("settings") is not None + + +def test_worker_rejects_non_postgres_queue_before_observability(runner, monkeypatch): + """A non-postgres queue fails fast, before any metrics server is started.""" + monkeypatch.setattr( + "nextcloud_mcp_server.cli.get_settings", + lambda: _fake_settings(ingest_queue="memory"), + ) + + called = {} + monkeypatch.setattr( + "nextcloud_mcp_server.cli._init_worker_observability", + lambda settings: called.setdefault("init", True), + ) + + result = runner.invoke(worker, []) + assert result.exit_code != 0 + assert "INGEST_QUEUE=postgres" in result.output + assert "init" not in called