Skip to content
Merged
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
2 changes: 1 addition & 1 deletion .github/lambda-layer-publish.toml
Original file line number Diff line number Diff line change
@@ -1,2 +1,2 @@
[layer]
sdk-version = "2.0.0"
sdk-version = "2.1.0"
16 changes: 15 additions & 1 deletion .github/scripts/js_examples/run.py
Original file line number Diff line number Diff line change
Expand Up @@ -61,15 +61,24 @@
TESTING_SOURCE = REPO_ROOT / "packages/aws-durable-execution-sdk-python-testing/src"
JS_SDK_URL = "https://github.com/aws/aws-durable-execution-sdk-js.git"
EXAMPLES_REL = Path("packages/aws-durable-execution-sdk-js-examples")
# The examples depend on these workspaces by "*". The root "npm run build"
# The examples depend on these local workspaces. The root "npm run build"
# also builds the insight tools and the VS Code extension, which the harness
# does not need, so the script builds only these, in dependency order.
JS_BUILD_WORKSPACES = (
"packages/aws-durable-execution-sdk-js",
"packages/aws-durable-execution-sdk-js-testing",
"packages/aws-durable-execution-sdk-js-otel",
"packages/aws-durable-execution-sdk-js-extras",
"packages/aws-durable-execution-sdk-js-microvm-worker",
"packages/aws-durable-execution-sdk-js-examples",
)
# Older JS refs predate these packages; their examples do not need them.
OPTIONAL_JS_BUILD_WORKSPACES = frozenset(
{
"packages/aws-durable-execution-sdk-js-extras",
"packages/aws-durable-execution-sdk-js-microvm-worker",
}
)
# otel examples export spans to an OpenTelemetry collector. The harness does
# not run one, so these examples cannot pass here. They are not selected.
EXCLUDED_DIRS = ("/otel/",)
Expand Down Expand Up @@ -202,6 +211,11 @@ def build_js_sdk(js_dir: Path, *, force: bool) -> None:
[npm, "ci", "--no-audit", "--no-fund"], cwd=js_dir, env=env, check=True
)
for workspace in JS_BUILD_WORKSPACES:
if (
workspace in OPTIONAL_JS_BUILD_WORKSPACES
and not (js_dir / workspace / "package.json").is_file()
):
continue
log(f"npm run build -w {workspace}")
subprocess.run(
[npm, "run", "build", "-w", workspace], cwd=js_dir, env=env, check=True
Expand Down
52 changes: 52 additions & 0 deletions .github/scripts/tests/test_js_examples.py
Original file line number Diff line number Diff line change
Expand Up @@ -623,6 +623,58 @@ def test_source_fingerprint_changes_with_local_edits(tmp_path: Path) -> None:
assert run.source_fingerprint(tmp_path) != untracked


@pytest.mark.parametrize("has_microvm_packages", [False, True])
def test_build_produces_example_workspace_dependencies(
tmp_path: Path,
monkeypatch: pytest.MonkeyPatch,
has_microvm_packages: bool,
) -> None:
core = "packages/aws-durable-execution-sdk-js"
testing = "packages/aws-durable-execution-sdk-js-testing"
otel = "packages/aws-durable-execution-sdk-js-otel"
extras = "packages/aws-durable-execution-sdk-js-extras"
worker = "packages/aws-durable-execution-sdk-js-microvm-worker"
examples_package = "packages/aws-durable-execution-sdk-js-examples"
packages = [core, testing, otel, examples_package]
if has_microvm_packages:
packages.extend([extras, worker])
for package in packages:
directory = tmp_path / package
directory.mkdir(parents=True)
(directory / "package.json").write_text("{}")

built: list[str] = []

def execute(args: list[str], **_kwargs: Any) -> subprocess.CompletedProcess[str]:
if args[1] == "ci":
return subprocess.CompletedProcess(args, 0)
workspace = args[-1]
assert (tmp_path / workspace / "package.json").is_file()
if workspace == worker:
assert testing in built
if workspace == examples_package and has_microvm_packages:
# The generator imports extras/microvm and typechecks the worker.
# A successful examples build requires both artifacts to exist.
assert extras in built
assert worker in built
built.append(workspace)
return subprocess.CompletedProcess(args, 0)

stamp = tmp_path / "built-stamp"
monkeypatch.setattr(run, "_head", lambda _: "test-head")
monkeypatch.setattr(run, "state_file", lambda *_: stamp)
monkeypatch.setattr(run, "source_fingerprint", lambda _: "fingerprint")
monkeypatch.setattr(run.shutil, "which", lambda _: "/fake/npm")
monkeypatch.setattr(run.subprocess, "run", execute)

run.build_js_sdk(tmp_path, force=True)

assert built[-1] == examples_package
assert stamp.read_text() == "fingerprint\n"
if not has_microvm_packages:
assert built == [core, testing, otel, examples_package]


def test_proxy_creates_the_dump_directory(tmp_path: Path) -> None:
dump_dir = tmp_path / "new" / "dump"
proxy = invoke_proxy.ProxyServer(0, "http://127.0.0.1:1", dump_dir)
Expand Down
42 changes: 37 additions & 5 deletions .github/scripts/tests/test_opentelemetry_conformance_workflow.py
Original file line number Diff line number Diff line change
@@ -1,12 +1,12 @@
from pathlib import Path

import yaml


WORKFLOW_PATH = (
Path(__file__).parents[2] / "workflows" / "opentelemetry-conformance-tests.yml"
)
EXAMPLES_DIR = (
".build/durable-sdk/packages/aws-durable-execution-sdk-python-conformance-tests-otel"
)
EXAMPLES_DIR = ".build/durable-sdk/packages/aws-durable-execution-sdk-python-conformance-tests-otel"


def test_opentelemetry_conformance_caller_uses_current_workflow_contract() -> None:
Expand All @@ -24,6 +24,7 @@ def test_opentelemetry_conformance_caller_uses_current_workflow_contract() -> No

for configuration in (
"language: python",
"runs_on: codebuild-github-actions-runner-${{ github.run_id }}-${{ github.run_attempt }}",
"resource_prefix: p",
"sdk_repository: aws/aws-durable-execution-sdk-python",
"sdk_ref: ${{ github.event.pull_request.head.sha || github.sha }}",
Expand Down Expand Up @@ -66,8 +67,39 @@ def test_opentelemetry_conformance_runs_when_the_handlers_change() -> None:
workflow = WORKFLOW_PATH.read_text()

trigger_path = (
" - "
'"packages/aws-durable-execution-sdk-python-conformance-tests-otel/**"'
' - "packages/aws-durable-execution-sdk-python-conformance-tests-otel/**"'
)
# Once for pull_request, once for push.
assert workflow.count(trigger_path) == 2


def test_opentelemetry_conformance_queues_complete_runs() -> None:
workflow = yaml.safe_load(WORKFLOW_PATH.read_text())

assert workflow["concurrency"] == {
"group": "otel-conformance-tests",
"cancel-in-progress": False,
"queue": "max",
}


def test_cloud_tests_queue_each_shared_runtime_stack() -> None:
workflow = yaml.safe_load(WORKFLOW_PATH.with_name("cloud-tests.yml").read_text())

assert workflow["jobs"]["example-tests"]["concurrency"] == {
"group": "cloud-tests-${{ matrix.python-prefix }}",
"cancel-in-progress": False,
"queue": "max",
}


def test_general_conformance_keeps_all_pending_runs() -> None:
workflow = yaml.safe_load(
WORKFLOW_PATH.with_name("conformance-tests.yml").read_text()
)

assert workflow["concurrency"] == {
"group": "conformance-tests-global",
"cancel-in-progress": False,
"queue": "max",
}
2 changes: 2 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -84,6 +84,8 @@ jobs:
run: hatch run types:check
- name: Run tests + coverage
run: hatch run test:cov
- name: Verify supported legacy core compatibility
run: hatch run test-pypi-otel-legacy:test
- name: Build distribution
run: |
for pkg in packages/*/; do
Expand Down
1 change: 1 addition & 0 deletions .github/workflows/cloud-tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,7 @@ jobs:
concurrency:
group: cloud-tests-${{ matrix.python-prefix }}
cancel-in-progress: false
queue: max
strategy:
fail-fast: false
matrix:
Expand Down
1 change: 1 addition & 0 deletions .github/workflows/conformance-tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@ concurrency:
# run against the shared stack.
group: conformance-tests-global
cancel-in-progress: false
queue: max

permissions:
contents: read
Expand Down
10 changes: 9 additions & 1 deletion .github/workflows/opentelemetry-conformance-tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,13 @@ on:
default: main
type: string

# Backend stacks are shared across PRs. Queue whole runs so reusable
# backend jobs do not displace another PR from their pending slots.
concurrency:
group: otel-conformance-tests
cancel-in-progress: false
queue: max

This comment was marked as outdated.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Checked against the current official Actions workflow schema: concurrency-mapping includes queue, with allowed values single and max. Both workflow-level and job-level concurrency use that mapping.

Runtime validation agrees: Cloud tests and Conformance Tests succeeded on the reviewed commit 68166eb; all six OTel test jobs also succeeded with these same concurrency declarations before the subsequent JS harness-only commit. These workflows are therefore accepted and executing. Keeping queue: max with cancel-in-progress: false preserves pending runs without canceling the active tests.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Codex AI review · Finding arf_v1_bsc5lxoxpqaaru6j2vtvgojbll

[P1] Remove the unsupported queue key. GitHub Actions concurrency accepts group and cancel-in-progress; queue: max makes this workflow fail validation before any job runs. The same key was added to cloud-tests.yml and conformance-tests.yml, disabling those workflows too. Remove all three keys, and use an external locking or dispatch mechanism if every pending run must be retained.


permissions: {}

jobs:
Expand All @@ -59,9 +66,10 @@ jobs:
actions: write
contents: read
id-token: write
uses: aws/aws-durable-execution-conformance-tests/.github/workflows/opentelemetry-orchestrator.yml@a628f5589bbf067a441696c792ab3023c0d0899b
uses: aws/aws-durable-execution-conformance-tests/.github/workflows/opentelemetry-orchestrator.yml@a66037abbbfa55fde97f714e30f0bc262edefd63
with:
language: python
runs_on: codebuild-github-actions-runner-${{ github.run_id }}-${{ github.run_attempt }}
resource_prefix: p
sdk_repository: aws/aws-durable-execution-sdk-python
sdk_ref: ${{ github.event.pull_request.head.sha || github.sha }}
Expand Down
3 changes: 2 additions & 1 deletion CONTRIBUTING.md
Original file line number Diff line number Diff line change
Expand Up @@ -75,7 +75,8 @@ hatch run dev-examples:test # run examples tests only
To verify packages work against the published PyPI version of the core SDK (rather than the local workspace):

```bash
hatch run test-pypi-otel:test # test otel against PyPI core SDK
hatch run test-pypi-otel:test # test new OTel capabilities against capable installed core
hatch run test-pypi-otel-legacy:test # valid registrations/lifecycles on supported core 2.0.x
hatch run test-pypi-examples:test # test examples against PyPI core SDK
```

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -8,8 +8,8 @@ version = "0.0.0"
description = "OpenTelemetry conformance test handlers for the AWS Durable Execution SDK for Python, exercised by the aws-durable-execution-conformance-tests OTel suites."
requires-python = ">=3.11"
dependencies = [
"aws-durable-execution-sdk-python==2.0.1",
"aws-durable-execution-sdk-python-otel==1.0.0",
"aws-durable-execution-sdk-python==2.1.0",
"aws-durable-execution-sdk-python-otel==1.1.0",
]

[tool.hatch.build.targets.wheel]
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@ version = "0.0.0"
description = "Cross-SDK conformance test handlers for the AWS Durable Execution SDK for Python, exercised by the aws-durable-execution-conformance-tests runner."
requires-python = ">=3.11"
dependencies = [
"aws-durable-execution-sdk-python==2.0.1",
"aws-durable-execution-sdk-python==2.1.0",
]

[tool.hatch.build.targets.wheel]
Expand Down
26 changes: 25 additions & 1 deletion packages/aws-durable-execution-sdk-python-otel/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,30 @@ DURABLE_EXECUTION_PLUGINS=otel-execution
cold start, so the handler does not need to import or explicitly register the
plugin.

Automatic mutual-exclusion validation requires core SDK 2.1.0 or later together
with OTel 1.1.0 or later. OTel 1.1 remains compatible with core 2.0.x for existing
valid registrations; those older cores do not enforce the new group metadata.
Configure only one OTel view on every core version. `InvocationOtelPlugin`
shows work within each Lambda invocation; `ExecutionOtelPlugin` shows logical
operations across the whole execution. They emit overlapping telemetry and
manage competing active contexts, so the coordinated core 2.1+/OTel 1.1+ pair
rejects both at cold start with `PluginLoadError` naming both views. Keep only one across the combined
`plugins=[...]` and `DURABLE_EXECUTION_PLUGINS` configuration. Unrelated plugins
may run alongside the selected view; using no OTel plugin is also supported.

The bundled views explicitly declare `__durable_registration_api__ = 1`.
Custom subclasses inherit their group validation and registration cleanup.
The closest marker gates registration and selects the callback contract; an
explicit value other than the integer `1` shadows the inherited opt-in.
Once enabled, each MRO class explicitly declaring the integer `1` contributes
its resolved `exclusive_group`. A subclass's own group adds to the inherited
view group rather than replacing it; `None` cannot remove an inherited group.
Repeated groups count once per plugin. Unmarked classes' unrelated fields/helpers
named `exclusive_group` or `on_registration_result` remain inert. Repeat the marker
to add a group or customize the callback, which is still notified once per plugin.
Re-accepting a previously rejected bundled instance
restores its constructor-time logger enrichment immediately.

### 1. ADOT Lambda Layer

This plugin requires the [AWS Distro for OpenTelemetry (ADOT) Lambda layer](https://aws-otel.github.io/docs/getting-started/lambda) to export traces from your Lambda function.
Expand Down Expand Up @@ -406,7 +430,7 @@ setups.
## Requirements

- Python >= 3.11
- `aws-durable-execution-sdk-python` >= 2.0.0
- `aws-durable-execution-sdk-python` >= 2.0.0 (core >= 2.1.0 with OTel >= 1.1.0 for automatic view-exclusivity validation)
- An ADOT/community OpenTelemetry Lambda layer, or the `standalone` extra

## License
Expand Down
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
# SPDX-FileCopyrightText: 2025-present Amazon.com, Inc. or its affiliates.
#
# SPDX-License-Identifier: Apache-2.0
__version__ = "1.0.0"
__version__ = "1.1.0"
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,7 @@
import datetime
import logging
import threading
from typing import Any
from typing import Any, ClassVar

from aws_durable_execution_sdk_python.plugin import (
DurableInstrumentationPlugin,
Expand Down Expand Up @@ -91,7 +91,10 @@
canonical_trace_id,
)
from aws_durable_execution_sdk_python_otel.otel_plugin_config import OtelPluginConfig
from aws_durable_execution_sdk_python_otel.log_filter import install_log_filter
from aws_durable_execution_sdk_python_otel.log_filter import (
install_log_filter,
uninstall_log_filter,
)
from aws_durable_execution_sdk_python_otel.provider import create_tracer_provider


Expand Down Expand Up @@ -124,6 +127,9 @@ class ExecutionOtelPlugin(DurableInstrumentationPlugin):
span).
"""

__durable_registration_api__: ClassVar[int] = 1
exclusive_group: ClassVar[str | None] = "aws-durable-execution-otel-view"
Comment thread
zhongkechen marked this conversation as resolved.
Comment thread
zhongkechen marked this conversation as resolved.

def __init__(self, config: OtelPluginConfig | None = None) -> None:
self._config = config or OtelPluginConfig()
self._context_extractor: ContextExtractor = (
Expand Down Expand Up @@ -165,9 +171,19 @@ def __init__(self, config: OtelPluginConfig | None = None) -> None:
self._lock = threading.RLock()
self._tracing_enabled = False

self._registration_accepted = False
if self._config.enrich_logger:
install_log_filter(self)

def on_registration_result(self, registered: bool) -> None:
"""Preserve accepted resources, releasing only a discarded constructor."""
if registered:
self._registration_accepted = True
if self._config.enrich_logger:
install_log_filter(self)
elif not self._registration_accepted:
uninstall_log_filter(self)

def _bind_sdk_tracer(self) -> bool:
"""Bind to an SDK tracer, retrying a deferred global provider."""
self._sampling_delegate = None
Expand Down Expand Up @@ -414,6 +430,7 @@ def _with_sampling(self, parent_context: Context) -> Context:
# ------------------------------------------------------------------
def on_invocation_start(self, info: InvocationStartInfo) -> None:
logger.debug("Durable invocation started: %s", info)
self._registration_accepted = True
self._reset_state()
if info.execution_start_time is None:
logger.warning(
Expand Down Expand Up @@ -488,6 +505,10 @@ def on_invocation_start(self, info: InvocationStartInfo) -> None:
),
)

# Cover handlers installed after construction as well.
if self._config.enrich_logger:
install_log_filter(self)

def _start_workflow_span(self, info: InvocationStartInfo) -> None:
"""Install a non-recording placeholder for the execution-scoped Workflow span.

Expand Down
Loading
Loading