Добавлен OTLP-провайдер, реализовано отбрасывание метрик и трейсов в observability + добавлен перезапуск nginx при пересборке контейнеров (ошибка, когда докер меняет адреса сервисов)
This commit is contained in:
@@ -9,6 +9,8 @@ from datetime import UTC, datetime, timedelta
|
||||
|
||||
import httpx
|
||||
import structlog
|
||||
from opentelemetry import trace
|
||||
from prometheus_client import start_http_server
|
||||
from sqlalchemy import and_, func, or_, select, update
|
||||
|
||||
from app.db import Database, SendStatus, SmsOutboundMessage
|
||||
@@ -23,6 +25,7 @@ from app.metrics import (
|
||||
from app.provider import IdgtlClient, IdgtlConfig
|
||||
from app.service import RuntimeSettings, load_runtime_settings
|
||||
from app.settings import Settings, get_settings
|
||||
from app.telemetry import add_trace_context, init_telemetry
|
||||
|
||||
log = structlog.get_logger()
|
||||
MAX_CONNECT_ATTEMPTS = 3
|
||||
@@ -32,6 +35,8 @@ def configure_logging(level: str) -> None:
|
||||
logging.basicConfig(level=level, format="%(message)s")
|
||||
structlog.configure(
|
||||
processors=[
|
||||
structlog.contextvars.merge_contextvars,
|
||||
add_trace_context,
|
||||
structlog.processors.TimeStamper(fmt="iso", utc=True, key="timestamp"),
|
||||
structlog.stdlib.add_log_level,
|
||||
structlog.processors.JSONRenderer(),
|
||||
@@ -169,7 +174,14 @@ async def update_queue_metrics(db: Database) -> None:
|
||||
|
||||
async def worker_loop(stop: asyncio.Event) -> None:
|
||||
settings = get_settings()
|
||||
telemetry = init_telemetry("sms-worker")
|
||||
configure_logging(settings.log_level)
|
||||
structlog.contextvars.bind_contextvars(**{"service.name": "sms-worker"})
|
||||
metrics_server, metrics_thread = start_http_server(
|
||||
settings.metrics_port,
|
||||
addr="0.0.0.0", # noqa: S104 - internal Docker-network listener
|
||||
)
|
||||
tracer = trace.get_tracer("han.sms.worker")
|
||||
db = Database(settings.database_url)
|
||||
async with httpx.AsyncClient() as http:
|
||||
try:
|
||||
@@ -179,16 +191,30 @@ async def worker_loop(stop: asyncio.Event) -> None:
|
||||
async with db.sessions() as session:
|
||||
runtime = await load_runtime_settings(session)
|
||||
SETTINGS_VALID.set(1)
|
||||
claim_started_ns = time.time_ns()
|
||||
message = await lease_message(db, runtime)
|
||||
claim_finished_ns = time.time_ns()
|
||||
if message is None:
|
||||
await update_queue_metrics(db)
|
||||
await asyncio.wait_for(stop.wait(), timeout=runtime.poll_interval_ms / 1000)
|
||||
continue
|
||||
client = IdgtlClient(http, provider_config(settings, runtime))
|
||||
started = time.monotonic()
|
||||
result = await client.send(message)
|
||||
PROVIDER_LATENCY.labels("idgtl").observe(time.monotonic() - started)
|
||||
await save_result(db, message.id, result, message.attempt_count)
|
||||
process_span = tracer.start_span("sms.process", start_time=claim_started_ns)
|
||||
with trace.use_span(process_span, end_on_exit=True):
|
||||
claim_span = tracer.start_span("sms.claim", start_time=claim_started_ns)
|
||||
claim_span.set_attribute("sms.claimed", True)
|
||||
claim_span.end(end_time=claim_finished_ns)
|
||||
span = trace.get_current_span()
|
||||
span.set_attribute("messaging.operation.name", "send")
|
||||
span.set_attribute("messaging.system", "idgtl")
|
||||
span.set_attribute("sms.attempt", message.attempt_count)
|
||||
client = IdgtlClient(http, provider_config(settings, runtime))
|
||||
with tracer.start_as_current_span("sms.provider"):
|
||||
started = time.monotonic()
|
||||
result = await client.send(message)
|
||||
PROVIDER_LATENCY.labels("idgtl").observe(time.monotonic() - started)
|
||||
span.set_attribute("sms.outcome", result.send_status.value)
|
||||
with tracer.start_as_current_span("sms.save_result"):
|
||||
await save_result(db, message.id, result, message.attempt_count)
|
||||
except TimeoutError:
|
||||
continue
|
||||
except Exception:
|
||||
@@ -200,6 +226,10 @@ async def worker_loop(stop: asyncio.Event) -> None:
|
||||
pass
|
||||
finally:
|
||||
await db.close()
|
||||
metrics_server.shutdown()
|
||||
metrics_thread.join(timeout=5)
|
||||
if telemetry:
|
||||
telemetry.shutdown()
|
||||
|
||||
|
||||
def run() -> None:
|
||||
|
||||
Reference in New Issue
Block a user