Skip to content
Draft
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
2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -205,7 +205,7 @@ 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|_support|notifications|\\.)).*"
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
261 changes: 182 additions & 79 deletions scripts/gen_nexus_system_api.py
Original file line number Diff line number Diff line change
@@ -1,11 +1,11 @@
import ast
import os
import re
import shutil
import subprocess
import sys
import tempfile
from importlib.util import module_from_spec, spec_from_file_location
from pathlib import Path
from typing import cast

import gen_protos

Expand All @@ -22,27 +22,37 @@
/ "api_upstream"
/ "nexus"
)
wit_path = wit_input_dir / "workflow-service.wit"
wit_deps_dir = wit_input_dir / "deps"
python_support_path = base_dir / "scripts" / "nex_gen_support.py"
output_dir = base_dir / "temporalio" / "nexus" / "system" / "workflow_service"
workflow_output_dir = base_dir / "temporalio" / "nexus" / "system" / "workflow_service"
notification_output_dir = base_dir / "temporalio" / "nexus" / "notifications"
support_output_dir = base_dir / "temporalio" / "nexus" / "_support"
workflow_init_path = base_dir / "temporalio" / "workflow" / "__init__.py"
workflowservice_request_response_proto = (
NEX_GEN_VERSION = "0.2.7"
proto_files = [
gen_protos.api_proto_dir
/ "temporal"
/ "api"
/ "workflowservice"
/ "v1"
/ "request_response.proto"
)
NEX_GEN_VERSION = "0.2.4"
/ "request_response.proto",
gen_protos.api_proto_dir
/ "temporal"
/ "api"
/ "notificationservice"
/ "v1"
/ "request_response.proto",
]


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

if shutil.which("nexgen") is None:
if (
shutil.which("nexgen") is None
or subprocess.check_output(["nexgen", "--version"], text=True).strip()
!= f"nexgen {NEX_GEN_VERSION}"
):
subprocess.check_call(
[
"cargo",
Expand All @@ -68,24 +78,99 @@ def build_descriptor_set(descriptor_path: Path) -> None:
f"--proto_path={gen_protos.proto_dir}",
"--include_imports",
f"--descriptor_set_out={descriptor_path}",
str(workflowservice_request_response_proto),
*map(str, proto_files),
]
)


def generate_workflow_exports() -> None:
spec = spec_from_file_location(
"temporalio_nexus_system_workflow_service_exports",
output_dir / "__init__.py",
submodule_search_locations=[str(output_dir)],
def generate_package(
command: list[str], wit_name: str, output_dir: Path, *, native_api: bool
) -> None:
args = [*command, "python", str(wit_input_dir / wit_name), str(wit_deps_dir)]
if native_api:
args.append("--native-api")
subprocess.check_call(
[
*args,
"--system-nexus",
"--support-file",
str(python_support_path),
"--descriptors",
str(output_dir.parent / "temporal_api.bin"),
"--output",
str(output_dir),
]
)
if spec is None or spec.loader is None:
raise RuntimeError(f"Cannot load generated workflow service from {output_dir}")
module = module_from_spec(spec)
sys.modules[spec.name] = module
spec.loader.exec_module(module)
exports = cast(list[str], module.__all__)


def merge_support_trees(support_dirs: list[Path], destination: Path) -> None:
merged_files: dict[Path, Path] = {}
for support_dir in support_dirs:
if not support_dir.is_dir():
raise RuntimeError(
f"generator did not produce support directory: {support_dir}"
)
for source in support_dir.rglob("*"):
if not source.is_file():
continue
relative_path = source.relative_to(support_dir)
if previous := merged_files.get(relative_path):
if previous.read_bytes() != source.read_bytes():
raise RuntimeError(
f"generated support files differ at {relative_path}: {previous} and {source}"
)
else:
merged_files[relative_path] = source
destination.mkdir(parents=True, exist_ok=True)
for relative_path, source in merged_files.items():
target = destination / relative_path
target.parent.mkdir(parents=True, exist_ok=True)
shutil.copy2(source, target)


def rewrite_support_imports(output_dir: Path) -> None:
for source in output_dir.rglob("*.py"):
if "_support" in source.relative_to(output_dir).parts:
continue
content = source.read_text()
rewritten = re.sub(
r"from \._support(\.[A-Za-z_][A-Za-z0-9_]*)? import ",
r"from temporalio.nexus._support\1 import ",
content,
)
if rewritten != content:
source.write_text(rewritten)


def workflow_exports() -> list[str]:

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I assume we need something like this for notification service too.

tree = ast.parse((workflow_output_dir / "__init__.py").read_text())
for statement in tree.body:
if not isinstance(statement, ast.Assign) or not any(
isinstance(target, ast.Name) and target.id == "__all__"
for target in statement.targets
):
continue
value = ast.literal_eval(statement.value)
if not isinstance(value, list) or not all(
isinstance(item, str) for item in value
):
raise RuntimeError(
"generated workflow package __all__ must be a list of strings"
)
return value
raise RuntimeError("generated workflow package does not define __all__")


def replace_marker_block(
content: str, begin: str, end: str, replacement: list[str]
) -> str:
start = content.index(begin)
finish = content.index("\n", content.index(end, start)) + 1
return content[:start] + "".join(replacement) + content[finish:]


def generate_workflow_exports() -> None:
exports = workflow_exports()
import_block = [
"# BEGIN GENERATED NEXUS SYSTEM EXPORTS\n",
"from temporalio.nexus.system.workflow_service import (\n",
Expand All @@ -99,73 +184,91 @@ def generate_workflow_exports() -> None:
" # END GENERATED NEXUS SYSTEM __ALL__\n",
]
content = workflow_init_path.read_text()
start = content.index("# BEGIN GENERATED NEXUS SYSTEM EXPORTS")
end = content.index("# END GENERATED NEXUS SYSTEM EXPORTS", start)
end = content.index("\n", end) + 1
content = content[:start] + "".join(import_block) + content[end:]
start = content.index(" # BEGIN GENERATED NEXUS SYSTEM __ALL__")
end = content.index(" # END GENERATED NEXUS SYSTEM __ALL__", start)
end = content.index("\n", end) + 1
workflow_init_path.write_text(content[:start] + "".join(all_block) + content[end:])
content = replace_marker_block(
content, import_block[0].strip(), import_block[-1].strip(), import_block
)
workflow_init_path.write_text(
replace_marker_block(
content,
" # BEGIN GENERATED NEXUS SYSTEM __ALL__",
" # END GENERATED NEXUS SYSTEM __ALL__",
all_block,
)
)


def generate_nexus_system_api() -> None:
if not wit_path.exists():
raise RuntimeError(f"missing WIT source: {wit_path}")
if not wit_deps_dir.exists():
raise RuntimeError(f"missing WIT dependency directory: {wit_deps_dir}")
if not python_support_path.exists():
raise RuntimeError(f"missing Python support source: {python_support_path}")
def publish_generated_packages(
staged_workflow: Path, staged_notification: Path
) -> None:
staged_support_dirs = [
staged_workflow / "_support",
staged_notification / "_support",
]
staged_support = staged_workflow.parent / "_support"
merge_support_trees(staged_support_dirs, staged_support)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I don't think we should be merging the independent support directories into one. The support files should live with the service they are defined by.

for output_dir in (staged_workflow, staged_notification):
rewrite_support_imports(output_dir)
shutil.rmtree(output_dir / "_support")
for output_dir in (
workflow_output_dir,
notification_output_dir,
support_output_dir,
):
shutil.rmtree(output_dir, ignore_errors=True)
workflow_output_dir.parent.mkdir(parents=True, exist_ok=True)
shutil.copytree(staged_workflow, workflow_output_dir)
notification_output_dir.mkdir(parents=True)
shutil.copy2(
staged_notification / "models.py", notification_output_dir / "models.py"
)
notification_output_dir.joinpath("__init__.py").write_text(
"from .models import OnCompleteRequest, OnCompleteRequestResult, OnCompleteResponse\n\n"
"__all__ = [\n"
' "OnCompleteRequest",\n'
' "OnCompleteRequestResult",\n'
' "OnCompleteResponse",\n'
"]\n"
)
shutil.copytree(staged_support, support_output_dir)
workflow_output_dir.parent.joinpath("__init__.py").touch()


def generate_nexus_system_api() -> None:
required_paths = [
wit_input_dir / "workflow-service.wit",
wit_input_dir / "notification-service.wit",
wit_deps_dir,
python_support_path,
*proto_files,
]
for path in required_paths:
if not path.exists():
raise RuntimeError(f"missing generator input: {path}")
with tempfile.TemporaryDirectory(dir=base_dir) as temp_dir:
descriptor_path = Path(temp_dir) / "temporal_api.bin"
staging_dir = Path(temp_dir)
descriptor_path = staging_dir / "temporal_api.bin"
build_descriptor_set(descriptor_path)
command = nex_gen_command()

shutil.rmtree(output_dir, ignore_errors=True)
output_dir.parent.mkdir(parents=True, exist_ok=True)
subprocess.check_call(
[
*command,
"python",
str(wit_path),
str(wit_deps_dir),
"--native-api",
"--system-nexus",
"--support-file",
str(python_support_path),
"--descriptors",
str(descriptor_path),
"--output",
str(output_dir),
]
staged_workflow = staging_dir / "workflow_service"
staged_notification = staging_dir / "notifications"
generate_package(
command, "workflow-service.wit", staged_workflow, native_api=True
)

(output_dir.parent / "__init__.py").touch()
generate_package(
command, "notification-service.wit", staged_notification, native_api=False
)
publish_generated_packages(staged_workflow, staged_notification)
generate_workflow_exports()
format_paths = [
str(workflow_output_dir),
str(notification_output_dir),
str(support_output_dir),
str(workflow_init_path),
]
subprocess.check_call(
[
sys.executable,
"-m",
"ruff",
"check",
"--select",
"I",
"--fix",
str(output_dir),
str(workflow_init_path),
]
)
subprocess.check_call(
[
sys.executable,
"-m",
"ruff",
"format",
str(output_dir),
str(workflow_init_path),
]
[sys.executable, "-m", "ruff", "check", "--select", "I", "--fix", *format_paths]
)
subprocess.check_call([sys.executable, "-m", "ruff", "format", *format_paths])


if __name__ == "__main__":
Expand Down
24 changes: 23 additions & 1 deletion scripts/nex_gen_support.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@

import temporalio.api.common.v1.message_pb2 as common_pb2
import temporalio.api.enums.v1.workflow_pb2 as workflow_enums_pb2
import temporalio.api.failure.v1.message_pb2 as failure_pb2
import temporalio.api.taskqueue.v1.message_pb2 as taskqueue_pb2
import temporalio.api.workflow.v1
import temporalio.common
Expand Down Expand Up @@ -142,8 +143,11 @@ def _payload_to_value(payload: common_pb2.Payload) -> object:

def payload_from_proto(
proto: common_pb2.Payload,
type_hint: typing.Any = None,
) -> object:
return _payload_to_value(proto)
return temporalio.nexus.system._current_user_payload_converter().from_payloads(
[proto], [type_hint] if type_hint is not None else None
)[0]


def payload_to_proto(
Expand All @@ -152,6 +156,24 @@ def payload_to_proto(
return _value_to_payload(payload)


def failure_from_proto(
proto: failure_pb2.Failure,
) -> BaseException:
return temporalio.nexus.system._current_user_failure_converter().from_failure(
proto, temporalio.nexus.system._current_user_payload_converter()
)


def failure_to_proto(
failure: BaseException,
) -> failure_pb2.Failure:
proto = failure_pb2.Failure()
temporalio.nexus.system._current_user_failure_converter().to_failure(
failure, temporalio.nexus.system._current_user_payload_converter(), proto
)
return proto


def memo_from_proto(
proto: common_pb2.Memo,
) -> collections.abc.Mapping[str, object]:
Expand Down
2 changes: 1 addition & 1 deletion temporalio/bridge/sdk-core

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Just a reminder it seems like we need to put this in rust main/API before completing.

Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
# Generated by nexgen v0.2.4. DO NOT EDIT!
# Generated by nexgen v0.2.7. DO NOT EDIT!

from __future__ import annotations

Expand Down
Loading
Loading