feat(generator): добавлено восстановление state по модельному времени

- Зачем:
  - рестарт генератора должен продолжать поток от модельной точки без дублей и смешивания state разных настроек.
- Что:
  - state v2 хранит модельную и настенную метки, скорость, timezone, T0 и seed.
  - live-восстановление считает модельную точку по настенной дельте и проверяет совместимость config.
  - добавлен путь восстановления от T_end и тесты короткого и долгого простоя.
- Проверка:
  - make generator-test.
  - ClickHouse-сценарии короткого и долгого восстановления state.
  - reviewer gate issue 04 пройден после исправления совместимости state.
This commit is contained in:
2026-06-14 18:47:34 +03:00
parent 589e2321b3
commit fcb06da16e
7 changed files with 540 additions and 35 deletions
+18 -5
View File
@@ -175,6 +175,12 @@ class TickStreamGenerator:
rng_state: tuple,
last_batch_id: str,
last_timestamp: datetime,
model_timestamp: datetime | None = None,
wall_timestamp: datetime | None = None,
model_time_speed: float = 1.0,
model_timezone: str = "UTC",
model_t0: datetime | None = None,
gen_seed: int | None = None,
) -> GeneratorState:
"""Возвращает JSON-сериализуемый снимок тикового слоя."""
return GeneratorState(
@@ -182,6 +188,12 @@ class TickStreamGenerator:
rng_state=rng_state,
last_batch_id=last_batch_id,
last_timestamp=last_timestamp,
model_timestamp=model_timestamp,
wall_timestamp=wall_timestamp,
model_time_speed=model_time_speed,
model_timezone=model_timezone,
model_t0=model_t0,
gen_seed=gen_seed,
population=[
{
"user_domain_id": user.user_domain_id,
@@ -220,6 +232,7 @@ class TickStreamGenerator:
def restore_state(
self,
state: GeneratorState,
resume_model_at: datetime | None = None,
restarted_at: datetime | None = None,
) -> None:
"""Восстанавливает популяцию и активные визиты из state v2."""
@@ -231,17 +244,17 @@ class TickStreamGenerator:
self.population.users = users
self.active_visits = []
restarted_time = (
_normalize_tick_time(restarted_at)
if restarted_at is not None
model_resume_time = (
_normalize_tick_time(resume_model_at or restarted_at)
if resume_model_at is not None or restarted_at is not None
else None
)
for item in state.active_visits:
visit = self._visit_from_state(item, users_by_id)
if self._is_overdue_after_restart(visit, restarted_time):
if self._is_overdue_after_restart(visit, model_resume_time):
self.population.finish_visit(
visit.user,
self._last_released_at(visit, state.last_timestamp),
self._last_released_at(visit, state.model_timestamp),
)
continue
self.active_visits.append(visit)
+77 -7
View File
@@ -73,14 +73,13 @@ class GeneratorService:
restored_state = self.state_manager.load()
if restored_state:
try:
self.generator.rng.setstate(restored_state.rng_state)
self.stream.restore_state(
self._restore_live_state(
restored_state,
restarted_at=datetime.now(timezone.utc),
wall_now_utc=datetime.now(timezone.utc),
)
self._tick = restored_state.tick
logger.info(
f"Restored state: continuing from tick {self._tick}, "
f"model_time={self._model_time.isoformat()}, "
f"last_batch_id={restored_state.last_batch_id}"
)
except Exception as e:
@@ -115,21 +114,92 @@ class GeneratorService:
if self.state_manager:
self.state_manager.close()
def restore_from_startup_history(
self,
state,
model_t_end: datetime,
) -> None:
"""Восстанавливает слепок стартовой истории ровно от T_end."""
self._restore_state_snapshot(state, resume_model_at=model_t_end)
def _restore_live_state(self, state, wall_now_utc: datetime) -> None:
"""Восстанавливает live-state с учётом прошедшего настенного времени."""
self._validate_live_state_config(state)
resume_model_at = self._calculate_live_resume_model_at(
state,
wall_now_utc=wall_now_utc,
)
self._restore_state_snapshot(state, resume_model_at=resume_model_at)
def _restore_state_snapshot(self, state, resume_model_at: datetime) -> None:
"""Применяет state к генератору и тиковому слою."""
resume_model_at = self._as_aware_utc(resume_model_at)
self.generator.rng.setstate(state.rng_state)
self.stream.restore_state(state, resume_model_at=resume_model_at)
self._tick = state.tick
self._model_time = resume_model_at
def _calculate_live_resume_model_at(
self,
state,
wall_now_utc: datetime,
) -> datetime:
"""Считает модельную точку live-восстановления по state v2."""
wall_now_utc = self._as_aware_utc(wall_now_utc)
wall_saved_at = self._as_aware_utc(state.wall_timestamp)
idle_seconds = max(0.0, (wall_now_utc - wall_saved_at).total_seconds())
return self._as_aware_utc(state.model_timestamp) + timedelta(
seconds=idle_seconds * state.model_time_speed,
)
def _validate_live_state_config(self, state) -> None:
"""Проверяет, что state относится к текущей конфигурации live-запуска."""
mismatches = []
if state.gen_seed != self.config.seed:
mismatches.append("gen_seed")
if self._as_aware_utc(state.model_t0) != self.config.model_t0:
mismatches.append("model_t0")
if state.model_timezone != self.config.model_timezone:
mismatches.append("model_timezone")
if abs(state.model_time_speed - self.config.model_time_speed) > 1e-9:
mismatches.append("model_time_speed")
if mismatches:
raise ValueError(
"state config mismatch: " + ", ".join(mismatches)
)
@staticmethod
def _as_aware_utc(value: datetime) -> datetime:
if value.tzinfo is None:
return value.replace(tzinfo=timezone.utc)
return value.astimezone(timezone.utc)
def _save_state(self, batch_id: str) -> None:
"""Сохраняет текущее состояние генератора."""
if not self.state_manager or not self.config.state_enabled:
return
try:
wall_now = datetime.now(timezone.utc)
state = self.stream.to_state(
tick=self._tick,
rng_state=self.generator.rng.getstate(),
last_batch_id=batch_id,
last_timestamp=datetime.now(timezone.utc),
last_timestamp=self._model_time,
model_timestamp=self._model_time,
wall_timestamp=wall_now,
model_time_speed=self.config.model_time_speed,
model_timezone=self.config.model_timezone,
model_t0=self.config.model_t0,
gen_seed=self.config.seed,
)
self.state_manager.save(state)
self.state_manager.flush()
logger.debug(f"Saved state: tick={self._tick}, batch_id={batch_id}")
logger.debug(
f"Saved state: tick={self._tick}, "
f"model_time={self._model_time.isoformat()}, batch_id={batch_id}"
)
except Exception as e:
logger.warning(f"Failed to save state: {e}")
METRICS_ERRORS_TOTAL.labels(topic="state").inc()
@@ -179,8 +249,8 @@ class GeneratorService:
if status in ("success", "partial"):
METRICS_LAST_SUCCESS.set_to_current_time()
self._save_state(batch_id)
self._advance_model_time()
self._save_state(batch_id)
self.publisher.flush()
pub_duration = time.time() - pub_start
+90 -2
View File
@@ -3,7 +3,8 @@
import logging
import random
from dataclasses import dataclass, field
from datetime import datetime
from datetime import datetime, timezone
from zoneinfo import ZoneInfo, ZoneInfoNotFoundError
logger = logging.getLogger("generator")
@@ -25,7 +26,56 @@ def _require_keys(item: dict, keys: tuple[str, ...], label: str) -> None:
raise ValueError(f"{label} missing fields: {', '.join(missing)}")
def _parse_aware_utc(value: str, label: str) -> datetime:
if not isinstance(value, str) or not value:
raise ValueError(f"{label} must be a non-empty string")
timestamp = datetime.fromisoformat(value.replace("Z", "+00:00"))
if timestamp.tzinfo is None:
raise ValueError(f"{label} must include timezone")
return timestamp.astimezone(timezone.utc)
def _validate_resume_fields(data: dict) -> None:
_require_keys(
data,
(
"model_timestamp",
"wall_timestamp",
"model_time_speed",
"model_timezone",
"model_t0",
"gen_seed",
),
"state",
)
_parse_aware_utc(data["model_timestamp"], "model_timestamp")
_parse_aware_utc(data["wall_timestamp"], "wall_timestamp")
_parse_aware_utc(data["model_t0"], "model_t0")
model_time_speed = data["model_time_speed"]
if (
isinstance(model_time_speed, bool)
or not isinstance(model_time_speed, int | float)
or model_time_speed <= 0
):
raise ValueError("model_time_speed must be a positive number")
model_timezone = data["model_timezone"]
if not isinstance(model_timezone, str) or not model_timezone:
raise ValueError("model_timezone must be a string")
try:
ZoneInfo(model_timezone)
except ZoneInfoNotFoundError as e:
raise ValueError(f"unknown model_timezone: {model_timezone}") from e
gen_seed = data["gen_seed"]
if gen_seed is not None:
if isinstance(gen_seed, bool) or not isinstance(gen_seed, int):
raise ValueError("gen_seed must be an integer or null")
def _validate_v2_payload(data: dict) -> None:
_validate_resume_fields(data)
population = data.get("population")
active_visits = data.get("active_visits")
pending_visit_births = data.get("pending_visit_births", 0.0)
@@ -129,10 +179,30 @@ class GeneratorState:
last_batch_id: str
last_timestamp: datetime
version: str = STATE_VERSION
model_timestamp: datetime | None = None
wall_timestamp: datetime | None = None
model_time_speed: float = 1.0
model_timezone: str = "UTC"
model_t0: datetime | None = None
gen_seed: int | None = None
population: list[dict] = field(default_factory=list)
active_visits: list[dict] = field(default_factory=list)
pending_visit_births: float = 0.0
def __post_init__(self) -> None:
if self.model_timestamp is None:
self.model_timestamp = self._as_aware_utc(self.last_timestamp)
if self.wall_timestamp is None:
self.wall_timestamp = self._as_aware_utc(self.last_timestamp)
if self.model_t0 is None:
self.model_t0 = self.model_timestamp
@staticmethod
def _as_aware_utc(value: datetime) -> datetime:
if value.tzinfo is None:
return value.replace(tzinfo=timezone.utc)
return value.astimezone(timezone.utc)
def to_dict(self) -> dict:
"""Конвертирует в словарь для JSON-сериализации."""
return {
@@ -140,6 +210,12 @@ class GeneratorState:
"rng_state": self.rng_state,
"last_batch_id": self.last_batch_id,
"last_timestamp": self.last_timestamp.isoformat(),
"model_timestamp": self.model_timestamp.isoformat(),
"wall_timestamp": self.wall_timestamp.isoformat(),
"model_time_speed": self.model_time_speed,
"model_timezone": self.model_timezone,
"model_t0": self.model_t0.isoformat(),
"gen_seed": self.gen_seed,
"version": self.version,
"population": self.population,
"active_visits": self.active_visits,
@@ -160,6 +236,12 @@ class GeneratorState:
)
raise ValueError(f"unsupported state version: {version}")
_validate_v2_payload(data)
model_timestamp = _parse_aware_utc(
data["model_timestamp"],
"model_timestamp",
)
wall_timestamp = _parse_aware_utc(data["wall_timestamp"], "wall_timestamp")
model_t0 = _parse_aware_utc(data["model_t0"], "model_t0")
rng_state_raw = data.get("rng_state")
if not rng_state_raw:
@@ -184,9 +266,15 @@ class GeneratorState:
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")
data.get("last_timestamp", data["model_timestamp"])
),
version=version,
model_timestamp=model_timestamp,
wall_timestamp=wall_timestamp,
model_time_speed=float(data["model_time_speed"]),
model_timezone=data["model_timezone"],
model_t0=model_t0,
gen_seed=data["gen_seed"],
population=data.get("population", []),
active_visits=data.get("active_visits", []),
pending_visit_births=data.get("pending_visit_births", 0.0),