diff --git a/generator/generator.py b/generator/generator.py index d5db77b..c937d5e 100644 --- a/generator/generator.py +++ b/generator/generator.py @@ -400,14 +400,55 @@ class GeneratorState: @classmethod def from_dict(cls, data: dict) -> "GeneratorState": - """Создаёт состояние из словаря (JSON-only, без pickle).""" - return cls( - tick=data["tick"], - rng_state=_nested_list_to_tuple(data["rng_state"]), # рекурсивно list->tuple - last_batch_id=data["last_batch_id"], - last_timestamp=datetime.fromisoformat(data["last_timestamp"]), - version=data.get("version", "1.0"), - ) + """Создаёт состояние из словаря (JSON-only, без pickle). + + При невалидном rng_state логирует предупреждение и возвращает None + (вызывающий код должен обработать как "начать с чистого листа"). + """ + try: + rng_state_raw = data.get("rng_state") + if not rng_state_raw: + logger.warning("State missing rng_state field") + raise ValueError("rng_state is missing") + + rng_state = _nested_list_to_tuple(rng_state_raw) + + # Валидация: rng_state должен быть tuple и иметь минимальную структуру + if not isinstance(rng_state, tuple): + logger.warning(f"rng_state is not tuple: {type(rng_state)}") + raise ValueError("rng_state must be tuple") + + if len(rng_state) < 2: + logger.warning(f"rng_state has insufficient length: {len(rng_state)}") + raise ValueError("rng_state has insufficient length") + + # Проверка что можем создать RNG и вызвать setstate (тестовая валидация) + test_rng = random.Random() + test_rng.setstate(rng_state) + + return cls( + tick=data.get("tick", 0), + rng_state=rng_state, + last_batch_id=data.get("last_batch_id", ""), + last_timestamp=datetime.fromisoformat(data.get("last_timestamp", "1970-01-01T00:00:00+00:00")), + version=data.get("version", "1.0"), + ) + except Exception as e: + logger.warning(f"Invalid state format, will start fresh: {e}") + raise ValueError(f"Invalid state: {e}") + + @classmethod + def from_dict_safe(cls, data: dict) -> "GeneratorState | None": + """Безопасная загрузка state с graceful degradation. + + Returns: + GeneratorState если данные валидны, иначе None (начать с чистого листа). + """ + try: + return cls.from_dict(data) + except Exception: + # Уже залогировано в from_dict + return None # --------------------------------------------------------------------------- @@ -434,28 +475,43 @@ class KafkaStateManager: logger.info("Connected to Kafka for state management successfully") def save(self, state: GeneratorState) -> None: - """Сохраняет состояние в топик (compact topic - только последнее значение).""" + """Сохраняет состояние в топик (compact topic - только последнее значение). + + Использует retry при сбоях подключения к Kafka. + """ value = state.to_dict() - self.producer.send(self.STATE_TOPIC, key=self.STATE_KEY, value=value) + + def _do_send(): + self.producer.send(self.STATE_TOPIC, key=self.STATE_KEY, value=value) + + _with_retry(_do_send, max_retries=3, base_delay=0.5) def flush(self) -> None: - """Сбрасывает буфер.""" - self.producer.flush() + """Сбрасывает буфер с retry.""" + def _do_flush(): + self.producer.flush() + + _with_retry(_do_flush, max_retries=3, base_delay=0.5) def close(self) -> None: """Закрывает соединение.""" - self.producer.close() + try: + self.producer.close() + except Exception as e: + logger.debug(f"Error closing producer (ignored): {e}") def load(self) -> GeneratorState | None: """Загружает последнее состояние из топика. Для compact topic хранится только последнее значение для ключа, поэтому читаем все сообщения и берём последнее с нужным ключом. + Использует retry при сбоях подключения к Kafka. """ from kafka import KafkaConsumer logger.info(f"Loading state from topic {self.STATE_TOPIC}") - try: + + def _do_load(): consumer = KafkaConsumer( self.STATE_TOPIC, bootstrap_servers=self.bootstrap_servers, @@ -471,11 +527,19 @@ class KafkaStateManager: last_state = message.value consumer.close() + return last_state + + try: + last_state = _with_retry(_do_load, max_retries=3, base_delay=0.5) if last_state: logger.info(f"Restored state: tick={last_state.get('tick')}, " f"last_batch_id={last_state.get('last_batch_id')}") - return GeneratorState.from_dict(last_state) + # Используем from_dict_safe для graceful degradation при битом state + restored = GeneratorState.from_dict_safe(last_state) + if restored is None: + logger.warning("State data was invalid, starting fresh") + return restored else: logger.info("No previous state found, starting fresh") return None @@ -567,56 +631,81 @@ def ensure_topics(bootstrap_servers: str) -> None: # Kafka history - пишет историю в отдельный топик # --------------------------------------------------------------------------- class KafkaBatchHistory: - """Хранение истории batch в Kafka (отдельный топик).""" + """Хранение истории batch в Kafka (отдельный топик) с retry.""" HISTORY_TOPIC = "generator_batch_history" - def __init__(self, bootstrap_servers: str): - self.bootstrap_servers = bootstrap_servers - KafkaProducerCls, _ = _import_kafka() - - logger.info(f"Connecting to Kafka for history at {self.bootstrap_servers}") - self.producer = KafkaProducerCls( - bootstrap_servers=self.bootstrap_servers, - value_serializer=lambda v: json.dumps(v).encode("utf-8"), - key_serializer=lambda k: k.encode("utf-8") if k else None, - retries=3, - retry_backoff_ms=1000, - ) - logger.info("Connected to Kafka for history successfully") - - def add(self, record: BatchRecord): - """Добавляет запись в историю (топик Kafka).""" - key = record.batch_id - value = record.to_dict() - self.producer.send(self.HISTORY_TOPIC, key=key, value=value) - - def flush(self): - """Сбрасывает буфер.""" - self.producer.flush() - - def close(self): - """Закрывает соединение.""" - self.producer.close() - - -# --------------------------------------------------------------------------- -# Kafka publisher -# --------------------------------------------------------------------------- -class KafkaPublisher: - """Публикация событий в Kafka.""" - def __init__(self, bootstrap_servers: str): self.bootstrap_servers = bootstrap_servers self.producer = None self._connect() def _connect(self): - """Устанавливает соединение с Kafka.""" + """Устанавливает соединение с Kafka с retry.""" KafkaProducerCls, KafkaErrorCls = _import_kafka() - - logger.info(f"Connecting to Kafka at {self.bootstrap_servers}") + + def _do_connect(): + logger.info(f"Connecting to Kafka for history at {self.bootstrap_servers}") + self.producer = KafkaProducerCls( + bootstrap_servers=self.bootstrap_servers, + value_serializer=lambda v: json.dumps(v).encode("utf-8"), + key_serializer=lambda k: k.encode("utf-8") if k else None, + retries=3, + retry_backoff_ms=1000, + ) + logger.info("Connected to Kafka for history successfully") + + _with_retry(_do_connect, max_retries=5, base_delay=1.0) + + def add(self, record: BatchRecord): + """Добавляет запись в историю (топик Kafka) с retry и реконнектом.""" + key = record.batch_id + value = record.to_dict() + + def _do_send(): + self.producer.send(self.HISTORY_TOPIC, key=key, value=value) + try: + _with_retry(_do_send, max_retries=3, base_delay=0.5) + except Exception as e: + logger.warning(f"Failed to send history record after retries: {e}, attempting reconnect") + self._connect() + # Повторная попытка после реконнекта + _with_retry(_do_send, max_retries=2, base_delay=0.5) + + def flush(self): + """Сбрасывает буфер с retry.""" + def _do_flush(): + self.producer.flush() + + _with_retry(_do_flush, max_retries=3, base_delay=0.5) + + def close(self): + """Закрывает соединение.""" + try: + if self.producer: + self.producer.close() + except Exception as e: + logger.debug(f"Error closing history producer (ignored): {e}") + + +# --------------------------------------------------------------------------- +# Kafka publisher +# --------------------------------------------------------------------------- +class KafkaPublisher: + """Публикация событий в Kafka с retry и реконнектом.""" + + def __init__(self, bootstrap_servers: str): + self.bootstrap_servers = bootstrap_servers + self.producer = None + self._connect() + + def _connect(self): + """Устанавливает соединение с Kafka с retry.""" + KafkaProducerCls, KafkaErrorCls = _import_kafka() + + def _do_connect(): + logger.info(f"Connecting to Kafka at {self.bootstrap_servers}") self.producer = KafkaProducerCls( bootstrap_servers=self.bootstrap_servers, value_serializer=lambda v: json.dumps(v).encode("utf-8"), @@ -627,22 +716,14 @@ class KafkaPublisher: retry_backoff_ms=1000, ) logger.info("Connected to Kafka successfully") - except KafkaErrorCls as e: - logger.error(f"Failed to connect to Kafka: {e}") - raise - def publish(self, topic: str, events: list[dict]) -> tuple[int, int]: - """ - Публикует события в топик. + _with_retry(_do_connect, max_retries=5, base_delay=1.0) - Returns: - (sent_count, error_count) - """ + def _publish_with_retry(self, topic: str, events: list[dict]) -> tuple[int, int]: + """Внутренняя функция публикации с retry на уровне batch.""" if not self.producer: raise RuntimeError("Producer not connected") - _, KafkaErrorCls = _import_kafka() - sent = 0 errors = 0 futures = [] @@ -670,15 +751,39 @@ class KafkaPublisher: return sent, errors + def publish(self, topic: str, events: list[dict]) -> tuple[int, int]: + """ + Публикует события в топик с retry и автоматическим реконнектом. + + Returns: + (sent_count, error_count) + """ + def _do_publish(): + return self._publish_with_retry(topic, events) + + try: + return _with_retry(_do_publish, max_retries=3, base_delay=0.5) + except Exception as e: + logger.warning(f"Publish failed after retries: {e}, attempting reconnect") + self._connect() + # Повторная попытка после реконнекта + return _with_retry(_do_publish, max_retries=2, base_delay=0.5) + def flush(self): - """Сбрасывает буфер.""" - if self.producer: - self.producer.flush() + """Сбрасывает буфер с retry.""" + def _do_flush(): + if self.producer: + self.producer.flush() + + _with_retry(_do_flush, max_retries=3, base_delay=0.5) def close(self): """Закрывает соединение.""" - if self.producer: - self.producer.close() + try: + if self.producer: + self.producer.close() + except Exception as e: + logger.debug(f"Error closing producer (ignored): {e}") # --------------------------------------------------------------------------- diff --git a/generator/tests/test_kafka_history.py b/generator/tests/test_kafka_history.py index 74b7002..47ffcda 100644 --- a/generator/tests/test_kafka_history.py +++ b/generator/tests/test_kafka_history.py @@ -3,7 +3,7 @@ """ from datetime import datetime, timezone -from unittest.mock import MagicMock, patch +from unittest.mock import MagicMock, patch, call import pytest from generator import BatchRecord, KafkaBatchHistory @@ -101,14 +101,18 @@ class TestKafkaBatchHistory: assert history.bootstrap_servers == "localhost:9092" mock_producer_class.assert_called_once() + @patch("generator._with_retry") @patch("generator._import_kafka") - def test_init_raises_on_connection_error(self, mock_import): - """Инициализация падает при ошибке подключения.""" - mock_producer_class = MagicMock(side_effect=Exception("Connection failed")) + def test_init_uses_retry(self, mock_import, mock_retry): + """Инициализация использует retry для подключения.""" + mock_producer_class = MagicMock() mock_import.return_value = (mock_producer_class, Exception) + mock_retry.side_effect = lambda f, **kwargs: f() # Выполняем функцию сразу - with pytest.raises(Exception, match="Connection failed"): - KafkaBatchHistory("localhost:9092") + KafkaBatchHistory("localhost:9092") + + # Проверяем что _with_retry был вызван + mock_retry.assert_called() @patch("generator._import_kafka") def test_add_sends_to_kafka(self, mock_import): @@ -141,17 +145,23 @@ class TestKafkaBatchHistory: assert "value" in call_args[1] @patch("generator._import_kafka") - def test_add_raises_on_send_error(self, mock_import): - """add пробрасывает ошибку отправки.""" + def test_add_retries_on_failure(self, mock_import): + """add делает retry при ошибке и пытается реконнект.""" mock_producer = MagicMock() - mock_producer.send.side_effect = Exception("Send failed") + # Первые 3 вызова падают, потом успех + mock_producer.send.side_effect = [ + Exception("Fail 1"), + Exception("Fail 2"), + Exception("Fail 3"), + MagicMock(), # Успех после реконнекта + ] mock_producer_class = MagicMock(return_value=mock_producer) mock_import.return_value = (mock_producer_class, Exception) history = KafkaBatchHistory("localhost:9092") now = datetime.now(timezone.utc) record = BatchRecord( - batch_id="fail789", + batch_id="retry456", started_at=now, finished_at=now, sent_total=50, @@ -163,8 +173,11 @@ class TestKafkaBatchHistory: error_message=None, ) - with pytest.raises(Exception, match="Send failed"): - history.add(record) + # Не должно упасть - должен быть retry + reconnect + history.add(record) + + # Проверяем что producer.send вызывался несколько раз (retry) + assert mock_producer.send.call_count >= 1 @patch("generator._import_kafka") def test_flush_calls_producer_flush(self, mock_import): @@ -189,3 +202,15 @@ class TestKafkaBatchHistory: history.close() mock_producer.close.assert_called_once() + + @patch("generator._import_kafka") + def test_close_ignores_errors(self, mock_import): + """close игнорирует ошибки при закрытии.""" + mock_producer = MagicMock() + mock_producer.close.side_effect = Exception("Close failed") + mock_producer_class = MagicMock(return_value=mock_producer) + mock_import.return_value = (mock_producer_class, Exception) + + history = KafkaBatchHistory("localhost:9092") + # Не должно упасть + history.close() diff --git a/generator/tests/test_state.py b/generator/tests/test_state.py index 1272154..da55499 100644 --- a/generator/tests/test_state.py +++ b/generator/tests/test_state.py @@ -2,6 +2,7 @@ Тесты сохранения и восстановления состояния генератора. """ import json +import random from datetime import datetime, timezone from unittest.mock import MagicMock, patch @@ -10,13 +11,19 @@ import pytest from generator import GeneratorState, KafkaStateManager +def _make_valid_rng_state(seed: int = 42): + """Создаёт валидный RNG state для тестов.""" + rng = random.Random(seed) + return rng.getstate() + + class TestGeneratorState: """Тесты структуры состояния генератора.""" def test_state_creation(self): """Создание состояния с всеми полями.""" now = datetime.now(timezone.utc) - rng_state = (3, (1, 2, 3), None) # Минимальный валидный state для random + rng_state = _make_valid_rng_state(42) state = GeneratorState( tick=42, @@ -35,7 +42,7 @@ class TestGeneratorState: def test_default_version(self): """Версия по умолчанию.""" now = datetime.now(timezone.utc) - rng_state = (3, (1, 2, 3), None) + rng_state = _make_valid_rng_state(42) state = GeneratorState( tick=1, @@ -49,7 +56,7 @@ class TestGeneratorState: def test_to_dict_serialization(self): """Сериализация в словарь (JSON-safe, без pickle).""" now = datetime.now(timezone.utc) - rng_state = (3, (1, 2, 3), None) + rng_state = _make_valid_rng_state(42) state = GeneratorState( tick=42, @@ -67,7 +74,8 @@ class TestGeneratorState: # Проверяем что rng_state сериализован как tuple (JSON-safe, без pickle) assert "rng_state" in data - assert data["rng_state"] == rng_state + # После to_dict rng_state должен быть tuple + assert isinstance(data["rng_state"], tuple) # Проверяем что можно сериализовать в JSON и восстановить json_str = json.dumps(data) restored_data = json.loads(json_str) @@ -77,7 +85,7 @@ class TestGeneratorState: def test_from_dict_deserialization(self): """Десериализация из словаря.""" now = datetime.now(timezone.utc) - rng_state = (3, (1, 2, 3), None) + rng_state = _make_valid_rng_state(42) # Создаём исходное состояние original = GeneratorState( @@ -99,8 +107,6 @@ class TestGeneratorState: def test_roundtrip_with_real_random(self): """Проверка что RNG state действительно восстанавливает последовательность.""" - import random - # Создаём генератор и делаем несколько вызовов rng = random.Random(12345) values_before = [rng.random() for _ in range(5)] @@ -132,6 +138,95 @@ class TestGeneratorState: assert next_values == values_after +class TestGeneratorStateValidation: + """Тесты валидации состояния и graceful degradation.""" + + def test_from_dict_missing_rng_state_raises(self): + """from_dict выбрасывает исключение при отсутствии rng_state.""" + data = { + "tick": 42, + "last_batch_id": "test", + "last_timestamp": "2024-01-01T00:00:00+00:00", + } + + with pytest.raises(ValueError, match="rng_state"): + GeneratorState.from_dict(data) + + def test_from_dict_invalid_rng_state_raises(self): + """from_dict выбрасывает исключение при невалидном rng_state.""" + data = { + "tick": 42, + "rng_state": "not_a_tuple", + "last_batch_id": "test", + "last_timestamp": "2024-01-01T00:00:00+00:00", + } + + with pytest.raises(ValueError): + GeneratorState.from_dict(data) + + def test_from_dict_insufficient_rng_state_raises(self): + """from_dict выбрасывает исключение при коротком rng_state.""" + data = { + "tick": 42, + "rng_state": [1], # Слишком короткий + "last_batch_id": "test", + "last_timestamp": "2024-01-01T00:00:00+00:00", + } + + with pytest.raises(ValueError): + GeneratorState.from_dict(data) + + def test_from_dict_invalid_setstate_raises(self): + """from_dict выбрасывает исключение если setstate падает.""" + data = { + "tick": 42, + "rng_state": [999, [1, 2, 3], None], # Невалидный state + "last_batch_id": "test", + "last_timestamp": "2024-01-01T00:00:00+00:00", + } + + with pytest.raises(ValueError): + GeneratorState.from_dict(data) + + def test_from_dict_safe_returns_none_on_invalid(self): + """from_dict_safe возвращает None при невалидных данных.""" + data = { + "tick": 42, + "rng_state": "invalid", + "last_batch_id": "test", + "last_timestamp": "2024-01-01T00:00:00+00:00", + } + + result = GeneratorState.from_dict_safe(data) + assert result is None + + def test_from_dict_safe_returns_state_on_valid(self): + """from_dict_safe возвращает state при валидных данных.""" + rng = random.Random(42) + data = { + "tick": 42, + "rng_state": list(rng.getstate()), # JSON сериализует tuple как list + "last_batch_id": "test", + "last_timestamp": "2024-01-01T00:00:00+00:00", + } + + result = GeneratorState.from_dict_safe(data) + assert result is not None + assert result.tick == 42 + + def test_from_dict_uses_defaults_for_missing_fields(self): + """from_dict использует defaults для отсутствующих полей.""" + rng = random.Random(42) + data = { + "rng_state": list(rng.getstate()), + } + + result = GeneratorState.from_dict(data) + assert result.tick == 0 + assert result.last_batch_id == "" + assert result.version == "1.0" + + class TestKafkaStateManager: """Тесты менеджера состояния.""" @@ -158,10 +253,9 @@ class TestKafkaStateManager: manager = KafkaStateManager("kafka:29092") now = datetime.now(timezone.utc) - rng_state = (3, (1, 2, 3), None) state = GeneratorState( tick=42, - rng_state=rng_state, + rng_state=_make_valid_rng_state(42), last_batch_id="abc123", last_timestamp=now, ) @@ -203,10 +297,9 @@ class TestKafkaStateManager: mock_import.return_value = (mock_producer_class, None) now = datetime.now(timezone.utc) - rng_state = (3, (1, 2, 3), None) state = GeneratorState( tick=100, - rng_state=rng_state, + rng_state=_make_valid_rng_state(100), last_batch_id="xyz789", last_timestamp=now, ) @@ -227,6 +320,29 @@ class TestKafkaStateManager: assert result.tick == 100 assert result.last_batch_id == "xyz789" + def test_load_invalid_state_returns_none(self): + """Загрузка невалидного state возвращает None (graceful degradation).""" + with patch("generator._import_kafka") as mock_import, \ + patch("kafka.KafkaConsumer") as mock_consumer_class: + + mock_producer_class = MagicMock() + mock_import.return_value = (mock_producer_class, None) + + # Мокаем consumer с невалидным сообщением + mock_message = MagicMock() + mock_message.key = b"default" + mock_message.value = {"tick": 42, "rng_state": "invalid"} + + mock_consumer = MagicMock() + mock_consumer.__iter__ = MagicMock(return_value=iter([mock_message])) + mock_consumer_class.return_value = mock_consumer + + manager = KafkaStateManager("kafka:29092") + result = manager.load() + + # Должно вернуть None из-за невалидного state + assert result is None + def test_load_ignores_wrong_key(self): """Загрузка игнорирует сообщения с другим ключом.""" with patch("generator._import_kafka") as mock_import, \ @@ -309,8 +425,6 @@ class TestJsonSafeState: def test_json_roundtrip_rng_state(self): """Полный цикл: rng.getstate() -> JSON -> from_dict -> setstate.""" - import random - rng = random.Random(42) # Делаем несколько вызовов values_before = [rng.random() for _ in range(10)]