From e602e4152cd7b8ff913e5644c05dc2a3be93a408 Mon Sep 17 00:00:00 2001 From: Mohammad Dashti Date: Sat, 3 Oct 2026 03:03:09 -0700 Subject: [PATCH] Pinned Core at the external-stream repairs. This Core resumes reconstructed waits on wakes met during replay, stops stale waits from forcing an output task, and keeps a retained task with caching off. It also vendors the stream record envelope, which plain gen-protos now writes. --- .gitmodules | 2 +- CHANGELOG.md | 10 ++ temporalio/api/stream/__init__.py | 0 temporalio/api/stream/v1/__init__.py | 6 + temporalio/api/stream/v1/message_pb2.py | 67 ++++++++++ temporalio/api/stream/v1/message_pb2.pyi | 161 +++++++++++++++++++++++ temporalio/bridge/sdk-core | 2 +- 7 files changed, 246 insertions(+), 2 deletions(-) create mode 100644 temporalio/api/stream/__init__.py create mode 100644 temporalio/api/stream/v1/__init__.py create mode 100644 temporalio/api/stream/v1/message_pb2.py create mode 100644 temporalio/api/stream/v1/message_pb2.pyi diff --git a/.gitmodules b/.gitmodules index f3899a763..a7fdbacf9 100644 --- a/.gitmodules +++ b/.gitmodules @@ -1,4 +1,4 @@ [submodule "sdk-core"] path = temporalio/bridge/sdk-core url = https://github.com/moedash/sdk-rust.git - branch = moe/AI-198-ch-ext-core-2-machines + branch = moe/AI-198-if-ext-core-5-zero-cache-retention diff --git a/CHANGELOG.md b/CHANGELOG.md index fa3038107..4479e2b33 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -18,6 +18,16 @@ to include examples, links to docs, or any other relevant information. ## [Unreleased] +### Fixed + +- 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 diff --git a/temporalio/api/stream/__init__.py b/temporalio/api/stream/__init__.py new file mode 100644 index 000000000..e69de29bb diff --git a/temporalio/api/stream/v1/__init__.py b/temporalio/api/stream/v1/__init__.py new file mode 100644 index 000000000..fc6311caa --- /dev/null +++ b/temporalio/api/stream/v1/__init__.py @@ -0,0 +1,6 @@ +from .message_pb2 import StreamRecord, StreamRecordKind + +__all__ = [ + "StreamRecord", + "StreamRecordKind", +] diff --git a/temporalio/api/stream/v1/message_pb2.py b/temporalio/api/stream/v1/message_pb2.py new file mode 100644 index 000000000..7ea5be780 --- /dev/null +++ b/temporalio/api/stream/v1/message_pb2.py @@ -0,0 +1,67 @@ +# -*- coding: utf-8 -*- +# Generated by the protocol buffer compiler. DO NOT EDIT! +# source: temporal/api/stream/v1/message.proto +"""Generated protocol buffer code.""" + +from google.protobuf import descriptor as _descriptor +from google.protobuf import descriptor_pool as _descriptor_pool +from google.protobuf import message as _message +from google.protobuf import reflection as _reflection +from google.protobuf import symbol_database as _symbol_database +from google.protobuf.internal import enum_type_wrapper + +# @@protoc_insertion_point(imports) + +_sym_db = _symbol_database.Default() + + +from temporalio.api.common.v1 import ( + message_pb2 as temporal_dot_api_dot_common_dot_v1_dot_message__pb2, +) + +DESCRIPTOR = _descriptor_pool.Default().AddSerializedFile( + b'\n$temporal/api/stream/v1/message.proto\x12\x16temporal.api.stream.v1\x1a$temporal/api/common/v1/message.proto"\xd4\x02\n\x0cStreamRecord\x12-\n\x04\x62ody\x18\x01 \x01(\x0b\x32\x1f.temporal.api.common.v1.Payload\x12\x44\n\x08metadata\x18\x02 \x03(\x0b\x32\x32.temporal.api.stream.v1.StreamRecord.MetadataEntry\x12\r\n\x05topic\x18\x03 \x01(\t\x12\x36\n\x04kind\x18\x04 \x01(\x0e\x32(.temporal.api.stream.v1.StreamRecordKind\x12\x13\n\x0bproducer_id\x18\x05 \x01(\t\x12\x0f\n\x07\x61ttempt\x18\x06 \x01(\x03\x12\x10\n\x08sequence\x18\x07 \x01(\x03\x1aP\n\rMetadataEntry\x12\x0b\n\x03key\x18\x01 \x01(\t\x12.\n\x05value\x18\x02 \x01(\x0b\x32\x1f.temporal.api.common.v1.Payload:\x02\x38\x01*r\n\x10StreamRecordKind\x12"\n\x1eSTREAM_RECORD_KIND_UNSPECIFIED\x10\x00\x12\x1b\n\x17STREAM_RECORD_KIND_DATA\x10\x01\x12\x1d\n\x19STREAM_RECORD_KIND_FINISH\x10\x02\x42\x89\x01\n\x19io.temporal.api.stream.v1B\x0cMessageProtoP\x01Z#go.temporal.io/api/stream/v1;stream\xaa\x02\x18Temporalio.Api.Stream.V1\xea\x02\x1bTemporalio::Api::Stream::V1b\x06proto3' +) + +_STREAMRECORDKIND = DESCRIPTOR.enum_types_by_name["StreamRecordKind"] +StreamRecordKind = enum_type_wrapper.EnumTypeWrapper(_STREAMRECORDKIND) +STREAM_RECORD_KIND_UNSPECIFIED = 0 +STREAM_RECORD_KIND_DATA = 1 +STREAM_RECORD_KIND_FINISH = 2 + + +_STREAMRECORD = DESCRIPTOR.message_types_by_name["StreamRecord"] +_STREAMRECORD_METADATAENTRY = _STREAMRECORD.nested_types_by_name["MetadataEntry"] +StreamRecord = _reflection.GeneratedProtocolMessageType( + "StreamRecord", + (_message.Message,), + { + "MetadataEntry": _reflection.GeneratedProtocolMessageType( + "MetadataEntry", + (_message.Message,), + { + "DESCRIPTOR": _STREAMRECORD_METADATAENTRY, + "__module__": "temporalio.api.stream.v1.message_pb2", + # @@protoc_insertion_point(class_scope:temporal.api.stream.v1.StreamRecord.MetadataEntry) + }, + ), + "DESCRIPTOR": _STREAMRECORD, + "__module__": "temporalio.api.stream.v1.message_pb2", + # @@protoc_insertion_point(class_scope:temporal.api.stream.v1.StreamRecord) + }, +) +_sym_db.RegisterMessage(StreamRecord) +_sym_db.RegisterMessage(StreamRecord.MetadataEntry) + +if _descriptor._USE_C_DESCRIPTORS == False: + DESCRIPTOR._options = None + DESCRIPTOR._serialized_options = b"\n\031io.temporal.api.stream.v1B\014MessageProtoP\001Z#go.temporal.io/api/stream/v1;stream\252\002\030Temporalio.Api.Stream.V1\352\002\033Temporalio::Api::Stream::V1" + _STREAMRECORD_METADATAENTRY._options = None + _STREAMRECORD_METADATAENTRY._serialized_options = b"8\001" + _STREAMRECORDKIND._serialized_start = 445 + _STREAMRECORDKIND._serialized_end = 559 + _STREAMRECORD._serialized_start = 103 + _STREAMRECORD._serialized_end = 443 + _STREAMRECORD_METADATAENTRY._serialized_start = 363 + _STREAMRECORD_METADATAENTRY._serialized_end = 443 +# @@protoc_insertion_point(module_scope) diff --git a/temporalio/api/stream/v1/message_pb2.pyi b/temporalio/api/stream/v1/message_pb2.pyi new file mode 100644 index 000000000..a553b3bf2 --- /dev/null +++ b/temporalio/api/stream/v1/message_pb2.pyi @@ -0,0 +1,161 @@ +""" +@generated by mypy-protobuf. Do not edit manually! +isort:skip_file +""" + +import builtins +import collections.abc +import sys +import typing + +import google.protobuf.descriptor +import google.protobuf.internal.containers +import google.protobuf.internal.enum_type_wrapper +import google.protobuf.message + +import temporalio.api.common.v1.message_pb2 + +if sys.version_info >= (3, 10): + import typing as typing_extensions +else: + import typing_extensions + +DESCRIPTOR: google.protobuf.descriptor.FileDescriptor + +class _StreamRecordKind: + ValueType = typing.NewType("ValueType", builtins.int) + V: typing_extensions.TypeAlias = ValueType + +class _StreamRecordKindEnumTypeWrapper( + google.protobuf.internal.enum_type_wrapper._EnumTypeWrapper[ + _StreamRecordKind.ValueType + ], + builtins.type, +): # noqa: F821 + DESCRIPTOR: google.protobuf.descriptor.EnumDescriptor + STREAM_RECORD_KIND_UNSPECIFIED: _StreamRecordKind.ValueType # 0 + """Read as DATA.""" + STREAM_RECORD_KIND_DATA: _StreamRecordKind.ValueType # 1 + """A value the producer published; `body` carries it.""" + STREAM_RECORD_KIND_FINISH: _StreamRecordKind.ValueType # 2 + """The producer named by `producer_id` writes nothing more on `topic`. + Says nothing about that producer's outcome and does not end the stream. + """ + +class StreamRecordKind(_StreamRecordKind, metaclass=_StreamRecordKindEnumTypeWrapper): + """What a record means to a reader. Kept on the record itself so every store + and every language reads it the same way without a private envelope. + """ + +STREAM_RECORD_KIND_UNSPECIFIED: StreamRecordKind.ValueType # 0 +"""Read as DATA.""" +STREAM_RECORD_KIND_DATA: StreamRecordKind.ValueType # 1 +"""A value the producer published; `body` carries it.""" +STREAM_RECORD_KIND_FINISH: StreamRecordKind.ValueType # 2 +"""The producer named by `producer_id` writes nothing more on `topic`. +Says nothing about that producer's outcome and does not end the stream. +""" +global___StreamRecordKind = StreamRecordKind + +class StreamRecord(google.protobuf.message.Message): + """One entry in a stream. The record is the wire format: stores keep it + serialized as is and readers in every language decode the same bytes. + """ + + DESCRIPTOR: google.protobuf.descriptor.Descriptor + + class MetadataEntry(google.protobuf.message.Message): + DESCRIPTOR: google.protobuf.descriptor.Descriptor + + KEY_FIELD_NUMBER: builtins.int + VALUE_FIELD_NUMBER: builtins.int + key: builtins.str + @property + def value(self) -> temporalio.api.common.v1.message_pb2.Payload: ... + def __init__( + self, + *, + key: builtins.str = ..., + value: temporalio.api.common.v1.message_pb2.Payload | None = ..., + ) -> None: ... + def HasField( + self, field_name: typing_extensions.Literal["value", b"value"] + ) -> builtins.bool: ... + def ClearField( + self, + field_name: typing_extensions.Literal["key", b"key", "value", b"value"], + ) -> None: ... + + BODY_FIELD_NUMBER: builtins.int + METADATA_FIELD_NUMBER: builtins.int + TOPIC_FIELD_NUMBER: builtins.int + KIND_FIELD_NUMBER: builtins.int + PRODUCER_ID_FIELD_NUMBER: builtins.int + ATTEMPT_FIELD_NUMBER: builtins.int + SEQUENCE_FIELD_NUMBER: builtins.int + @property + def body(self) -> temporalio.api.common.v1.message_pb2.Payload: + """The value the producer published, stored as sent. A stream provider + writes and reads it as part of the record, so applying a payload codec + to it is the provider's job. + """ + @property + def metadata( + self, + ) -> google.protobuf.internal.containers.MessageMap[ + builtins.str, temporalio.api.common.v1.message_pb2.Payload + ]: + """Producer-supplied provenance, stored as sent.""" + topic: builtins.str + """Producer-supplied grouping label, stored as sent.""" + kind: global___StreamRecordKind.ValueType + """How to read this record. Unspecified is read as DATA.""" + producer_id: builtins.str + """Who wrote the record. Empty when the owning Workflow did.""" + attempt: builtins.int + """The producer's attempt. Readers treat a later attempt by the same + producer as superseding what the earlier one wrote. + """ + sequence: builtins.int + """The producer's position within its attempt, zero when it does not number + its records. Stored as sent; the server does not assign, validate or + order by it, and the stream's own offsets are what order a read. + """ + def __init__( + self, + *, + body: temporalio.api.common.v1.message_pb2.Payload | None = ..., + metadata: collections.abc.Mapping[ + builtins.str, temporalio.api.common.v1.message_pb2.Payload + ] + | None = ..., + topic: builtins.str = ..., + kind: global___StreamRecordKind.ValueType = ..., + producer_id: builtins.str = ..., + attempt: builtins.int = ..., + sequence: builtins.int = ..., + ) -> None: ... + def HasField( + self, field_name: typing_extensions.Literal["body", b"body"] + ) -> builtins.bool: ... + def ClearField( + self, + field_name: typing_extensions.Literal[ + "attempt", + b"attempt", + "body", + b"body", + "kind", + b"kind", + "metadata", + b"metadata", + "producer_id", + b"producer_id", + "sequence", + b"sequence", + "topic", + b"topic", + ], + ) -> None: ... + +global___StreamRecord = StreamRecord diff --git a/temporalio/bridge/sdk-core b/temporalio/bridge/sdk-core index b5fe6f798..5a0e3e008 160000 --- a/temporalio/bridge/sdk-core +++ b/temporalio/bridge/sdk-core @@ -1 +1 @@ -Subproject commit b5fe6f798d2883f155ee6bf457aabcccc2464305 +Subproject commit 5a0e3e008f5d159c1106f0ef25bdcb6466ca14e4