ВМ1: реализован сбор логов
This commit is contained in:
@@ -206,6 +206,7 @@ class SafetyTask(Base):
|
||||
message_id: Mapped[uuid.UUID] = mapped_column(ForeignKey(f"{SCHEMA}.messages.id"), unique=True)
|
||||
attachment_id: Mapped[uuid.UUID | None] = mapped_column(UUID(as_uuid=True))
|
||||
quarantine_object_key: Mapped[str | None] = mapped_column(String(1024))
|
||||
traceparent: Mapped[str | None] = mapped_column(String(55))
|
||||
status: Mapped[str] = mapped_column(String(16))
|
||||
processing_mode: Mapped[str | None] = mapped_column(String(16))
|
||||
config_version: Mapped[int | None] = mapped_column(BigInteger)
|
||||
@@ -228,6 +229,7 @@ class DeliveryOutbox(Base):
|
||||
message_id: Mapped[uuid.UUID] = mapped_column(ForeignKey(f"{SCHEMA}.messages.id"), unique=True)
|
||||
external_chat_id: Mapped[uuid.UUID] = mapped_column(UUID(as_uuid=True))
|
||||
payload_json: Mapped[dict[str, Any]] = mapped_column(JSON)
|
||||
traceparent: Mapped[str | None] = mapped_column(String(55))
|
||||
status: Mapped[str] = mapped_column(String(16), default="pending")
|
||||
attempt_count: Mapped[int] = mapped_column(Integer, default=0)
|
||||
next_attempt_at: Mapped[datetime] = mapped_column(DateTime(timezone=True))
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import re
|
||||
from collections.abc import Mapping
|
||||
from collections.abc import Mapping, MutableMapping
|
||||
from typing import Any
|
||||
|
||||
REDACTED = "[REDACTED]"
|
||||
@@ -42,6 +42,6 @@ def sanitize_value(value: Any) -> Any:
|
||||
def redact_event(
|
||||
_logger: Any,
|
||||
_method_name: str,
|
||||
event_dict: dict[str, Any],
|
||||
) -> dict[str, Any]:
|
||||
event_dict: MutableMapping[str, Any],
|
||||
) -> MutableMapping[str, Any]:
|
||||
return sanitize_value(event_dict)
|
||||
|
||||
@@ -89,11 +89,22 @@ from app.services import (
|
||||
start_session,
|
||||
)
|
||||
from app.settings import get_settings
|
||||
from app.telemetry import add_trace_context, current_trace_id, init_telemetry, instrument_fastapi
|
||||
from app.telemetry import (
|
||||
TelemetryRuntime,
|
||||
add_trace_context,
|
||||
current_trace_id,
|
||||
init_telemetry,
|
||||
instrument_fastapi,
|
||||
)
|
||||
|
||||
|
||||
def configure_logging(level: str) -> None:
|
||||
def configure_logging(level: str, telemetry: TelemetryRuntime | None = None) -> None:
|
||||
logging.basicConfig(level=level, format="%(message)s")
|
||||
if telemetry and telemetry.logging_handler not in logging.getLogger().handlers:
|
||||
telemetry.logging_handler.addFilter(
|
||||
lambda record: not record.name.startswith("opentelemetry")
|
||||
)
|
||||
logging.getLogger().addHandler(telemetry.logging_handler)
|
||||
structlog.configure(
|
||||
processors=[
|
||||
structlog.contextvars.merge_contextvars,
|
||||
@@ -102,7 +113,10 @@ def configure_logging(level: str) -> None:
|
||||
structlog.processors.TimeStamper(fmt="iso", utc=True, key="timestamp"),
|
||||
structlog.stdlib.add_log_level,
|
||||
structlog.processors.JSONRenderer(),
|
||||
]
|
||||
],
|
||||
logger_factory=structlog.stdlib.LoggerFactory(),
|
||||
wrapper_class=structlog.stdlib.BoundLogger,
|
||||
cache_logger_on_first_use=True,
|
||||
)
|
||||
|
||||
|
||||
@@ -135,7 +149,7 @@ async def refresh_jwks_cache(app: FastAPI) -> None:
|
||||
async def lifespan(app: FastAPI):
|
||||
settings = get_settings()
|
||||
telemetry = init_telemetry()
|
||||
configure_logging(settings.log_level)
|
||||
configure_logging(settings.log_level, telemetry)
|
||||
app.state.settings = settings
|
||||
app.state.db = Database(settings.database_url)
|
||||
app.state.http = httpx.AsyncClient()
|
||||
@@ -178,7 +192,7 @@ async def lifespan(app: FastAPI):
|
||||
telemetry.shutdown()
|
||||
|
||||
|
||||
EXPECTED_API_DB_REVISION = "0012_safety_v2_checkpoint"
|
||||
EXPECTED_API_DB_REVISION = "0013_delivery_trace_context"
|
||||
|
||||
|
||||
app = FastAPI(
|
||||
|
||||
@@ -58,6 +58,7 @@ from app.schemas import (
|
||||
encode_cursor,
|
||||
)
|
||||
from app.settings import Settings
|
||||
from app.telemetry import current_traceparent
|
||||
|
||||
MESSAGE_SAFETY_REPLIES = {
|
||||
"text": (
|
||||
@@ -95,6 +96,7 @@ async def ensure_delivery_outbox(
|
||||
external_chat_id: uuid.UUID,
|
||||
payload_json: dict[str, Any],
|
||||
next_attempt_at: datetime,
|
||||
traceparent: str | None = None,
|
||||
) -> DeliveryOutbox:
|
||||
"""Create the per-message outbox row or return the concurrent winner."""
|
||||
statement = (
|
||||
@@ -104,6 +106,7 @@ async def ensure_delivery_outbox(
|
||||
message_id=message_id,
|
||||
external_chat_id=external_chat_id,
|
||||
payload_json=payload_json,
|
||||
traceparent=traceparent,
|
||||
next_attempt_at=next_attempt_at,
|
||||
)
|
||||
.on_conflict_do_nothing(index_elements=[DeliveryOutbox.message_id])
|
||||
@@ -1039,6 +1042,7 @@ async def send_message(
|
||||
message_id=message.id,
|
||||
attachment_id=attachment.id if attachment else None,
|
||||
quarantine_object_key=attachment.quarantine_object_key if attachment else None,
|
||||
traceparent=current_traceparent(),
|
||||
status="polling",
|
||||
processing_mode=verdict["processing_mode"],
|
||||
config_version=verdict["config_version"],
|
||||
@@ -1133,6 +1137,7 @@ async def send_message(
|
||||
# delivery worker race the synchronous first attempt.
|
||||
next_attempt_at=datetime.now(UTC)
|
||||
+ timedelta(seconds=settings.bitrix_local_app_http_timeout_sec + 5),
|
||||
traceparent=current_traceparent(),
|
||||
)
|
||||
await session.commit()
|
||||
await publish_message_status(fanout, message, settings)
|
||||
|
||||
@@ -1,11 +1,15 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import os
|
||||
from collections.abc import MutableMapping
|
||||
from dataclasses import dataclass
|
||||
from typing import Any
|
||||
|
||||
from fastapi import FastAPI
|
||||
from opentelemetry import metrics, trace
|
||||
from opentelemetry._logs import set_logger_provider
|
||||
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
|
||||
@@ -13,7 +17,9 @@ 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.propagate import extract, inject, 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
|
||||
@@ -27,8 +33,11 @@ from opentelemetry.trace.propagation.tracecontext import TraceContextTextMapProp
|
||||
class TelemetryRuntime:
|
||||
tracer_provider: TracerProvider
|
||||
meter_provider: MeterProvider
|
||||
logger_provider: LoggerProvider
|
||||
logging_handler: LoggingHandler
|
||||
|
||||
def shutdown(self) -> None:
|
||||
self.logger_provider.shutdown()
|
||||
self.meter_provider.shutdown()
|
||||
self.tracer_provider.shutdown()
|
||||
|
||||
@@ -79,11 +88,29 @@ def init_telemetry(service_name: str | None = None) -> TelemetryRuntime | None:
|
||||
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,
|
||||
schedule_delay_millis=5000,
|
||||
max_export_batch_size=512,
|
||||
export_timeout_millis=3000,
|
||||
)
|
||||
)
|
||||
set_logger_provider(logger_provider)
|
||||
logging_handler = LoggingHandler(level=logging.NOTSET, logger_provider=logger_provider)
|
||||
|
||||
HTTPXClientInstrumentor().instrument()
|
||||
SQLAlchemyInstrumentor().instrument(enable_commenter=False)
|
||||
RedisInstrumentor().instrument()
|
||||
BotocoreInstrumentor().instrument()
|
||||
_runtime = TelemetryRuntime(tracer_provider, meter_provider)
|
||||
_runtime = TelemetryRuntime(
|
||||
tracer_provider,
|
||||
meter_provider,
|
||||
logger_provider,
|
||||
logging_handler,
|
||||
)
|
||||
return _runtime
|
||||
|
||||
|
||||
@@ -97,8 +124,8 @@ def instrument_fastapi(app: FastAPI) -> None:
|
||||
def add_trace_context(
|
||||
_logger: Any,
|
||||
_method_name: str,
|
||||
event_dict: dict[str, Any],
|
||||
) -> dict[str, Any]:
|
||||
event_dict: MutableMapping[str, Any],
|
||||
) -> MutableMapping[str, Any]:
|
||||
context = trace.get_current_span().get_span_context()
|
||||
if context.is_valid:
|
||||
event_dict["trace_id"] = format(context.trace_id, "032x")
|
||||
@@ -109,3 +136,17 @@ def add_trace_context(
|
||||
def current_trace_id() -> str | None:
|
||||
context = trace.get_current_span().get_span_context()
|
||||
return format(context.trace_id, "032x") if context.is_valid else None
|
||||
|
||||
|
||||
def current_traceparent() -> str | None:
|
||||
carrier: dict[str, str] = {}
|
||||
inject(carrier)
|
||||
return carrier.get("traceparent")
|
||||
|
||||
|
||||
def origin_links(traceparent: str | None) -> list[trace.Link]:
|
||||
if not traceparent:
|
||||
return []
|
||||
context = extract({"traceparent": traceparent})
|
||||
span_context = trace.get_current_span(context).get_span_context()
|
||||
return [trace.Link(span_context)] if span_context.is_valid else []
|
||||
|
||||
@@ -1,12 +1,14 @@
|
||||
import asyncio
|
||||
import logging
|
||||
import uuid
|
||||
from collections.abc import AsyncIterator
|
||||
from collections.abc import AsyncIterator, Awaitable, Callable
|
||||
from contextlib import asynccontextmanager
|
||||
from datetime import UTC, datetime, timedelta
|
||||
|
||||
import httpx
|
||||
import redis.asyncio as redis
|
||||
import structlog
|
||||
from opentelemetry import trace
|
||||
from sqlalchemy import delete, select
|
||||
|
||||
from app.db import (
|
||||
@@ -25,6 +27,7 @@ from app.integrations import (
|
||||
SafetyClient,
|
||||
fresh_openlines_payload,
|
||||
)
|
||||
from app.logging_security import redact_event
|
||||
from app.notification_models import ClientUploadDraft
|
||||
from app.notification_service import expire_notifications
|
||||
from app.realtime import RealtimeFanout
|
||||
@@ -37,10 +40,45 @@ from app.services import (
|
||||
publish_message_status,
|
||||
)
|
||||
from app.settings import Settings, get_settings
|
||||
from app.telemetry import TelemetryRuntime, add_trace_context, init_telemetry, origin_links
|
||||
|
||||
log = structlog.get_logger()
|
||||
|
||||
|
||||
def configure_logging(level: str, telemetry: TelemetryRuntime | None = None) -> None:
|
||||
logging.basicConfig(level=level, format="%(message)s")
|
||||
if telemetry and telemetry.logging_handler not in logging.getLogger().handlers:
|
||||
telemetry.logging_handler.addFilter(
|
||||
lambda record: not record.name.startswith("opentelemetry")
|
||||
)
|
||||
logging.getLogger().addHandler(telemetry.logging_handler)
|
||||
structlog.configure(
|
||||
processors=[
|
||||
structlog.contextvars.merge_contextvars,
|
||||
add_trace_context,
|
||||
redact_event,
|
||||
structlog.processors.TimeStamper(fmt="iso", utc=True, key="timestamp"),
|
||||
structlog.stdlib.add_log_level,
|
||||
structlog.processors.JSONRenderer(),
|
||||
],
|
||||
logger_factory=structlog.stdlib.LoggerFactory(),
|
||||
wrapper_class=structlog.stdlib.BoundLogger,
|
||||
cache_logger_on_first_use=True,
|
||||
)
|
||||
|
||||
|
||||
def run_worker(service_name: str, target: Callable[[], Awaitable[None]]) -> None:
|
||||
settings = get_settings()
|
||||
telemetry = init_telemetry(service_name)
|
||||
configure_logging(settings.log_level, telemetry)
|
||||
structlog.contextvars.bind_contextvars(**{"service.name": service_name})
|
||||
try:
|
||||
asyncio.run(target())
|
||||
finally:
|
||||
if telemetry:
|
||||
telemetry.shutdown()
|
||||
|
||||
|
||||
@asynccontextmanager
|
||||
async def worker_http_clients(
|
||||
settings: Settings,
|
||||
@@ -61,6 +99,7 @@ async def delivery_once(
|
||||
worker_id: str,
|
||||
batch_size: int = 20,
|
||||
) -> int:
|
||||
tracer = trace.get_tracer("han.api.delivery-worker")
|
||||
async with db.sessions() as session:
|
||||
rows = (
|
||||
(
|
||||
@@ -88,33 +127,43 @@ async def delivery_once(
|
||||
row = await session.get(DeliveryOutbox, row_id, with_for_update=True)
|
||||
if row is None:
|
||||
continue
|
||||
message = None
|
||||
dialog = None
|
||||
try:
|
||||
payload = await fresh_openlines_payload(row.payload_json, s3)
|
||||
await client.send(row.message_id, payload, f"worker-{worker_id}")
|
||||
row.status = "delivered"
|
||||
message = await session.get(Message, row.message_id)
|
||||
with (
|
||||
tracer.start_as_current_span(
|
||||
"delivery.process",
|
||||
links=origin_links(row.traceparent),
|
||||
),
|
||||
structlog.contextvars.bound_contextvars(
|
||||
delivery_outbox_id=str(row.id),
|
||||
message_id=str(row.message_id),
|
||||
),
|
||||
):
|
||||
message = None
|
||||
dialog = None
|
||||
try:
|
||||
payload = await fresh_openlines_payload(row.payload_json, s3)
|
||||
await client.send(row.message_id, payload, f"worker-{worker_id}")
|
||||
row.status = "delivered"
|
||||
message = await session.get(Message, row.message_id)
|
||||
if message:
|
||||
message.delivery_status = "delivered"
|
||||
dialog = await session.get(Dialog, message.dialog_id)
|
||||
if dialog:
|
||||
dialog.status = "waiting_for_company"
|
||||
dialog.last_message_at = datetime.now(UTC)
|
||||
except DependencyFailure:
|
||||
row.attempt_count += 1
|
||||
row.status = "dead_letter" if row.attempt_count >= 12 else "retry"
|
||||
row.next_attempt_at = datetime.now(UTC) + timedelta(
|
||||
seconds=min(3600, 2**row.attempt_count)
|
||||
)
|
||||
row.last_error_code = "dependency_unavailable"
|
||||
row.locked_at = None
|
||||
row.locked_by = None
|
||||
await session.commit()
|
||||
if message:
|
||||
message.delivery_status = "delivered"
|
||||
dialog = await session.get(Dialog, message.dialog_id)
|
||||
if dialog:
|
||||
dialog.status = "waiting_for_company"
|
||||
dialog.last_message_at = datetime.now(UTC)
|
||||
except DependencyFailure:
|
||||
row.attempt_count += 1
|
||||
row.status = "dead_letter" if row.attempt_count >= 12 else "retry"
|
||||
row.next_attempt_at = datetime.now(UTC) + timedelta(
|
||||
seconds=min(3600, 2**row.attempt_count)
|
||||
)
|
||||
row.last_error_code = "dependency_unavailable"
|
||||
row.locked_at = None
|
||||
row.locked_by = None
|
||||
await session.commit()
|
||||
if message:
|
||||
await publish_message_status(fanout, message, settings)
|
||||
if dialog:
|
||||
await publish_dialog_status(fanout, dialog)
|
||||
await publish_message_status(fanout, message, settings)
|
||||
if dialog:
|
||||
await publish_dialog_status(fanout, dialog)
|
||||
return len(ids)
|
||||
|
||||
|
||||
@@ -127,6 +176,7 @@ async def safety_once(
|
||||
worker_id: str,
|
||||
batch_size: int = 20,
|
||||
) -> int:
|
||||
tracer = trace.get_tracer("han.api.safety-recovery-worker")
|
||||
async with db.sessions() as session:
|
||||
rows = (
|
||||
(
|
||||
@@ -155,7 +205,14 @@ async def safety_once(
|
||||
continue
|
||||
message = None
|
||||
try:
|
||||
verdict = await safety.poll(task.poll_location, f"worker-{worker_id}")
|
||||
with tracer.start_as_current_span(
|
||||
"safety.recovery.poll",
|
||||
links=origin_links(task.traceparent),
|
||||
):
|
||||
verdict = await safety.poll(
|
||||
task.poll_location,
|
||||
f"worker-{worker_id}",
|
||||
)
|
||||
message = await session.get(Message, task.message_id)
|
||||
attachment = (
|
||||
await session.get(MessageAttachment, task.attachment_id)
|
||||
@@ -204,6 +261,7 @@ async def safety_once(
|
||||
attachment=attachment,
|
||||
),
|
||||
next_attempt_at=datetime.now(UTC),
|
||||
traceparent=task.traceparent,
|
||||
)
|
||||
elif verdict["_status"] == 403 and message:
|
||||
message.safety_processing_mode = verdict["processing_mode"]
|
||||
@@ -365,20 +423,20 @@ async def notification_draft_cleanup_loop() -> None:
|
||||
|
||||
|
||||
def delivery_main() -> None:
|
||||
asyncio.run(loop("delivery"))
|
||||
run_worker("delivery-worker", lambda: loop("delivery"))
|
||||
|
||||
|
||||
def safety_main() -> None:
|
||||
asyncio.run(loop("safety"))
|
||||
run_worker("safety-recovery-worker", lambda: loop("safety"))
|
||||
|
||||
|
||||
def cleanup_main() -> None:
|
||||
asyncio.run(loop("cleanup"))
|
||||
run_worker("cleanup-worker", lambda: loop("cleanup"))
|
||||
|
||||
|
||||
def notification_expire_main() -> None:
|
||||
asyncio.run(notification_expire_loop())
|
||||
run_worker("notification-expire-worker", notification_expire_loop)
|
||||
|
||||
|
||||
def notification_draft_cleanup_main() -> None:
|
||||
asyncio.run(notification_draft_cleanup_loop())
|
||||
run_worker("notification-draft-cleanup-worker", notification_draft_cleanup_loop)
|
||||
|
||||
Reference in New Issue
Block a user