Files

495 lines
18 KiB
Python

from __future__ import annotations
import asyncio
import logging
import os
from datetime import datetime, timezone
from redis.asyncio import Redis
from app.audit_logger import AuditLogger
from app.config import settings
from app.correlation import CorrelationRegistry
from app.llm_remediation import LLMRemediationAdapter
from app.llm_triage import LLMTriageAdapter
from app.mail_notifier import MailNotifier
from app.matrix_notifier import MatrixNotifier
from app.matrix_token_manager import MatrixTokenManager
from app.notifications.dispatcher import NotificationDispatcher
from app.processor_service import ProcessorService
from app.queue_repo import RedisQueueRepository
from app.redis_repo import RedisStateRepository
from app.zabbix_client import ZabbixApiClient
from app.zabbix_enricher import ZabbixEnricher
from app.llm_correlation import LLMCorrelationAdapter
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s %(levelname)s %(name)s %(message)s",
)
logger = logging.getLogger(__name__)
def _utc_now() -> str:
return datetime.now(timezone.utc).isoformat()
def _worker_status_key() -> str:
return f"{settings.redis_key_prefix}:worker:status"
async def _publish_worker_status(
redis_client: Redis,
**fields: str | int | bool | None,
) -> None:
key = _worker_status_key()
mapping: dict[str, str] = {}
for k, v in fields.items():
if v is None:
continue
if isinstance(v, bool):
mapping[k] = "true" if v else "false"
else:
mapping[k] = str(v)
mapping["heartbeat_at"] = _utc_now()
if mapping:
await redis_client.hset(key, mapping=mapping)
async def _increment_worker_counter(
redis_client: Redis,
field_name: str,
) -> None:
key = _worker_status_key()
await redis_client.hincrby(key, field_name, 1)
await redis_client.hset(key, mapping={"heartbeat_at": _utc_now()})
async def run_worker() -> None:
if not settings.redis_enabled:
raise RuntimeError("REDIS_ENABLED must be true for worker mode")
if not settings.queue_enabled:
raise RuntimeError("QUEUE_ENABLED must be true for worker mode")
redis_client = Redis.from_url(
settings.redis_url,
encoding="utf-8",
decode_responses=False,
)
queue_repo = RedisQueueRepository(
client=redis_client,
queue_name=settings.queue_name,
processing_name=settings.queue_processing_name,
deadletter_name=settings.queue_deadletter_name,
dedup_ttl_seconds=settings.queue_dedup_ttl_seconds,
)
await queue_repo.ping()
audit_logger = None
if settings.audit_enabled:
audit_logger = AuditLogger(
client=redis_client,
key_prefix=settings.audit_key_prefix,
ttl_seconds=settings.audit_ttl_seconds,
max_stage_records=settings.audit_max_stage_records,
)
await _publish_worker_status(
redis_client,
state="starting",
started_at=_utc_now(),
pid=os.getpid(),
queue_name=settings.queue_name,
queue_processing_name=settings.queue_processing_name,
queue_deadletter_name=settings.queue_deadletter_name,
matrix_enabled=settings.matrix_enabled,
mail_enabled=settings.mail_enabled,
zabbix_api_enabled=settings.zabbix_api_enabled,
llm_enabled=settings.llm_enabled,
llm_triage_enabled=settings.llm_triage_enabled,
audit_enabled=settings.audit_enabled,
matrix_initialized=False,
mail_initialized=False,
zabbix_initialized=False,
llm_remediation_initialized=False,
llm_triage_initialized=False,
llm_correlation_initialized=False,
processed_count=0,
failed_count=0,
requeued_count=0,
deadletter_count=0,
)
if settings.queue_requeue_processing_on_startup:
moved = await queue_repo.requeue_processing_messages()
if moved:
logger.info("Requeued %s messages from processing list back to main queue", moved)
await _publish_worker_status(
redis_client,
processing_requeued_on_startup=moved,
last_requeue_at=_utc_now(),
)
redis_repo = RedisStateRepository(
client=redis_client,
key_prefix=settings.redis_key_prefix,
fingerprint_ttl_seconds=settings.redis_fingerprint_ttl_seconds,
event_ttl_seconds=settings.redis_event_ttl_seconds,
suppress_window_seconds=settings.suppress_window_seconds,
flap_window_seconds=settings.flap_window_seconds,
flap_threshold=settings.flap_threshold,
triage_cache_ttl_seconds=settings.llm_triage_cache_ttl_seconds,
)
matrix_token_manager = None
matrix_notifier = None
if settings.matrix_enabled:
matrix_token_manager = MatrixTokenManager(
token_endpoint=settings.matrix_oauth_token_endpoint,
client_id=settings.matrix_oauth_client_id,
client_secret=settings.matrix_oauth_client_secret or None,
initial_access_token=settings.matrix_access_token,
initial_refresh_token=settings.matrix_refresh_token,
initial_expires_in_seconds=settings.matrix_access_token_expires_in_seconds,
refresh_margin_seconds=settings.matrix_refresh_margin_seconds,
state_file=settings.matrix_token_state_file,
timeout_seconds=settings.matrix_request_timeout_seconds,
verify_tls=settings.matrix_verify_tls,
)
matrix_token_manager.start_background_refresh()
matrix_notifier = MatrixNotifier(
homeserver_url=settings.matrix_homeserver_url,
room_id=settings.matrix_room_id,
token_manager=matrix_token_manager,
message_type=settings.matrix_message_type,
timeout_seconds=settings.matrix_request_timeout_seconds,
verify_tls=settings.matrix_verify_tls,
)
logger.info("Worker: Matrix notifier initialized")
await _publish_worker_status(redis_client, matrix_initialized=True)
mail_notifier = None
if settings.mail_enabled:
mail_notifier = MailNotifier(
smtp_host=settings.mail_smtp_host,
smtp_port=settings.mail_smtp_port,
username=settings.mail_smtp_username,
password=settings.mail_smtp_password,
from_addr=settings.mail_from,
to_addr=settings.mail_to,
use_starttls=settings.mail_use_starttls,
use_tls=settings.mail_use_tls,
timeout_seconds=settings.mail_timeout_seconds,
)
logger.info("Worker: Mail notifier initialized")
await _publish_worker_status(redis_client, mail_initialized=True)
notification_dispatcher = NotificationDispatcher(
matrix_notifier=matrix_notifier,
mail_notifier=mail_notifier,
)
zabbix_api_client = None
zabbix_enricher = None
if settings.zabbix_api_enabled:
zabbix_api_client = ZabbixApiClient(
api_url=settings.zabbix_api_url,
api_token=settings.zabbix_api_token,
timeout_seconds=settings.zabbix_api_timeout_seconds,
verify_tls=settings.zabbix_api_verify_tls,
)
zabbix_enricher = ZabbixEnricher(
client=zabbix_api_client,
web_url=settings.zabbix_web_url,
graph_period_hours=settings.zabbix_graph_period_hours,
graph_timezone=settings.zabbix_graph_timezone,
)
logger.info("Worker: Zabbix enricher initialized")
await _publish_worker_status(redis_client, zabbix_initialized=True)
llm_remediation_adapter = None
if settings.llm_enabled:
llm_remediation_adapter = LLMRemediationAdapter(
base_url=settings.llm_base_url,
model=settings.llm_model,
timeout_seconds=settings.llm_timeout_seconds,
verify_tls=settings.llm_verify_tls,
temperature=settings.llm_temperature,
max_steps=settings.llm_max_steps,
max_commands=settings.llm_max_commands,
)
logger.info("Worker: LLM remediation adapter initialized")
await _publish_worker_status(redis_client, llm_remediation_initialized=True)
llm_triage_adapter = None
if settings.llm_enabled and settings.llm_triage_enabled:
llm_triage_adapter = LLMTriageAdapter(
base_url=settings.llm_base_url,
model=settings.llm_model,
timeout_seconds=settings.llm_timeout_seconds,
verify_tls=settings.llm_verify_tls,
temperature=settings.llm_temperature,
)
logger.info("Worker: LLM triage adapter initialized")
await _publish_worker_status(redis_client, llm_triage_initialized=True)
llm_correlation_adapter = None
if settings.llm_enabled and settings.llm_correlation_enabled:
llm_correlation_adapter = LLMCorrelationAdapter(
base_url=settings.llm_base_url,
model=settings.llm_model,
timeout_seconds=settings.llm_timeout_seconds,
verify_tls=settings.llm_verify_tls,
temperature=settings.llm_temperature,
)
logger.info("Worker: LLM correlation adapter initialized")
await _publish_worker_status(redis_client, llm_correlation_initialized=True)
correlation_registry = None
if settings.correlation_enabled:
correlation_registry = CorrelationRegistry.from_yaml_files(
event_kind_rules_path=settings.correlation_kind_rules_path,
root_cause_map_path=settings.correlation_root_cause_path,
)
logger.info("Worker: Correlation registry initialized")
processor_service = ProcessorService(
redis_repo=redis_repo,
notification_dispatcher=notification_dispatcher,
zabbix_enricher=zabbix_enricher,
llm_remediation_adapter=llm_remediation_adapter,
llm_triage_adapter=llm_triage_adapter,
correlation_registry=correlation_registry,
audit_logger=audit_logger,
llm_correlation_adapter=llm_correlation_adapter,
)
logger.info("Worker started: queue=%s", settings.queue_name)
await _publish_worker_status(redis_client, state="running")
try:
while True:
await _publish_worker_status(redis_client, state="idle")
message = await queue_repo.claim_event(
timeout_seconds=settings.queue_block_timeout_seconds
)
if message is None:
continue
identity = queue_repo.build_processing_identity(message.envelope)
await _publish_worker_status(
redis_client,
state="processing",
current_job_id=message.job_id,
current_identity=identity,
current_attempt=message.attempt,
current_correlation_id=message.envelope.event.correlation_id,
current_event_id=message.envelope.event.event_id or "",
)
if audit_logger is not None:
await audit_logger.log_worker_started(
envelope=message.envelope,
job_id=message.job_id,
attempt=message.attempt,
identity=identity,
)
if await queue_repo.is_processed(identity):
logger.info(
"Worker skipped already processed message: job_id=%s identity=%s",
message.job_id,
identity,
)
await queue_repo.ack_message(message.raw_message)
if audit_logger is not None:
await audit_logger.log_worker_outcome(
envelope=message.envelope,
state="duplicate_skipped",
attempt=message.attempt,
)
await _publish_worker_status(
redis_client,
state="idle",
last_skipped_job_id=message.job_id,
last_skipped_identity=identity,
current_job_id="",
current_identity="",
current_attempt="",
current_correlation_id="",
current_event_id="",
)
continue
try:
await processor_service.process(message.envelope)
await queue_repo.mark_processed(identity)
await queue_repo.ack_message(message.raw_message)
await _increment_worker_counter(redis_client, "processed_count")
if audit_logger is not None:
await audit_logger.log_worker_outcome(
envelope=message.envelope,
state="processed",
attempt=message.attempt,
)
logger.info(
"Worker processed message successfully: job_id=%s identity=%s attempt=%s",
message.job_id,
identity,
message.attempt,
)
await _publish_worker_status(
redis_client,
state="idle",
last_processed_at=_utc_now(),
last_job_id=message.job_id,
last_identity=identity,
last_attempt=message.attempt,
last_correlation_id=message.envelope.event.correlation_id,
last_event_id=message.envelope.event.event_id or "",
last_error="",
current_job_id="",
current_identity="",
current_attempt="",
current_correlation_id="",
current_event_id="",
)
except Exception as exc:
await _increment_worker_counter(redis_client, "failed_count")
logger.exception(
"Worker failed to process message: job_id=%s identity=%s attempt=%s error=%s",
message.job_id,
identity,
message.attempt,
exc,
)
if audit_logger is not None:
await audit_logger.log_stage(
correlation_id=message.envelope.event.correlation_id,
event_id=message.envelope.event.event_id,
stage="worker_exception",
status="error",
details={
"attempt": message.attempt,
"identity": identity,
"error": str(exc),
},
)
await _publish_worker_status(
redis_client,
state="error",
last_error=str(exc),
last_error_at=_utc_now(),
last_failed_job_id=message.job_id,
last_failed_identity=identity,
last_failed_attempt=message.attempt,
)
if message.attempt + 1 < settings.queue_max_attempts:
await queue_repo.requeue_message(message, error=str(exc))
await _increment_worker_counter(redis_client, "requeued_count")
if audit_logger is not None:
await audit_logger.log_worker_outcome(
envelope=message.envelope,
state="requeued",
error=str(exc),
attempt=message.attempt + 1,
)
logger.warning(
"Worker requeued message: job_id=%s next_attempt=%s",
message.job_id,
message.attempt + 1,
)
await _publish_worker_status(
redis_client,
state="idle",
last_requeue_at=_utc_now(),
current_job_id="",
current_identity="",
current_attempt="",
current_correlation_id="",
current_event_id="",
)
else:
await queue_repo.deadletter_message(message, error=str(exc))
await _increment_worker_counter(redis_client, "deadletter_count")
if audit_logger is not None:
await audit_logger.log_worker_outcome(
envelope=message.envelope,
state="deadletter",
error=str(exc),
attempt=message.attempt + 1,
)
logger.error(
"Worker sent message to deadletter: job_id=%s attempts=%s",
message.job_id,
message.attempt + 1,
)
await _publish_worker_status(
redis_client,
state="idle",
last_deadletter_at=_utc_now(),
current_job_id="",
current_identity="",
current_attempt="",
current_correlation_id="",
current_event_id="",
)
finally:
await _publish_worker_status(redis_client, state="stopping")
if zabbix_api_client is not None:
await zabbix_api_client.close()
if llm_remediation_adapter is not None:
await llm_remediation_adapter.close()
if llm_triage_adapter is not None:
await llm_triage_adapter.close()
if llm_correlation_adapter is not None:
await llm_correlation_adapter.close()
if matrix_notifier is not None:
await matrix_notifier.close()
if matrix_token_manager is not None:
await matrix_token_manager.close()
await _publish_worker_status(redis_client, state="stopped")
await redis_client.close()
def main() -> None:
try:
asyncio.run(run_worker())
except KeyboardInterrupt:
logger.info("Worker stopped by user")
if __name__ == "__main__":
main()