From be4063bd07c79e591dacc5e3e2f2bd9bd8d286bb Mon Sep 17 00:00:00 2001 From: MatveyVarfolomeev Date: Fri, 25 Sep 2026 22:55:18 +0300 Subject: [PATCH 1/2] feat: add assisted blend editor and harden release serving --- DESIGN.md | 6 +- README.md | 6 +- config/release_manifest.json | 10 +- docs/DEMO_SCRIPT.md | 12 +- global_tests/test_blend_assist.py | 115 ++++++ global_tests/test_stage2_data.py | 65 +++- global_tests/test_stage2_features.py | 60 +++ global_tests/test_stage6_acceptance.py | 10 + global_tests/test_stage7_history_serving.py | 108 +++++- global_tests/test_ui.py | 16 +- global_tests/test_ui_interactive.py | 40 +- scripts/build_release.ps1 | 8 +- source/acceptance.py | 6 + source/agents/blend_assist.py | 192 ++++++++++ source/agents/optimizer.py | 25 ++ source/agents/quality.py | 15 +- source/config.py | 20 +- source/contracts.py | 13 +- source/data/prepare.py | 146 +++++-- source/data/state.py | 35 +- source/desktop.py | 10 +- source/journal.py | 30 +- source/main.py | 40 +- source/ml/features.py | 182 ++++++++- source/orchestrator.py | 72 +++- source/ui.py | 184 ++++++--- source/ui_charts.py | 22 +- source/ui_data.py | 26 +- source/ui_what_if.py | 405 ++++++++++++++++++-- 29 files changed, 1694 insertions(+), 185 deletions(-) create mode 100644 global_tests/test_blend_assist.py create mode 100644 source/agents/blend_assist.py diff --git a/DESIGN.md b/DESIGN.md index 7cc234c..4a9b79a 100644 --- a/DESIGN.md +++ b/DESIGN.md @@ -99,7 +99,6 @@ flowchart TD ```text HackathonPetrolCode/ ├── README.md # запуск и ограничения готовой версии -├── IMPLEMENTATION_PLAN.md # сроки, этапы и критерии готовности ├── DESIGN.md # архитектура и контракты ├── pyproject.toml # Ruff и pytest; заменяет config.toml ├── requirements.txt # проверенные зависимости Python 3.11 @@ -568,12 +567,15 @@ point/upper и reason codes. В отдельных экранах доступн action-shadow и контроль ПАК--ЛИМС; ни один из этих результатов не разрешает изменение реальных уставок. -В v1.1.0 редактор `source/ui_what_if.py` создаёт проверенную копию +Редактор `source/ui_what_if.py` создаёт проверенную копию `ScenarioConfig` только для `model_demo` без реальных controls. Свойства A/B, запасы, доли, масса партии и синтетическая кривая присадки передаются через `run_model_demo(scenario_override=...)` существующему циклу решения. Preset на диске не изменяется; изменённые входы попадают в журнал и JSON-экспорт. Сброс восстанавливает поля preset и требует явного пересчёта. +С v1.2.0 он подбирает допустимую рецептуру двух компонентов при закреплённых +значениях долей или присадки. Предложения проходят тот же фильтр ограничений, +что и обычный расчёт; исходное и показанное значения записываются в журнал. `source/ui_charts.py` строит график исторического point/upper и лимита серы. Ось времени соответствует цели прогноза (as_of + 60 минут); пропуски сохраняются, diff --git a/README.md b/README.md index d41eb9b..8d9dd64 100644 --- a/README.md +++ b/README.md @@ -6,7 +6,7 @@ ## Запуск за три шага -В v1.1.0 добавлены **синтетический what-if** (свойства, запасы, присадка, пересчёт и сброс) и **график исторического прогноза**. Технические пути на экране истории скрыты в дополнительных настройках. Инструкция и материалы для жюри обновлены. +В v1.2.0 синтетический what-if подбирает допустимую смесь при изменении доли или присадки: можно закрепить введённое значение, сравнить варианты и отменить автоподбор. Исторический прогноз показывает причины отказа и сохраняет статус запуска в журнале. 1. Скачайте **Neftekod-Windows-x64.zip** по ссылке выше. 2. Распакуйте **весь архив** в доступную для записи папку, например `D:\Нефтекод`. @@ -22,7 +22,7 @@ Python, установка библиотек, обучение моделей | Раздел | Что доступно | | --- | --- | -| Обзор и Смесь | Шесть сценариев и редактор «Синтетический what-if»: свойства компонентов, запасы и присадка | +| Обзор и Смесь | Шесть сценариев и редактор «Синтетический what-if»: автоподбор рецептуры, свойства компонентов, запасы и присадка | | АВТ и Гидроочистка | Исторические сигналы, единицы, источники и свежесть | | История/ML | График прогноза и верхней оценки, выбор периода, прогресс, отмена и экспорт | | Hybrid и v2 | Условное смешение по сере и исследовательские предупреждения | @@ -55,6 +55,8 @@ Python, установка библиотек, обучение моделей ## Для разработчика +Редактор синтетической смеси предлагает допустимые доли после изменения поля, сохраняет закреплённые значения и показывает варианты «ближе», «дешевле» и «меньше присадки». Для проверки откройте «Обзор» → «Смесь» → «Синтетический what-if», измените долю компонента и нажмите Enter или «Подобрать смесь». «Отменить автоподбор» возвращает введённые значения; переключатель «Автоподбор» оставляет ручной режим. + В `main` находится сдаваемая версия без дневников разработки. Планы этапов, старые ревью и журнал проблем сохранены в [ветке dev](https://github.com/bug00n/HackathonPetrolCode/tree/dev); конкретная поставка закреплена тегом релиза. В Git хранятся исходники, конфигурации, тесты и технические материалы. EXE, презентация, Word-документ, подготовленные данные и модели поставляются через Releases. Уведомления сторонних компонентов сохранены в [third_party/licenses](third_party/licenses). [Запуск исходников, тесты и сборка](docs/DEVELOPMENT.md) · [Архитектура и контракты](DESIGN.md) · [Уточнения организаторов](docs/technical/QA_CLARIFICATIONS.md) diff --git a/config/release_manifest.json b/config/release_manifest.json index 97a23e8..f90024a 100644 --- a/config/release_manifest.json +++ b/config/release_manifest.json @@ -1,7 +1,15 @@ { "schema_version": "1.0", - "release_id": "neftekod-v1.1.0", + "release_id": "neftekod-v1.2.0", "prepared_dataset": "data/processed/66bdfcbb23b4", + "legacy_config_sha256": "00a45650caa0f415e78037d56bc49de7161f55f9eecd24675e8874d64efb63e2", + "preparation_config_sha256": "7f8dab1f4befa7ed4d14108146d344eeeb4341711880144ad0195d163a55c4dd", + "prepared_sha256": { + "telemetry.csv.gz": "20b434e7fd01e49ad45ec82635ea14940677267f58d0d3f70f63812d8c17c84d", + "quality.csv.gz": "fbef330af12854052c7622f10c561460af3d1a708d59b87892ad49d82f0be9e4", + "issues.csv.gz": "bffe17ad1fcfdfaa4496405ce1b7eb2d0f318c67103bfbaf63a0e45d17ee9fb9", + "feature_order.json": "fb7b26c4991e029f246baa8efc99f971f07fc765bbb5e59d521794a3ed33932e" + }, "forecast_artifact": "artifacts/models/sulfur-upper-33efba3c3141", "v2_artifact": "artifacts/models/sulfur-v2-shadow-61f972d07181", "action_artifact": null, diff --git a/docs/DEMO_SCRIPT.md b/docs/DEMO_SCRIPT.md index 83eee67..ade2db0 100644 --- a/docs/DEMO_SCRIPT.md +++ b/docs/DEMO_SCRIPT.md @@ -97,19 +97,23 @@ Hybrid проверяет только серу. Он не подтвержда Нажмите «Синтетический what-if» на «Обзоре» или «Смеси». За основу берётся выбранный готовый сценарий. При повторном открытии после успешного пересчёта редактор показывает последние изменённые входы этого сценария. Это эксперимент с модельными допущениями, а не правка истории предприятия. -На вкладке «Компоненты и текущая рецептура» для A и B задаются сера и её верхняя оценка (мг/кг), T95 и его верхняя оценка (°C), цетановое число и его нижняя оценка, запас (т), стоимость (условный индекс/т), риск от 0 до 1 и текущая массовая доля (%). Пустое поле качества означает неизвестное значение и может привести к отказу; запас и другие обязательные числа оставлять пустыми нельзя. +На первой вкладке «Рецептура» задайте доли A, B, присадки и массу партии. После Enter или выхода из поля редактор подбирает допустимые значения остальных долей, оставляя последнее изменённое значение неизменным. Галочка «Зафиксировать» удерживает долю и при следующих правках. В списке вариантов можно выбрать смесь ближе к введённой, с меньшей стоимостью или меньшим количеством присадки. Под списком показаны изменения и проверки серы, T95 и цетанового числа. Если смесь невозможна при закреплённых полях или запасах, редактор объясняет причину. -На вкладке «Партия и присадка» задаются положительная масса партии (т), текущая доля присадки (%), максимум присадки (1, 2 или 3%), запас (т), стоимость (условный индекс/т) и прирост цетанового числа при 1, 2 и 3%. Кривая эффекта остаётся синтетической; прирост не должен уменьшаться с дозой. +На вкладке «Свойства и запасы» для A и B задаются сера и её верхняя оценка (мг/кг), T95 и его верхняя оценка (°C), цетановое число и его нижняя оценка, запас (т), стоимость (условный индекс/т) и риск от 0 до 1. Пустое поле качества означает неизвестное значение и может привести к отказу; запас и другие обязательные числа оставлять пустыми нельзя. + +На вкладке «Параметры присадки» задаются максимум присадки (1, 2 или 3%), запас (т), стоимость (условный индекс/т) и прирост цетанового числа при 1, 2 и 3%. Кривая эффекта остаётся синтетической; прирост не должен уменьшаться с дозой. | Кнопка | Действие | | --- | --- | -| Пересчитать what-if | Проверяет поля, закрывает редактор и запускает существующий расчёт; результат появляется на «Смеси» с отметкой what-if | +| Подобрать смесь | Пересчитывает допустимые варианты с учётом закреплённых значений | +| Отменить автоподбор | Возвращает введённые перед подбором значения; ручной режим также доступен через переключатель «Автоподбор» | +| Пересчитать what-if | Проверяет поля, закрывает редактор и оценивает именно показанную смесь; результат появляется на «Смеси» с отметкой what-if | | Сбросить к сценарию | Восстанавливает исходные поля; затем нажмите «Пересчитать what-if», чтобы обновить результат | | Закрыть без изменений | Закрывает редактор без применения введённых значений; крестик действует так же | Сумма текущих долей A, B и присадки должна равняться 100%. Верхняя оценка не может быть ниже точки, нижняя — выше. Числа должны быть конечными и неотрицательными; допустимы точка и запятая. Ошибка ввода остаётся в окне. Нехватка запасов, пропущенное качество или отсутствие допустимой смеси дают объяснённый отказ, а не вымышленную рекомендацию. -Оптимизатор сохраняет прежний дискретный перебор: соотношение A/B с шагом 5% без присадки, присадка с шагом 1%. Редактор не обещает непрерывный математический оптимум. Изменённый сценарий сохраняется в журнале `input.json` и в экспорте; исходный preset на диске остаётся прежним. +Обычный расчёт сохраняет прежний дискретный перебор и добавляет допустимые граничные смеси, если сетка пропускает узкий допустимый диапазон. Редактор ищет смесь в пределах заданных свойств, запасов и закреплённых полей; он не обещает непрерывный математический оптимум. Изменённый сценарий и введённые значения сохраняются в журнале `input.json` и в экспорте; исходный preset на диске остаётся прежним. Для дополнительной минуты показа выберите «Повышенная сера», откройте редактор и измените серу A с 6 до 1, верхнюю оценку с 7 до 2 мг/кг. Нажмите пересчёт и сравните качество исходной смеси. Затем снова откройте редактор, сбросьте поля и пересчитайте: исходный пример вернётся. diff --git a/global_tests/test_blend_assist.py b/global_tests/test_blend_assist.py new file mode 100644 index 0000000..a1f4a7c --- /dev/null +++ b/global_tests/test_blend_assist.py @@ -0,0 +1,115 @@ +"""The synthetic blend helper must agree with the decision engine.""" + +import json +from pathlib import Path + +import pytest + +from source.agents.blend_assist import feasible_blends +from source.agents.optimizer import evaluate_candidates, generate_candidates +from source.config import load_runtime_config, load_scenario +from source.contracts import ScenarioConfig +from source.main import run_model_demo +from source.orchestrator import _model_demo_state +from source.ui_what_if import assist_editor_values, editor_values, scenario_from_editor + + +def test_editor_keeps_changed_share_and_only_repairs_recipe() -> None: + scenario = load_scenario("config/scenarios/blend_risk.json") + values = editor_values(scenario) + values["B.fraction"] = "10" + options = assist_editor_values(scenario, values, frozenset({"B.fraction"})) + assert options + for option, updated in options: + assert updated["B.fraction"] == "10" + assert updated["A.sulfur.value"] == values["A.sulfur.value"] + assert updated["total_mass_t"] == values["total_mass_t"] + assert scenario_from_editor(scenario, updated) + assert option.evaluation.feasible + + +def test_editor_explains_impossible_fixed_stock() -> None: + scenario = load_scenario("config/scenarios/blend_risk.json") + values = editor_values(scenario) + values["B.fraction"] = "40" + with pytest.raises(ValueError, match="доступно 30"): + assist_editor_values(scenario, values, frozenset({"B.fraction"})) + + +def test_lowering_additive_cap_repairs_the_current_dose() -> None: + scenario = load_scenario("config/scenarios/blend_risk.json") + values = editor_values(scenario) | {"dose": "3", "max_dose": "1"} + options = assist_editor_values(scenario, values, frozenset()) + assert options + assert all(float(updated["dose"]) <= 1 for _, updated in options) + + +def test_continuous_search_finds_narrow_feasible_window() -> None: + base = load_scenario("config/scenarios/blend_risk.json") + payload = base.model_dump(mode="json") + payload["blend_components"][0]["available_mass_t"] = 100 + payload["blend_components"][1]["available_mass_t"] = 100 + payload["blend_components"][0]["sulfur"].update(value=0, upper=0) + payload["blend_components"][1]["sulfur"].update(value=100, upper=100) + payload["blend_components"][0]["cetane_number"].update(value=0, lower=0) + payload["blend_components"][1]["cetane_number"].update(value=100, lower=100) + payload["constraints"] = [ + payload["constraints"][0] | {"upper": 48.9}, + payload["constraints"][1] | {"upper": 400}, + payload["constraints"][2] | {"lower": 48.1}, + ] + payload["cetane_additive"]["available_mass_t"] = 0 + scenario = ScenarioConfig.model_validate(payload) + from datetime import UTC, datetime + + state = _model_demo_state(datetime.now(UTC), scenario) + found = feasible_blends(scenario, dict(scenario.current_blend_mass_fractions), 0) + assert found + assert all(0.481 <= item.fractions["B"] <= 0.489 for item in found) + assert any( + item.feasible + for item in evaluate_candidates(state, generate_candidates(state, scenario), scenario) + ) + bounded = load_runtime_config(Path("config/runtime.toml")).model_copy( + update={"max_candidates": 85} + ) + candidates = generate_candidates(state, scenario, bounded) + assert len(candidates) == 85 + assert any(item.feasible for item in evaluate_candidates(state, candidates, scenario)) + + +def test_fixed_recipe_and_editor_context_are_journaled(tmp_path: Path) -> None: + preset = load_scenario("config/scenarios/blend_risk.json") + values = editor_values(preset) | {"B.fraction": "10"} + _, updated = assist_editor_values(preset, values, frozenset({"B.fraction"}))[0] + scenario = scenario_from_editor(preset, updated) + context = {"original_values": values, "displayed_values": updated, "assisted": True} + result = run_model_demo( + preset.id, + run_dir=tmp_path, + scenario_override=scenario, + recipe_mode="evaluate", + editor_context=context, + ) + assert result.selected is not None + assert result.selected.candidate.kind.value == "hold" + recorded = json.loads((tmp_path / result.run_id / "input.json").read_text()) + assert recorded["editor_context"] == context + assert ( + recorded["scenario"]["current_blend_mass_fractions"] + == scenario.current_blend_mass_fractions + ) + status = json.loads((tmp_path / result.run_id / "status.json").read_text()) + assert status["status"] == "completed" + + +def test_failed_run_has_durable_status(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> None: + def fail(*_args, **_kwargs): + raise RuntimeError("deliberate failure") + + monkeypatch.setattr("source.orchestrator._run_cycle", fail) + with pytest.raises(RuntimeError, match="deliberate failure"): + run_model_demo("blend_risk", run_dir=tmp_path) + statuses = [json.loads(path.read_text()) for path in tmp_path.glob("*/status.json")] + assert len(statuses) == 1 + assert statuses[0]["status"] == "failed" diff --git a/global_tests/test_stage2_data.py b/global_tests/test_stage2_data.py index 8463ca7..8452b6f 100644 --- a/global_tests/test_stage2_data.py +++ b/global_tests/test_stage2_data.py @@ -2,6 +2,7 @@ from __future__ import annotations +import io import json import tarfile from datetime import datetime @@ -69,6 +70,39 @@ def _write_minimal_materials(root: Path) -> None: ).to_excel(root / "ЛИМСы 01.01.2023 - н.в_ (2).xlsx", header=False, index=False) +@pytest.mark.parametrize("external_tar", [False, True]) +def test_archive_extraction_rejects_symlinked_data_directory( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch, external_tar: bool +) -> None: + """Archive metadata must not redirect CSV extraction outside the temporary directory.""" + archive_path = tmp_path / "malicious.tar" + with tarfile.open(archive_path, "w") as archive: + link = tarfile.TarInfo("data") + link.type = tarfile.SYMTYPE + link.linkname = "../outside" + archive.addfile(link) + for name in ("avt_tags.csv", "242000_tags.csv"): + payload = b"date,value\n2026-01-15,1\n" + member = tarfile.TarInfo(f"data/{name}") + member.size = len(payload) + archive.addfile(member, io.BytesIO(payload)) + + outside = tmp_path / "outside" + outside.mkdir() + destination = tmp_path / "extracted" + destination.mkdir() + if external_tar: + + def reject_tar_archive(path: Path) -> None: + raise tarfile.ReadError(path) + + monkeypatch.setattr(prepare_module.tarfile, "open", reject_tar_archive) + + with pytest.raises(ValueError, match="unsafe archive member"): + prepare_module._extract_telemetry(archive_path, destination) + assert not list(outside.iterdir()) + + def test_prepared_dataset_roundtrip(tmp_path: Path) -> None: """Verify a written prepared dataset can be loaded back without contract drift.""" data = _fixture_prepared_data() @@ -76,13 +110,42 @@ def test_prepared_dataset_roundtrip(tmp_path: Path) -> None: dataset_path = write_prepared_dataset(data, tmp_path) loaded = load_prepared_dataset(dataset_path) - assert loaded.manifest == data.manifest + assert loaded.manifest.model_dump(exclude={"prepared_sha256"}) == data.manifest.model_dump( + exclude={"prepared_sha256"} + ) + assert loaded.manifest.prepared_sha256 is not None assert loaded.feature_order == data.feature_order pd.testing.assert_frame_equal(loaded.telemetry, data.telemetry, check_dtype=False) pd.testing.assert_frame_equal(loaded.quality, data.quality, check_dtype=False) pd.testing.assert_frame_equal(loaded.issues, data.issues, check_dtype=False) +def test_prepared_dataset_memory_cache_isolated_and_invalidated( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + dataset_path = write_prepared_dataset(_fixture_prepared_data(), tmp_path) + read_csv = pd.read_csv + reads: list[str] = [] + + def counted_read_csv(path: Path, *args: object, **kwargs: object) -> pd.DataFrame: + reads.append(Path(path).name) + return read_csv(path, *args, **kwargs) + + monkeypatch.setattr(prepare_module.pd, "read_csv", counted_read_csv) + first = load_prepared_dataset(dataset_path) + original = first.quality.loc[0, "value"] + first.quality.loc[0, "value"] = 999 + second = load_prepared_dataset(dataset_path) + assert second.quality.loc[0, "value"] == original + assert len(reads) == 3 + + second.quality.loc[0, "value"] = 998 + second.quality.to_csv(dataset_path / "quality.csv.gz", index=False, compression="gzip") + with pytest.raises(ValueError, match="checksum mismatch"): + load_prepared_dataset(dataset_path) + assert len(reads) == 6 + + def test_prepared_dataset_publish_retries_transient_permission_error( tmp_path: Path, monkeypatch: pytest.MonkeyPatch ) -> None: diff --git a/global_tests/test_stage2_features.py b/global_tests/test_stage2_features.py index 23033e2..cceb381 100644 --- a/global_tests/test_stage2_features.py +++ b/global_tests/test_stage2_features.py @@ -16,6 +16,7 @@ baseline_feature_name, build_features, build_supervised_dataset, + prepare_feature_batch, ) FIXTURES = Path(__file__).parent / "fixtures" @@ -249,6 +250,65 @@ def test_training_and_serving_use_identical_features_and_model_order() -> None: pd.testing.assert_frame_equal(served, expected) +def test_batched_history_features_match_single_points_with_delayed_lims() -> None: + signal = "ht:2:Mg.Sulfur" + data = _prepared( + pd.DataFrame( + { + "timestamp": pd.date_range("2026-01-01", periods=8, freq="1h", tz="UTC"), + "ht:F1": np.arange(8, dtype=float), + } + ), + [ + _observation( + "old", + "2026-01-01T00:00:00Z", + 5.0, + source="lims", + signal_id=signal, + available_at="2026-01-01T02:00:00Z", + ), + _observation( + "new", + "2026-01-01T04:00:00Z", + 7.0, + source="lims", + signal_id=signal, + available_at="2026-01-01T06:00:00Z", + ), + ], + ("ht:F1",), + ) + model = SimpleNamespace( + metadata=SimpleNamespace( + feature_names=( + "ht:F1__lag_0m", + baseline_feature_name(signal, "lims"), + f"{signal}__lims_mean_60m", + ), + target_signal=signal, + target_source="lims", + target_unit="mg/kg", + ) + ) + queries = pd.date_range("2026-01-01T01:00:00Z", periods=4, freq="2h") + batch = prepare_feature_batch(data, queries, model) + for query in queries: + state = ProcessState( + state_id="fixture", + as_of=query.to_pydatetime(), + dataset_id=data.manifest.dataset_id, + mode=OperationMode.HISTORY, + signals={}, + ) + single = build_features(data, query.to_pydatetime(), state, model) + batched = build_features(data, query.to_pydatetime(), state, model, batch=batch) + pd.testing.assert_frame_equal(single, batched) + assert np.isnan(batch.frame.iloc[0][baseline_feature_name(signal, "lims")]) + assert batch.frame.iloc[1][baseline_feature_name(signal, "lims")] == 5.0 + assert batch.frame.iloc[-1][baseline_feature_name(signal, "lims")] == 7.0 + + def test_build_features_rejects_naive_time_and_state_mismatch() -> None: data = _prepared( pd.DataFrame({"timestamp": [pd.Timestamp("2026-01-01T00:00:00Z")]}), diff --git a/global_tests/test_stage6_acceptance.py b/global_tests/test_stage6_acceptance.py index 05320af..425d0a4 100644 --- a/global_tests/test_stage6_acceptance.py +++ b/global_tests/test_stage6_acceptance.py @@ -13,6 +13,7 @@ from source.acceptance import ( JOURNAL_FILES, + export_journals, load_episode_specs, run_acceptance_suite, sha256_file, @@ -83,6 +84,15 @@ def test_acceptance_suite_reproduces_decisions_and_exports_full_journals( assert f"sulfur_risk/{filename}" in names +@pytest.mark.parametrize("name", ("../outside", r"..\outside")) +def test_journal_export_rejects_path_components(name: str, tmp_path: Path) -> None: + """A catalog id must not create a ZIP entry outside its own journal directory.""" + destination = tmp_path / "journals.zip" + with pytest.raises(ValueError, match="single path components"): + export_journals([(name, tmp_path / "source")], destination) + assert not destination.exists() + + def test_acceptance_decision_fingerprints_repeat(tmp_path: Path) -> None: reports = [ run_acceptance_suite( diff --git a/global_tests/test_stage7_history_serving.py b/global_tests/test_stage7_history_serving.py index 8076617..3c424e7 100644 --- a/global_tests/test_stage7_history_serving.py +++ b/global_tests/test_stage7_history_serving.py @@ -13,7 +13,7 @@ from source.config import config_fingerprint, load_runtime_config from source.contracts import DatasetManifest, Recommendation -from source.data.prepare import PreparedData, write_prepared_dataset +from source.data.prepare import PreparedData, load_prepared_dataset, write_prepared_dataset from source.main import main, run_history_command from source.ml.artifacts import LastValueRegressor, feature_schema_hash, save_model @@ -114,6 +114,112 @@ def _write_point_model(directory: Path, tag_dictionary_sha256: str) -> Path: return artifact_dir +def test_reused_feature_sources_match_uncached_batches(tmp_path: Path) -> None: + from source.ml.artifacts import load_model + from source.ml.features import prepare_feature_batch, prepare_feature_sources + + data = _prepared_data() + model = load_model( + _write_point_model(tmp_path / "models", data.manifest.tag_dictionary_sha256), + trusted=True, + ) + queries = pd.date_range("2026-01-15T08:40:00Z", periods=4, freq="10min") + sources = prepare_feature_sources(data, model) + for chunk in (queries[:2], queries[2:]): + expected = prepare_feature_batch(data, chunk, model).frame + actual = prepare_feature_batch(data, chunk, model, sources=sources).frame + pd.testing.assert_frame_equal(expected, actual) + + +def test_feature_recipe_mismatch_fails_before_forecast(tmp_path: Path) -> None: + from source.ml.artifacts import load_model + from source.ml.features import prepare_feature_sources + + data = _prepared_data() + model = load_model( + _write_point_model(tmp_path / "models", data.manifest.tag_dictionary_sha256), + trusted=True, + ) + metadata = model.metadata.model_dump(mode="python") + metadata["processing"]["feature_definition"]["lags_minutes"] = [0, 15] + from types import SimpleNamespace + + with pytest.raises(ValueError, match="feature recipe"): + prepare_feature_sources(data, SimpleNamespace(metadata=metadata)) + + +def test_feature_recipe_accepts_telemetry_order_from_model() -> None: + from dataclasses import replace + from types import SimpleNamespace + + from source.ml.features import prepare_feature_sources + + data = _prepared_data() + telemetry = data.telemetry.copy() + telemetry["ht:F2"] = telemetry["ht:F1"] + 1 + data = replace(data, telemetry=telemetry, feature_order=("ht:F2", "ht:F1", TARGET_SIGNAL)) + metadata = { + "feature_names": ("ht:F1__lag_0m", "ht:F2__lag_0m", FEATURE_NAME), + "target_signal": TARGET_SIGNAL, + "target_source": "pak", + "target_unit": TARGET_UNIT, + "horizon_minutes": 60, + "processing": {"feature_definition": {"telemetry_signals": ["ht:F1", "ht:F2"]}}, + } + sources = prepare_feature_sources(data, SimpleNamespace(metadata=metadata)) + assert sources.feature_spec[1] == ("ht:F2", "ht:F1") + + +def test_runtime_output_path_does_not_invalidate_new_prepared_data(tmp_path: Path) -> None: + from dataclasses import replace + + from source.main import PROJECT_ROOT, _validate_current_prepared_dataset + + config = load_runtime_config(Path("config/runtime.toml")) + data = _prepared_data() + manifest = data.manifest.model_copy( + update={ + "preparation_version": "2", + "config_sha256": hashlib.sha256( + config_fingerprint(config, preparation_only=True).encode() + ).hexdigest(), + } + ) + data = replace(data, manifest=manifest) + path = write_prepared_dataset(data, tmp_path / "processed") + data = load_prepared_dataset(path) + changed_output = config.model_copy(update={"runs_dir": Path("other-runs")}) + _validate_current_prepared_dataset(data, changed_output, PROJECT_ROOT, path) + changed_preparation = config.model_copy(update={"source_timezone": "UTC"}) + with pytest.raises(ValueError, match="config_sha256"): + _validate_current_prepared_dataset(data, changed_preparation, PROJECT_ROOT, path) + + +def test_pinned_legacy_dataset_tolerates_output_path_change(tmp_path: Path) -> None: + from source.main import _validate_current_prepared_dataset + + config = load_runtime_config(Path("config/runtime.toml")) + for relative in ("config/tags.csv", "config/telemetry_rules.json"): + destination = tmp_path / relative + destination.parent.mkdir(parents=True, exist_ok=True) + destination.write_bytes(Path(relative).read_bytes()) + path = write_prepared_dataset(_prepared_data(), tmp_path / "data/processed") + data = load_prepared_dataset(path) + release = { + "prepared_dataset": f"data/processed/{path.name}", + "legacy_config_sha256": data.manifest.config_sha256, + "preparation_config_sha256": hashlib.sha256( + config_fingerprint(config, preparation_only=True).encode() + ).hexdigest(), + } + (tmp_path / "config/release_manifest.json").write_text(json.dumps(release), encoding="utf-8") + changed_output = config.model_copy(update={"runs_dir": Path("other-runs")}) + _validate_current_prepared_dataset(data, changed_output, tmp_path, path) + changed_preparation = config.model_copy(update={"source_timezone": "UTC"}) + with pytest.raises(ValueError, match="config_sha256"): + _validate_current_prepared_dataset(data, changed_preparation, tmp_path, path) + + def test_run_history_cli_uses_trusted_artifact_and_writes_journal( tmp_path: Path, capsys: pytest.CaptureFixture[str], diff --git a/global_tests/test_ui.py b/global_tests/test_ui.py index 25d59d7..835ee99 100644 --- a/global_tests/test_ui.py +++ b/global_tests/test_ui.py @@ -294,7 +294,7 @@ def release_context_root(tmp_path: Path) -> Path: def test_release_manifest_pins_one_live_context(release_context_root: Path) -> None: context = ui_data_module.discover_ui_context(release_context_root) - assert context.release_id == "neftekod-v1.1.0" + assert context.release_id == "neftekod-v1.2.0" assert context.latest_dataset is not None assert context.latest_dataset.endswith("66bdfcbb23b4") assert tuple(item.model_id for item in context.forecast_artifacts) == ( @@ -306,6 +306,13 @@ def test_release_manifest_pins_one_live_context(release_context_root: Path) -> N assert context.action_artifacts == () +def test_corrupt_release_manifest_cannot_fall_back_to_latest(release_context_root: Path) -> None: + release_file = release_context_root / "config/release_manifest.json" + release_file.write_text("{broken", encoding="utf-8") + with pytest.raises(ValueError, match="invalid release manifest"): + ui_data_module.discover_ui_context(release_context_root) + + def test_missing_release_pins_do_not_select_other_available_inputs( tmp_path: Path, monkeypatch: pytest.MonkeyPatch, @@ -376,7 +383,12 @@ def test_history_interval_ui_retains_and_exports_completed_result( monkeypatch.setattr(ui_module.PetrolCodeApp, "calculate", lambda self: None) monkeypatch.setattr(ui_module, "history_interval_command", lambda *args, **kwargs: payload) monkeypatch.setattr( - ui_module.threading, "Thread", lambda target, **kwargs: SimpleNamespace(start=target) + ui_module, + "ThreadPoolExecutor", + lambda **kwargs: SimpleNamespace( + submit=lambda target, *args: target(*args), + shutdown=lambda **kwargs: None, + ), ) monkeypatch.setattr( ui_module.filedialog, "asksaveasfilename", lambda **kwargs: str(destination) diff --git a/global_tests/test_ui_interactive.py b/global_tests/test_ui_interactive.py index 94ce74b..4ef29f8 100644 --- a/global_tests/test_ui_interactive.py +++ b/global_tests/test_ui_interactive.py @@ -146,7 +146,7 @@ def test_editor_reset_validation_and_chart_on_real_tk() -> None: buttons["Пересчитать what-if"].invoke() assert not received and dialog.winfo_exists() buttons["Сбросить к сценарию"].invoke() - assert first.get() == "6" + assert first.get() == editor_values(preset)["A.fraction"] buttons["Пересчитать what-if"].invoke() assert len(received) == 1 and received[0].id.endswith("_what_if") chart = HistoryChart(root) @@ -159,3 +159,41 @@ def test_editor_reset_validation_and_chart_on_real_tk() -> None: assert not chart.find_withtag("point") finally: root.destroy() + + +def test_editor_assist_and_undo_on_real_tk() -> None: + try: + root = tk.Tk() + except tk.TclError as exc: + pytest.skip(f"Tk display unavailable: {exc}") + root.withdraw() + try: + preset = load_scenario("config/scenarios/blend_risk.json") + dialog = open_editor(root, preset, preset, lambda _scenario: None) + dialog.withdraw() + entries = [widget for widget in descendants(dialog) if isinstance(widget, ttk.Entry)] + b_share = entries[1] + b_share.delete(0, "end") + b_share.insert(0, "10") + locks = [ + widget + for widget in descendants(dialog) + if isinstance(widget, ttk.Checkbutton) and widget.cget("text") == "Зафиксировать" + ] + locks[1].invoke() + buttons = { + widget.cget("text"): widget + for widget in descendants(dialog) + if isinstance(widget, ttk.Button) + } + buttons["Подобрать смесь"].invoke() + assert b_share.get() == "10" + assert float(entries[0].get()) + float(b_share.get()) + float( + entries[2].get() + ) == pytest.approx(100) + buttons["Отменить автоподбор"].invoke() + assert b_share.get() == "10" + assert entries[0].get() == editor_values(preset)["A.fraction"] + dialog.destroy() + finally: + root.destroy() diff --git a/scripts/build_release.ps1 b/scripts/build_release.ps1 index a971c99..5bb64a6 100644 --- a/scripts/build_release.ps1 +++ b/scripts/build_release.ps1 @@ -15,13 +15,10 @@ if (Test-Path -LiteralPath $destination) { $manifestPath = Join-Path $root "config\release_manifest.json" $manifest = Get-Content -Raw -LiteralPath $manifestPath | ConvertFrom-Json $required = @( - "README.md", "requirements.txt", "requirements.lock.txt", + "README.md", "DESIGN.md", "requirements.txt", "requirements.lock.txt", "pyproject.toml", "source", "config", "global_tests", "materials", "scripts", "docs", "third_party" ) -$required += @(Get-ChildItem -LiteralPath $root -Filter "*.md" -File | - Where-Object { $_.Name -notin @("README.md", "AGENTS.md") } | - ForEach-Object { $_.Name }) $selected = @( [string]$manifest.prepared_dataset, [string]$manifest.forecast_artifact, @@ -112,7 +109,8 @@ if ($BinaryDirectory) { if ($process.ExitCode -ne 0 -or -not (Test-Path -LiteralPath $smokeReport)) { throw "Portable EXE verification failed" } $smoke = Get-Content -Raw -LiteralPath $smokeReport | ConvertFrom-Json if ($smoke.error -or $smoke.tk -ne "ok" -or $smoke.history.points -ne 2 -or - -not $smoke.history_chart.rendered -or $smoke.what_if.editor -ne "ok") { + -not $smoke.history_chart.rendered -or $smoke.what_if.editor -ne "ok" -or + $smoke.what_if.assist -ne "ok") { throw "Portable smoke report is incomplete" } } finally { diff --git a/source/acceptance.py b/source/acceptance.py index 6a34893..dcd1893 100644 --- a/source/acceptance.py +++ b/source/acceptance.py @@ -212,6 +212,12 @@ def export_journals( raise ValueError("at least one journal is required") if len({name for name, _ in sources}) != len(sources): raise ValueError("journal export names must be unique") + # ZIP member names must stay inside one episode directory on extraction. + if any( + not name or name in {".", ".."} or "/" in name or "\\" in name or "\0" in name + for name, _ in sources + ): + raise ValueError("journal export names must be single path components") destination.parent.mkdir(parents=True, exist_ok=True) manifest: dict[str, Any] = {"schema_version": "1.0", "journals": []} file_descriptor, temporary_name = tempfile.mkstemp( diff --git a/source/agents/blend_assist.py b/source/agents/blend_assist.py new file mode 100644 index 0000000..682ba06 --- /dev/null +++ b/source/agents/blend_assist.py @@ -0,0 +1,192 @@ +"""Find feasible synthetic two-component recipes without a coarse share grid.""" + +from __future__ import annotations + +from dataclasses import dataclass +from datetime import UTC, datetime + +from source.agents.optimizer import evaluate_candidates +from source.contracts import CandidateAction, CandidateEvaluation, CandidateKind, ScenarioConfig +from source.orchestrator import _model_demo_state + + +@dataclass(frozen=True) +class BlendOption: + label: str + fractions: dict[str, float] + dose: float + evaluation: CandidateEvaluation + + +def _bound( + interval: tuple[float, float], slope: float, intercept: float, limit: float, *, upper: bool +) -> tuple[float, float]: + """Intersect lo <= share_B <= hi with one linear quality inequality.""" + lo, hi = interval + if abs(slope) < 1e-12: + valid = intercept <= limit + 1e-10 if upper else intercept >= limit - 1e-10 + return (lo, hi) if valid else (1.0, 0.0) + edge = (limit - intercept) / slope + if (slope > 0) == upper: + return lo, min(hi, edge) + return max(lo, edge), hi + + +def _gain(scenario: ScenarioConfig, dose: float) -> float: + additive = scenario.cetane_additive + if additive is None or dose == 0: + return 0.0 + points = additive.response_curve + for left, right in zip(points, points[1:], strict=False): + if left.mass_fraction <= dose <= right.mass_fraction: + share = (dose - left.mass_fraction) / (right.mass_fraction - left.mass_fraction) + return left.cetane_gain + share * (right.cetane_gain - left.cetane_gain) + raise ValueError("Доза вне заданной кривой присадки") + + +def _shares(scenario: ScenarioConfig, dose: float, fixed: dict[str, float]) -> tuple[float, float]: + a, b = scenario.blend_components + mass = scenario.total_mass_t + assert mass is not None and mass > 0 + diesel = 1.0 - dose + lo = max(0.0, diesel - a.available_mass_t / mass) + hi = min(diesel, b.available_mass_t / mass) + if a.id in fixed: + lo = max(lo, diesel - fixed[a.id]) + hi = min(hi, diesel - fixed[a.id]) + if b.id in fixed: + lo = max(lo, fixed[b.id]) + hi = min(hi, fixed[b.id]) + for check in scenario.constraints: + if check.metric not in {"sulfur", "t95", "cetane_number"}: + raise ValueError(f"Автоподбор не поддерживает ограничение {check.metric}") + left_metric = getattr(a, check.metric) + right_metric = getattr(b, check.metric) + if left_metric is None or right_metric is None: + return 1.0, 0.0 + field = ( + "upper" + if check.use_upper_estimate + else "lower" + if check.use_lower_estimate + else "value" + ) + left = getattr(left_metric, field) + right = getattr(right_metric, field) + if left is None or right is None: + return 1.0, 0.0 + if check.metric == "sulfur": + intercept, slope, correction = diesel * left, right - left, 0.0 + else: + intercept, slope = left, (right - left) / diesel + correction = _gain(scenario, dose) if check.metric == "cetane_number" else 0.0 + if check.upper is not None: + lo, hi = _bound((lo, hi), slope, intercept + correction, check.upper, upper=True) + if check.lower is not None: + lo, hi = _bound((lo, hi), slope, intercept + correction, check.lower, upper=False) + return lo, hi + + +def assist_blend( + scenario: ScenarioConfig, + desired: dict[str, float], + desired_dose: float, + locked: frozenset[str] = frozenset(), +) -> tuple[BlendOption, ...]: + """Return three distinct useful recipes; desired and locked fractions use 0..1 units.""" + found = feasible_blends(scenario, desired, desired_dose, locked) + if not found: + return () + a, b = scenario.blend_components + + def distance(item: BlendOption) -> float: + return sum(abs(item.fractions[key] - desired[key]) for key in desired) + abs( + item.dose - desired_dose + ) + + def cost(item: BlendOption) -> float: + for assessment in item.evaluation.assessments: + metric = assessment.metrics.get("cost_proxy") + if metric and metric.value is not None: + return metric.value + return float("inf") + + rankings = ( + ("Ближе к введённому", lambda item: (distance(item), cost(item))), + ("Ниже стоимость", lambda item: (cost(item), distance(item))), + ("Меньше присадки", lambda item: (item.dose, distance(item), cost(item))), + ) + chosen: list[BlendOption] = [] + for label, key in rankings: + best = min(found, key=lambda item: (*key(item), item.dose, item.fractions[b.id])) + if not any( + abs(best.dose - other.dose) < 1e-10 + and abs(best.fractions[b.id] - other.fractions[b.id]) < 1e-10 + for other in chosen + ): + chosen.append(BlendOption(label, best.fractions, best.dose, best.evaluation)) + return tuple(chosen) + + +def feasible_blends( + scenario: ScenarioConfig, + desired: dict[str, float], + desired_dose: float, + locked: frozenset[str] = frozenset(), +) -> tuple[BlendOption, ...]: + """Evaluate interval boundaries and nearest points through the shared hard checks.""" + if scenario.mode.value != "model_demo" or len(scenario.blend_components) != 2: + raise ValueError("Автоподбор доступен только для модельной смеси из двух компонентов") + a, b = scenario.blend_components + if set(desired) != {a.id, b.id}: + raise ValueError("Укажите доли обоих компонентов") + if any(value < 0 or value > 1 for value in (*desired.values(), desired_dose)): + raise ValueError("Доли должны быть в пределах 0–100%") + additive = scenario.cetane_additive + doses: tuple[float, ...] + if additive is None: + doses = (0.0,) + elif "dose" in locked: + doses = (desired_dose,) + else: + steps = round(additive.max_mass_fraction / additive.fraction_step) + doses = tuple( + sorted({desired_dose, *(i * additive.fraction_step for i in range(steps + 1))}) + ) + state = _model_demo_state(datetime.now(UTC), scenario) + assert scenario.total_mass_t is not None + found: list[BlendOption] = [] + fixed = {name: desired[name] for name in (a.id, b.id) if name in locked} + for dose in doses: + if dose < 0 or dose > (additive.max_mass_fraction if additive else 0.0) + 1e-9: + continue + if additive and dose * scenario.total_mass_t > additive.available_mass_t + 1e-9: + continue + lo, hi = _shares(scenario, dose, fixed) + if lo > hi + 1e-9: + continue + # Endpoints, closest shares, and the former grid's useful interior points. + points = { + max(lo, min(hi, value)) + for value in ( + lo, + hi, + desired[b.id], + 1 - dose - desired[a.id], + (lo + hi) / 2, + ) + } + for share_b in sorted(points): + fractions = {a.id: 1 - dose - share_b, b.id: share_b} + candidate = CandidateAction( + id=f"assisted:{a.id}={fractions[a.id]:.12g},{b.id}={share_b:.12g},dose={dose:.12g}", + kind=CandidateKind.BLEND, + blend_mass_fractions=fractions, + additive_mass_fraction=dose, + horizon_minutes=60, + is_model_scenario=True, + ) + evaluation = evaluate_candidates(state, (candidate,), scenario)[0] + if evaluation.feasible: + found.append(BlendOption("", fractions, dose, evaluation)) + return tuple(found) diff --git a/source/agents/optimizer.py b/source/agents/optimizer.py index 60d75ac..8d093fd 100644 --- a/source/agents/optimizer.py +++ b/source/agents/optimizer.py @@ -98,6 +98,31 @@ def generate_candidates( is_model_scenario=scenario.mode is OperationMode.MODEL_DEMO, ) ) + # Preserve established scenario choices; search continuously if the grid + # falsely reports that no admissible recipe exists. + if any(item.feasible for item in evaluate_candidates(state, tuple(candidates), scenario)): + return tuple(candidates) + from source.agents.blend_assist import feasible_blends + + seen = { + ( + round(item.blend_mass_fractions.get(component_b, -1), 9), + round(item.additive_mass_fraction, 9), + ) + for item in candidates + } + for option in feasible_blends( + scenario, + dict(scenario.current_blend_mass_fractions), + scenario.current_additive_mass_fraction, + ): + key = (round(option.fractions[component_b], 9), round(option.dose, 9)) + if key not in seen: + if len(candidates) >= max_candidates: + candidates[-1] = option.evaluation.candidate + break + candidates.append(option.evaluation.candidate) + seen.add(key) return tuple(candidates) diff --git a/source/agents/quality.py b/source/agents/quality.py index 19ee047..b563f2d 100644 --- a/source/agents/quality.py +++ b/source/agents/quality.py @@ -171,11 +171,24 @@ def _forecast_history_quality( applicability = check_applicability(features) if not getattr(applicability, "available", False): code = str(getattr(applicability, "reason_code", "OUT_OF_DOMAIN")) + detail = "Forecast features are missing or outside the validated training domain." + if code == "OUT_OF_DOMAIN": + age_features = [ + name + for name in getattr(applicability, "violations", ()) + if str(name).endswith("_age_minutes") + ] + if age_features: + detail = ( + "The saved model has not been validated for this measurement age. " + "Choose a historical measurement time or use a model " + "validated for delayed data." + ) return _unavailable( state, code, target_signal, - "Forecast features are missing or outside the validated training domain.", + detail, ) prediction = float(predict(features)[0]) upper: float | None = None diff --git a/source/config.py b/source/config.py index 5f13924..c6600f9 100644 --- a/source/config.py +++ b/source/config.py @@ -96,11 +96,21 @@ def load_telemetry_rules(path: str | Path) -> dict[str, dict[str, object]]: return result -def config_fingerprint(config: RuntimeConfig) -> str: - """Stable JSON used by the preparation manifest hash.""" - return json.dumps( - config.model_dump(mode="json"), ensure_ascii=False, sort_keys=True, separators=(",", ":") - ) +def config_fingerprint(config: RuntimeConfig, *, preparation_only: bool = False) -> str: + """Stable preparation identity; preserve the v1.0 full-config fingerprint.""" + payload = config.model_dump(mode="json") + if preparation_only: + payload = { + name: payload[name] + for name in ( + "source_timezone", + "lims_delay_hours", + "tag_dictionary_path", + "telemetry_rules_path", + "materials_dir", + ) + } + return json.dumps(payload, ensure_ascii=False, sort_keys=True, separators=(",", ":")) __all__ = [ diff --git a/source/contracts.py b/source/contracts.py index e61ad78..efb4a71 100644 --- a/source/contracts.py +++ b/source/contracts.py @@ -600,7 +600,7 @@ class SourceArtifact(ContractModel): class DatasetManifest(ContractModel): - schema_version: Literal["1.0"] = SCHEMA_VERSION + schema_version: Literal["1.0", "1.1"] = SCHEMA_VERSION dataset_id: Annotated[str, Field(min_length=12, max_length=12)] preparation_version: str created_at: AwareDatetime @@ -612,6 +612,17 @@ class DatasetManifest(ContractModel): telemetry_rules_sha256: Annotated[str, Field(pattern=r"^[0-9a-f]{64}$")] | None = None row_counts: dict[str, Annotated[int, Field(ge=0)]] time_ranges: dict[str, tuple[AwareDatetime, AwareDatetime] | None] + prepared_sha256: dict[str, Annotated[str, Field(pattern=r"^[0-9a-f]{64}$")]] | None = None + + @model_validator(mode="after") + def validate_prepared_hashes(self) -> "DatasetManifest": + if self.schema_version == "1.1" and ( + self.prepared_sha256 is None + or set(self.prepared_sha256) + != {"telemetry.csv.gz", "quality.csv.gz", "issues.csv.gz", "feature_order.json"} + ): + raise ValueError("schema 1.1 requires hashes of every prepared data file") + return self __all__ = [ diff --git a/source/data/prepare.py b/source/data/prepare.py index b3223e2..001dbcd 100644 --- a/source/data/prepare.py +++ b/source/data/prepare.py @@ -8,9 +8,11 @@ import subprocess import tarfile import tempfile +import threading import time from dataclasses import dataclass from datetime import UTC, datetime +from functools import lru_cache from pathlib import Path from typing import Any from zoneinfo import ZoneInfo @@ -30,10 +32,18 @@ Validity, ) -PREPARATION_VERSION = "1" +PREPARATION_VERSION = "2" _ARCHIVE_MEMBERS = {"data/avt_tags.csv", "data/242000_tags.csv"} _PUBLISH_ATTEMPTS = 5 _PUBLISH_RETRY_SECONDS = 0.05 +_PREPARED_CACHE_LOCK = threading.Lock() +_PREPARED_FILES = ( + "manifest.json", + "telemetry.csv.gz", + "quality.csv.gz", + "issues.csv.gz", + "feature_order.json", +) RAW_UNIT_MAP: dict[str, Unit] = { "°С": Unit.CELSIUS, @@ -168,24 +178,44 @@ def _publish_directory(temporary: Path, target: Path) -> None: time.sleep(_PUBLISH_RETRY_SECONDS * (2**attempt)) +def _validate_archive_entries(entries: list[tuple[str, bool, bool]]) -> None: + """Accept only the two CSV files and their real parent directory.""" + expected = _ARCHIVE_MEMBERS | {"data"} + seen: set[str] = set() + for name, is_directory, is_file in entries: + # A single leading ./ is harmless; stripping arbitrary dots would hide ../. + normalized = name.removeprefix("./") + if normalized == "data/": + normalized = "data" + if ( + "\\" in name + or normalized not in expected + or normalized in seen + or (normalized == "data" and not is_directory) + or (normalized != "data" and not is_file) + ): + raise ValueError(f"unsafe archive member: {name}") + seen.add(normalized) + if seen != expected: + raise ValueError(f"unexpected archive members: {sorted(seen)}") + + def _extract_telemetry(archive: Path, destination: Path) -> Path: """Validate archive members before extraction to a temporary directory.""" try: with tarfile.open(archive) as stream: - members = { - member.name.replace("\\", "/").lstrip("./").rstrip("/") - for member in stream.getmembers() - } - if members != _ARCHIVE_MEMBERS | {"data"}: - raise ValueError(f"unexpected archive members: {sorted(members)}") + members = stream.getmembers() + _validate_archive_entries( + [(member.name, member.isdir(), member.isfile()) for member in members] + ) destination_root = destination.resolve() - for member in stream.getmembers(): + for member in members: target = (destination / member.name).resolve() try: target.relative_to(destination_root) except ValueError as exc: raise ValueError(f"unsafe archive member: {member.name}") from exc - stream.extractall(destination) + stream.extractall(destination, members=members) return destination / "data" except tarfile.TarError: pass @@ -193,9 +223,17 @@ def _extract_telemetry(archive: Path, destination: Path) -> Path: listing = subprocess.run( ["tar", "-tf", str(archive)], check=True, capture_output=True, text=True ).stdout.splitlines() - members = {name.replace("\\", "/").lstrip("./").rstrip("/") for name in listing} - if members != _ARCHIVE_MEMBERS | {"data"}: - raise ValueError(f"unexpected archive members: {sorted(members)}") + details = subprocess.run( + ["tar", "-tvf", str(archive)], check=True, capture_output=True, text=True + ).stdout.splitlines() + if len(listing) != len(details): + raise ValueError("archive member listings disagree") + _validate_archive_entries( + [ + (name, detail.startswith("d"), detail.startswith("-")) + for name, detail in zip(listing, details, strict=True) + ] + ) subprocess.run(["tar", "-xf", str(archive), "-C", str(destination)], check=True) return destination / "data" @@ -288,7 +326,9 @@ def prepare_dataset(materials_dir: Path, config: RuntimeConfig) -> PreparedData: _source_artifact(path, materials_dir.parent) for path in (archive, pak_path, lims_path, tag_path, rules_path) ) - config_hash = hashlib.sha256(config_fingerprint(config).encode()).hexdigest() + config_hash = hashlib.sha256( + config_fingerprint(config, preparation_only=True).encode() + ).hexdigest() tag_hash = _sha256(tag_path) telemetry_rules_hash = _sha256(rules_path) identity_parts = sorted(item.sha256 for item in sources) + [ @@ -352,9 +392,9 @@ def write_prepared_dataset(data: PreparedData, output_root: Path) -> Path: } if manifest_path.exists(): existing = DatasetManifest.model_validate_json(manifest_path.read_text(encoding="utf-8")) - if existing.model_dump(exclude={"created_at"}) != data.manifest.model_dump( - exclude={"created_at"} - ): + if existing.model_dump( + exclude={"created_at", "schema_version", "prepared_sha256"} + ) != data.manifest.model_dump(exclude={"created_at", "schema_version", "prepared_sha256"}): raise FileExistsError(f"dataset_id collision at {target}") missing = [ path.name @@ -363,20 +403,28 @@ def write_prepared_dataset(data: PreparedData, output_root: Path) -> Path: ] if missing: raise FileExistsError(f"incomplete prepared dataset at {target}: missing {missing}") + _verify_prepared_files(target, existing) return target temporary = Path(tempfile.mkdtemp(prefix=f".{target.name}.", dir=output_root)) try: for name, frame in files.items(): frame.to_csv(temporary / name, index=False, compression="gzip") - (temporary / "manifest.json").write_text( - data.manifest.model_dump_json(indent=2), encoding="utf-8", newline="\n" - ) (temporary / "feature_order.json").write_text( json.dumps(list(data.feature_order), ensure_ascii=False, indent=2), encoding="utf-8", newline="\n", ) + hashes = {name: _sha256(temporary / name) for name in _PREPARED_FILES[1:]} + published_manifest = data.manifest.model_copy( + update={ + "schema_version": "1.1" if data.manifest.preparation_version == "2" else "1.0", + "prepared_sha256": hashes, + } + ) + (temporary / "manifest.json").write_text( + published_manifest.model_dump_json(indent=2), encoding="utf-8", newline="\n" + ) required = ( temporary / "manifest.json", temporary / "feature_order.json", @@ -391,24 +439,60 @@ def write_prepared_dataset(data: PreparedData, output_root: Path) -> Path: return target +def _verify_prepared_files(dataset_path: Path, manifest: DatasetManifest) -> None: + expected = manifest.prepared_sha256 + if expected is None: + release_root = dataset_path.parent.parent.parent + release_file = release_root / "config/release_manifest.json" + if release_file.is_file(): + release = json.loads(release_file.read_text(encoding="utf-8")) + if (release_root / release.get("prepared_dataset", "")).resolve() == dataset_path: + expected = release.get("prepared_sha256") + if not isinstance(expected, dict): + raise ValueError("pinned legacy dataset has no trusted prepared hashes") + if expected is None: + return # Unpinned legacy development fixtures cannot be retroactively authenticated. + if set(expected) != set(_PREPARED_FILES[1:]): + raise ValueError("prepared file hash manifest is incomplete") + for name, digest in expected.items(): + if not isinstance(digest, str) or _sha256(dataset_path / name) != digest: + raise ValueError(f"prepared dataset checksum mismatch: {name}") + + +@lru_cache(maxsize=1) +def _read_prepared_dataset(dataset_path: Path, file_hashes: tuple[str, ...]) -> PreparedData: + """Keep one verified, unmodified parsed dataset in process memory.""" + required = tuple(dataset_path / name for name in _PREPARED_FILES) + manifest_path, telemetry_path, quality_path, issues_path, feature_order_path = required + data = PreparedData( + telemetry=pd.read_csv(telemetry_path), + quality=pd.read_csv(quality_path), + issues=pd.read_csv(issues_path), + manifest=DatasetManifest.model_validate_json(manifest_path.read_text(encoding="utf-8")), + feature_order=tuple(json.loads(feature_order_path.read_text(encoding="utf-8"))), + ) + if tuple(_sha256(path) for path in required) != file_hashes: + raise OSError(f"prepared dataset changed while loading: {dataset_path}") + return data + + def load_prepared_dataset(dataset_path: Path) -> PreparedData: - """Load a prepared dataset written by :func:`write_prepared_dataset`.""" + """Load prepared data, reusing a content-checked in-memory parse when possible.""" dataset_path = dataset_path.resolve() - manifest_path = dataset_path / "manifest.json" - telemetry_path = dataset_path / "telemetry.csv.gz" - quality_path = dataset_path / "quality.csv.gz" - issues_path = dataset_path / "issues.csv.gz" - feature_order_path = dataset_path / "feature_order.json" - required = (manifest_path, telemetry_path, quality_path, issues_path, feature_order_path) + required = tuple(dataset_path / name for name in _PREPARED_FILES) missing = [path.name for path in required if not path.is_file()] if missing: raise FileNotFoundError(f"{dataset_path}: missing prepared dataset files: {missing}") + file_hashes = tuple(_sha256(path) for path in required) + with _PREPARED_CACHE_LOCK: + cached = _read_prepared_dataset(dataset_path, file_hashes) + _verify_prepared_files(dataset_path, cached.manifest) return PreparedData( - telemetry=pd.read_csv(telemetry_path), - quality=pd.read_csv(quality_path), - issues=pd.read_csv(issues_path), - manifest=DatasetManifest.model_validate_json(manifest_path.read_text(encoding="utf-8")), - feature_order=tuple(json.loads(feature_order_path.read_text(encoding="utf-8"))), + telemetry=cached.telemetry.copy(), + quality=cached.quality.copy(), + issues=cached.issues.copy(), + manifest=cached.manifest.model_copy(deep=True), + feature_order=cached.feature_order, ) diff --git a/source/data/state.py b/source/data/state.py index 64d7516..c519469 100644 --- a/source/data/state.py +++ b/source/data/state.py @@ -4,6 +4,7 @@ import hashlib import json +from dataclasses import replace from datetime import UTC, datetime import pandas as pd @@ -56,6 +57,31 @@ def _stage_from_signal_id(signal_id: str) -> Stage: return Stage.HYDROTREATMENT +def _utc_frame(frame: pd.DataFrame, columns: tuple[str, ...]) -> pd.DataFrame: + """Normalize timestamp columns once; reuse frames already normalized to UTC.""" + for column in columns: + dtype = frame[column].dtype + if not isinstance(dtype, pd.DatetimeTZDtype) or dtype.tz != UTC: + break + else: + return frame + normalized = frame.copy() + for column in columns: + normalized[column] = pd.to_datetime(normalized[column], utc=True) + return normalized + + +def prepare_state_data(data: PreparedData) -> PreparedData: + """Return a per-request dataset with UTC timestamps for repeated state builds.""" + quality = _utc_frame(data.quality, ("measured_at", "available_at")) + telemetry = ( + _utc_frame(data.telemetry, ("timestamp",)) + if "timestamp" in data.telemetry.columns + else data.telemetry + ) + return replace(data, quality=quality, telemetry=telemetry) + + def _telemetry_candidates( frame: pd.DataFrame, signal_id: str, @@ -98,13 +124,10 @@ def build_state( if as_of.tzinfo is None or as_of.utcoffset() is None: raise ValueError("as_of must be timezone-aware") as_of = as_of.astimezone(UTC) - frame = data.quality.copy() - frame["measured_at"] = pd.to_datetime(frame["measured_at"], utc=True) - frame["available_at"] = pd.to_datetime(frame["available_at"], utc=True) + prepared = prepare_state_data(data) + frame = prepared.quality visible = frame[(frame["measured_at"] <= as_of) & (frame["available_at"] <= as_of)] - telemetry = data.telemetry.copy() - if "timestamp" in telemetry.columns: - telemetry["timestamp"] = pd.to_datetime(telemetry["timestamp"], utc=True) + telemetry = prepared.telemetry control_units = {control.signal_id: control.unit for control in scenario.controls} snapshots: dict[str, SignalSnapshot] = {} diff --git a/source/desktop.py b/source/desktop.py index 87c9807..be15e3f 100644 --- a/source/desktop.py +++ b/source/desktop.py @@ -25,7 +25,7 @@ def main() -> None: from source.ui import PetrolCodeApp from source.ui_charts import HistoryChart from source.ui_data import discover_ui_context, ui_hybrid_snapshot, ui_stage_snapshot - from source.ui_what_if import editor_values, scenario_from_editor + from source.ui_what_if import assist_editor_values, editor_values, scenario_from_editor if len(sys.argv) == 3 and sys.argv[1] == "--smoke-report": destination = Path(sys.argv[2]) @@ -86,10 +86,18 @@ def descendants(widget: tk.Misc) -> Iterator[tk.Misc]: edited = run_model_demo(preset.id, scenario_override=scenario) if edited.scenario_id != "blend_risk_what_if" or edited.selected is None: raise RuntimeError("Packaged what-if calculation failed") + options = assist_editor_values( + preset, + editor_values(preset) | {"B.fraction": "10"}, + frozenset({"B.fraction"}), + ) + if not options or options[0][1]["B.fraction"] != "10": + raise RuntimeError("Packaged blend assistance failed") report["what_if"] = { "scenario_id": edited.scenario_id, "status": edited.status.value, "editor": "ok", + "assist": "ok", } hybrid = ui_hybrid_snapshot(dataset, model, at) report["hybrid"] = asdict(hybrid) diff --git a/source/journal.py b/source/journal.py index 834ac45..4546da0 100644 --- a/source/journal.py +++ b/source/journal.py @@ -35,6 +35,32 @@ def _json_line(payload: object) -> str: return json.dumps(payload, ensure_ascii=False, sort_keys=True) + "\n" +def write_run_status(run_dir: Path, run_id: str, status: str, detail: str = "") -> None: + """Publish a durable lifecycle marker for success, failure, or cancellation.""" + if status not in {"started", "completed", "failed", "cancelled", "interrupted"}: + raise ValueError("invalid run lifecycle status") + _atomic_write( + run_dir / run_id / "status.json", + json.dumps({"run_id": run_id, "status": status, "detail": detail}, ensure_ascii=False), + ) + + +def recover_interrupted_runs(run_dir: Path) -> int: + """Mark abandoned started runs after restarting the local application.""" + count = 0 + for path in run_dir.glob("*/status.json"): + try: + payload = json.loads(path.read_text(encoding="utf-8")) + except (OSError, json.JSONDecodeError): + continue + if payload.get("status") == "started": + write_run_status( + run_dir, path.parent.name, "interrupted", "process stopped before completion" + ) + count += 1 + return count + + def write_run_journal( run_dir: Path, state: ProcessState, @@ -46,6 +72,7 @@ def write_run_journal( selection_reason: str = "unknown", rejection_summary: dict[str, int] | None = None, features: pd.DataFrame | None = None, + editor_context: dict[str, object] | None = None, ) -> Path: """Write a compact reproducible record of one successful run.""" rejection_summary = rejection_summary or {} @@ -75,6 +102,7 @@ def write_run_journal( "state": state.model_dump(mode="json"), "scenario": scenario.model_dump(mode="json"), "context": context.model_dump(mode="json"), + "editor_context": editor_context, }, ensure_ascii=False, indent=2, @@ -104,4 +132,4 @@ def write_run_journal( return target -__all__ = ["write_run_journal"] +__all__ = ["write_run_journal", "write_run_status", "recover_interrupted_runs"] diff --git a/source/main.py b/source/main.py index 1ef277e..da5d0aa 100644 --- a/source/main.py +++ b/source/main.py @@ -38,6 +38,7 @@ prepare_dataset, write_prepared_dataset, ) +from source.data.state import prepare_state_data if TYPE_CHECKING: from source.ml.action_effects import ActionEnabledForecastModel @@ -130,6 +131,8 @@ def run_model_demo( run_dir: str | Path | None = None, *, scenario_override: ScenarioConfig | None = None, + recipe_mode: str = "optimize", + editor_context: dict[str, object] | None = None, ) -> Recommendation: """Run one deterministic stage-1 model-demo scenario on fixture data.""" from source.orchestrator import run_cycle @@ -152,6 +155,8 @@ def run_model_demo( config=config, context=DecisionContext(), run_dir=_resolve_path(run_dir if run_dir is not None else config.runs_dir, root), + recipe_mode=recipe_mode, + editor_context=editor_context, ) @@ -173,10 +178,30 @@ def _validate_current_prepared_dataset( data: PreparedData, config: RuntimeConfig, root: Path, + dataset_path: Path, ) -> None: """Fail early when a prepared dataset was built from older runtime inputs.""" + preparation_hash = hashlib.sha256( + config_fingerprint(config, preparation_only=True).encode() + ).hexdigest() + config_hash = hashlib.sha256( + config_fingerprint(config, preparation_only=False).encode() + ).hexdigest() + if data.manifest.schema_version == "1.1": + config_hash = preparation_hash + else: + release_path = root / "config/release_manifest.json" + if release_path.is_file(): + release = json.loads(release_path.read_text(encoding="utf-8")) + pinned = _resolve_path(release.get("prepared_dataset", ""), root).resolve() + if ( + pinned == dataset_path.resolve() + and release.get("legacy_config_sha256") == data.manifest.config_sha256 + ): + if release.get("preparation_config_sha256") == preparation_hash: + config_hash = data.manifest.config_sha256 expected = { - "config_sha256": hashlib.sha256(config_fingerprint(config).encode()).hexdigest(), + "config_sha256": config_hash, "tag_dictionary_sha256": _sha256_file(_resolve_path(config.tag_dictionary_path, root)), "telemetry_rules_sha256": _sha256_file(_resolve_path(config.telemetry_rules_path, root)), } @@ -204,7 +229,7 @@ def _load_current_prepared_dataset( root: Path, ) -> PreparedData: data = load_prepared_dataset(_resolve_path(dataset, root)) - _validate_current_prepared_dataset(data, config, root) + _validate_current_prepared_dataset(data, config, root, _resolve_path(dataset, root)) return data @@ -927,6 +952,7 @@ def history_interval_command( cancelled: Callable[[], bool] | None = None, ) -> dict[str, object]: """Replay a user-selected historical interval without fitting or tuning.""" + from source.ml.features import prepare_feature_batch, prepare_feature_sources from source.orchestrator import run_cycle if start.tzinfo is None or end.tzinfo is None: @@ -943,6 +969,7 @@ def history_interval_command( config = load_runtime_config(_resolve_path(config_path, root)) data = _load_current_prepared_dataset(dataset, config, root) + data = prepare_state_data(data) forecast_model = _load_trusted_model(model, data, root) scenario = load_scenario(root / "config/scenarios/history.json") output_dir = _resolve_path(run_dir if run_dir is not None else config.runs_dir, root) @@ -957,9 +984,15 @@ def history_interval_command( counts: dict[str, int] = {} reason_counts: dict[str, int] = {} available_forecasts = 0 - for timestamp in timestamps: + feature_batch = None + feature_sources = prepare_feature_sources(data, forecast_model) + for index, timestamp in enumerate(timestamps): if cancelled is not None and cancelled(): break + if index % 24 == 0: + feature_batch = prepare_feature_batch( + data, timestamps[index : index + 24], forecast_model, sources=feature_sources + ) result = run_cycle( data=data, model=forecast_model, @@ -968,6 +1001,7 @@ def history_interval_command( config=config, context=DecisionContext(), run_dir=output_dir, + feature_batch=feature_batch, ) carrier = result.selected or result.baseline sulfur = None diff --git a/source/ml/features.py b/source/ml/features.py index 533f399..81122d5 100644 --- a/source/ml/features.py +++ b/source/ml/features.py @@ -47,6 +47,28 @@ class SupervisedDataset: excluded_counts: dict[str, int] +@dataclass(frozen=True, eq=False) +class FeatureBatch: + """Features for one bounded history interval chunk and its exact inputs.""" + + data: PreparedData + model: object + queries: pd.DatetimeIndex + frame: FeatureFrame + + +@dataclass(frozen=True, eq=False) +class FeatureSources: + """Validated and sorted source tables shared by history replay chunks.""" + + data: PreparedData + model: object + telemetry: pd.DataFrame + target_history: pd.DataFrame + published_history: pd.DataFrame + feature_spec: tuple[tuple[str, ...], tuple[str, ...], str, SourceKind, str] + + def _enum_value(value: object) -> str: """Return a stable string for plain strings and string enums.""" raw = getattr(value, "value", value) @@ -250,7 +272,10 @@ def _window_values( def _last_available_quality( - observations: pd.DataFrame, queries: pd.DatetimeIndex + observations: pd.DataFrame, + queries: pd.DatetimeIndex, + *, + events_sorted: bool = False, ) -> tuple[np.ndarray[Any, np.dtype[np.float64]], np.ndarray[Any, np.dtype[np.float64]]]: """Select the latest measured observation among those published by each as_of.""" values = np.full(len(queries), np.nan, dtype=float) @@ -258,9 +283,13 @@ def _last_available_quality( if observations.empty or queries.empty: return values, ages - events = observations.sort_values( - ["available_at", "measured_at", "observation_id"], kind="stable" - ).reset_index(drop=True) + events = ( + observations + if events_sorted + else observations.sort_values( + ["available_at", "measured_at", "observation_id"], kind="stable" + ).reset_index(drop=True) + ) available = pd.DatetimeIndex(events["available_at"]).as_unit("ns").astype("int64").to_numpy() measured = pd.DatetimeIndex(events["measured_at"]).as_unit("ns").astype("int64").to_numpy() query_times = queries.as_unit("ns").astype("int64").to_numpy() @@ -323,15 +352,30 @@ def _feature_matrix( target_signal_id: str, target_source: SourceKind, target_unit: str, + *, + sources: FeatureSources | None = None, ) -> pd.DataFrame: """Build the common training/serving matrix for already-normalized UTC times.""" - telemetry = _telemetry(data, telemetry_signals) + telemetry = sources.telemetry if sources is not None else _telemetry(data, telemetry_signals) lag_values = _lag_values(telemetry, queries, telemetry_signals) window_values = _window_values(telemetry, queries, telemetry_signals) - target_history = _quality_rows(data, target_signal_id, target_source, target_unit) - last_values, ages = _last_available_quality(target_history, queries) + target_history = ( + sources.target_history + if sources is not None + else _quality_rows(data, target_signal_id, target_source, target_unit) + ) + published_history = ( + sources.published_history + if sources is not None + else target_history.sort_values( + ["available_at", "measured_at", "observation_id"], kind="stable" + ).reset_index(drop=True) + ) + last_values, ages = _last_available_quality(published_history, queries, events_sorted=True) quality_lags = { - lag: _last_available_quality(target_history, queries - pd.Timedelta(minutes=lag))[0] + lag: _last_available_quality( + published_history, queries - pd.Timedelta(minutes=lag), events_sorted=True + )[0] for lag in LAG_MINUTES } quality_windows = _quality_window_values(target_history, queries) @@ -442,19 +486,10 @@ def _metadata_value(metadata: object, *names: str) -> object: raise ValueError(f"model metadata is missing {names[0]}") -def build_features( - data: PreparedData, - as_of: datetime, - state: ProcessState, - model: object, -) -> FeatureFrame: - """Build one serving row in the exact feature order stored by the model.""" - query = _as_utc(as_of, "as_of") - if _as_utc(state.as_of, "state.as_of") != query: - raise ValueError("state.as_of must match feature as_of") - if state.dataset_id != data.manifest.dataset_id: - raise ValueError("state and prepared data dataset_id mismatch") - +def _model_feature_spec( + data: PreparedData, model: object +) -> tuple[tuple[str, ...], tuple[str, ...], str, SourceKind, str]: + """Resolve the frozen model's feature contract once per feature build.""" metadata = getattr(model, "metadata", None) if metadata is None: raise ValueError("model needs metadata") @@ -483,6 +518,107 @@ def build_features( if signal in data.telemetry.columns and any(name.startswith(f"{signal}__") for name in feature_names) ) + if isinstance(processing, dict) and isinstance(processing.get("feature_definition"), dict): + definition = processing["feature_definition"] + expected_recipe = { + "horizon_minutes": _metadata_value(metadata, "horizon_minutes"), + "lags_minutes": list(LAG_MINUTES), + "rolling_windows_minutes": list(WINDOW_MINUTES), + "telemetry_max_age_minutes": TELEMETRY_MAX_AGE_MINUTES, + "telemetry_signals": list(telemetry_signals), + "quality_history_source": target_source.value, + "autoregressive_target_signal": target_signal_id, + } + mismatched = [ + key + for key, value in expected_recipe.items() + if key in definition + and ( + sorted(definition[key]) != sorted(telemetry_signals) + if key == "telemetry_signals" and isinstance(definition[key], list) + else definition[key] != value + ) + ] + if mismatched: + raise ValueError("model feature recipe is incompatible with serving code") + return feature_names, telemetry_signals, target_signal_id, target_source, target_unit + + +def prepare_feature_sources(data: PreparedData, model: object) -> FeatureSources: + """Parse the large source tables once for a historical interval.""" + spec = _model_feature_spec(data, model) + _, telemetry_signals, target_signal_id, target_source, target_unit = spec + target_history = _quality_rows(data, target_signal_id, target_source, target_unit) + published_history = target_history.sort_values( + ["available_at", "measured_at", "observation_id"], kind="stable" + ).reset_index(drop=True) + return FeatureSources( + data, + model, + _telemetry(data, telemetry_signals), + target_history, + published_history, + spec, + ) + + +def prepare_feature_batch( + data: PreparedData, + queries: pd.DatetimeIndex, + model: object, + *, + sources: FeatureSources | None = None, +) -> FeatureBatch: + """Build feature rows together for a bounded, UTC-aware replay chunk.""" + if queries.empty or queries.tz is None: + raise ValueError("feature batch requires timezone-aware, nonempty queries") + utc_queries = queries.tz_convert("UTC") + if sources is not None and (sources.data is not data or sources.model is not model): + raise ValueError("feature sources inputs do not match this replay") + feature_names, telemetry_signals, target_signal_id, target_source, target_unit = ( + sources.feature_spec if sources is not None else _model_feature_spec(data, model) + ) + generated = _feature_matrix( + data, + utc_queries, + telemetry_signals, + target_signal_id, + target_source, + target_unit, + sources=sources, + ) + missing = set(feature_names).difference(generated.columns) + if missing: + raise ValueError(f"model requests unknown features: {sorted(missing)}") + return FeatureBatch(data, model, utc_queries, generated.loc[:, list(feature_names)]) + + +def build_features( + data: PreparedData, + as_of: datetime, + state: ProcessState, + model: object, + *, + batch: FeatureBatch | None = None, +) -> FeatureFrame: + """Build one serving row in the exact feature order stored by the model.""" + query = _as_utc(as_of, "as_of") + if _as_utc(state.as_of, "state.as_of") != query: + raise ValueError("state.as_of must match feature as_of") + if state.dataset_id != data.manifest.dataset_id: + raise ValueError("state and prepared data dataset_id mismatch") + feature_names, telemetry_signals, target_signal_id, target_source, target_unit = ( + _model_feature_spec(data, model) + ) + if batch is not None: + if batch.data is not data or batch.model is not model: + raise ValueError("feature batch inputs do not match this replay") + if tuple(batch.frame.columns) != feature_names: + raise ValueError("feature batch columns do not match model") + position = batch.queries.get_indexer(pd.DatetimeIndex([query]))[0] + if position < 0: + raise ValueError("feature batch does not contain as_of") + return batch.frame.iloc[[position]].reset_index(drop=True) generated = _feature_matrix( data, pd.DatetimeIndex([query]), @@ -498,12 +634,16 @@ def build_features( __all__ = [ + "FeatureBatch", "FeatureFrame", + "FeatureSources", "LAG_MINUTES", "SUPERVISED_COLUMNS", "SupervisedDataset", "WINDOW_MINUTES", "baseline_feature_name", "build_features", + "prepare_feature_batch", + "prepare_feature_sources", "build_supervised_dataset", ] diff --git a/source/orchestrator.py b/source/orchestrator.py index 32b87ba..9ea4790 100644 --- a/source/orchestrator.py +++ b/source/orchestrator.py @@ -38,10 +38,10 @@ Unit, ) from source.explain import build_explanation -from source.journal import write_run_journal +from source.journal import write_run_journal, write_run_status from source.ml.action_effects import ActionOutcome, VerifiedActionEffectModel from source.ml.controls import generate_setpoint_candidates -from source.ml.features import build_features +from source.ml.features import FeatureBatch, build_features from source.ml.policy import PolicyParameters, assess_change_policy @@ -438,7 +438,7 @@ def _recheck_selected( return rechecked -def run_cycle( +def _run_cycle( data: Any, as_of: datetime, model: Any, @@ -446,6 +446,10 @@ def run_cycle( config: RuntimeConfig, context: DecisionContext, run_dir: Path, + feature_batch: FeatureBatch | None = None, + recipe_mode: str = "optimize", + editor_context: dict[str, object] | None = None, + run_id: str | None = None, ) -> Recommendation: """Run one backend recommendation cycle and persist its journal.""" if as_of.tzinfo is None or as_of.utcoffset() is None: @@ -463,7 +467,11 @@ def run_cycle( state = _build_cycle_state(data, as_of, effective_scenario, config) features: pd.DataFrame | None = None if model is not None and effective_scenario.mode is OperationMode.HISTORY: - features = build_features(data, as_of, state, model) + features = ( + build_features(data, as_of, state, model) + if feature_batch is None + else build_features(data, as_of, state, model, batch=feature_batch) + ) action_model = _action_model(model) if effective_scenario.mode is OperationMode.HISTORY and action_model is not None: candidates = generate_setpoint_candidates( @@ -484,7 +492,22 @@ def run_cycle( model=model, ) else: - candidates = generate_candidates(state, effective_scenario, config) + if recipe_mode not in {"optimize", "evaluate"}: + raise ValueError("recipe_mode must be optimize or evaluate") + if recipe_mode == "evaluate" and effective_scenario.mode is not OperationMode.MODEL_DEMO: + raise ValueError("fixed recipe evaluation is available only in model_demo") + candidates = ( + ( + CandidateAction( + id="hold", + kind=CandidateKind.HOLD, + horizon_minutes=config.horizon_minutes, + is_model_scenario=True, + ), + ) + if recipe_mode == "evaluate" + else generate_candidates(state, effective_scenario, config) + ) evaluations = evaluate_candidates( state, candidates, effective_scenario, features=features, model=model ) @@ -514,7 +537,7 @@ def run_cycle( ) ) result = Recommendation( - run_id=str(uuid4()), + run_id=run_id or str(uuid4()), state_id=state.state_id, as_of=state.as_of, scenario_id=scenario.id, @@ -548,8 +571,45 @@ def run_cycle( selection_reason=selection_reason, rejection_summary=_rejection_summary(evaluations), features=features, + editor_context=editor_context, ) return result +def run_cycle( + data: Any, + as_of: datetime, + model: Any, + scenario: ScenarioConfig, + config: RuntimeConfig, + context: DecisionContext, + run_dir: Path, + feature_batch: FeatureBatch | None = None, + recipe_mode: str = "optimize", + editor_context: dict[str, object] | None = None, +) -> Recommendation: + """Run and record a complete decision lifecycle, including failed attempts.""" + run_id = str(uuid4()) + write_run_status(run_dir, run_id, "started") + try: + result = _run_cycle( + data, + as_of, + model, + scenario, + config, + context, + run_dir, + feature_batch, + recipe_mode, + editor_context, + run_id, + ) + write_run_status(run_dir, run_id, "completed") + return result + except Exception as exc: + write_run_status(run_dir, run_id, "failed", f"{type(exc).__name__}: {exc}") + raise + + __all__ = ["run_cycle"] diff --git a/source/ui.py b/source/ui.py index 6ab5983..e086eea 100644 --- a/source/ui.py +++ b/source/ui.py @@ -12,6 +12,7 @@ import sys import threading import tkinter as tk +from concurrent.futures import ThreadPoolExecutor from dataclasses import asdict, dataclass from datetime import datetime from functools import partial @@ -19,7 +20,7 @@ from tkinter import filedialog, messagebox, ttk from typing import Any, Callable, Sequence -from source.config import load_scenario +from source.config import load_runtime_config, load_scenario from source.contracts import ( CandidateEvaluation, ConstraintStatus, @@ -27,6 +28,7 @@ RecommendationStatus, ScenarioConfig, ) +from source.journal import recover_interrupted_runs from source.main import ( PROJECT_ROOT, action_shadow_estimate_command, @@ -294,7 +296,18 @@ def journal_entries(run_dir: Path) -> tuple[Path, ...]: """Return newest journal result files first without trusting partial directories.""" if not run_dir.exists(): return () - files = [path for path in run_dir.glob("*/result.json") if path.is_file()] + files = [] + for path in run_dir.glob("*/result.json"): + if not path.is_file(): + continue + status_path = path.parent / "status.json" + if status_path.is_file(): + try: + if json.loads(status_path.read_text(encoding="utf-8")).get("status") != "completed": + continue + except (OSError, json.JSONDecodeError): + continue + files.append(path) return tuple(sorted(files, key=lambda path: path.stat().st_mtime, reverse=True)) @@ -468,7 +481,11 @@ def __init__( self.status_var = tk.StringVar(value="Готово к расчёту") self._ui_events: queue.SimpleQueue[Callable[[], None]] = queue.SimpleQueue() self._closing = False + recover_interrupted_runs( + PROJECT_ROOT / load_runtime_config(PROJECT_ROOT / "config/runtime.toml").runs_dir + ) self._interval_cancel = threading.Event() + self._history_executor = ThreadPoolExecutor(max_workers=1, thread_name_prefix="history") self._shared_dataset = tk.StringVar(value=discover_ui_context().latest_dataset or "") self._shared_as_of = tk.StringVar(value=RELEASE_DEMO_AS_OF) self._result: Recommendation | None = None @@ -507,6 +524,7 @@ def _drain_ui_events(self) -> None: def destroy(self) -> None: self._closing = True self._interval_cancel.set() + self._history_executor.shutdown(wait=False, cancel_futures=True) if hasattr(self, "_poll_id"): self.after_cancel(self._poll_id) super().destroy() @@ -601,7 +619,10 @@ def _clear_body(self) -> None: def _set_nav(self, page: str) -> None: for key, button in self._nav_buttons.items(): - button.configure(fg="white" if key == page else "#C0C8CD") + button.configure( + fg="white" if key == page else "#C0C8CD", + bg="#31454B" if key == page else HEADER, + ) def show_page(self, page: str) -> None: page = PAGE_ALIASES.get(page, page) @@ -698,11 +719,26 @@ def _secondary_button(parent: tk.Misc, text: str, command: Any) -> tk.Button: def _field( parent: tk.Misc, label: str, variable: tk.StringVar, width: int | None = None ) -> None: - tk.Label(parent, text=label, bg=BG, fg=TEXT, font=("Segoe UI", 10, "bold")).pack(anchor="w") + tk.Label( + parent, + text=label, + bg=parent.cget("bg"), + fg=MUTED, + font=("Segoe UI", 10, "bold"), + ).pack(anchor="w") entry = tk.Entry( - parent, textvariable=variable, bg=SURFACE, fg=TEXT, bd=1, width=width or 20 + parent, + textvariable=variable, + bg=SURFACE, + fg=TEXT, + relief="flat", + bd=0, + highlightthickness=1, + highlightbackground=BORDER, + highlightcolor=TEAL, + width=width or 20, ) - entry.pack(fill="x", pady=(4, 10), ipady=6) + entry.pack(fill="x", pady=(6, 8), ipady=8) @staticmethod def _stage_tone(status: str) -> str: @@ -1306,13 +1342,57 @@ def settings(parent: tk.Misc) -> None: summary = tk.StringVar( value="Укажите время и запустите расчёт. Реальные действия отключены." ) + result_panel = self._surface(page) + result_panel.pack(fill="x", pady=(0, 14)) tk.Label( - page, textvariable=summary, bg=BG, fg=TEXT, anchor="w", justify="left", wraplength=1100 - ).pack(fill="x", pady=8) - chart = HistoryChart(page) - chart.pack(fill="x", pady=(0, 10)) - output = tk.Text(page, height=7, bg=SURFACE, fg=TEXT, bd=1, wrap="word") - output.pack(fill="both", expand=True) + result_panel, + text="ПОСЛЕДНИЙ РЕЗУЛЬТАТ", + bg=SURFACE, + fg=MUTED, + font=("Segoe UI", 9, "bold"), + ).pack(anchor="w", padx=20, pady=(14, 4)) + tk.Label( + result_panel, + textvariable=summary, + bg=SURFACE, + fg=TEXT, + anchor="w", + justify="left", + wraplength=1100, + font=("Segoe UI", 12, "bold"), + ).pack(fill="x", padx=20, pady=(0, 14)) + chart_panel = self._surface(page) + chart_panel.pack(fill="x", pady=(0, 14)) + tk.Label( + chart_panel, + text="Динамика серы", + bg=SURFACE, + fg=TEXT, + font=("Segoe UI", 14, "bold"), + ).pack(anchor="w", padx=20, pady=(16, 0)) + chart = HistoryChart(chart_panel) + chart.pack(fill="x", padx=12, pady=(0, 8)) + details = self._surface(page) + details.pack(fill="x") + tk.Label( + details, + text="Детали расчёта", + bg=SURFACE, + fg=TEXT, + font=("Segoe UI", 14, "bold"), + ).pack(anchor="w", padx=20, pady=(16, 10)) + output = tk.Text( + details, + height=8, + bg="#F8FAFA", + fg=TEXT, + bd=0, + wrap="word", + font=("Consolas", 10), + padx=14, + pady=10, + ) + output.pack(fill="x", padx=20, pady=(0, 18)) def render(view: UiHistoryReplayView) -> None: output.configure(state="normal") @@ -1335,6 +1415,8 @@ def calculate( as_of: str, action_model: str | None, ) -> None: + if request_id != self._history_request_id: + return view = ui_history_snapshot(dataset, model, as_of, action_model) def finish() -> None: @@ -1364,17 +1446,14 @@ def start_calculation() -> None: output.delete("1.0", "end") output.insert("1.0", "Расчёт исторического прогноза…") output.configure(state="disabled") - threading.Thread( - target=calculate, - args=( - self._history_request_id, - dataset_var.get().strip() or None, - model_var.get().strip() or None, - as_of_var.get(), - action_var.get().strip() or None, - ), - daemon=True, - ).start() + self._history_executor.submit( + calculate, + self._history_request_id, + dataset_var.get().strip() or None, + model_var.get().strip() or None, + as_of_var.get(), + action_var.get().strip() or None, + ) def calculate_interval() -> None: self._history_request_id += 1 @@ -1425,6 +1504,8 @@ def show_progress() -> None: progress(0, 0) def work() -> None: + if request_id != self._history_request_id or cancel.is_set(): + return payload = None try: payload = history_interval_command( @@ -1466,32 +1547,27 @@ def finish() -> None: except tk.TclError: return - threading.Thread(target=work, daemon=True).start() + self._history_executor.submit(work) actions = tk.Frame(page, bg=BG) - actions.pack(fill="x", pady=(0, 12), before=chart) - self._primary_button( - actions, - "Рассчитать на выбранный момент", - start_calculation, - ).pack(side="left") - self._secondary_button(actions, "Журнал", lambda: self.show_page("journal")).pack( - side="left", padx=12 + actions.pack(fill="x", pady=(0, 14), before=chart_panel) + for column in range(3): + actions.grid_columnconfigure(column, weight=1) + self._primary_button(actions, "Рассчитать на выбранный момент", start_calculation).grid( + row=0, column=0, sticky="ew", padx=(0, 8), pady=(0, 8) ) - self._secondary_button(actions, "Рассчитать интервал", calculate_interval).pack(side="left") + interval_button = self._secondary_button(actions, "Рассчитать интервал", calculate_interval) + interval_button.configure(fg=TEAL, font=("Segoe UI", 11, "bold")) + interval_button.grid(row=0, column=1, sticky="ew", padx=(0, 8), pady=(0, 8)) self._secondary_button( actions, "Отменить интервал", lambda: self._interval_cancel.set() - ).pack(side="left") + ).grid(row=0, column=2, sticky="ew", pady=(0, 8)) self._secondary_button( actions, "Экспортировать этот результат", self.export_history_result - ).pack(side="left", padx=12) - buttons = [child for child in actions.winfo_children() if isinstance(child, tk.Button)] - for button in buttons: - button.pack_forget() - for index, button in enumerate(buttons): - button.grid(row=index // 3, column=index % 3, sticky="ew", padx=4, pady=4) - for column in range(3): - actions.grid_columnconfigure(column, weight=1) + ).grid(row=1, column=0, sticky="ew", padx=(0, 8)) + self._secondary_button(actions, "Журнал", lambda: self.show_page("journal")).grid( + row=1, column=1, sticky="ew", padx=(0, 8) + ) render( self._history_view or UiHistoryReplayView( @@ -2067,11 +2143,22 @@ def open_what_if_dialog(self) -> None: ) def recalculate(scenario: ScenarioConfig) -> None: - self.calculate(scenario) + self.calculate(scenario, recipe_mode="evaluate") + + def recalculate_with_context(scenario: ScenarioConfig, context: dict[str, object]) -> None: + self.calculate(scenario, recipe_mode="evaluate", editor_context=context) - open_editor(self, preset, current, recalculate) + open_editor( + self, preset, current, recalculate, calculate_with_context=recalculate_with_context + ) - def calculate(self, scenario_override: ScenarioConfig | None = None) -> None: + def calculate( + self, + scenario_override: ScenarioConfig | None = None, + *, + recipe_mode: str = "optimize", + editor_context: dict[str, object] | None = None, + ) -> None: self._calculation_request_id += 1 request_id = self._calculation_request_id scenario_id = SCENARIO_LABELS[self.scenario_var.get()] @@ -2085,7 +2172,12 @@ def work() -> None: scenario = scenario_override or load_scenario( PROJECT_ROOT / f"config/scenarios/{scenario_id}.json" ) - result = run_model_demo(scenario_id, scenario_override=scenario_override) + result = run_model_demo( + scenario_id, + scenario_override=scenario_override, + recipe_mode=recipe_mode, + editor_context=editor_context, + ) except Exception as exc: # UI boundary: render backend failure without crashing Tk. self._post_ui(partial(self._calculation_failed, request_id, exc)) return diff --git a/source/ui_charts.py b/source/ui_charts.py index b3b47b4..5b0639a 100644 --- a/source/ui_charts.py +++ b/source/ui_charts.py @@ -80,13 +80,21 @@ def draw(self) -> None: self.delete("all") width, height = max(self.winfo_width(), 600), max(self.winfo_height(), 290) left, right, top, bottom = 62.0, float(width - 24), 52.0, float(height - 64) - self.create_text( - left, - 16, - anchor="w", - fill="#007D78", - text="● Прогноз серы ● Верхняя оценка (синяя) — Предел 10 мг/кг", - ) + for legend_x, color, label, dashed in ( + (left, "#007D78", "Прогноз серы", False), + (left + 164, "#3366BB", "Верхняя оценка", False), + (left + 354, "#D88400", "Предел 10 мг/кг", True), + ): + self.create_line( + legend_x, + 18, + legend_x + 22, + 18, + fill=color, + width=3, + dash=(6, 4) if dashed else (), + ) + self.create_text(legend_x + 30, 18, anchor="w", fill="#38485A", text=label) self.create_text( right, height - 10, diff --git a/source/ui_data.py b/source/ui_data.py index 74fa2b5..fb7260f 100644 --- a/source/ui_data.py +++ b/source/ui_data.py @@ -170,13 +170,22 @@ def list_prepared_datasets(root: Path = PROJECT_ROOT) -> tuple[Path, ...]: def _release_manifest(root: Path) -> dict[str, Any]: - """Read the optional release pin without making it a runtime dependency.""" + """Allow discovery only when no release manifest exists at all.""" path = root / "config/release_manifest.json" + if not path.exists(): + return {} try: payload = json.loads(path.read_text(encoding="utf-8")) - except (OSError, json.JSONDecodeError): - return {} - return payload if isinstance(payload, dict) else {} + except (OSError, json.JSONDecodeError) as exc: + raise ValueError(f"invalid release manifest: {path}") from exc + if not isinstance(payload, dict) or payload.get("schema_version") != "1.0": + raise ValueError("unsupported release manifest schema") + for key in ("release_id", "prepared_dataset", "forecast_artifact", "v2_artifact"): + if not isinstance(payload.get(key), str) or not payload[key]: + raise ValueError(f"release manifest needs {key}") + if "action_artifact" not in payload: + raise ValueError("release manifest needs action_artifact, null if disabled") + return payload def _manifest_dataset_id(path: Path) -> str | None: @@ -292,11 +301,7 @@ def configured_artifact( forecast_artifacts=forecast, action_artifacts=action, v2_artifacts=v2, - release_id=( - str(_release_manifest(root).get("release_id")) - if _release_manifest(root).get("release_id") - else None - ), + release_id=str(manifest["release_id"]) if manifest.get("release_id") else None, ) @@ -343,6 +348,9 @@ def ui_stage_snapshot( try: config = load_runtime_config(root / "config/runtime.toml") data = load_prepared_dataset(data_path) + data.telemetry["timestamp"] = pd.to_datetime( + data.telemetry["timestamp"], utc=True, errors="coerce" + ) timestamp = _parse_as_of(as_of) if as_of is not None else _default_as_of(data) tags = _tags_by_signal_id(root) groups = AVT_SIGNAL_GROUPS if canonical_page == "avt" else HYDROTREATING_SIGNAL_GROUPS diff --git a/source/ui_what_if.py b/source/ui_what_if.py index f593d79..61ac97d 100644 --- a/source/ui_what_if.py +++ b/source/ui_what_if.py @@ -5,8 +5,10 @@ import math import tkinter as tk from collections.abc import Callable +from functools import partial from tkinter import ttk +from source.agents.blend_assist import BlendOption, assist_blend from source.contracts import ScenarioConfig COMPONENT_FIELDS = ( @@ -128,50 +130,375 @@ def number(key: str) -> float: return ScenarioConfig.model_validate(data) +def assist_editor_values( + base: ScenarioConfig, values: dict[str, str], locked: frozenset[str] +) -> tuple[tuple[BlendOption, dict[str, str]], ...]: + """Parse physical inputs, then repair only the freely adjustable recipe fields.""" + if len(base.blend_components) != 2: + raise ValueError("Автоподбор поддерживает два компонента") + component_ids = tuple(item.id for item in base.blend_components) + + def fraction(key: str) -> float: + try: + result = float(values[key].strip().replace(",", ".")) / 100 + except (KeyError, ValueError) as exc: + raise ValueError(f"{key}: введите число от 0 до 100") from exc + if not math.isfinite(result) or result < 0 or result > 1: + raise ValueError(f"{key}: требуется число от 0 до 100") + return result + + desired = {key: fraction(f"{key}.fraction") for key in component_ids} + dose = fraction("dose") + temporary = dict(values) + # Only the recipe is temporarily balanced so scenario_from_editor can validate + # unchanged component properties before the actual constrained search. + try: + max_dose = float(values["max_dose"].strip().replace(",", ".")) / 100 + except (KeyError, ValueError) as exc: + raise ValueError("max_dose: введите максимум присадки 1, 2 или 3%") from exc + temporary_dose = min(dose, max_dose) if math.isfinite(max_dose) else dose + temporary["dose"] = f"{temporary_dose * 100:.12g}" + temporary[f"{component_ids[0]}.fraction"] = f"{(1 - temporary_dose) * 50:.12g}" + temporary[f"{component_ids[1]}.fraction"] = f"{(1 - temporary_dose) * 50:.12g}" + scenario = scenario_from_editor(base, temporary) + fixed = frozenset( + key.removesuffix(".fraction") for key in locked if key.endswith(".fraction") + ) | (frozenset({"dose"}) if "dose" in locked else frozenset()) + options = assist_blend(scenario, desired, dose, fixed) + if not options: + mass = scenario.total_mass_t or 0 + for component in scenario.blend_components: + if component.id in fixed and desired[component.id] * mass > component.available_mass_t: + needed = desired[component.id] * mass + raise ValueError( + f"Компонента {component.id} нужно {needed:.3g} т, " + f"доступно {component.available_mass_t:.3g} т. " + "Снимите фиксацию или измените долю." + ) + raise ValueError( + "Подходящей смеси нет при заданных свойствах, запасах и закреплённых долях. " + "Снимите фиксацию или измените указанное свойство или массу партии." + ) + result = [] + for option in options: + updated = dict(values) + for key, share in option.fractions.items(): + updated[f"{key}.fraction"] = f"{share * 100:.12g}" + updated["dose"] = f"{option.dose * 100:.12g}" + scenario_from_editor(base, updated) + result.append((option, updated)) + return tuple(result) + + def open_editor( parent: tk.Tk, preset: ScenarioConfig, current: ScenarioConfig, calculate: Callable[[ScenarioConfig], None], + *, + calculate_with_context: Callable[[ScenarioConfig, dict[str, object]], None] | None = None, ) -> tk.Toplevel: dialog = tk.Toplevel(parent) dialog.title("Синтетический what-if — не реальные данные") - dialog.geometry("940x680") - dialog.minsize(800, 620) + dialog.geometry("940x700") + dialog.minsize(800, 700) dialog.transient(parent) + dialog.configure(bg="#F5F7F7") + style = ttk.Style(dialog) + style.configure("WhatIf.TNotebook", background="#F5F7F7", borderwidth=0) + style.configure( + "WhatIf.TNotebook.Tab", + background="#EDF1F3", + foreground="#38485A", + padding=(16, 7), + font=("Segoe UI", 10, "bold"), + ) + style.map( + "WhatIf.TNotebook.Tab", + background=[("selected", "#FFFFFF")], + foreground=[("selected", "#007D78")], + ) + style.configure("WhatIf.TFrame", background="#FFFFFF") + style.configure("WhatIf.Footer.TFrame", background="#F5F7F7") + style.configure("WhatIf.TLabelframe", background="#FFFFFF", bordercolor="#D8E0E4") + style.configure( + "WhatIf.TLabelframe.Label", + background="#FFFFFF", + foreground="#101B28", + font=("Segoe UI", 11, "bold"), + ) + style.configure( + "WhatIf.TLabel", background="#FFFFFF", foreground="#38485A", font=("Segoe UI", 10) + ) + style.configure("WhatIf.TEntry", fieldbackground="#FFFFFF", foreground="#101B28", padding=3) + style.configure( + "WhatIf.TButton", + background="#FFFFFF", + foreground="#101B28", + bordercolor="#D8E0E4", + padding=(16, 9), + font=("Segoe UI", 10), + ) + style.map("WhatIf.TButton", background=[("active", "#EDF1F3")]) + style.configure( + "WhatIf.Primary.TButton", + background="#007D78", + foreground="#FFFFFF", + bordercolor="#007D78", + padding=(18, 9), + font=("Segoe UI", 10, "bold"), + ) + style.map("WhatIf.Primary.TButton", background=[("active", "#006965")]) + heading = tk.Frame(dialog, bg="#F5F7F7") + heading.pack(fill="x", padx=24, pady=(10, 6)) tk.Label( - dialog, - text="Синтетические свойства и запасы. Реальное управление отключено.", - fg="#A55E00", - font=("Segoe UI", 12, "bold"), - ).pack(pady=(14, 6)) + heading, + text="Синтетический what-if", + bg="#F5F7F7", + fg="#101B28", + font=("Segoe UI", 18, "bold"), + ).pack(anchor="w") tk.Label( - dialog, text="Пустое свойство означает неизвестное, а не ноль. Доли с присадкой = 100%." - ).pack() - notebook = ttk.Notebook(dialog) - notebook.pack(fill="both", expand=True, padx=16, pady=10) + heading, + text="Реальное управление отключено. Доли компонентов и присадки должны составлять 100 %.", + bg="#F5F7F7", + fg="#617082", + font=("Segoe UI", 10), + ).pack(anchor="w", pady=(2, 0)) + warning = tk.Label( + dialog, + text="Пустое свойство означает неизвестное значение, а не ноль.", + bg="#FFF8E8", + fg="#8C5400", + anchor="w", + font=("Segoe UI", 10, "bold"), + padx=16, + pady=6, + highlightbackground="#E9B850", + highlightthickness=1, + ) + warning.pack(fill="x", padx=24, pady=(0, 8)) + notebook = ttk.Notebook(dialog, style="WhatIf.TNotebook") + notebook.pack(fill="both", expand=True, padx=24, pady=(0, 8)) variables = { key: tk.StringVar(dialog, value=value) for key, value in editor_values(current).items() } - components = ttk.Frame(notebook, padding=12) - notebook.add(components, text="Компоненты и текущая рецептура") + recipe = ttk.Frame(notebook, padding=20, style="WhatIf.TFrame") + notebook.add(recipe, text="Рецептура") + ttk.Label( + recipe, + text="Введите значение. Свободные доли подберутся после Enter или выхода из поля.", + style="WhatIf.TLabel", + ).grid(row=0, column=0, columnspan=3, sticky="w", pady=(0, 14)) + locks: dict[str, tk.BooleanVar] = {} + entries: dict[str, ttk.Entry] = {} + recipe_keys = [f"{item.id}.fraction" for item in preset.blend_components] + ["dose"] + for row, key in enumerate(recipe_keys, 1): + label = "Присадка, %" if key == "dose" else f"Компонент {key.split('.')[0]}, %" + ttk.Label(recipe, text=label, style="WhatIf.TLabel").grid( + row=row, column=0, sticky="w", pady=7 + ) + entry = ttk.Entry(recipe, textvariable=variables[key], width=18, style="WhatIf.TEntry") + entry.grid(row=row, column=1, padx=12, sticky="w") + entries[key] = entry + locks[key] = tk.BooleanVar(dialog, value=False) + ttk.Checkbutton(recipe, text="Зафиксировать", variable=locks[key]).grid( + row=row, column=2, sticky="w" + ) + ttk.Label(recipe, text="Масса партии, т", style="WhatIf.TLabel").grid( + row=4, column=0, sticky="w", pady=7 + ) + entries["total_mass_t"] = ttk.Entry( + recipe, textvariable=variables["total_mass_t"], width=18, style="WhatIf.TEntry" + ) + entries["total_mass_t"].grid(row=4, column=1, padx=12, sticky="w") + preview = tk.StringVar( + dialog, value="Исходный сценарий. Измените поле или нажмите «Подобрать смесь»." + ) + ttk.Label( + recipe, textvariable=preview, style="WhatIf.TLabel", wraplength=790, justify="left" + ).grid(row=5, column=0, columnspan=3, sticky="w", pady=(18, 6)) + variants = ttk.Combobox(recipe, state="readonly", width=38) + variants.grid(row=6, column=0, columnspan=2, sticky="w", pady=5) + variants.grid_remove() + components = ttk.Frame(notebook, padding=10, style="WhatIf.TFrame") + notebook.add(components, text="Свойства и запасы") for col, component in enumerate(preset.blend_components): - frame = ttk.LabelFrame(components, text=f"Компонент {component.id}", padding=10) + frame = ttk.LabelFrame( + components, + text=f"Компонент {component.id}", + padding=8, + style="WhatIf.TLabelframe", + ) frame.grid(row=0, column=col, sticky="nsew", padx=8) components.columnconfigure(col, weight=1) for row, (field, label) in enumerate(COMPONENT_FIELDS): - ttk.Label(frame, text=label).grid(row=row, column=0, sticky="w", pady=6) - ttk.Entry(frame, textvariable=variables[f"{component.id}.{field}"], width=14).grid( - row=row, column=1, padx=10 + if field == "fraction": + continue + ttk.Label(frame, text=label, style="WhatIf.TLabel").grid( + row=row, column=0, sticky="w", pady=2 + ) + entry = ttk.Entry( + frame, + textvariable=variables[f"{component.id}.{field}"], + width=14, + style="WhatIf.TEntry", ) - additive = ttk.Frame(notebook, padding=20) - notebook.add(additive, text="Партия и присадка") + entry.grid(row=row, column=1, padx=10) + entries[f"{component.id}.{field}"] = entry + additive = ttk.Frame(notebook, padding=20, style="WhatIf.TFrame") + notebook.add(additive, text="Параметры присадки") for row, (key, label) in enumerate(ADDITIVE_FIELDS): - ttk.Label(additive, text=label).grid(row=row, column=0, sticky="w", pady=7) - ttk.Entry(additive, textvariable=variables[key], width=20).grid(row=row, column=1, padx=16) + if key in {"total_mass_t", "dose"}: + continue + ttk.Label(additive, text=label, style="WhatIf.TLabel").grid( + row=row, column=0, sticky="w", pady=7 + ) + entry = ttk.Entry(additive, textvariable=variables[key], width=20, style="WhatIf.TEntry") + entry.grid(row=row, column=1, padx=16) + entries[key] = entry error = tk.StringVar(dialog) - tk.Label(dialog, textvariable=error, fg="#B3473C", wraplength=880, justify="left").pack( - fill="x", padx=20 + tk.Label( + dialog, + textvariable=error, + bg="#F5F7F7", + fg="#B3473C", + wraplength=880, + justify="left", + ).pack(fill="x", padx=24) + + style.configure("WhatIf.Changed.TEntry", fieldbackground="#E2F6EF", padding=3) + assistant_on = tk.BooleanVar(dialog, value=True) + before_assist: dict[str, str] | None = None + options: tuple[tuple[BlendOption, dict[str, str]], ...] = () + pending: str | None = None + updating = False + + def values_now() -> dict[str, str]: + return {key: var.get() for key, var in variables.items()} + + def show_option(position: int) -> None: + nonlocal updating + option, updated = options[position] + old = values_now() + updating = True + try: + for key in recipe_keys: + variables[key].set(updated[key]) + entries[key].configure( + style=("WhatIf.Changed.TEntry" if updated[key] != old[key] else "WhatIf.TEntry") + ) + finally: + updating = False + names = { + **{f"{item.id}.fraction": f"Компонент {item.id}" for item in preset.blend_components}, + "dose": "Присадка", + } + changes = [ + f"{names[key]}: {old[key]}% → {updated[key]}%" + for key in recipe_keys + if old[key] != updated[key] + ] + quality = next( + (item.metrics for item in option.evaluation.assessments if "sulfur" in item.metrics), + {}, + ) + checks = [] + for constraint in preset.constraints: + metric = quality.get(constraint.metric) + if metric is None: + continue + bound = ( + "upper" + if constraint.use_upper_estimate + else "lower" + if constraint.use_lower_estimate + else "value" + ) + value = getattr(metric, bound) + limit = constraint.upper if constraint.upper is not None else constraint.lower + if value is not None and limit is not None: + sign = "≤" if constraint.upper is not None else "≥" + label = {"sulfur": "Сера", "t95": "T95", "cetane_number": "ЦЧ"}.get( + constraint.metric, constraint.metric + ) + checks.append(f"{label}: {value:.2f} {sign} {limit:g}") + preview.set( + f"{option.label}. " + + ("; ".join(changes) if changes else "Смесь уже подходит.") + + ("\nПроверка: " + "; ".join(checks) if checks else "") + ) + error.set("") + + def pick(_: tk.Event[tk.Misc] | None = None) -> None: + if variants.current() >= 0: + show_option(variants.current()) + + variants.bind("<>", pick) + + def suggest(changed: str | None = None) -> None: + nonlocal options, before_assist, pending + pending = None + if updating or not assistant_on.get(): + return + raw = values_now() + fixed = frozenset(key for key, locked in locks.items() if locked.get()) + if changed in locks: + fixed |= {changed} + try: + options = assist_editor_values(preset, raw, fixed) + except ValueError as exc: + options = () + variants.grid_remove() + error.set(str(exc)) + preview.set("Автоподбор не смог составить допустимую смесь.") + return + before_assist = raw + variants["values"] = tuple(item.label for item, _ in options) + variants.current(0) + variants.grid() + show_option(0) + + def schedule(changed: str) -> None: + nonlocal pending + if pending is not None: + dialog.after_cancel(pending) + pending = dialog.after(400, lambda: suggest(changed)) + + def on_return(_event: tk.Event[tk.Misc], *, name: str) -> None: + suggest(name) + + def on_focus_out(_event: tk.Event[tk.Misc], *, name: str) -> None: + schedule(name) + + for key, entry in entries.items(): + entry.bind("", partial(on_return, name=key)) + entry.bind("", partial(on_focus_out, name=key)) + + def undo_assist() -> None: + nonlocal updating, before_assist + if before_assist is None: + return + updating = True + try: + for key in recipe_keys: + if key not in locks or not locks[key].get(): + variables[key].set(before_assist[key]) + entries[key].configure(style="WhatIf.TEntry") + finally: + updating = False + before_assist = None + preview.set("Автоподбор отменён. Введённое значение сохранено.") + variants.grid_remove() + + ttk.Checkbutton(recipe, text="Автоподбор", variable=assistant_on).grid( + row=7, column=0, sticky="w", pady=(12, 0) + ) + ttk.Button(recipe, text="Подобрать смесь", command=suggest).grid( + row=7, column=1, sticky="w", pady=(12, 0) + ) + ttk.Button(recipe, text="Отменить автоподбор", command=undo_assist).grid( + row=7, column=2, sticky="w", pady=(12, 0) ) def apply() -> None: @@ -182,20 +509,42 @@ def apply() -> None: except ValueError as exc: error.set(str(exc)) return - calculate(scenario) + if calculate_with_context is None: + calculate(scenario) + else: + calculate_with_context( + scenario, + { + "original_values": before_assist, + "displayed_values": values_now(), + "locked_fields": [key for key, value in locks.items() if value.get()], + "assisted": before_assist is not None, + }, + ) dialog.destroy() def reset() -> None: + nonlocal before_assist for key, value in editor_values(preset).items(): variables[key].set(value) + for lock_var in locks.values(): + lock_var.set(False) + before_assist = None + variants.grid_remove() error.set( "Восстановлены исходные поля сценария. " "Нажмите «Пересчитать what-if» для нового результата." ) - buttons = ttk.Frame(dialog, padding=16) + buttons = ttk.Frame(dialog, padding=(24, 8), style="WhatIf.Footer.TFrame") buttons.pack(fill="x") - ttk.Button(buttons, text="Пересчитать what-if", command=apply).pack(side="left") - ttk.Button(buttons, text="Сбросить к сценарию", command=reset).pack(side="left", padx=10) - ttk.Button(buttons, text="Закрыть без изменений", command=dialog.destroy).pack(side="right") + ttk.Button( + buttons, text="Пересчитать what-if", command=apply, style="WhatIf.Primary.TButton" + ).pack(side="left") + ttk.Button(buttons, text="Сбросить к сценарию", command=reset, style="WhatIf.TButton").pack( + side="left", padx=10 + ) + ttk.Button( + buttons, text="Закрыть без изменений", command=dialog.destroy, style="WhatIf.TButton" + ).pack(side="right") return dialog From 059dd147d0a7197de85f7cc6d6109634003d34c8 Mon Sep 17 00:00:00 2001 From: MatveyVarfolomeev Date: Fri, 25 Sep 2026 23:10:36 +0300 Subject: [PATCH 2/2] ci: provide required aggregate lint-and-test check --- .github/workflows/ci.yml | 14 ++++++++++++-- 1 file changed, 12 insertions(+), 2 deletions(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 059c42d..06a815b 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -6,8 +6,9 @@ on: pull_request: branches: [ main, dev ] -jobs: - lint-and-test: +jobs: + test-matrix: + name: lint-and-test (${{ matrix.os }}, ${{ matrix.python }}) strategy: fail-fast: false matrix: @@ -53,3 +54,12 @@ jobs: run: | python -m source.main validate-stage0 python -m source.main acceptance --output reports/ci-acceptance + + lint-and-test: + name: lint-and-test + runs-on: ubuntu-latest + needs: test-matrix + if: always() + steps: + - name: Require both platform checks + run: test '${{ needs.test-matrix.result }}' = 'success'