Skip to content

Pinned Core to the stream head and added the stream client. - #8

Closed
moedash wants to merge 45 commits into
moe/AI-198-streams-provider-workflow-streamsfrom
moe/AI-198-py-05-native-wire
Closed

moedash wants to merge 45 commits into
moe/AI-198-streams-provider-workflow-streamsfrom
moe/AI-198-py-05-native-wire

Conversation

@moedash

@moedash moedash commented Sep 25, 2026 •

Copy link
Copy Markdown
Owner

This PR pins Core to the stream head and adds a client for the server's stream service.

What changed?

  • temporalio/bridge/sdk-core pins the moe/AI-198-core-3-delivery head. It carries the stream activation job, the subscribe_stream and append_stream_records commands, the subscribe start position, and the notification channel: the SubscribeNotificationChannel and UnsubscribeNotificationChannel command machines and the NotificationsReceived job folded from the task's scheduled events, with each notification's linked_to carried through as the server sent it, a common.v1.Execution naming a workflow or a standalone activity. The live workflow channel cases from Added the workflow-side stream reader and writer. #12, the unsubscribe one included, run from this layer, since the Core below it only carries the protos.
  • The bridge protos, temporalio/bridge/_visitor.py and the Nexus system API are regenerated from that pin. The payload visitor reaches each record's body and metadata.
  • temporalio/api/streamservice/v1/ holds the vendored stream service protos. scripts/gen_stream_protos.py builds them from a server checkout with a protobuf 3 protoc, and mypy-protobuf writes the stubs. The script rewrites the stubs' temporal.api.* type references onto temporalio.
  • temporalio/client_stream.py holds StreamClient, StreamHandle and WorkflowStreamHandle over describe, append and poll. A failure surfaces as StreamNotFoundError, StreamProducerError or RPCError, never as a grpc type. The producer conflict is recognised under the FAILED_PRECONDITION the server sends with its reason token, and under the INVALID_ARGUMENT a server built before the tokens sent.
  • A record crosses between StreamRecord and the stored shape by descriptor. A field with no counterpart on the other side raises.
  • StreamClient.activity_stream() opens a stream an activity owns. A first read takes start=, and a standalone handle has truncate().

This layer wires nothing into a workflow or a provider. #9 adds the provider and #14 adds replay.

Part of AI-198 (epic AI-37).

Why?

sdk-core doesn't know the stream service, so the client opens its own grpcio channel. On this layer that channel is plaintext, and #9 builds it from the client's connection. The pin, the generated code and the client sit together so the next two layers are Python a reviewer can read without a rebuild.

How did you test it?

uv run poe build-develop, then ruff check --select I, ruff format --check, mypy, pyright, basedpyright and pydocstyle on a cold .mypy_cache, where the stub rewrite bug shows. tests/test_client_stream_records.py and tests/test_client_stream_owner.py run without a server. tests/test_client_stream.py runs live with TEMPORAL_STREAM_TARGET on a server from the moedash/temporal chain. It covers append, poll, describe, the deduplicated repeat, the long poll, both producer refusals and start positions on a truncated stream. The channel cases from #12, the workflow one included, passed against that server with -E host:port. TLS and API keys come with #9.

  • Unit Tests
  • Staging
  • End to End Tests

The pinned Core carries the stream activation job and the two stream commands, so the generated bridge modules and the payload visitor have to match it.
The stream service is still defined in the server rather than in the api submodule, so its generated code is vendored here; the generator rewrites imports and stubs onto temporalio so a cold mypy run resolves them.
Outside code needs a way to append to and read a server-side stream. The client translates every transport failure into StreamNotFoundError or RPCError and applies the client's payload codec to bodies, so both sides of a namespace with a codec agree.
The server refuses a producer sequence it already holds, which is exactly what StreamProducerError is documented for; it reached the caller as a bare argument failure instead.
Six hand-copied fields each way meant a field added to StreamRecord would be dropped in silence on the way out and on the way back.
The wire no longer carries a sentinel for an unnumbered record, so the
memory producer starts its own numbering at one and leaves zero to the
producers that never set it.
The pinned Core names the stream command fields after what they carry, so the generated bridge modules have to match it.
…e-wire

# Conflicts:
#	temporalio/bridge/sdk-core
The merged Core renumbered its crates and moved the payload warning options behind the client's experimental feature.
The upstream Core merge carried newer API protos and the lang proto changes that go with them.
@moedash
moedash changed the base branch from moe/AI-198-streams-interface to moe/AI-198-py-04-accessors September 26, 2026 03:16
Every native live module needs a server with the record kind rule, so one gate for one case only hid the requirement.
The server's owned-stream calls take one owner reference for a workflow, a standalone activity, or a workflow's activity.
A workflow owner keeps going out in the workflow fields, so a server without owner support still routes it; only an activity sends the owner reference.
…e-wire

# Conflicts:
#	temporalio/bridge/sdk-core
…tos.

The bridge's subscribe command and the server's poll and subscribe inputs carry a start position, so the lang side can ask for the earliest record, the tail or the last N.
A reader with no offset had to send zero, which a truncated stream refuses. The first poll now carries the position and later polls continue from the offset it returns; truncate() is there so a caller and the tests can reach a truncated stream.
The Core repin at this layer updated the nexus WIT doc strings and marked user metadata experimental, but only the protobuf output was regenerated. The check-protos job regenerates with nexgen 0.2.4 and saw these two files drift.
The server's StreamLifecycle gained max_bytes beside max_items, and
StreamState reports held_bytes, so the standalone policy can bound bytes.
Generated with the same protobuf 3 toolchain as before, so the modules still
load on the oldest runtime this package supports.
moedash added 18 commits October 1, 2026 11:37
Core moves to the delivery branch head that carries the Wake message, the
wakes field on the workflow task poll response and the WakeWorkflowExecution
call, so the vendored tree and the bridge client regenerate from the pin.
… moe/AI-198-py-05-native-wire

# Conflicts:
#	temporalio/api/workflow/v1/message_pb2.py
#	temporalio/api/workflowservice/v1/request_response_pb2.py
The move only touches the C bridge, so the vendored protos and the bridge client regenerate unchanged.
…ients.

Core moves to the core-3 commit that carries the subscribe-notification-channel command machine and the NotificationsReceived job folded from the task's scheduled events. The vendored tree and the payload visitor regenerate from the pin, and the live workflow channel case now runs on this layer.
The server reports a producer conflict as FAILED_PRECONDITION with a reason token, and this layer only knew the INVALID_ARGUMENT a server built before the tokens sent, so the typed error was lost on the current server.
…e-wire

# Conflicts:
#	temporalio/api/workflow/v1/message_pb2.py
#	temporalio/api/workflowservice/v1/request_response_pb2.py
#	temporalio/bridge/sdk-core
…e-wire

# Conflicts:
#	temporalio/api/enums/v1/failed_cause_pb2.py
#	temporalio/api/history/v1/message_pb2.py
#	temporalio/api/workflow/v1/message_pb2.py
#	temporalio/api/workflowservice/v1/request_response_pb2.py
#	temporalio/bridge/proto/workflow_commands/__init__.py
#	temporalio/bridge/proto/workflow_commands/workflow_commands_pb2.py
#	temporalio/bridge/proto/workflow_commands/workflow_commands_pb2.pyi
#	temporalio/bridge/sdk-core
The delivery Core carries the unsubscribe command machine, so a workflow on this layer can end a subscription and replay through its event.
…AI-198-py-05-native-wire

# Conflicts:
#	temporalio/api/workflowservice/v1/request_response_pb2.py
#	temporalio/bridge/sdk-core
The core-3 head carries the api change: the five channel requests and both
`linked_to` fields take `common.v1.Execution`, so a notification names a
standalone activity owner as well as a workflow.
@moedash
moedash changed the base branch from moe/AI-198-py-04-accessors to moe/AI-198-streams-provider-workflow-streams October 3, 2026 01:07
…AI-198-py-05-native-wire

# Conflicts:
#	temporalio/api/workflowservice/v1/request_response_pb2.py
#	temporalio/bridge/sdk-core
…numbers.

The core-3 head carries the renumbered api: the five channel requests and
both `linked_to` fields sit on the numbers the workflow fields had, with
nothing reserved. Names and types are unchanged.
@moedash

moedash commented Oct 3, 2026

Copy link
Copy Markdown
Owner Author

Replaced by #19, #20, #21, #22, #23, #24, #25, #26, #27, #28, #29, #30, #31, #32, #33, #34, #35, #36, #37, #38, #39, #40, #41, #42, #43, #44, #45, #46, #47, #48, #49, #50, #51, #52, #53, #54 and #55.

Same content, split into 37 PRs in the v3 series: the notification channel first, then the streaming interface, then native streams and the rest. The branch stays as a pin.

@moedash moedash closed this Oct 3, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

skip-changelog Changelog entry rides another PR

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant