mirror of
https://github.com/langchain-ai/langgraph.git
synced 2026-09-08 18:57:52 +02:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
5854de5dc5 | ||
|
|
8d1acadbc9 |
@@ -32,6 +32,8 @@ def test(config: pathlib.Path, port: int, tag: str, verbose: bool):
|
||||
docker_compose=None,
|
||||
port=port,
|
||||
watch=False,
|
||||
debugger_port=None,
|
||||
debugger_base_url=f"http://127.0.0.1:{port}",
|
||||
postgres_uri=None,
|
||||
api_version=None,
|
||||
image=tag,
|
||||
@@ -171,5 +173,5 @@ if __name__ == "__main__":
|
||||
except BaseException:
|
||||
logger.exception("Test failed")
|
||||
raise
|
||||
|
||||
|
||||
logger.info("Test execution finished")
|
||||
|
||||
@@ -77,10 +77,6 @@ __pypackages__/
|
||||
# Environments
|
||||
.env
|
||||
.envrc
|
||||
*.crt
|
||||
*.key
|
||||
*.pem
|
||||
credentials.json
|
||||
.venv
|
||||
.venvs
|
||||
env/
|
||||
@@ -102,7 +98,6 @@ dmypy.json
|
||||
|
||||
.vercel
|
||||
.turbo
|
||||
node_modules/
|
||||
.editorconfig
|
||||
.scratch
|
||||
.worktrees/
|
||||
|
||||
@@ -48,6 +48,9 @@ def get_anonymized_params(
|
||||
if kwargs.get("docker_compose"):
|
||||
params["docker_compose"] = True
|
||||
|
||||
if kwargs.get("debugger_port"):
|
||||
params["debugger_port"] = True
|
||||
|
||||
if kwargs.get("postgres_uri"):
|
||||
params["postgres_uri"] = True
|
||||
|
||||
|
||||
@@ -5,7 +5,6 @@ import pathlib
|
||||
import shutil
|
||||
import sys
|
||||
from collections.abc import Sequence
|
||||
from urllib.parse import SplitResult, urlencode, urlsplit, urlunsplit
|
||||
|
||||
import click
|
||||
import click.exceptions
|
||||
@@ -141,6 +140,17 @@ OPT_VERBOSE = click.option(
|
||||
help="Show more output from the server logs",
|
||||
)
|
||||
OPT_WATCH = click.option("--watch", is_flag=True, help="Restart on file changes")
|
||||
OPT_DEBUGGER_PORT = click.option(
|
||||
"--debugger-port",
|
||||
type=int,
|
||||
help="Pull the debugger image locally and serve the UI on specified port",
|
||||
)
|
||||
OPT_DEBUGGER_BASE_URL = click.option(
|
||||
"--debugger-base-url",
|
||||
type=str,
|
||||
help="URL used by the debugger to access LangGraph API. Defaults to http://127.0.0.1:[PORT]",
|
||||
)
|
||||
|
||||
OPT_POSTGRES_URI = click.option(
|
||||
"--postgres-uri",
|
||||
help="Postgres URI to use for the database. Defaults to launching a local database",
|
||||
@@ -232,94 +242,18 @@ cli.add_command(deploy)
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
def _validated_http_url(value: str, option_name: str) -> SplitResult:
|
||||
try:
|
||||
parsed = urlsplit(value)
|
||||
hostname = parsed.hostname
|
||||
_ = parsed.port
|
||||
except ValueError as exc:
|
||||
raise click.UsageError(
|
||||
f"{option_name} must be a valid HTTP(S) URL without credentials."
|
||||
) from exc
|
||||
|
||||
if (
|
||||
value != value.strip()
|
||||
or parsed.scheme not in {"http", "https"}
|
||||
or not parsed.netloc
|
||||
or not hostname
|
||||
or parsed.username is not None
|
||||
or parsed.password is not None
|
||||
):
|
||||
raise click.UsageError(
|
||||
f"{option_name} must be a valid HTTP(S) URL without credentials."
|
||||
)
|
||||
return parsed
|
||||
|
||||
|
||||
def _studio_link(
|
||||
*,
|
||||
port: int,
|
||||
studio_url: str | None,
|
||||
api_url: str | None,
|
||||
debugger_base_url: str | None,
|
||||
) -> str:
|
||||
if debugger_base_url is not None:
|
||||
if api_url is not None and api_url != debugger_base_url:
|
||||
raise click.UsageError(
|
||||
"--api-url and --debugger-base-url cannot specify different URLs."
|
||||
)
|
||||
click.echo(
|
||||
"Warning: --debugger-base-url is deprecated; use --api-url instead.",
|
||||
err=True,
|
||||
)
|
||||
api_url = debugger_base_url
|
||||
|
||||
studio_url = "https://smith.langchain.com" if studio_url is None else studio_url
|
||||
api_url = f"http://127.0.0.1:{port}" if api_url is None else api_url
|
||||
studio_parts = _validated_http_url(studio_url, "--studio-url")
|
||||
_validated_http_url(api_url, "--api-url")
|
||||
if studio_parts.query or studio_parts.fragment:
|
||||
raise click.UsageError(
|
||||
"--studio-url must not include a query string or fragment."
|
||||
)
|
||||
|
||||
studio_path = f"{studio_parts.path.rstrip('/')}/studio/"
|
||||
return urlunsplit(
|
||||
studio_parts._replace(
|
||||
path=studio_path,
|
||||
query=urlencode({"baseUrl": api_url}),
|
||||
)
|
||||
)
|
||||
|
||||
|
||||
@OPT_RECREATE
|
||||
@OPT_PULL
|
||||
@OPT_PORT
|
||||
@OPT_DOCKER_COMPOSE
|
||||
@OPT_CONFIG
|
||||
@OPT_VERBOSE
|
||||
@OPT_DEBUGGER_PORT
|
||||
@OPT_DEBUGGER_BASE_URL
|
||||
@OPT_WATCH
|
||||
@OPT_POSTGRES_URI
|
||||
@OPT_API_VERSION
|
||||
@OPT_ENGINE_RUNTIME_MODE
|
||||
@click.option(
|
||||
"--studio-url",
|
||||
type=str,
|
||||
default=None,
|
||||
help="URL of the LangGraph Studio instance. Defaults to https://smith.langchain.com",
|
||||
)
|
||||
@click.option(
|
||||
"--api-url",
|
||||
type=str,
|
||||
default=None,
|
||||
help="URL that LangGraph Studio uses to access the API. Defaults to http://127.0.0.1:[PORT]",
|
||||
)
|
||||
@click.option(
|
||||
"--debugger-base-url",
|
||||
type=str,
|
||||
default=None,
|
||||
hidden=True,
|
||||
)
|
||||
@click.option(
|
||||
"--image",
|
||||
type=str,
|
||||
@@ -350,21 +284,14 @@ def up(
|
||||
watch: bool,
|
||||
wait: bool,
|
||||
verbose: bool,
|
||||
debugger_port: int | None,
|
||||
debugger_base_url: str | None,
|
||||
postgres_uri: str | None,
|
||||
api_version: str | None,
|
||||
engine_runtime_mode: str,
|
||||
studio_url: str | None,
|
||||
api_url: str | None,
|
||||
debugger_base_url: str | None,
|
||||
image: str | None,
|
||||
base_image: str | None,
|
||||
):
|
||||
studio_link = _studio_link(
|
||||
port=port,
|
||||
studio_url=studio_url,
|
||||
api_url=api_url,
|
||||
debugger_base_url=debugger_base_url,
|
||||
)
|
||||
click.secho("Starting LangGraph API server...", fg="green")
|
||||
click.secho(
|
||||
"""For local dev, requires env var LANGSMITH_API_KEY with access to LangSmith Deployment.
|
||||
@@ -381,6 +308,8 @@ For production use, requires a license key in env var LANGGRAPH_CLOUD_LICENSE_KE
|
||||
pull=pull,
|
||||
watch=watch,
|
||||
verbose=verbose,
|
||||
debugger_port=debugger_port,
|
||||
debugger_base_url=debugger_base_url,
|
||||
postgres_uri=postgres_uri,
|
||||
api_version=api_version,
|
||||
engine_runtime_mode=engine_runtime_mode,
|
||||
@@ -408,12 +337,20 @@ For production use, requires a license key in env var LANGGRAPH_CLOUD_LICENSE_KE
|
||||
if "unpacking to docker.io" in line:
|
||||
set("Starting...")
|
||||
elif "Application startup complete" in line:
|
||||
debugger_origin = (
|
||||
f"http://localhost:{debugger_port}"
|
||||
if debugger_port
|
||||
else "https://smith.langchain.com"
|
||||
)
|
||||
debugger_base_url_query = (
|
||||
debugger_base_url or f"http://127.0.0.1:{port}"
|
||||
)
|
||||
set("")
|
||||
sys.stdout.write(
|
||||
f"""Ready!
|
||||
- API: http://localhost:{port}
|
||||
- Docs: http://localhost:{port}/docs
|
||||
- LangGraph Studio: {studio_link}
|
||||
- LangGraph Studio: {debugger_origin}/studio/?baseUrl={debugger_base_url_query}
|
||||
"""
|
||||
)
|
||||
sys.stdout.flush()
|
||||
@@ -998,6 +935,8 @@ def prepare_args_and_stdin(
|
||||
docker_compose: pathlib.Path | None,
|
||||
port: int,
|
||||
watch: bool,
|
||||
debugger_port: int | None = None,
|
||||
debugger_base_url: str | None = None,
|
||||
postgres_uri: str | None = None,
|
||||
api_version: str | None = None,
|
||||
engine_runtime_mode: str = "combined_queue_worker",
|
||||
@@ -1011,6 +950,8 @@ def prepare_args_and_stdin(
|
||||
stdin = langgraph_cli.docker.compose(
|
||||
capabilities,
|
||||
port=port,
|
||||
debugger_port=debugger_port,
|
||||
debugger_base_url=debugger_base_url,
|
||||
postgres_uri=postgres_uri,
|
||||
image=image,
|
||||
base_image=base_image,
|
||||
@@ -1048,6 +989,8 @@ def prepare(
|
||||
pull: bool,
|
||||
watch: bool,
|
||||
verbose: bool,
|
||||
debugger_port: int | None = None,
|
||||
debugger_base_url: str | None = None,
|
||||
postgres_uri: str | None = None,
|
||||
api_version: str | None = None,
|
||||
engine_runtime_mode: str = "combined_queue_worker",
|
||||
@@ -1089,6 +1032,8 @@ def prepare(
|
||||
docker_compose=docker_compose,
|
||||
port=port,
|
||||
watch=watch,
|
||||
debugger_port=debugger_port,
|
||||
debugger_base_url=debugger_base_url or f"http://127.0.0.1:{port}",
|
||||
postgres_uri=postgres_uri,
|
||||
api_version=api_version,
|
||||
engine_runtime_mode=engine_runtime_mode,
|
||||
|
||||
@@ -142,6 +142,29 @@ def check_capabilities(runner) -> DockerCapabilities:
|
||||
)
|
||||
|
||||
|
||||
def debugger_compose(*, port: int | None = None, base_url: str | None = None) -> dict:
|
||||
if port is None:
|
||||
return ""
|
||||
|
||||
config = {
|
||||
"langgraph-debugger": {
|
||||
"image": "langchain/langgraph-debugger",
|
||||
"restart": "on-failure",
|
||||
"depends_on": {
|
||||
"langgraph-postgres": {"condition": "service_healthy"},
|
||||
},
|
||||
"ports": [f'"{port}:3968"'],
|
||||
}
|
||||
}
|
||||
|
||||
if base_url:
|
||||
config["langgraph-debugger"]["environment"] = {
|
||||
"VITE_STUDIO_LOCAL_GRAPH_URL": base_url
|
||||
}
|
||||
|
||||
return config
|
||||
|
||||
|
||||
# Function to convert dictionary to YAML
|
||||
def dict_to_yaml(d: dict, *, indent: int = 0) -> str:
|
||||
"""Convert a dictionary to a YAML string."""
|
||||
@@ -168,6 +191,8 @@ def compose_as_dict(
|
||||
capabilities: DockerCapabilities,
|
||||
*,
|
||||
port: int,
|
||||
debugger_port: int | None = None,
|
||||
debugger_base_url: str | None = None,
|
||||
# postgres://user:password@host:port/database?option=value
|
||||
postgres_uri: str | None = None,
|
||||
# If you are running against an already-built image, you can pass it here
|
||||
@@ -228,6 +253,12 @@ def compose_as_dict(
|
||||
else:
|
||||
services["langgraph-postgres"]["healthcheck"]["interval"] = "5s"
|
||||
|
||||
# Add optional debugger service if debugger_port is specified
|
||||
if debugger_port:
|
||||
services["langgraph-debugger"] = debugger_compose(
|
||||
port=debugger_port, base_url=debugger_base_url
|
||||
)["langgraph-debugger"]
|
||||
|
||||
# Add langgraph-api service
|
||||
api_environment = {
|
||||
"REDIS_URI": "redis://langgraph-redis:6379",
|
||||
@@ -258,7 +289,7 @@ def compose_as_dict(
|
||||
"test": "python /api/healthcheck.py",
|
||||
"interval": "60s",
|
||||
"start_interval": "1s",
|
||||
"start_period": "60s",
|
||||
"start_period": "10s",
|
||||
}
|
||||
|
||||
# Final compose dictionary with volumes included if needed
|
||||
@@ -274,6 +305,8 @@ def compose(
|
||||
capabilities: DockerCapabilities,
|
||||
*,
|
||||
port: int,
|
||||
debugger_port: int | None = None,
|
||||
debugger_base_url: str | None = None,
|
||||
# postgres://user:password@host:port/database?option=value
|
||||
postgres_uri: str | None = None,
|
||||
image: str | None = None,
|
||||
@@ -285,6 +318,8 @@ def compose(
|
||||
compose_content = compose_as_dict(
|
||||
capabilities,
|
||||
port=port,
|
||||
debugger_port=debugger_port,
|
||||
debugger_base_url=debugger_base_url,
|
||||
postgres_uri=postgres_uri,
|
||||
image=image,
|
||||
base_image=base_image,
|
||||
|
||||
@@ -970,20 +970,21 @@ def python_config_to_docker_uv_lock(
|
||||
f"{uv_export_project_dir}/uv.lock",
|
||||
)
|
||||
)
|
||||
for package_root in sorted(
|
||||
plan.all_workspace_roots,
|
||||
key=lambda root: root.as_posix(),
|
||||
):
|
||||
if package_root == plan.project_root:
|
||||
continue
|
||||
package_relative_path = pathlib.PurePosixPath(
|
||||
package_root.relative_to(plan.project_root).as_posix()
|
||||
# `uv export --package` resolves a member by reading that member's own
|
||||
# manifest, so the root manifest names it but is not enough to find it. Only
|
||||
# the manifests are copied; member sources arrive later, per member.
|
||||
#
|
||||
# Restricted to the closure rather than every member. Members outside it are
|
||||
# not needed for the export, and copying them would reach for paths that
|
||||
# `.dockerignore` may have removed from the build context.
|
||||
for member_root in sorted(set(plan.container_roots) - {plan.project_root}):
|
||||
member_relative = pathlib.PurePosixPath(
|
||||
member_root.relative_to(plan.project_root).as_posix()
|
||||
)
|
||||
package_pyproject_path = package_relative_path / "pyproject.toml"
|
||||
docker_plan.add_raw(
|
||||
copy_from_project_root(
|
||||
package_pyproject_path,
|
||||
f"{uv_export_project_dir}/{package_pyproject_path.as_posix()}",
|
||||
member_relative / "pyproject.toml",
|
||||
f"{uv_export_project_dir}/{member_relative}/pyproject.toml",
|
||||
)
|
||||
)
|
||||
docker_plan.add_instruction("WORKDIR", uv_export_project_dir)
|
||||
|
||||
@@ -8,11 +8,10 @@ from contextlib import contextmanager
|
||||
from pathlib import Path
|
||||
|
||||
import click
|
||||
import pytest
|
||||
from click.testing import CliRunner
|
||||
|
||||
import langgraph_cli.deploy as deploy_module
|
||||
from langgraph_cli.cli import _studio_link, cli, prepare_args_and_stdin
|
||||
from langgraph_cli.cli import cli, prepare_args_and_stdin
|
||||
from langgraph_cli.config import Config, _get_pip_cleanup_lines, validate_config
|
||||
from langgraph_cli.docker import DEFAULT_POSTGRES_URI, DockerCapabilities, Version
|
||||
from langgraph_cli.util import clean_empty_lines
|
||||
@@ -57,6 +56,8 @@ def test_prepare_args_and_stdin() -> None:
|
||||
Config(dependencies=[".", "../../.."], graphs={"agent": "agent.py:graph"})
|
||||
)
|
||||
port = 8000
|
||||
debugger_port = 8001
|
||||
debugger_graph_url = f"http://127.0.0.1:{port}"
|
||||
|
||||
actual_args, actual_stdin = prepare_args_and_stdin(
|
||||
capabilities=DEFAULT_DOCKER_CAPABILITIES,
|
||||
@@ -64,6 +65,8 @@ def test_prepare_args_and_stdin() -> None:
|
||||
config=config,
|
||||
docker_compose=pathlib.Path("custom-docker-compose.yml"),
|
||||
port=port,
|
||||
debugger_port=debugger_port,
|
||||
debugger_base_url=debugger_graph_url,
|
||||
watch=True,
|
||||
)
|
||||
|
||||
@@ -107,6 +110,16 @@ services:
|
||||
retries: 5
|
||||
interval: 60s
|
||||
start_interval: 1s
|
||||
langgraph-debugger:
|
||||
image: langchain/langgraph-debugger
|
||||
restart: on-failure
|
||||
depends_on:
|
||||
langgraph-postgres:
|
||||
condition: service_healthy
|
||||
ports:
|
||||
- "{debugger_port}:3968"
|
||||
environment:
|
||||
VITE_STUDIO_LOCAL_GRAPH_URL: {debugger_graph_url}
|
||||
langgraph-api:
|
||||
ports:
|
||||
- "8000:8000"
|
||||
@@ -122,7 +135,7 @@ services:
|
||||
test: python /api/healthcheck.py
|
||||
interval: 60s
|
||||
start_interval: 1s
|
||||
start_period: 60s
|
||||
start_period: 10s
|
||||
|
||||
pull_policy: build
|
||||
build:
|
||||
@@ -165,6 +178,8 @@ def test_prepare_args_and_stdin_with_image() -> None:
|
||||
Config(dependencies=[".", "../../.."], graphs={"agent": "agent.py:graph"})
|
||||
)
|
||||
port = 8000
|
||||
debugger_port = 8001
|
||||
debugger_graph_url = f"http://127.0.0.1:{port}"
|
||||
|
||||
actual_args, actual_stdin = prepare_args_and_stdin(
|
||||
capabilities=DEFAULT_DOCKER_CAPABILITIES,
|
||||
@@ -172,6 +187,8 @@ def test_prepare_args_and_stdin_with_image() -> None:
|
||||
config=config,
|
||||
docker_compose=pathlib.Path("custom-docker-compose.yml"),
|
||||
port=port,
|
||||
debugger_port=debugger_port,
|
||||
debugger_base_url=debugger_graph_url,
|
||||
watch=True,
|
||||
image="my-cool-image",
|
||||
)
|
||||
@@ -216,6 +233,16 @@ services:
|
||||
retries: 5
|
||||
interval: 60s
|
||||
start_interval: 1s
|
||||
langgraph-debugger:
|
||||
image: langchain/langgraph-debugger
|
||||
restart: on-failure
|
||||
depends_on:
|
||||
langgraph-postgres:
|
||||
condition: service_healthy
|
||||
ports:
|
||||
- "{debugger_port}:3968"
|
||||
environment:
|
||||
VITE_STUDIO_LOCAL_GRAPH_URL: {debugger_graph_url}
|
||||
langgraph-api:
|
||||
ports:
|
||||
- "8000:8000"
|
||||
@@ -232,7 +259,7 @@ services:
|
||||
test: python /api/healthcheck.py
|
||||
interval: 60s
|
||||
start_interval: 1s
|
||||
start_period: 60s
|
||||
start_period: 10s
|
||||
|
||||
|
||||
develop:
|
||||
@@ -262,82 +289,6 @@ def test_version_option() -> None:
|
||||
)
|
||||
|
||||
|
||||
def test_up_help_shows_hosted_studio_options() -> None:
|
||||
result = CliRunner().invoke(cli, ["up", "--help"])
|
||||
|
||||
assert result.exit_code == 0, result.output
|
||||
assert "--studio-url" in result.output
|
||||
assert "--api-url" in result.output
|
||||
assert "--debugger-port" not in result.output
|
||||
assert "--debugger-base-url" not in result.output
|
||||
|
||||
|
||||
def test_studio_link_defaults_to_hosted_studio() -> None:
|
||||
assert _studio_link(
|
||||
port=8123,
|
||||
studio_url=None,
|
||||
api_url=None,
|
||||
debugger_base_url=None,
|
||||
) == ("https://smith.langchain.com/studio/?baseUrl=http%3A%2F%2F127.0.0.1%3A8123")
|
||||
|
||||
|
||||
def test_studio_link_supports_self_hosted_and_remote_urls() -> None:
|
||||
assert _studio_link(
|
||||
port=8123,
|
||||
studio_url="https://langsmith.example.com/prefix/",
|
||||
api_url="https://api.example.com/graph?tenant=a®ion=eu",
|
||||
debugger_base_url=None,
|
||||
) == (
|
||||
"https://langsmith.example.com/prefix/studio/"
|
||||
"?baseUrl=https%3A%2F%2Fapi.example.com%2Fgraph%3Ftenant%3Da%26region%3Deu"
|
||||
)
|
||||
|
||||
|
||||
def test_studio_link_supports_deprecated_debugger_base_url(capsys) -> None:
|
||||
assert _studio_link(
|
||||
port=8123,
|
||||
studio_url=None,
|
||||
api_url=None,
|
||||
debugger_base_url="https://api.example.com",
|
||||
).endswith("?baseUrl=https%3A%2F%2Fapi.example.com")
|
||||
assert "--debugger-base-url is deprecated; use --api-url" in capsys.readouterr().err
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("studio_url", "api_url"),
|
||||
[
|
||||
("javascript:alert(1)", None),
|
||||
("https://user:password@example.com", None),
|
||||
("https://smith.langchain.com?workspace=test", None),
|
||||
(None, "file:///tmp/langgraph.sock"),
|
||||
(None, "https://user:password@example.com"),
|
||||
],
|
||||
)
|
||||
def test_studio_link_rejects_unsafe_urls(
|
||||
studio_url: str | None, api_url: str | None
|
||||
) -> None:
|
||||
with pytest.raises(click.UsageError):
|
||||
_studio_link(
|
||||
port=8123,
|
||||
studio_url=studio_url,
|
||||
api_url=api_url,
|
||||
debugger_base_url=None,
|
||||
)
|
||||
|
||||
|
||||
def test_studio_link_rejects_conflicting_api_url_aliases() -> None:
|
||||
with pytest.raises(
|
||||
click.UsageError,
|
||||
match="cannot specify different URLs",
|
||||
):
|
||||
_studio_link(
|
||||
port=8123,
|
||||
studio_url=None,
|
||||
api_url="https://api.example.com",
|
||||
debugger_base_url="https://other.example.com",
|
||||
)
|
||||
|
||||
|
||||
def test_top_level_help_shows_deploy_subcommands() -> None:
|
||||
runner = CliRunner()
|
||||
|
||||
|
||||
@@ -1403,19 +1403,6 @@ def test_config_to_docker_uv_lock():
|
||||
"COPY --from=uv-workspace-root uv.lock /tmp/uv_export/project/uv.lock"
|
||||
in docker
|
||||
)
|
||||
workspace_pyprojects = [
|
||||
"apps/agent/pyproject.toml",
|
||||
"libs/extra/pyproject.toml",
|
||||
"libs/shared/pyproject.toml",
|
||||
]
|
||||
export_instruction = "RUN uv export --package agent"
|
||||
for pyproject_path in workspace_pyprojects:
|
||||
copy_instruction = (
|
||||
"COPY --from=uv-workspace-root "
|
||||
f"{pyproject_path} /tmp/uv_export/project/{pyproject_path}"
|
||||
)
|
||||
assert copy_instruction in docker
|
||||
assert docker.index(copy_instruction) < docker.index(export_instruction)
|
||||
assert additional_contexts == {"uv-workspace-root": str(project_root.resolve())}
|
||||
|
||||
assert (
|
||||
|
||||
@@ -16,7 +16,7 @@ DEFAULT_DOCKER_CAPABILITIES = DockerCapabilities(
|
||||
)
|
||||
|
||||
|
||||
def test_compose_with_custom_db():
|
||||
def test_compose_with_no_debugger_and_custom_db():
|
||||
port = 8123
|
||||
custom_postgres_uri = "custom_postgres_uri"
|
||||
actual_compose_str = compose(
|
||||
@@ -42,7 +42,7 @@ def test_compose_with_custom_db():
|
||||
assert clean_empty_lines(actual_compose_str) == expected_compose_str
|
||||
|
||||
|
||||
def test_compose_with_custom_db_and_healthcheck():
|
||||
def test_compose_with_no_debugger_and_custom_db_with_healthcheck():
|
||||
port = 8123
|
||||
custom_postgres_uri = "custom_postgres_uri"
|
||||
actual_compose_str = compose(
|
||||
@@ -71,11 +71,39 @@ def test_compose_with_custom_db_and_healthcheck():
|
||||
test: python /api/healthcheck.py
|
||||
interval: 60s
|
||||
start_interval: 1s
|
||||
start_period: 60s"""
|
||||
start_period: 10s"""
|
||||
assert clean_empty_lines(actual_compose_str) == expected_compose_str
|
||||
|
||||
|
||||
def test_compose_with_default_db():
|
||||
def test_compose_with_debugger_and_custom_db():
|
||||
port = 8123
|
||||
custom_postgres_uri = "custom_postgres_uri"
|
||||
actual_compose_str = compose(
|
||||
DEFAULT_DOCKER_CAPABILITIES,
|
||||
port=port,
|
||||
postgres_uri=custom_postgres_uri,
|
||||
)
|
||||
expected_compose_str = f"""services:
|
||||
langgraph-redis:
|
||||
image: redis:6
|
||||
healthcheck:
|
||||
test: redis-cli ping
|
||||
interval: 5s
|
||||
timeout: 1s
|
||||
retries: 5
|
||||
langgraph-api:
|
||||
ports:
|
||||
- "{port}:8000"
|
||||
depends_on:
|
||||
langgraph-redis:
|
||||
condition: service_healthy
|
||||
environment:
|
||||
REDIS_URI: redis://langgraph-redis:6379
|
||||
POSTGRES_URI: {custom_postgres_uri}"""
|
||||
assert clean_empty_lines(actual_compose_str) == expected_compose_str
|
||||
|
||||
|
||||
def test_compose_with_debugger_and_default_db():
|
||||
port = 8123
|
||||
actual_compose_str = compose(DEFAULT_DOCKER_CAPABILITIES, port=port)
|
||||
expected_compose_str = f"""volumes:
|
||||
@@ -274,6 +302,72 @@ def test_compose_with_api_version_and_custom_postgres():
|
||||
assert clean_empty_lines(actual_compose_str) == expected_compose_str
|
||||
|
||||
|
||||
def test_compose_with_api_version_and_debugger():
|
||||
"""Test compose function with api_version and debugger port."""
|
||||
port = 8123
|
||||
debugger_port = 8001
|
||||
api_version = "0.2.74"
|
||||
|
||||
actual_compose_str = compose(
|
||||
DEFAULT_DOCKER_CAPABILITIES,
|
||||
port=port,
|
||||
api_version=api_version,
|
||||
debugger_port=debugger_port,
|
||||
)
|
||||
|
||||
expected_compose_str = f"""volumes:
|
||||
langgraph-data:
|
||||
driver: local
|
||||
services:
|
||||
langgraph-redis:
|
||||
image: redis:6
|
||||
healthcheck:
|
||||
test: redis-cli ping
|
||||
interval: 5s
|
||||
timeout: 1s
|
||||
retries: 5
|
||||
langgraph-postgres:
|
||||
image: pgvector/pgvector:pg16
|
||||
ports:
|
||||
- "5433:5432"
|
||||
environment:
|
||||
POSTGRES_DB: postgres
|
||||
POSTGRES_USER: postgres
|
||||
POSTGRES_PASSWORD: postgres
|
||||
command:
|
||||
- postgres
|
||||
- -c
|
||||
- shared_preload_libraries=vector
|
||||
volumes:
|
||||
- langgraph-data:/var/lib/postgresql/data
|
||||
healthcheck:
|
||||
test: pg_isready -U postgres
|
||||
start_period: 10s
|
||||
timeout: 1s
|
||||
retries: 5
|
||||
interval: 5s
|
||||
langgraph-debugger:
|
||||
image: langchain/langgraph-debugger
|
||||
restart: on-failure
|
||||
depends_on:
|
||||
langgraph-postgres:
|
||||
condition: service_healthy
|
||||
ports:
|
||||
- "{debugger_port}:3968"
|
||||
langgraph-api:
|
||||
ports:
|
||||
- "{port}:8000"
|
||||
depends_on:
|
||||
langgraph-redis:
|
||||
condition: service_healthy
|
||||
langgraph-postgres:
|
||||
condition: service_healthy
|
||||
environment:
|
||||
REDIS_URI: redis://langgraph-redis:6379
|
||||
POSTGRES_URI: {DEFAULT_POSTGRES_URI}"""
|
||||
assert clean_empty_lines(actual_compose_str) == expected_compose_str
|
||||
|
||||
|
||||
def test_compose_distributed_mode_with_custom_db():
|
||||
"""Test compose with engine_runtime_mode='distributed' adds N_JOBS_PER_WORKER=0."""
|
||||
port = 8123
|
||||
|
||||
Generated
+6
-6
@@ -266,20 +266,20 @@ wheels = [
|
||||
|
||||
[[package]]
|
||||
name = "langgraph-checkpoint"
|
||||
version = "4.2.0"
|
||||
version = "4.0.1"
|
||||
source = { registry = "https://pypi.org/simple" }
|
||||
dependencies = [
|
||||
{ name = "langchain-core" },
|
||||
{ name = "ormsgpack" },
|
||||
]
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/dc/e1/089c4c9e0a2fec7f883f82ae8e6a727138d50074cfeb6644bc2d13b1019b/langgraph_checkpoint-4.2.0.tar.gz", hash = "sha256:51a593b6bee684b0818e5d6e58e28ab340c6db7794575056ce7bd1b746a84ed7", size = 180239, upload-time = "2026-08-07T20:05:03.756Z" }
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/b1/44/a8df45d1e8b4637e29789fa8bae1db022c953cc7ac80093cfc52e923547e/langgraph_checkpoint-4.0.1.tar.gz", hash = "sha256:b433123735df11ade28829e40ce25b9be614930cd50245ff2af60629234befd9", size = 158135, upload-time = "2026-02-27T21:06:16.092Z" }
|
||||
wheels = [
|
||||
{ url = "https://files.pythonhosted.org/packages/05/71/3b475f09bd57d3a5649792c66353312b4432afd843f301739dfcebd157f0/langgraph_checkpoint-4.2.0-py3-none-any.whl", hash = "sha256:0547fd228935a0b758865de3a3d6d7a2537c308895d0f9ab092ce9151b5da942", size = 56833, upload-time = "2026-08-07T20:05:02.655Z" },
|
||||
{ url = "https://files.pythonhosted.org/packages/65/4c/09a4a0c42f5d2fc38d6c4d67884788eff7fd2cfdf367fdf7033de908b4c0/langgraph_checkpoint-4.0.1-py3-none-any.whl", hash = "sha256:e3adcd7a0e0166f3b48b8cf508ce0ea366e7420b5a73aa81289888727769b034", size = 50453, upload-time = "2026-02-27T21:06:14.293Z" },
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "langgraph-checkpoint-postgres"
|
||||
version = "3.1.1"
|
||||
version = "3.0.5"
|
||||
source = { registry = "https://pypi.org/simple" }
|
||||
dependencies = [
|
||||
{ name = "langgraph-checkpoint" },
|
||||
@@ -287,9 +287,9 @@ dependencies = [
|
||||
{ name = "psycopg" },
|
||||
{ name = "psycopg-pool" },
|
||||
]
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/06/92/1e8959f8cd1b56e672fde3227f6fd642be85af6c5fd662d73921074aa39d/langgraph_checkpoint_postgres-3.1.1.tar.gz", hash = "sha256:d320e147ddad8c374cd546df0b52b532dd54d0541dd9fd23fc738cbd5de76f41", size = 150413, upload-time = "2026-07-30T19:15:39.014Z" }
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/95/7a/8f439966643d32111248a225e6cb33a182d07c90de780c4dbfc1e0377832/langgraph_checkpoint_postgres-3.0.5.tar.gz", hash = "sha256:a8fd7278a63f4f849b5cbc7884a15ca8f41e7d5f7467d0a66b31e8c24492f7eb", size = 127856, upload-time = "2026-03-18T21:25:29.785Z" }
|
||||
wheels = [
|
||||
{ url = "https://files.pythonhosted.org/packages/03/32/ba457698a48a0e18d786caa770033067049fbe36d6846f8e50f13b594b51/langgraph_checkpoint_postgres-3.1.1-py3-none-any.whl", hash = "sha256:6e353aecd8150de144fef8e51a49076f58b7d6830d4cf51392b7ad4d79832ba7", size = 50778, upload-time = "2026-07-30T19:15:37.405Z" },
|
||||
{ url = "https://files.pythonhosted.org/packages/e8/87/b0f98b33a67204bca9d5619bcd9574222f6b025cf3c125eedcec9a50ecbc/langgraph_checkpoint_postgres-3.0.5-py3-none-any.whl", hash = "sha256:86d7040a88fd70087eaafb72251d796696a0a2d856168f5c11ef620771411552", size = 42907, upload-time = "2026-03-18T21:25:28.75Z" },
|
||||
]
|
||||
|
||||
[[package]]
|
||||
|
||||
@@ -1,10 +1,11 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import dis
|
||||
import ast
|
||||
import inspect
|
||||
import re
|
||||
import textwrap
|
||||
from collections.abc import Callable, Sequence
|
||||
from functools import partial
|
||||
from types import CodeType, FunctionType
|
||||
from typing import Any
|
||||
|
||||
from langchain_core.runnables import (
|
||||
@@ -16,6 +17,7 @@ from langchain_core.runnables import (
|
||||
from langchain_core.runnables.base import RunnableBindingBase
|
||||
from langchain_core.runnables.config import run_in_executor
|
||||
from langgraph.checkpoint.base import ChannelVersions
|
||||
from typing_extensions import override
|
||||
|
||||
from langgraph._internal._runnable import RunnableCallable, RunnableSeq
|
||||
from langgraph._internal._timeout import sync_timeout_unsupported
|
||||
@@ -135,87 +137,155 @@ def validate_timeout_supported(runnable: Runnable, *, name: str) -> None:
|
||||
raise sync_timeout_unsupported(name)
|
||||
|
||||
|
||||
# Values treated as dead ends when deciding whether to walk a function's
|
||||
# bytecode. A container can hold a graph, but `find_subgraph_pregel` does not
|
||||
# look inside one, so skipping it costs nothing while that holds. Matched by
|
||||
# exact type, since a subclass of a builtin can carry attributes.
|
||||
_LEAF_TYPES = frozenset(
|
||||
{
|
||||
int,
|
||||
float,
|
||||
complex,
|
||||
bool,
|
||||
str,
|
||||
bytes,
|
||||
bytearray,
|
||||
list,
|
||||
tuple,
|
||||
dict,
|
||||
set,
|
||||
frozenset,
|
||||
type(None),
|
||||
}
|
||||
)
|
||||
|
||||
|
||||
def get_function_nonlocals(func: Callable) -> list[Any]:
|
||||
"""Get the values a function reaches from outside its own scope.
|
||||
"""Get the nonlocal variables accessed by a function.
|
||||
|
||||
Args:
|
||||
func: The function to check.
|
||||
|
||||
Returns:
|
||||
Every captured cell value, the globals the function names, and each
|
||||
value along an attribute path it loads. Over-approximates: a value can
|
||||
come back without the function reaching it at runtime.
|
||||
List[Any]: The nonlocal variables accessed by the function.
|
||||
"""
|
||||
func = getattr(func, "__func__", func) # bound method -> function
|
||||
wrapped = getattr(func, "__wrapped__", None)
|
||||
if callable(wrapped):
|
||||
func = getattr(wrapped, "__func__", wrapped)
|
||||
if not isinstance(func, FunctionType):
|
||||
try:
|
||||
code = inspect.getsource(func)
|
||||
tree = ast.parse(textwrap.dedent(code))
|
||||
visitor = FunctionNonLocals()
|
||||
visitor.visit(tree)
|
||||
values: list[Any] = []
|
||||
closure = (
|
||||
inspect.getclosurevars(func.__wrapped__)
|
||||
if hasattr(func, "__wrapped__") and callable(func.__wrapped__)
|
||||
else inspect.getclosurevars(func)
|
||||
)
|
||||
candidates = {**closure.globals, **closure.nonlocals}
|
||||
for k, v in candidates.items():
|
||||
if k in visitor.nonlocals:
|
||||
values.append(v)
|
||||
for kk in visitor.nonlocals:
|
||||
if "." in kk and kk.startswith(k):
|
||||
vv = v
|
||||
for part in kk.split(".")[1:]:
|
||||
if vv is None:
|
||||
break
|
||||
else:
|
||||
try:
|
||||
vv = getattr(vv, part)
|
||||
except AttributeError:
|
||||
break
|
||||
else:
|
||||
values.append(vv)
|
||||
except (SyntaxError, TypeError, OSError, SystemError):
|
||||
return []
|
||||
code = func.__code__
|
||||
|
||||
cells: dict[str, Any] = {}
|
||||
for name, cell in zip(code.co_freevars, func.__closure__ or ()):
|
||||
try:
|
||||
cells[name] = cell.cell_contents
|
||||
except ValueError:
|
||||
continue # empty cell: a recursive def not yet bound
|
||||
|
||||
# Every captured value counts, referenced or not: over-declaring costs an
|
||||
# introspection entry, under-declaring drops the subgraph's checkpoints and
|
||||
# stream events. Checking each cell against the bytecode would cost more and
|
||||
# only trade the cheap error for the expensive one.
|
||||
values: list[Any] = list(cells.values())
|
||||
global_ns = func.__globals__
|
||||
globals_ = {name: global_ns[name] for name in code.co_names if name in global_ns}
|
||||
if all(type(v) in _LEAF_TYPES for v in (*cells.values(), *globals_.values())):
|
||||
return values
|
||||
|
||||
# Nested code objects hold the references made by inner defs, lambdas and
|
||||
# comprehensions, which resolve against the namespaces gathered above.
|
||||
codes = [code]
|
||||
for c in codes:
|
||||
codes.extend(k for k in c.co_consts if isinstance(k, CodeType))
|
||||
value: Any = None
|
||||
for instruction in dis.get_instructions(c):
|
||||
opname = instruction.opname
|
||||
if opname == "LOAD_GLOBAL":
|
||||
value = globals_.get(instruction.argval)
|
||||
elif opname == "LOAD_DEREF":
|
||||
value = cells.get(instruction.argval)
|
||||
elif opname in ("LOAD_ATTR", "LOAD_METHOD"):
|
||||
value = getattr(value, instruction.argval, None)
|
||||
else:
|
||||
value = None # anything else ends the chain: `a, b.c` is not `a.c`
|
||||
continue
|
||||
if value is not None:
|
||||
values.append(value)
|
||||
return values
|
||||
|
||||
|
||||
class FunctionNonLocals(ast.NodeVisitor):
|
||||
"""Get the nonlocal variables accessed of a function."""
|
||||
|
||||
def __init__(self) -> None:
|
||||
self.nonlocals: set[str] = set()
|
||||
|
||||
@override
|
||||
def visit_FunctionDef(self, node: ast.FunctionDef) -> Any:
|
||||
"""Visit a function definition.
|
||||
|
||||
Args:
|
||||
node: The node to visit.
|
||||
|
||||
Returns:
|
||||
Any: The result of the visit.
|
||||
"""
|
||||
visitor = NonLocals()
|
||||
visitor.visit(node)
|
||||
self.nonlocals.update(visitor.loads - visitor.stores)
|
||||
|
||||
@override
|
||||
def visit_AsyncFunctionDef(self, node: ast.AsyncFunctionDef) -> Any:
|
||||
"""Visit an async function definition.
|
||||
|
||||
Args:
|
||||
node: The node to visit.
|
||||
|
||||
Returns:
|
||||
Any: The result of the visit.
|
||||
"""
|
||||
visitor = NonLocals()
|
||||
visitor.visit(node)
|
||||
self.nonlocals.update(visitor.loads - visitor.stores)
|
||||
|
||||
@override
|
||||
def visit_Lambda(self, node: ast.Lambda) -> Any:
|
||||
"""Visit a lambda function.
|
||||
|
||||
Args:
|
||||
node: The node to visit.
|
||||
|
||||
Returns:
|
||||
Any: The result of the visit.
|
||||
"""
|
||||
visitor = NonLocals()
|
||||
visitor.visit(node)
|
||||
self.nonlocals.update(visitor.loads - visitor.stores)
|
||||
|
||||
|
||||
class NonLocals(ast.NodeVisitor):
|
||||
"""Get nonlocal variables accessed."""
|
||||
|
||||
def __init__(self) -> None:
|
||||
self.loads: set[str] = set()
|
||||
self.stores: set[str] = set()
|
||||
|
||||
@override
|
||||
def visit_Name(self, node: ast.Name) -> Any:
|
||||
"""Visit a name node.
|
||||
|
||||
Args:
|
||||
node: The node to visit.
|
||||
|
||||
Returns:
|
||||
Any: The result of the visit.
|
||||
"""
|
||||
if isinstance(node.ctx, ast.Load):
|
||||
self.loads.add(node.id)
|
||||
elif isinstance(node.ctx, ast.Store):
|
||||
self.stores.add(node.id)
|
||||
|
||||
@override
|
||||
def visit_Attribute(self, node: ast.Attribute) -> Any:
|
||||
"""Visit an attribute node.
|
||||
|
||||
Args:
|
||||
node: The node to visit.
|
||||
|
||||
Returns:
|
||||
Any: The result of the visit.
|
||||
"""
|
||||
if isinstance(node.ctx, ast.Load):
|
||||
parent = node.value
|
||||
attr_expr = node.attr
|
||||
while isinstance(parent, ast.Attribute):
|
||||
attr_expr = parent.attr + "." + attr_expr
|
||||
parent = parent.value
|
||||
if isinstance(parent, ast.Name):
|
||||
self.loads.add(parent.id + "." + attr_expr)
|
||||
self.loads.discard(parent.id)
|
||||
elif isinstance(parent, ast.Call):
|
||||
if isinstance(parent.func, ast.Name):
|
||||
self.loads.add(parent.func.id)
|
||||
else:
|
||||
parent = parent.func
|
||||
attr_expr = ""
|
||||
while isinstance(parent, ast.Attribute):
|
||||
if attr_expr:
|
||||
attr_expr = parent.attr + "." + attr_expr
|
||||
else:
|
||||
attr_expr = parent.attr
|
||||
parent = parent.value
|
||||
if isinstance(parent, ast.Name):
|
||||
self.loads.add(parent.id + "." + attr_expr)
|
||||
|
||||
|
||||
def is_xxh3_128_hexdigest(value: str) -> bool:
|
||||
"""Check if the given string matches the format of xxh3_128_hexdigest."""
|
||||
return bool(re.fullmatch(r"[0-9a-f]{32}", value))
|
||||
|
||||
@@ -1,286 +0,0 @@
|
||||
"""Tests for subgraph auto-detection (`pregel/_utils.py`).
|
||||
|
||||
Detection failing is silent — the graph still runs, only introspection goes
|
||||
quiet — so every shape a node can hold a graph in is pinned here. The expected
|
||||
values are what the source-parsing implementation this replaced produced for
|
||||
the same shapes, except for `sourceless`, whose source it could not read,
|
||||
`empty_closure_cell`, on which it raised, and `unreachable_attribute_chain`,
|
||||
where it reported a graph that dropped code could never invoke.
|
||||
"""
|
||||
|
||||
import functools
|
||||
import operator
|
||||
from typing import Annotated, Any
|
||||
|
||||
import pytest
|
||||
from typing_extensions import TypedDict
|
||||
|
||||
from langgraph.graph import END, START, StateGraph
|
||||
from langgraph.pregel._utils import get_function_nonlocals
|
||||
|
||||
|
||||
class State(TypedDict):
|
||||
log: Annotated[list, operator.add]
|
||||
|
||||
|
||||
def _leaf(tag: str) -> Any:
|
||||
"""Return a compiled graph that reports itself as `tag`."""
|
||||
builder = StateGraph(State)
|
||||
builder.add_node(tag, lambda s: {"log": [tag]})
|
||||
builder.add_edge(START, tag)
|
||||
builder.add_edge(tag, END)
|
||||
compiled = builder.compile()
|
||||
compiled.name = tag
|
||||
return compiled
|
||||
|
||||
|
||||
def _detect(node: Any) -> str | None:
|
||||
"""Return the name of the subgraph detected for `node`, or None."""
|
||||
builder = StateGraph(State)
|
||||
builder.add_node("n", node)
|
||||
builder.add_edge(START, "n")
|
||||
builder.add_edge("n", END)
|
||||
subgraphs = builder.compile().nodes["n"].subgraphs
|
||||
return getattr(subgraphs[0], "name", "?") if subgraphs else None
|
||||
|
||||
|
||||
class _Box:
|
||||
def __init__(self, payload: Any) -> None:
|
||||
self.payload = payload
|
||||
|
||||
|
||||
class _ListSubclass(list):
|
||||
pass
|
||||
|
||||
|
||||
class _MethodHolder:
|
||||
def __init__(self) -> None:
|
||||
self.graph = _leaf("via_self")
|
||||
|
||||
def as_node(self, state: State) -> Any:
|
||||
return self.graph.invoke(state)
|
||||
|
||||
|
||||
MODULE_GRAPH = _leaf("module_global")
|
||||
CHAIN = _Box(_Box(_leaf("attr_chain")))
|
||||
GRAPH_IN_PLAIN_LIST = [_leaf("in_list")]
|
||||
METHOD_HOLDER = _MethodHolder()
|
||||
|
||||
|
||||
def closure_capture() -> Any:
|
||||
sub = _leaf("closure")
|
||||
|
||||
def node(state: State) -> Any:
|
||||
return sub.invoke(state)
|
||||
|
||||
return node
|
||||
|
||||
|
||||
def module_global() -> Any:
|
||||
def node(state: State) -> Any:
|
||||
return MODULE_GRAPH.invoke(state)
|
||||
|
||||
return node
|
||||
|
||||
|
||||
def attribute_chain() -> Any:
|
||||
def node(state: State) -> Any:
|
||||
return CHAIN.payload.payload.invoke(state)
|
||||
|
||||
return node
|
||||
|
||||
|
||||
def nested_def_captured_attribute() -> Any:
|
||||
"""A chain on a captured holder, named only inside a nested code object.
|
||||
|
||||
The captured value is the holder, not the graph, so the chain itself has to
|
||||
be recovered from the nested scope.
|
||||
"""
|
||||
holder = _Box(_leaf("nested_captured"))
|
||||
|
||||
def node(state: State) -> Any:
|
||||
def inner() -> Any:
|
||||
return holder.payload.invoke(state)
|
||||
|
||||
return inner()
|
||||
|
||||
return node
|
||||
|
||||
|
||||
def unreachable_branch() -> Any:
|
||||
"""A captured graph referenced only from code the compiler removes."""
|
||||
sub = _leaf("unreachable")
|
||||
|
||||
def node(state: State) -> Any:
|
||||
if False:
|
||||
sub.invoke(state)
|
||||
return {"log": []}
|
||||
|
||||
return node
|
||||
|
||||
|
||||
def unreachable_attribute_chain() -> Any:
|
||||
"""A graph named only along an attribute path the compiler dropped.
|
||||
|
||||
The closure keeps `holder`, but the `.payload` load is gone. The source
|
||||
parser reported this one; dropped code cannot invoke anything, so that was
|
||||
a phantom rather than a detection.
|
||||
"""
|
||||
holder = _Box(_leaf("unreachable_attr"))
|
||||
|
||||
def node(state: State) -> Any:
|
||||
if False:
|
||||
holder.payload.invoke(state)
|
||||
return {"log": []}
|
||||
|
||||
return node
|
||||
|
||||
|
||||
def wrapper_referencing_nothing() -> Any:
|
||||
"""A wrapper whose own scope holds nothing, so only `__wrapped__` leads on.
|
||||
|
||||
`functools.wraps` would leave the wrapper closing over the inner function;
|
||||
setting the attribute by hand does not.
|
||||
"""
|
||||
sub = _leaf("via_wrapped")
|
||||
|
||||
def inner(state: State) -> Any:
|
||||
return sub.invoke(state)
|
||||
|
||||
def wrapper(state: State) -> Any:
|
||||
return {"log": []}
|
||||
|
||||
wrapper.__wrapped__ = inner
|
||||
return wrapper
|
||||
|
||||
|
||||
def captured_list_subclass() -> Any:
|
||||
"""A `list` subclass is not a leaf: it can carry a graph as an attribute."""
|
||||
holder = _ListSubclass()
|
||||
holder.payload = _leaf("list_subclass")
|
||||
|
||||
def node(state: State) -> Any:
|
||||
return holder.payload.invoke(state)
|
||||
|
||||
return node
|
||||
|
||||
|
||||
def empty_closure_cell() -> Any:
|
||||
"""An unassigned closure variable leaves a cell that cannot be read."""
|
||||
sub = _leaf("beside_empty_cell")
|
||||
|
||||
def node(state: State) -> Any:
|
||||
return unassigned, sub.invoke(state)
|
||||
|
||||
return node
|
||||
unassigned = 1 # never runs, so the cell it creates is never filled
|
||||
|
||||
|
||||
def sourceless() -> Any:
|
||||
"""A node compiled without a source file, which `getsource` could not read."""
|
||||
namespace: dict[str, Any] = {"SOURCELESS": _leaf("sourceless")}
|
||||
exec(
|
||||
compile(
|
||||
"def node(state):\n return SOURCELESS.invoke(state)", "<test>", "exec"
|
||||
),
|
||||
namespace,
|
||||
)
|
||||
return namespace["node"]
|
||||
|
||||
|
||||
async def _async_node(state: State) -> Any:
|
||||
return await MODULE_GRAPH.ainvoke(state)
|
||||
|
||||
|
||||
def async_node() -> Any:
|
||||
return _async_node
|
||||
|
||||
|
||||
def no_subgraph() -> Any:
|
||||
"""Nothing but leaf values in reach, so the bytecode walk is skipped."""
|
||||
|
||||
def node(state: State) -> Any:
|
||||
return {"log": [len("abc") + 1]}
|
||||
|
||||
return node
|
||||
|
||||
|
||||
def recombined_names() -> Any:
|
||||
"""Loads `CHAIN.payload` and `local.payload`, never `CHAIN.payload.payload`."""
|
||||
|
||||
def node(state: State) -> Any:
|
||||
local = _Box("not a graph")
|
||||
return {"log": [CHAIN.payload, local.payload]}
|
||||
|
||||
return node
|
||||
|
||||
|
||||
def broken_attribute_chain() -> Any:
|
||||
holder = _Box("a string, so `.payload.missing` cannot resolve")
|
||||
|
||||
def node(state: State) -> Any:
|
||||
return holder.payload.missing.invoke(state)
|
||||
|
||||
return node
|
||||
|
||||
|
||||
def nested_def_global() -> Any:
|
||||
"""A global named only in a nested code object: out of reach, as before."""
|
||||
|
||||
def node(state: State) -> Any:
|
||||
def inner() -> Any:
|
||||
return MODULE_GRAPH.invoke(state)
|
||||
|
||||
return inner()
|
||||
|
||||
return node
|
||||
|
||||
|
||||
def graph_in_plain_list() -> Any:
|
||||
def node(state: State) -> Any:
|
||||
return GRAPH_IN_PLAIN_LIST[0].invoke(state)
|
||||
|
||||
return node
|
||||
|
||||
|
||||
def bound_method_self() -> Any:
|
||||
return METHOD_HOLDER.as_node
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("factory", "expected"),
|
||||
[
|
||||
(closure_capture, "closure"),
|
||||
(module_global, "module_global"),
|
||||
(attribute_chain, "attr_chain"),
|
||||
(nested_def_captured_attribute, "nested_captured"),
|
||||
(unreachable_branch, "unreachable"),
|
||||
(wrapper_referencing_nothing, "via_wrapped"),
|
||||
(captured_list_subclass, "list_subclass"),
|
||||
(empty_closure_cell, "beside_empty_cell"),
|
||||
(sourceless, "sourceless"),
|
||||
(async_node, "module_global"),
|
||||
# Shapes no reference chain reaches: a subscript, an instance attribute
|
||||
# of `self`, a global named only in a nested scope, and an attribute
|
||||
# path the compiler dropped.
|
||||
(no_subgraph, None),
|
||||
(recombined_names, None),
|
||||
(broken_attribute_chain, None),
|
||||
(nested_def_global, None),
|
||||
(graph_in_plain_list, None),
|
||||
(bound_method_self, None),
|
||||
(unreachable_attribute_chain, None),
|
||||
],
|
||||
ids=lambda value: value.__name__ if callable(value) else str(value),
|
||||
)
|
||||
def test_subgraph_detection(factory: Any, expected: str | None) -> None:
|
||||
assert _detect(factory()) == expected
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
"candidate",
|
||||
[functools.partial(lambda state, extra: {"log": [extra]}, extra="x"), len],
|
||||
ids=["partial", "builtin"],
|
||||
)
|
||||
def test_callables_without_a_code_object_are_handled(candidate: Any) -> None:
|
||||
assert get_function_nonlocals(candidate) == []
|
||||
@@ -1,15 +1,8 @@
|
||||
from langgraph_sdk.auth import Auth
|
||||
from langgraph_sdk.client import get_client, get_sync_client
|
||||
from langgraph_sdk.encryption import Encryption
|
||||
from langgraph_sdk.encryption.types import DecryptResult, EncryptionContext
|
||||
from langgraph_sdk.encryption.types import EncryptionContext
|
||||
|
||||
__version__ = "0.4.3"
|
||||
__version__ = "0.4.2"
|
||||
|
||||
__all__ = [
|
||||
"Auth",
|
||||
"DecryptResult",
|
||||
"Encryption",
|
||||
"EncryptionContext",
|
||||
"get_client",
|
||||
"get_sync_client",
|
||||
]
|
||||
__all__ = ["Auth", "Encryption", "EncryptionContext", "get_client", "get_sync_client"]
|
||||
|
||||
@@ -18,9 +18,6 @@ import warnings
|
||||
|
||||
from langgraph_sdk.encryption import types
|
||||
|
||||
_BlobDecryptorT = typing.TypeVar("_BlobDecryptorT", bound=types.BlobDecryptor)
|
||||
_JsonDecryptorT = typing.TypeVar("_JsonDecryptorT", bound=types.JsonDecryptor)
|
||||
|
||||
|
||||
class LangGraphBetaWarning(UserWarning):
|
||||
"""Warning for beta features in LangGraph SDK."""
|
||||
@@ -144,7 +141,7 @@ class _DecryptDecorators:
|
||||
def __init__(self, parent: Encryption):
|
||||
self._parent = parent
|
||||
|
||||
def blob(self, fn: _BlobDecryptorT) -> _BlobDecryptorT:
|
||||
def blob(self, fn: types.BlobDecryptor) -> types.BlobDecryptor:
|
||||
"""Register a blob decryption handler.
|
||||
|
||||
The handler will be called to decrypt opaque data like checkpoint blobs.
|
||||
@@ -152,9 +149,7 @@ class _DecryptDecorators:
|
||||
Example:
|
||||
```python
|
||||
@encryption.decrypt.blob
|
||||
async def decrypt_blob(
|
||||
ctx: EncryptionContext, blob: bytes
|
||||
) -> bytes | DecryptResult[bytes]:
|
||||
async def decrypt_blob(ctx: EncryptionContext, blob: bytes) -> bytes:
|
||||
# Decrypt the blob using your encryption service
|
||||
return decrypted_blob
|
||||
```
|
||||
@@ -175,15 +170,13 @@ class _DecryptDecorators:
|
||||
self._parent._blob_decryptor = fn
|
||||
return fn
|
||||
|
||||
def json(self, fn: _JsonDecryptorT) -> _JsonDecryptorT:
|
||||
def json(self, fn: types.JsonDecryptor) -> types.JsonDecryptor:
|
||||
"""Register the JSON decryption handler.
|
||||
|
||||
Example:
|
||||
```python
|
||||
@encryption.decrypt.json
|
||||
async def decrypt_json(
|
||||
ctx: EncryptionContext, data: dict
|
||||
) -> dict | DecryptResult[dict]:
|
||||
async def decrypt_json(ctx: EncryptionContext, data: dict) -> dict:
|
||||
# Decrypt the data
|
||||
return decrypt_data(data)
|
||||
```
|
||||
@@ -376,7 +369,7 @@ class Encryption:
|
||||
"""Reference to encryption type definitions.
|
||||
|
||||
Provides access to all type definitions used in the encryption system,
|
||||
including EncryptionContext, DecryptResult, BlobEncryptor, BlobDecryptor,
|
||||
including EncryptionContext, BlobEncryptor, BlobDecryptor,
|
||||
JsonEncryptor, and JsonDecryptor.
|
||||
"""
|
||||
|
||||
|
||||
@@ -9,30 +9,10 @@ from __future__ import annotations
|
||||
|
||||
import typing
|
||||
from collections.abc import Awaitable, Callable
|
||||
from dataclasses import dataclass
|
||||
|
||||
Json = dict[str, typing.Any]
|
||||
"""JSON-serializable dictionary type for structured data encryption."""
|
||||
|
||||
T = typing.TypeVar("T")
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class DecryptResult(typing.Generic[T]):
|
||||
"""Decrypted data and optional replacement ciphertext.
|
||||
|
||||
Return this from a decrypt handler when encrypted data should be replaced,
|
||||
such as after rotating its encryption key. Returning plaintext directly
|
||||
remains supported when no replacement is needed.
|
||||
|
||||
Attributes:
|
||||
plaintext: Decrypted data returned to the caller
|
||||
replacement: New encrypted data to persist in place of the input
|
||||
"""
|
||||
|
||||
plaintext: T
|
||||
replacement: T | None = None
|
||||
|
||||
|
||||
class EncryptionContext:
|
||||
"""Context passed to encryption/decryption handlers.
|
||||
@@ -77,9 +57,7 @@ Returns:
|
||||
Awaitable that resolves to encrypted bytes
|
||||
"""
|
||||
|
||||
BlobDecryptor = Callable[
|
||||
[EncryptionContext, bytes], Awaitable[bytes | DecryptResult[bytes]]
|
||||
]
|
||||
BlobDecryptor = Callable[[EncryptionContext, bytes], Awaitable[bytes]]
|
||||
"""Handler for decrypting opaque blob data like checkpoints.
|
||||
|
||||
Note: Must be an async function. Decryption typically involves I/O operations
|
||||
@@ -90,8 +68,7 @@ Args:
|
||||
blob: The encrypted bytes to decrypt
|
||||
|
||||
Returns:
|
||||
Awaitable that resolves to decrypted bytes, or a DecryptResult containing
|
||||
decrypted bytes and replacement ciphertext
|
||||
Awaitable that resolves to decrypted bytes
|
||||
"""
|
||||
|
||||
JsonEncryptor = Callable[[EncryptionContext, Json], Awaitable[Json]]
|
||||
@@ -124,9 +101,7 @@ Returns:
|
||||
Awaitable that resolves to encrypted JSON dictionary
|
||||
"""
|
||||
|
||||
JsonDecryptor = Callable[
|
||||
[EncryptionContext, Json], Awaitable[Json | DecryptResult[Json]]
|
||||
]
|
||||
JsonDecryptor = Callable[[EncryptionContext, Json], Awaitable[Json]]
|
||||
"""Handler for decrypting structured JSON data.
|
||||
|
||||
Note: Must be an async function. Decryption typically involves I/O operations
|
||||
@@ -140,8 +115,7 @@ Args:
|
||||
data: The encrypted JSON dictionary
|
||||
|
||||
Returns:
|
||||
Awaitable that resolves to a decrypted JSON dictionary, or a DecryptResult
|
||||
containing decrypted JSON and replacement ciphertext
|
||||
Awaitable that resolves to decrypted JSON dictionary
|
||||
"""
|
||||
|
||||
if typing.TYPE_CHECKING:
|
||||
|
||||
@@ -1,40 +1,8 @@
|
||||
from collections.abc import Awaitable, Callable
|
||||
|
||||
import pytest
|
||||
|
||||
from langgraph_sdk import DecryptResult, EncryptionContext
|
||||
from langgraph_sdk.encryption import DuplicateHandlerError, Encryption
|
||||
|
||||
|
||||
def test_decrypt_result():
|
||||
result = DecryptResult(plaintext=b"plain", replacement=b"rotated")
|
||||
|
||||
assert result.plaintext == b"plain"
|
||||
assert result.replacement == b"rotated"
|
||||
assert DecryptResult(plaintext={"plain": True}).replacement is None
|
||||
|
||||
|
||||
def test_decrypt_decorators_preserve_return_types():
|
||||
encryption = Encryption()
|
||||
|
||||
@encryption.decrypt.blob
|
||||
async def blob_dec(_ctx: EncryptionContext, data: bytes) -> bytes:
|
||||
return data
|
||||
|
||||
@encryption.decrypt.json
|
||||
async def json_dec(
|
||||
_ctx: EncryptionContext, data: dict[str, object]
|
||||
) -> dict[str, object]:
|
||||
return data
|
||||
|
||||
blob_handler: Callable[[EncryptionContext, bytes], Awaitable[bytes]] = blob_dec
|
||||
json_handler: Callable[
|
||||
[EncryptionContext, dict[str, object]], Awaitable[dict[str, object]]
|
||||
] = json_dec
|
||||
assert blob_handler is blob_dec
|
||||
assert json_handler is json_dec
|
||||
|
||||
|
||||
class TestHandlerValidation:
|
||||
"""Test duplicate handler and signature validation."""
|
||||
|
||||
|
||||
Reference in New Issue
Block a user