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()