Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
34 commits
Select commit Hold shift + click to select a range
71dd547
Aligned the external stream input and output replay schedules.
moedash Sep 25, 2026
e7e593a
Gave the external stream tests the shapes the real path has.
moedash Sep 25, 2026
781f8b1
Made the external stream docstrings build under pydoctor.
moedash Sep 25, 2026
186987b
Backported the upstream fixes for the latest dependency set.
moedash Sep 25, 2026
b7cb0b1
Adopted main's generator scripts and regenerated their output.
moedash Sep 25, 2026
becb881
Repinned Core to the reviewed external repairs head.
moedash Sep 25, 2026
a62917c
Answered a legacy query without the stream snapshot beside it.
moedash Sep 25, 2026
8e2deaa
Pointed the legacy query id at Core's own constant.
moedash Sep 25, 2026
891d514
Kept a replay marker from earning a drain of its own.
moedash Sep 25, 2026
69987f2
Held only the server-backed stream cases off the skipping clock.
moedash Sep 25, 2026
506475a
Repinned Core to the external repairs head.
moedash Sep 26, 2026
85f3618
Repinned Core for the wake protos.
moedash Oct 1, 2026
d33b6af
Sent external-stream wakes over the wake call, with a Signal fallback.
moedash Oct 1, 2026
b3c3e1f
Kept the producer's doc comment free of a pydoctor link.
moedash Oct 1, 2026
74e0b75
Allowed external-stream publishes from the workflow constructor.
moedash Oct 1, 2026
f57fd80
Satisfied basedpyright on the wake transport and the new external-str…
moedash Oct 2, 2026
9c19f23
Pinned Core with the notification channel and regenerated the clients.
moedash Oct 2, 2026
2ce4cd9
Notified the stream's channel ahead of the wake call and the Signal.
moedash Oct 2, 2026
500f494
Subscribed external-stream readers to their streams' channels.
moedash Oct 2, 2026
65b96d0
Covered the channel transport, the subscription and the notified reader.
moedash Oct 2, 2026
bb7e6fb
Added the notification channel surface to the workflow and client.
moedash Oct 2, 2026
8c40b07
Derived a positionless wake's counter from the backend's clock rule.
moedash Oct 2, 2026
1e84b51
Skipped the retained-first-task cases on servers with channels.
moedash Oct 2, 2026
3aaa09d
Pinned Core without the point-to-point wake and with the linked chann…
moedash Oct 2, 2026
12fa935
Removed the point-to-point wake transport.
moedash Oct 2, 2026
c14a88d
Listened on the linked channel for the streams a run owns.
moedash Oct 2, 2026
93f9ac1
Addressed the client's channel calls to a workflow's linked channel.
moedash Oct 2, 2026
10ca9cb
Matched the main chain's linked channel surface and the listener list.
moedash Oct 2, 2026
1ce5766
Repinned Core to the external repairs head.
moedash Oct 2, 2026
68cdb9e
Repinned Core to the unsubscribe and channel-report head.
moedash Oct 2, 2026
4795652
Reported the readers' channels instead of subscribing to them.
moedash Oct 2, 2026
49faf69
Pinned Core at the protos that address a channel by execution.
moedash Oct 3, 2026
a59c3a8
Addressed a linked channel by execution.
moedash Oct 3, 2026
97ddcd8
Pinned Core at the protos that reuse the old numbers for execution.
moedash Oct 3, 2026
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
4 changes: 2 additions & 2 deletions .gitmodules
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
[submodule "sdk-core"]
path = temporalio/bridge/sdk-core
url = https://github.com/mfateev/sdk-core.git
branch = task/python-sdk-streaming
url = https://github.com/moedash/sdk-rust.git
branch = moe/AI-198-external-core-repairs
93 changes: 93 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,21 @@ to include examples, links to docs, or any other relevant information.

## [Unreleased]

### Fixed

- Preserve empty activations in the shared External Workflow Streams input and
output replay schedule. Workflows that read input, publish decisions, and
schedule Activities now reproduce that schedule during replay. Inconsistent
prerelease markers are rejected explicitly rather than guessing where omitted
activations belonged.
- Resume external input waits when cold replay encounters a wake in an already
loaded History page, including Workers with workflow caching disabled.
- Avoid an unnecessary output replacement Workflow Task after stream input has
resumed the Workflow and it is waiting on an Activity or timer.
- Keep an incomplete retained external stream task alive when workflow caching
is disabled; evict it after its normal task boundary instead of repeatedly
interrupting input readiness with shutdown markers.

### Added

- Added experimental External Workflow Streams in
Expand All @@ -33,6 +48,84 @@ to include examples, links to docs, or any other relevant information.
through `ExternalOutputStreamClient`. Workflow output is staged outside
History and becomes readable only after its compact Workflow Task marker is
committed.
### Changed

- Standalone Activities are now generally available (GA). (Standalone Activities as Nexus operations
and Standalone Activities operator commands remain experimental. Operator commands are `pause`,
`unpause`, `updateOptions`, `restoreOriginal`.)
- System Nexus Signal-with-Start Workflow operations now use the typed
`WorkflowOutboundInterceptor.start_signal_with_start_workflow` interception point instead of
the generic `WorkflowOutboundInterceptor.start_nexus_operation` method.
- System Nexus Signal-with-Start Workflow operations now invoke
`WorkflowOutboundInterceptor.start_system_nexus_operation` after their typed interception
point. They continue not to invoke `WorkflowOutboundInterceptor.start_nexus_operation`.
- The experimental `GetNexusOperationResultInput` now includes the Nexus endpoint, service, and
operation.

### :boom: Breaking Changes

- Experimental external storage: `ExternalStorage.driver_selector` is now called with a
`StorageDriverSelectContext` instead of a `StorageDriverStoreContext`. Update the annotation;
the new type carries the same `target` field. Since selectors are plain callables, a stale
annotation fails type checking rather than at runtime.
- `client.ActivityExecution` and `client.ActivityExecutionDescription` had some fields removed or renamed
to match RPC API.
- Dataclass parameters for these types were changed to `frozen=True, eq=False, kw_only=True`.
- `scheduled_time` was renamed `schedule_time`.
- `last_failure` was changed from field to method that runs data converter on demand.
- `state_transition_count`, `eager_execution_requested`, `paused` and `long_poll_token` were removed.
- ActivityHandle.describe() long-poll token was removed. The functionality can still be used manually
through raw gRPC API.

### Fixed

- `temporalio.contrib.google_genai` now requires `google-genai` 2.21.0 or later
and supports its file download API, including video inputs and download
destinations.
- `temporalio.contrib.deepagents` no longer dedups repeated identical tool,
model, and backend-op calls: each dispatch runs its own Activity, and the
continue-as-new result cache is retired for new executions (a continued run
resumes from the carried transcript and never re-executes prior dispatches,
so a carried cache entry could only serve stale results). Patch-gated
(`deepagents.retire-result-cache`), so histories recorded before this change
replay unchanged; note that deferring the patch keeps the full legacy dedup
cache — including the stale-result behavior this entry describes — and that
a chain upgraded mid-continue-as-new re-executes rather than reuses a
repeated identical call (the conservative direction).
- `contrib.deepagents`: summarization middleware configured with a model name string now routes its LLM calls through Activities instead of running them in the Workflow.

- **Experimental**: External storage metrics now report the wall-clock time storage was in flight.
Previously each batch's duration was summed, over-reporting the time whenever storage operations
ran concurrently.
- System Nexus Signal-with-Start workflow operations now give custom payload
converters the target workflow's serialization context when encoding their
inner request payloads.
- Cancelling an activity from a signal while the workflow itself is cancelled
no longer causes a nondeterminism error from duplicate activity-cancellation
commands.
- `StrandsPlugin` now disables Botocore retries for its default Bedrock model so
model request retries are handled exclusively by Temporal.
- `temporalio.contrib.openai_agents` now honors the `retry-after-ms` and
`retry-after` headers when OpenAI returns `x-should-retry: true`. Previously
the delay the server asked for was discarded on that path and the activity
retried on its configured interval instead.
- Nexus-context workflow/activity starts no longer set `on_conflict_options` when there are no links
or callbacks to attach.
- The workflow sandbox now passes `pydantic_core` through by default, alongside `pydantic`.

## [1.32.0] - 2026-08-24

### Added

- Added `temporalio.converter.create_payload_validation_error` to create the
non-retryable application error used when a converted payload fails validation.
- Added experimental `temporalio.contrib.opentelemetry.ReplaySafeMeterProvider` and
`ReplaySafeLoggerProvider` (and exported `ReplaySafeTracerProvider`): wrap an
OpenTelemetry provider so metrics and log events recorded from workflow code (e.g. by
Google ADK) are not duplicated on replay. `GoogleAdkPlugin` warns when a global OTel
provider is not replay-safe.
- Added `LoggingConfig.format` to select compact, pretty, or newline-delimited JSON output for
Core logs written to the console.

- Added the `Runtime(disable_environment_info=...)` option to control whether
runtime, hosting, and platform information is included in worker heartbeats.
Expand Down
7 changes: 6 additions & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -45,7 +45,7 @@ lambda-worker-otel = [
"opentelemetry-sdk-extension-aws>=2.0.0,<3",
]
aioboto3 = ["aioboto3>=10.4.0", "types-aioboto3[s3]>=10.4.0"]
google-genai = ["google-genai>=2.10.0,<3.0.0"]
google-genai = ["google-genai>=2.21.0,<3.0.0"]
strands-agents = ["strands-agents>=1.39.0"]

[project.urls]
Expand Down Expand Up @@ -188,8 +188,13 @@ exclude = [
# Ignore generated code
'temporalio/api',
'temporalio/bridge/proto',
'temporalio/nexus/system/workflow_service',
]

[[tool.mypy.overrides]]
module = "temporalio.nexus.system.workflow_service.*"
ignore_errors = true

[tool.pydocstyle]
convention = "google"
# https://github.com/PyCQA/pydocstyle/issues/363#issuecomment-625563088
Expand Down
25 changes: 18 additions & 7 deletions scripts/gen_nexus_system_api.py
Original file line number Diff line number Diff line change
Expand Up @@ -35,15 +35,28 @@
/ "v1"
/ "request_response.proto"
)
NEX_GEN_VERSION = "0.2.4"


def nex_gen_command() -> list[str]:
if bin_path := os.environ.get("NEX_GEN_BIN"):
return [bin_path]

if shutil.which("nex-gen") is None:
subprocess.check_call(["cargo", "install", "--locked", "nex-gen", "--force"])
return ["nex-gen"]
if shutil.which("nexgen") is None:
subprocess.check_call(
[
"cargo",
"install",
"--locked",
"nexgen",
"--version",
NEX_GEN_VERSION,
"--features",
"advanced",
"--force",
]
)
return ["nexgen"]


def build_descriptor_set(descriptor_path: Path) -> None:
Expand Down Expand Up @@ -114,13 +127,11 @@ def generate_nexus_system_api() -> None:
subprocess.check_call(
[
*command,
"generate",
"--lang",
"python",
"--input",
str(wit_path),
"--input",
str(wit_deps_dir),
"--native-api",
"--system-nexus",
"--support-file",
str(python_support_path),
"--descriptors",
Expand Down
37 changes: 25 additions & 12 deletions scripts/gen_payload_visitor.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@
import sys
from importlib.util import module_from_spec, spec_from_file_location
from pathlib import Path
from typing import cast
from typing import cast, get_args, get_origin, get_type_hints

import google.protobuf.message
import nexusrpc
Expand All @@ -19,14 +19,15 @@
from temporalio.bridge.proto.workflow_completion.workflow_completion_pb2 import (
WorkflowActivationCompletion,
)
from temporalio.converter._payload_converter import _get_transfer_type_converter


def discover_system_nexus_roots() -> list[Descriptor]:
module_path = (
base_dir / "temporalio" / "nexus" / "system" / "workflow_service" / "service.py"
)
module_path = base_dir / "temporalio" / "nexus" / "system" / "workflow_service"
spec = spec_from_file_location(
"temporalio_nexus_system_workflow_service", module_path
"temporalio_nexus_system_workflow_service",
module_path / "__init__.py",
submodule_search_locations=[str(module_path)],
)
if spec is None or spec.loader is None:
raise RuntimeError(f"Cannot load generated system service from {module_path}")
Expand All @@ -35,10 +36,14 @@ def discover_system_nexus_roots() -> list[Descriptor]:
spec.loader.exec_module(module)

roots: list[Descriptor] = []
for operation in vars(module.WorkflowService).values():
if not isinstance(operation, nexusrpc.Operation):
for annotation in get_type_hints(module._services.WorkflowService).values():
if get_origin(annotation) is not nexusrpc.Operation:
continue
for proto_type in (operation.input_type, operation.output_type):
for operation_type in get_args(annotation):
converter = _get_transfer_type_converter(operation_type)
proto_type = (
converter.transfer_type if converter is not None else operation_type
)
if isinstance(proto_type, type) and issubclass(
proto_type, google.protobuf.message.Message
):
Expand Down Expand Up @@ -119,12 +124,20 @@ def generate(self, roots: list[Descriptor]) -> str:

The generated code defines async visitor functions for each reachable
protobuf message type starting from WorkflowActivation, including support
for repeated fields and map entries, and a convenience entrypoint
function `visit`.
for repeated fields and map entries. Payload-free roots get no-op methods
so the `visit` entrypoint recognizes them as supported.
"""

for r in roots:
self.walk(r)
for root in roots:
if not self.walk(root):
self.methods.append(
f"""\
async def _visit_{name_for(root)}(
self, fs: VisitorFunctions, o: Any
) -> None:
pass
"""
)

header = """
from __future__ import annotations
Expand Down
101 changes: 101 additions & 0 deletions scripts/gen_stream_api_protos.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,101 @@
"""Regenerate the stream additions to the vendored public API protos.

The stream protos live on a branch of the api repo that the Core submodule
does not pin yet. Regenerating everything from that branch would also pull in
whatever else moved on the api main line since the pin, so this stages the
API protos the submodule pins, applies the branch's own diff on top, and
regenerates only the files that diff touches. Every other module under
``temporalio/api`` stays byte-identical.

uv run --python 3.10 --no-project --with "grpcio-tools==1.48.2" \\
--with "mypy-protobuf==3.3.0" --with "protobuf<4" \\
scripts/gen_stream_api_protos.py /path/to/api origin/main..origin/<branch>

The toolchain pins match ``scripts/_proto/Dockerfile``. Generated code refuses
to load on a protobuf runtime older than the one it was built against, and
these modules end up in applications that pin protobuf themselves. Run
``uv run poe format`` afterwards, as the ``gen-protos`` task does.
"""

from __future__ import annotations

import shutil
import subprocess
import sys
import tempfile
from pathlib import Path

BASE = Path(__file__).parent.parent
sys.path.insert(0, str(BASE / "scripts"))

import gen_protos # noqa: E402

API_OUT = BASE / "temporalio" / "api"


def main() -> None:
if len(sys.argv) != 3:
sys.exit("usage: gen_stream_api_protos.py /path/to/api-checkout <base>..<ref>")
api = Path(sys.argv[1]).resolve()
revisions = sys.argv[2]
diff = subprocess.run(
["git", "-C", str(api), "diff", revisions, "--", "temporal/"],
check=True,
capture_output=True,
text=True,
).stdout
touched = subprocess.run(
["git", "-C", str(api), "diff", "--name-only", revisions, "--", "temporal/"],
check=True,
capture_output=True,
text=True,
).stdout.split()
if not touched:
sys.exit(f"{revisions} touches no proto under temporal/")

with tempfile.TemporaryDirectory() as tmp:
stage = Path(tmp) / "stage"
shutil.copytree(gen_protos.api_proto_dir / "temporal", stage / "temporal")
subprocess.run(
["patch", "-p1", "--silent"],
check=True,
cwd=stage,
input=diff,
text=True,
)
out = Path(tmp) / "out"
out.mkdir()
subprocess.check_call(
[
sys.executable,
"-mgrpc_tools.protoc",
f"--proto_path={stage}",
f"--python_out={out}",
f"--mypy_out={out}",
*touched,
]
)
packages: set[Path] = set()
for proto in touched:
relative = Path(proto).relative_to("temporal/api").with_suffix("")
for suffix in ("_pb2.py", "_pb2.pyi"):
generated = (
out
/ "temporal"
/ "api"
/ relative.with_name(relative.name + suffix)
)
target = API_OUT / relative.with_name(relative.name + suffix)
target.parent.mkdir(parents=True, exist_ok=True)
shutil.copyfile(generated, target)
print(f"wrote {target.relative_to(BASE)}")
packages.add(target.parent)

for package in sorted(packages):
(package.parent / "__init__.py").touch()
gen_protos.fix_generated_output(package)
print(f"rewrote {(package / '__init__.py').relative_to(BASE)}")


if __name__ == "__main__":
main()
Loading
Loading