Skip to content
Open
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
4 changes: 4 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 @@ -17,11 +18,13 @@
StartChildWorkflowExecutionCommandAttributes,
StartTimerCommandAttributes,
SubscribeNotificationChannelCommandAttributes,
SubscribeStreamCommandAttributes,
UnsubscribeNotificationChannelCommandAttributes,
UpsertWorkflowSearchAttributesCommandAttributes,
)

__all__ = [
"AppendStreamRecordsCommandAttributes",
"CancelTimerCommandAttributes",
"CancelWorkflowExecutionCommandAttributes",
"Command",
Expand All @@ -40,6 +43,7 @@
"StartChildWorkflowExecutionCommandAttributes",
"StartTimerCommandAttributes",
"SubscribeNotificationChannelCommandAttributes",
"SubscribeStreamCommandAttributes",
"UnsubscribeNotificationChannelCommandAttributes",
"UpsertWorkflowSearchAttributesCommandAttributes",
]
125 changes: 80 additions & 45 deletions temporalio/api/command/v1/message_pb2.py

Large diffs are not rendered by default.

130 changes: 130 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 @@ -1121,6 +1122,8 @@ class Command(google.protobuf.message.Message):
REQUEST_CANCEL_NEXUS_OPERATION_COMMAND_ATTRIBUTES_FIELD_NUMBER: builtins.int
SUBSCRIBE_NOTIFICATION_CHANNEL_COMMAND_ATTRIBUTES_FIELD_NUMBER: builtins.int
UNSUBSCRIBE_NOTIFICATION_CHANNEL_COMMAND_ATTRIBUTES_FIELD_NUMBER: builtins.int
APPEND_STREAM_RECORDS_COMMAND_ATTRIBUTES_FIELD_NUMBER: builtins.int
SUBSCRIBE_STREAM_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 @@ -1219,6 +1222,14 @@ class Command(google.protobuf.message.Message):
def unsubscribe_notification_channel_command_attributes(
self,
) -> global___UnsubscribeNotificationChannelCommandAttributes: ...
@property
def append_stream_records_command_attributes(
self,
) -> global___AppendStreamRecordsCommandAttributes: ...
@property
def subscribe_stream_command_attributes(
self,
) -> global___SubscribeStreamCommandAttributes: ...
def __init__(
self,
*,
Expand Down Expand Up @@ -1267,10 +1278,16 @@ class Command(google.protobuf.message.Message):
| None = ...,
unsubscribe_notification_channel_command_attributes: global___UnsubscribeNotificationChannelCommandAttributes
| None = ...,
append_stream_records_command_attributes: global___AppendStreamRecordsCommandAttributes
| None = ...,
subscribe_stream_command_attributes: global___SubscribeStreamCommandAttributes
| 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 @@ -1307,6 +1324,8 @@ class Command(google.protobuf.message.Message):
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",
Expand All @@ -1318,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 @@ -1358,6 +1379,8 @@ class Command(google.protobuf.message.Message):
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",
Expand Down Expand Up @@ -1389,6 +1412,8 @@ class Command(google.protobuf.message.Message):
"request_cancel_nexus_operation_command_attributes",
"subscribe_notification_channel_command_attributes",
"unsubscribe_notification_channel_command_attributes",
"append_stream_records_command_attributes",
"subscribe_stream_command_attributes",
]
| None
): ...
Expand Down Expand Up @@ -1443,3 +1468,108 @@ class UnsubscribeNotificationChannelCommandAttributes(google.protobuf.message.Me
global___UnsubscribeNotificationChannelCommandAttributes = (
UnsubscribeNotificationChannelCommandAttributes
)

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
6 changes: 4 additions & 2 deletions temporalio/api/enums/v1/command_type_pb2.py

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

4 changes: 4 additions & 0 deletions temporalio/api/enums/v1/command_type_pb2.pyi
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,8 @@ class _CommandTypeEnumTypeWrapper(
COMMAND_TYPE_REQUEST_CANCEL_NEXUS_OPERATION: _CommandType.ValueType # 18
COMMAND_TYPE_SUBSCRIBE_NOTIFICATION_CHANNEL: _CommandType.ValueType # 19
COMMAND_TYPE_UNSUBSCRIBE_NOTIFICATION_CHANNEL: _CommandType.ValueType # 20
COMMAND_TYPE_APPEND_STREAM_RECORDS: _CommandType.ValueType # 21
COMMAND_TYPE_SUBSCRIBE_STREAM: _CommandType.ValueType # 22

class CommandType(_CommandType, metaclass=_CommandTypeEnumTypeWrapper):
"""Whenever this list of command types is changed do change the function shouldBufferEvent in mutableStateBuilder.go to make sure to do the correct event ordering."""
Expand All @@ -70,4 +72,6 @@ COMMAND_TYPE_SCHEDULE_NEXUS_OPERATION: CommandType.ValueType # 17
COMMAND_TYPE_REQUEST_CANCEL_NEXUS_OPERATION: CommandType.ValueType # 18
COMMAND_TYPE_SUBSCRIBE_NOTIFICATION_CHANNEL: CommandType.ValueType # 19
COMMAND_TYPE_UNSUBSCRIBE_NOTIFICATION_CHANNEL: CommandType.ValueType # 20
COMMAND_TYPE_APPEND_STREAM_RECORDS: CommandType.ValueType # 21
COMMAND_TYPE_SUBSCRIBE_STREAM: CommandType.ValueType # 22
global___CommandType = CommandType
Loading
Loading