mirror of
https://github.com/langchain-ai/langgraph.git
synced 2026-09-05 17:27:47 +02:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
f3de8ed373 |
@@ -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.
|
||||
|
||||
|
||||
@@ -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"
|
||||
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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()
|
||||
|
||||
Generated
+3
-3
@@ -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]]
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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])
|
||||
|
||||
Generated
+6
-7
@@ -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]]
|
||||
|
||||
@@ -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[
|
||||
|
||||
Reference in New Issue
Block a user