Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .github/scripts/deploy-cloud-run-simulation-entry.sh
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
32 changes: 32 additions & 0 deletions docs/engineering/skills/observability.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion libs/policyengine-simulation-observability/pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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]
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
StdoutLogDestination,
configure,
instrument_fastapi,
process_instance_id,
)

Platform = Literal["google_cloud_run", "modal", "local", "other"]
Expand Down Expand Up @@ -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()
Expand Down Expand Up @@ -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),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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,
Expand Down
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import inspect
import re

from fastapi import FastAPI
from fastapi.testclient import TestClient
Expand Down Expand Up @@ -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",
Expand All @@ -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(
Expand All @@ -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",
Expand All @@ -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:
Expand Down
20 changes: 20 additions & 0 deletions libs/policyengine-simulation-observability/tests/test_stages.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
8 changes: 4 additions & 4 deletions libs/policyengine-simulation-observability/uv.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion projects/policyengine-simulation-entry/pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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"
)
Expand Down
10 changes: 5 additions & 5 deletions projects/policyengine-simulation-entry/uv.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

4 changes: 2 additions & 2 deletions projects/policyengine-simulation-executor/pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand All @@ -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.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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
Expand All @@ -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)
Expand Down
Loading
Loading