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
2 changes: 1 addition & 1 deletion .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -232,7 +232,7 @@ jobs:
- run: uv add --dev --python 3.10 "googleapis-common-protos==1.70.0"
- run: uv add --python 3.10 "protobuf<4"
- run: uv sync --all-extras
- run: cargo install --locked nexgen --version 0.2.4 --features advanced --force
- run: cargo install --locked nexgen --version 0.2.7 --features advanced --force
- run: poe build-develop
- run: poe gen-protos
- name: Check generation unchanged
Expand Down
8 changes: 8 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,14 @@ to include examples, links to docs, or any other relevant information.
instead of from an instance (i.e. its first parameter is an unbound `self`).
- `workflow.new_random()` accepts an optional `name` that is mixed into the seed, so differently
named generators, and `workflow.random()`, produce different sequences.
- **Experimental**: Standalone Nexus operations can report their outcome to a Nexus
service with completion callbacks. Pass `completion_callbacks` to
`NexusClient.start_operation` or `NexusClient.execute_operation`. Create the
callbacks with `temporalio.nexus.create_completion_callback`. When the
operation finishes, the server calls an on-complete operation on a worker for the
task queue of the callback. The call carries the result or the failure, and a source
context that you supply. The new `temporalio.nexus.notifications` module has the
input and output types of on-complete operations.

### Changed

Expand Down
4 changes: 4 additions & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -292,3 +292,7 @@ exclude = ["temporalio/bridge/target/**/*", "temporalio/bridge/sdk-core/.git"]
# Prevent uv commands from building the package by default
package = false
exclude-newer = "2 weeks"

[tool.uv.sources]
# TODO: Remove once a nexus-rpc release records the service of each operation
nexus-rpc = { git = "https://github.com/nexus-rpc/sdk-python", rev = "e1d70abbb9378a8398fd472bd466c5dfa9f1aca7" }
6 changes: 5 additions & 1 deletion scripts/gen_payload_visitor.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@
sys.path.insert(0, str(base_dir))

from temporalio.api.common.v1.message_pb2 import Payload, Payloads, SearchAttributes
from temporalio.api.notificationservice.v1 import OnCompleteRequest, OnCompleteResponse
from temporalio.bridge.proto.nexus import NexusTaskCompletion
from temporalio.bridge.proto.workflow_activation.workflow_activation_pb2 import (
WorkflowActivation,
Expand Down Expand Up @@ -440,11 +441,14 @@ def write_bridge_visitors() -> None:
out_path = base_dir / "temporalio" / "bridge" / "_visitor.py"

# Build root descriptors: WorkflowActivation, WorkflowActivationCompletion,
# NexusTaskCompletion, and the system Nexus operation roots.
# NexusTaskCompletion, the system Nexus operation roots, and the notification
# service messages delivered to Nexus handler callbacks.
roots: list[Descriptor] = [
WorkflowActivation.DESCRIPTOR,
WorkflowActivationCompletion.DESCRIPTOR,
NexusTaskCompletion.DESCRIPTOR,
OnCompleteRequest.DESCRIPTOR,
OnCompleteResponse.DESCRIPTOR,
] + discover_system_nexus_roots()

code = VisitorGenerator().generate(roots)
Expand Down
15 changes: 15 additions & 0 deletions temporalio/bridge/_visitor.py
Original file line number Diff line number Diff line change
Expand Up @@ -624,6 +624,21 @@ async def _visit_coresdk_nexus_NexusTaskCompletion(
elif o.HasField("failure"):
await self._visit_temporal_api_failure_v1_Failure(fs, o.failure)

async def _visit_temporal_api_notificationservice_v1_OnCompleteRequest(
self, fs: VisitorFunctions, o: Any
):
if o.HasField("source_context"):
await self._visit_temporal_api_common_v1_Payload(fs, o.source_context)
if o.HasField("success"):
await self._visit_temporal_api_common_v1_Payload(fs, o.success)
elif o.HasField("failure"):
await self._visit_temporal_api_failure_v1_Failure(fs, o.failure)

async def _visit_temporal_api_notificationservice_v1_OnCompleteResponse(
self, fs: VisitorFunctions, o: Any
) -> None:
pass

async def _visit_temporal_api_common_v1_Header(self, fs: VisitorFunctions, o: Any):
for v in o.fields.values():
await self._visit_temporal_api_common_v1_Payload(fs, v)
Expand Down
5 changes: 5 additions & 0 deletions temporalio/client/_impl.py
Original file line number Diff line number Diff line change
Expand Up @@ -1580,6 +1580,11 @@ async def start_nexus_operation(
"temporalio.api.enums.v1.NexusOperationIdConflictPolicy.ValueType",
int(input.id_conflict_policy),
),
# Callback source contexts use the started operation's serialization
# context, like the rest of the request.
completion_callbacks=[
await cb._to_proto(data_converter) for cb in input.completion_callbacks
],
)

if input.schedule_to_close_timeout is not None:
Expand Down
2 changes: 2 additions & 0 deletions temporalio/client/_interceptor.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
import temporalio.api.common.v1
import temporalio.api.workflowservice.v1
import temporalio.common
import temporalio.nexus
from temporalio.converter import DataConverter

if TYPE_CHECKING:
Expand Down Expand Up @@ -605,6 +606,7 @@ class StartNexusOperationInput:
headers: Mapping[str, str]
rpc_metadata: Mapping[str, str | bytes]
rpc_timeout: timedelta | None
completion_callbacks: Sequence[temporalio.nexus.CompletionCallback[Any]] = ()


@dataclass
Expand Down
Loading
Loading