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
4 changes: 4 additions & 0 deletions temporalio/client/_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@
import temporalio.converter
import temporalio.runtime
import temporalio.service
import temporalio.streams
import temporalio.workflow
from temporalio.service import (
ConnectConfig,
Expand Down Expand Up @@ -279,6 +280,7 @@ def __init__(
default_workflow_query_reject_condition: None
| (temporalio.common.QueryRejectCondition) = None,
header_codec_behavior: HeaderCodecBehavior = HeaderCodecBehavior.NO_CODEC,
stream_provider: temporalio.streams.StreamProvider | None = None,
):
"""Create a Temporal client from a service client.

Expand All @@ -293,6 +295,7 @@ def __init__(
interceptors=interceptors,
default_workflow_query_reject_condition=default_workflow_query_reject_condition,
header_codec_behavior=header_codec_behavior,
stream_provider=stream_provider,
)
self._initial_config = config.copy()

Expand Down Expand Up @@ -3306,6 +3309,7 @@ class ClientConfig(TypedDict, total=False):
temporalio.common.QueryRejectCondition | None
]
header_codec_behavior: Required[HeaderCodecBehavior]
stream_provider: temporalio.streams.StreamProvider | None


def _channel_execution(
Expand Down
162 changes: 162 additions & 0 deletions temporalio/streams/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,162 @@
"""Streams: a channel a workflow reads, decides on, and writes.

.. warning::
This module is experimental and may change in future versions. The
design is meant to be the shape that goes GA; the label is the SDK's
release convention, not a licence to break it.

The contract, in five statements:

1. **A workflow publishes only to topics of its own stream, and it publishes
transactionally.** ``temporalio.workflow.StreamWriter.publish``
returns at once. The record is visible when the Workflow Task is accepted,
and never if the task fails, so no reader can see a decision the workflow
did not commit.
2. **Reading is an observation, and the SDK records it.** What
``temporalio.workflow.StreamReader`` handed to workflow code,
including the boundary where it found nothing, is committed with the
commands that reading produced. Recovery re-supplies the same records in
the same order.
3. **Anything that does I/O publishes on its own account.** An activity or an
outside process writes through a :class:`StreamProducer` with a producer
id, an attempt and a sequence, and its records are visible as soon as the
store accepts them. Those three let a reader tell a retry from a new
generation. A retry that carries the same content is written once; one that
carries different content at the same sequence is refused with
:class:`StreamProducerError`, so a divergent retry is never dropped in
silence.
4. **A cursor is opaque and belongs to its provider.** Hand it back to resume
strictly after the record it names; :meth:`StreamHandle.latest` positions a
follower. Do not compare two cursors or do arithmetic on one. A read with
no cursor yet starts at :data:`BEGINNING`, at :data:`END`, or at the last
``N`` records with ``last=N``.
5. **A workflow addresses its streams relative to itself, by topic.** A topic
can be written by the workflow and by outside producers, and read by the
workflow and by outside consumers; which of those happen is the
application's business. A topic is defined once with :func:`topic`, with
the type its records decode to, and that definition is shared by the
workflow, its activities and the backend; a plain string names a topic
decided at runtime. A call that names no topic addresses the workflow's
default topic, :data:`DEFAULT_TOPIC`.

A provider is an object, registered once as a plugin:
``Client.connect(plugins=[provider])``; workers built from that client inherit
it, and ``Worker(plugins=[provider])`` or ``Replayer(plugins=[provider])``
registers it on a worker alone. Each context then asks for its stream the
same way. Workflow code uses ``temporalio.workflow.stream_reader`` and
``temporalio.workflow.stream_writer``. An activity uses
``temporalio.activity.stream_handle``, which is its own workflow pinned
to its run unless told otherwise. Any process holding a client uses
``temporalio.client.Client.get_stream_handle``, which mirrors
``get_workflow_handle``. The explicit form,
``provider.get_stream_handle(client, workflow_id)``, stays for a process that
talks to two stores. This module keeps the shared types, the errors and the
protocols a provider implements; nothing here that workflow code imports does
I/O.

A handle is bound to its client and provider, so a stream is handed to another
process as a :class:`StreamRef`: the owner and the topic as plain data, with
no cursor and no provider name. :meth:`StreamHandle.ref` makes one, the
default data converter carries it as JSON, and the receiver opens it with
``client.get_stream_handle(ref)`` or ``activity.stream_handle(ref)`` on
whatever provider its client has.

A stream can also stand alone, with an id of its own and no owner.
``client.create_stream(stream_id, retention=...)`` creates it with a retention
policy and returns its handle, ``client.get_stream_handle(stream_id=...)``
reaches an existing one, and the handle's ``close()`` seals it, after which
appends are refused with :class:`StreamClosedError` and the retained records
stay readable. A provider whose store cannot hold an ownerless stream raises
:class:`StreamUnsupportedError` for both.

What the contract does not promise: that a :attr:`RecordKind.FINISH` record
means the writing activity succeeded, that a superseded attempt's records can
be withdrawn, or that a stream outlives the retention its provider is
configured for. Reading somebody else's stream is out of scope for this
release.

The record on the wire is ``temporal.api.stream.v1.StreamRecord`` on every
provider, with the user's value in ``body`` as an ordinary payload, so a
reader in any language decodes the same bytes and a payload codec applies. A
provider owes that body what the SDK gives every payload it sends: it encodes
it through the client's data converter, so the codec and the
:class:`temporalio.converter.ExternalStorage` drivers apply, it takes the
retry fingerprint over the converted bytes before either runs and leaves the
plaintext hash on the record under :data:`CONTENT_HASH_KEY`, and it offloads a
workflow's own publish off the workflow thread. :func:`encode_body`,
:func:`decode_body` and :func:`content_fingerprint` are the shared code for
that; :class:`StreamProvider` states the rule.
"""

from __future__ import annotations

from temporalio.streams._body import (
CONTENT_HASH_KEY,
content_fingerprint,
content_hash,
decode_body,
encode_body,
)
from temporalio.streams._errors import (
StreamClosedError,
StreamCursorError,
StreamError,
StreamNotFoundError,
StreamProducerError,
StreamUnsupportedError,
)
from temporalio.streams._provider import (
ReadSource,
StreamHandle,
StreamProducer,
StreamProvider,
WorkflowStreamProvider,
WriteSink,
)
from temporalio.streams._record import (
BEGINNING,
END,
Cursor,
RecordKind,
StreamRecord,
Supersession,
)
from temporalio.streams._ref import StreamOwnerKind, StreamRef
from temporalio.streams._topic import (
DEFAULT_TOPIC,
StreamTopic,
resolve_topic,
topic,
)

__all__ = [
"BEGINNING",
"CONTENT_HASH_KEY",
"DEFAULT_TOPIC",
"END",
"Cursor",
"ReadSource",
"RecordKind",
"StreamClosedError",
"StreamCursorError",
"StreamError",
"StreamHandle",
"StreamNotFoundError",
"StreamOwnerKind",
"StreamProducer",
"StreamProducerError",
"StreamProvider",
"StreamRecord",
"StreamRef",
"StreamTopic",
"StreamUnsupportedError",
"Supersession",
"WorkflowStreamProvider",
"WriteSink",
"content_fingerprint",
"content_hash",
"decode_body",
"encode_body",
"resolve_topic",
"topic",
]
116 changes: 116 additions & 0 deletions temporalio/streams/_body.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,116 @@
"""What a provider owes a record's body between the converter and its store.

:func:`temporalio.streams._wire.to_wire` converts a value into the body with
the payload converter and stops there. Every other payload the SDK sends then
passes through the payload codec and external storage, and a stream body owes
the same, or a codec-protected deployment would leak plaintext through its
streams and a claim-check deployment would push oversized bodies at its store.
A provider runs the body through :func:`encode_body` before it stores or ships
a record and through :func:`decode_body` after it reads one back, off the
workflow thread in both directions.

The order inside :func:`encode_body` is the point. The plaintext hash is taken
first and stamped on the record, and :func:`content_fingerprint` is taken over
the converted records too, before the codec runs, because a codec that
encrypts with a fresh nonce makes every retry's bytes differ, and a store that
fingerprinted those bytes would refuse the retry as a divergent write. The
store keeps the plaintext hash under :data:`CONTENT_HASH_KEY` and can compare
retries by it without ever seeing the plaintext.
"""

from __future__ import annotations

import hashlib
from collections.abc import Sequence

from temporalio.api.common.v1 import Payload
from temporalio.api.stream.v1 import StreamRecord as WireRecord
from temporalio.converter import DataConverter

__all__ = [
"CONTENT_HASH_KEY",
"content_fingerprint",
"content_hash",
"decode_body",
"encode_body",
]

CONTENT_HASH_KEY = "temporal.io/content-hash"
"""The record metadata key the plaintext hash of the body is stored under.

Its value is a payload with ``encoding`` ``binary/plain`` whose data is the
hex digest :func:`content_hash` returns. A ``FINISH`` record carries no body
and no hash.
"""

_HASH_ENCODING = b"binary/plain"


def content_hash(payload: Payload) -> str:
"""The hex SHA-256 of ``payload`` as the converter produced it.

Taken over the deterministic serialization of the whole payload, metadata
included, so two payloads that differ only in their encoding hash apart.
"""
return hashlib.sha256(payload.SerializeToString(deterministic=True)).hexdigest()


def content_fingerprint(records: Sequence[WireRecord]) -> bytes:
"""The identity of one append, taken over its converted records.

Length-delimited, so a batch split differently cannot collide with this
one. Take it before :func:`encode_body`, while the bodies are still what
the converter produced; that is what makes a retry through a
nondeterministic codec match its original.
"""
digest = hashlib.sha256()
for record in records:
body = record.SerializeToString(deterministic=True)
digest.update(len(body).to_bytes(8, "big"))
digest.update(body)
return digest.digest()


async def encode_body(converter: DataConverter, record: WireRecord) -> WireRecord:
"""Stamp the plaintext hash on ``record`` and encode its body for the store.

In place, and returned for convenience. The hash goes under
:data:`CONTENT_HASH_KEY` first; then the body passes through
``converter``'s payload codec and external storage in the order
:meth:`temporalio.converter.DataConverter.encode` uses, so a body above
the external storage threshold is replaced by a claim and the claim is
what the store holds. A record without a body is returned untouched.
"""
if not record.HasField("body"):
return record
record.metadata[CONTENT_HASH_KEY].CopyFrom(
Payload(
metadata={"encoding": _HASH_ENCODING},
data=content_hash(record.body).encode(),
)
)
encoded = await converter._encode_payload_sequence([record.body])
stored = await converter._external_store_payload_sequence(encoded)
record.body.CopyFrom(stored[0])
return record


async def decode_body(converter: DataConverter, record: WireRecord) -> WireRecord:
"""Undo :func:`encode_body` on a record read back from the store.

In place, and returned for convenience. The body is retrieved from
external storage when it is a claim and then run through the payload
codec, in the order :meth:`temporalio.converter.DataConverter.decode`
uses, leaving the payload the converter can turn back into a value. The
hash stays on the record.

Raises:
RuntimeError: The body is a claim and ``converter`` has no external
storage to redeem it with.
"""
if not record.HasField("body"):
return record
retrieved = await converter._external_retrieve_payload_sequence([record.body])
decoded = await converter._decode_payload_sequence(retrieved)
record.body.CopyFrom(decoded[0])
return record
48 changes: 48 additions & 0 deletions temporalio/streams/_errors.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,48 @@
"""The errors a stream call raises.

Every stream condition is a :class:`StreamError`, so a caller can catch by
meaning the way it catches other :class:`temporalio.exceptions.TemporalError`
subclasses. Argument mistakes stay ``ValueError``. A provider's transport
failure surfaces as :class:`temporalio.service.RPCError`, never as the
transport's own exception type.
"""

from __future__ import annotations

import temporalio.exceptions

__all__ = [
"StreamClosedError",
"StreamCursorError",
"StreamError",
"StreamNotFoundError",
"StreamProducerError",
"StreamUnsupportedError",
]


class StreamError(temporalio.exceptions.TemporalError):
"""Base for stream conditions."""


class StreamNotFoundError(StreamError):
"""The workflow, chain or topic does not exist or is past retention."""


class StreamCursorError(StreamError):
"""The cursor was minted by another provider or names a record no longer retained."""


class StreamProducerError(StreamError):
"""The producer attempt or sequence conflicts with what the store holds."""


class StreamClosedError(StreamError):
"""The standalone stream was sealed, so it takes no more records.

Its retained records stay readable; only appends are refused.
"""


class StreamUnsupportedError(StreamError):
"""This provider does not offer the requested capability."""
26 changes: 26 additions & 0 deletions temporalio/streams/_ids.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,26 @@
"""The store key a provider derives from a workflow id and a topic.

A workflow id may contain any character, ``:`` included, so joining the pair
with a bare ``:`` is ambiguous: ``("a:b", "c")`` and ``("a", "b:c")`` would
land in one store. Every provider that keys a store by the pair goes through
:func:`topic_key`, so they all agree and none of them collides.
"""

from __future__ import annotations

__all__ = ["topic_key"]


def _escape(component: str) -> str:
# Percent first, so an escaped component cannot be mistaken for one that
# already contained the escape.
return component.replace("%", "%25").replace(":", "%3A")


def topic_key(workflow_id: str, topic: str) -> str:
"""The store key for ``topic`` of ``workflow_id``'s stream.

Both components are percent-encoded before joining, so the only bare ``:``
in the result is the separator.
"""
return f"{_escape(workflow_id)}:{_escape(topic)}"
Loading
Loading