Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
70 commits
Select commit Hold shift + click to select a range
75de955
Align external stream input and output replay schedules
mdashti Sep 6, 2026
c00339e
Preserve stream task retention with zero workflow cache
mdashti Sep 6, 2026
632e16f
Put the shared stream interface on the repaired external SDK.
moedash Sep 14, 2026
66f353c
Let a terminal record land on a consumer that is closing.
moedash Sep 6, 2026
f06271e
Added the temporalio.streams interface with a provider registry.
moedash Sep 15, 2026
3dd101f
Added the in-memory reference provider and the conformance tests.
moedash Sep 15, 2026
62be804
Added the shared demo loop.
moedash Sep 15, 2026
7964e49
Moved the client-side binding into the provider registry.
moedash Sep 15, 2026
af2483d
Added the prepare and drain hooks the Workflow Streams transport needs.
moedash Sep 16, 2026
0741345
Made cursors exclusive, added latest(), and allowed topic appends.
moedash Sep 16, 2026
fb74e79
Made cursors exclusive, added latest(), and allowed topic appends.
moedash Sep 16, 2026
295fe3f
Applied an external stream replay marker before the task's first drain.
moedash Sep 16, 2026
668e432
Sorted and formatted the streams package for ruff.
moedash Sep 16, 2026
0877fe4
Made the streams package pass the type and doc linters.
moedash Sep 16, 2026
29cc1cb
Made the Redis provider pass the linters.
moedash Sep 16, 2026
89c1a37
Gave the external stream test fakes the shapes the real path has.
moedash Sep 16, 2026
97849a8
Snapshotted the markers where the abandoned task begins.
moedash Sep 16, 2026
d963173
Added the changelog entry for the shared stream interface.
moedash Sep 16, 2026
5bf0094
Made the demo call the provider lifecycle hooks.
moedash Sep 16, 2026
105ad9b
Support google-genai 2.21 file downloads (#1865)
brianstrauch Sep 15, 2026
e69ad05
Fix tests with latest dependencies (#1798)
brianstrauch Aug 31, 2026
914ee40
Adopted main's Nexus system API generator script.
moedash Sep 17, 2026
14b5042
Loosened the status error's response annotation.
moedash Sep 17, 2026
c4bed86
Held the external stream suite off the time-skipping server.
moedash Sep 17, 2026
86dc631
Gave each mock tool call its own id.
moedash Sep 17, 2026
c27c620
Repinned Core to the finalized external repairs branch.
moedash Sep 17, 2026
6b34091
Regenerated the Nexus system API from the repinned Core.
moedash Sep 17, 2026
9437d75
Adopted main's payload visitor generator.
moedash Sep 17, 2026
e6007d3
Made the external stream docstrings build under pydoctor.
moedash Sep 17, 2026
a3f50cd
Woke outside readers across threads and ignored idle_timeout in memory.
moedash Sep 18, 2026
528b811
Named the last record in append's cursor and allowed None for it.
moedash Sep 18, 2026
217497f
Resolved producer identity once in the package.
moedash Sep 18, 2026
8495140
Required an explicit provider and exported instance.
moedash Sep 18, 2026
ad177b2
Typed the handles and records honestly.
moedash Sep 18, 2026
952adfb
Kept topics off inbound frames and skipped unreadable frames with a w…
moedash Sep 18, 2026
276cb15
Let a broken provider import surface instead of vanishing from the re…
moedash Sep 18, 2026
785402d
Added the workflow-side conformance tests.
moedash Sep 18, 2026
7325643
Replaced stale wording and renamed the demo types.
moedash Sep 18, 2026
0c7616c
Keyed stores by one helper that survives a colon in the workflow id.
moedash Sep 19, 2026
8368735
Added an async close hook next to prepare and drain.
moedash Sep 19, 2026
e3522ae
Parametrised the conformance suite over providers.
moedash Sep 19, 2026
31f6fc0
Let a provider setup own the workflow the conformance cases address.
moedash Sep 19, 2026
f507eef
Recorded whether a provider reports a dropped repeat.
moedash Sep 19, 2026
17bbcfa
Dropped the committed demo output and corrected two comments.
moedash Sep 19, 2026
f71a030
Adapted the Redis provider to the interface changes.
moedash Sep 19, 2026
6fb7abe
Registered the Redis provider in the conformance suite.
moedash Sep 19, 2026
05ac20d
Ran the interface loop and the replay query in one Redis live module.
moedash Sep 19, 2026
e1d1b22
Renamed the store key helper to inbound_stream_id for the provider br…
moedash Sep 19, 2026
8b9f110
Repinned Core to the reviewed moe/AI-198-external-core-repairs head.
moedash Sep 19, 2026
362b643
Added the StreamError family for stream conditions.
moedash Sep 21, 2026
600f6d5
Vendored the stream additions to the public API protos.
moedash Sep 21, 2026
4362c31
Split the stream provider in two and made the proto the record.
moedash Sep 21, 2026
e05e651
Carried the stream provider from the worker into the workflow runtime.
moedash Sep 21, 2026
d37f8ad
Moved the stream conformance suite onto the new surface.
moedash Sep 21, 2026
1ea907c
Moved the stream demo and changelog entry onto the new surface.
moedash Sep 21, 2026
2659a17
Rewrote the Redis provider with an input and an output stream per topic.
moedash Sep 21, 2026
88a6298
Kept the stream finish hook off evicted runs and collected coroutines.
moedash Sep 21, 2026
d11aa01
Typed the hook test's workflow input as awaitable.
moedash Sep 21, 2026
4995cf1
Added activity.stream_handle and Client.get_stream_handle.
moedash Sep 21, 2026
dc32947
Registered the demo's provider on the client and reported from the ac…
moedash Sep 21, 2026
a667640
Registered the Redis conformance provider on the client.
moedash Sep 21, 2026
be743fd
Added typed topic definitions.
moedash Sep 21, 2026
c37b351
Took topic definitions on the Redis handle.
moedash Sep 21, 2026
a9bda0f
Matched the Redis live module to typed topic definitions.
moedash Sep 21, 2026
62b57e3
Seeded a Redis workflow reader from a cursor.
moedash Sep 21, 2026
c7e30af
Passed a backend's integrity loss through the replay read.
moedash Sep 21, 2026
d87787d
Trimmed the Redis streams by retention on every append.
moedash Sep 21, 2026
4b6d572
Listed stream_provider in ClientConnectConfig.
moedash Sep 22, 2026
6ec392f
Answered a legacy query without the stream snapshot beside it.
moedash Sep 22, 2026
d3aeb56
Repinned Core to the repairs head that carries the stream protos.
moedash Sep 22, 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
3 changes: 3 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -15,3 +15,6 @@ temporalio/bridge/temporal_sdk_bridge*
tags
/.claude
tmpclaude-*

# Demo run output, written per provider; the numbers live in the doc.
streams_demo/results-*/
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
110 changes: 110 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -18,8 +18,40 @@ 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

- **Experimental**: `temporalio.streams` defines one stream interface a workflow
can read, decide on, and write. A provider is registered once as a plugin,
`Client.connect(plugins=[provider])`, and workers built from that client
inherit it; each context then asks for its stream the same way:
`workflow.stream_reader()` and `workflow.stream_writer()` in workflow code,
`activity.stream_handle()` in an activity, and `client.get_stream_handle()`
anywhere a client is held. A topic is a typed definition,
`streams.topic("inputs", Token)`, shared by workflow, activity and client
code; a plain string names a topic decided at runtime. The record on the wire
is `temporal.api.stream.v1.StreamRecord` on every provider.
`temporalio.streams.providers.memory.MemoryStreams` is the in-memory
reference provider the conformance tests run against, and
`temporalio.streams.providers.redis.RedisStreams` serves the same interface
over External Workflow Streams, one topic as an input and an output stream.
- `ExternalStreamSubscription.records()` yields each value with the provider
offset it was read from, for a reader that has to name where it got to.

- Added experimental External Workflow Streams in
`temporalio.contrib.external_workflow_streams`. Workflow stream payloads are
stored in a configured external backend instead of Temporal History, with a
Expand All @@ -33,6 +65,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