Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
45 commits
Select commit Hold shift + click to select a range
95438f3
Pinned Core to the stream head and regenerated the bridge protos.
moedash Sep 25, 2026
1774972
Added the vendored stream service protos and their generator.
moedash Sep 25, 2026
a34074f
Added the stream service client.
moedash Sep 25, 2026
b09e3f1
Named a producer conflict as a producer error.
moedash Sep 25, 2026
c18e031
Copied a stream record by descriptor in both directions.
moedash Sep 25, 2026
ced05ac
Freed zero to mean a producer that does not number.
moedash Sep 25, 2026
1ebd3c1
Repinned Core to the renamed stream command fields.
moedash Sep 25, 2026
6dbda10
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-py-05-nativ…
moedash Sep 26, 2026
c3e1f58
Built the bridge against the 1.0 Core crates.
moedash Sep 26, 2026
397ac1e
Regenerated the protos from the merged Core.
moedash Sep 26, 2026
dab098b
Refused the new callback link variant like the batch job one.
moedash Sep 26, 2026
491193d
Split the unset-kind default into its own gated stream test.
moedash Sep 26, 2026
3117bfd
Dropped the server-age gate from the unset-kind stream test.
moedash Sep 26, 2026
7947fc6
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-py-05-nativ…
moedash Sep 28, 2026
4428e7a
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-py-05-nativ…
moedash Sep 28, 2026
52edf6d
Regenerated the vendored stream protos with the owner reference.
moedash Sep 28, 2026
9addbc6
Opened activity-owned streams from the wire client.
moedash Sep 28, 2026
f6ba1b1
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-py-05-nativ…
moedash Sep 28, 2026
ad73179
Repinned Core to the subscribe start position and regenerated its pro…
moedash Sep 28, 2026
cf13959
Let a first read on the stream client name a start position.
moedash Sep 29, 2026
515b591
Regenerated the Nexus system API for the repinned Core WIT.
moedash Sep 29, 2026
fb38837
Merge remote-tracking branch 'origin/moe/AI-198-py-04-accessors' into…
moedash Sep 30, 2026
9d26e5c
Merge remote-tracking branch 'origin/moe/AI-198-py-04-accessors' into…
moedash Sep 30, 2026
7868260
Regenerated the vendored stream service protos for the byte cap.
moedash Oct 1, 2026
41b05d8
Repinned Core for the wake protos.
moedash Oct 1, 2026
17ef4dd
Merge remote-tracking branch 'origin/moe/AI-198-py-04-accessors' into…
moedash Oct 1, 2026
19216f9
Repinned Core to the delivery head with the C bridge wake dispatch.
moedash Oct 1, 2026
1b4b23f
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-py-05-nativ…
moedash Oct 2, 2026
7e427c5
Pinned Core at the channel-round delivery head and regenerated the cl…
moedash Oct 2, 2026
89eafa3
Mapped the producer conflict under its failed-precondition status.
moedash Oct 2, 2026
dc3cab3
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-py-05-nativ…
moedash Oct 2, 2026
09f28f9
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-py-05-nativ…
moedash Oct 2, 2026
2a18a2d
Pinned Core at the linked-channel head and regenerated the clients.
moedash Oct 2, 2026
de9228a
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-py-05-nativ…
moedash Oct 2, 2026
394faeb
Flipped the linked-channel skip on the delivery Core pin.
moedash Oct 2, 2026
8378092
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-py-05-nativ…
moedash Oct 2, 2026
7ee4c7b
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-py-05-nativ…
moedash Oct 2, 2026
ce77f32
Pinned Core at the unsubscribe head and regenerated the clients.
moedash Oct 2, 2026
4fcf9aa
Flipped the unsubscribe skip on the delivery Core pin.
moedash Oct 2, 2026
da557f4
Merge branch 'moe/AI-198-py-04-accessors' into moe/AI-198-py-05-nativ…
moedash Oct 2, 2026
ecaa55d
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Oct 3, 2026
7516568
Pinned the delivery Core with a linked channel addressed by execution.
moedash Oct 3, 2026
94aa0a8
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Oct 3, 2026
7657155
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Oct 3, 2026
38aa617
Pinned the delivery Core with the execution fields on their original …
moedash Oct 3, 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
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()
4 changes: 2 additions & 2 deletions temporalio/api/callback/v1/message_pb2.py

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

8 changes: 8 additions & 0 deletions temporalio/api/callback/v1/message_pb2.pyi
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@ class CallbackInfo(google.protobuf.message.Message):
LAST_ATTEMPT_FAILURE_FIELD_NUMBER: builtins.int
NEXT_ATTEMPT_SCHEDULE_TIME_FIELD_NUMBER: builtins.int
BLOCKED_REASON_FIELD_NUMBER: builtins.int
REQUEST_ID_FIELD_NUMBER: builtins.int
@property
def callback(self) -> temporalio.api.common.v1.message_pb2.Callback:
"""Information on how this callback should be invoked (e.g. its URL and type)."""
Expand All @@ -57,6 +58,10 @@ class CallbackInfo(google.protobuf.message.Message):
"""The time when the next attempt is scheduled."""
blocked_reason: builtins.str
"""If the state is BLOCKED, blocked reason provides additional information."""
request_id: builtins.str
"""Server-generated request ID used as an idempotency token when invoking callbacks.
It has no relation to caller-side request_id sent in operations like StartNexusOperationExecutionRequest.
"""
def __init__(
self,
*,
Expand All @@ -71,6 +76,7 @@ class CallbackInfo(google.protobuf.message.Message):
next_attempt_schedule_time: google.protobuf.timestamp_pb2.Timestamp
| None = ...,
blocked_reason: builtins.str = ...,
request_id: builtins.str = ...,
) -> None: ...
def HasField(
self,
Expand Down Expand Up @@ -104,6 +110,8 @@ class CallbackInfo(google.protobuf.message.Message):
b"next_attempt_schedule_time",
"registration_time",
b"registration_time",
"request_id",
b"request_id",
"state",
b"state",
],
Expand Down
110 changes: 68 additions & 42 deletions temporalio/api/common/v1/message_pb2.py

Large diffs are not rendered by default.

Loading
Loading