Repository navigation
Update system nexus generation to include notification service #1925
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
80df531
db6195f
7edf5e6
e783534
2552adb
e81a1d8
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| 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 | ||
|
|
||
|
|
@@ -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", | ||
|
|
@@ -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]: | ||
| 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", | ||
|
|
@@ -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) | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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__": | ||
|
|
||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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. |
| +34 −0 | crates/protos/protos/api_upstream/nexus/notification-service.wit |
There was a problem hiding this comment.
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.