Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
58 commits
Select commit Hold shift + click to select a range
f641ea8
Added the changelog entry for the Workflow Streams provider.
moedash Sep 16, 2026
f35717e
Added the Nexus front that hides the store behind one endpoint.
moedash Sep 15, 2026
832b8c3
Made cursors exclusive, added latest(), and allowed topic appends.
moedash Sep 16, 2026
50b1dab
Made the Nexus stream front pass the linters.
moedash Sep 16, 2026
1ee4c7a
Added the changelog entry for the Nexus stream front.
moedash Sep 16, 2026
93f87cf
Bounded the Nexus read with a timeout 3.10 has.
moedash Sep 16, 2026
b6637c9
Added the IDL contract and generation for the stream endpoint.
moedash Sep 18, 2026
b90346a
Moved the nexus provider onto the generated contract.
moedash Sep 18, 2026
299a90b
Applied the payload codec on the nexus provider's caller side.
moedash Sep 18, 2026
534bd59
Made the stream contract generator refuse a nexgen other than the pin…
moedash Sep 19, 2026
3cb7d1a
Regenerated the stream contract with byte payloads, read bounds and t…
moedash Sep 19, 2026
c54d750
Reworked the Nexus front's dedupe, reads, error mapping and endpoint …
moedash Sep 19, 2026
769cfd0
Covered the Nexus front's refusals, retries, subscription reuse and d…
moedash Sep 19, 2026
ff52038
Carried the StreamRecord proto and producer identity on the stream co…
moedash Sep 21, 2026
46ae99b
Rewrote the Nexus front on the plugin surface.
moedash Sep 21, 2026
95be1ae
Covered the Nexus front on the new surface.
moedash Sep 21, 2026
8996d41
Let the Nexus front register on a client like any provider.
moedash Sep 21, 2026
40348b9
Took topic definitions on the Nexus front's handle.
moedash Sep 21, 2026
807af3e
Merged the workflow-streams branch with the connect config fix and Co…
moedash Sep 22, 2026
2d7f0c9
Made the stream handler hold one batch and one read at a time.
moedash Sep 25, 2026
f6ef1a7
Kept the endpoint's transport failures off the caller.
moedash Sep 25, 2026
eb49d81
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Sep 26, 2026
e629ae1
Encoded the Nexus contract with the internal payload converter.
moedash Sep 26, 2026
526532b
Expected the one-based sequence from the memory producer.
moedash Sep 26, 2026
ca69a20
Merge remote-tracking branch 'origin/moe/AI-198-streams-provider-work…
moedash Sep 29, 2026
bc93c6b
Merge remote-tracking branch 'origin/moe/AI-198-streams-provider-work…
moedash Sep 29, 2026
6563558
Absorbed the current stream contract on the Nexus front.
moedash Sep 29, 2026
fa7c994
Addressed a stream by reference on the Nexus front.
moedash Sep 30, 2026
913d13d
Merge remote-tracking branch 'origin/moe/AI-198-py-04-accessors' into…
moedash Sep 30, 2026
b8ca5b8
Handed a stream across a Nexus operation as a StreamRef.
moedash Sep 30, 2026
5c6e0f4
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Sep 30, 2026
f7fcf03
Used a standalone activity as the owner Workflow Streams refuses.
moedash Sep 30, 2026
a4e6567
Merge remote-tracking branch 'origin/moe/AI-198-streams-provider-work…
moedash Sep 30, 2026
85126ae
Typed the absent client in the ref-opening test as Any.
moedash Sep 30, 2026
85bf45d
Merge remote-tracking branch 'origin/moe/AI-198-streams-provider-work…
moedash Sep 30, 2026
9820335
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Oct 1, 2026
ee0ab1b
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Oct 1, 2026
a7804c2
Declared the reference's optional members nullable on the wire.
moedash Oct 1, 2026
369b792
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Oct 1, 2026
355cbcf
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Oct 1, 2026
0b9bf57
Pointed the Nexus provider docstring at the call that opens a reference.
moedash Oct 1, 2026
d751e93
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Oct 2, 2026
83c2792
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Oct 2, 2026
f7b3e31
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Oct 2, 2026
c053c1a
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Oct 2, 2026
15ff000
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Oct 2, 2026
2f100a4
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Oct 2, 2026
223f727
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Oct 2, 2026
c051962
Added the Nexus operation that consumes a stream through its channel.
moedash Oct 2, 2026
560e278
Tested the stream consumer against a channel server.
moedash Oct 2, 2026
65df92f
Named the consumer's channel the way the server names a stream's.
moedash Oct 2, 2026
0acc115
Merge branch 'moe/AI-198-streams-provider-workflow-streams' of github…
moedash Oct 2, 2026
22db179
Took the channel address from the client's stream_channel rule.
moedash Oct 2, 2026
4b06358
Merge branch 'moe/AI-198-streams-provider-workflow-streams' of github…
moedash Oct 2, 2026
b110ed0
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Oct 3, 2026
c85be89
Registered the stream consumer by execution.
moedash Oct 3, 2026
5d31c9e
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
moedash Oct 3, 2026
995c4b4
Merge branch 'moe/AI-198-streams-provider-workflow-streams' into moe/…
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
16 changes: 16 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -62,6 +62,22 @@ to include examples, links to docs, or any other relevant information.
worker plugin, so a workflow reads and publishes through
`temporalio.contrib.workflow_streams` without naming it. Records are the
`StreamRecord` proto inside the shipped item payload, and a handle without a
run id follows continue-as-new run by run.
- **Experimental**: `temporalio.streams.providers.workflow_streams` serves the
stream interface over the shipped Workflow Streams transport, so a workflow
reads and publishes through `temporalio.contrib.workflow_streams` without
naming it.
- **Experimental**: `temporalio.streams.providers.nexus.NexusStreams` puts one
Nexus endpoint in front of a storage provider, so a caller reaches a stream
through the endpoint and never names the store, and
`TemporalStreamsHandler` serves that endpoint by fronting the provider's own
handles. Its contract is defined in `temporal_streams.nexusrpc.yaml` and the
bindings are generated from it; a record crosses as the serialized
`StreamRecord` proto, and both operations address a stream by a `StreamRef`
naming its owner (a workflow, an activity or a standalone stream) and topic,
which the handler maps onto the store's accessor for that owner. Configure
the front with `data_converter=` to run a payload codec on the caller side,
so records are encoded before they leave the process.
run id follows continue-as-new run by run and a reset into the run reset to.
An outside publish is an Update that answers with the batch's position and
refuses a conflicting repeat, falling back to the shipped Signal on a
Expand Down
11 changes: 10 additions & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -125,16 +125,19 @@ format = [
]
gen-docs = "uv run scripts/gen_docs.py"
gen-nexus-system-api = "uv run scripts/gen_nexus_system_api.py"
gen-streams-nexus-api = "uv run scripts/gen_streams_nexus_api.py"
gen-protos = [
{ cmd = "uv run scripts/gen_protos.py" },
{ ref = "gen-nexus-system-api" },
{ ref = "gen-streams-nexus-api" },
{ cmd = "uv run scripts/gen_payload_visitor.py" },
{ cmd = "uv run scripts/gen_bridge_client.py" },
{ ref = "format" },
]
gen-protos-docker = [
{ cmd = "uv run scripts/gen_protos_docker.py" },
{ ref = "gen-nexus-system-api" },
{ ref = "gen-streams-nexus-api" },
{ cmd = "uv run scripts/gen_payload_visitor.py" },
{ cmd = "uv run scripts/gen_bridge_client.py" },
{ ref = "format" },
Expand Down Expand Up @@ -198,16 +201,21 @@ exclude = [
'temporalio/api',
'temporalio/bridge/proto',
'temporalio/nexus/system/workflow_service',
'temporalio/streams/providers/_nexus_generated',
]

[[tool.mypy.overrides]]
module = "temporalio.nexus.system.workflow_service.*"
ignore_errors = true

[[tool.mypy.overrides]]
module = "temporalio.streams.providers._nexus_generated.*"
ignore_errors = true

[tool.pydocstyle]
convention = "google"
# https://github.com/PyCQA/pydocstyle/issues/363#issuecomment-625563088
match_dir = "^(?!(docs|scripts|tests|api|proto|system|\\.)).*"
match_dir = "^(?!(docs|scripts|tests|api|proto|system|_nexus_generated|\\.)).*"
add_ignore = [
# We like to wrap at a certain number of chars, even long summary sentences.
# https://github.com/PyCQA/pydocstyle/issues/184
Expand Down Expand Up @@ -241,6 +249,7 @@ privacy = [
"HIDDEN:temporalio.worker.workflow_sandbox.importer",
"HIDDEN:temporalio.worker.workflow_sandbox.in_sandbox",
"HIDDEN:**.*_pb2*",
"HIDDEN:temporalio.streams.providers._nexus_generated._definitions",
]
project-name = "Temporal Python"
sidebar-expand-depth = 2
Expand Down
95 changes: 95 additions & 0 deletions scripts/gen_streams_nexus_api.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,95 @@
import os
import shutil
import subprocess
import sys
from pathlib import Path

base_dir = Path(__file__).parent.parent
providers_dir = base_dir / "temporalio" / "streams" / "providers"
contract_path = providers_dir / "temporal_streams.nexusrpc.yaml"
output_dir = providers_dir / "_nexus_generated"
# Pinned to the version CI installs, because nexgen's output changes between
# releases and check-protos compares against what is committed here.
NEX_GEN_VERSION = "0.2.4"


def nex_gen_command() -> list[str]:
if bin_path := os.environ.get("NEX_GEN_BIN"):
return [bin_path]

if shutil.which("nexgen") is None:
subprocess.check_call(
[
"cargo",
"install",
"--locked",
"nexgen",
"--version",
NEX_GEN_VERSION,
# Same build as the system API script installs, so one binary
# serves both and neither can overwrite the other's output.
"--features",
"advanced",
"--force",
]
)
return ["nexgen"]


def check_version(command: list[str]) -> None:
# A different release on PATH would regenerate different code locally
# and the drift would only show in CI, so refuse before writing anything.
reported = subprocess.check_output([*command, "--version"], text=True).strip()
found = reported.split()[-1] if reported else ""
if found != NEX_GEN_VERSION:
raise SystemExit(
f"found nexgen {found or '?'} at {command[0]}, but the stream contract is "
f"generated with {NEX_GEN_VERSION}. Install it with `cargo install --locked "
f"nexgen --version {NEX_GEN_VERSION} --features advanced` or point "
"NEX_GEN_BIN at that binary."
)


def generate_streams_nexus_api() -> None:
if not contract_path.exists():
raise RuntimeError(f"missing stream contract: {contract_path}")

command = nex_gen_command()
check_version(command)
shutil.rmtree(output_dir, ignore_errors=True)
subprocess.check_call(
[
*command,
"python",
str(contract_path),
"--output",
str(output_dir),
]
)
subprocess.check_call(
[
sys.executable,
"-m",
"ruff",
"check",
"--select",
"I",
"--fix",
str(output_dir),
]
)
subprocess.check_call(
[
sys.executable,
"-m",
"ruff",
"format",
str(output_dir),
]
)


if __name__ == "__main__":
print("Generating stream endpoint Nexus API...", file=sys.stderr)
generate_streams_nexus_api()
print("Done", file=sys.stderr)
25 changes: 25 additions & 0 deletions temporalio/streams/providers/_nexus_generated/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,25 @@
# Generated by nexgen v0.2.4. DO NOT EDIT!

from __future__ import annotations

from ._definitions import Violation
from .models import (
AppendInput,
AppendOutput,
ReadInput,
ReadOutput,
RecordWire,
StreamRef,
)
from .services import TemporalStreams

__all__ = [
"Violation",
"AppendInput",
"AppendOutput",
"ReadInput",
"ReadOutput",
"RecordWire",
"StreamRef",
"TemporalStreams",
]
Loading
Loading