Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
84 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
ca9c2e4
Let a pushed history carry stream slices.
moedash Sep 25, 2026
3ebd9f5
Fetched a consuming workflow's stream records for the Replayer.
moedash Sep 25, 2026
a4be6ac
Followed a workflow reset on the handle and in the Replayer.
moedash Sep 25, 2026
2c5ce1b
Carried a workflow's stream records with its exported history.
moedash Sep 25, 2026
207fbd3
Tested replay, queries, a reset and a codec on the native provider.
moedash Sep 25, 2026
7280f8d
Added the changelog entry for server-side streams.
moedash Sep 25, 2026
46deb20
Tied the series to the original branch head.
moedash Sep 25, 2026
451002d
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Sep 25, 2026
c927584
Named the replayer's stream client limit and closed its channel.
moedash Sep 25, 2026
d77a597
Stopped parsing an exported history twice to look for slices.
moedash Sep 25, 2026
27cad8d
Followed a chain of two resets to its end.
moedash Sep 25, 2026
693e1f2
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Sep 25, 2026
2bd2273
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Sep 25, 2026
b396523
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Sep 26, 2026
9cf995c
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Sep 28, 2026
345c64d
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Sep 28, 2026
420bbd8
Matched the activity handle to the replay layer's successor lookup.
moedash Sep 28, 2026
73e6585
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Sep 29, 2026
0286204
Tested a pinned read of a reset run from BEGINNING.
moedash Sep 29, 2026
fc6babf
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Sep 29, 2026
351950b
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Sep 29, 2026
68f3f6b
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Sep 30, 2026
9cdb227
Derived the replayer's stream channel from its client.
moedash Sep 30, 2026
af134ca
Filled a task's short stream range from the stream service before it …
moedash Sep 30, 2026
d2b2de7
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Sep 30, 2026
daabc87
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Sep 30, 2026
a1676e8
Expected the reset run to read the reset-point task's input again.
moedash Sep 30, 2026
98622d6
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Sep 30, 2026
6621c5a
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Oct 1, 2026
a81a8c1
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Oct 1, 2026
176ad2a
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Oct 1, 2026
6979d7a
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Oct 1, 2026
a254ad9
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Oct 1, 2026
781a4b1
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Oct 2, 2026
41b898d
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Oct 2, 2026
4448201
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Oct 2, 2026
83820c3
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Oct 2, 2026
3352512
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Oct 2, 2026
25e36fc
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Oct 2, 2026
c2aba51
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
moedash Oct 3, 2026
e473391
Merge branch 'moe/AI-198-py-06-native-provider' into moe/AI-198-py-07…
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
21 changes: 21 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,27 @@ to include examples, links to docs, or any other relevant information.
converter, so a payload codec and external storage apply to them.
`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.
- **Experimental**: `temporalio.streams.providers.workflow_streams.WorkflowStreamsProvider`
serves the stream interface over the shipped Workflow Streams transport as a
worker plugin, so a workflow reads and publishes through
Expand Down
13 changes: 12 additions & 1 deletion temporalio/bridge/src/worker.rs
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ use temporalio_common::protos::coresdk::{
nexus::NexusTaskCompletion, ActivityHeartbeat, ActivityTaskCompletion,
};
use temporalio_common::protos::temporal::api::history::v1::History;
use temporalio_common::protos::temporal::api::stream::v1::StreamSlice;
use temporalio_common::protos::temporal::api::worker::v1::{PluginInfo, StorageDriverInfo};
use temporalio_sdk_core::replay::{HistoryForReplay, ReplayWorkerInput};
use temporalio_sdk_core::{
Expand Down Expand Up @@ -950,14 +951,24 @@ impl HistoryPusher {

#[pymethods]
impl HistoryPusher {
/// Feed one history to the replay worker. `stream_slices` are serialized
/// `temporal.api.stream.v1.StreamSlice` messages carrying the records the
/// history's completed tasks consumed, which History itself never holds.
#[pyo3(signature = (workflow_id, history_proto, stream_slices = Vec::new()))]
fn push_history<'p>(
&self,
py: Python<'p>,
workflow_id: &str,
history_proto: &Bound<'_, PyBytes>,
stream_slices: Vec<Bound<'p, PyBytes>>,
) -> PyResult<Bound<'p, PyAny>> {
let history = History::decode(history_proto.as_bytes())
.map_err(|err| PyValueError::new_err(format!("Invalid proto: {err}")))?;
let slices = stream_slices
.iter()
.map(|slice| StreamSlice::decode(slice.as_bytes()))
.collect::<Result<Vec<_>, _>>()
.map_err(|err| PyValueError::new_err(format!("Invalid stream slice proto: {err}")))?;
let wfid = workflow_id.to_string();
let tx = if let Some(tx) = self.tx.as_ref() {
tx.clone()
Expand All @@ -968,7 +979,7 @@ impl HistoryPusher {
};
// We accept this doesn't have logging/tracing
self.runtime.future_into_py(py, async move {
tx.send(HistoryForReplay::new(history, wfid))
tx.send(HistoryForReplay::new(history, wfid).with_stream_slices(slices))
.await
.map_err(|_| {
PyRuntimeError::new_err(
Expand Down
63 changes: 55 additions & 8 deletions temporalio/client/_workflow.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@

import asyncio
import functools
import json
import warnings
from asyncio import Future
from collections.abc import (
Expand Down Expand Up @@ -31,6 +32,7 @@
import temporalio.api.common.v1
import temporalio.api.enums.v1
import temporalio.api.history.v1
import temporalio.api.stream.v1
import temporalio.api.update.v1
import temporalio.api.workflow.v1
import temporalio.api.workflowservice.v1
Expand Down Expand Up @@ -1698,6 +1700,20 @@ class WorkflowHistory:
events: Sequence[temporalio.api.history.v1.HistoryEvent]
"""History events for the workflow."""

stream_slices: Sequence[temporalio.api.stream.v1.StreamSlice] = ()
"""The stream records the workflow's completed tasks consumed, if captured.

History records only the offsets each Workflow Task consumed from a
server-side stream; the records live in the stream. A history whose tasks
consumed records and that carries none here cannot be replayed without the
server that still holds the stream. One slice per recorded range, tagged
with the ``WorkflowTaskCompleted`` event that recorded it, in the shape the
server puts on a poll response.
:py:meth:`temporalio.worker.Replayer.fetch_stream_slices` fills it while
the stream is retained, and :py:meth:`to_json` and :py:meth:`from_json`
carry it, so an exported history file is the whole replay input.
"""

@property
def run_id(self) -> str:
"""Run ID extracted from the first event."""
Expand All @@ -1719,31 +1735,62 @@ def from_json(workflow_id: str, history: str | dict[str, Any]) -> WorkflowHistor
Args:
workflow_id: The workflow's ID
history: A string or parsed-to-dict representation of workflow
history
history. A ``streamSlices`` list beside ``events``, as
:py:meth:`to_json` writes one, is read into
:py:attr:`stream_slices`.

Returns:
Workflow history
"""
parsed = _history_from_json(history)
return WorkflowHistory(workflow_id, parsed.events)
raw: Any = []
if isinstance(history, dict):
raw = history.get("streamSlices") or history.get("stream_slices") or []
elif '"streamSlices"' in history or '"stream_slices"' in history:
# Read again only when the text mentions them. Almost every
# history has none, and handing the parsed dict to
# _history_from_json instead would make it deep-copy an export
# that can be very large.
decoded = json.loads(history)
if isinstance(decoded, dict):
raw = decoded.get("streamSlices") or decoded.get("stream_slices") or []
slices: list[temporalio.api.stream.v1.StreamSlice] = [
google.protobuf.json_format.ParseDict(
entry,
temporalio.api.stream.v1.StreamSlice(),
ignore_unknown_fields=True,
)
for entry in raw
]
return WorkflowHistory(workflow_id, parsed.events, slices)

def to_json(self) -> str:
"""Convert this history to JSON.

Note, this does not include the workflow ID.
Note, this does not include the workflow ID. The stream slices, when
there are any, are written as a ``streamSlices`` list beside the
events; without them the output is the history proto's JSON alone.
"""
return google.protobuf.json_format.MessageToJson(
temporalio.api.history.v1.History(events=self.events)
)
if not self.stream_slices:
return google.protobuf.json_format.MessageToJson(
temporalio.api.history.v1.History(events=self.events)
)
return json.dumps(self.to_json_dict(), indent=2)

def to_json_dict(self) -> dict[str, Any]:
"""Convert this history to JSON-compatible dict.

Note, this does not include the workflow ID.
Note, this does not include the workflow ID. See :py:meth:`to_json`
for how the stream slices are carried.
"""
return google.protobuf.json_format.MessageToDict(
out = google.protobuf.json_format.MessageToDict(
temporalio.api.history.v1.History(events=self.events)
)
if self.stream_slices:
out["streamSlices"] = [
google.protobuf.json_format.MessageToDict(s) for s in self.stream_slices
]
return out


@dataclass
Expand Down
79 changes: 70 additions & 9 deletions temporalio/streams/providers/native.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,11 @@
A cursor names the run as well as the offset, because an owned stream belongs
to one run and a successor's starts over at zero. A handle without a run id
reads run after run, learning from the poll that a run's stream is closed and
from the run's close event who came next; with a run id it is pinned.
from the run's close event who came next; with a run id it is pinned. A run
that was reset is followed too: its close event does not name the run reset
from it, describe does, and that run's streams continue the base run's offset
space where it inherited a subscription, so the read resumes at the floor the
stream reports rather than at zero.

Prototype support for AI-198. It needs a server built from that branch and
reaches the stream service on a channel of its own, opened with the client's
Expand Down Expand Up @@ -411,6 +415,10 @@ def read(
pinned one: an earlier run of a chain has ended and holds neither the
tail nor the newest records. The server resolves each on the first
poll, in the read that serves it.

The chain is followed across continue-as-new and across a reset: a
run reset from a closed one is read next, from the floor its stream
reports. A handle pinned to a run that was reset ends with that run.
"""
check_read_start(after, last)
topic, result_type = resolve_topic(topic, result_type)
Expand Down Expand Up @@ -469,10 +477,10 @@ async def _read(
break
if self._run_id is not None:
return
successor = await self._successor(run_id)
if successor is None:
following = await self._successor(topic, run_id)
if following is None:
return
run_id, offset = successor, 0
run_id, offset = following

@staticmethod
def _previous(
Expand Down Expand Up @@ -581,31 +589,84 @@ async def _predecessor(self, run_id: str) -> str | None:
try:
async for event in handle.fetch_history_events(page_size=1):
attributes = event.workflow_execution_started_event_attributes
return attributes.continued_execution_run_id or None
if attributes.continued_execution_run_id:
return attributes.continued_execution_run_id
# A reset run's start event is the base run's, copied, and the
# original run id it carries is kept across resets, so it names
# the run the chain of resets began from.
original = attributes.original_execution_run_id
if original and original != run_id:
return original
return None
except RPCError as error:
if error.status != RPCStatusCode.NOT_FOUND:
raise
# The run's History is gone: the chain's retained part starts here.
return None

async def _successor(self, run_id: str) -> str | None:
async def _successor(self, topic: str, run_id: str) -> tuple[str, int] | None:
"""The run that carries on after ``run_id``, and where its stream starts.

A continue-as-new names its successor in the close event, and the
successor's streams start at zero. A reset does not: the base run is
closed with no word of the reset in its own History, and only describe
names the run reset from it. That run reads on streams of its own that
continue the base run's offset space where it inherited a subscription,
and such a stream refuses any offset below its floor, so the read
resumes at the floor the stream reports.
"""
handle = self._client.get_workflow_handle(self._workflow_id, run_id=run_id)
events = handle.fetch_history_events(
event_filter_type=WorkflowHistoryEventFilterType.CLOSE_EVENT
)
closed = False
try:
async for event in events:
closed = True
if event.HasField(
"workflow_execution_continued_as_new_event_attributes"
):
attributes = (
event.workflow_execution_continued_as_new_event_attributes
)
return attributes.new_execution_run_id or None
if not attributes.new_execution_run_id:
return None
return attributes.new_execution_run_id, 0
except RPCError as error:
if error.status != RPCStatusCode.NOT_FOUND:
raise
return None
return None
if not closed:
return None
reset_run = await self._reset_run(run_id)
if reset_run is None:
return None
return reset_run, await self._floor(topic, reset_run)

async def _reset_run(self, run_id: str) -> str | None:
"""The run ``run_id`` was reset into, which only describe reports."""
handle = self._client.get_workflow_handle(self._workflow_id, run_id=run_id)
try:
description = await handle.describe()
except RPCError as error:
if error.status == RPCStatusCode.NOT_FOUND:
return None
raise
extended = description.raw_description.workflow_extended_info
return extended.reset_run_id or None

async def _floor(self, topic: str, run_id: str) -> int:
"""The first offset a run's stream holds.

A stream the run inherited a subscription to exists from the reset on
and starts at the inherited offset. One the base run only published to
is created on the run's first publish, at zero, and does not exist
before that.
"""
try:
return (await self._stream(topic, run_id).describe()).base_offset
except StreamNotFoundError:
return 0


class NativeActivityStreamHandle(NativeStreamHandle):
Expand Down Expand Up @@ -685,7 +746,7 @@ async def _first_run(self) -> str:
async def _predecessor(self, run_id: str) -> str | None:
return None

async def _successor(self, run_id: str) -> str | None:
async def _successor(self, topic: str, run_id: str) -> tuple[str, int] | None:
return None


Expand Down
Loading
Loading