Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
43 commits
Select commit Hold shift + click to select a range
9343163
Pinned Core to the stream revision and regenerated bridge protos.
moedash Sep 14, 2026
b4a5c9a
Added the vendored stream API protos.
moedash Sep 14, 2026
ac3ff22
Added the stream client and the workflow stream surface.
moedash Sep 14, 2026
42589d8
Added the Workflow Streams surface over a server-side stream.
moedash Sep 14, 2026
303461b
Added the native provider module.
moedash Sep 15, 2026
372974d
Made cursors exclusive, added latest(), and allowed topic appends.
moedash Sep 16, 2026
bb21548
Read the delivered message shape the buffer now hands back.
moedash Sep 16, 2026
87f6006
Regenerated the protos from the pinned Core.
moedash Sep 16, 2026
d4fe1c8
Sorted and formatted the server-side stream modules.
moedash Sep 16, 2026
7a0eda3
Made the server-side stream modules pass the type and doc linters.
moedash Sep 16, 2026
facbf71
Added the changelog entry for server-side streams.
moedash Sep 16, 2026
014c37c
Regenerated the vendored stream protos for protobuf 3.
moedash Sep 16, 2026
96f8b6c
Repinned Core to the finalized replay-slice branch.
moedash Sep 17, 2026
aee2eb3
Qualified the LangGraph docstring's stream references.
moedash Sep 17, 2026
eb3f3f8
Repinned Core to the WIT-aligned replay-slice head.
moedash Sep 17, 2026
23138d9
Regenerated the bridge protos from the repinned Core.
moedash Sep 17, 2026
bf7039e
Adapted the native provider to the interface changes.
moedash Sep 19, 2026
07755d9
Pinned native stream handles to a run and closed their channels.
moedash Sep 19, 2026
d8c2efd
Stopped the flusher without cancelling it and pinned the contrib client.
moedash Sep 19, 2026
aa803f2
Checked stream ranges for continuity and mirrored the byte limits.
moedash Sep 19, 2026
a6eedbc
Corrected the subscribe docstring and two stale test comments.
moedash Sep 19, 2026
aa5f2c0
Repinned Core to the reviewed moe/AI-198-fix-replay-slice-lookahead h…
moedash Sep 19, 2026
200a5c6
Regenerated the bridge protos from the repinned Core.
moedash Sep 19, 2026
837404c
Repinned Core to the reviewed replay-slice head.
moedash Sep 21, 2026
6deaad6
Regenerated the bridge protos from the repinned Core.
moedash Sep 21, 2026
f734908
Regenerated the vendored stream service protos from the record server.
moedash Sep 21, 2026
0bbf54b
Moved the workflow stream runtime to records and batched a task's pub…
moedash Sep 21, 2026
3d22dc5
Moved the stream client and the contrib surface to records and SDK er…
moedash Sep 21, 2026
d3243e9
Rewrote the native provider with one owned stream per topic.
moedash Sep 21, 2026
45bd0ab
Opened the native tests' handles through client.get_stream_handle.
moedash Sep 21, 2026
d5e8f4b
Took topic definitions on the native handle.
moedash Sep 21, 2026
ebba767
Let a pushed history carry stream slices and repinned Core.
moedash Sep 21, 2026
4f9fa21
Tested a payload codec on both halves of the native provider.
moedash Sep 21, 2026
e5d1e20
Fetched a consuming workflow's stream records for the Replayer.
moedash Sep 21, 2026
dcab142
Repinned Core so a task owed records it was not sent fails first.
moedash Sep 21, 2026
5c46cec
Followed a workflow reset on the native handle.
moedash Sep 21, 2026
fe80330
Fetched each era of a reset run's history from the run that holds it.
moedash Sep 21, 2026
437be0c
Repinned Core so a sticky legacy query owed records goes unanswered.
moedash Sep 21, 2026
43504a8
Tested queries and a reset against the native provider.
moedash Sep 21, 2026
f44aa3d
Carried a workflow's stream records with its exported history.
moedash Sep 22, 2026
9038d8a
Rewrote the vendored stream stubs' type references onto temporalio.
moedash Sep 22, 2026
d665abf
Merged the interface branch's connect config fix and Core repin.
moedash Sep 22, 2026
5590dbc
Marked the reset event name as a literal in the eras docstring.
moedash Sep 22, 2026
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,3 +1,3 @@
[submodule "sdk-core"]
path = temporalio/bridge/sdk-core
url = https://github.com/temporalio/sdk-rust.git
url = https://github.com/moedash/sdk-rust.git
21 changes: 21 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,27 @@ to include examples, links to docs, or any other relevant information.
is `temporal.api.stream.v1.StreamRecord` on every provider.
`temporalio.streams.providers.memory.MemoryStreams` is the in-memory
reference provider the conformance tests run against.
- **Experimental**: server-side streams. A workflow publishes to a stream it
owns with a command the server applies in its Workflow Task's commit, and
reads the ranges the server delivers on its Workflow Tasks, through
`temporalio.workflow.append_stream_records`, `subscribe_stream` and
`read_stream_records`. `temporalio.client_stream` and
`temporalio.contrib.server_streams` reach the same stream from outside a
workflow, and `temporalio.streams.providers.native.NativeStreams` puts it
behind the shared stream interface with one owned stream per topic. Requires
a server that serves the stream service. `Replayer(stream_client=)` replays a
workflow that read such a stream while the server still holds it: History
records only the offsets each task consumed, so the replayer fetches the
records from the stream service and hands them to the replay with the
history. A range the stream no longer holds fails the replay with
`StreamNotFoundError`. A handle without a run id follows a workflow reset as
it follows a continue-as-new, reading the reset run from the floor its stream
reports, and the replayer fetches the ranges recorded before a reset point
from the run the workflow was reset from. For offline replay,
`Replayer.fetch_stream_slices(client, history)` attaches the records to a
`WorkflowHistory` while the stream is retained, `to_json()` and `from_json()`
carry them as `streamSlices` beside the events, and a history that carries
them replays with no server.

### Changed

Expand Down
193 changes: 193 additions & 0 deletions scripts/gen_stream_protos.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,193 @@
"""Regenerate the vendored server-side stream protos.

The stream service is still defined in the server rather than in the public
API repo, so its protos are copied in and compiled here. Run this whenever the
server-side definitions change:

uv run python scripts/gen_stream_protos.py /path/to/temporal

Run it with an interpreter whose ``grpcio-tools`` matches the oldest protobuf
runtime this package has to load on. Generated code refuses to load on a
runtime older than the one it was built against, and these modules end up in
applications that pin protobuf themselves. The stubs come from
``mypy-protobuf`` rather than protoc's own ``--pyi_out``, which that
generation of protoc does not have.

Two rewrites make the copies compile against this SDK. The server's own
routing and API-category options are dropped, because they mean nothing to a
client and their proto files are not vendored. Paths move under
``temporalio/`` so the generated modules import each other and the public API
messages by the names this package actually publishes.
"""

from __future__ import annotations

import os
import re
import shutil
import subprocess
import sys
import tempfile
from pathlib import Path

BASE = Path(__file__).parent.parent
OUT = BASE / "temporalio" / "api" / "streamservice" / "v1"
# The stream protos reference the public API messages. They normally come from
# the sdk-core submodule; TEMPORAL_API_PROTOS points at a `temporalio/api`
# checkout instead when the submodule is not initialized.
UPSTREAM_API = (
BASE
/ "temporalio"
/ "bridge"
/ "sdk-core"
/ "crates"
/ "protos"
/ "protos"
/ "api_upstream"
)


def api_proto_path() -> Path:
override = os.environ.get("TEMPORAL_API_PROTOS")
if override:
return Path(override).resolve()
if (
UPSTREAM_API / "temporal" / "api" / "common" / "v1" / "message.proto"
).is_file():
return UPSTREAM_API
raise SystemExit(
"no api protos found: initialize the sdk-core submodule or set "
"TEMPORAL_API_PROTOS to a temporalio/api checkout"
)


SERVER_PROTO_DIR = Path("chasm/lib/stream/proto/v1")
STAGE_PROTO_DIR = Path("temporalio/api/stream/v1")
FILES = [
"message.proto",
"stream_state.proto",
"request_response.proto",
"service.proto",
]

# Server-only options. They carry shard routing and rate-limit category, which
# a client neither reads nor can resolve, since the files defining them are not
# vendored.
DROP_IMPORT = re.compile(r'^import "temporal/server/.*?";\n', re.M)
DROP_OPTION = re.compile(r"^\s*option \(temporal\.server\.api\..*?\n", re.M)

# protoc names a python module after the proto path. These put the generated
# modules where this package publishes them.
FIX_PY = [
(
re.compile(r"from temporalio\.api\.stream\.v1 import"),
"from temporalio.api.streamservice.v1 import",
),
(re.compile(r"temporalio\.api\.stream\.v1\."), "temporalio.api.streamservice.v1."),
(re.compile(r"from temporal\.api\."), "from temporalio.api."),
(re.compile(r"import temporal\.api\."), "import temporalio.api."),
# mypy-protobuf also writes the public API types fully qualified in the
# stubs, `temporal.api.common.v1.message_pb2.Payload`, which the import
# rewrites above do not reach. The `_pb2` suffix keeps this off the proto
# package names inside the serialized descriptors, which stay `temporal.api`.
# After the streamservice rewrite on purpose, so the public
# `temporal.api.stream.v1` package is not sent to the vendored one.
(
re.compile(r"\btemporal\.api\.(\w+)\.v1\.(\w+_pb2)\b"),
r"temporalio.api.\1.v1.\2",
),
]


def stage(server: Path, into: Path) -> list[str]:
dest = into / STAGE_PROTO_DIR
dest.mkdir(parents=True)
for name in FILES:
text = (server / SERVER_PROTO_DIR / name).read_text()
text = text.replace(str(SERVER_PROTO_DIR) + "/", str(STAGE_PROTO_DIR) + "/")
text = DROP_IMPORT.sub("", text)
text = DROP_OPTION.sub("", text)
(dest / name).write_text(text)
return [str(STAGE_PROTO_DIR / name) for name in FILES]


def main() -> None:
if len(sys.argv) != 2:
sys.exit("usage: gen_stream_protos.py /path/to/temporal-server-checkout")
server = Path(sys.argv[1]).resolve()
if not (server / SERVER_PROTO_DIR).is_dir():
sys.exit(f"no stream protos under {server / SERVER_PROTO_DIR}")

with tempfile.TemporaryDirectory() as tmp:
work = Path(tmp)
targets = stage(server, work)
out = work / "out"
out.mkdir()
subprocess.run(
[
sys.executable,
"-m",
"grpc_tools.protoc",
f"--proto_path={work}",
f"--proto_path={api_proto_path()}",
f"--python_out={out}",
f"--grpc_python_out={out}",
f"--mypy_out={out}",
f"--mypy_grpc_out={out}",
*targets,
],
check=True,
cwd=work,
)

generated = out / STAGE_PROTO_DIR
for path in sorted(generated.iterdir()):
if path.name == "__init__.py":
continue
text = path.read_text()
for pattern, repl in FIX_PY:
text = pattern.sub(repl, text)
(OUT / path.name).write_text(text)
print(f"wrote {OUT / path.name}")

shutil.rmtree(OUT / "__pycache__", ignore_errors=True)
write_init()


MODULES = ["message_pb2", "stream_state_pb2", "request_response_pb2"]

INIT_HEADER = """\
# Generated by scripts/gen_stream_protos.py. Do not edit.
#
# Vendored rather than taken from the api submodule, because the stream service
# is still defined in the server. See temporalio/client_stream.py.
"""


def write_init() -> None:
"""Re-export every generated message and enum, so the package surface
cannot drift from the protos it was built from."""
import importlib

sys.path.insert(0, str(BASE))
blocks, exported = [], []
for module in MODULES:
mod = importlib.import_module(f"temporalio.api.streamservice.v1.{module}")
names = list(mod.DESCRIPTOR.message_types_by_name)
for enum in mod.DESCRIPTOR.enum_types_by_name.values():
names.append(enum.name)
names.extend(value.name for value in enum.values)
names.sort()
exported.extend(names)
joined = "".join(f" {name},\n" for name in names)
blocks.append(f"from .{module} import (\n{joined})\n")

listed = "".join(f' "{name}",\n' for name in sorted(exported))
(OUT / "__init__.py").write_text(
INIT_HEADER + "\n" + "\n".join(blocks) + f"\n__all__ = [\n{listed}]\n"
)
print(f"wrote {OUT / '__init__.py'}")


if __name__ == "__main__":
main()
Empty file.
146 changes: 146 additions & 0 deletions temporalio/api/streamservice/v1/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,146 @@
# Generated by scripts/gen_stream_protos.py. Do not edit.
#
# Vendored rather than taken from the api submodule, because the stream service
# is still defined in the server. See temporalio/client_stream.py.

from .message_pb2 import (
StreamRecord,
StreamRecordBatch,
)
from .request_response_pb2 import (
AddMessagesInput,
AddMessagesOutput,
AddMessagesRequest,
AddMessagesResponse,
AddWorkflowMessagesInput,
AddWorkflowMessagesRequest,
AddWorkflowMessagesResponse,
AdvanceConsumerHeadInput,
AdvanceConsumerHeadOutput,
AdvanceConsumerHeadRequest,
AdvanceConsumerHeadResponse,
CloseStreamInput,
CloseStreamOutput,
CloseStreamRequest,
CloseStreamResponse,
CreateStreamInput,
CreateStreamOutput,
CreateStreamRequest,
CreateStreamResponse,
DeleteStreamInput,
DeleteStreamOutput,
DeleteStreamRequest,
DeleteStreamResponse,
DescribeStreamInput,
DescribeStreamOutput,
DescribeStreamRequest,
DescribeStreamResponse,
DescribeWorkflowStreamInput,
DescribeWorkflowStreamRequest,
DescribeWorkflowStreamResponse,
FinishWritingInput,
FinishWritingOutput,
FinishWritingRequest,
FinishWritingResponse,
ListStreamsInput,
ListStreamsOutput,
ListStreamsRequest,
ListStreamsResponse,
PollMessagesInput,
PollMessagesOutput,
PollMessagesRequest,
PollMessagesResponse,
PollWorkflowMessagesInput,
PollWorkflowMessagesRequest,
PollWorkflowMessagesResponse,
RegisterStreamConsumerInput,
RegisterStreamConsumerOutput,
RegisterStreamConsumerRequest,
RegisterStreamConsumerResponse,
StreamListEntry,
SubscribeWorkflowInput,
SubscribeWorkflowOutput,
SubscribeWorkflowRequest,
SubscribeWorkflowResponse,
TruncateStreamInput,
TruncateStreamOutput,
TruncateStreamRequest,
TruncateStreamResponse,
)
from .stream_state_pb2 import (
ConsumerCursor,
ProducerCursor,
StreamBudget,
StreamLifecycle,
StreamState,
WorkflowStreamCursor,
)

__all__ = [
"AddMessagesInput",
"AddMessagesOutput",
"AddMessagesRequest",
"AddMessagesResponse",
"AddWorkflowMessagesInput",
"AddWorkflowMessagesRequest",
"AddWorkflowMessagesResponse",
"AdvanceConsumerHeadInput",
"AdvanceConsumerHeadOutput",
"AdvanceConsumerHeadRequest",
"AdvanceConsumerHeadResponse",
"CloseStreamInput",
"CloseStreamOutput",
"CloseStreamRequest",
"CloseStreamResponse",
"ConsumerCursor",
"CreateStreamInput",
"CreateStreamOutput",
"CreateStreamRequest",
"CreateStreamResponse",
"DeleteStreamInput",
"DeleteStreamOutput",
"DeleteStreamRequest",
"DeleteStreamResponse",
"DescribeStreamInput",
"DescribeStreamOutput",
"DescribeStreamRequest",
"DescribeStreamResponse",
"DescribeWorkflowStreamInput",
"DescribeWorkflowStreamRequest",
"DescribeWorkflowStreamResponse",
"FinishWritingInput",
"FinishWritingOutput",
"FinishWritingRequest",
"FinishWritingResponse",
"ListStreamsInput",
"ListStreamsOutput",
"ListStreamsRequest",
"ListStreamsResponse",
"PollMessagesInput",
"PollMessagesOutput",
"PollMessagesRequest",
"PollMessagesResponse",
"PollWorkflowMessagesInput",
"PollWorkflowMessagesRequest",
"PollWorkflowMessagesResponse",
"ProducerCursor",
"RegisterStreamConsumerInput",
"RegisterStreamConsumerOutput",
"RegisterStreamConsumerRequest",
"RegisterStreamConsumerResponse",
"StreamBudget",
"StreamLifecycle",
"StreamListEntry",
"StreamRecord",
"StreamRecordBatch",
"StreamState",
"SubscribeWorkflowInput",
"SubscribeWorkflowOutput",
"SubscribeWorkflowRequest",
"SubscribeWorkflowResponse",
"TruncateStreamInput",
"TruncateStreamOutput",
"TruncateStreamRequest",
"TruncateStreamResponse",
"WorkflowStreamCursor",
]
Loading
Loading