Skip to content
Closed
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
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()
8 changes: 8 additions & 0 deletions temporalio/api/command/v1/__init__.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
from .message_pb2 import (
AppendStreamRecordsCommandAttributes,
CancelTimerCommandAttributes,
CancelWorkflowExecutionCommandAttributes,
Command,
Expand All @@ -16,10 +17,14 @@
SignalExternalWorkflowExecutionCommandAttributes,
StartChildWorkflowExecutionCommandAttributes,
StartTimerCommandAttributes,
SubscribeNotificationChannelCommandAttributes,
SubscribeStreamCommandAttributes,
UnsubscribeNotificationChannelCommandAttributes,
UpsertWorkflowSearchAttributesCommandAttributes,
)

__all__ = [
"AppendStreamRecordsCommandAttributes",
"CancelTimerCommandAttributes",
"CancelWorkflowExecutionCommandAttributes",
"Command",
Expand All @@ -37,5 +42,8 @@
"SignalExternalWorkflowExecutionCommandAttributes",
"StartChildWorkflowExecutionCommandAttributes",
"StartTimerCommandAttributes",
"SubscribeNotificationChannelCommandAttributes",
"SubscribeStreamCommandAttributes",
"UnsubscribeNotificationChannelCommandAttributes",
"UpsertWorkflowSearchAttributesCommandAttributes",
]
153 changes: 112 additions & 41 deletions temporalio/api/command/v1/message_pb2.py

Large diffs are not rendered by default.

203 changes: 203 additions & 0 deletions temporalio/api/command/v1/message_pb2.pyi
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ import temporalio.api.enums.v1.workflow_pb2
import temporalio.api.failure.v1.message_pb2
import temporalio.api.sdk.v1.event_group_marker_pb2
import temporalio.api.sdk.v1.user_metadata_pb2
import temporalio.api.stream.v1.message_pb2
import temporalio.api.taskqueue.v1.message_pb2
import temporalio.api.workflow.v1.message_pb2

Expand Down Expand Up @@ -1119,6 +1120,10 @@ class Command(google.protobuf.message.Message):
MODIFY_WORKFLOW_PROPERTIES_COMMAND_ATTRIBUTES_FIELD_NUMBER: builtins.int
SCHEDULE_NEXUS_OPERATION_COMMAND_ATTRIBUTES_FIELD_NUMBER: builtins.int
REQUEST_CANCEL_NEXUS_OPERATION_COMMAND_ATTRIBUTES_FIELD_NUMBER: builtins.int
APPEND_STREAM_RECORDS_COMMAND_ATTRIBUTES_FIELD_NUMBER: builtins.int
SUBSCRIBE_STREAM_COMMAND_ATTRIBUTES_FIELD_NUMBER: builtins.int
SUBSCRIBE_NOTIFICATION_CHANNEL_COMMAND_ATTRIBUTES_FIELD_NUMBER: builtins.int
UNSUBSCRIBE_NOTIFICATION_CHANNEL_COMMAND_ATTRIBUTES_FIELD_NUMBER: builtins.int
command_type: temporalio.api.enums.v1.command_type_pb2.CommandType.ValueType
@property
def user_metadata(self) -> temporalio.api.sdk.v1.user_metadata_pb2.UserMetadata:
Expand Down Expand Up @@ -1209,6 +1214,22 @@ class Command(google.protobuf.message.Message):
def request_cancel_nexus_operation_command_attributes(
self,
) -> global___RequestCancelNexusOperationCommandAttributes: ...
@property
def append_stream_records_command_attributes(
self,
) -> global___AppendStreamRecordsCommandAttributes: ...
@property
def subscribe_stream_command_attributes(
self,
) -> global___SubscribeStreamCommandAttributes: ...
@property
def subscribe_notification_channel_command_attributes(
self,
) -> global___SubscribeNotificationChannelCommandAttributes: ...
@property
def unsubscribe_notification_channel_command_attributes(
self,
) -> global___UnsubscribeNotificationChannelCommandAttributes: ...
def __init__(
self,
*,
Expand Down Expand Up @@ -1253,10 +1274,20 @@ class Command(google.protobuf.message.Message):
| None = ...,
request_cancel_nexus_operation_command_attributes: global___RequestCancelNexusOperationCommandAttributes
| None = ...,
append_stream_records_command_attributes: global___AppendStreamRecordsCommandAttributes
| None = ...,
subscribe_stream_command_attributes: global___SubscribeStreamCommandAttributes
| None = ...,
subscribe_notification_channel_command_attributes: global___SubscribeNotificationChannelCommandAttributes
| None = ...,
unsubscribe_notification_channel_command_attributes: global___UnsubscribeNotificationChannelCommandAttributes
| None = ...,
) -> None: ...
def HasField(
self,
field_name: typing_extensions.Literal[
"append_stream_records_command_attributes",
b"append_stream_records_command_attributes",
"attributes",
b"attributes",
"cancel_timer_command_attributes",
Expand Down Expand Up @@ -1291,6 +1322,12 @@ class Command(google.protobuf.message.Message):
b"start_child_workflow_execution_command_attributes",
"start_timer_command_attributes",
b"start_timer_command_attributes",
"subscribe_notification_channel_command_attributes",
b"subscribe_notification_channel_command_attributes",
"subscribe_stream_command_attributes",
b"subscribe_stream_command_attributes",
"unsubscribe_notification_channel_command_attributes",
b"unsubscribe_notification_channel_command_attributes",
"upsert_workflow_search_attributes_command_attributes",
b"upsert_workflow_search_attributes_command_attributes",
"user_metadata",
Expand All @@ -1300,6 +1337,8 @@ class Command(google.protobuf.message.Message):
def ClearField(
self,
field_name: typing_extensions.Literal[
"append_stream_records_command_attributes",
b"append_stream_records_command_attributes",
"attributes",
b"attributes",
"cancel_timer_command_attributes",
Expand Down Expand Up @@ -1338,6 +1377,12 @@ class Command(google.protobuf.message.Message):
b"start_child_workflow_execution_command_attributes",
"start_timer_command_attributes",
b"start_timer_command_attributes",
"subscribe_notification_channel_command_attributes",
b"subscribe_notification_channel_command_attributes",
"subscribe_stream_command_attributes",
b"subscribe_stream_command_attributes",
"unsubscribe_notification_channel_command_attributes",
b"unsubscribe_notification_channel_command_attributes",
"upsert_workflow_search_attributes_command_attributes",
b"upsert_workflow_search_attributes_command_attributes",
"user_metadata",
Expand Down Expand Up @@ -1365,8 +1410,166 @@ class Command(google.protobuf.message.Message):
"modify_workflow_properties_command_attributes",
"schedule_nexus_operation_command_attributes",
"request_cancel_nexus_operation_command_attributes",
"append_stream_records_command_attributes",
"subscribe_stream_command_attributes",
"subscribe_notification_channel_command_attributes",
"unsubscribe_notification_channel_command_attributes",
]
| None
): ...

global___Command = Command

class AppendStreamRecordsCommandAttributes(google.protobuf.message.Message):
"""Appends records to a stream the Workflow owns. Applied inside the Workflow
Task's own commit. Produces one `WorkflowStreamRecordsAppended` event
carrying the offset range and none of the payload; it schedules no further
work.
"""

DESCRIPTOR: google.protobuf.descriptor.Descriptor

STREAM_NAME_FIELD_NUMBER: builtins.int
RECORDS_FIELD_NUMBER: builtins.int
stream_name: builtins.str
"""Name of a stream this Workflow owns, scoped to the Workflow. Created on
first use. Empty means the Workflow's default output stream. A Workflow
cannot append to a stream in another execution, so this is never the id
of a standalone stream.
"""
@property
def records(
self,
) -> google.protobuf.internal.containers.RepeatedCompositeFieldContainer[
temporalio.api.stream.v1.message_pb2.StreamRecord
]:
"""Stored in order. The server sets `producer_id` to empty on each record,
because the owning Workflow is the producer here.
"""
def __init__(
self,
*,
stream_name: builtins.str = ...,
records: collections.abc.Iterable[
temporalio.api.stream.v1.message_pb2.StreamRecord
]
| None = ...,
) -> None: ...
def ClearField(
self,
field_name: typing_extensions.Literal[
"records", b"records", "stream_name", b"stream_name"
],
) -> None: ...

global___AppendStreamRecordsCommandAttributes = AppendStreamRecordsCommandAttributes

class SubscribeStreamCommandAttributes(google.protobuf.message.Message):
"""Subscribe this Workflow to a stream, so later Workflow Tasks carry the ranges
it has not consumed yet.

The stream's addressing is resolved by the server rather than supplied here.
A Workflow cannot look it up without doing I/O, and a value it carried would
be a reading rather than a fact, so it could differ on replay.
"""

DESCRIPTOR: google.protobuf.descriptor.Descriptor

STREAM_NAME_OR_ID_FIELD_NUMBER: builtins.int
START_OFFSET_FIELD_NUMBER: builtins.int
START_POSITION_FIELD_NUMBER: builtins.int
stream_name_or_id: builtins.str
"""Stream to consume, named either way round: a stream this Workflow owns
by the name it appends under, a stream in another execution by its id.
The server tries them in that order, so a Workflow that owns a stream
under this name cannot reach a standalone stream with the same id. When
neither exists the Workflow gets a stream of its own by that name, which
is how a reader subscribes before the first record is written.
"""
start_offset: builtins.int
"""Where to start, as an absolute offset. Read only when `start_position`
is unset. A negative value is refused: the head of the stream is asked
for with `start_position.tail`.
"""
@property
def start_position(
self,
) -> temporalio.api.stream.v1.message_pb2.StreamStartPosition:
"""Where to start. The server resolves it once, when it registers the
subscription, and records the resolved absolute offset on the subscribed
event, so replay does not resolve it again. Setting it together with a
non-zero `start_offset` fails the command.
"""
def __init__(
self,
*,
stream_name_or_id: builtins.str = ...,
start_offset: builtins.int = ...,
start_position: temporalio.api.stream.v1.message_pb2.StreamStartPosition
| None = ...,
) -> None: ...
def HasField(
self, field_name: typing_extensions.Literal["start_position", b"start_position"]
) -> builtins.bool: ...
def ClearField(
self,
field_name: typing_extensions.Literal[
"start_offset",
b"start_offset",
"start_position",
b"start_position",
"stream_name_or_id",
b"stream_name_or_id",
],
) -> None: ...

global___SubscribeStreamCommandAttributes = SubscribeStreamCommandAttributes

class SubscribeNotificationChannelCommandAttributes(google.protobuf.message.Message):
"""Makes the Workflow a listener of a notification channel for this run. The
next notifications on the channel arrive on the scheduled event of a Workflow
Task. The subscription ends with the run, and a successor subscribes again.
"""

DESCRIPTOR: google.protobuf.descriptor.Descriptor

CHANNEL_FIELD_NUMBER: builtins.int
channel: builtins.str
"""The channel to listen on, as the writers name it."""
def __init__(
self,
*,
channel: builtins.str = ...,
) -> None: ...
def ClearField(
self, field_name: typing_extensions.Literal["channel", b"channel"]
) -> None: ...

global___SubscribeNotificationChannelCommandAttributes = (
SubscribeNotificationChannelCommandAttributes
)

class UnsubscribeNotificationChannelCommandAttributes(google.protobuf.message.Message):
"""Ends the run's subscription to a notification channel. Notifications already
recorded on a scheduled event still reach that Workflow Task; later ones do
not. A command naming a channel the run is not subscribed to records its
event and changes nothing, so replay matches every command to an event.
"""

DESCRIPTOR: google.protobuf.descriptor.Descriptor

CHANNEL_FIELD_NUMBER: builtins.int
channel: builtins.str
"""The channel to stop listening on, as the writers name it."""
def __init__(
self,
*,
channel: builtins.str = ...,
) -> None: ...
def ClearField(
self, field_name: typing_extensions.Literal["channel", b"channel"]
) -> None: ...

global___UnsubscribeNotificationChannelCommandAttributes = (
UnsubscribeNotificationChannelCommandAttributes
)
Loading
Loading