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
2 changes: 1 addition & 1 deletion .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/moedash/sdk-rust.git
branch = moe/AI-198-ch-ext-core-2-machines
branch = moe/AI-198-if-ext-core-5-zero-cache-retention
10 changes: 10 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Empty file.
6 changes: 6 additions & 0 deletions temporalio/api/stream/v1/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
from .message_pb2 import StreamRecord, StreamRecordKind

__all__ = [
"StreamRecord",
"StreamRecordKind",
]
67 changes: 67 additions & 0 deletions temporalio/api/stream/v1/message_pb2.py

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

161 changes: 161 additions & 0 deletions temporalio/api/stream/v1/message_pb2.pyi
Original file line number Diff line number Diff line change
@@ -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
Loading