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
21 changes: 21 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,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