Compare commits

..
Author SHA1 Message Date
hari-dhanushkodiandopen-swe[bot] <open-swe@users.noreply.github.com> f3de8ed373 fix(cli): add monorepo commands to dockerfile
Co-authored-by: open-swe[bot] <open-swe@users.noreply.github.com>
2026-07-20 17:26:01 +00:00
14 changed files with 153 additions and 392 deletions
@@ -238,12 +238,6 @@ class BasePostgresStore(Generic[C]):
conn: C
_deserializer: Callable[[bytes | orjson.Fragment], dict[str, Any]] | None
index_config: PostgresIndexConfig | None
ttl_config: TTLConfig | None
@property
def _omit_expired(self) -> bool:
"""Whether expired-but-unswept rows should be filtered from reads."""
return bool(self.ttl_config and self.ttl_config.get("omit_expired"))
def _get_batch_GET_ops_queries(
self,
@@ -264,17 +258,12 @@ class BasePostgresStore(Generic[C]):
namespace_groups[op.namespace].append((idx, op.key))
refresh_ttls[op.namespace].append(op.refresh_ttl)
omit_expired = self._omit_expired
expiry_clause = (
"AND (s.expires_at IS NULL OR s.expires_at > NOW())" if omit_expired else ""
)
results = []
for namespace, items in namespace_groups.items():
_, keys = zip(*items, strict=False)
this_refresh_ttls = refresh_ttls[namespace]
query = f"""
query = """
WITH passed_in AS (
SELECT unnest(%s::text[]) AS key,
unnest(%s::bool[]) AS do_refresh
@@ -287,14 +276,12 @@ class BasePostgresStore(Generic[C]):
AND s.key = p.key
AND p.do_refresh = TRUE
AND s.ttl_minutes IS NOT NULL
{expiry_clause}
RETURNING s.key
)
SELECT s.key, s.value, s.created_at, s.updated_at
FROM store s
JOIN passed_in p ON s.key = p.key
WHERE s.prefix = %s
{expiry_clause}
"""
ns_text = _namespace_to_text(namespace)
params = (
@@ -435,13 +422,6 @@ class BasePostgresStore(Generic[C]):
- embedding_requests: list of (original_index_in_search_ops, text_query)
"""
omit_expired = self._omit_expired
search_expiry_clause = (
"AND (store.expires_at IS NULL OR store.expires_at > NOW())"
if omit_expired
else ""
)
queries = []
embedding_requests = []
for idx, (_, op) in enumerate(search_ops):
@@ -511,7 +491,7 @@ class BasePostgresStore(Generic[C]):
{score_operator} AS neg_score
FROM store
JOIN store_vectors sv ON store.prefix = sv.prefix AND store.key = sv.key
WHERE {ns_condition} {extra_filters} {search_expiry_clause}
WHERE {ns_condition} {extra_filters}
ORDER BY {score_operator} ASC
LIMIT %s
"""
@@ -547,7 +527,7 @@ class BasePostgresStore(Generic[C]):
base_query = f"""
SELECT store.prefix, store.key, store.value, store.created_at, store.updated_at, NULL AS score
FROM store
WHERE {ns_condition} {extra_filters} {search_expiry_clause}
WHERE {ns_condition} {extra_filters}
ORDER BY store.updated_at DESC
LIMIT %s
OFFSET %s
@@ -611,10 +591,7 @@ class BasePostgresStore(Generic[C]):
"""
params: list[Any] = [op.max_depth, op.max_depth]
omit_expired = self._omit_expired
conditions = []
if omit_expired:
conditions.append("(expires_at IS NULL OR expires_at > NOW())")
if op.match_conditions:
for condition in op.match_conditions:
if condition.match_type == "prefix":
@@ -739,135 +739,3 @@ async def test_store_ttl(store):
# Now has been (TTL_SECONDS-2)*2 > TTL_SECONDS + TTL_SECONDS/2
results = await store.asearch(ns, query="bar", refresh_ttl=False)
assert len(results) == 0
async def _aexpire_now(
store: AsyncPostgresStore, ns: tuple[str, ...], key: str
) -> None:
"""Backdate a row's expires_at into the past without deleting it (unswept)."""
async with store._cursor() as cur:
await cur.execute(
"UPDATE store SET expires_at = NOW() - INTERVAL '1 minute' "
"WHERE prefix = %s AND key = %s",
(".".join(ns), key),
)
async def _arow_exists(
store: AsyncPostgresStore, ns: tuple[str, ...], key: str
) -> bool:
async with store._cursor() as cur:
await cur.execute(
"SELECT COUNT(*) AS n FROM store WHERE prefix = %s AND key = %s",
(".".join(ns), key),
)
return (await cur.fetchone())["n"] == 1
async def _astored_expires_at(store: AsyncPostgresStore, ns: tuple[str, ...], key: str):
async with store._cursor() as cur:
await cur.execute(
"SELECT expires_at FROM store WHERE prefix = %s AND key = %s",
(".".join(ns), key),
)
return (await cur.fetchone())["expires_at"]
async def test_omit_expired_filters_read_paths(store: AsyncPostgresStore) -> None:
await store.stop_ttl_sweeper() # deterministic: no background deletion
store.ttl_config["omit_expired"] = True
expired_ns = ("omit", "expired")
control_ns = ("omit", "control")
await store.aput(expired_ns, "e", {"data": "gone"}, ttl=TTL_MINUTES)
await store.aput(control_ns, "c", {"data": "keep"}, ttl=None)
await _aexpire_now(store, expired_ns, "e")
# The row is expired but physically still present (unswept).
assert await _arow_exists(store, expired_ns, "e")
# aget omits it; the never-expiring control is still returned.
assert await store.aget(expired_ns, "e") is None
assert await store.aget(control_ns, "c") is not None
# asearch omits it but returns the control.
assert await store.asearch(expired_ns) == []
assert [i.key for i in await store.asearch(control_ns)] == ["c"]
# alist_namespaces drops the expired-only namespace, keeps the control.
namespaces = await store.alist_namespaces(prefix=("omit",))
assert expired_ns not in namespaces
assert control_ns in namespaces
@pytest.mark.parametrize("omit", [None, False], ids=["default", "explicit-false"])
async def test_omit_expired_disabled_preserves_expired_rows(
store: AsyncPostgresStore, omit
) -> None:
await store.stop_ttl_sweeper()
if omit is not None:
store.ttl_config["omit_expired"] = omit
ns = ("keep",)
await store.aput(ns, "k", {"data": "still-here"}, ttl=TTL_MINUTES)
await _aexpire_now(store, ns, "k")
assert await store.aget(ns, "k", refresh_ttl=False) is not None
assert [i.key for i in await store.asearch(ns, refresh_ttl=False)] == ["k"]
assert ns in await store.alist_namespaces(prefix=("keep",))
async def test_omit_expired_refresh_ttl_only_refreshes_live_rows(
store: AsyncPostgresStore,
) -> None:
await store.stop_ttl_sweeper()
store.ttl_config["omit_expired"] = True
ns = ("refresh",)
await store.aput(ns, "expired", {"n": 0}, ttl=TTL_MINUTES)
await store.aput(ns, "live_get", {"n": 1}, ttl=TTL_MINUTES)
await store.aput(ns, "live_search", {"n": 2}, ttl=TTL_MINUTES)
await _aexpire_now(store, ns, "expired")
expired_before = await _astored_expires_at(store, ns, "expired")
get_before = await _astored_expires_at(store, ns, "live_get")
search_before = await _astored_expires_at(store, ns, "live_search")
# refresh_ttl=True must NOT resurrect the expired row (via aget or asearch)...
assert await store.aget(ns, "expired", refresh_ttl=True) is None
live_keys = [i.key for i in await store.asearch(ns, refresh_ttl=True)]
assert "expired" not in live_keys
assert await _astored_expires_at(store, ns, "expired") == expired_before
# ...but must still extend the live rows that were read.
assert await store.aget(ns, "live_get", refresh_ttl=True) is not None
assert await _astored_expires_at(store, ns, "live_get") > get_before
assert await _astored_expires_at(store, ns, "live_search") > search_before
async def test_omit_expired_search_pagination(store: AsyncPostgresStore) -> None:
await store.stop_ttl_sweeper()
store.ttl_config["omit_expired"] = True
ns = ("page",)
for k in ("a", "b", "c"):
await store.aput(ns, k, {"k": k}, ttl=TTL_MINUTES)
await store.aput(ns, "expired", {"k": "x"}, ttl=TTL_MINUTES)
await _aexpire_now(store, ns, "expired")
seconds_ago = {"a": 1, "expired": 2, "b": 3, "c": 4}
# updated_at DESC orders these a, expired, b, c, so the expired row sits inside
# the first limit=2 window. Correct (pre-LIMIT) filtering yields live pages
# [a, b] then [c]; post-LIMIT filtering would underfill page 1 to just [a].
async with store._cursor() as cur:
for key, secs in seconds_ago.items():
await cur.execute(
"UPDATE store SET updated_at = NOW() - (%s * INTERVAL '1 second') "
"WHERE prefix = %s AND key = %s",
(secs, ".".join(ns), key),
)
page1 = await store.asearch(ns, limit=2, offset=0)
page2 = await store.asearch(ns, limit=2, offset=2)
assert [i.key for i in page1] == ["a", "b"]
assert [i.key for i in page2] == ["c"]
@@ -863,133 +863,6 @@ def test_store_ttl(store):
assert len(res) == 0
def _expire_now(store: PostgresStore, ns: tuple[str, ...], key: str) -> None:
"""Backdate a row's expires_at into the past without deleting it (unswept)."""
with store._cursor() as cur:
cur.execute(
"UPDATE store SET expires_at = NOW() - INTERVAL '1 minute' "
"WHERE prefix = %s AND key = %s",
(".".join(ns), key),
)
def _row_exists(store: PostgresStore, ns: tuple[str, ...], key: str) -> bool:
with store._cursor() as cur:
cur.execute(
"SELECT COUNT(*) AS n FROM store WHERE prefix = %s AND key = %s",
(".".join(ns), key),
)
return cur.fetchone()["n"] == 1
def _stored_expires_at(store: PostgresStore, ns: tuple[str, ...], key: str):
with store._cursor() as cur:
cur.execute(
"SELECT expires_at FROM store WHERE prefix = %s AND key = %s",
(".".join(ns), key),
)
return cur.fetchone()["expires_at"]
def test_omit_expired_filters_read_paths(store: PostgresStore) -> None:
store.stop_ttl_sweeper() # deterministic: no background deletion
store.ttl_config["omit_expired"] = True
expired_ns = ("omit", "expired")
control_ns = ("omit", "control")
store.put(expired_ns, "e", {"data": "gone"}, ttl=TTL_MINUTES)
store.put(control_ns, "c", {"data": "keep"}, ttl=None)
_expire_now(store, expired_ns, "e")
# The row is expired but physically still present (unswept).
assert _row_exists(store, expired_ns, "e")
# get omits it; the never-expiring control is still returned.
assert store.get(expired_ns, "e") is None
assert store.get(control_ns, "c") is not None
# search omits it but returns the control.
assert store.search(expired_ns) == []
assert [i.key for i in store.search(control_ns)] == ["c"]
# list_namespaces drops the expired-only namespace, keeps the control.
namespaces = store.list_namespaces(prefix=("omit",))
assert expired_ns not in namespaces
assert control_ns in namespaces
@pytest.mark.parametrize("omit", [None, False], ids=["default", "explicit-false"])
def test_omit_expired_disabled_preserves_expired_rows(
store: PostgresStore, omit
) -> None:
store.stop_ttl_sweeper()
if omit is not None:
store.ttl_config["omit_expired"] = omit
ns = ("keep",)
store.put(ns, "k", {"data": "still-here"}, ttl=TTL_MINUTES)
_expire_now(store, ns, "k")
assert store.get(ns, "k", refresh_ttl=False) is not None
assert [i.key for i in store.search(ns, refresh_ttl=False)] == ["k"]
assert ns in store.list_namespaces(prefix=("keep",))
def test_omit_expired_refresh_ttl_only_refreshes_live_rows(
store: PostgresStore,
) -> None:
store.stop_ttl_sweeper()
store.ttl_config["omit_expired"] = True
ns = ("refresh",)
store.put(ns, "expired", {"n": 0}, ttl=TTL_MINUTES)
store.put(ns, "live_get", {"n": 1}, ttl=TTL_MINUTES)
store.put(ns, "live_search", {"n": 2}, ttl=TTL_MINUTES)
_expire_now(store, ns, "expired")
expired_before = _stored_expires_at(store, ns, "expired")
get_before = _stored_expires_at(store, ns, "live_get")
search_before = _stored_expires_at(store, ns, "live_search")
# refresh_ttl=True must NOT resurrect the expired row (via get or search)...
assert store.get(ns, "expired", refresh_ttl=True) is None
assert "expired" not in [i.key for i in store.search(ns, refresh_ttl=True)]
assert _stored_expires_at(store, ns, "expired") == expired_before
# ...but must still extend the live rows that were read.
assert store.get(ns, "live_get", refresh_ttl=True) is not None
assert _stored_expires_at(store, ns, "live_get") > get_before
assert _stored_expires_at(store, ns, "live_search") > search_before
def test_omit_expired_search_pagination(store: PostgresStore) -> None:
store.stop_ttl_sweeper()
store.ttl_config["omit_expired"] = True
ns = ("page",)
for k in ("a", "b", "c"):
store.put(ns, k, {"k": k}, ttl=TTL_MINUTES)
store.put(ns, "expired", {"k": "x"}, ttl=TTL_MINUTES)
_expire_now(store, ns, "expired")
seconds_ago = {"a": 1, "expired": 2, "b": 3, "c": 4}
# updated_at DESC orders these a, expired, b, c, so the expired row sits inside
# the first limit=2 window. Correct (pre-LIMIT) filtering yields live pages
# [a, b] then [c]; post-LIMIT filtering would underfill page 1 to just [a].
with store._cursor() as cur:
for key, secs in seconds_ago.items():
cur.execute(
"UPDATE store SET updated_at = NOW() - (%s * INTERVAL '1 second') "
"WHERE prefix = %s AND key = %s",
(secs, ".".join(ns), key),
)
page1 = store.search(ns, limit=2, offset=0)
page2 = store.search(ns, limit=2, offset=2)
assert [i.key for i in page1] == ["a", "b"]
assert [i.key for i in page2] == ["c"]
@pytest.mark.parametrize(
"vector_type,distance_type",
[
@@ -552,14 +552,6 @@ class TTLConfig(TypedDict, total=False):
This can be overridden per-operation by explicitly setting `refresh_ttl`.
Defaults to `True` if not configured.
"""
omit_expired: bool
"""Whether to omit expired items from read operations.
If `True`, and if the store supports this option, expired items will not be
returned by `GET` or `SEARCH`, or included in namespace listings, even before
a TTL sweep deletes them.
Defaults to `False` if not configured.
"""
default_ttl: float | None
"""Default TTL (time-to-live) in minutes for new items.
+3 -3
View File
@@ -1276,9 +1276,9 @@ js-tiktoken@^1.0.12:
base64-js "^1.5.1"
js-yaml@^4.1.1:
version "4.3.0"
resolved "https://registry.yarnpkg.com/js-yaml/-/js-yaml-4.3.0.tgz#d1900572a7f7cf0b5f540c83673e60bad3436592"
integrity sha512-1td788aAnnZ5qs7V2QIRl1owjtYpbKt749Y3xauqQgwIIGF/xXWz1wMTEBx5O3LK3lXLVuqXPdPxj2BoFHaW9Q==
version "4.2.0"
resolved "https://registry.yarnpkg.com/js-yaml/-/js-yaml-4.2.0.tgz#2bd9e85682dd91bd469afb809d816043b3d49524"
integrity sha512-ePWsvanv0DWuDRsW8dnt+R4jQ31SCRCQ7hhNcPXZPsoBZiemuZNYGf7adZdqX2D86j6rvKp3RpCxVTSb8WQlOw==
dependencies:
argparse "^2.0.1"
+37 -1
View File
@@ -15,7 +15,11 @@ from langgraph_cli.analytics import log_command
from langgraph_cli.config import Config
from langgraph_cli.constants import DEFAULT_CONFIG, DEFAULT_PORT
from langgraph_cli.deploy import deploy
from langgraph_cli.docker import DockerCapabilities, build_docker_image
from langgraph_cli.docker import (
DockerCapabilities,
build_docker_image,
get_config_to_docker_build_context,
)
from langgraph_cli.exec import Runner, subp_exec
from langgraph_cli.progress import Progress
from langgraph_cli.templates import TEMPLATE_HELP_STRING, create_new
@@ -546,6 +550,14 @@ tests
)
@OPT_API_VERSION
@OPT_ENGINE_RUNTIME_MODE
@click.option(
"--install-command",
help="Custom install command to run from the build context root. If not provided, auto-detects based on package manager files.",
)
@click.option(
"--build-command",
help="Custom build command to run from the langgraph.json directory. If not provided, uses default build process.",
)
@log_command
def dockerfile(
save_path: str,
@@ -554,9 +566,24 @@ def dockerfile(
base_image: str | None = None,
api_version: str | None = None,
engine_runtime_mode: str = "combined_queue_worker",
install_command: str | None = None,
build_command: str | None = None,
) -> None:
from click import secho
if install_command and langgraph_cli.config.has_disallowed_build_command_content(
install_command
):
raise click.UsageError(
"install_command contains disallowed characters or patterns."
)
if build_command and langgraph_cli.config.has_disallowed_build_command_content(
build_command
):
raise click.UsageError(
"build_command contains disallowed characters or patterns."
)
save_path = pathlib.Path(save_path).absolute()
secho(f"🔍 Validating configuration at path: {config}", fg="yellow")
config_json = langgraph_cli.config.validate_config_file(config)
@@ -569,12 +596,21 @@ def dockerfile(
config_json, engine_runtime_mode=engine_runtime_mode
)
build_context = get_config_to_docker_build_context(
config_json,
install_command=install_command,
build_command=build_command,
)
secho(f"📝 Generating Dockerfile at {save_path}", fg="yellow")
dockerfile_content, additional_contexts = langgraph_cli.config.config_to_docker(
config_path=config,
config=config_json,
base_image=effective_base_image,
api_version=api_version,
install_command=install_command,
build_command=build_command,
build_context=build_context,
)
with open(str(save_path), "w", encoding="utf-8") as f:
f.write(dockerfile_content)
+20 -10
View File
@@ -330,6 +330,20 @@ def compose(
return compose_str
def get_config_to_docker_build_context(
config_json: dict,
install_command: str | None = None,
build_command: str | None = None,
) -> str | None:
"""Return the build context passed to config_to_docker, if overridden."""
is_js_project = config_json.get("node_version") and not config_json.get(
"python_version"
)
if is_js_project and (build_command or install_command):
return str(pathlib.Path.cwd())
return None
def build_docker_image(
runner,
set: Callable[[str], None],
@@ -365,16 +379,12 @@ def build_docker_image(
"-t",
tag,
]
# determine build context: use current directory for JS projects, config parent for Python
is_js_project = config_json.get("node_version") and not config_json.get(
"python_version"
config_to_docker_build_context = get_config_to_docker_build_context(
config_json,
install_command=install_command,
build_command=build_command,
)
# build/install commands only apply to JS projects for now
# without install/build command, JS projects will follow the old behavior
if is_js_project and (build_command or install_command):
build_context = str(pathlib.Path.cwd())
else:
build_context = str(config.parent)
build_context = config_to_docker_build_context or str(config.parent)
# Deep copy to avoid mutating the caller's config (config_to_docker
# rewrites graph paths to container-internal paths in place).
@@ -386,7 +396,7 @@ def build_docker_image(
api_version=api_version,
install_command=install_command,
build_command=build_command,
build_context=build_context,
build_context=config_to_docker_build_context,
)
# add additional_contexts
if additional_contexts:
+76
View File
@@ -1088,6 +1088,82 @@ def test_dockerfile_command_with_api_version_nodejs() -> None:
assert "FROM langchain/langgraphjs-api:0.2.74-node20" in dockerfile
def test_dockerfile_command_nodejs_monorepo_commands() -> None:
runner = CliRunner()
with runner.isolated_filesystem():
root = pathlib.Path.cwd()
config_dir = root / "apps" / "agent"
graph_path = config_dir / "src" / "agent.ts"
graph_path.parent.mkdir(parents=True)
graph_path.touch()
(root / "package.json").write_text(
json.dumps({"packageManager": "pnpm@10.0.0"})
)
(root / "pnpm-lock.yaml").touch()
(root / "pnpm-workspace.yaml").write_text('packages:\n - "apps/*"\n')
config_path = config_dir / "langgraph.json"
config_path.write_text(
json.dumps(
{
"node_version": "20",
"graphs": {"agent": "src/agent.ts:graph"},
"image_distro": "wolfi",
}
)
)
save_path = root / "Dockerfile"
result = runner.invoke(
cli,
[
"dockerfile",
str(save_path),
"--config",
str(config_path),
"--install-command",
"pnpm install --frozen-lockfile",
"--build-command",
"pnpm run build",
],
)
assert result.exit_code == 0, result.output
assert save_path.exists()
dockerfile = save_path.read_text()
container_root = f"/deps/{root.name}"
assert f"ADD . {container_root}" in dockerfile
assert f"WORKDIR {container_root}" in dockerfile
assert "RUN pnpm install --frozen-lockfile" in dockerfile
assert f"WORKDIR {container_root}/apps/agent" in dockerfile
assert "RUN pnpm run build" in dockerfile
def test_dockerfile_command_rejects_disallowed_nodejs_commands() -> None:
runner = CliRunner()
config_content = {
"node_version": "20",
"graphs": {"agent": "agent.js:graph"},
}
with temporary_config_folder(config_content) as temp_dir:
(temp_dir / "agent.js").touch()
for option in ("--install-command", "--build-command"):
result = runner.invoke(
cli,
[
"dockerfile",
str(temp_dir / "Dockerfile"),
"--config",
str(temp_dir / "config.json"),
option,
"npm install; echo bad",
],
)
assert result.exit_code != 0
assert "contains disallowed characters or patterns" in result.output
def test_build_command_with_api_version() -> None:
"""Test the 'build' command with --api-version flag."""
runner = CliRunner()
+3 -3
View File
@@ -2031,11 +2031,11 @@ wheels = [
[[package]]
name = "setuptools"
version = "83.0.0"
version = "82.0.1"
source = { registry = "https://pypi.org/simple" }
sdist = { url = "https://files.pythonhosted.org/packages/34/26/f5d29e25ffdb535afef2d35cdb55b325298f96debd670da4c325e08d70f4/setuptools-83.0.0.tar.gz", hash = "sha256:025bccbbf0fa05b6192bc64ae1e7b16e001fd6d6d4d5de03c97b1c1ade523bef", size = 1154254, upload-time = "2026-07-04T15:31:22.699Z" }
sdist = { url = "https://files.pythonhosted.org/packages/4f/db/cfac1baf10650ab4d1c111714410d2fbb77ac5a616db26775db562c8fab2/setuptools-82.0.1.tar.gz", hash = "sha256:7d872682c5d01cfde07da7bccc7b65469d3dca203318515ada1de5eda35efbf9", size = 1152316, upload-time = "2026-03-09T12:47:17.221Z" }
wheels = [
{ url = "https://files.pythonhosted.org/packages/5d/40/e1e72872c6354b306daef1703549e8e83b4d43cfea356311bf722a043752/setuptools-83.0.0-py3-none-any.whl", hash = "sha256:29b23c360f22f414dc7336bb39178cc7bcbf6021ed2733cde173f09dba19abb3", size = 1008090, upload-time = "2026-07-04T15:31:20.885Z" },
{ url = "https://files.pythonhosted.org/packages/9d/76/f789f7a86709c6b087c5a2f52f911838cad707cc613162401badc665acfe/setuptools-82.0.1-py3-none-any.whl", hash = "sha256:a59e362652f08dcd477c78bb6e7bd9d80a7995bc73ce773050228a348ce2e5bb", size = 1006223, upload-time = "2026-03-09T12:47:15.026Z" },
]
[[package]]
+4 -4
View File
@@ -3512,7 +3512,7 @@ class Pregel(
control: RunControl | None = None,
transformers: Sequence[Callable[[tuple[str, ...]], Any]] | None = None,
**kwargs: Any,
) -> GraphRunStream:
) -> Any:
"""Internal v3 sync streaming implementation. Public entry: stream_events(version='v3').
Extra keyword arguments are forwarded to the underlying ``stream(...)``
@@ -3568,7 +3568,7 @@ class Pregel(
control: RunControl | None = None,
transformers: Sequence[Callable[[tuple[str, ...]], Any]] | None = None,
**kwargs: Any,
) -> AsyncGraphRunStream:
) -> Any:
"""Internal v3 async streaming implementation. Public entry: astream_events(version='v3').
Extra keyword arguments are forwarded to the underlying ``astream(...)``
@@ -3633,7 +3633,7 @@ class Pregel(
control: RunControl | None = None,
transformers: Sequence[Callable[[tuple[str, ...]], Any]] | None = None,
**kwargs: Any,
) -> GraphRunStream: ...
) -> Any: ...
def stream_events(
self,
@@ -3738,7 +3738,7 @@ class Pregel(
control: RunControl | None = None,
transformers: Sequence[Callable[[tuple[str, ...]], Any]] | None = None,
**kwargs: Any,
) -> Awaitable[AsyncGraphRunStream]: ...
) -> Awaitable[Any]: ...
def astream_events(
self,
+1 -27
View File
@@ -12,13 +12,7 @@ from langgraph.stream._mux import StreamMux
from langgraph.stream._types import ProtocolEvent
if TYPE_CHECKING:
from langchain_core.language_models.chat_model_stream import (
AsyncChatModelStream,
ChatModelStream,
)
from langgraph.stream.stream_channel import StreamChannel
from langgraph.stream.transformers import LifecyclePayload, SubgraphStatus
from langgraph.stream.transformers import SubgraphStatus
def _drive_until_done(pump: Callable[[], bool]) -> None:
@@ -54,16 +48,6 @@ class GraphRunStream:
experimental and may change.
"""
# Native projections always registered by `stream_events(version="v3")`.
# Attached dynamically by the `setattr` loop in `__init__`; declared here
# so type checkers see them. Opt-in native projections (`updates`,
# `custom`, `checkpoints`, `debug`, `tasks`) are only present when their
# transformer is registered, so they are reached via `extensions[...]`.
values: StreamChannel[dict[str, Any]]
messages: StreamChannel[ChatModelStream]
lifecycle: StreamChannel[LifecyclePayload]
subgraphs: StreamChannel[SubgraphRunStream]
def __init__(
self,
graph_iter: Iterator[Any] | None,
@@ -345,16 +329,6 @@ class AsyncGraphRunStream:
experimental and may change.
"""
# Native projections always registered by `astream_events(version="v3")`.
# Attached dynamically by the `setattr` loop in `__init__`; declared here
# so type checkers see them. Opt-in native projections (`updates`,
# `custom`, `checkpoints`, `debug`, `tasks`) are only present when their
# transformer is registered, so they are reached via `extensions[...]`.
values: StreamChannel[dict[str, Any]]
messages: StreamChannel[AsyncChatModelStream]
lifecycle: StreamChannel[LifecyclePayload]
subgraphs: StreamChannel[AsyncSubgraphRunStream]
def __init__(
self,
graph_aiter: AsyncIterator[Any] | None,
@@ -12,10 +12,6 @@ from dataclasses import dataclass
from typing import Annotated, Any, TypeVar
import pytest
from langchain_core.language_models.chat_model_stream import (
AsyncChatModelStream,
ChatModelStream,
)
from langchain_core.messages import AIMessage, BaseMessage
from langgraph.checkpoint.memory import InMemorySaver
from pydantic import BaseModel, ValidationError
@@ -28,14 +24,6 @@ from langgraph.func import entrypoint
from langgraph.graph import StateGraph
from langgraph.graph.message import MessagesState
from langgraph.runtime import RunControl
from langgraph.stream import (
AsyncGraphRunStream,
AsyncSubgraphRunStream,
GraphRunStream,
LifecyclePayload,
StreamChannel,
SubgraphRunStream,
)
from langgraph.types import (
CheckpointPayload,
CheckpointStreamPart,
@@ -1190,32 +1178,3 @@ def _check_type_narrowing(part: StreamPart[_StateT, _OutputT]) -> None:
assert_type(part, DebugStreamPart[_StateT])
assert_type(part["data"], DebugPayload[_StateT])
assert_type(part["ns"], tuple[str, ...])
# --- v3 stream_events return / projection typing checks ---
# These functions are never called at runtime; `ty` validates the
# assert_type calls. They pin the public typing surface of
# stream_events(version="v3") / astream_events(version="v3"): the handle
# type and the always-registered native projections.
def _check_stream_events_v3_typing() -> None:
"""Compile-time checks for sync v3 typing — never called at runtime."""
graph = _make_simple_graph().compile()
run = graph.stream_events(_SIMPLE_INPUT, version="v3")
assert_type(run, GraphRunStream)
assert_type(run.values, StreamChannel[dict[str, Any]])
assert_type(run.messages, StreamChannel[ChatModelStream])
assert_type(run.lifecycle, StreamChannel[LifecyclePayload])
assert_type(run.subgraphs, StreamChannel[SubgraphRunStream])
async def _check_astream_events_v3_typing() -> None:
"""Compile-time checks for async v3 typing — never called at runtime."""
graph = _make_simple_graph().compile()
run = await graph.astream_events(_SIMPLE_INPUT, version="v3")
assert_type(run, AsyncGraphRunStream)
assert_type(run.values, StreamChannel[dict[str, Any]])
assert_type(run.messages, StreamChannel[AsyncChatModelStream])
assert_type(run.lifecycle, StreamChannel[LifecyclePayload])
assert_type(run.subgraphs, StreamChannel[AsyncSubgraphRunStream])
+6 -7
View File
@@ -1345,7 +1345,7 @@ wheels = [
[[package]]
name = "jupyterlab"
version = "4.5.10"
version = "4.5.9"
source = { registry = "https://pypi.org/simple" }
dependencies = [
{ name = "async-lru" },
@@ -1362,11 +1362,10 @@ dependencies = [
{ name = "tomli", marker = "python_full_version < '3.11'" },
{ name = "tornado" },
{ name = "traitlets" },
{ name = "typing-extensions", marker = "python_full_version < '3.12'" },
]
sdist = { url = "https://files.pythonhosted.org/packages/74/24/621aa20ec0d2fe72f52095bda0fc1be7738ac21aabe4129ff623140d5cdf/jupyterlab-4.5.10.tar.gz", hash = "sha256:77e8d80b78be59b2eaba2154562e21caa6e79c2f1281d6f486584f7144ee2f47", size = 23998879, upload-time = "2026-07-21T12:43:27.324Z" }
sdist = { url = "https://files.pythonhosted.org/packages/e8/52/a8d4895bef501ffeb6af448e8bf7079541c7772978211963aa653518c2d9/jupyterlab-4.5.9.tar.gz", hash = "sha256:dd79a073fecae7a39066ea99e4627ed6c76269ac926e95a810e1e1df6358d865", size = 23994445, upload-time = "2026-06-17T15:42:16.406Z" }
wheels = [
{ url = "https://files.pythonhosted.org/packages/9f/c9/940f95f17ee4e413ad252bf8d4f2ee9a341f18cfeda87775fef3d7847321/jupyterlab-4.5.10-py3-none-any.whl", hash = "sha256:5967ca61e692e67a2f30b5a2b901c941dc6ce56c0b0e357bc6d34fed5ec095f6", size = 12452502, upload-time = "2026-07-21T12:43:23.542Z" },
{ url = "https://files.pythonhosted.org/packages/c6/bb/2f9b425062416fba58f580c9b89c3b07277ccdf0a292501fedbca8ea00ea/jupyterlab-4.5.9-py3-none-any.whl", hash = "sha256:5ff0f908e8ac0afbed32b106fdef360f101c0a6654d1bf4a81e98a293ae1b336", size = 12449803, upload-time = "2026-06-17T15:42:12.18Z" },
]
[[package]]
@@ -3384,11 +3383,11 @@ wheels = [
[[package]]
name = "setuptools"
version = "83.0.0"
version = "80.9.0"
source = { registry = "https://pypi.org/simple" }
sdist = { url = "https://files.pythonhosted.org/packages/34/26/f5d29e25ffdb535afef2d35cdb55b325298f96debd670da4c325e08d70f4/setuptools-83.0.0.tar.gz", hash = "sha256:025bccbbf0fa05b6192bc64ae1e7b16e001fd6d6d4d5de03c97b1c1ade523bef", size = 1154254, upload-time = "2026-07-04T15:31:22.699Z" }
sdist = { url = "https://files.pythonhosted.org/packages/18/5d/3bf57dcd21979b887f014ea83c24ae194cfcd12b9e0fda66b957c69d1fca/setuptools-80.9.0.tar.gz", hash = "sha256:f36b47402ecde768dbfafc46e8e4207b4360c654f1f3bb84475f0a28628fb19c", size = 1319958, upload-time = "2025-05-27T00:56:51.443Z" }
wheels = [
{ url = "https://files.pythonhosted.org/packages/5d/40/e1e72872c6354b306daef1703549e8e83b4d43cfea356311bf722a043752/setuptools-83.0.0-py3-none-any.whl", hash = "sha256:29b23c360f22f414dc7336bb39178cc7bcbf6021ed2733cde173f09dba19abb3", size = 1008090, upload-time = "2026-07-04T15:31:20.885Z" },
{ url = "https://files.pythonhosted.org/packages/a3/dc/17031897dae0efacfea57dfd3a82fdd2a2aeb58e0ff71b77b87e44edc772/setuptools-80.9.0-py3-none-any.whl", hash = "sha256:062d34222ad13e0cc312a4c02d73f059e86a4acbfbdea8f8f76b28c99f306922", size = 1201486, upload-time = "2025-05-27T00:56:49.664Z" },
]
[[package]]
-3
View File
@@ -378,8 +378,6 @@ class Run(TypedDict):
"""The run metadata."""
multitask_strategy: MultitaskStrategy
"""Strategy to handle concurrent runs on the same thread."""
langsmith_session_name: NotRequired[str | None]
"""The LangSmith session name associated with this run."""
class Cron(TypedDict):
@@ -490,7 +488,6 @@ RunSelectField = Literal[
"metadata",
"kwargs",
"multitask_strategy",
"langsmith_session_name",
]
CronSelectField = Literal[