diff --git a/VM2_services/codebase/services/.env.example b/VM2_services/codebase/services/.env.example index 20216a3..b7b1444 100644 --- a/VM2_services/codebase/services/.env.example +++ b/VM2_services/codebase/services/.env.example @@ -1,6 +1,7 @@ # Non-secret VM2 deployment manifest. Never add DSNs, passwords, tokens or keys. APP_ENV=production-like RELEASE_VERSION= +LOG_LEVEL=INFO SECRETS_SOURCE=selectel MESSAGE_SAFETY_IMAGE=/han-message-safety@sha256: BITRIX_SYNC_IMAGE=/han-bitrix-sync@sha256: @@ -8,6 +9,8 @@ NGINX_IMAGE=nginxinc/nginx-unprivileged@sha256: REDIS_IMAGE=redis@sha256: CLAMAV_IMAGE=clamav/clamav@sha256: OTEL_COLLECTOR_IMAGE=otel/opentelemetry-collector-contrib@sha256: +REDIS_EXPORTER_IMAGE=oliver006/redis_exporter@sha256: +NGINX_EXPORTER_IMAGE=nginx/nginx-prometheus-exporter@sha256: PROCESSING_PUBLIC_HOST= PROCESSING_PRIVATE_BIND_ADDRESS= diff --git a/VM2_services/codebase/services/bitrix-sync/.env.example b/VM2_services/codebase/services/bitrix-sync/.env.example index 0939bab..006b427 100644 --- a/VM2_services/codebase/services/bitrix-sync/.env.example +++ b/VM2_services/codebase/services/bitrix-sync/.env.example @@ -15,3 +15,8 @@ BITRIX_SYNC_CONTACT_SOURCE=WEB BITRIX_SYNC_WEBHOOK_ALLOWED_CIDRS=203.0.113.0/24 BITRIX_SYNC_HTTP_TIMEOUT_SEC=10 BITRIX_SYNC_DB_POOL_SIZE=5 +# Optional: without this endpoint the service keeps JSON stdout and runs normally. +OTEL_EXPORTER_OTLP_ENDPOINT= +RELEASE_VERSION=unknown +APP_ENV=production-like +LOG_LEVEL=INFO diff --git a/VM2_services/codebase/services/bitrix-sync/README.md b/VM2_services/codebase/services/bitrix-sync/README.md index d29fa6e..6960a38 100644 --- a/VM2_services/codebase/services/bitrix-sync/README.md +++ b/VM2_services/codebase/services/bitrix-sync/README.md @@ -16,6 +16,22 @@ Disabled mode требует только `BITRIX_SYNC_ENABLED=false` и валидирует весь каталог secret files, portal identity, custom fields, HTTPS host lock и непустой CIDR allow-list до startup. +## Наблюдаемость + +JSON stdout включён всегда. Если задан стандартный +`OTEL_EXPORTER_OTLP_ENDPOINT`, traces, metrics и логи дополнительно отправляются +напрямую по OTLP/gRPC через bounded batch queues. Ошибка инициализации, экспорта +или остановки телеметрии не меняет readiness и не останавливает API, worker или +reconciliation. Процессы различаются как `bitrix-sync-api`, +`bitrix-sync-worker` и `bitrix-sync-reconciliation`; namespace остаётся +`han-chat`. + +FastAPI (кроме `/health/live`), SQLAlchemy и HTTPX инструментированы. Перед +stdout и OTLP выполняется application redaction: query/token/credential URL, +headers, payload, PII, object keys и SQL не экспортируются. Бизнес-метрики +используют только закрытые множества labels; UUID и внешние идентификаторы в +labels/spans не записываются. + ## Границы безопасности - CRM credential URL используется как единый секрет; redirect выключен, TLS diff --git a/VM2_services/codebase/services/bitrix-sync/app/crm.py b/VM2_services/codebase/services/bitrix-sync/app/crm.py index 4bf3f87..f097728 100644 --- a/VM2_services/codebase/services/bitrix-sync/app/crm.py +++ b/VM2_services/codebase/services/bitrix-sync/app/crm.py @@ -5,10 +5,14 @@ import ssl from dataclasses import dataclass from datetime import UTC, datetime from enum import StrEnum +from time import monotonic from typing import Any from urllib.parse import urljoin, urlsplit import httpx +from opentelemetry.trace import SpanKind + +from app.telemetry import record_crm, record_retry, safe_span class CrmOutcome(StrEnum): @@ -49,6 +53,28 @@ class CrmClient: await self._client.aclose() async def call(self, method: str, params: dict[str, Any], *, mutating: bool) -> CrmResult: + started = monotonic() + outcome = "error" + operation = method.rsplit(".", 1)[-1] + try: + with safe_span( + "bitrix_sync.crm.request", + kind=SpanKind.CLIENT, + attributes={"crm.operation": operation}, + ): + result = await self._call(method, params, mutating=mutating) + outcome = result.outcome.value + if result.outcome in { + CrmOutcome.RETRY, + CrmOutcome.UNCERTAIN, + CrmOutcome.RATE_LIMITED, + }: + record_retry("crm") + return result + finally: + record_crm(method, outcome, monotonic() - started) + + async def _call(self, method: str, params: dict[str, Any], *, mutating: bool) -> CrmResult: if method not in ALLOWED_METHODS: raise ValueError("unapproved CRM method") url = urljoin(self._base_url, method + ".json") diff --git a/VM2_services/codebase/services/bitrix-sync/app/engine.py b/VM2_services/codebase/services/bitrix-sync/app/engine.py index ed7f751..fa70777 100644 --- a/VM2_services/codebase/services/bitrix-sync/app/engine.py +++ b/VM2_services/codebase/services/bitrix-sync/app/engine.py @@ -20,6 +20,13 @@ from app.domain import ( validate_phone, ) from app.repository import LeasedTask, LeasedWebhook, Profile, Repository +from app.telemetry import ( + record_business_alert, + record_dead_letter, + record_limiter, + record_retry, + record_transition, +) INSERT_CRM_COMMAND = text( """ @@ -81,6 +88,7 @@ class WorkflowEngine: await self._manual(workflow_id, exc.code, task.user_id) await self.repository.complete_task(task, workflow_id) except RetryableWorkflow as exc: + record_retry("task") delay = exc.retry_after or full_jitter_delay( task.attempt_count + 1, self.settings.retry_base_seconds, @@ -508,10 +516,12 @@ class WorkflowEngine: ) if limiter_delay <= 0: break + record_limiter(limiter_delay) await asyncio.sleep(limiter_delay) async with self._in_flight: result = await self.crm.call(method, params, mutating=mutating) status = result.outcome.value + record_transition("command", status) async with self.repository.transaction() as connection: await connection.execute( UPDATE_CRM_COMMAND, @@ -580,6 +590,7 @@ class WorkflowEngine: *, selected_external_id: str | None = None, ) -> None: + record_business_alert(alert_type) fingerprint = safe_hash(f"{alert_type}:{user_id}") assert fingerprint is not None async with self.repository.transaction() as connection: @@ -727,6 +738,7 @@ class WorkflowEngine: ) async def _technical_failure(self, workflow_id: uuid.UUID, code: str) -> None: + record_dead_letter("crm_command") async with self.repository.transaction() as connection: await connection.execute( text( diff --git a/VM2_services/codebase/services/bitrix-sync/app/main.py b/VM2_services/codebase/services/bitrix-sync/app/main.py index 8134ef9..764ede0 100644 --- a/VM2_services/codebase/services/bitrix-sync/app/main.py +++ b/VM2_services/codebase/services/bitrix-sync/app/main.py @@ -12,20 +12,38 @@ from sqlalchemy import text from app.config import Settings, load_settings from app.repository import Repository from app.security import WebhookValidationError, parse_bounded_form, validate_webhook +from app.telemetry import ( + init_telemetry, + instrument_fastapi, + log_event, + record_webhook, + shutdown_telemetry, +) @asynccontextmanager async def lifespan(app: FastAPI): - settings = load_settings() - app.state.settings = settings - app.state.repository = ( - Repository(settings.database_url.get_secret_value(), settings.db_pool_size) - if settings.enabled and settings.database_url - else None - ) - yield - if app.state.repository: - await app.state.repository.close() + init_telemetry("bitrix-sync-api") + app.state.repository = None + try: + settings = load_settings() + app.state.settings = settings + app.state.repository = ( + Repository(settings.database_url.get_secret_value(), settings.db_pool_size) + if settings.enabled and settings.database_url + else None + ) + log_event( + "service.started", + "Bitrix sync API started", + attributes={"mode": settings.mode}, + ) + yield + finally: + if app.state.repository: + await app.state.repository.close() + log_event("service.stopped", "Bitrix sync API stopped") + shutdown_telemetry() app = FastAPI( @@ -35,6 +53,7 @@ app = FastAPI( redoc_url=None, lifespan=lifespan, ) +instrument_fastapi(app) def settings(request: Request) -> Settings: @@ -107,29 +126,30 @@ async def alert_webhook( async def _receive(request: Request, receiver: str, query: dict[str, str]) -> Response: - config = settings(request) - if not config.enabled: - raise HTTPException(status_code=503, detail="sync_disabled") - if request.headers.get("content-type", "").split(";", 1)[0].lower() != ( - "application/x-www-form-urlencoded" - ): - raise HTTPException(status_code=400, detail="invalid_content_type") - content_length = request.headers.get("content-length") - if content_length and ( - not content_length.isdigit() or int(content_length) > config.webhook_max_body_bytes - ): - raise HTTPException(status_code=413, detail="body_too_large") - body = await request.body() - if len(body) > config.webhook_max_body_bytes: - raise HTTPException(status_code=413, detail="body_too_large") - form = parse_bounded_form(body, max_fields=config.webhook_max_fields) - # The container is reachable only from the trusted VM2 nginx network. - # nginx overwrites X-Real-IP from the TCP peer after its CIDR check. - source_ip = request.headers.get("x-real-ip") or (request.client.host if request.client else "") - alert_entity_type_id = ( - await _alert_entity_type(repository(request)) if receiver == "alert" else None - ) try: + config = settings(request) + if not config.enabled: + raise HTTPException(status_code=503, detail="sync_disabled") + if request.headers.get("content-type", "").split(";", 1)[0].lower() != ( + "application/x-www-form-urlencoded" + ): + raise HTTPException(status_code=400, detail="invalid_content_type") + content_length = request.headers.get("content-length") + if content_length and ( + not content_length.isdigit() or int(content_length) > config.webhook_max_body_bytes + ): + raise HTTPException(status_code=413, detail="body_too_large") + body = await request.body() + if len(body) > config.webhook_max_body_bytes: + raise HTTPException(status_code=413, detail="body_too_large") + form = parse_bounded_form(body, max_fields=config.webhook_max_fields) + # nginx overwrites X-Real-IP from the TCP peer after its CIDR check. + source_ip = request.headers.get("x-real-ip") or ( + request.client.host if request.client else "" + ) + alert_entity_type_id = ( + await _alert_entity_type(repository(request)) if receiver == "alert" else None + ) event = validate_webhook( receiver, query, @@ -139,12 +159,21 @@ async def _receive(request: Request, receiver: str, query: dict[str, str]) -> Re alert_entity_type_id=alert_entity_type_id, ) except PermissionError as exc: + record_webhook(receiver, "rejected") raise HTTPException(status_code=403, detail="forbidden") from exc except WebhookValidationError as exc: + record_webhook(receiver, "rejected") raise HTTPException(status_code=400, detail="malformed_webhook") from exc + except HTTPException: + record_webhook(receiver, "rejected") + raise + except Exception: + record_webhook(receiver, "error") + raise await repository(request).insert_webhook( event.receiver_type, event.event_type, event.entity_id, event.source_ip ) + record_webhook(receiver, "accepted") return Response(status_code=202) @@ -170,4 +199,5 @@ def run() -> None: host="0.0.0.0", # noqa: S104 - container-only port, not host-published port=8080, proxy_headers=False, + access_log=False, ) diff --git a/VM2_services/codebase/services/bitrix-sync/app/reconciliation.py b/VM2_services/codebase/services/bitrix-sync/app/reconciliation.py index 2055cea..50ca283 100644 --- a/VM2_services/codebase/services/bitrix-sync/app/reconciliation.py +++ b/VM2_services/codebase/services/bitrix-sync/app/reconciliation.py @@ -2,6 +2,7 @@ from __future__ import annotations import asyncio from datetime import UTC, datetime, timedelta +from time import monotonic from typing import Any from sqlalchemy import text @@ -9,6 +10,13 @@ from sqlalchemy import text from app.config import load_settings from app.crm import CrmClient, CrmOutcome from app.repository import Repository +from app.telemetry import ( + init_telemetry, + log_event, + record_reconciliation, + safe_span, + shutdown_telemetry, +) class IncrementalReconciler: @@ -127,25 +135,44 @@ class IncrementalReconciler: async def reconciliation_main() -> None: - settings = load_settings() - if not settings.enabled: - return - assert settings.database_url and settings.crm_rest_webhook_url and settings.portal_host - assert settings.contact_registered_field - repository = Repository(settings.database_url.get_secret_value(), settings.db_pool_size) - crm = CrmClient( - settings.crm_rest_webhook_url.get_secret_value(), - settings.portal_host, - settings.http_timeout_sec, - ) - reconciler = IncrementalReconciler( - repository, crm, settings.rest_field_name(settings.contact_registered_field) - ) + init_telemetry("bitrix-sync-reconciliation") + repository: Repository | None = None + crm: CrmClient | None = None + started = monotonic() + outcome = "error" + count = 0 try: - await reconciler.run_once(settings.reconciliation_overlap_seconds) + settings = load_settings() + if not settings.enabled: + outcome = "skipped" + log_event("reconciliation.disabled", "Bitrix sync reconciliation is disabled") + return + assert settings.database_url and settings.crm_rest_webhook_url and settings.portal_host + assert settings.contact_registered_field + repository = Repository(settings.database_url.get_secret_value(), settings.db_pool_size) + crm = CrmClient( + settings.crm_rest_webhook_url.get_secret_value(), + settings.portal_host, + settings.http_timeout_sec, + ) + reconciler = IncrementalReconciler( + repository, crm, settings.rest_field_name(settings.contact_registered_field) + ) + with safe_span("bitrix_sync.reconciliation.run"): + count = await reconciler.run_once(settings.reconciliation_overlap_seconds) + outcome = "success" + log_event( + "reconciliation.completed", + "Bitrix sync reconciliation completed", + attributes={"scanned_count": count}, + ) finally: - await crm.close() - await repository.close() + record_reconciliation(outcome, count, started) + if crm is not None: + await crm.close() + if repository is not None: + await repository.close() + shutdown_telemetry() def run() -> None: diff --git a/VM2_services/codebase/services/bitrix-sync/app/repository.py b/VM2_services/codebase/services/bitrix-sync/app/repository.py index a7e8ea8..d08d025 100644 --- a/VM2_services/codebase/services/bitrix-sync/app/repository.py +++ b/VM2_services/codebase/services/bitrix-sync/app/repository.py @@ -13,6 +13,7 @@ from sqlalchemy import text from sqlalchemy.ext.asyncio import AsyncConnection, AsyncEngine, create_async_engine from app.domain import safe_hash +from app.telemetry import record_status_snapshot def postgres_ssl_context() -> ssl.SSLContext: @@ -541,6 +542,10 @@ class Repository: "SELECT status, count(*) count " "FROM bitrix_sync.crm_commands GROUP BY status" ), + "queue_oldest_age_seconds": ( + "SELECT coalesce(extract(epoch from now()-min(created_at)),0) " + "FROM han_app.sync_queue WHERE status IN ('pending','retry_wait','leased')" + ), "webhook_lag_seconds": ( "SELECT coalesce(extract(epoch from now()-min(received_at)),0) " "FROM bitrix_sync.webhook_inbox WHERE status IN ('received','retry_wait')" @@ -559,4 +564,5 @@ class Repository: else: output[key] = result.scalar_one_or_none() output["generated_at"] = datetime.now(UTC).isoformat() + record_status_snapshot(output) return output diff --git a/VM2_services/codebase/services/bitrix-sync/app/telemetry.py b/VM2_services/codebase/services/bitrix-sync/app/telemetry.py new file mode 100644 index 0000000..dc9cf13 --- /dev/null +++ b/VM2_services/codebase/services/bitrix-sync/app/telemetry.py @@ -0,0 +1,514 @@ +from __future__ import annotations + +import json +import logging +import os +import re +import sys +from collections.abc import Iterator, Mapping +from contextlib import contextmanager +from dataclasses import dataclass +from datetime import UTC, datetime +from time import monotonic +from typing import Any + +from fastapi import FastAPI +from opentelemetry import metrics, trace +from opentelemetry.exporter.otlp.proto.grpc._log_exporter import OTLPLogExporter +from opentelemetry.exporter.otlp.proto.grpc.metric_exporter import OTLPMetricExporter +from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporter +from opentelemetry.instrumentation.fastapi import FastAPIInstrumentor +from opentelemetry.instrumentation.httpx import HTTPXClientInstrumentor +from opentelemetry.instrumentation.sqlalchemy import SQLAlchemyInstrumentor +from opentelemetry.propagate import set_global_textmap +from opentelemetry.sdk._logs import LoggerProvider, LoggingHandler +from opentelemetry.sdk._logs.export import BatchLogRecordProcessor +from opentelemetry.sdk.metrics import MeterProvider +from opentelemetry.sdk.metrics.export import PeriodicExportingMetricReader +from opentelemetry.sdk.resources import Resource +from opentelemetry.sdk.trace import TracerProvider +from opentelemetry.sdk.trace.export import BatchSpanProcessor +from opentelemetry.trace import SpanKind, Status, StatusCode +from opentelemetry.trace.propagation.tracecontext import TraceContextTextMapPropagator + +SERVICE_NAMES = frozenset( + {"bitrix-sync-api", "bitrix-sync-worker", "bitrix-sync-reconciliation"} +) +_SECRET_KEYS = re.compile( + r"(authorization|cookie|token|secret|password|phone|email|full.?name|" + r"message|payload|body|query|url|dsn|statement|object.?key|filename)", + re.IGNORECASE, +) +_SENSITIVE_VALUES = ( + re.compile(r"(?i)\bbearer\s+\S+"), + re.compile(r"(?i)\b(?:token|password|secret|code)=\S+"), + re.compile(r"https?://[^\s?#]+[?#]\S+"), + re.compile(r"\b[A-Za-z0-9_-]{20,}\.[A-Za-z0-9_-]{20,}\.[A-Za-z0-9_-]{10,}\b"), + re.compile(r"\b[\w.+-]+@[\w.-]+\.[A-Za-z]{2,}\b"), + re.compile(r"(? Any: + if value is None or isinstance(value, (bool, int, float)): + return value + text = str(value) + for pattern in _SENSITIVE_VALUES: + text = pattern.sub("[REDACTED]", text) + return text[:512] + + +def redact(value: Any, *, key: str = "") -> Any: + """Return a bounded telemetry-safe copy without mutating application data.""" + if key and _SECRET_KEYS.search(key): + return "[REDACTED]" + if isinstance(value, Mapping): + return { + str(item_key)[:64]: redact(item, key=str(item_key)) + for item_key, item in value.items() + } + if isinstance(value, (list, tuple, set, frozenset)): + return [redact(item) for item in list(value)[:32]] + return _safe_scalar(value) + + +def _trace_fields() -> dict[str, str]: + context = trace.get_current_span().get_span_context() + if not context.is_valid: + return {} + return { + "trace_id": format(context.trace_id, "032x"), + "span_id": format(context.span_id, "016x"), + } + + +class JsonFormatter(logging.Formatter): + def format(self, record: logging.LogRecord) -> str: + event = { + "timestamp": ( + datetime.now(UTC).isoformat(timespec="milliseconds").replace("+00:00", "Z") + ), + "level": record.levelname, + "service.name": _configured_service_name, + "service.version": os.getenv("RELEASE_VERSION", "unknown"), + "environment": os.getenv("APP_ENV", "production-like"), + "module": record.name, + "event": getattr( + record, + "event_name", + "log", + ), + "message": record.getMessage(), + **_trace_fields(), + } + attributes = getattr(record, "telemetry_attributes", None) + if isinstance(attributes, Mapping): + event.update(attributes) + if record.exc_info: + event["error.type"] = record.exc_info[0].__name__ + return json.dumps(redact(event), ensure_ascii=False, separators=(",", ":"), default=str) + + +class RedactionFilter(logging.Filter): + def filter(self, record: logging.LogRecord) -> bool: + record.msg = redact(record.getMessage()) if hasattr(record, "event_name") else "[REDACTED]" + record.args = () + attributes = getattr(record, "telemetry_attributes", None) + safe_attributes = dict(attributes) if isinstance(attributes, Mapping) else {} + if record.exc_info: + safe_attributes["error.type"] = record.exc_info[0].__name__ + record.exc_info = None + record.exc_text = None + record.telemetry_attributes = redact(safe_attributes) + return True + + +class RedactingBatchSpanProcessor(BatchSpanProcessor): + """Remove sensitive auto-instrumentation attributes before queueing/export.""" + + def on_end(self, span: Any) -> None: + attributes = getattr(span, "_attributes", None) + if isinstance(attributes, Mapping): + sanitized = dict(attributes) + for key in list(sanitized): + key_text = str(key) + if _SECRET_KEYS.search(key_text) or key_text in { + "http.target", + "url.full", + "db.statement", + }: + sanitized[key] = "[REDACTED]" + else: + sanitized[key] = _safe_scalar(sanitized[key]) + span._attributes = sanitized + super().on_end(span) + + +@dataclass(slots=True) +class TelemetryRuntime: + tracer_provider: TracerProvider | None = None + meter_provider: MeterProvider | None = None + logger_provider: LoggerProvider | None = None + + def shutdown(self) -> None: + for provider in (self.logger_provider, self.meter_provider, self.tracer_provider): + if provider is None: + continue + try: + provider.shutdown() + except Exception: + logging.getLogger(__name__).warning( + "Telemetry shutdown failed", + extra={"event_name": "telemetry.shutdown.failed"}, + ) + + +def _resource(service_name: str) -> Resource: + return Resource.create( + { + "service.name": service_name, + "service.namespace": "han-chat", + "service.version": os.getenv("RELEASE_VERSION", "unknown"), + "deployment.environment": os.getenv("APP_ENV", "production-like"), + } + ) + + +def _configure_stdout(service_name: str) -> None: + global _configured_service_name, _logging_configured + _configured_service_name = service_name + if _logging_configured: + return + handler = logging.StreamHandler(sys.stdout) + handler.setFormatter(JsonFormatter()) + handler.addFilter(RedactionFilter()) + root = logging.getLogger() + root.addHandler(handler) + level_name = os.getenv("LOG_LEVEL", "INFO").upper() + root.setLevel(getattr(logging, level_name, logging.INFO)) + _logging_configured = True + + +def init_telemetry(service_name: str) -> TelemetryRuntime: + """Initialize all signals; exporter failures never stop the business process.""" + global _runtime + if service_name not in SERVICE_NAMES: + raise ValueError("unregistered bitrix-sync process service.name") + _configure_stdout(service_name) + if _runtime is not None: + return _runtime + runtime = TelemetryRuntime() + _runtime = runtime + endpoint = os.getenv("OTEL_EXPORTER_OTLP_ENDPOINT", "").strip() + if not endpoint: + return runtime + + try: + resource = _resource(service_name) + insecure = endpoint.startswith("http://") + set_global_textmap(TraceContextTextMapPropagator()) + + tracer_provider = TracerProvider(resource=resource) + tracer_provider.add_span_processor( + RedactingBatchSpanProcessor( + OTLPSpanExporter(endpoint=endpoint, insecure=insecure, timeout=3), + max_queue_size=2048, + schedule_delay_millis=5000, + max_export_batch_size=512, + export_timeout_millis=3000, + ) + ) + trace.set_tracer_provider(tracer_provider) + runtime.tracer_provider = tracer_provider + + metric_reader = PeriodicExportingMetricReader( + OTLPMetricExporter(endpoint=endpoint, insecure=insecure, timeout=3), + export_interval_millis=30_000, + export_timeout_millis=3000, + ) + meter_provider = MeterProvider(resource=resource, metric_readers=[metric_reader]) + metrics.set_meter_provider(meter_provider) + runtime.meter_provider = meter_provider + + logger_provider = LoggerProvider(resource=resource) + logger_provider.add_log_record_processor( + BatchLogRecordProcessor( + OTLPLogExporter(endpoint=endpoint, insecure=insecure, timeout=3), + max_queue_size=2048, + schedule_delay_millis=5000, + max_export_batch_size=512, + export_timeout_millis=3000, + ) + ) + otlp_handler = LoggingHandler(level=logging.NOTSET, logger_provider=logger_provider) + otlp_handler.addFilter(RedactionFilter()) + logging.getLogger().addHandler(otlp_handler) + runtime.logger_provider = logger_provider + + def sanitize_httpx_request(span: trace.Span, request: Any) -> None: + if not span.is_recording(): + return + url = request[1] + host = getattr(url, "host", "") + scheme = getattr(url, "scheme", "https") + safe_url = f"{scheme}://{host}/[REDACTED]" + span.set_attribute("url.full", safe_url) + span.set_attribute("http.url", safe_url) + + HTTPXClientInstrumentor().instrument(request_hook=sanitize_httpx_request) + SQLAlchemyInstrumentor().instrument(enable_commenter=False) + except Exception: + logging.getLogger(__name__).exception( + "Telemetry initialization failed; stdout fallback remains active", + extra={"event_name": "telemetry.init.failed"}, + ) + return runtime + + +def instrument_fastapi(app: FastAPI) -> None: + try: + FastAPIInstrumentor.instrument_app( + app, + excluded_urls="/health/live", + http_capture_headers_server_request=[], + http_capture_headers_server_response=[], + ) + except Exception: + logging.getLogger(__name__).exception( + "FastAPI instrumentation failed", + extra={"event_name": "telemetry.fastapi.failed"}, + ) + + +def shutdown_telemetry() -> None: + global _runtime + if _runtime is not None: + _runtime.shutdown() + _runtime = None + + +def log_event( + event: str, + message: str, + *, + level: int = logging.INFO, + attributes: Mapping[str, Any] | None = None, +) -> None: + logging.getLogger("bitrix_sync").log( + level, + message, + extra={"event_name": event, "telemetry_attributes": redact(attributes or {})}, + ) + + +@contextmanager +def safe_span( + name: str, + *, + kind: SpanKind = SpanKind.INTERNAL, + attributes: Mapping[str, Any] | None = None, +) -> Iterator[trace.Span]: + safe_attributes: dict[str, bool | int | float | str] = {} + for key, value in redact(attributes or {}).items(): + if not isinstance(value, (bool, int, float, str)): + continue + key_text = str(key)[:64] + if key_text == "workflow.type": + value = _bounded(str(value), _SAFE_TASK_TYPES) + elif key_text == "crm.operation": + value = _bounded(str(value), _SAFE_CRM_OPERATIONS) + elif key_text == "receiver": + value = _bounded(str(value), _SAFE_RECEIVERS, "unknown") + safe_attributes[key_text] = value + with trace.get_tracer("han.bitrix_sync").start_as_current_span( + name[:128], kind=kind, attributes=safe_attributes + ) as span: + try: + yield span + except Exception as exc: + span.set_status(Status(StatusCode.ERROR, type(exc).__name__)) + raise + + +def _bounded(value: str, allowed: frozenset[str], fallback: str = "other") -> str: + return value if value in allowed else fallback + + +def _metric_fail_open(function: Any) -> Any: + def wrapped(*args: Any, **kwargs: Any) -> None: + try: + function(*args, **kwargs) + except Exception: + return + + return wrapped + + +@_metric_fail_open +def record_claim(queue: str, count: int) -> None: + queue_name = _bounded(queue, frozenset({"task", "webhook", "rebind"})) + metrics.get_meter("han.bitrix_sync").create_counter( + "bitrix_sync_claimed_items", unit="{item}" + ).add(max(0, count), {"queue": queue_name}) + + +@_metric_fail_open +def record_workflow(task_type: str, outcome: str, duration_seconds: float) -> None: + attrs = { + "workflow.type": _bounded(task_type, _SAFE_TASK_TYPES), + "outcome": _bounded(outcome, _SAFE_CODES), + } + meter = metrics.get_meter("han.bitrix_sync") + meter.create_counter("bitrix_sync_workflows", unit="{workflow}").add(1, attrs) + meter.create_histogram("bitrix_sync_workflow_duration", unit="s").record( + max(0.0, duration_seconds), attrs + ) + + +@_metric_fail_open +def record_webhook(receiver: str, outcome: str) -> None: + metrics.get_meter("han.bitrix_sync").create_counter( + "bitrix_sync_webhooks", unit="{webhook}" + ).add( + 1, + { + "receiver": _bounded(receiver, _SAFE_RECEIVERS, "unknown"), + "outcome": _bounded(outcome, frozenset({"accepted", "rejected", "error"})), + }, + ) + + +@_metric_fail_open +def record_crm(method: str, outcome: str, duration_seconds: float) -> None: + operation = method.rsplit(".", 1)[-1] + operation = _bounded(operation, _SAFE_CRM_OPERATIONS) + attrs = {"operation": operation, "outcome": _bounded(outcome, _SAFE_CODES)} + meter = metrics.get_meter("han.bitrix_sync") + meter.create_counter("bitrix_sync_crm_calls", unit="{call}").add(1, attrs) + meter.create_histogram("bitrix_sync_crm_duration", unit="s").record( + max(0.0, duration_seconds), attrs + ) + + +@_metric_fail_open +def record_limiter(delay_seconds: float) -> None: + metrics.get_meter("han.bitrix_sync").create_histogram( + "bitrix_sync_limiter_delay", unit="s" + ).record(max(0.0, delay_seconds)) + + +@_metric_fail_open +def record_transition(kind: str, state: str) -> None: + metrics.get_meter("han.bitrix_sync").create_counter( + "bitrix_sync_transitions", unit="{transition}" + ).add( + 1, + { + "kind": _bounded(kind, frozenset({"workflow", "command", "webhook"})), + "state": _bounded( + state, + frozenset( + { + "succeeded", + "retry", + "uncertain", + "permanent", + "rate_limited", + "waiting_manual", + "processed", + "error", + } + ), + ), + }, + ) + + +@_metric_fail_open +def record_retry(kind: str) -> None: + metrics.get_meter("han.bitrix_sync").create_counter( + "bitrix_sync_retries", unit="{retry}" + ).add(1, {"kind": _bounded(kind, frozenset({"task", "webhook", "rebind", "crm"}))}) + + +@_metric_fail_open +def record_dead_letter(operation: str) -> None: + metrics.get_meter("han.bitrix_sync").create_counter( + "bitrix_sync_dead_letters", unit="{item}" + ).add(1, {"operation": _bounded(operation, frozenset({"crm_command"}))}) + + +@_metric_fail_open +def record_business_alert(alert_type: str) -> None: + metrics.get_meter("han.bitrix_sync").create_counter( + "bitrix_sync_business_alerts", unit="{alert}" + ).add(1, {"alert.type": _bounded(alert_type, _SAFE_ALERT_TYPES)}) + + +@_metric_fail_open +def record_status_snapshot(status: Mapping[str, Any]) -> None: + meter = metrics.get_meter("han.bitrix_sync") + depth = meter.create_histogram("bitrix_sync_queue_depth", unit="{item}") + for queue_name in ("queue", "workflows", "commands"): + values = status.get(queue_name) + if not isinstance(values, Mapping): + continue + for state, count in values.items(): + if isinstance(count, int): + depth.record( + max(0, count), + { + "queue": queue_name, + "state": _bounded(str(state), _SAFE_QUEUE_STATES), + }, + ) + for key in ("queue_oldest_age_seconds", "webhook_lag_seconds"): + value = status.get(key) + if isinstance(value, (int, float)): + meter.create_histogram(f"bitrix_sync_{key}", unit="s").record(max(0.0, value)) + + +@_metric_fail_open +def record_reconciliation(outcome: str, count: int, started: float) -> None: + attrs = {"outcome": _bounded(outcome, frozenset({"success", "error", "skipped"}))} + meter = metrics.get_meter("han.bitrix_sync") + meter.create_counter("bitrix_sync_reconciliation_runs", unit="{run}").add(1, attrs) + meter.create_histogram("bitrix_sync_reconciliation_duration", unit="s").record( + max(0.0, monotonic() - started), attrs + ) + meter.create_histogram("bitrix_sync_reconciliation_items", unit="{item}").record( + max(0, count), attrs + ) diff --git a/VM2_services/codebase/services/bitrix-sync/app/worker.py b/VM2_services/codebase/services/bitrix-sync/app/worker.py index 22e0527..0f250ea 100644 --- a/VM2_services/codebase/services/bitrix-sync/app/worker.py +++ b/VM2_services/codebase/services/bitrix-sync/app/worker.py @@ -4,43 +4,62 @@ import asyncio import signal import socket import uuid +from time import monotonic from app.config import load_settings from app.crm import CrmClient from app.domain import full_jitter_delay from app.engine import RetryableWorkflow, WorkflowEngine from app.repository import Repository +from app.telemetry import ( + init_telemetry, + log_event, + record_claim, + record_retry, + record_workflow, + safe_span, + shutdown_telemetry, +) async def worker_main() -> None: - settings = load_settings() - if not settings.enabled: - return - assert settings.database_url and settings.crm_rest_webhook_url and settings.portal_host - repository = Repository(settings.database_url.get_secret_value(), settings.db_pool_size) - crm = CrmClient( - settings.crm_rest_webhook_url.get_secret_value(), - settings.portal_host, - settings.http_timeout_sec, - ) - engine = WorkflowEngine(repository, crm, settings) - stop = asyncio.Event() - loop = asyncio.get_running_loop() - for event in (signal.SIGINT, signal.SIGTERM): - try: - loop.add_signal_handler(event, stop.set) - except NotImplementedError: - pass - worker_id = f"{socket.gethostname()}:{uuid.uuid4()}" + init_telemetry("bitrix-sync-worker") + repository: Repository | None = None + crm: CrmClient | None = None try: + settings = load_settings() + if not settings.enabled: + log_event("worker.disabled", "Bitrix sync worker is disabled") + return + assert settings.database_url and settings.crm_rest_webhook_url and settings.portal_host + repository = Repository(settings.database_url.get_secret_value(), settings.db_pool_size) + crm = CrmClient( + settings.crm_rest_webhook_url.get_secret_value(), + settings.portal_host, + settings.http_timeout_sec, + ) + engine = WorkflowEngine(repository, crm, settings) + stop = asyncio.Event() + loop = asyncio.get_running_loop() + for event in (signal.SIGINT, signal.SIGTERM): + try: + loop.add_signal_handler(event, stop.set) + except NotImplementedError: + pass + worker_id = f"{socket.gethostname()}:{uuid.uuid4()}" + log_event("worker.started", "Bitrix sync worker started") while not stop.is_set(): - tasks = await repository.claim_tasks( - worker_id, settings.claim_size, settings.lease_seconds - ) - webhooks = await repository.claim_webhooks( - worker_id, settings.claim_size, settings.lease_seconds - ) - rebind_ids = await repository.pending_rebind_ids(settings.claim_size) + with safe_span("bitrix_sync.worker.claim"): + tasks = await repository.claim_tasks( + worker_id, settings.claim_size, settings.lease_seconds + ) + webhooks = await repository.claim_webhooks( + worker_id, settings.claim_size, settings.lease_seconds + ) + rebind_ids = await repository.pending_rebind_ids(settings.claim_size) + record_claim("task", len(tasks)) + record_claim("webhook", len(webhooks)) + record_claim("rebind", len(rebind_ids)) if not tasks and not webhooks and not rebind_ids: try: await asyncio.wait_for(stop.wait(), timeout=1) @@ -49,29 +68,68 @@ async def worker_main() -> None: for task in tasks: if stop.is_set(): break - await engine.process(task) + started = monotonic() + outcome = "success" + try: + with safe_span( + "bitrix_sync.worker.process", + attributes={"workflow.type": task.task_type}, + ): + await engine.process(task) + except Exception: + outcome = "error" + raise + finally: + record_workflow(task.task_type, outcome, monotonic() - started) for item in webhooks: if stop.is_set(): break + started = monotonic() + outcome = "success" try: - await engine.process_webhook(item) + with safe_span( + "bitrix_sync.worker.webhook", + attributes={"receiver": item.receiver_type}, + ): + await engine.process_webhook(item) except RetryableWorkflow as exc: + outcome = "retry" + record_retry("webhook") delay = exc.retry_after or full_jitter_delay( item.attempt_count + 1, settings.retry_base_seconds, settings.retry_max_seconds, ) await repository.retry_webhook(item, exc.code, delay) + except Exception: + outcome = "error" + raise + finally: + record_workflow("contact.webhook", outcome, monotonic() - started) for request_id in rebind_ids: if stop.is_set(): break + started = monotonic() + outcome = "success" try: - await engine.process_rebind(request_id) + with safe_span("bitrix_sync.worker.rebind"): + await engine.process_rebind(request_id) except RetryableWorkflow: + outcome = "retry" + record_retry("rebind") continue + except Exception: + outcome = "error" + raise + finally: + record_workflow("rebind", outcome, monotonic() - started) finally: - await crm.close() - await repository.close() + if crm is not None: + await crm.close() + if repository is not None: + await repository.close() + log_event("worker.stopped", "Bitrix sync worker stopped") + shutdown_telemetry() def run() -> None: diff --git a/VM2_services/codebase/services/bitrix-sync/pyproject.toml b/VM2_services/codebase/services/bitrix-sync/pyproject.toml index ccc7171..ba75cd4 100644 --- a/VM2_services/codebase/services/bitrix-sync/pyproject.toml +++ b/VM2_services/codebase/services/bitrix-sync/pyproject.toml @@ -8,6 +8,12 @@ dependencies = [ "asyncpg>=0.30,<1", "fastapi>=0.116,<1", "httpx>=0.28,<1", + "opentelemetry-api>=1.36,<2", + "opentelemetry-exporter-otlp-proto-grpc>=1.36,<2", + "opentelemetry-instrumentation-fastapi>=0.57b0,<1", + "opentelemetry-instrumentation-httpx>=0.57b0,<1", + "opentelemetry-instrumentation-sqlalchemy>=0.57b0,<1", + "opentelemetry-sdk>=1.36,<2", "pydantic-settings>=2.10,<3", "python-multipart>=0.0.20,<1", "sqlalchemy[asyncio]>=2.0.41,<3", diff --git a/VM2_services/codebase/services/bitrix-sync/tests/test_telemetry.py b/VM2_services/codebase/services/bitrix-sync/tests/test_telemetry.py new file mode 100644 index 0000000..34ccbaa --- /dev/null +++ b/VM2_services/codebase/services/bitrix-sync/tests/test_telemetry.py @@ -0,0 +1,172 @@ +from __future__ import annotations + +import json +import logging +from typing import Any + +import pytest +from opentelemetry.sdk.trace import ReadableSpan +from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter + +from app import telemetry + + +@pytest.mark.parametrize( + "canary", + [ + "Bearer canary-authorization-value", + "person@example.test", + "+7 (999) 123-45-67", + "https://portal.example/path?token=canary", + "aaaaaaaaaaaaaaaaaaaa.bbbbbbbbbbbbbbbbbbbb.cccccccccc", + ], +) +def test_redaction_removes_canaries_from_json_stdout(canary: str) -> None: + record = logging.LogRecord( + "test", + logging.INFO, + __file__, + 1, + "safe event %s", + (canary,), + None, + ) + record.event_name = "test.canary" + record.telemetry_attributes = { + "authorization": canary, + "nested": {"email": canary, "safe": canary}, + } + + telemetry.RedactionFilter().filter(record) + output = telemetry.JsonFormatter().format(record) + + assert canary not in output + assert json.loads(output)["event"] == "test.canary" + + +def test_redaction_blocks_sensitive_keys_and_bounds_collections() -> None: + result = telemetry.redact( + { + "db.statement": "SELECT secret FROM users", + "object_key": "private/file.txt", + "safe": list(range(100)), + } + ) + + assert result["db.statement"] == "[REDACTED]" + assert result["object_key"] == "[REDACTED]" + assert len(result["safe"]) == 32 + + +def test_redaction_drops_exception_text_before_otlp() -> None: + try: + raise RuntimeError("secret-token-in-exception") + except RuntimeError: + record = logging.LogRecord( + "test", + logging.ERROR, + __file__, + 1, + "operation failed", + (), + exc_info=__import__("sys").exc_info(), + ) + + telemetry.RedactionFilter().filter(record) + + assert record.exc_info is None + assert record.telemetry_attributes["error.type"] == "RuntimeError" + assert "secret-token-in-exception" not in telemetry.JsonFormatter().format(record) + + +def test_span_processor_replaces_immutable_sdk_attributes() -> None: + exporter = InMemorySpanExporter() + processor = telemetry.RedactingBatchSpanProcessor(exporter) + span = ReadableSpan( + name="crm.request", + attributes={"url.full": "https://crm.example/?token=secret", "safe.outcome": "success"}, + ) + + processor.on_end(span) + processor.shutdown() + + assert span.attributes["url.full"] == "[REDACTED]" + assert span.attributes["safe.outcome"] == "success" + + +def test_init_without_endpoint_is_backward_compatible(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.delenv("OTEL_EXPORTER_OTLP_ENDPOINT", raising=False) + monkeypatch.setattr(telemetry, "_runtime", None) + + runtime = telemetry.init_telemetry("bitrix-sync-worker") + + assert runtime.tracer_provider is None + assert runtime.meter_provider is None + assert runtime.logger_provider is None + telemetry.shutdown_telemetry() + + +def test_each_process_has_distinct_service_resource() -> None: + names = { + telemetry._resource(name).attributes["service.name"] # noqa: SLF001 + for name in telemetry.SERVICE_NAMES + } + + assert names == { + "bitrix-sync-api", + "bitrix-sync-worker", + "bitrix-sync-reconciliation", + } + + +def test_init_is_fail_open_when_exporter_construction_fails( + monkeypatch: pytest.MonkeyPatch, +) -> None: + monkeypatch.setenv("OTEL_EXPORTER_OTLP_ENDPOINT", "http://collector:4317") + monkeypatch.setattr(telemetry, "_runtime", None) + + def fail_exporter(**_kwargs: Any) -> None: + raise RuntimeError("collector unavailable") + + monkeypatch.setattr(telemetry, "OTLPSpanExporter", fail_exporter) + + runtime = telemetry.init_telemetry("bitrix-sync-reconciliation") + + assert runtime.tracer_provider is None + telemetry.shutdown_telemetry() + + +def test_business_metric_attributes_are_bounded(monkeypatch: pytest.MonkeyPatch) -> None: + captured: list[dict[str, str]] = [] + + class Instrument: + def add(self, _value: int, attributes: dict[str, str]) -> None: + captured.append(attributes) + + def record(self, _value: float, attributes: dict[str, str]) -> None: + captured.append(attributes) + + class Meter: + def create_counter(self, *_args: Any, **_kwargs: Any) -> Instrument: + return Instrument() + + def create_histogram(self, *_args: Any, **_kwargs: Any) -> Instrument: + return Instrument() + + monkeypatch.setattr(telemetry.metrics, "get_meter", lambda _name: Meter()) + + telemetry.record_workflow("user-2f3c0a5e-identifier", "novel-state", 0.5) + telemetry.record_crm("unknown.dynamic.method", "novel-state", 0.1) + + assert captured + assert all("user-2f3c0a5e-identifier" not in values.values() for values in captured) + assert all("novel-state" not in values.values() for values in captured) + + +def test_business_metrics_are_fail_open(monkeypatch: pytest.MonkeyPatch) -> None: + def unavailable(_name: str) -> None: + raise RuntimeError("metrics unavailable") + + monkeypatch.setattr(telemetry.metrics, "get_meter", unavailable) + + telemetry.record_claim("task", 1) diff --git a/VM2_services/codebase/services/deployment/RUNBOOK.md b/VM2_services/codebase/services/deployment/RUNBOOK.md index 01b5d52..9b7837d 100644 --- a/VM2_services/codebase/services/deployment/RUNBOOK.md +++ b/VM2_services/codebase/services/deployment/RUNBOOK.md @@ -34,8 +34,11 @@ The full step-by-step procedure with gates and copy-paste commands is in expose unauthenticated `PING` only for health and a password-protected `safety` user limited to required `han:safety:*` keys/commands. The password in `MESSAGE_SAFETY_REDIS_URL` must match. Start from - `redis/redis-safety.acl.template`, replace - `REPLACE_WITH_LONG_RANDOM_PASSWORD`, and never commit the password. + `redis/redis-safety.acl.template`, replace both placeholders, and add the + `exporter` user limited to `PING`/`INFO`. Its ACL password must exactly match + the separate `REDIS_EXPORTER_PASSWORD` secret. That secret is a JSON password + map, `{"redis://redis-safety:6379":""}`, not a raw password + string. Never commit either password. 8. Provision distinct runtime and migration DB credentials. `MESSAGE_SAFETY_CONFIG_ADMIN_DATABASE_URL` may migrate/activate policy while `MESSAGE_SAFETY_DATABASE_URL` cannot; `BITRIX_SYNC_MIGRATION_DATABASE_URL` @@ -76,13 +79,15 @@ Compose first: [`arch-10-deployment.md`](../../../../architectory/arch-10-deployment.md) §6. 3. **Images** — build and push `han-message-safety`, `han-bitrix-sync`; record immutable digests for every `*_IMAGE` in `.env.example` (nginx, redis, clamav, - otel-collector). + otel-collector, Redis exporter and nginx exporter). 4. **Selectel Secrets Manager** — populate all remote names from `deployment/secrets/config.example.json` (DSNs, tokens, S3 read-only keys, - `REDIS_SAFETY_ACL`, internal TLS PEM for `8443`). Dedicated VM2 IAM principal - with read-only access to those names only. + `REDIS_SAFETY_ACL`, `REDIS_EXPORTER_PASSWORD`, internal TLS PEM for `8443`). + Dedicated VM2 IAM principal with read-only access to those names only. 5. **S3 quarantine bucket** and SigNoz OTLP endpoint — non-secret values in - `.env`. + `.env`. Current self-hosted SigNoz accepts private plaintext OTLP without + authentication; do not provision a fake auth secret or mandatory non-empty + header. 6. **Internal TLS** — internal-CA certificate with SAN = VM2 private DNS; PEM stored in Secrets Manager, not in the release tree. @@ -136,6 +141,11 @@ Copy `.env.example` → `.env`, install loader config as password with `systemd-creds`, edit nginx allow-lists. Details: [`RUNBOOK.ru.md`](RUNBOOK.ru.md) §4. +The local Collector receives logs directly over OTLP (no `filelog`), scrapes +only itself plus the Redis/nginx exporters, and collects host metrics through +read-only `/hostfs`. Exporters and nginx `stub_status` use internal networks and +`expose` only; they have no host-published ports. + ## 5. PostgreSQL CA and initial public TLS Install managed PostgreSQL CA under `/etc/han/ca`, issue Let's Encrypt cert for diff --git a/VM2_services/codebase/services/deployment/RUNBOOK.ru.md b/VM2_services/codebase/services/deployment/RUNBOOK.ru.md index 21cb522..db3ec89 100644 --- a/VM2_services/codebase/services/deployment/RUNBOOK.ru.md +++ b/VM2_services/codebase/services/deployment/RUNBOOK.ru.md @@ -37,8 +37,12 @@ deployment-артефакты: `VM2_services/codebase/services/`. Локальн 7. `REDIS_SAFETY_ACL` — полный ACL-файл, а не просто пароль. Он должен открывать неаутентифицированный `PING` только для health и защищённого паролем пользователя `safety`, ограниченного необходимыми - ключами/командами `han:safety:*`. - (Пароль в `MESSAGE_SAFETY_REDIS_URL` должен совпадать. Используйте `redis/redis-safety.acl.template`, заменив `REPLACE_WITH_LONG_RANDOM_PASSWORD) + ключами/командами `han:safety:*`. Добавьте пользователя `exporter` только с + `PING`/`INFO`; его пароль в ACL должен в точности совпадать с отдельным + `REDIS_EXPORTER_PASSWORD`. Этот secret хранится в формате JSON password map: + `{"redis://redis-safety:6379":"<ТОТ_ЖЕ_ПАРОЛЬ>"}`, а не как голая строка. + Пароль `safety` в `MESSAGE_SAFETY_REDIS_URL` также должен совпадать с ACL. Используйте + `redis/redis-safety.acl.template`, заменив оба плейсхолдера. 8. Выделите отдельные учётные данные БД для runtime и миграций. `MESSAGE_SAFETY_CONFIG_ADMIN_DATABASE_URL` может мигрировать/активировать политику, а `MESSAGE_SAFETY_DATABASE_URL` — нет; `BITRIX_SYNC_MIGRATION_DATABASE_URL` @@ -83,12 +87,14 @@ deployment-артефакты: `VM2_services/codebase/services/`. Локальн [`arch-10-deployment.md`](../../../../architectory/arch-10-deployment.md) §6. 3. **Образы** — собрать и push `han-message-safety`, `han-bitrix-sync`; получить immutable digest для всех `*_IMAGE` в `.env.example` (nginx, redis, - clamav, otel-collector). + clamav, otel-collector, Redis exporter, nginx exporter). 4. **Selectel Secrets Manager** — заполнить все remote names из `deployment/secrets/config.example.json` (DSN, tokens, S3 read-only keys, - `REDIS_SAFETY_ACL`, internal TLS PEM для `8443`). Отдельный IAM principal - VM2 с read-only доступом только к этим именам. + `REDIS_SAFETY_ACL`, `REDIS_EXPORTER_PASSWORD`, internal TLS PEM для `8443`). + Отдельный IAM principal VM2 с read-only доступом только к этим именам. 5. **S3 quarantine bucket** и SigNoz OTLP endpoint — значения в `.env`. + Текущий self-hosted SigNoz принимает private plaintext OTLP без auth, поэтому + не создавайте фиктивный auth-secret или обязательный непустой header. 6. **Internal TLS** — сертификат внутренней CA с SAN = private DNS VM2; PEM хранится в Secrets Manager, не в каталоге релиза. @@ -301,6 +307,13 @@ DSN, token, password, access/secret key туда не записываются. создайте отдельный VM2 IAM principal с read-only доступом только к remote names из mapping. +Локальный Collector принимает traces, metrics и logs напрямую по OTLP; чтение +Docker JSON через `filelog` не используется. Prometheus receiver собирает только +метрики самого Collector, `redis-exporter` и `nginx-exporter`, а `hostmetrics` — +метрики VM через read-only `/hostfs`. Exporter-контейнеры имеют только `expose` +во внутренних сетях и не публикуют host ports. В nginx endpoint +`/stub_status` слушает только внутренний `8081`. + Зашифруйте пароль Selectel service user через systemd credentials, не помещая его в аргументы или history: diff --git a/VM2_services/codebase/services/deployment/preflight.sh b/VM2_services/codebase/services/deployment/preflight.sh index 594c9ac..ac0b86c 100644 --- a/VM2_services/codebase/services/deployment/preflight.sh +++ b/VM2_services/codebase/services/deployment/preflight.sh @@ -59,7 +59,7 @@ if [ -f "$ENV_FILE" ]; then true|false) ;; *) fail "OTEL_REMOTE_TLS_INSECURE must be exactly true or false" ;; esac - for image_key in MESSAGE_SAFETY_IMAGE BITRIX_SYNC_IMAGE NGINX_IMAGE REDIS_IMAGE CLAMAV_IMAGE OTEL_COLLECTOR_IMAGE; do + for image_key in MESSAGE_SAFETY_IMAGE BITRIX_SYNC_IMAGE NGINX_IMAGE REDIS_IMAGE CLAMAV_IMAGE OTEL_COLLECTOR_IMAGE REDIS_EXPORTER_IMAGE NGINX_EXPORTER_IMAGE; do image=$(/usr/bin/awk -F= -v key="$image_key" '$1 == key {print substr($0, index($0, "=") + 1)}' "$ENV_FILE") echo "$image" | /usr/bin/grep -Eq '@sha256:[0-9a-f]{64}$' || fail "$image_key must be pinned by sha256 digest" @@ -70,6 +70,25 @@ if [ -f "$ENV_FILE" ]; then esac fi +collector_config="$ROOT/observability/otel-collector.yaml" +[ -f "$collector_config" ] || fail "OpenTelemetry Collector config is missing" +if [ -f "$collector_config" ]; then + /usr/bin/grep -Fq 'service.namespace, value: han-chat' "$collector_config" || + fail "Collector must enforce service.namespace=han-chat" + /usr/bin/grep -Fq 'receivers: [otlp, prometheus, hostmetrics]' "$collector_config" || + fail "Collector metrics pipeline is incomplete" + /usr/bin/grep -Fq 'receivers: [otlp]' "$collector_config" || + fail "Collector direct OTLP logs pipeline is missing" + ! /usr/bin/grep -Fq 'filelog' "$collector_config" || + fail "Collector filelog receiver is forbidden for direct OTLP logging" + /usr/bin/grep -Fq 'redis-exporter:9121' "$collector_config" || + fail "Collector Redis exporter scrape target is missing" + /usr/bin/grep -Fq 'nginx-exporter:9113' "$collector_config" || + fail "Collector nginx exporter scrape target is missing" + /usr/bin/grep -Fq 'tail_sampling:' "$collector_config" || + fail "Collector tail sampling is missing" +fi + bitrix_allowlist="$ROOT/nginx/allowlists/bitrix-webhook-allowlist.conf" private_allowlist="$ROOT/nginx/allowlists/private-caller-allowlist.conf" for allowlist in "$bitrix_allowlist" "$private_allowlist"; do @@ -106,7 +125,8 @@ BITRIX_SYNC_CRM_REST_WEBHOOK_URL BITRIX_SYNC_CONTACT_RECEIVER_TOKEN BITRIX_SYNC_ALERT_RECEIVER_TOKEN BITRIX_SYNC_SERVICE_TOKEN -REDIS_SAFETY_ACL' +REDIS_SAFETY_ACL +REDIS_EXPORTER_PASSWORD' if [ -f "$MANIFEST" ]; then old_ifs=$IFS @@ -123,6 +143,37 @@ if [ -f "$MANIFEST" ]; then done IFS=$old_ifs + redis_acl=$(/usr/bin/awk -F= \ + '$1 == "REDIS_SAFETY_ACL" {print substr($0, index($0, "=") + 1)}' \ + "$MANIFEST") + redis_exporter_password_file=$(/usr/bin/awk -F= \ + '$1 == "REDIS_EXPORTER_PASSWORD" {print substr($0, index($0, "=") + 1)}' \ + "$MANIFEST") + if [ -f "$redis_acl" ] && [ -f "$redis_exporter_password_file" ]; then + /usr/bin/grep -Eq '^user exporter reset on >[^[:space:]]+ -@all \+ping \+info$' "$redis_acl" || + fail "Redis ACL must contain the restricted exporter user" + if exporter_password=$(/usr/bin/python3 -c ' +import json +import sys + +target = "redis://redis-safety:6379" +with open(sys.argv[1], encoding="utf-8") as source: + values = json.load(source) +if set(values) != {target}: + raise SystemExit("password map must contain exactly redis://redis-safety:6379") +password = values[target] +if not isinstance(password, str) or not password or any(char.isspace() for char in password): + raise SystemExit("password must be a non-empty whitespace-free string") +print(password, end="") +' "$redis_exporter_password_file"); then + /usr/bin/grep -Fq -- ">$exporter_password " "$redis_acl" || + fail "Redis exporter password map does not match REDIS_SAFETY_ACL" + unset exporter_password + else + fail "REDIS_EXPORTER_PASSWORD must be a valid redis_exporter JSON password map" + fi + fi + internal_cert=$(/usr/bin/awk -F= \ '$1 == "VM2_INTERNAL_TLS_CERTIFICATE" {print substr($0, index($0, "=") + 1)}' \ "$MANIFEST") diff --git a/VM2_services/codebase/services/deployment/secrets/config.example.json b/VM2_services/codebase/services/deployment/secrets/config.example.json index 41df51a..f65b19e 100644 --- a/VM2_services/codebase/services/deployment/secrets/config.example.json +++ b/VM2_services/codebase/services/deployment/secrets/config.example.json @@ -90,6 +90,11 @@ "remote": "vm2/REDIS_SAFETY_ACL", "consumers": ["redis-safety"], "max_bytes": 4096 + }, + "REDIS_EXPORTER_PASSWORD": { + "remote": "vm2/REDIS_EXPORTER_PASSWORD", + "consumers": ["redis-exporter"], + "max_bytes": 1024 } } } diff --git a/VM2_services/codebase/services/docker-compose.yml b/VM2_services/codebase/services/docker-compose.yml index cdf19f9..fb723c3 100644 --- a/VM2_services/codebase/services/docker-compose.yml +++ b/VM2_services/codebase/services/docker-compose.yml @@ -16,6 +16,8 @@ x-postgres-ca-volume: &postgres-ca-volume x-message-safety-environment: &message-safety-environment APP_ENV: ${APP_ENV:?set APP_ENV} + RELEASE_VERSION: ${RELEASE_VERSION:?set RELEASE_VERSION} + LOG_LEVEL: ${LOG_LEVEL:-INFO} MESSAGE_SAFETY_HOST: ${MESSAGE_SAFETY_HOST:-0.0.0.0} MESSAGE_SAFETY_PORT: ${MESSAGE_SAFETY_PORT:-8080} MESSAGE_SAFETY_WORKER_CONCURRENCY: ${MESSAGE_SAFETY_WORKER_CONCURRENCY:-5} @@ -36,6 +38,8 @@ x-message-safety-environment: &message-safety-environment x-bitrix-sync-environment: &bitrix-sync-environment APP_ENV: ${APP_ENV:?set APP_ENV} + RELEASE_VERSION: ${RELEASE_VERSION:?set RELEASE_VERSION} + LOG_LEVEL: ${LOG_LEVEL:-INFO} BITRIX_SYNC_ENABLED: ${BITRIX_SYNC_ENABLED:-false} BITRIX_SYNC_MODE: ${BITRIX_SYNC_MODE:-disabled} BITRIX_SYNC_PORTAL_HOST: ${BITRIX_SYNC_PORTAL_HOST:?set approved portal} @@ -91,7 +95,8 @@ services: nofile: soft: 4096 hard: 4096 - networks: [public, backend] + networks: [public, backend, observability] + expose: ["8081"] depends_on: message-safety-api: condition: service_started @@ -130,6 +135,40 @@ services: mem_limit: 640m cpus: 1.0 + redis-exporter: + <<: *hardening + image: ${REDIS_EXPORTER_IMAGE:?set immutable Redis exporter image digest} + environment: + REDIS_ADDR: redis://redis-safety:6379 + REDIS_USER: exporter + REDIS_PASSWORD_FILE: /run/secrets/redis_exporter_password + secrets: + - redis_exporter_password + networks: [backend, observability] + expose: ["9121"] + depends_on: + redis-safety: + condition: service_healthy + tmpfs: + - /tmp:rw,noexec,nosuid,nodev,size=16m + pids_limit: 50 + mem_limit: 64m + cpus: 0.25 + + nginx-exporter: + <<: *hardening + image: ${NGINX_EXPORTER_IMAGE:?set immutable nginx exporter image digest} + command: + - --nginx.scrape-uri=http://nginx:8081/stub_status + networks: [observability] + expose: ["9113"] + depends_on: + nginx: + condition: service_healthy + pids_limit: 50 + mem_limit: 64m + cpus: 0.25 + clamd: <<: *hardening image: ${CLAMAV_IMAGE:?set immutable ClamAV image digest} @@ -203,6 +242,7 @@ services: volumes: - ./observability/otel-collector.yaml:/etc/otelcol-contrib/config.yaml:ro - otel-queue:/var/lib/otelcol/queue + - /:/hostfs:ro networks: observability: {} telemetry-egress: @@ -210,9 +250,9 @@ services: depends_on: otel-queue-init: condition: service_completed_successfully - expose: ["4317", "4318"] + expose: ["4317", "4318", "13133", "8888"] healthcheck: - test: ["CMD", "/otelcol-contrib", "components"] + test: ["CMD", "/otelcol-contrib", "validate", "--config=/etc/otelcol-contrib/config.yaml"] interval: 30s timeout: 5s retries: 3 @@ -465,3 +505,5 @@ secrets: file: /run/han-chat/secrets/BITRIX_SYNC_SERVICE_TOKEN redis_safety_acl: file: /run/han-chat/secrets/REDIS_SAFETY_ACL + redis_exporter_password: + file: /run/han-chat/secrets/REDIS_EXPORTER_PASSWORD diff --git a/VM2_services/codebase/services/message-safety/alembic/versions/0002_task_trace_context.py b/VM2_services/codebase/services/message-safety/alembic/versions/0002_task_trace_context.py new file mode 100644 index 0000000..10433f4 --- /dev/null +++ b/VM2_services/codebase/services/message-safety/alembic/versions/0002_task_trace_context.py @@ -0,0 +1,42 @@ +"""store safe async trace context on safety tasks + +Revision ID: 0002_task_trace_context +Revises: 0001_message_safety_v2 +""" + +from alembic import op + +revision = "0002_task_trace_context" +down_revision = "0001_message_safety_v2" +branch_labels = None +depends_on = None + + +def upgrade() -> None: + op.execute( + "ALTER TABLE message_safety.safety_tasks " + "ADD COLUMN IF NOT EXISTS origin_trace_id bytea NULL" + ) + op.execute( + "ALTER TABLE message_safety.safety_tasks " + "ADD COLUMN IF NOT EXISTS origin_span_id bytea NULL" + ) + op.execute( + "ALTER TABLE message_safety.safety_tasks " + "ADD COLUMN IF NOT EXISTS origin_trace_flags integer NULL" + ) + op.execute( + "ALTER TABLE message_safety.safety_tasks " + "ADD COLUMN IF NOT EXISTS origin_tracestate varchar(512) NULL" + ) + + +def downgrade() -> None: + op.execute( + "ALTER TABLE message_safety.safety_tasks DROP COLUMN IF EXISTS origin_tracestate" + ) + op.execute( + "ALTER TABLE message_safety.safety_tasks DROP COLUMN IF EXISTS origin_trace_flags" + ) + op.execute("ALTER TABLE message_safety.safety_tasks DROP COLUMN IF EXISTS origin_span_id") + op.execute("ALTER TABLE message_safety.safety_tasks DROP COLUMN IF EXISTS origin_trace_id") diff --git a/VM2_services/codebase/services/message-safety/app/adapters.py b/VM2_services/codebase/services/message-safety/app/adapters.py index 79d2599..bde7084 100644 --- a/VM2_services/codebase/services/message-safety/app/adapters.py +++ b/VM2_services/codebase/services/message-safety/app/adapters.py @@ -8,11 +8,14 @@ import boto3 import dns.asyncresolver from botocore.config import Config from botocore.exceptions import BotoCoreError, ClientError +from opentelemetry import trace from app.contracts import Attachment from app.file_pipeline import DependencyFailure, ObjectChanged from app.url_policy import DnsError, DnsNxDomain +tracer = trace.get_tracer("message-safety.dependencies") + class TrustedDnsResolver: def __init__(self, nameservers: list[str]) -> None: @@ -22,6 +25,14 @@ class TrustedDnsResolver: async def resolve( self, hostname: str + ) -> tuple[ipaddress.IPv4Address | ipaddress.IPv6Address, ...]: + with tracer.start_as_current_span( + "message_safety.dns.resolve", attributes={"server.address.type": "domain"} + ): + return await self._resolve(hostname) + + async def _resolve( + self, hostname: str ) -> tuple[ipaddress.IPv4Address | ipaddress.IPv6Address, ...]: found: list[ipaddress.IPv4Address | ipaddress.IPv6Address] = [] try: diff --git a/VM2_services/codebase/services/message-safety/app/api.py b/VM2_services/codebase/services/message-safety/app/api.py index 1b3670e..d80b4cb 100644 --- a/VM2_services/codebase/services/message-safety/app/api.py +++ b/VM2_services/codebase/services/message-safety/app/api.py @@ -3,6 +3,7 @@ from __future__ import annotations import hmac import json import uuid +from time import perf_counter from typing import Annotated from fastapi import Depends, FastAPI, Header, Request @@ -14,6 +15,7 @@ from app.db import TaskStatus from app.rate_limit import RateLimited from app.repository import ConflictError from app.service import CapabilityUnavailable, SafetyService, TaskFailed +from app.telemetry import logger, record_poll CHECK_ADAPTER = TypeAdapter(CheckRequest) MAX_BODY = 16_384 @@ -53,14 +55,39 @@ def create_app(service: SafetyService, token: str) -> FastAPI: @app.middleware("http") async def request_context(request: Request, call_next): + started = perf_counter() supplied = request.headers.get("X-Request-ID") try: request.state.request_id = str(uuid.UUID(supplied)) if supplied else str(uuid.uuid4()) except ValueError: request.state.request_id = str(uuid.uuid4()) - response = await call_next(request) + try: + response = await call_next(request) + except Exception: + logger.error( + "Safety request failed", + extra={ + "event": "http_request_failed", + "request_id": request.state.request_id, + "route": _route_template(request), + "method": request.method, + }, + ) + raise response.headers["X-Request-ID"] = request.state.request_id response.headers["Cache-Control"] = "no-store" + if request.url.path != "/health/live": + logger.info( + "Safety request completed", + extra={ + "event": "http_request_completed", + "request_id": request.state.request_id, + "route": _route_template(request), + "method": request.method, + "status_code": response.status_code, + "duration_ms": round((perf_counter() - started) * 1000, 3), + }, + ) return response @app.get("/health/live") @@ -197,8 +224,10 @@ def create_app(service: SafetyService, token: str) -> FastAPI: ) task = await service.repository.task(parsed) if not task: + record_poll("not_found") return error(404, "task_not_found", request.state.request_id) if task.status == TaskStatus.failed: + record_poll("failed") return error( 503, "task_failed", @@ -211,6 +240,7 @@ def create_app(service: SafetyService, token: str) -> FastAPI: }, ) if task.status in {TaskStatus.pending, TaskStatus.processing}: + record_poll(task.status.value) result: Pending | Verdict = Pending( config_version=task.config_version, task_id=task.id, @@ -218,6 +248,7 @@ def create_app(service: SafetyService, token: str) -> FastAPI: rules_version=task.rules_version, ) else: + record_poll("terminal") result = service._verdict( task.status == TaskStatus.allowed, task.processing_mode, @@ -233,3 +264,8 @@ def create_app(service: SafetyService, token: str) -> FastAPI: return response return app + + +def _route_template(request: Request) -> str: + route = request.scope.get("route") + return getattr(route, "path", "unmatched") diff --git a/VM2_services/codebase/services/message-safety/app/db.py b/VM2_services/codebase/services/message-safety/app/db.py index a5d31d1..f12bf35 100644 --- a/VM2_services/codebase/services/message-safety/app/db.py +++ b/VM2_services/codebase/services/message-safety/app/db.py @@ -145,6 +145,10 @@ class SafetyTask(Base): detector_version: Mapped[str] = mapped_column(String(128)) scanner_engine: Mapped[str] = mapped_column(String(32)) signatures_version: Mapped[str] = mapped_column(String(128)) + origin_trace_id: Mapped[bytes | None] = mapped_column(LargeBinary(16)) + origin_span_id: Mapped[bytes | None] = mapped_column(LargeBinary(8)) + origin_trace_flags: Mapped[int | None] = mapped_column(Integer) + origin_tracestate: Mapped[str | None] = mapped_column(String(512)) created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True)) updated_at: Mapped[datetime] = mapped_column(DateTime(timezone=True)) finished_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True)) diff --git a/VM2_services/codebase/services/message-safety/app/file_pipeline.py b/VM2_services/codebase/services/message-safety/app/file_pipeline.py index 7aa3b4f..fbfd656 100644 --- a/VM2_services/codebase/services/message-safety/app/file_pipeline.py +++ b/VM2_services/codebase/services/message-safety/app/file_pipeline.py @@ -10,12 +10,14 @@ from dataclasses import dataclass from pathlib import Path from typing import Protocol +from opentelemetry import trace from PIL import Image, UnidentifiedImageError from pillow_heif import register_heif_opener from app.contracts import Attachment register_heif_opener() +tracer = trace.get_tracer("message-safety.dependencies") class ObjectChanged(RuntimeError): @@ -122,6 +124,13 @@ class ClamAvInstream: self.host, self.port, self.timeout = host, port, timeout async def scan(self, chunks: AsyncIterator[bytes]) -> str | None: + with tracer.start_as_current_span( + "message_safety.clamav.scan", + attributes={"server.address.type": "clamav"}, + ): + return await self._scan(chunks) + + async def _scan(self, chunks: AsyncIterator[bytes]) -> str | None: async def operation() -> str | None: reader, writer = await asyncio.open_connection(self.host, self.port) try: @@ -148,6 +157,10 @@ class ClamAvInstream: raise DependencyFailure("ClamAV unavailable") from exc async def signatures_version(self) -> str: + with tracer.start_as_current_span("message_safety.clamav.version"): + return await self._signatures_version() + + async def _signatures_version(self) -> str: try: reader, writer = await asyncio.wait_for( asyncio.open_connection(self.host, self.port), 2.0 diff --git a/VM2_services/codebase/services/message-safety/app/main.py b/VM2_services/codebase/services/message-safety/app/main.py index 318929a..6a68a5d 100644 --- a/VM2_services/codebase/services/message-safety/app/main.py +++ b/VM2_services/codebase/services/message-safety/app/main.py @@ -12,6 +12,7 @@ from app.file_pipeline import ClamAvInstream, DependencyFailure from app.repository import Repository from app.service import SafetyService from app.settings import BootstrapSettings, EmergencyMode +from app.telemetry import init_telemetry, instrument_fastapi, shutdown_telemetry async def build_runtime() -> tuple[object, object]: @@ -48,19 +49,24 @@ async def build_runtime() -> tuple[object, object]: signatures_version=signatures_version, ) app = create_app(service, settings.service_token.get_secret_value()) + instrument_fastapi(app) return app, engine async def serve() -> None: settings = BootstrapSettings() - app, engine = await build_runtime() + init_telemetry("api") + engine = None try: + app, engine = await build_runtime() server = uvicorn.Server( uvicorn.Config(app, host=settings.host, port=settings.port, proxy_headers=False) ) await server.serve() finally: - await engine.dispose() + if engine is not None: + await engine.dispose() + shutdown_telemetry() def run() -> None: diff --git a/VM2_services/codebase/services/message-safety/app/service.py b/VM2_services/codebase/services/message-safety/app/service.py index 87bd7f8..2d4445b 100644 --- a/VM2_services/codebase/services/message-safety/app/service.py +++ b/VM2_services/codebase/services/message-safety/app/service.py @@ -2,6 +2,7 @@ from __future__ import annotations import uuid from datetime import UTC, datetime, timedelta +from time import perf_counter from app.config import ActiveConfig from app.contracts import CheckRequest, FileCheck, Pending, TextCheck, Verdict @@ -19,6 +20,14 @@ from app.normalization import normalize_text from app.rate_limit import ConservativeRateLimiter from app.repository import QueueFull, Repository from app.settings import EmergencyMode +from app.telemetry import ( + current_trace_fields, + record_cache, + record_check, + record_dependency, + record_file_rejection, + record_runtime_state, +) from app.url_policy import DnsError, Resolver, canonicalize, check_url, extract_urls @@ -54,6 +63,7 @@ class SafetyService: self.rate_limiter = ConservativeRateLimiter( config.document["rate"]["text_rps"], config.document["rate"]["file_rps"] ) + record_runtime_state("mock" if mode.mock else "standard", config.version) def _verdict( self, @@ -74,6 +84,22 @@ class SafetyService: ) async def check(self, request: CheckRequest) -> Verdict | Pending: + started = perf_counter() + outcome = "error" + try: + result = await self._check(request) + outcome = "pending" if isinstance(result, Pending) else result.verdict + return result + finally: + record_check( + request.content_kind, + "mock" if self.mode.mock else "standard", + outcome, + (perf_counter() - started) * 1000, + self.config.version, + ) + + async def _check(self, request: CheckRequest) -> Verdict | Pending: digest = fingerprint(request) existing = await self.repository.get_request(request.message_id) if existing: @@ -101,6 +127,7 @@ class SafetyService: cache = await self.repository.text_cache( normalized.analysis_sha256, self.config.rules_version ) + record_cache("text_rules", cache is not None) if cache: deny_rule = cache.deny_rule_id else: @@ -145,6 +172,7 @@ class SafetyService: cached_link = await self.repository.link_cache( canonical.digest, self.config.rules_version, self.config.version ) + record_cache("link_verdict", cached_link is not None) if cached_link and cached_link.verdict == "deny": return await self._persist_sync( request, @@ -160,7 +188,9 @@ class SafetyService: _, rule = await check_url( raw, self.resolver, self.config.document["link"]["dns_lookup_timeout_sec"] ) + record_dependency("dns", "resolve", "success") except DnsError as exc: + record_dependency("dns", "resolve", "error") raise CapabilityUnavailable("dns") from exc if rule != "url.nxdomain" and not cached_link: now = datetime.now(UTC) @@ -200,6 +230,7 @@ class SafetyService: set(self.config.document["file_policy"]["enabled_mime_types"]), ) if rule: + record_file_rejection(rule) return await self._persist_sync( request, digest, self._verdict(False, "standard", rule, self.config.rules_version) ) @@ -207,6 +238,7 @@ class SafetyService: cached = await self.repository.file_cache( content_digest, self.config, self.signatures_version ) + record_cache("file_verdict", cached is not None) if cached: return await self._persist_sync( request, @@ -220,6 +252,7 @@ class SafetyService: ) now = datetime.now(UTC) task_id = uuid.uuid4() + trace_id, span_id, trace_flags, tracestate = current_trace_fields() task = SafetyTask( id=task_id, message_id=request.message_id, @@ -243,6 +276,10 @@ class SafetyService: detector_version=self.config.detector.version, scanner_engine="clamav", signatures_version=self.signatures_version, + origin_trace_id=trace_id, + origin_span_id=span_id, + origin_trace_flags=trace_flags, + origin_tracestate=tracestate, created_at=now, updated_at=now, ) @@ -262,6 +299,7 @@ class SafetyService: row, task, max_pending=self.config.document["task"]["max_pending"] ) except QueueFull as exc: + record_dependency("postgresql", "queue_capacity", "full") raise CapabilityUnavailable("queue_capacity") from exc if created: await self._audit( diff --git a/VM2_services/codebase/services/message-safety/app/telemetry.py b/VM2_services/codebase/services/message-safety/app/telemetry.py new file mode 100644 index 0000000..a411d21 --- /dev/null +++ b/VM2_services/codebase/services/message-safety/app/telemetry.py @@ -0,0 +1,377 @@ +from __future__ import annotations + +import json +import logging +import os +import re +import sys +from collections.abc import Mapping, MutableMapping +from dataclasses import dataclass +from datetime import UTC, datetime +from typing import Any + +from fastapi import FastAPI +from opentelemetry import metrics, trace +from opentelemetry.exporter.otlp.proto.grpc._log_exporter import OTLPLogExporter +from opentelemetry.exporter.otlp.proto.grpc.metric_exporter import OTLPMetricExporter +from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporter +from opentelemetry.instrumentation.botocore import BotocoreInstrumentor +from opentelemetry.instrumentation.fastapi import FastAPIInstrumentor +from opentelemetry.instrumentation.httpx import HTTPXClientInstrumentor +from opentelemetry.instrumentation.redis import RedisInstrumentor +from opentelemetry.instrumentation.sqlalchemy import SQLAlchemyInstrumentor +from opentelemetry.propagate import set_global_textmap +from opentelemetry.sdk._logs import LoggerProvider, LoggingHandler +from opentelemetry.sdk._logs.export import BatchLogRecordProcessor +from opentelemetry.sdk.metrics import MeterProvider +from opentelemetry.sdk.metrics.export import PeriodicExportingMetricReader +from opentelemetry.sdk.resources import Resource +from opentelemetry.sdk.trace import TracerProvider +from opentelemetry.sdk.trace.export import BatchSpanProcessor +from opentelemetry.trace import Link, SpanContext, TraceFlags, TraceState +from opentelemetry.trace.propagation.tracecontext import TraceContextTextMapPropagator + +SERVICE_NAME = "message-safety" +SERVICE_NAMESPACE = "han-chat" +_SENSITIVE_KEY = re.compile( + r"(authorization|cookie|token|secret|password|phone|email|message|text|filename|" + r"object[_-]?key|url|dsn|statement)", + re.IGNORECASE, +) +_SENSITIVE_VALUE = re.compile( + r"(?:https?://\S+\?\S+)|(?:[\w.+-]+@[\w.-]+\.[A-Za-z]{2,})|" + r"(?:\+?\d[\d ()-]{8,}\d)|(?:bearer\s+\S+)", + re.IGNORECASE, +) +_STANDARD_LOG_RECORD = set(logging.makeLogRecord({}).__dict__) +_SAFE_STRUCTURED_KEYS = { + "deployment.environment", + "duration_ms", + "environment", + "event", + "level", + "method", + "outcome", + "process_role", + "request_id", + "route", + "service.name", + "service.version", + "span_id", + "status_code", + "timestamp", + "trace_id", +} +_runtime: TelemetryRuntime | None = None + + +def _redact(value: Any, key: str = "") -> Any: + if key in _SAFE_STRUCTURED_KEYS: + return value + if _SENSITIVE_KEY.search(key): + return "[REDACTED]" + if isinstance(value, dict): + return {str(k)[:64]: _redact(v, str(k)) for k, v in list(value.items())[:64]} + if isinstance(value, (list, tuple)): + return [_redact(item) for item in value[:32]] + if isinstance(value, str): + return _SENSITIVE_VALUE.sub("[REDACTED]", value) + return value + + +class RedactionFilter(logging.Filter): + def filter(self, record: logging.LogRecord) -> bool: + record.msg = _redact(record.getMessage()) + record.args = () + for key in tuple(record.__dict__): + if key not in _STANDARD_LOG_RECORD: + record.__dict__[key] = _redact(record.__dict__[key], key) + return True + + +def _redact_span_attributes(attributes: MutableMapping[str, Any]) -> None: + for key in tuple(attributes): + key_text = str(key) + if _SENSITIVE_KEY.search(key_text) or key_text in { + "http.target", + "http.url", + "url.full", + "url.query", + "db.statement", + "db.query.text", + }: + attributes[key] = "[REDACTED]" + else: + attributes[key] = _redact(attributes[key], key_text) + + +class RedactingBatchSpanProcessor(BatchSpanProcessor): + """Redact auto-instrumentation attributes before they enter the export queue.""" + + def on_end(self, span: Any) -> None: + attributes = getattr(span, "_attributes", None) + if isinstance(attributes, Mapping): + sanitized = dict(attributes) + _redact_span_attributes(sanitized) + span._attributes = sanitized + super().on_end(span) + + +class JsonFormatter(logging.Formatter): + def format(self, record: logging.LogRecord) -> str: + context = trace.get_current_span().get_span_context() + payload: dict[str, Any] = { + "timestamp": datetime.now(UTC) + .isoformat(timespec="milliseconds") + .replace("+00:00", "Z"), + "level": record.levelname, + "service.name": SERVICE_NAME, + "service.version": os.getenv("RELEASE_VERSION", "unknown"), + "environment": os.getenv("APP_ENV", "development"), + "event": getattr(record, "event", "application_log"), + "message": record.getMessage(), + } + if context.is_valid: + payload["trace_id"] = f"{context.trace_id:032x}" + payload["span_id"] = f"{context.span_id:016x}" + for key in ( + "request_id", + "route", + "method", + "status_code", + "duration_ms", + "process_role", + "outcome", + ): + value = getattr(record, key, None) + if value is not None: + payload[key] = value + return json.dumps(_redact(payload), ensure_ascii=False, separators=(",", ":")) + + +def configure_logging() -> logging.Logger: + logger = logging.getLogger("message_safety") + logger.setLevel(os.getenv("LOG_LEVEL", "INFO").upper()) + logger.propagate = False + if not any(getattr(handler, "_message_safety_stdout", False) for handler in logger.handlers): + handler = logging.StreamHandler(sys.stdout) + handler._message_safety_stdout = True # type: ignore[attr-defined] + handler.addFilter(RedactionFilter()) + handler.setFormatter(JsonFormatter()) + logger.addHandler(handler) + return logger + + +logger = configure_logging() + + +@dataclass(slots=True) +class TelemetryRuntime: + tracer_provider: TracerProvider + meter_provider: MeterProvider + logger_provider: LoggerProvider + + def shutdown(self) -> None: + for provider in (self.logger_provider, self.meter_provider, self.tracer_provider): + try: + provider.shutdown() + except Exception: + logger.error( + "Telemetry provider shutdown failed", + extra={"event": "telemetry_shutdown_failed"}, + ) + + +def _resource(role: str) -> Resource: + return Resource.create( + { + "service.name": SERVICE_NAME, + "service.namespace": SERVICE_NAMESPACE, + "service.version": os.getenv("RELEASE_VERSION", "unknown"), + "deployment.environment": os.getenv("APP_ENV", "development"), + "process.role": role, + } + ) + + +def init_telemetry(role: str) -> TelemetryRuntime | None: + global _runtime + if _runtime is not None: + return _runtime + endpoint = os.getenv("OTEL_EXPORTER_OTLP_ENDPOINT", "").strip() + if not endpoint: + logger.info( + "OTLP endpoint is not configured; using JSON stdout", + extra={"event": "telemetry_stdout_only", "process_role": role}, + ) + return None + try: + resource = _resource(role) + insecure = endpoint.startswith("http://") + set_global_textmap(TraceContextTextMapPropagator()) + + tracer_provider = TracerProvider(resource=resource) + tracer_provider.add_span_processor( + RedactingBatchSpanProcessor( + OTLPSpanExporter(endpoint=endpoint, insecure=insecure, timeout=3), + max_queue_size=2048, + max_export_batch_size=512, + export_timeout_millis=3000, + ) + ) + trace.set_tracer_provider(tracer_provider) + + metric_reader = PeriodicExportingMetricReader( + OTLPMetricExporter(endpoint=endpoint, insecure=insecure, timeout=3), + export_interval_millis=30_000, + export_timeout_millis=3000, + ) + meter_provider = MeterProvider(resource=resource, metric_readers=[metric_reader]) + metrics.set_meter_provider(meter_provider) + + logger_provider = LoggerProvider(resource=resource) + logger_provider.add_log_record_processor( + BatchLogRecordProcessor( + OTLPLogExporter(endpoint=endpoint, insecure=insecure, timeout=3), + max_queue_size=2048, + max_export_batch_size=512, + export_timeout_millis=3000, + ) + ) + otlp_handler = LoggingHandler(level=logging.INFO, logger_provider=logger_provider) + otlp_handler.addFilter(RedactionFilter()) + logger.addHandler(otlp_handler) + + HTTPXClientInstrumentor().instrument() + SQLAlchemyInstrumentor().instrument(enable_commenter=False) + RedisInstrumentor().instrument() + BotocoreInstrumentor().instrument() + _runtime = TelemetryRuntime(tracer_provider, meter_provider, logger_provider) + logger.info( + "OpenTelemetry initialized", + extra={"event": "telemetry_initialized", "process_role": role}, + ) + return _runtime + except Exception: + logger.error( + "OpenTelemetry initialization failed; continuing with JSON stdout", + extra={"event": "telemetry_init_failed", "process_role": role}, + ) + return None + + +def instrument_fastapi(app: FastAPI) -> None: + if _runtime is None: + return + try: + FastAPIInstrumentor.instrument_app(app, excluded_urls="/health/live,/health/ready") + except Exception: + logger.error( + "FastAPI instrumentation failed; continuing", + extra={"event": "telemetry_instrumentation_failed"}, + ) + + +def shutdown_telemetry() -> None: + global _runtime + if _runtime is not None: + _runtime.shutdown() + _runtime = None + + +_meter = metrics.get_meter("message-safety.business") +checks = _meter.create_counter("message_safety.checks", description="Safety checks by outcome") +check_duration = _meter.create_histogram( + "message_safety.check.duration", unit="ms", description="Safety check duration" +) +cache_access = _meter.create_counter( + "message_safety.cache.access", description="Safety cache accesses" +) +worker_tasks = _meter.create_counter( + "message_safety.worker.tasks", description="Worker task outcomes" +) +worker_age = _meter.create_histogram( + "message_safety.worker.task.age", unit="s", description="Age of claimed worker tasks" +) +dependency_calls = _meter.create_counter( + "message_safety.dependency.calls", description="Dependency call outcomes" +) +polls = _meter.create_counter("message_safety.polls", description="Task poll outcomes") +file_rejections = _meter.create_counter( + "message_safety.file.rejections", description="File rejection categories" +) +runtime_state = _meter.create_up_down_counter( + "message_safety.runtime.state", description="Current runtime mode and config" +) + + +def record_check( + kind: str, mode: str, outcome: str, duration_ms: float, config_version: int +) -> None: + attributes = { + "content_kind": kind, + "processing_mode": mode, + "outcome": outcome, + "config_version": config_version, + } + checks.add(1, attributes) + check_duration.record(duration_ms, attributes) + + +def record_cache(cache: str, hit: bool) -> None: + cache_access.add(1, {"cache": cache, "result": "hit" if hit else "miss"}) + + +def record_worker(outcome: str, age_seconds: float | None = None) -> None: + worker_tasks.add(1, {"outcome": outcome}) + if age_seconds is not None: + worker_age.record(max(age_seconds, 0), {"outcome": outcome}) + + +def record_dependency(dependency: str, operation: str, outcome: str) -> None: + dependency_calls.add( + 1, {"dependency": dependency, "operation": operation, "outcome": outcome} + ) + + +def record_poll(outcome: str) -> None: + polls.add(1, {"outcome": outcome}) + + +def record_file_rejection(reason: str) -> None: + file_rejections.add(1, {"reason": reason}) + + +def record_runtime_state(mode: str, config_version: int) -> None: + runtime_state.add(1, {"processing_mode": mode, "config_version": config_version}) + + +def current_trace_fields() -> tuple[bytes | None, bytes | None, int | None, str | None]: + context = trace.get_current_span().get_span_context() + if not context.is_valid: + return None, None, None, None + return ( + context.trace_id.to_bytes(16, "big"), + context.span_id.to_bytes(8, "big"), + int(context.trace_flags), + str(context.trace_state) or None, + ) + + +def task_link(task: Any) -> Link | None: + trace_id = getattr(task, "origin_trace_id", None) + span_id = getattr(task, "origin_span_id", None) + if not trace_id or not span_id: + return None + try: + tracestate = getattr(task, "origin_tracestate", None) + context = SpanContext( + trace_id=int.from_bytes(trace_id, "big"), + span_id=int.from_bytes(span_id, "big"), + is_remote=True, + trace_flags=TraceFlags(getattr(task, "origin_trace_flags", 0) or 0), + trace_state=TraceState.from_header([tracestate] if tracestate else []), + ) + return Link(context) if context.is_valid else None + except (TypeError, ValueError): + return None diff --git a/VM2_services/codebase/services/message-safety/app/worker.py b/VM2_services/codebase/services/message-safety/app/worker.py index c9b49ce..00d043c 100644 --- a/VM2_services/codebase/services/message-safety/app/worker.py +++ b/VM2_services/codebase/services/message-safety/app/worker.py @@ -5,6 +5,8 @@ import socket import uuid from datetime import UTC, datetime, timedelta +from opentelemetry import trace + from app.adapters import S3VersionReader from app.config import validate_config from app.contracts import Attachment @@ -19,6 +21,15 @@ from app.file_pipeline import ( ) from app.repository import Repository from app.settings import BootstrapSettings +from app.telemetry import ( + init_telemetry, + record_dependency, + record_worker, + shutdown_telemetry, + task_link, +) + +tracer = trace.get_tracer("message-safety.worker") class Worker: @@ -34,14 +45,30 @@ class Worker: self.owner = f"{socket.gethostname()}:{uuid.uuid4()}" async def once(self) -> bool: - task = await self.repository.claim(self.owner) + with tracer.start_as_current_span("message_safety.worker.claim"): + task = await self.repository.claim(self.owner) if not task: return False + link = task_link(task) + with tracer.start_as_current_span( + "message_safety.worker.process", + links=[link] if link else (), + attributes={ + "messaging.operation.type": "process", + "message_safety.attempt": task.attempt_count, + }, + ): + await self._process(task) + return True + + async def _process(self, task) -> None: + task_age = (datetime.now(UTC) - task.created_at).total_seconds() row = await self.repository.config_version(task.config_version) rules, detector, digest = validate_config(row.config, self.artifacts) if digest != row.config_sha256: await self.repository.retry_or_fail(task, row.config["task"]["max_attempts"]) - return True + record_worker("config_error", task_age) + return stop = asyncio.Event() heartbeat = asyncio.create_task( self._heartbeat( @@ -65,18 +92,22 @@ class Worker: body, _ = await collect_and_hash( self.reader, attachment, max_size=row.config["file_policy"]["max_size_bytes"] ) + record_dependency("s3", "get_object", "success") rule = detect_format(body, attachment.mime_type) if not rule: malware = await self.antivirus.scan(one_chunk(body)) + record_dependency("clamav", "scan", "success") rule = "file.malware_detected" if malware else None - finished = await self.repository.finish( - task.id, - self.owner, - task.lease_generation, - allow=rule is None, - rule_id=rule or "safety.all_checks_passed", - ) + with tracer.start_as_current_span("message_safety.worker.finalize"): + finished = await self.repository.finish( + task.id, + self.owner, + task.lease_generation, + allow=rule is None, + rule_id=rule or "safety.all_checks_passed", + ) if finished: + record_worker("allow" if rule is None else "deny", task_age) now = datetime.now(UTC) await self.repository.put_file_cache( FileVerdictCache( @@ -111,19 +142,23 @@ class Worker: ) ) except ObjectChanged: - await self.repository.finish( - task.id, - self.owner, - task.lease_generation, - allow=False, - rule_id="file.object_changed", - ) + record_dependency("s3", "get_object", "object_changed") + with tracer.start_as_current_span("message_safety.worker.finalize"): + await self.repository.finish( + task.id, + self.owner, + task.lease_generation, + allow=False, + rule_id="file.object_changed", + ) + record_worker("deny", task_age) except DependencyFailure: + record_dependency("file_pipeline", "scan", "error") await self.repository.retry_or_fail(task, row.config["task"]["max_attempts"]) + record_worker("retry_or_fail", task_age) finally: stop.set() await heartbeat - return True async def _heartbeat( self, task_id, generation: int, interval: int, lease: int, stop: asyncio.Event @@ -144,25 +179,29 @@ class Worker: async def serve() -> None: settings = BootstrapSettings() - assert settings.database_url and settings.s3_access_key and settings.s3_secret_key - engine, sessions = engine_and_sessions(settings.database_url.get_secret_value()) - worker = Worker( - Repository(sessions), - S3VersionReader( - settings.s3_endpoint_url, - settings.s3_bucket, - settings.s3_access_key.get_secret_value(), - settings.s3_secret_key.get_secret_value(), - ), - ClamAvInstream(settings.clamav_host, settings.clamav_port), - settings.artifacts_dir, - ) + init_telemetry("worker") + engine = None try: + assert settings.database_url and settings.s3_access_key and settings.s3_secret_key + engine, sessions = engine_and_sessions(settings.database_url.get_secret_value()) + worker = Worker( + Repository(sessions), + S3VersionReader( + settings.s3_endpoint_url, + settings.s3_bucket, + settings.s3_access_key.get_secret_value(), + settings.s3_secret_key.get_secret_value(), + ), + ClamAvInstream(settings.clamav_host, settings.clamav_port), + settings.artifacts_dir, + ) async with asyncio.TaskGroup() as group: for _ in range(settings.worker_concurrency): group.create_task(worker.loop()) finally: - await engine.dispose() + if engine is not None: + await engine.dispose() + shutdown_telemetry() def run() -> None: diff --git a/VM2_services/codebase/services/message-safety/pyproject.toml b/VM2_services/codebase/services/message-safety/pyproject.toml index ab9f5ed..36c58d4 100644 --- a/VM2_services/codebase/services/message-safety/pyproject.toml +++ b/VM2_services/codebase/services/message-safety/pyproject.toml @@ -11,6 +11,14 @@ dependencies = [ "httpx>=0.27", "idna>=3.7", "jsonschema>=4.23", + "opentelemetry-api>=1.36,<2", + "opentelemetry-exporter-otlp-proto-grpc>=1.36,<2", + "opentelemetry-instrumentation-botocore>=0.57b0,<1", + "opentelemetry-instrumentation-fastapi>=0.57b0,<1", + "opentelemetry-instrumentation-httpx>=0.57b0,<1", + "opentelemetry-instrumentation-redis>=0.57b0,<1", + "opentelemetry-instrumentation-sqlalchemy>=0.57b0,<1", + "opentelemetry-sdk>=1.36,<2", "pillow>=10.4", "pillow-heif>=0.18", "pydantic-settings>=2.5", diff --git a/VM2_services/codebase/services/message-safety/tests/test_telemetry.py b/VM2_services/codebase/services/message-safety/tests/test_telemetry.py new file mode 100644 index 0000000..92f746c --- /dev/null +++ b/VM2_services/codebase/services/message-safety/tests/test_telemetry.py @@ -0,0 +1,131 @@ +from __future__ import annotations + +import json +import logging +from types import SimpleNamespace + +from opentelemetry import trace +from opentelemetry.sdk.trace import ReadableSpan +from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter +from opentelemetry.trace import NonRecordingSpan, SpanContext, TraceFlags, TraceState + +from app import telemetry + + +def test_telemetry_is_fail_open_without_endpoint(monkeypatch) -> None: + monkeypatch.delenv("OTEL_EXPORTER_OTLP_ENDPOINT", raising=False) + monkeypatch.setattr(telemetry, "_runtime", None) + + assert telemetry.init_telemetry("api") is None + + +def test_resource_uses_canonical_namespace(monkeypatch) -> None: + monkeypatch.setenv("APP_ENV", "test") + resource = telemetry._resource("worker") + + assert resource.attributes["service.name"] == "message-safety" + assert resource.attributes["service.namespace"] == "han-chat" + assert resource.attributes["deployment.environment"] == "test" + assert resource.attributes["process.role"] == "worker" + + +def test_json_stdout_redacts_sensitive_data_and_adds_trace_context() -> None: + record = logging.LogRecord( + "message_safety", + logging.INFO, + __file__, + 1, + "Bearer top-secret user@example.org +79991234567 https://host/path?token=x", + (), + None, + ) + record.event = "redaction_canary" + record.authorization = "secret" + telemetry.RedactionFilter().filter(record) + context = SpanContext( + trace_id=1, + span_id=2, + is_remote=False, + trace_flags=TraceFlags(1), + trace_state=TraceState(), + ) + with trace.use_span(NonRecordingSpan(context)): + payload = json.loads(telemetry.JsonFormatter().format(record)) + + serialized = json.dumps(payload) + assert "top-secret" not in serialized + assert "user@example.org" not in serialized + assert "79991234567" not in serialized + assert "?token=x" not in serialized + assert payload["trace_id"] == f"{1:032x}" + assert payload["span_id"] == f"{2:016x}" + assert payload["service.name"] == "message-safety" + + +def test_span_attributes_are_redacted_before_export() -> None: + attributes = { + "url.full": "https://storage.example/private?token=secret", + "db.statement": "SELECT private_value FROM safety_tasks", + "safe.outcome": "allow", + } + + telemetry._redact_span_attributes(attributes) + + assert attributes["url.full"] == "[REDACTED]" + assert attributes["db.statement"] == "[REDACTED]" + assert attributes["safe.outcome"] == "allow" + + +def test_span_processor_replaces_immutable_sdk_attributes() -> None: + exporter = InMemorySpanExporter() + processor = telemetry.RedactingBatchSpanProcessor(exporter) + span = ReadableSpan( + name="database.connect", + attributes={"db.statement": "SELECT secret", "safe.outcome": "allow"}, + ) + + processor.on_end(span) + processor.shutdown() + + assert span.attributes["db.statement"] == "[REDACTED]" + assert span.attributes["safe.outcome"] == "allow" + + +def test_task_link_rebuilds_valid_remote_context() -> None: + task = SimpleNamespace( + origin_trace_id=(123).to_bytes(16, "big"), + origin_span_id=(456).to_bytes(8, "big"), + origin_trace_flags=1, + origin_tracestate=None, + ) + + link = telemetry.task_link(task) + + assert link is not None + assert link.context.is_remote + assert link.context.trace_id == 123 + assert link.context.span_id == 456 + + +def test_current_trace_fields_round_trip_into_task_link() -> None: + context = SpanContext( + trace_id=123, + span_id=456, + is_remote=False, + trace_flags=TraceFlags(1), + trace_state=TraceState(), + ) + with trace.use_span(NonRecordingSpan(context)): + trace_id, span_id, flags, tracestate = telemetry.current_trace_fields() + + link = telemetry.task_link( + SimpleNamespace( + origin_trace_id=trace_id, + origin_span_id=span_id, + origin_trace_flags=flags, + origin_tracestate=tracestate, + ) + ) + assert link is not None + assert link.context.trace_id == context.trace_id + assert link.context.span_id == context.span_id diff --git a/VM2_services/codebase/services/nginx/templates/10-vm2.conf.template b/VM2_services/codebase/services/nginx/templates/10-vm2.conf.template index 248391e..31e763f 100644 --- a/VM2_services/codebase/services/nginx/templates/10-vm2.conf.template +++ b/VM2_services/codebase/services/nginx/templates/10-vm2.conf.template @@ -14,6 +14,10 @@ server { access_log off; return 200 "ok\n"; } + location = /stub_status { + access_log off; + stub_status; + } location / { return 404; } diff --git a/VM2_services/codebase/services/observability/otel-collector.yaml b/VM2_services/codebase/services/observability/otel-collector.yaml index 361aa57..b0d0672 100644 --- a/VM2_services/codebase/services/observability/otel-collector.yaml +++ b/VM2_services/codebase/services/observability/otel-collector.yaml @@ -11,6 +11,43 @@ receivers: endpoint: 0.0.0.0:4317 http: endpoint: 0.0.0.0:4318 + prometheus: + config: + scrape_configs: + - job_name: otel-collector + scrape_interval: 30s + static_configs: + - targets: ["127.0.0.1:8888"] + labels: + service_name: otel-collector + - job_name: redis-safety + scrape_interval: 30s + static_configs: + - targets: ["redis-exporter:9121"] + labels: + service_name: redis + - job_name: nginx + scrape_interval: 30s + static_configs: + - targets: ["nginx-exporter:9113"] + labels: + service_name: nginx + hostmetrics: + root_path: /hostfs + collection_interval: 30s + scrapers: + cpu: + disk: + filesystem: + exclude_mount_points: + mount_points: + - /hostfs/(dev|proc|sys|run)($|/) + match_type: regexp + load: + memory: + network: + paging: + processes: processors: memory_limiter: @@ -19,7 +56,7 @@ processors: spike_limit_mib: 96 resource/vm2: attributes: - - {key: service.namespace, value: han-processing, action: upsert} + - {key: service.namespace, value: han-chat, action: upsert} - {key: deployment.environment, value: "${env:APP_ENV}", action: upsert} - {key: service.version, value: "${env:RELEASE_VERSION}", action: upsert} attributes/redact: @@ -39,6 +76,32 @@ processors: - {key: messaging.message.body, action: delete} - {key: aws.s3.key, action: delete} - {key: s3.object.key, action: delete} + filter/noise: + error_mode: ignore + traces: + span: + - 'attributes["http.route"] == "/health/live"' + - 'attributes["http.route"] == "/nginx-health/live"' + logs: + log_record: + - 'severity_number < SEVERITY_NUMBER_INFO' + tail_sampling: + decision_wait: 10s + num_traces: 10000 + expected_new_traces_per_sec: 50 + policies: + - name: errors + type: status_code + status_code: + status_codes: [ERROR] + - name: slow + type: latency + latency: + threshold_ms: 1000 + - name: baseline + type: probabilistic + probabilistic: + sampling_percentage: 10 batch: timeout: 5s send_batch_size: 1024 @@ -64,13 +127,16 @@ service: pipelines: traces: receivers: [otlp] - processors: [memory_limiter, resource/vm2, attributes/redact, batch] + processors: [memory_limiter, resource/vm2, attributes/redact, filter/noise, tail_sampling, batch] exporters: [otlp/remote] metrics: - receivers: [otlp] + receivers: [otlp, prometheus, hostmetrics] processors: [memory_limiter, resource/vm2, attributes/redact, batch] exporters: [otlp/remote] logs: receivers: [otlp] - processors: [memory_limiter, resource/vm2, attributes/redact, batch] + processors: [memory_limiter, resource/vm2, attributes/redact, filter/noise, batch] exporters: [otlp/remote] + telemetry: + metrics: + address: 0.0.0.0:8888 diff --git a/VM2_services/codebase/services/redis/redis-safety.acl.template b/VM2_services/codebase/services/redis/redis-safety.acl.template index bc5fe94..4540df5 100644 --- a/VM2_services/codebase/services/redis/redis-safety.acl.template +++ b/VM2_services/codebase/services/redis/redis-safety.acl.template @@ -3,3 +3,6 @@ user default reset on nopass -@all +ping # Runtime cache user: replace the placeholder with a long random password. user safety reset on >REPLACE_WITH_LONG_RANDOM_PASSWORD ~han:safety:* -@all +ping +get +set + +# Metrics user: use the password stored in the REDIS_EXPORTER_PASSWORD JSON map. +user exporter reset on >REPLACE_WITH_EXPORTER_RANDOM_PASSWORD -@all +ping +info diff --git a/VM2_services/codebase/services/tests/test_observability_infrastructure.py b/VM2_services/codebase/services/tests/test_observability_infrastructure.py new file mode 100644 index 0000000..66643d2 --- /dev/null +++ b/VM2_services/codebase/services/tests/test_observability_infrastructure.py @@ -0,0 +1,121 @@ +from __future__ import annotations + +import json +import re +from pathlib import Path + + +ROOT = Path(__file__).resolve().parents[1] + + +def read(relative: str) -> str: + return (ROOT / relative).read_text(encoding="utf-8") + + +def service_block(compose: str, name: str) -> str: + match = re.search( + rf"^ {re.escape(name)}:\n(?P(?:^(?: |\s*$).*\n?)*)", + compose, + re.MULTILINE, + ) + assert match, f"missing Compose service {name}" + return match.group("body") + + +def test_collector_has_only_approved_prometheus_targets() -> None: + config = read("observability/otel-collector.yaml") + + targets = re.findall(r'targets: \["([^"]+)"\]', config) + assert targets == [ + "127.0.0.1:8888", + "redis-exporter:9121", + "nginx-exporter:9113", + ] + assert "service_name: otel-collector" in config + assert "service_name: redis" in config + assert "service_name: nginx" in config + assert "hostmetrics:" in config + assert "service.namespace, value: han-chat" in config + assert "filter/noise:" in config + assert "tail_sampling:" in config + + +def test_logs_are_direct_otlp_without_filelog() -> None: + config = read("observability/otel-collector.yaml") + logs = config.split(" logs:\n", 1)[1] + + assert "receivers: [otlp]" in logs + assert "otlp/remote" in logs + assert "filelog" not in config + + +def test_apps_receive_release_and_logging_context() -> None: + compose = read("docker-compose.yml") + + assert compose.count("RELEASE_VERSION: ${RELEASE_VERSION:?set RELEASE_VERSION}") >= 3 + assert compose.count("LOG_LEVEL: ${LOG_LEVEL:-INFO}") == 2 + + +def test_collector_exposes_health_and_self_metrics_and_validates_config() -> None: + compose = read("docker-compose.yml") + collector = service_block(compose, "otel-collector") + + assert 'expose: ["4317", "4318", "13133", "8888"]' in collector + assert '"/otelcol-contrib", "validate"' in collector + + +def test_exporters_are_internal_only_and_digest_pinned() -> None: + compose = read("docker-compose.yml") + + for name, port in (("redis-exporter", "9121"), ("nginx-exporter", "9113")): + block = service_block(compose, name) + assert f'expose: ["{port}"]' in block + assert "\n ports:" not in block + assert "observability" in block + assert "${REDIS_EXPORTER_IMAGE:?set immutable Redis exporter image digest}" in compose + assert "${NGINX_EXPORTER_IMAGE:?set immutable nginx exporter image digest}" in compose + + +def test_nginx_status_is_internal_and_query_free() -> None: + template = read("nginx/templates/10-vm2.conf.template") + + status = template.split("location = /stub_status", 1)[1].split("}", 1)[0] + assert "stub_status;" in status + assert "access_log off;" in status + assert "listen 8081;" in template + assert ":8081:8081" not in read("docker-compose.yml") + + +def test_redis_exporter_uses_restricted_acl_and_secret_file() -> None: + acl = read("redis/redis-safety.acl.template") + compose = read("docker-compose.yml") + secret_config = json.loads(read("deployment/secrets/config.example.json")) + + exporter_line = next(line for line in acl.splitlines() if line.startswith("user exporter ")) + assert "-@all +ping +info" in exporter_line + assert "~*" not in exporter_line + assert "REDIS_PASSWORD_FILE: /run/secrets/redis_exporter_password" in compose + assert secret_config["secrets"]["REDIS_EXPORTER_PASSWORD"]["consumers"] == [ + "redis-exporter" + ] + + +def test_self_hosted_signoz_does_not_require_fake_auth() -> None: + collector = read("observability/otel-collector.yaml") + env_example = read(".env.example") + + assert "OTEL_REMOTE_AUTH" not in collector + assert "OTEL_REMOTE_AUTH" not in env_example + assert "headers:" not in collector + + +def test_preflight_covers_observability_images_and_acl_pair() -> None: + preflight = read("deployment/preflight.sh") + + assert "REDIS_EXPORTER_IMAGE NGINX_EXPORTER_IMAGE" in preflight + assert "REDIS_EXPORTER_PASSWORD" in preflight + assert 'target = "redis://redis-safety:6379"' in preflight + assert "Redis exporter password map does not match REDIS_SAFETY_ACL" in preflight + assert "Collector must enforce service.namespace=han-chat" in preflight + assert "Collector direct OTLP logs pipeline is missing" in preflight + assert "Collector filelog receiver is forbidden" in preflight