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
11 changes: 11 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -76,6 +76,17 @@ to include examples, links to docs, or any other relevant information.
refuses a conflicting repeat, falling back to the shipped Signal on a
workflow whose worker predates it. A workflow's activity keeps its own
streams in the workflow's log under `activity/<id>/<name>`.
- **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.

### Changed

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