diff --git a/.github/scripts/deploy-cloud-run-simulation-entry.sh b/.github/scripts/deploy-cloud-run-simulation-entry.sh index e33f41a59..5b208c1ac 100755 --- a/.github/scripts/deploy-cloud-run-simulation-entry.sh +++ b/.github/scripts/deploy-cloud-run-simulation-entry.sh @@ -154,6 +154,7 @@ jq -n ' OLD_GATEWAY_AUTH_CLIENT_ID: env.OLD_GATEWAY_AUTH_CLIENT_ID_VALUE, OBSERVABILITY_SERVICE_NAMESPACE: env.OBSERVABILITY_SERVICE_NAMESPACE, OBSERVABILITY_TRACE_PROJECT_ID: env.OBSERVABILITY_TRACE_PROJECT_ID, + CLOUD_RUN_REGION: env.REGION, OTEL_EXPORTER_OTLP_ENDPOINT: env.OTEL_EXPORTER_OTLP_ENDPOINT, OTEL_EXPORTER_OTLP_PROTOCOL: "grpc", OTEL_TRACES_EXPORTER: "otlp", diff --git a/docs/engineering/skills/observability.md b/docs/engineering/skills/observability.md index 853244ede..488f16822 100644 --- a/docs/engineering/skills/observability.md +++ b/docs/engineering/skills/observability.md @@ -86,6 +86,38 @@ When adding a run configuration or runtime stage: 3. import the plan in runtime code; 4. add focused tests for the new stage and identifier propagation. +The Stage 12 coordinator measures output preparation with these nested +operations: + +| Stage | Work measured | +|---|---| +| `stage12_coordinator_preparation` | The complete interval between acquiring report execution ownership and starting child-state persistence | +| `stage12_output_planning` | The complete coordinator interval that builds the shared output plan | +| `stage12_country_model_load` | Loading the country model and its Python modules | +| `stage12_output_configuration` | Creating planning simulations and configuring requested output groups | +| `stage12_output_variable_resolution` | Resolving entity variables and constructing the output schema | +| `stage12_child_input_planning` | Attaching the output plan to one baseline or reform child input | + +`stage12_child_input_planning` runs once for each simulation and includes a +`simulation_role` attribute. The other planning operations include the +`country` attribute. + +## Metric resource identity + +Every process that emits OpenTelemetry metrics must have a unique +`service.instance.id`. A Cloud Run revision name identifies deployed code and +a Modal task identifier identifies one task; neither value alone identifies a +Python worker process. Build the resource value with +`policyengine_observability.process_instance_id`, using the Cloud Run revision +or Modal task identifier as the platform prefix. The helper adds the process +ID and a random process value, and returns the same result for later calls in +that process. + +Cloud Run deployments must also provide `CLOUD_RUN_REGION`. Modal provides +`MODAL_REGION` at runtime. These values populate each workload's +`cloud.region`. The collector separately adds the central Google Monitoring +location to metrics before export. + ## Traces and polling HTTP instrumentation carries W3C trace context across synchronous service diff --git a/libs/policyengine-simulation-observability/pyproject.toml b/libs/policyengine-simulation-observability/pyproject.toml index c5d81c16b..fa8a65533 100644 --- a/libs/policyengine-simulation-observability/pyproject.toml +++ b/libs/policyengine-simulation-observability/pyproject.toml @@ -10,7 +10,7 @@ requires-python = ">=3.13" dependencies = [ "fastapi>=0.115.0", "pydantic>=2.0", - "policyengine-observability[fastapi,google,otlp-grpc]>=3.0.1,<4", + "policyengine-observability[fastapi,google,otlp-grpc]>=3.0.2,<4", ] [project.optional-dependencies] diff --git a/libs/policyengine-simulation-observability/src/policyengine_simulation_observability/observability.py b/libs/policyengine-simulation-observability/src/policyengine_simulation_observability/observability.py index ca50b765c..7195be901 100644 --- a/libs/policyengine-simulation-observability/src/policyengine_simulation_observability/observability.py +++ b/libs/policyengine-simulation-observability/src/policyengine_simulation_observability/observability.py @@ -18,6 +18,7 @@ StdoutLogDestination, configure, instrument_fastapi, + process_instance_id, ) Platform = Literal["google_cloud_run", "modal", "local", "other"] @@ -63,6 +64,14 @@ def build_runtime( ) -> ObservabilityRuntime: """Create one explicitly owned simulation observability runtime.""" + platform_instance_id = ( + os.getenv("MODAL_TASK_ID") if platform == "modal" else os.getenv("K_REVISION") + ) + region = ( + os.getenv("MODAL_REGION") + if platform == "modal" + else os.getenv("GOOGLE_CLOUD_REGION") or os.getenv("CLOUD_RUN_REGION") + ) trace_project = os.getenv("OBSERVABILITY_TRACE_PROJECT_ID", "").strip() logging_project = os.getenv("OBSERVABILITY_LOGGING_PROJECT_ID", "").strip() log_name = os.getenv("OBSERVABILITY_LOG_NAME", "").strip() @@ -93,8 +102,8 @@ def build_runtime( deployment=DeploymentIdentity( environment=environment, platform=platform, - region=os.getenv("GOOGLE_CLOUD_REGION") or os.getenv("CLOUD_RUN_REGION"), - instance_id=os.getenv("K_REVISION") or os.getenv("MODAL_TASK_ID"), + region=region, + instance_id=process_instance_id(service_name, platform_instance_id), ), logging=LoggingConfig( destinations=tuple(destinations), diff --git a/libs/policyengine-simulation-observability/src/policyengine_simulation_observability/stages.py b/libs/policyengine-simulation-observability/src/policyengine_simulation_observability/stages.py index 18982c592..d923c3379 100644 --- a/libs/policyengine-simulation-observability/src/policyengine_simulation_observability/stages.py +++ b/libs/policyengine-simulation-observability/src/policyengine_simulation_observability/stages.py @@ -80,6 +80,12 @@ class Stage(StrEnum): STAGE12_ENTRY_DISPATCH = "stage12_entry_dispatch" STAGE12_COORDINATOR_EXECUTION = "stage12_coordinator_execution" STAGE12_COORDINATOR_CLAIM = "stage12_coordinator_claim" + STAGE12_COORDINATOR_PREPARATION = "stage12_coordinator_preparation" + STAGE12_OUTPUT_PLANNING = "stage12_output_planning" + STAGE12_COUNTRY_MODEL_LOAD = "stage12_country_model_load" + STAGE12_OUTPUT_CONFIGURATION = "stage12_output_configuration" + STAGE12_OUTPUT_VARIABLE_RESOLUTION = "stage12_output_variable_resolution" + STAGE12_CHILD_INPUT_PLANNING = "stage12_child_input_planning" STAGE12_CHILD_STATE_CREATE = "stage12_child_state_create" STAGE12_CHILD_DISPATCH = "stage12_child_dispatch" STAGE12_CHILD_WAIT = "stage12_child_wait" @@ -183,6 +189,12 @@ def name(self, stage: Stage) -> str: Stage.STAGE12_ENTRY_DISPATCH, Stage.STAGE12_COORDINATOR_EXECUTION, Stage.STAGE12_COORDINATOR_CLAIM, + Stage.STAGE12_COORDINATOR_PREPARATION, + Stage.STAGE12_OUTPUT_PLANNING, + Stage.STAGE12_COUNTRY_MODEL_LOAD, + Stage.STAGE12_OUTPUT_CONFIGURATION, + Stage.STAGE12_OUTPUT_VARIABLE_RESOLUTION, + Stage.STAGE12_CHILD_INPUT_PLANNING, Stage.STAGE12_CHILD_STATE_CREATE, Stage.STAGE12_CHILD_DISPATCH, Stage.STAGE12_CHILD_WAIT, diff --git a/libs/policyengine-simulation-observability/tests/test_observability.py b/libs/policyengine-simulation-observability/tests/test_observability.py index 429103f8b..f0755284c 100644 --- a/libs/policyengine-simulation-observability/tests/test_observability.py +++ b/libs/policyengine-simulation-observability/tests/test_observability.py @@ -1,4 +1,5 @@ import inspect +import re from fastapi import FastAPI from fastapi.testclient import TestClient @@ -26,6 +27,9 @@ def test_modal_runtime_has_explicit_identity_and_remote_logging(monkeypatch): monkeypatch.setenv("OBSERVABILITY_TRACE_PROJECT_ID", "trace-project") monkeypatch.setenv("OBSERVABILITY_LOGGING_PROJECT_ID", "log-project") monkeypatch.setenv("OBSERVABILITY_LOG_NAME", "simulation-modal") + monkeypatch.setenv("MODAL_REGION", "us-east") + monkeypatch.setenv("MODAL_TASK_ID", "task-123") + monkeypatch.setenv("K_REVISION", "irrelevant-cloud-run-revision") runtime = build_runtime( service_name="policyengine-simulation-py5-2-0", service_role="simulation_worker", @@ -38,6 +42,11 @@ def test_modal_runtime_has_explicit_identity_and_remote_logging(monkeypatch): assert runtime.config.service.role == "simulation_worker" assert runtime.config.deployment.environment == "main" assert runtime.config.deployment.platform == "modal" + assert runtime.config.deployment.region == "us-east" + assert re.fullmatch( + r"task-123:\d+:[0-9a-f]{32}", + runtime.config.deployment.instance_id or "", + ) assert runtime.config.otel.sampling_ratio == 1.0 assert runtime.config.application_attribute_keys is None assert runtime.config.dispatch_attribute_keys == frozenset( @@ -56,6 +65,9 @@ def test_cloud_run_runtime_uses_stdout_without_direct_remote_logging(monkeypatch monkeypatch.setenv("OTEL_SDK_DISABLED", "true") monkeypatch.delenv("OBSERVABILITY_LOGGING_PROJECT_ID", raising=False) monkeypatch.delenv("OBSERVABILITY_LOG_NAME", raising=False) + monkeypatch.setenv("CLOUD_RUN_REGION", "us-central1") + monkeypatch.setenv("K_REVISION", "entry-00123-abc") + monkeypatch.setenv("MODAL_TASK_ID", "irrelevant-modal-task") runtime = init_process_observability( service_name="policyengine-simulation-entry-prod", service_role="simulation_entry", @@ -64,6 +76,11 @@ def test_cloud_run_runtime_uses_stdout_without_direct_remote_logging(monkeypatch ) try: assert runtime.config.logging.capture_standard_library is True + assert runtime.config.deployment.region == "us-central1" + assert re.fullmatch( + r"entry-00123-abc:\d+:[0-9a-f]{32}", + runtime.config.deployment.instance_id or "", + ) assert len(runtime.config.logging.destinations) == 1 assert isinstance(runtime.config.logging.destinations[0], StdoutLogDestination) finally: diff --git a/libs/policyengine-simulation-observability/tests/test_stages.py b/libs/policyengine-simulation-observability/tests/test_stages.py index d6203e503..9481c909b 100644 --- a/libs/policyengine-simulation-observability/tests/test_stages.py +++ b/libs/policyengine-simulation-observability/tests/test_stages.py @@ -29,3 +29,23 @@ def test_stage_plan_rejects_a_stage_from_another_configuration() -> None: with pytest.raises(ValueError, match="not registered"): annual.name(Stage.STAGE12_AGGREGATION) + + +@pytest.mark.parametrize( + "configuration", + [ + RunConfiguration.STAGE12_SHADOW_REPORT, + RunConfiguration.STAGE12_CANONICAL_REPORT, + ], +) +def test_stage12_report_plans_include_output_planning_breakdown( + configuration: RunConfiguration, +) -> None: + stages = RUN_STAGE_REGISTRY[configuration].stages + + assert Stage.STAGE12_COORDINATOR_PREPARATION in stages + assert Stage.STAGE12_OUTPUT_PLANNING in stages + assert Stage.STAGE12_COUNTRY_MODEL_LOAD in stages + assert Stage.STAGE12_OUTPUT_CONFIGURATION in stages + assert Stage.STAGE12_OUTPUT_VARIABLE_RESOLUTION in stages + assert Stage.STAGE12_CHILD_INPUT_PLANNING in stages diff --git a/libs/policyengine-simulation-observability/uv.lock b/libs/policyengine-simulation-observability/uv.lock index c07f95a5c..4afa35fa9 100644 --- a/libs/policyengine-simulation-observability/uv.lock +++ b/libs/policyengine-simulation-observability/uv.lock @@ -720,11 +720,11 @@ wheels = [ [[package]] name = "policyengine-observability" -version = "3.0.1" +version = "3.0.2" source = { registry = "https://pypi.org/simple" } -sdist = { url = "https://files.pythonhosted.org/packages/00/a6/f5d3e49523bf3d3231e9d2297fb03a4ae0704ae246d86153356dcb7eee7c/policyengine_observability-3.0.1.tar.gz", hash = "sha256:35e934b31843545a13570d6ba2e6a4b0cb002488f2c7bb1840351f6f90f30104", size = 124370, upload-time = "2026-09-28T16:48:19.916Z" } +sdist = { url = "https://files.pythonhosted.org/packages/74/2a/0c261c8ca693bb3a13d7fe837d3f89f7fdda245496d4b3e92eefed3c972e/policyengine_observability-3.0.2.tar.gz", hash = "sha256:554f907e43d8eeb274c1298fbe990b1626ac45f194c49995ddde091ddcdc2a64", size = 125718, upload-time = "2026-09-30T21:36:02.028Z" } wheels = [ - { url = "https://files.pythonhosted.org/packages/27/6c/879ecc01a618885e4bc741d8abb35e601d638f878e3145cda677b75c46b7/policyengine_observability-3.0.1-py3-none-any.whl", hash = "sha256:b7630d4d2d91c183b733d2e3470d3c0cbad67fd59a8c5e2f2bcc75872eea3c02", size = 40859, upload-time = "2026-09-28T16:48:18.661Z" }, + { url = "https://files.pythonhosted.org/packages/a5/9c/b48fb64c1bebf83bd5fdeebe914188806c3513cf2de0fc95ee1d95082fc4/policyengine_observability-3.0.2-py3-none-any.whl", hash = "sha256:d84c7564922f6a3294268e89549065924edde932c31adf8b7c714558e107cb31", size = 42029, upload-time = "2026-09-30T21:36:00.841Z" }, ] [package.optional-dependencies] @@ -768,7 +768,7 @@ requires-dist = [ { name = "black", marker = "extra == 'build'", specifier = ">=25.1.0" }, { name = "fastapi", specifier = ">=0.115.0" }, { name = "httpx2", marker = "extra == 'test'" }, - { name = "policyengine-observability", extras = ["fastapi", "google", "otlp-grpc"], specifier = ">=3.0.1,<4" }, + { name = "policyengine-observability", extras = ["fastapi", "google", "otlp-grpc"], specifier = ">=3.0.2,<4" }, { name = "pydantic", specifier = ">=2.0" }, { name = "pyright", marker = "extra == 'build'", specifier = ">=1.1.401" }, { name = "pytest", marker = "extra == 'test'", specifier = ">=8.3.4" }, diff --git a/projects/policyengine-simulation-entry/pyproject.toml b/projects/policyengine-simulation-entry/pyproject.toml index 25ea8d6f9..0c8bb1766 100644 --- a/projects/policyengine-simulation-entry/pyproject.toml +++ b/projects/policyengine-simulation-entry/pyproject.toml @@ -16,7 +16,7 @@ dependencies = [ "fastapi>=0.115.0,<1", "httpx>=0.28.0,<1", "modal>=1.4,<2", - "policyengine-observability[fastapi,google,httpx,otlp-grpc]>=3.0.1,<4", + "policyengine-observability[fastapi,google,httpx,otlp-grpc]>=3.0.2,<4", "policyengine-simulation-contract", "policyengine-simulation-observability", "policyengine-stage12-persistence", diff --git a/projects/policyengine-simulation-entry/tests/test_deployment_assets.py b/projects/policyengine-simulation-entry/tests/test_deployment_assets.py index f9f5f04e9..77d118cdf 100644 --- a/projects/policyengine-simulation-entry/tests/test_deployment_assets.py +++ b/projects/policyengine-simulation-entry/tests/test_deployment_assets.py @@ -485,6 +485,7 @@ def test_cloud_run_deployment_escapes_environment_and_pins_secret_versions(tmp_p assert runtime_environment["OBSERVABILITY_TRACE_PROJECT_ID"] == ( "policyengine-observability" ) + assert runtime_environment["CLOUD_RUN_REGION"] == "us-central1" assert runtime_environment["OTEL_EXPORTER_OTLP_ENDPOINT"] == ( "https://collector.example" ) diff --git a/projects/policyengine-simulation-entry/uv.lock b/projects/policyengine-simulation-entry/uv.lock index 428255c8f..6f3d0bf24 100644 --- a/projects/policyengine-simulation-entry/uv.lock +++ b/projects/policyengine-simulation-entry/uv.lock @@ -840,11 +840,11 @@ wheels = [ [[package]] name = "policyengine-observability" -version = "3.0.1" +version = "3.0.2" source = { registry = "https://pypi.org/simple" } -sdist = { url = "https://files.pythonhosted.org/packages/00/a6/f5d3e49523bf3d3231e9d2297fb03a4ae0704ae246d86153356dcb7eee7c/policyengine_observability-3.0.1.tar.gz", hash = "sha256:35e934b31843545a13570d6ba2e6a4b0cb002488f2c7bb1840351f6f90f30104", size = 124370, upload-time = "2026-09-28T16:48:19.916Z" } +sdist = { url = "https://files.pythonhosted.org/packages/74/2a/0c261c8ca693bb3a13d7fe837d3f89f7fdda245496d4b3e92eefed3c972e/policyengine_observability-3.0.2.tar.gz", hash = "sha256:554f907e43d8eeb274c1298fbe990b1626ac45f194c49995ddde091ddcdc2a64", size = 125718, upload-time = "2026-09-30T21:36:02.028Z" } wheels = [ - { url = "https://files.pythonhosted.org/packages/27/6c/879ecc01a618885e4bc741d8abb35e601d638f878e3145cda677b75c46b7/policyengine_observability-3.0.1-py3-none-any.whl", hash = "sha256:b7630d4d2d91c183b733d2e3470d3c0cbad67fd59a8c5e2f2bcc75872eea3c02", size = 40859, upload-time = "2026-09-28T16:48:18.661Z" }, + { url = "https://files.pythonhosted.org/packages/a5/9c/b48fb64c1bebf83bd5fdeebe914188806c3513cf2de0fc95ee1d95082fc4/policyengine_observability-3.0.2-py3-none-any.whl", hash = "sha256:d84c7564922f6a3294268e89549065924edde932c31adf8b7c714558e107cb31", size = 42029, upload-time = "2026-09-30T21:36:00.841Z" }, ] [package.optional-dependencies] @@ -924,7 +924,7 @@ requires-dist = [ { name = "httpx", specifier = ">=0.28.0,<1" }, { name = "modal", specifier = ">=1.4,<2" }, { name = "openapi-python-client", marker = "extra == 'build'", specifier = ">=0.21.6" }, - { name = "policyengine-observability", extras = ["fastapi", "google", "httpx", "otlp-grpc"], specifier = ">=3.0.1,<4" }, + { name = "policyengine-observability", extras = ["fastapi", "google", "httpx", "otlp-grpc"], specifier = ">=3.0.2,<4" }, { name = "policyengine-simulation-contract", editable = "../../libs/policyengine-simulation-contract" }, { name = "policyengine-simulation-observability", editable = "../../libs/policyengine-simulation-observability" }, { name = "policyengine-stage12-persistence", editable = "../../libs/policyengine-stage12-persistence" }, @@ -952,7 +952,7 @@ requires-dist = [ { name = "black", marker = "extra == 'build'", specifier = ">=25.1.0" }, { name = "fastapi", specifier = ">=0.115.0" }, { name = "httpx2", marker = "extra == 'test'" }, - { name = "policyengine-observability", extras = ["fastapi", "google", "otlp-grpc"], specifier = ">=3.0.1,<4" }, + { name = "policyengine-observability", extras = ["fastapi", "google", "otlp-grpc"], specifier = ">=3.0.2,<4" }, { name = "pydantic", specifier = ">=2.0" }, { name = "pyright", marker = "extra == 'build'", specifier = ">=1.1.401" }, { name = "pytest", marker = "extra == 'test'", specifier = ">=8.3.4" }, diff --git a/projects/policyengine-simulation-executor/pyproject.toml b/projects/policyengine-simulation-executor/pyproject.toml index 7265a279d..d67ee1e1d 100644 --- a/projects/policyengine-simulation-executor/pyproject.toml +++ b/projects/policyengine-simulation-executor/pyproject.toml @@ -24,7 +24,7 @@ dependencies = [ "policyengine[models]==6.2.1", "tables>=3.10.2", "modal>=0.73.0", - "policyengine-observability[fastapi,google,otlp-grpc]>=3.0.1,<4", + "policyengine-observability[fastapi,google,otlp-grpc]>=3.0.2,<4", # Imported directly by the artifact store client; do not rely on the # transitive pin via policyengine. "google-cloud-storage>=2", @@ -47,7 +47,7 @@ modal-simulation-image = [ "policyengine[models]==6.2.1", "fastapi>=0.115.0", "tables>=3.10.2", - "policyengine-observability[fastapi,google,otlp-grpc]>=3.0.1,<4", + "policyengine-observability[fastapi,google,otlp-grpc]>=3.0.2,<4", # The artifact fetch layer and store client read GCS. The lib is also a # transitive dependency of policyengine, but the image must not depend # on an upstream package keeping it. diff --git a/projects/policyengine-simulation-executor/src/policyengine_simulation_executor/stage12_runtime/coordination.py b/projects/policyengine-simulation-executor/src/policyengine_simulation_executor/stage12_runtime/coordination.py index e6d285916..ae0f867e5 100644 --- a/projects/policyengine-simulation-executor/src/policyengine_simulation_executor/stage12_runtime/coordination.py +++ b/projects/policyengine-simulation-executor/src/policyengine_simulation_executor/stage12_runtime/coordination.py @@ -2,7 +2,7 @@ from __future__ import annotations -from collections.abc import Callable +from collections.abc import Callable, Mapping from contextlib import nullcontext from datetime import UTC, datetime from hashlib import sha256 @@ -198,9 +198,8 @@ def coordinate_report( artifacts: Stage12ArtifactStore | None = None, invoker: ChildInvoker | None = None, aggregator: Callable[..., dict[str, Any]] = build_aggregate_report, - output_plan_resolver: Callable[ - [ReportExecutionInput], Stage12OutputPlan - ] = resolve_report_output_plan, + output_plan_resolver: Callable[[ReportExecutionInput], Stage12OutputPlan] + | None = None, runtime: ObservabilityRuntime | None = None, ) -> dict[str, Any]: report = ReportExecutionInput.model_validate(payload) @@ -242,26 +241,53 @@ def coordinate_report( "evaluation_id": str(parent.evaluation_id), "status": parent.status.value, } - context = context.model_copy( - update={ - "artifact_prefix": ( - f"stage-12-runs/{parent.environment}/" - f"{parent.created_at:%Y}/{parent.created_at:%m}/" - f"{parent.evaluation_id}" - ), - "created_at": parent.created_at, - "retention_expires_at": parent.retention_expires_at, - } - ) - function_name = context.simulation_callable descriptors: dict[str, SimulationArtifactDescriptor] = {} calls: dict[str, ChildCall] = {} try: - output_plan = output_plan_resolver(report) - simulations = ( - plan_simulation_input(report.baseline, output_plan), - plan_simulation_input(report.reform, output_plan), - ) + + def stage_scope( + stage: Stage, + attributes: Mapping[str, object] | None = None, + ): + if runtime is None: + return nullcontext() + return runtime.operation( + stage_plan.name(stage), + attributes=attributes, + ) + + with stage_scope(Stage.STAGE12_COORDINATOR_PREPARATION): + context = context.model_copy( + update={ + "artifact_prefix": ( + f"stage-12-runs/{parent.environment}/" + f"{parent.created_at:%Y}/{parent.created_at:%m}/" + f"{parent.evaluation_id}" + ), + "created_at": parent.created_at, + "retention_expires_at": parent.retention_expires_at, + } + ) + function_name = context.simulation_callable + with stage_scope( + Stage.STAGE12_OUTPUT_PLANNING, + {"country": report.baseline.geography.country}, + ): + if output_plan_resolver is None: + output_plan = resolve_report_output_plan( + report, + stage_scope=stage_scope, + ) + else: + output_plan = output_plan_resolver(report) + simulations = [] + for simulation in (report.baseline, report.reform): + with stage_scope( + Stage.STAGE12_CHILD_INPUT_PLANNING, + {"simulation_role": simulation.role.value}, + ): + simulations.append(plan_simulation_input(simulation, output_plan)) + planned_simulations = tuple(simulations) child_state_span = ( runtime.span(stage_plan.name(Stage.STAGE12_CHILD_STATE_CREATE)) if runtime is not None @@ -276,9 +302,9 @@ def coordinate_report( function_name=function_name, ) ).record - for simulation in simulations + for simulation in planned_simulations } - for simulation in simulations: + for simulation in planned_simulations: child = children[simulation.role] if child.status is ComparisonRunLifecycleStatus.SUCCEEDED: descriptor = descriptor_from_record(child, simulation) diff --git a/projects/policyengine-simulation-executor/src/policyengine_simulation_executor/stage12_runtime/output_planning.py b/projects/policyengine-simulation-executor/src/policyengine_simulation_executor/stage12_runtime/output_planning.py index c6dea0d2b..4d899d6c6 100644 --- a/projects/policyengine-simulation-executor/src/policyengine_simulation_executor/stage12_runtime/output_planning.py +++ b/projects/policyengine-simulation-executor/src/policyengine_simulation_executor/stage12_runtime/output_planning.py @@ -2,7 +2,8 @@ from __future__ import annotations -from collections.abc import Mapping +from collections.abc import Callable, Mapping +from contextlib import AbstractContextManager, nullcontext from typing import Protocol, cast import pandas as pd @@ -22,6 +23,7 @@ SimulationExecutionInput, Stage12OutputPlan, ) +from policyengine_simulation_observability.stages import Stage UK_GEOGRAPHIC_DATASET_VARIABLES: dict[str, tuple[str, ...]] = { "household": ("constituency_code_oa", "la_code_oa"), @@ -37,6 +39,19 @@ def resolve_entity_variables( ) -> dict[str, list[str]]: ... +StageScope = Callable[ + [Stage, Mapping[str, object] | None], + AbstractContextManager[object], +] + + +def _unmonitored_stage_scope( + _stage: Stage, + _attributes: Mapping[str, object] | None = None, +) -> AbstractContextManager[object]: + return nullcontext() + + def _country_model(country: CountryId) -> OutputVariableModel: if country == "us": from policyengine.tax_benefit_models import us @@ -119,57 +134,65 @@ def _configure_country_outputs( return labor_supply_active -def resolve_report_output_plan(report: ReportExecutionInput) -> Stage12OutputPlan: +def resolve_report_output_plan( + report: ReportExecutionInput, + *, + stage_scope: StageScope = _unmonitored_stage_scope, +) -> Stage12OutputPlan: """Resolve one country-owned output schema for both report simulations.""" country = report.baseline.geography.country - model = _country_model(country) - baseline = _planning_simulation(report.baseline, model=model) - reform = _planning_simulation(report.reform, model=model) - include_cliff_impacts = _include_cliff_impacts(report) - labor_supply_active = _configure_country_outputs( - country=country, - baseline=baseline, - reform=reform, - include_cliff_impacts=include_cliff_impacts, - ) - requirements = ReportOutputRequirements( - aggregates=report.requested_aggregates, - include_cliff_impacts=include_cliff_impacts, - labor_supply_response_active=labor_supply_active, - ) - baseline_variables = model.resolve_entity_variables(baseline) - reform_variables = model.resolve_entity_variables(reform) - dataset_variables = _required_dataset_variables( - country=country, - aggregates=requirements.aggregates, - ) - entities = [] - for entity in sorted( - set(baseline_variables) | set(reform_variables) | set(dataset_variables) - ): - entity_dataset_variables = dataset_variables.get(entity, ()) - materialized = tuple( - sorted( - set(baseline_variables.get(entity, ())) - | set(reform_variables.get(entity, ())) - | set(entity_dataset_variables) - ) + stage_attributes = {"country": country} + with stage_scope(Stage.STAGE12_COUNTRY_MODEL_LOAD, stage_attributes): + model = _country_model(country) + with stage_scope(Stage.STAGE12_OUTPUT_CONFIGURATION, stage_attributes): + baseline = _planning_simulation(report.baseline, model=model) + reform = _planning_simulation(report.reform, model=model) + include_cliff_impacts = _include_cliff_impacts(report) + labor_supply_active = _configure_country_outputs( + country=country, + baseline=baseline, + reform=reform, + include_cliff_impacts=include_cliff_impacts, ) - additional = tuple( - sorted( - set((baseline.extra_variables or {}).get(entity, ())) - | set((reform.extra_variables or {}).get(entity, ())) - ) + requirements = ReportOutputRequirements( + aggregates=report.requested_aggregates, + include_cliff_impacts=include_cliff_impacts, + labor_supply_response_active=labor_supply_active, ) - entities.append( - EntityOutputPlan( - entity=entity, - materialized_variables=materialized, - additional_variables=additional, - dataset_variables=entity_dataset_variables, - ) + with stage_scope(Stage.STAGE12_OUTPUT_VARIABLE_RESOLUTION, stage_attributes): + baseline_variables = model.resolve_entity_variables(baseline) + reform_variables = model.resolve_entity_variables(reform) + dataset_variables = _required_dataset_variables( + country=country, + aggregates=requirements.aggregates, ) + entities = [] + for entity in sorted( + set(baseline_variables) | set(reform_variables) | set(dataset_variables) + ): + entity_dataset_variables = dataset_variables.get(entity, ()) + materialized = tuple( + sorted( + set(baseline_variables.get(entity, ())) + | set(reform_variables.get(entity, ())) + | set(entity_dataset_variables) + ) + ) + additional = tuple( + sorted( + set((baseline.extra_variables or {}).get(entity, ())) + | set((reform.extra_variables or {}).get(entity, ())) + ) + ) + entities.append( + EntityOutputPlan( + entity=entity, + materialized_variables=materialized, + additional_variables=additional, + dataset_variables=entity_dataset_variables, + ) + ) return Stage12OutputPlan( country=country, requirements=requirements, diff --git a/projects/policyengine-simulation-executor/tests/test_stage12_output_planning.py b/projects/policyengine-simulation-executor/tests/test_stage12_output_planning.py index 02a857ea9..58ca4eef4 100644 --- a/projects/policyengine-simulation-executor/tests/test_stage12_output_planning.py +++ b/projects/policyengine-simulation-executor/tests/test_stage12_output_planning.py @@ -2,6 +2,7 @@ from __future__ import annotations +from contextlib import contextmanager from uuid import UUID import pandas as pd @@ -22,6 +23,7 @@ Stage12OutputPlan, ) from pydantic import JsonValue +from policyengine_simulation_observability.stages import Stage from policyengine_simulation_executor.stage12_runtime.output_planning import ( apply_output_plan, @@ -196,6 +198,23 @@ def reject_ensure(_simulation): assert first == second +def test_output_planning_records_each_expensive_substage() -> None: + observed: list[tuple[Stage, dict[str, object]]] = [] + + @contextmanager + def record_stage(stage: Stage, attributes): + observed.append((stage, dict(attributes or {}))) + yield + + resolve_report_output_plan(_report(), stage_scope=record_stage) + + assert observed == [ + (Stage.STAGE12_COUNTRY_MODEL_LOAD, {"country": "us"}), + (Stage.STAGE12_OUTPUT_CONFIGURATION, {"country": "us"}), + (Stage.STAGE12_OUTPUT_VARIABLE_RESOLUTION, {"country": "us"}), + ] + + def test_incomplete_aggregate_profile_is_rejected_before_dispatch() -> None: with pytest.raises(ValueError, match="complete aggregate profile"): resolve_report_output_plan(_report(aggregates=(ReportAggregate.BUDGET,))) diff --git a/projects/policyengine-simulation-executor/tests/test_stage12_runtime.py b/projects/policyengine-simulation-executor/tests/test_stage12_runtime.py index 3963458b8..6a7992118 100644 --- a/projects/policyengine-simulation-executor/tests/test_stage12_runtime.py +++ b/projects/policyengine-simulation-executor/tests/test_stage12_runtime.py @@ -2,9 +2,11 @@ from __future__ import annotations -from contextlib import nullcontext +import io +import json import time from concurrent.futures import ThreadPoolExecutor +from contextlib import contextmanager, nullcontext from datetime import UTC, datetime, timedelta from hashlib import sha256 from threading import Lock @@ -12,6 +14,15 @@ import pandas as pd import pytest +from policyengine_observability import ( + DeploymentIdentity, + LoggingConfig, + ObservabilityConfig, + OTelConfig, + ServiceIdentity, + StdoutLogDestination, + configure, +) from policyengine_simulation_contract.stage12_execution import ( ArtifactMediaType, ArtifactReference, @@ -44,7 +55,6 @@ UKLocalAuthorityBoundaryVersion, UKLocalAuthorityMetadata, ) - from policyengine_simulation_executor.stage12_artifacts import ( canonical_json_bytes, serialize_simulation_frames, @@ -53,10 +63,12 @@ SimulationCalculation, _build_spm_result, build_aggregate_report, - coordinate_report as coordinate_report_impl, run_single_simulation, simulation_input_sha256, ) +from policyengine_simulation_executor.stage12_runtime import ( + coordinate_report as coordinate_report_impl, +) from policyengine_simulation_executor.stage12_runtime.aggregation import ( validate_uk_local_authority_metadata, ) @@ -151,6 +163,22 @@ def coordinate_report(*args, **kwargs): return coordinate_report_impl(*args, **kwargs) +class RecordingRuntime: + def __init__(self) -> None: + self.operations: list[tuple[str, dict[str, object]]] = [] + + @contextmanager + def operation(self, name, *, attributes=None, **_kwargs): + self.operations.append((name, dict(attributes or {}))) + yield + + def span(self, *_args, **_kwargs): + return nullcontext() + + def capture_context(self): + return {} + + def _context() -> Stage12InvocationContext: return Stage12InvocationContext( request_id="request-1", @@ -986,6 +1014,108 @@ def test_duplicate_coordinator_submission_does_not_start_duplicate_children() -> assert duplicate_invoker.events == [] +def test_coordinator_measures_output_and_child_input_planning() -> None: + store = FakeStore() + artifacts = FakeArtifacts() + runtime = RecordingRuntime() + + coordinate_report( + _report().model_dump(mode="json"), + _context().model_dump(mode="json"), + _parent().model_dump(mode="json"), + application_name=_context().modal_application, + coordinator_invocation_id="coordinator-1", + store=store, + artifacts=artifacts, + invoker=ConcurrentInvoker(artifacts), + aggregator=lambda **_: {"result": "complete"}, + runtime=runtime, + ) + + planning_operations = [ + operation + for operation in runtime.operations + if operation[0] + in { + "stage12_coordinator_preparation", + "stage12_output_planning", + "stage12_child_input_planning", + } + ] + assert planning_operations == [ + ("stage12_coordinator_preparation", {}), + ("stage12_output_planning", {"country": "us"}), + ("stage12_child_input_planning", {"simulation_role": "baseline"}), + ("stage12_child_input_planning", {"simulation_role": "reform"}), + ] + + +def test_coordinator_planning_operations_preserve_observability_id() -> None: + output = io.StringIO() + runtime = configure( + ObservabilityConfig( + service=ServiceIdentity( + name="stage12-test", + namespace="policyengine.api-v1", + version="test", + role="coordinator", + ), + deployment=DeploymentIdentity( + environment="test", + platform="local", + region="us-central1", + instance_id="test-process", + ), + logging=LoggingConfig( + destinations=(StdoutLogDestination(),), + ), + otel=OTelConfig(enabled=False), + dispatch_attribute_keys=frozenset({"observability_id"}), + ) + ) + runtime._delivery._stdout = output + observability_id = "00000000-0000-4000-8000-000000000004" + artifacts = FakeArtifacts() + + try: + with runtime.operation( + "stage12_report", + remote_context={ + "captured_at": NOW.isoformat(), + "observability_id": observability_id, + }, + ): + coordinate_report( + _report().model_dump(mode="json"), + _context().model_dump(mode="json"), + _parent().model_dump(mode="json"), + application_name=_context().modal_application, + coordinator_invocation_id="coordinator-1", + store=FakeStore(), + artifacts=artifacts, + invoker=ConcurrentInvoker(artifacts), + aggregator=lambda **_: {"result": "complete"}, + runtime=runtime, + ) + finally: + runtime.shutdown() + + records = [json.loads(line) for line in output.getvalue().splitlines()] + planning_names = { + "stage12_coordinator_preparation", + "stage12_output_planning", + "stage12_child_input_planning", + } + planning_records = [ + record for record in records if record.get("operation.name") in planning_names + ] + assert len(planning_records) == 4 + assert all( + record["attributes"]["observability_id"] == observability_id + for record in planning_records + ) + + def test_coordinator_compares_automatic_run_after_successful_aggregation() -> None: store = FakeStore() parent = _automatic_parent() diff --git a/projects/policyengine-simulation-executor/uv.lock b/projects/policyengine-simulation-executor/uv.lock index ca5d39198..f944b36a9 100644 --- a/projects/policyengine-simulation-executor/uv.lock +++ b/projects/policyengine-simulation-executor/uv.lock @@ -1786,11 +1786,11 @@ provides-extras = ["test", "build"] [[package]] name = "policyengine-observability" -version = "3.0.1" +version = "3.0.2" source = { registry = "https://pypi.org/simple" } -sdist = { url = "https://files.pythonhosted.org/packages/00/a6/f5d3e49523bf3d3231e9d2297fb03a4ae0704ae246d86153356dcb7eee7c/policyengine_observability-3.0.1.tar.gz", hash = "sha256:35e934b31843545a13570d6ba2e6a4b0cb002488f2c7bb1840351f6f90f30104", size = 124370, upload-time = "2026-09-28T16:48:19.916Z" } +sdist = { url = "https://files.pythonhosted.org/packages/74/2a/0c261c8ca693bb3a13d7fe837d3f89f7fdda245496d4b3e92eefed3c972e/policyengine_observability-3.0.2.tar.gz", hash = "sha256:554f907e43d8eeb274c1298fbe990b1626ac45f194c49995ddde091ddcdc2a64", size = 125718, upload-time = "2026-09-30T21:36:02.028Z" } wheels = [ - { url = "https://files.pythonhosted.org/packages/27/6c/879ecc01a618885e4bc741d8abb35e601d638f878e3145cda677b75c46b7/policyengine_observability-3.0.1-py3-none-any.whl", hash = "sha256:b7630d4d2d91c183b733d2e3470d3c0cbad67fd59a8c5e2f2bcc75872eea3c02", size = 40859, upload-time = "2026-09-28T16:48:18.661Z" }, + { url = "https://files.pythonhosted.org/packages/a5/9c/b48fb64c1bebf83bd5fdeebe914188806c3513cf2de0fc95ee1d95082fc4/policyengine_observability-3.0.2-py3-none-any.whl", hash = "sha256:d84c7564922f6a3294268e89549065924edde932c31adf8b7c714558e107cb31", size = 42029, upload-time = "2026-09-30T21:36:00.841Z" }, ] [package.optional-dependencies] @@ -1894,7 +1894,7 @@ requires-dist = [ { name = "opentelemetry-instrumentation-sqlalchemy", specifier = ">=0.65b0,<0.66" }, { name = "policyengine", extras = ["models"], specifier = "==6.2.1" }, { name = "policyengine-fastapi", editable = "../../libs/policyengine-fastapi" }, - { name = "policyengine-observability", extras = ["fastapi", "google", "otlp-grpc"], specifier = ">=3.0.1,<4" }, + { name = "policyengine-observability", extras = ["fastapi", "google", "otlp-grpc"], specifier = ">=3.0.2,<4" }, { name = "policyengine-simulation-contract", extras = ["modal"], editable = "../../libs/policyengine-simulation-contract" }, { name = "policyengine-simulation-observability", editable = "../../libs/policyengine-simulation-observability" }, { name = "policyengine-stage12-persistence", editable = "../../libs/policyengine-stage12-persistence" }, @@ -1914,7 +1914,7 @@ modal-simulation-image = [ { name = "fastapi", specifier = ">=0.115.0" }, { name = "google-cloud-storage", specifier = ">=2" }, { name = "policyengine", extras = ["models"], specifier = "==6.2.1" }, - { name = "policyengine-observability", extras = ["fastapi", "google", "otlp-grpc"], specifier = ">=3.0.1,<4" }, + { name = "policyengine-observability", extras = ["fastapi", "google", "otlp-grpc"], specifier = ">=3.0.2,<4" }, { name = "psycopg", extras = ["binary"], specifier = ">=3.2,<4" }, { name = "pyarrow", specifier = ">=20,<24" }, { name = "sqlalchemy", specifier = ">=2,<3" }, @@ -1944,7 +1944,7 @@ requires-dist = [ { name = "modal", specifier = ">=1.4,<2" }, { name = "openapi-python-client", marker = "extra == 'build'", specifier = ">=0.21.6" }, { name = "packaging", specifier = ">=24" }, - { name = "policyengine-observability", extras = ["fastapi", "google", "otlp-grpc"], specifier = ">=3.0.1,<4" }, + { name = "policyengine-observability", extras = ["fastapi", "google", "otlp-grpc"], specifier = ">=3.0.2,<4" }, { name = "pydantic", specifier = ">=2.0" }, { name = "pyjwt", specifier = ">=2.10.1,<3.0.0" }, { name = "pyright", marker = "extra == 'build'", specifier = ">=1.1.401" }, @@ -1976,7 +1976,7 @@ requires-dist = [ { name = "black", marker = "extra == 'build'", specifier = ">=25.1.0" }, { name = "fastapi", specifier = ">=0.115.0" }, { name = "httpx2", marker = "extra == 'test'" }, - { name = "policyengine-observability", extras = ["fastapi", "google", "otlp-grpc"], specifier = ">=3.0.1,<4" }, + { name = "policyengine-observability", extras = ["fastapi", "google", "otlp-grpc"], specifier = ">=3.0.2,<4" }, { name = "pydantic", specifier = ">=2.0" }, { name = "pyright", marker = "extra == 'build'", specifier = ">=1.1.401" }, { name = "pytest", marker = "extra == 'test'", specifier = ">=8.3.4" }, diff --git a/projects/policyengine-simulation-gateway/pyproject.toml b/projects/policyengine-simulation-gateway/pyproject.toml index 0c1ca6c6d..615322049 100644 --- a/projects/policyengine-simulation-gateway/pyproject.toml +++ b/projects/policyengine-simulation-gateway/pyproject.toml @@ -25,7 +25,7 @@ dependencies = [ "pyjwt>=2.10.1,<3.0.0", # JWTDecoder (policyengine-fastapi lib) needs cryptography at runtime. "cryptography>=41.0.0", - "policyengine-observability[fastapi,google,otlp-grpc]>=3.0.1,<4", + "policyengine-observability[fastapi,google,otlp-grpc]>=3.0.2,<4", "modal>=1.4,<2", ] diff --git a/projects/policyengine-simulation-gateway/uv.lock b/projects/policyengine-simulation-gateway/uv.lock index cb99ed062..23e623e4e 100644 --- a/projects/policyengine-simulation-gateway/uv.lock +++ b/projects/policyengine-simulation-gateway/uv.lock @@ -1261,11 +1261,11 @@ provides-extras = ["test", "build"] [[package]] name = "policyengine-observability" -version = "3.0.1" +version = "3.0.2" source = { registry = "https://pypi.org/simple" } -sdist = { url = "https://files.pythonhosted.org/packages/00/a6/f5d3e49523bf3d3231e9d2297fb03a4ae0704ae246d86153356dcb7eee7c/policyengine_observability-3.0.1.tar.gz", hash = "sha256:35e934b31843545a13570d6ba2e6a4b0cb002488f2c7bb1840351f6f90f30104", size = 124370, upload-time = "2026-09-28T16:48:19.916Z" } +sdist = { url = "https://files.pythonhosted.org/packages/74/2a/0c261c8ca693bb3a13d7fe837d3f89f7fdda245496d4b3e92eefed3c972e/policyengine_observability-3.0.2.tar.gz", hash = "sha256:554f907e43d8eeb274c1298fbe990b1626ac45f194c49995ddde091ddcdc2a64", size = 125718, upload-time = "2026-09-30T21:36:02.028Z" } wheels = [ - { url = "https://files.pythonhosted.org/packages/27/6c/879ecc01a618885e4bc741d8abb35e601d638f878e3145cda677b75c46b7/policyengine_observability-3.0.1-py3-none-any.whl", hash = "sha256:b7630d4d2d91c183b733d2e3470d3c0cbad67fd59a8c5e2f2bcc75872eea3c02", size = 40859, upload-time = "2026-09-28T16:48:18.661Z" }, + { url = "https://files.pythonhosted.org/packages/a5/9c/b48fb64c1bebf83bd5fdeebe914188806c3513cf2de0fc95ee1d95082fc4/policyengine_observability-3.0.2-py3-none-any.whl", hash = "sha256:d84c7564922f6a3294268e89549065924edde932c31adf8b7c714558e107cb31", size = 42029, upload-time = "2026-09-30T21:36:00.841Z" }, ] [package.optional-dependencies] @@ -1355,7 +1355,7 @@ requires-dist = [ { name = "modal", specifier = ">=1.4,<2" }, { name = "openapi-python-client", marker = "extra == 'build'", specifier = ">=0.21.6" }, { name = "packaging", specifier = ">=24" }, - { name = "policyengine-observability", extras = ["fastapi", "google", "otlp-grpc"], specifier = ">=3.0.1,<4" }, + { name = "policyengine-observability", extras = ["fastapi", "google", "otlp-grpc"], specifier = ">=3.0.2,<4" }, { name = "pydantic", specifier = ">=2.0" }, { name = "pyjwt", specifier = ">=2.10.1,<3.0.0" }, { name = "pyright", marker = "extra == 'build'", specifier = ">=1.1.401" }, @@ -1387,7 +1387,7 @@ requires-dist = [ { name = "black", marker = "extra == 'build'", specifier = ">=25.1.0" }, { name = "fastapi", specifier = ">=0.115.0" }, { name = "httpx2", marker = "extra == 'test'" }, - { name = "policyengine-observability", extras = ["fastapi", "google", "otlp-grpc"], specifier = ">=3.0.1,<4" }, + { name = "policyengine-observability", extras = ["fastapi", "google", "otlp-grpc"], specifier = ">=3.0.2,<4" }, { name = "pydantic", specifier = ">=2.0" }, { name = "pyright", marker = "extra == 'build'", specifier = ">=1.1.401" }, { name = "pytest", marker = "extra == 'test'", specifier = ">=8.3.4" },