diff --git a/src/carbonfactor_parser/persistence/ingestion_run_history_mapping.py b/src/carbonfactor_parser/persistence/ingestion_run_history_mapping.py index 88d4d2c..798fbc8 100644 --- a/src/carbonfactor_parser/persistence/ingestion_run_history_mapping.py +++ b/src/carbonfactor_parser/persistence/ingestion_run_history_mapping.py @@ -18,7 +18,7 @@ ParserIngestionSourceResultRecord, ) if TYPE_CHECKING: - from carbonfactor_parser.pipeline.configured_cycle_runner import ConfiguredCycleResult + from carbonfactor_parser.pipeline.configured_cycle_models import ConfiguredCycleResult def build_ingestion_run_history_command_from_configured_cycle( diff --git a/src/carbonfactor_parser/pipeline/configured_cycle_history.py b/src/carbonfactor_parser/pipeline/configured_cycle_history.py new file mode 100644 index 0000000..6e670cd --- /dev/null +++ b/src/carbonfactor_parser/pipeline/configured_cycle_history.py @@ -0,0 +1,87 @@ +"""Configured ingestion cycle run-history persistence helpers.""" + +from __future__ import annotations + +from datetime import datetime +from typing import Callable + +from carbonfactor_parser.diagnostics.redaction import redact_sensitive_text +from carbonfactor_parser.persistence.ingestion_run_history import ( + ParserIngestionRunHistoryRepository, + ParserIngestionRunHistoryStatus, +) +from carbonfactor_parser.persistence.ingestion_run_history_mapping import ( + build_ingestion_run_history_command_from_configured_cycle, +) +from carbonfactor_parser.pipeline.configured_cycle_models import ConfiguredCycleResult + + +def persist_configured_cycle_history( + cycle: ConfiguredCycleResult, + *, + history_repository: ParserIngestionRunHistoryRepository, + started_at: datetime, + finished_at: datetime, + emit: Callable[[str], None] | None, +) -> ConfiguredCycleResult: + """Persist run-history for a configured cycle without affecting ingestion result.""" + + command = build_ingestion_run_history_command_from_configured_cycle( + cycle, + started_at=started_at, + finished_at=finished_at, + ) + try: + persist_result = history_repository.persist_ingestion_run_history(command) + except Exception as exc: # pragma: no cover - defensive boundary protection + safe_message = redact_sensitive_text(str(exc)) + if emit is not None: + emit( + "history_persistence " + f"status=failed run_id={cycle.run_id} " + "issue code=INGESTION_RUN_HISTORY_PERSISTENCE_EXCEPTION " + f"message={safe_message}" + ) + return ConfiguredCycleResult( + cycle_number=cycle.cycle_number, + run_id=cycle.run_id, + result=cycle.result, + history_persistence_status="failed", + history_persistence_issue_count=1, + ) + + issue_count = len(persist_result.issues) + if persist_result.status is ParserIngestionRunHistoryStatus.DECLARED: + if emit is not None: + emit(f"history_persistence status=declared run_id={cycle.run_id}") + else: + if emit is not None: + for issue in persist_result.issues or (): + safe_message = redact_sensitive_text(str(issue.message)) + emit( + "history_persistence " + f"status=failed run_id={cycle.run_id} " + f"issue code={issue.code} message={safe_message}" + ) + if not persist_result.issues: + emit( + "history_persistence " + f"status=failed run_id={cycle.run_id} " + "issue code=INGESTION_RUN_HISTORY_PERSISTENCE_FAILED " + "message=run history persistence failed" + ) + issue_count = 1 + return ConfiguredCycleResult( + cycle_number=cycle.cycle_number, + run_id=cycle.run_id, + result=cycle.result, + history_persistence_status=( + "declared" + if persist_result.status is ParserIngestionRunHistoryStatus.DECLARED + else "failed" + ), + history_persistence_issue_count=issue_count, + ) + + +__all__ = ("persist_configured_cycle_history",) diff --git a/src/carbonfactor_parser/pipeline/configured_cycle_models.py b/src/carbonfactor_parser/pipeline/configured_cycle_models.py new file mode 100644 index 0000000..4e74d5e --- /dev/null +++ b/src/carbonfactor_parser/pipeline/configured_cycle_models.py @@ -0,0 +1,36 @@ +"""Configured ingestion cycle result models.""" + +from __future__ import annotations + +from dataclasses import dataclass +from typing import TYPE_CHECKING + +from carbonfactor_parser.pipeline.production_e2e_year_orchestrator import ( + ProductionE2EYearOrchestratorResult, +) + +if TYPE_CHECKING: + from carbonfactor_parser.pipeline.configured_cycle_runner import ( + ConfiguredCycleRunnerStatus, + ) + + +@dataclass(frozen=True) +class ConfiguredCycleResult: + """One completed application cycle.""" + + cycle_number: int + run_id: str + result: ProductionE2EYearOrchestratorResult + history_persistence_status: str | None = None + history_persistence_issue_count: int = 0 + + +@dataclass(frozen=True) +class ConfiguredCycleRunnerResult: + """All cycles run by one application invocation.""" + + status: ConfiguredCycleRunnerStatus + cycles: tuple[ConfiguredCycleResult, ...] + schema_created_table_names: tuple[str, ...] + schema_missing_table_names: tuple[str, ...] diff --git a/src/carbonfactor_parser/pipeline/configured_cycle_runner.py b/src/carbonfactor_parser/pipeline/configured_cycle_runner.py index ad231f5..648732f 100644 --- a/src/carbonfactor_parser/pipeline/configured_cycle_runner.py +++ b/src/carbonfactor_parser/pipeline/configured_cycle_runner.py @@ -7,7 +7,6 @@ from __future__ import annotations -from dataclasses import dataclass from datetime import datetime, timezone from enum import Enum import time @@ -18,10 +17,6 @@ from carbonfactor_parser.persistence.ingestion_run_history import ( ParserIngestionRunHistoryRepository, - ParserIngestionRunHistoryStatus, -) -from carbonfactor_parser.persistence.ingestion_run_history_mapping import ( - build_ingestion_run_history_command_from_configured_cycle, ) from carbonfactor_parser.persistence.postgresql_ingestion_run_history_repository import ( PostgreSQLIngestionRunHistoryRepository, @@ -40,9 +35,18 @@ ConfiguredCycleValidationBoundary, build_configured_cycle_dependencies, ) +from carbonfactor_parser.pipeline.configured_cycle_history import ( + persist_configured_cycle_history, +) +from carbonfactor_parser.pipeline.configured_cycle_models import ( + ConfiguredCycleResult, + ConfiguredCycleRunnerResult, +) +from carbonfactor_parser.pipeline.configured_cycle_summary import ( + emit_configured_cycle_summary, +) from carbonfactor_parser.pipeline.production_e2e_year_orchestrator import ( ProductionE2EYearOrchestratorRequest, - ProductionE2EYearOrchestratorResult, ProductionE2EYearRunStatus, run_production_e2e_year_orchestrator, ) @@ -55,27 +59,6 @@ class ConfiguredCycleRunnerStatus(str, Enum): COMPLETED_WITH_FAILURES = "completed_with_failures" -@dataclass(frozen=True) -class ConfiguredCycleResult: - """One completed application cycle.""" - - cycle_number: int - run_id: str - result: ProductionE2EYearOrchestratorResult - history_persistence_status: str | None = None - history_persistence_issue_count: int = 0 - - -@dataclass(frozen=True) -class ConfiguredCycleRunnerResult: - """All cycles run by one application invocation.""" - - status: ConfiguredCycleRunnerStatus - cycles: tuple[ConfiguredCycleResult, ...] - schema_created_table_names: tuple[str, ...] - schema_missing_table_names: tuple[str, ...] - - def run_configured_cycle_runner( config: ConfiguredCycleRunnerConfig, *, @@ -121,7 +104,7 @@ def run_configured_cycle_runner( ) if emit is not None: emit_configured_cycle_summary(cycle, emit=emit) - cycle = _persist_configured_cycle_history( + cycle = persist_configured_cycle_history( cycle, history_repository=history_repository, started_at=started_at, @@ -152,118 +135,6 @@ def run_configured_cycle_runner( ) -def _persist_configured_cycle_history( - cycle: ConfiguredCycleResult, - *, - history_repository: ParserIngestionRunHistoryRepository, - started_at: datetime, - finished_at: datetime, - emit: Callable[[str], None] | None, -) -> ConfiguredCycleResult: - command = build_ingestion_run_history_command_from_configured_cycle( - cycle, - started_at=started_at, - finished_at=finished_at, - ) - try: - persist_result = history_repository.persist_ingestion_run_history(command) - except Exception as exc: # pragma: no cover - defensive boundary protection - safe_message = _redact_sensitive_text(str(exc)) - if emit is not None: - emit( - "history_persistence " - f"status=failed run_id={cycle.run_id} " - "issue code=INGESTION_RUN_HISTORY_PERSISTENCE_EXCEPTION " - f"message={safe_message}" - ) - return ConfiguredCycleResult( - cycle_number=cycle.cycle_number, - run_id=cycle.run_id, - result=cycle.result, - history_persistence_status="failed", - history_persistence_issue_count=1, - ) - - issue_count = len(persist_result.issues) - if persist_result.status is ParserIngestionRunHistoryStatus.DECLARED: - if emit is not None: - emit(f"history_persistence status=declared run_id={cycle.run_id}") - else: - if emit is not None: - for issue in persist_result.issues or (): - safe_message = _redact_sensitive_text(str(issue.message)) - emit( - "history_persistence " - f"status=failed run_id={cycle.run_id} " - f"issue code={issue.code} message={safe_message}" - ) - if not persist_result.issues: - emit( - "history_persistence " - f"status=failed run_id={cycle.run_id} " - "issue code=INGESTION_RUN_HISTORY_PERSISTENCE_FAILED " - "message=run history persistence failed" - ) - issue_count = 1 - return ConfiguredCycleResult( - cycle_number=cycle.cycle_number, - run_id=cycle.run_id, - result=cycle.result, - history_persistence_status=( - "declared" - if persist_result.status is ParserIngestionRunHistoryStatus.DECLARED - else "failed" - ), - history_persistence_issue_count=issue_count, - ) - - -def emit_configured_cycle_summary( - cycle: ConfiguredCycleResult, - *, - emit: Callable[[str], None] = print, -) -> None: - """Print user-readable summary output for one cycle.""" - - summary = cycle.result.summary - emit( - "cycle=" - f"{cycle.cycle_number} run_id={cycle.run_id} status={cycle.result.status.value}" - ) - emit( - "summary " - f"completed={summary.completed_family_count} " - f"no_available_source_year={summary.no_available_source_year_count} " - f"failed={summary.failed_family_count} " - f"parsed_rows={summary.parsed_row_count} " - f"inserted={summary.inserted_count} " - f"skipped_duplicates={summary.skipped_duplicate_count}" - ) - for family in cycle.result.family_results: - insert_summary = family.insert_summary - emit( - "source " - f"family={family.source_family} " - f"target_year={family.year_state.target_year} " - f"latest_year={family.year_state.latest_year} " - f"status={family.status.value} " - f"download_status={_download_status_value(family.download_result)} " - f"parse_status={_parse_status_value(family)} " - f"parsed_rows={family.parsed_row_count} " - f"master_inserted={getattr(insert_summary, 'master_inserted', 0)} " - f"master_skipped={getattr(insert_summary, 'master_skipped', 0)} " - f"detail_inserted={getattr(insert_summary, 'detail_inserted', 0)} " - f"detail_skipped={getattr(insert_summary, 'detail_skipped', 0)}" - ) - for failure in family.failures: - safe_message = _redact_sensitive_text(str(failure.message)) - emit( - "issue " - f"family={failure.source_family} stage={failure.stage} " - f"code={failure.code} message={safe_message}" - ) - - def _emit_startup_summary( config: ConfiguredCycleRunnerConfig, runtime: PostgreSQLRuntimeStartupResult, @@ -287,27 +158,6 @@ def _redact_sensitive_text(text: str) -> str: return redact_sensitive_text(text) -def _download_status_value(download_result: object | None) -> str: - if download_result is None: - return "not_run" - return str(getattr(getattr(download_result, "status", None), "value", "unknown")) - - -def _parse_status_value(family: object) -> str: - if getattr(family, "parsed_row_count", 0) > 0: - return "parsed" - failures = tuple(getattr(family, "failures", ())) - if any(getattr(failure, "stage", "") == "parser" for failure in failures): - return "failed" - download_result = getattr(family, "download_result", None) - if ( - download_result is None - or _download_status_value(download_result) != "downloaded" - ): - return "not_run" - return "no_rows" - - __all__ = ( "CONFIGURED_CYCLE_SOURCE_FAMILIES", "ConfiguredCycleResult", diff --git a/src/carbonfactor_parser/pipeline/configured_cycle_summary.py b/src/carbonfactor_parser/pipeline/configured_cycle_summary.py new file mode 100644 index 0000000..bb893e2 --- /dev/null +++ b/src/carbonfactor_parser/pipeline/configured_cycle_summary.py @@ -0,0 +1,75 @@ +"""Configured ingestion cycle summary output helpers.""" + +from __future__ import annotations + +from typing import Callable + +from carbonfactor_parser.diagnostics.redaction import redact_sensitive_text +from carbonfactor_parser.pipeline.configured_cycle_models import ConfiguredCycleResult + + +def emit_configured_cycle_summary( + cycle: ConfiguredCycleResult, + *, + emit: Callable[[str], None] = print, +) -> None: + """Print user-readable summary output for one cycle.""" + + summary = cycle.result.summary + emit( + "cycle=" + f"{cycle.cycle_number} run_id={cycle.run_id} status={cycle.result.status.value}" + ) + emit( + "summary " + f"completed={summary.completed_family_count} " + f"no_available_source_year={summary.no_available_source_year_count} " + f"failed={summary.failed_family_count} " + f"parsed_rows={summary.parsed_row_count} " + f"inserted={summary.inserted_count} " + f"skipped_duplicates={summary.skipped_duplicate_count}" + ) + for family in cycle.result.family_results: + insert_summary = family.insert_summary + emit( + "source " + f"family={family.source_family} " + f"target_year={family.year_state.target_year} " + f"latest_year={family.year_state.latest_year} " + f"status={family.status.value} " + f"download_status={_download_status_value(family.download_result)} " + f"parse_status={_parse_status_value(family)} " + f"parsed_rows={family.parsed_row_count} " + f"master_inserted={getattr(insert_summary, 'master_inserted', 0)} " + f"master_skipped={getattr(insert_summary, 'master_skipped', 0)} " + f"detail_inserted={getattr(insert_summary, 'detail_inserted', 0)} " + f"detail_skipped={getattr(insert_summary, 'detail_skipped', 0)}" + ) + for failure in family.failures: + safe_message = redact_sensitive_text(str(failure.message)) + emit( + "issue " + f"family={failure.source_family} stage={failure.stage} " + f"code={failure.code} message={safe_message}" + ) + + +def _download_status_value(download_result: object | None) -> str: + if download_result is None: + return "not_run" + return str(getattr(getattr(download_result, "status", None), "value", "unknown")) + + +def _parse_status_value(family: object) -> str: + if getattr(family, "parsed_row_count", 0) > 0: + return "parsed" + failures = tuple(getattr(family, "failures", ())) + if any(getattr(failure, "stage", "") == "parser" for failure in failures): + return "failed" + download_result = getattr(family, "download_result", None) + if download_result is None or _download_status_value(download_result) != "downloaded": + return "not_run" + return "no_rows" + + +__all__ = ("emit_configured_cycle_summary",) diff --git a/tests/test_configured_cycle_runner.py b/tests/test_configured_cycle_runner.py index 51fffa6..5effcd1 100644 --- a/tests/test_configured_cycle_runner.py +++ b/tests/test_configured_cycle_runner.py @@ -25,15 +25,23 @@ ConfiguredSourceYearArtifact as ExtractedConfiguredSourceYearArtifact, load_configured_cycle_runner_config as extracted_load_configured_cycle_runner_config, ) +from carbonfactor_parser.pipeline.configured_cycle_models import ( + ConfiguredCycleResult as ExtractedConfiguredCycleResult, +) +from carbonfactor_parser.pipeline.configured_cycle_summary import ( + emit_configured_cycle_summary as extracted_emit_configured_cycle_summary, +) from carbonfactor_parser.pipeline.configured_cycle_dependencies import ( ConfiguredCycleValidationBoundary as ExtractedConfiguredCycleValidationBoundary, build_configured_cycle_dependencies, ) from carbonfactor_parser.pipeline.configured_cycle_runner import ( + ConfiguredCycleResult, ConfiguredCycleRunnerConfig, ConfiguredCycleRunnerStatus, ConfiguredCycleValidationBoundary, ConfiguredSourceYearArtifact, + emit_configured_cycle_summary, load_configured_cycle_runner_config, run_configured_cycle_runner, ) @@ -63,6 +71,8 @@ def test_configured_cycle_runner_imports_remain_backward_compatible() -> None: assert ConfiguredSourceYearArtifact is ExtractedConfiguredSourceYearArtifact assert load_configured_cycle_runner_config is extracted_load_configured_cycle_runner_config assert ConfiguredCycleValidationBoundary is ExtractedConfiguredCycleValidationBoundary + assert ConfiguredCycleResult is ExtractedConfiguredCycleResult + assert emit_configured_cycle_summary is extracted_emit_configured_cycle_summary def test_build_configured_cycle_dependencies_wires_configured_runtime( @@ -699,6 +709,53 @@ def test_configured_cycle_runner_history_failure_does_not_fail_ingestion( assert "token=abc" not in captured.out +def test_configured_cycle_runner_history_exception_is_sanitized_and_non_fatal( + tmp_path: Path, + capsys: pytest.CaptureFixture[str], +) -> None: + connection = _FakeConnection() + repository = _ExceptionHistoryRepository() + config = ConfiguredCycleRunnerConfig( + postgresql_config_result=load_postgresql_runtime_config( + {"CARBONOPS_POSTGRESQL_DSN": "postgresql://user:pass@localhost/db"}, + ), + archive_root=tmp_path / "archive", + enabled_source_families=("ghg_protocol",), + initial_year=2024, + cycle_interval_seconds=0, + max_cycles=1, + source_years={"ghg_protocol": {2024: _artifact(2024, tmp_path / "ghg.csv", "v2024")}}, + ) + (tmp_path / "ghg.csv").write_text(_ghg_csv(2024), encoding="utf-8") + + result = run_configured_cycle_runner( + config, + startup=_startup(connection), + sleep=lambda _: None, + run_history_repository=repository, + ) + + captured = capsys.readouterr() + assert result.status is ConfiguredCycleRunnerStatus.COMPLETED + assert result.cycles[0].result.status.value == "completed" + assert result.cycles[0].history_persistence_status == "failed" + assert result.cycles[0].history_persistence_issue_count == 1 + assert "INGESTION_RUN_HISTORY_PERSISTENCE_EXCEPTION" in captured.out + assert "secret" not in captured.out + assert "token=abc" not in captured.out + + +class _ExceptionHistoryRepository: + @property + def provider_name(self) -> str: + return "fake" + + def persist_ingestion_run_history(self, command): + raise RuntimeError( + "database failed dsn=postgresql://user:secret@localhost/db token=abc" + ) + + class _FakeHistoryRepository: def __init__(self, *, status: str = "declared", issue_message: str = "") -> None: from carbonfactor_parser.persistence.ingestion_run_history import (