Compare commits

...
Author SHA1 Message Date
Elior Nataf Lackritz a8e732c879 fix(langgraph): give each bulk_update_state update its own task id
An update whose node has no pending task to reuse was stored under
uuid5(checkpoint_id, INTERRUPT), so every such update in one superstep
shared a task id. Savers keep one write per (task_id, idx), so all but the
first update's writes were dropped. Plain channels were unaffected, since
their value is stored in the new checkpoint, but a DeltaChannel replays
those writes and lost every update after the first.

The ith update now gets uuid5(checkpoint_id, f"{INTERRUPT}:{i}"). The first
keeps the old id, so a single update stores exactly what it did before.
2026-09-30 12:47:55 -04:00
eb69f67b65 fix(checkpoint-sqlite): walk delta ancestors by parent pointer (#8557)
## Summary

The sqlite delta history silently drops a parent checkpoint whose id sorts above its child's,
losing that parent's stored value and its pending writes. The channel hydrates short with no error.

Fixes #8550

## Problem

Stage 1 walked ancestors with:

```sql
WHERE thread_id = ? AND checkpoint_ns = ? AND checkpoint_id <= ?
ORDER BY checkpoint_id DESC
```

Ancestry is defined by the `parent_checkpoint_id` column. These two predicates add a second
requirement: that every child's id sorts above its parent's. The contract promises monotonic ids,
but that only holds within one process, so ids from processes with different clocks can break it.

When the requirement is violated the parent is excluded from the stream and its seed and writes go
with it. Dropping the range filter alone does not fix it: in `checkpoint_id DESC` order that parent
arrives *before* the target, so the walk streams past it before it has started.

## Fix

A recursive CTE anchored at the target, following `parent_checkpoint_id`:

```sql
WITH RECURSIVE ancestors(checkpoint_id, parent_checkpoint_id, type, checkpoint) AS (
    SELECT ... FROM checkpoints
    WHERE thread_id = ? AND checkpoint_ns = ? AND checkpoint_id = ?
    UNION ALL
    SELECT c.... FROM ancestors a CROSS JOIN checkpoints c
      ON c.checkpoint_id = a.parent_checkpoint_id
    WHERE c.thread_id = ? AND c.checkpoint_ns = ?
)
SELECT checkpoint_id, type, checkpoint FROM ancestors
```

Rows now arrive in walk order (target, parent, grandparent, ...), so `step_walk_with_row` no longer
needs its off-path skip or its `parent_cid` tracking; both are removed. The query reads only true
ancestors, where the old one read every row at or below the target including sibling branches.

`CROSS JOIN` pins the join order. The saver never runs `ANALYZE`, and with a plain `JOIN` sqlite put
`checkpoints` as the outer loop, scanning the whole thread on every recursion step. With `ancestors`
outside, each step is one primary key lookup. Through `get_delta_channel_history`:

| chain length | plain `JOIN` | `CROSS JOIN` |
| -- | -- | -- |
| 1000 | 0.032s | 0.001s |
| 2000 | 0.124s | 0.003s |
| 4000 | 0.475s | 0.006s |

## Cycle guard

Following pointers can loop where a bounded id scan could not, and a loop is reachable through
`put` alone: `put` writes with `INSERT OR REPLACE`, so re-putting an existing checkpoint id under a
descendant's config repoints that checkpoint at its own descendant. The walk stops on a repeated
`checkpoint_id` (one set insert per row, no depth ceiling that could truncate a long migrated
thread). sqlite yields recursive rows lazily, so abandoning the cursor ends the recursion.

`test_walk_terminates_when_put_makes_the_parent_chain_cycle` fails by hanging, not by asserting, if
the guard regresses (confirmed by deleting the guard locally). The package has no `pytest-timeout`,
so the CI job timeout is the backstop.

## Postgres

No equivalent change needed. It pages the whole thread with no id bound and follows parent pointers
in Python, and its upsert never rewrites `parent_checkpoint_id`, so it can neither miss this parent
nor form the loop. `BaseCheckpointSaver` and `InMemorySaver` also walk parent pointers.

## Test plan

New `libs/checkpoint-sqlite/tests/test_delta_parent_walk.py`:

- [x] Sync and async, parametrised over both id orders; the sync case also asserts equality with
      `BaseCheckpointSaver` on the same rows. `parent_id_sorts_above_child` is the bug,
      `parent_id_sorts_below_child` the control.
- [x] `test_walk_reaches_root_of_long_chain_with_descending_ids`: 40 checkpoints, only stored value
      at the root.
- [x] `test_walk_terminates_when_put_makes_the_parent_chain_cycle`.
- [x] `test_walk_step_looks_up_the_parent_by_primary_key`: asserts the recursive step's
      `EXPLAIN QUERY PLAN` is a key lookup, so a plain `JOIN` can't come back. Fails with it.
- [x] On `main`: 3 of the 6 walk tests fail (both `parent_id_sorts_above_child` cases and the long
      chain). The cycle test passes on `main` too, since the old bounded scan could not loop; it
      guards the new path.
- [x] #8550's repro returns `{'writes': [('task', 'ch', 'write-root')], 'seed': 'seed'}` sync and
      async (was `{'writes': []}` on `main`).
- [x] `libs/checkpoint-sqlite`: `make format`, `make lint` clean; full suite 125 passed, 2 skipped.
- [x] `libs/langgraph`: `-k "delta or sqlite"` 739 passed, 1 skipped.

Thanks to @lylelllll for the report, the minimal repro, the base-saver comparison that isolated it
to the fast path, and for suggesting the recursive CTE.




Co-authored-by: lylelllll <59271327+lylelllll@users.noreply.github.com>
2026-09-30 12:17:02 -04:00
c0279f0910 fix(checkpoint-postgres): derive the delta walk cursor once the target loads (#8556)
## Summary

`get_delta_channel_history` on Postgres returns an empty history for any `DeltaChannel` on a
target checkpoint that is not within the first stage-1 pagination page (1024 rows) of the thread.
No exception, no warning: the channel just hydrates empty.

Fixes #8448

## Problem

Stage 1 pages `checkpoints` newest-first from the head of the thread, and after each page
`_try_advance_walks` tries to move every not-yet-seeded channel's walk along the partial
`parent_of` map accumulated so far. The walk starts at the target's parent:

```python
if ch not in walk_cursor_by_ch:
    walk_cursor_by_ch[ch] = parent_of.get(target_id)
```

The target can be any checkpoint in the thread, not just the head, so on the first page
`parent_of` frequently has no row for it yet. `.get` then returns `None`, which is also what a
target with no parent returns, and the two are stored identically. Because the initialisation is
guarded by `ch not in walk_cursor_by_ch`, it never runs again: once the walk is parked at `None`
it stays there even after the target's real row and real parent load on a later page.

The result is an empty chain and no seed. Downstream `channels_from_checkpoint` does

```python
replay_ch = delta_spec.from_checkpoint(history.get("seed", MISSING))
replay_ch.replay_writes(history["writes"])
```

so `get_state`, `get_state_history` and `update_state` against an older checkpoint reconstruct a
`messages` channel as `[]` on a thread with hundreds of real messages.

## Fix

Start the walk only once `target_id` is actually present in `parent_of`, so "the target has not
loaded yet" stops sharing a representation with "the target is a root":

```python
if ch not in walk_cursor_by_ch:
    if target_id not in parent_of:
        continue
    walk_cursor_by_ch[ch] = parent_of[target_id]
```

`_try_advance_walks` is a static method on `BasePostgresSaver`, so `PostgresSaver` and
`AsyncPostgresSaver` are both covered by the one change.

## Why it's safe

`continue` leaves the channel exactly as it was, so a later page retries. The three existing
stop conditions are untouched: a channel that finds its seed still seeds, one that reaches a real
root still parks at `None`, and one waiting on an ancestor still keeps its cursor. Paging still
terminates on a short page, which is what ends the run for a target that really is a root.

## Long-term

The sibling sqlite implementation avoids this class of bug differently, by starting its stage-1
scan at the target (`checkpoint_id <= ?`) instead of at the head. Postgres could adopt the same
bound and would then never fetch a checkpoint newer than the target at all, which looks like the
bigger win on a long thread. It makes the read path depend on ancestors always sorting below their
descendants, though, which sqlite already assumes but the Postgres fast path currently does not.
#8550 now reports that assumption as a bug in sqlite, on the grounds that ancestry is defined by
`parent_checkpoint_id` and the contract does not require ids to be monotonic, so the bound is the
wrong direction to move Postgres in. Paging the full thread and following parent pointers is what
keeps this path correct when ids are not monotonic, and with this fix Postgres returns the right
history for #8550's scenario at every page size.

## Test plan

New `libs/checkpoint-postgres/tests/test_delta_pagination.py`. Page size is monkeypatched rather
than writing 1024+ real checkpoints per case, since the only thing that decides the behaviour is
which page the target lands on.

- [x] `test_async_target_older_than_the_first_page` and its sync twin, parametrised over page
      sizes `[_DELTA_PAGE_SIZE, 3, 2, 1]`. The thread has 8 checkpoints with a snapshot at step 1
      and the target at step 4, so every size at or below 3 leaves the target off the first page.
      The real page size is the control.
- [x] `test_root_target_has_no_history_and_still_terminates` covers the case where a `None` cursor
      is the correct answer, at page size 1 so the paging loop runs the length of the thread.
- [x] 6 of the 9 fail on `main` (`expected a snapshot seed, got '<missing>'`); the 3 that pass are
      the two controls and the root case.
- [x] `make format`, `make lint_package`, `make lint_tests` clean.
- [x] Full `libs/checkpoint-postgres` suite, rebased on current `main`: 279 passed, 3 skipped on Postgres 16.
- [x] Graph-level repro with `_DELTA_PAGE_SIZE = 5`: 10 invocations, then `get_state` on the 8th-newest
      checkpoint returns `[]` on `main` and the full history on this branch.

Thanks to @Navneet-Scaler for the report, the mechanism write-up, and the fix in #8453, which this
matches.




Co-authored-by: Navneet-Scaler <147032454+Navneet-Scaler@users.noreply.github.com>
2026-09-30 12:16:55 -04:00
ccurmeandGitHub f5804a5bf5 fix(ci): test locally-built wheel and publish to test pypi after pre-release checks (#9124) 2026-09-30 09:30:34 -04:00
John KennedyGitHubopen-swe[bot] <open-swe@users.noreply.github.com>
07b33185ea fix: reject credential-bearing Git dependencies (#8542)
## Description
Reject Git HTTP dependency URLs containing userinfo before Docker
generation so credentials cannot persist in Dockerfiles or image layers.
Validation now covers local requirement/package metadata and uv
pyproject/lock inputs while keeping errors token-free.

## Test Plan
- [x] Validate credentialed raw, local-manifest, and uv-managed Git URLs
are rejected without echoing secrets
- [x] Validate credential-free HTTPS and SSH Git URLs remain supported

Made by [Open
SWE](https://openswe.vercel.app/agents/81b07455-ece4-3ddc-9955-d7a5bea78d2c)

---------

Co-authored-by: open-swe[bot] <open-swe@users.noreply.github.com>
2026-09-27 21:34:53 +00:00
Hugo DURANDandGitHub 7daa3ab49d feat(cli): place self-hosted deployments on a listener (#9056)
Follow-up to #8482. `langgraph deploy --push-to` can now create a
deployment in a workspace that
deploys through a listener in the customer's own cluster, which is the
hybrid case. Before this,
creation in such a workspace was impossible from the CLI: the control
plane rejected it and the CLI
told the user to go and create the deployment in the UI first.

## Changes
- Smart Auto-Placement: The CLI now proactively checks your workspace.
If you only have one listener and one Kubernetes namespace configured
(and are using the managed cloud control plane), it automatically routes
your deployment there. No extra flags needed.
- New Disambiguation Flags: If your workspace has multiple listeners or
namespaces, the CLI will ask you to choose. You can now pass
--listener-id and --k8s-namespace to tell it exactly where to deploy.
- Failing Fast: The CLI now validates your listener and namespace
choices before it starts building and pushing the heavy Docker image. If
you provide an invalid ID, it stops immediately instead of wasting your
time and bandwidth.
- Fixed a Duplication Bug: Previously, if you had many deployments with
similar names, a pagination issue could hide your existing deployment
from the CLI, causing it to accidentally create a duplicate. The CLI now
queries the server for the exact deployment name to guarantee this
doesn't happen.
- Cleaner Errors: Error messages from the control plane are now stripped
of their clunky HTTP envelopes so you get clear, readable sentences when
something goes wrong.

## Testing

Deployment on 3 paths, hybrid, self-hosted, nominal
2026-09-23 13:56:01 -04:00
Sreekara YachamaneniandGitHub e868c3ccfd feat(cli): clarify agent flags and support env defaults (#9063)
Agent deployment options now print a private-beta notice. Rename
`--environment` to `--agent-environment` and accept `LANGSMITH_AGENT_ID`
/ `LANGSMITH_AGENT_ENVIRONMENT` as process-environment defaults for
deploy and list. Explicit flags take precedence, and the backend payload
is unchanged.

Validation: formatting and lint pass. A local smoke check verified
environment-only deployment, explicit flag precedence, list defaults,
and structured JSON output. Full CLI suite: 411 passed; the two known
Docker failures remain (`test_dockerfile_command_with_docker_compose`
and `test_build_generate_proper_build_context`). No new tests added; the
existing test invocation uses the renamed flag.
2026-09-23 17:53:46 +00:00
25 changed files with 2171 additions and 306 deletions
+23 -36
View File
@@ -139,23 +139,10 @@ jobs:
echo EOF echo EOF
} >> "$GITHUB_OUTPUT" } >> "$GITHUB_OUTPUT"
test-pypi-publish:
needs:
- build
- release-notes
permissions:
contents: read
id-token: write
uses: ./.github/workflows/_test_release.yml
with:
working-directory: ${{ inputs.working-directory }}
secrets: inherit
pre-release-checks: pre-release-checks:
needs: needs:
- build - build
- release-notes - release-notes
- test-pypi-publish
runs-on: ubuntu-latest runs-on: ubuntu-latest
steps: steps:
- uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1 - uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
@@ -180,31 +167,20 @@ jobs:
enable-cache: false enable-cache: false
working-directory: ${{ inputs.working-directory }} working-directory: ${{ inputs.working-directory }}
- name: Import published package - uses: actions/download-artifact@3e5f45b2cfb9172054b4087a40e8e0b5a5461e7c # v8.0.1
with:
name: dist
path: ${{ inputs.working-directory }}/dist/
- name: Import dist package
shell: bash shell: bash
working-directory: ${{ inputs.working-directory }} working-directory: ${{ inputs.working-directory }}
env: env:
PKG_NAME: ${{ needs.build.outputs.pkg-name }} PKG_NAME: ${{ needs.build.outputs.pkg-name }}
VERSION: ${{ needs.build.outputs.version }} VERSION: ${{ needs.build.outputs.version }}
# Here we use: # Install directly from the locally-built wheel (no index resolution needed).
# - The default regular PyPI index as the *primary* index, meaning
# that it takes priority (https://pypi.org/simple)
# - The test PyPI index as an extra index, so that any dependencies that
# are not found on test PyPI can be resolved and installed anyway.
# (https://test.pypi.org/simple). This will include the PKG_NAME==VERSION
# package because VERSION will not have been uploaded to regular PyPI yet.
# - attempt install again after 5 seconds if it fails because there is
# sometimes a delay in availability on test pypi
run: | run: |
uv run pip install \ uv run pip install dist/*.whl
--extra-index-url https://test.pypi.org/simple/ \
"$PKG_NAME==$VERSION" || \
( \
sleep 5 && \
uv run pip install \
--extra-index-url https://test.pypi.org/simple/ \
"$PKG_NAME==$VERSION" \
)
if [[ "$PKG_NAME" == *prebuilt* ]]; then if [[ "$PKG_NAME" == *prebuilt* ]]; then
uv run pip install langgraph uv run pip install langgraph
@@ -226,7 +202,7 @@ jobs:
run: uv sync --group test run: uv sync --group test
working-directory: ${{ inputs.working-directory }} working-directory: ${{ inputs.working-directory }}
# Overwrite the local version of the package with the test PyPI version. # Overwrite the local version of the package with the built version
- name: Import published package (again) - name: Import published package (again)
working-directory: ${{ inputs.working-directory }} working-directory: ${{ inputs.working-directory }}
shell: bash shell: bash
@@ -234,14 +210,25 @@ jobs:
PKG_NAME: ${{ needs.build.outputs.pkg-name }} PKG_NAME: ${{ needs.build.outputs.pkg-name }}
VERSION: ${{ needs.build.outputs.version }} VERSION: ${{ needs.build.outputs.version }}
run: | run: |
uv run pip install \ uv run pip install dist/*.whl
--extra-index-url https://test.pypi.org/simple/ \
"$PKG_NAME==$VERSION"
- name: Run unit tests - name: Run unit tests
run: make test run: make test
working-directory: ${{ inputs.working-directory }} working-directory: ${{ inputs.working-directory }}
test-pypi-publish:
needs:
- build
- release-notes
- pre-release-checks
permissions:
contents: read
id-token: write
uses: ./.github/workflows/_test_release.yml
with:
working-directory: ${{ inputs.working-directory }}
secrets: inherit
publish: publish:
needs: needs:
- build - build
@@ -448,11 +448,12 @@ class PostgresSaver(BasePostgresSaver):
Two-stage query, both stages cover ALL requested channels: Two-stage query, both stages cover ALL requested channels:
* Stage 1 (paged): dynamic SELECT over `checkpoints` with K parallel * Stage 1 (paged): dynamic SELECT over `checkpoints` with three
JSONB key lookups (one column pair per channel) — no subquery, no columns per channel: its version, an `EXISTS` probe for a stored
aggregation. Pages newest-first by `checkpoint_id` with a cursor; blob at that version, and its inline value. Pages newest-first by
page size is `_DELTA_PAGE_SIZE`. Stops paging when every channel `checkpoint_id` with a cursor; page size is `_DELTA_PAGE_SIZE`.
has found its seed or the chain is exhausted. Stops paging when every channel has found its seed or a page comes
back short.
* Stage 2 (per-channel UNION ALL): one branch per channel reading * Stage 2 (per-channel UNION ALL): one branch per channel reading
`checkpoint_writes` filtered to that channel's specific `checkpoint_writes` filtered to that channel's specific
@@ -172,30 +172,8 @@ class _DeltaStage2Row(TypedDict, total=False):
version: str | None # "b" rows only version: str | None # "b" rows only
# Multi-channel two-stage DeltaChannel reconstruction. # Delta history is rebuilt in two queries; `_build_delta_stage1_sql` and
# # `_build_delta_stage2_sql` document their shapes.
# Stage 1 scans checkpoint metadata (no blob bytes) and emits one row per
# checkpoint with K parallel JSONB key lookups (one column pair per
# requested delta channel: ver_i / hs_i). No subqueries, no aggregation.
# Python walks the parent chain once across all channels.
#
# Stage 2 fetches all writes and the seed blobs for ALL channels in a
# single roundtrip via `channel = ANY(%s)` and chain/seed-version
# filtering.
#
# Empirical comparison vs an alternative "ship full channel_versions /
# channel_values JSONB and let Python pick" form (1000 checkpoints,
# 8 total channels in graph, 3 delta channels requested):
#
# Postgres execution: A=0.24ms vs B=0.38ms (both negligible)
# End-to-end latency: A=6.83ms vs B=2.28ms (B is 3.0x faster)
# Wire payload: A=836KB vs B=330KB (61% smaller)
# Buffer hits: identical (167 blocks)
#
# B (this dynamic-columns design) wins because it avoids JSONB
# serialization on the wire and JSONB-to-dict deserialization in
# psycopg. Even at K=8 (8 delta channels = 16 dynamic columns), B
# still beats A end-to-end (4.2ms vs 6.8ms).
def _build_delta_stage1_sql(channels: Sequence[str], *, paged: bool) -> str: def _build_delta_stage1_sql(channels: Sequence[str], *, paged: bool) -> str:
@@ -335,10 +313,8 @@ def _build_delta_stage2_sql(
return " UNION ALL ".join(branches) return " UNION ALL ".join(branches)
# Stage 1 rows are dynamic-shape dicts: {checkpoint_id, parent_checkpoint_id, # Stage 1 rows are dicts keyed by the per-channel aliases
# ver_0, hs_0, ver_1, hs_1, ...}. Walking is parameterized by the channel # `_build_delta_stage1_sql` emits, so there is no static TypedDict.
# list to map indices back to channel names — no static TypedDict here.
# `dict[str, Any]` is the practical signature.
class BasePostgresSaver(BaseCheckpointSaver[str]): class BasePostgresSaver(BaseCheckpointSaver[str]):
@@ -431,9 +407,11 @@ class BasePostgresSaver(BaseCheckpointSaver[str]):
(a) it found a stored value for its channel — a blob or an inline (a) it found a stored value for its channel — a blob or an inline
primitive (channel becomes seeded), primitive (channel becomes seeded),
(b) it reached a real root (parent_of[cid] is None — fully (b) it reached a real root (parent_of[cid] is None — fully
materialized at this point), or materialized at this point),
(c) the next ancestor cid isn't in `parent_of` yet (waiting for (c) the next ancestor cid isn't in `parent_of` yet (waiting for
a later page; the cursor stays put). a later page; the cursor stays put), or
(d) the target's own row isn't in `parent_of` yet (the walk has
not started; no cursor is set, so a later page retries).
Mutates `chain_by_ch`, `seed_ver_by_ch`, `seed_inline_by_ch`, Mutates `chain_by_ch`, `seed_ver_by_ch`, `seed_inline_by_ch`,
`walk_cursor_by_ch`, and `seeded` in place. `walk_cursor_by_ch`, and `seeded` in place.
@@ -441,9 +419,12 @@ class BasePostgresSaver(BaseCheckpointSaver[str]):
for i, ch in enumerate(channels): for i, ch in enumerate(channels):
if ch in seeded: if ch in seeded:
continue continue
# First-time entry: cursor starts at the target's parent. # Pages start at the thread head, so the target may not have
# loaded yet; a `None` cursor would read as "target is a root".
if ch not in walk_cursor_by_ch: if ch not in walk_cursor_by_ch:
walk_cursor_by_ch[ch] = parent_of.get(target_id) if target_id not in parent_of:
continue
walk_cursor_by_ch[ch] = parent_of[target_id]
cur_cid = walk_cursor_by_ch[ch] cur_cid = walk_cursor_by_ch[ch]
ch_chain = chain_by_ch[ch] ch_chain = chain_by_ch[ch]
hb_i = hb_by_i_by_cid[i] hb_i = hb_by_i_by_cid[i]
@@ -0,0 +1,132 @@
from __future__ import annotations
from typing import Any
from uuid import uuid4
import pytest
from langgraph.checkpoint.base import (
Checkpoint,
DeltaChannelHistory,
empty_checkpoint,
)
from langgraph.checkpoint.base.id import uuid6
from langgraph.checkpoint.serde.types import _DeltaSnapshot
from langgraph.checkpoint.postgres import PostgresSaver
from langgraph.checkpoint.postgres.aio import AsyncPostgresSaver
from langgraph.checkpoint.postgres.base import _DELTA_PAGE_SIZE
from tests.conftest import DEFAULT_URI
CHANNEL = "items"
STEPS = 8
SEED_STEP = 1
SEED_VALUE = [10, 20]
TARGET_STEP = 4
# The real page size is the control; the rest leave the target off the first
# page (three checkpoints are newer than it).
PAGE_SIZES = [_DELTA_PAGE_SIZE, 3, 2, 1]
def _step_args(
thread_id: str, step: int, parent: dict | None
) -> tuple[dict, Checkpoint, dict[str, Any]]:
config: dict = {"configurable": {"thread_id": thread_id, "checkpoint_ns": ""}}
if parent is not None:
config["configurable"]["checkpoint_id"] = parent["configurable"][
"checkpoint_id"
]
checkpoint: Checkpoint = empty_checkpoint()
checkpoint["id"] = str(uuid6(clock_seq=step))
checkpoint["channel_versions"][CHANNEL] = f"v{step}"
if step == SEED_STEP:
checkpoint["channel_values"][CHANNEL] = _DeltaSnapshot(list(SEED_VALUE))
return config, checkpoint, {CHANNEL: f"v{step}"}
return config, checkpoint, {}
async def _abuild_chain(saver: AsyncPostgresSaver) -> list[dict]:
thread_id = str(uuid4())
parent: dict | None = None
configs: list[dict] = []
for step in range(STEPS):
config, checkpoint, new_versions = _step_args(thread_id, step, parent)
parent = await saver.aput(
config,
checkpoint,
{"source": "loop", "step": step, "parents": {}},
new_versions,
)
await saver.aput_writes(parent, [(CHANNEL, f"w{step}")], str(uuid4()))
configs.append(parent)
return configs
def _build_chain(saver: PostgresSaver) -> list[dict]:
thread_id = str(uuid4())
parent: dict | None = None
configs: list[dict] = []
for step in range(STEPS):
config, checkpoint, new_versions = _step_args(thread_id, step, parent)
parent = saver.put(
config,
checkpoint,
{"source": "loop", "step": step, "parents": {}},
new_versions,
)
saver.put_writes(parent, [(CHANNEL, f"w{step}")], str(uuid4()))
configs.append(parent)
return configs
def _assert_history(entry: DeltaChannelHistory, page_size: int) -> None:
seed = entry.get("seed")
assert isinstance(seed, _DeltaSnapshot), (
f"page_size={page_size}: expected a snapshot seed, "
f"got {entry.get('seed', '<missing>')!r}"
)
assert seed.value == SEED_VALUE
assert [w[2] for w in entry["writes"]] == ["w1", "w2", "w3"], (
f"page_size={page_size}: got {[w[2] for w in entry['writes']]}"
)
@pytest.mark.parametrize("page_size", PAGE_SIZES)
async def test_async_target_older_than_the_first_page(
page_size: int, monkeypatch: pytest.MonkeyPatch
) -> None:
monkeypatch.setattr("langgraph.checkpoint.postgres.aio._DELTA_PAGE_SIZE", page_size)
async with AsyncPostgresSaver.from_conn_string(DEFAULT_URI) as saver:
await saver.setup()
configs = await _abuild_chain(saver)
result = await saver.aget_delta_channel_history(
config=configs[TARGET_STEP], channels=[CHANNEL]
)
_assert_history(result[CHANNEL], page_size)
@pytest.mark.parametrize("page_size", PAGE_SIZES)
def test_sync_target_older_than_the_first_page(
page_size: int, monkeypatch: pytest.MonkeyPatch
) -> None:
monkeypatch.setattr("langgraph.checkpoint.postgres._DELTA_PAGE_SIZE", page_size)
with PostgresSaver.from_conn_string(DEFAULT_URI) as saver:
saver.setup()
configs = _build_chain(saver)
result = saver.get_delta_channel_history(
config=configs[TARGET_STEP], channels=[CHANNEL]
)
_assert_history(result[CHANNEL], page_size)
async def test_root_target_has_no_history_and_still_terminates(
monkeypatch: pytest.MonkeyPatch,
) -> None:
monkeypatch.setattr("langgraph.checkpoint.postgres.aio._DELTA_PAGE_SIZE", 1)
async with AsyncPostgresSaver.from_conn_string(DEFAULT_URI) as saver:
await saver.setup()
configs = await _abuild_chain(saver)
result = await saver.aget_delta_channel_history(
config=configs[0], channels=[CHANNEL]
)
assert result[CHANNEL] == {"writes": []}
@@ -507,13 +507,12 @@ class SqliteSaver(BaseCheckpointSaver[str]):
Two-stage query: Two-stage query:
* Stage 1 (paged): newest-first slice of `checkpoints` returning * Stage 1 (streamed): recursive CTE over `checkpoints` following
`(checkpoint_id, parent_checkpoint_id, type, checkpoint)` per `parent_checkpoint_id` from the target, returning
ancestor. Sqlite has no JSONB, so we ship the full serialized `(checkpoint_id, type, checkpoint)` per ancestor. Sqlite has no
checkpoint blob and inspect `channel_values` in Python. Pages JSONB, so we ship the full serialized checkpoint blob and inspect
newest-first by `checkpoint_id` with a `< cursor` predicate; `channel_values` in Python. Stops reading when every channel has
page size is `DELTA_PAGE_SIZE`. Stops paging when every channel found its seed or the chain is exhausted.
has found its seed or the chain is exhausted.
* Stage 2 (per-channel UNION ALL): one branch per channel reading * Stage 2 (per-channel UNION ALL): one branch per channel reading
`writes` filtered to that channel's specific `chain_cids`. No `writes` filtered to that channel's specific `chain_cids`. No
@@ -538,12 +537,14 @@ class SqliteSaver(BaseCheckpointSaver[str]):
seeded: set[str] = set() seeded: set[str] = set()
with self.cursor(transaction=False) as cur: with self.cursor(transaction=False) as cur:
cur.execute(DELTA_STAGE1_SQL, (thread_id, checkpoint_ns, checkpoint_id)) cur.execute(
DELTA_STAGE1_SQL,
(thread_id, checkpoint_ns, checkpoint_id, thread_id, checkpoint_ns),
)
for row in cur: for row in cur:
cid, parent_cid, type_tag, blob = row cid, type_tag, blob = row
if step_walk_with_row( if step_walk_with_row(
cid=cid, cid=cid,
parent_cid=parent_cid,
type_tag=type_tag, type_tag=type_tag,
blob=blob, blob=blob,
target_id=checkpoint_id, target_id=checkpoint_id,
@@ -26,16 +26,33 @@ from typing import Any
from langgraph.checkpoint.base import DeltaChannelHistory, PendingWrite from langgraph.checkpoint.base import DeltaChannelHistory, PendingWrite
# Stage 1 streams ancestors of `target_cid` newest-first. The `<=` # Stage 1 streams target, then its ancestors nearest-first, by following
# predicate keeps target itself in the stream so we can read its # `parent_checkpoint_id` rather than id order: ids are only monotonic within
# `parent_checkpoint_id` from the first row without a separate lookup; # one process, so a range scan by id can miss a parent whose id sorts above
# the caller skips target's own writes/seed (matches the # its child's. Target is the anchor row; its own writes/seed are skipped
# `BaseCheckpointSaver` contract). # (matches the `BaseCheckpointSaver` contract).
#
# `put` is `INSERT OR REPLACE`, so re-putting an existing id under a
# descendant's config makes the chain a loop. `step_walk_with_row` stops on a
# repeated id; sqlite yields recursive rows lazily, so abandoning the cursor
# ends the recursion.
#
# `CROSS JOIN` pins `ancestors` as the outer loop, so each step is one primary
# key lookup. With a plain `JOIN` and no `ANALYZE` stats, sqlite can put
# `checkpoints` outside and scan the whole thread per step.
DELTA_STAGE1_SQL = ( DELTA_STAGE1_SQL = (
"WITH RECURSIVE ancestors(checkpoint_id, parent_checkpoint_id, type, "
"checkpoint) AS ("
"SELECT checkpoint_id, parent_checkpoint_id, type, checkpoint " "SELECT checkpoint_id, parent_checkpoint_id, type, checkpoint "
"FROM checkpoints " "FROM checkpoints "
"WHERE thread_id = ? AND checkpoint_ns = ? AND checkpoint_id <= ? " "WHERE thread_id = ? AND checkpoint_ns = ? AND checkpoint_id = ? "
"ORDER BY checkpoint_id DESC" "UNION ALL "
"SELECT c.checkpoint_id, c.parent_checkpoint_id, c.type, c.checkpoint "
"FROM ancestors a CROSS JOIN checkpoints c "
"ON c.checkpoint_id = a.parent_checkpoint_id "
"WHERE c.thread_id = ? AND c.checkpoint_ns = ?"
") "
"SELECT checkpoint_id, type, checkpoint FROM ancestors"
) )
@@ -68,7 +85,6 @@ def build_delta_stage2_sql(*, chain_lens: Sequence[int]) -> str:
def step_walk_with_row( def step_walk_with_row(
*, *,
cid: str, cid: str,
parent_cid: str | None,
type_tag: str, type_tag: str,
blob: bytes, blob: bytes,
target_id: str, target_id: str,
@@ -81,36 +97,32 @@ def step_walk_with_row(
) -> bool: ) -> bool:
"""Process one streamed stage-1 row in the merged ancestor walk. """Process one streamed stage-1 row in the merged ancestor walk.
The cursor returns (cid, parent_cid, type, blob) rows in The cursor returns (cid, type, blob) rows in walk order starting at
`checkpoint_id` DESC order starting at target. The first row is target. The first row is target itself and is skipped (target's own
target itself; we read its parent_cid to seed the walk and otherwise writes/seed are not part of the contract).
skip it (target's own writes/seed are not part of the contract).
For each subsequent row, if `cid` matches the walk's current For each subsequent row we deserialize the blob, append the cid to
position, we deserialize the blob, append the cid to every every not-yet-seeded channel's chain, and check `channel_values` for
not-yet-seeded channel's chain, and check `channel_values` for
seeds. The deserialized checkpoint is dropped before advancing — no seeds. The deserialized checkpoint is dropped before advancing — no
cross-row cache, so peak in-flight is one deserialized checkpoint. cross-row cache, so peak in-flight is one deserialized checkpoint.
Off-path rows (different branch on the same thread) advance the Returns True when the caller can stop iterating and close the cursor:
cursor without doing any work. every requested channel is seeded, or the chain revisited a checkpoint.
Returns True when every requested channel is seeded — the caller
can stop iterating and close the cursor.
""" """
if "started" not in walk_state: if "started" not in walk_state:
if cid == target_id: if cid == target_id:
walk_state["started"] = True walk_state["started"] = True
walk_state["cur_cid"] = parent_cid
walk_state["active"] = {ch for ch in channels if ch not in seeded} walk_state["active"] = {ch for ch in channels if ch not in seeded}
walk_state["walked"] = {cid}
# Not target yet (or target not present): keep streaming. # Not target yet (or target not present): keep streaming.
return False return False
active: set[str] = walk_state["active"] active: set[str] = walk_state["active"]
if not active: if not active:
return True return True
if cid != walk_state["cur_cid"]: walked: set[str] = walk_state["walked"]
# Off-path row from a sibling branch — skip without deserializing. if cid in walked:
return False return True
walked.add(cid)
for ch in active: for ch in active:
chain_by_ch[ch].append(cid) chain_by_ch[ch].append(cid)
ckpt = serde.loads_typed((type_tag, blob)) ckpt = serde.loads_typed((type_tag, blob))
@@ -120,7 +132,6 @@ def step_walk_with_row(
seeded.add(ch) seeded.add(ch)
active.discard(ch) active.discard(ch)
del ckpt, channel_values del ckpt, channel_values
walk_state["cur_cid"] = parent_cid
return not active return not active
@@ -625,8 +625,8 @@ class AsyncSqliteSaver(BaseCheckpointSaver[str]):
"""Fast-path override of `BaseCheckpointSaver.aget_delta_channel_history`. """Fast-path override of `BaseCheckpointSaver.aget_delta_channel_history`.
See `SqliteSaver.get_delta_channel_history` for design notes; this See `SqliteSaver.get_delta_channel_history` for design notes; this
is the async equivalent using `aiosqlite` cursors. Stage 1 pages is the async equivalent using `aiosqlite` cursors. Stage 1 streams
the parent chain newest-first and Python-deserializes each the parent chain from the target and Python-deserializes each
checkpoint blob to find per-channel snapshots; stage 2 fetches checkpoint blob to find per-channel snapshots; stage 2 fetches
only the relevant writes via per-channel UNION ALL. only the relevant writes via per-channel UNION ALL.
""" """
@@ -650,13 +650,13 @@ class AsyncSqliteSaver(BaseCheckpointSaver[str]):
async with self.lock, self.conn.cursor() as cur: async with self.lock, self.conn.cursor() as cur:
await cur.execute( await cur.execute(
DELTA_STAGE1_SQL, (thread_id, checkpoint_ns, checkpoint_id) DELTA_STAGE1_SQL,
(thread_id, checkpoint_ns, checkpoint_id, thread_id, checkpoint_ns),
) )
async for row in cur: async for row in cur:
cid, parent_cid, type_tag, blob = row cid, type_tag, blob = row
if step_walk_with_row( if step_walk_with_row(
cid=cid, cid=cid,
parent_cid=parent_cid,
type_tag=type_tag, type_tag=type_tag,
blob=blob, blob=blob,
target_id=checkpoint_id, target_id=checkpoint_id,
@@ -0,0 +1,104 @@
from __future__ import annotations
from typing import Any
import pytest
from langgraph.checkpoint.base import (
BaseCheckpointSaver,
Checkpoint,
DeltaChannelHistory,
empty_checkpoint,
)
from langgraph.checkpoint.sqlite import SqliteSaver
from langgraph.checkpoint.sqlite._delta import DELTA_STAGE1_SQL
from langgraph.checkpoint.sqlite.aio import AsyncSqliteSaver
CHANNEL = "ch"
CONFIG: dict[str, Any] = {"configurable": {"thread_id": "t", "checkpoint_ns": ""}}
EXPECTED: DeltaChannelHistory = {
"writes": [("task", CHANNEL, "write-root")],
"seed": "seed",
}
def _checkpoint(checkpoint_id: str, values: dict[str, Any]) -> Checkpoint:
value = empty_checkpoint()
value["id"] = checkpoint_id
value["channel_values"] = values
return value
PARENT_ID_ORDERS = [
pytest.param("z-older", "a-newer", id="parent_id_sorts_above_child"),
pytest.param("a-older", "z-newer", id="parent_id_sorts_below_child"),
]
@pytest.mark.parametrize(("root_id", "child_id"), PARENT_ID_ORDERS)
def test_sync_walk_reaches_parent_whatever_the_id_order(
root_id: str, child_id: str
) -> None:
with SqliteSaver.from_conn_string(":memory:") as saver:
root = saver.put(CONFIG, _checkpoint(root_id, {CHANNEL: "seed"}), {}, {})
saver.put_writes(root, [(CHANNEL, "write-root")], "task")
child = saver.put(root, _checkpoint(child_id, {}), {}, {})
got = saver.get_delta_channel_history(config=child, channels=[CHANNEL])
reference = BaseCheckpointSaver.get_delta_channel_history(
saver, config=child, channels=[CHANNEL]
)
assert got[CHANNEL] == EXPECTED
assert got[CHANNEL] == reference[CHANNEL], "fast path disagrees with base"
@pytest.mark.parametrize(("root_id", "child_id"), PARENT_ID_ORDERS)
async def test_async_walk_reaches_parent_whatever_the_id_order(
root_id: str, child_id: str
) -> None:
async with AsyncSqliteSaver.from_conn_string(":memory:") as saver:
root = await saver.aput(CONFIG, _checkpoint(root_id, {CHANNEL: "seed"}), {}, {})
await saver.aput_writes(root, [(CHANNEL, "write-root")], "task")
child = await saver.aput(root, _checkpoint(child_id, {}), {}, {})
got = await saver.aget_delta_channel_history(config=child, channels=[CHANNEL])
assert got[CHANNEL] == EXPECTED
def test_walk_reaches_root_of_long_chain_with_descending_ids() -> None:
steps = 40
with SqliteSaver.from_conn_string(":memory:") as saver:
parent = saver.put(
CONFIG, _checkpoint(f"id-{steps:03d}", {CHANNEL: "seed"}), {}, {}
)
saver.put_writes(parent, [(CHANNEL, "write-root")], "task")
for step in range(steps - 1, 0, -1):
parent = saver.put(parent, _checkpoint(f"id-{step:03d}", {}), {}, {})
got = saver.get_delta_channel_history(config=parent, channels=[CHANNEL])
assert got[CHANNEL] == EXPECTED
def test_walk_terminates_when_put_makes_the_parent_chain_cycle() -> None:
with SqliteSaver.from_conn_string(":memory:") as saver:
a = saver.put(CONFIG, _checkpoint("cid-a", {}), {}, {})
b = saver.put(a, _checkpoint("cid-b", {}), {}, {})
repoint_a_under_b = _checkpoint("cid-a", {})
saver.put(b, repoint_a_under_b, {}, {})
got = saver.get_delta_channel_history(config=b, channels=[CHANNEL])
assert got[CHANNEL] == {"writes": []}
def test_walk_step_looks_up_the_parent_by_primary_key() -> None:
with SqliteSaver.from_conn_string(":memory:") as saver:
saver.setup()
plan = [
row[3]
for row in saver.conn.execute(
f"EXPLAIN QUERY PLAN {DELTA_STAGE1_SQL}", ("t", "", "id", "t", "")
)
]
assert any(
step.startswith("SEARCH c ") and "checkpoint_id=?" in step for step in plan
), f"recursive step should look up the parent by key, got {plan}"
+2
View File
@@ -103,6 +103,8 @@ The CLI uses a `langgraph.json` configuration file with these key settings:
} }
``` ```
Git dependencies should use credential-free URLs. The CLI conservatively scans direct `langgraph.json` dependencies, common Python package files, uv project and lock files, and common Node.js package and lock files for HTTP Git URLs with userinfo. This check is not exhaustive: generated Docker builds can copy other files, including nested requirement or constraint files, into image layers without scanning them. For private dependencies, provide short-lived credentials through your build environment's secret-backed Git credential helper. Do not store credentials in copied files such as `langgraph.json` or `pip_config_file`.
See the [full documentation](https://reference.langchain.com/python/langgraph-cli) for detailed configuration options. See the [full documentation](https://reference.langchain.com/python/langgraph-cli) for detailed configuration options.
## Development ## Development
+1 -1
View File
@@ -1 +1 @@
__version__ = "0.4.31" __version__ = "0.4.32"
+87 -3
View File
@@ -6,6 +6,7 @@ import re
import shlex import shlex
import textwrap import textwrap
from collections import Counter from collections import Counter
from collections.abc import Iterable
from typing import Literal, NamedTuple from typing import Literal, NamedTuple
import click import click
@@ -36,6 +37,10 @@ DISALLOWED_BUILD_COMMAND_CHARS = [
# This blocks background execution (cmd &) while allowing command # This blocks background execution (cmd &) while allowing command
# chaining (cmd1 && cmd2) which is common in build commands. # chaining (cmd1 && cmd2) which is common in build commands.
_SINGLE_AMPERSAND_RE = re.compile(r"(?<!&)&(?:&&)*(?!&)") _SINGLE_AMPERSAND_RE = re.compile(r"(?<!&)&(?:&&)*(?!&)")
_GIT_HTTP_AUTHORITY_RES = (
re.compile(r"git\+https?://(?P<authority>[^/\s\"']+)", re.I),
re.compile(r"\bgit\s*=\s*[\"']https?://(?P<authority>[^/\s\"']+)", re.I),
)
_API_VERSION_PATTERN = re.compile( _API_VERSION_PATTERN = re.compile(
r"^(?P<major>\d+)" r"^(?P<major>\d+)"
r"(?:\.(?P<minor>\d+))?" r"(?:\.(?P<minor>\d+))?"
@@ -78,6 +83,62 @@ def has_disallowed_build_command_content(command: str) -> bool:
return False return False
def _has_git_http_url_userinfo(dependency: str) -> bool:
"""Check whether a Git HTTP URL contains userinfo."""
return any(
"@" in match.group("authority")
for pattern in _GIT_HTTP_AUTHORITY_RES
for match in pattern.finditer(dependency)
)
def _validate_git_http_url_userinfo(
values: Iterable[str], *, source: pathlib.Path | None = None
) -> None:
"""Reject credential-bearing Git HTTP URLs without echoing their values."""
if not any(_has_git_http_url_userinfo(value) for value in values):
return
message = (
"Git dependency URLs must not contain credentials or other URL "
"userinfo because generated Dockerfiles and image layers can retain "
"them. Use a credential-free Git URL and provide short-lived "
"credentials through your build environment's secret-backed Git "
"credential helper."
)
if source is not None:
message += f" Found in: {source}"
raise click.UsageError(message)
def _validate_git_http_url_userinfo_files(paths: Iterable[pathlib.Path]) -> None:
"""Reject credential-bearing Git HTTP URLs in dependency files."""
for path in paths:
path = path.resolve()
if not path.is_file():
continue
try:
contents = path.read_text(encoding="utf-8", errors="replace")
except OSError:
raise click.UsageError(
f"Could not inspect dependency file for embedded credentials: {path}"
) from None
_validate_git_http_url_userinfo([contents], source=path)
def _validate_local_dependency_files(config_path: pathlib.Path, config: Config) -> None:
"""Validate dependency files copied into a non-uv Python image."""
paths: list[pathlib.Path] = []
for dependency in config["dependencies"]:
if not isinstance(dependency, str) or not dependency.startswith("."):
continue
root = (config_path.parent / dependency).resolve()
paths.extend(
root / name
for name in ("requirements.txt", "pyproject.toml", "setup.py", "setup.cfg")
)
_validate_git_http_url_userinfo_files(paths)
MIN_PYTHON_VERSION = "3.11" MIN_PYTHON_VERSION = "3.11"
DEFAULT_PYTHON_VERSION = "3.11" DEFAULT_PYTHON_VERSION = "3.11"
@@ -320,7 +381,9 @@ def _get_source_kind(config: Config) -> str | None:
return kind if isinstance(kind, str) else None return kind if isinstance(kind, str) else None
def validate_config(config: Config) -> Config: def validate_config(
config: Config, *, source_path: pathlib.Path | None = None
) -> Config:
"""Validate a configuration dictionary.""" """Validate a configuration dictionary."""
graphs = config.get("graphs", {}) graphs = config.get("graphs", {})
@@ -415,6 +478,15 @@ def validate_config(config: Config) -> Config:
' "source": {"kind": "uv", "root": ".."}' ' "source": {"kind": "uv", "root": ".."}'
) )
_validate_git_http_url_userinfo(
(
dependency
for dependency in config["dependencies"]
if isinstance(dependency, str)
),
source=source_path,
)
source = config.get("source") source = config.get("source")
source_kind = _get_source_kind(config) source_kind = _get_source_kind(config)
if source is not None and not isinstance(source, dict): if source is not None and not isinstance(source, dict):
@@ -609,7 +681,7 @@ def validate_config_file(config_path: pathlib.Path) -> Config:
"""Load and validate a configuration file.""" """Load and validate a configuration file."""
with open(config_path) as f: with open(config_path) as f:
config = json.load(f) config = json.load(f)
validated = validate_config(config) validated = validate_config(config, source_path=config_path.resolve())
# Enforce the package.json doesn't enforce an # Enforce the package.json doesn't enforce an
# incompatible Node.js version # incompatible Node.js version
if validated.get("node_version"): if validated.get("node_version"):
@@ -1280,6 +1352,7 @@ def python_config_to_docker(
api_version=api_version, api_version=api_version,
build_tools_to_uninstall=build_tools_to_uninstall, build_tools_to_uninstall=build_tools_to_uninstall,
) )
_validate_local_dependency_files(config_path, config)
if pip_installer == "auto": if pip_installer == "auto":
if _image_supports_uv(base_image): if _image_supports_uv(base_image):
pip_installer = "uv" pip_installer = "uv"
@@ -1490,7 +1563,18 @@ def node_config_to_docker(
) -> tuple[str, dict[str, str]]: ) -> tuple[str, dict[str, str]]:
# Calculate paths for monorepo support # Calculate paths for monorepo support
install_root = ( install_root = (
pathlib.Path(build_context).resolve() if build_context else config_path.parent pathlib.Path(build_context).resolve()
if build_context
else config_path.parent.resolve()
)
config_root = config_path.parent.resolve()
dependency_roots = (
(install_root, config_root) if install_root != config_root else (install_root,)
)
_validate_git_http_url_userinfo_files(
root / name
for root in dependency_roots
for name in ("package.json", "package-lock.json", "yarn.lock", "pnpm-lock.yaml")
) )
install_cmd = install_command or _get_node_pm_install_cmd(install_root) install_cmd = install_command or _get_node_pm_install_cmd(install_root)
if build_context: if build_context:
+314 -79
View File
@@ -26,6 +26,7 @@ from langgraph_cli.dependency_tracking import find_tracked_packages
from langgraph_cli.docker import build_docker_image, can_build_locally from langgraph_cli.docker import build_docker_image, can_build_locally
from langgraph_cli.exec import CommandRunner, Runner, subp_exec from langgraph_cli.exec import CommandRunner, Runner, subp_exec
from langgraph_cli.host_backend import ( from langgraph_cli.host_backend import (
MAX_PAGE_SIZE,
ControlPlaneEndpoints, ControlPlaneEndpoints,
HostBackendClient, HostBackendClient,
HostBackendError, HostBackendError,
@@ -101,15 +102,15 @@ _NATIVE_AMD64_MACHINE = "x86_64"
_PUSH_ATTEMPTS = 3 _PUSH_ATTEMPTS = 3
_LOCAL_BUILD_TAG_PREFIX = "langgraph-deploy-tmp" _LOCAL_BUILD_TAG_PREFIX = "langgraph-deploy-tmp"
_OPERATOR_DEFAULT_RESOURCE_SPEC: Mapping[str, object] = {} _OPERATOR_DEFAULT_RESOURCE_SPEC: Mapping[str, object] = {}
_LISTENER_REQUIRED_MARKER = "listener_id' is required"
_HYBRID_LISTENER_GUIDANCE = (
"This workspace deploys through a listener in your own cluster, and the "
"control plane needs a listener ID to create a deployment. Create the "
"deployment once in the LangSmith UI, choosing the listener and namespace, "
"then re-run with --deployment-id <id>."
)
_CUSTOMER_REGISTRY_SOURCE: SourceName = "external_docker" _CUSTOMER_REGISTRY_SOURCE: SourceName = "external_docker"
_LISTENER_REQUIRED_MARKER = "listener_id' is required"
_LISTENERS_SHOWN = 10
_LISTENER_NOT_FOUND_STATUSES = frozenset({404, 422})
_LISTENERS_DOCS_URL = "https://docs.langchain.com/langsmith/control-plane#listeners"
_NO_LISTENERS = (
"This workspace has no listeners, so --listener-id and --k8s-namespace "
"do not apply."
)
_TERMINAL_STATUSES = frozenset( _TERMINAL_STATUSES = frozenset(
@@ -161,6 +162,134 @@ class ByAgent:
DeploymentSelector = ById | ByName | ByAgent DeploymentSelector = ById | ByName | ByAgent
@dataclass(frozen=True, slots=True)
class Listener:
id: str
compute_id: str
namespaces: tuple[str, ...]
@classmethod
def from_resource(cls, resource: Mapping[str, object]) -> "Listener":
identifier = str(resource.get("id") or "")
if not identifier:
raise HostBackendError(
"The control plane returned a listener without an id."
)
compute_config = resource.get("compute_config")
namespaces = (
compute_config.get("k8s_namespaces")
if isinstance(compute_config, Mapping)
else None
)
return cls(
identifier,
str(resource.get("compute_id", "")),
tuple(str(namespace) for namespace in namespaces)
if isinstance(namespaces, list)
else (),
)
@dataclass(frozen=True, slots=True)
class Unplaced:
@property
def summary(self) -> str:
return ""
def source_config(self) -> dict[str, object]:
return {}
@dataclass(frozen=True, slots=True)
class OnListener:
listener_id: str
k8s_namespace: str
@property
def summary(self) -> str:
return (
f"Deploying through listener {self.listener_id} "
f"in namespace {self.k8s_namespace}"
)
def source_config(self) -> dict[str, object]:
return {
"listener_id": self.listener_id,
"listener_config": {"k8s_namespace": self.k8s_namespace},
}
Placement = Unplaced | OnListener
@dataclass(frozen=True, slots=True)
class RequestedPlacement:
listener_id: str | None = None
k8s_namespace: str | None = None
@property
def requested(self) -> bool:
return self.listener_id is not None or self.k8s_namespace is not None
def ensure_not_requested(self, deployment_id: str) -> None:
if self.requested:
raise click.UsageError(
"Listener and namespace are fixed when a deployment is created. "
f"Deployment {deployment_id} already exists, so drop --listener-id "
"and --k8s-namespace, or create a new deployment with a different "
"--name."
)
def on(self, listener: Listener) -> Placement:
return OnListener(listener.id, self._namespace(listener))
def among(self, listeners: Sequence[Listener]) -> Placement:
if not listeners:
if self.requested:
raise click.UsageError(_NO_LISTENERS)
return Unplaced()
if len(listeners) > 1:
raise click.UsageError(
"This workspace has several listeners. Choose one with "
f"--listener-id:\n{_describe_listeners(listeners)}"
)
return self.on(listeners[0])
def _namespace(self, listener: Listener) -> str:
if not listener.namespaces:
raise click.UsageError(
f"Listener {listener.id} serves no namespaces. Check its configuration."
)
if self.k8s_namespace is None:
if len(listener.namespaces) == 1:
return listener.namespaces[0]
raise click.UsageError(
f"Listener {listener.id} serves several namespaces. Choose one with "
f"--k8s-namespace: {', '.join(listener.namespaces)}"
)
if self.k8s_namespace not in listener.namespaces:
raise click.UsageError(
f"Listener {listener.id} does not serve namespace "
f"'{self.k8s_namespace}'. Choose one of: "
f"{', '.join(listener.namespaces)}"
)
return self.k8s_namespace
def _describe_listeners(listeners: Sequence[Listener]) -> str:
shown = listeners[:_LISTENERS_SHOWN]
lines = [
f" {listener.id} cluster {listener.compute_id} "
f"namespaces: {', '.join(listener.namespaces)}"
for listener in shown
]
if len(listeners) > len(shown):
lines.append(f" ... and {len(listeners) - len(shown)} more")
if len(listeners) == MAX_PAGE_SIZE:
lines.append(f" (only the first {MAX_PAGE_SIZE} listeners were read)")
return "\n".join(lines)
@dataclass(frozen=True, slots=True) @dataclass(frozen=True, slots=True)
class ExistingDeployment: class ExistingDeployment:
id: str id: str
@@ -379,15 +508,16 @@ def _source_of(resource: object) -> str | None:
def find_deployment_by_name( def find_deployment_by_name(
client: HostBackendClient, name: str client: HostBackendClient, name: str
) -> ExistingDeployment | None: ) -> ExistingDeployment | None:
listed = client.list_deployments(name_contains=name) listed = client.list_deployments(name=name, name_contains=name, limit=MAX_PAGE_SIZE)
resources = listed.get("resources", []) if isinstance(listed, dict) else [] for resource in listed:
for resource in resources: if resource.get("name") == name and resource.get("id"):
if (
isinstance(resource, dict)
and resource.get("name") == name
and resource.get("id")
):
return ExistingDeployment(str(resource["id"]), _source_of(resource)) return ExistingDeployment(str(resource["id"]), _source_of(resource))
if len(listed) >= MAX_PAGE_SIZE:
raise click.ClickException(
"This workspace has more deployments than the CLI can search, so it "
f"cannot tell whether '{name}' already exists. Pass --deployment-id to "
"update an existing deployment."
)
return None return None
@@ -683,14 +813,22 @@ def _find_deployment(
existing = _call_host_backend_with_optional_tenant( existing = _call_host_backend_with_optional_tenant(
client, client,
lambda c: c.list_deployments( lambda c: c.list_deployments(
agent_id=selector.agent_id, agent_environment=selector.environment agent_id=selector.agent_id,
agent_environment=selector.environment,
limit=MAX_PAGE_SIZE,
), ),
) )
if len(existing) > 1:
raise click.ClickException(
"This control plane does not filter deployments by agent, so the "
f"CLI cannot tell which one belongs to '{selector.agent_id}' in "
f"{selector.environment}. Deploy by --name instead."
)
found = next( found = next(
( (
ExistingDeployment(str(dep["id"]), _source_of(dep)) ExistingDeployment(str(dep["id"]), _source_of(dep))
for dep in existing.get("resources", []) for dep in existing
if not dep.get("is_preview") if dep.get("id") and not dep.get("is_preview")
), ),
None, None,
) )
@@ -758,21 +896,18 @@ def _create_deployment(
def _get_deployment_status_url( def _get_deployment_status_url(
updated: object, deployment_id: str, host_url: str updated: object, deployment_id: str, endpoints: ControlPlaneEndpoints
) -> str | None: ) -> str | None:
"""Compute the LangSmith dashboard URL for a deployment, if possible."""
tenant_id = updated.get("tenant_id") if isinstance(updated, dict) else None tenant_id = updated.get("tenant_id") if isinstance(updated, dict) else None
if not tenant_id: if not tenant_id:
return None return None
base = ControlPlaneEndpoints.from_control_plane_url(host_url).dashboard_url return f"{endpoints.dashboard_url}/o/{tenant_id}/host/deployments/{deployment_id}"
return f"{base}/o/{tenant_id}/host/deployments/{deployment_id}"
def _emit_deployment_status_url( def _emit_deployment_status_url(
updated: object, deployment_id: str, host_url: str updated: object, deployment_id: str, endpoints: ControlPlaneEndpoints
) -> str | None: ) -> str | None:
"""Emit the deployment status URL and return it.""" url = _get_deployment_status_url(updated, deployment_id, endpoints)
url = _get_deployment_status_url(updated, deployment_id, host_url)
if url: if url:
_get_emitter().status_url(url) _get_emitter().status_url(url)
return url return url
@@ -790,14 +925,11 @@ def _poll_revision_status(
) -> tuple[str, str | None]: ) -> tuple[str, str | None]:
"""Poll latest revision status until terminal status or timeout.""" """Poll latest revision status until terminal status or timeout."""
em = _get_emitter() em = _get_emitter()
revisions_resp = client.list_revisions(deployment_id, limit=1) revisions = client.list_revisions(deployment_id, limit=1)
resources = ( if not revisions:
revisions_resp.get("resources", []) if isinstance(revisions_resp, dict) else []
)
if not resources:
return "", None return "", None
revision_id = str(resources[0]["id"]) revision_id = str(revisions[0]["id"])
last_status = "" last_status = ""
deadline = time.time() + timeout_seconds deadline = time.time() + timeout_seconds
start_time = time.monotonic() start_time = time.monotonic()
@@ -1318,6 +1450,7 @@ def _run_remote_build(
@dataclass(frozen=True, slots=True) @dataclass(frozen=True, slots=True)
class DeployContext: class DeployContext:
client: HostBackendClient client: HostBackendClient
endpoints: ControlPlaneEndpoints
spec: BuildSpec spec: BuildSpec
verbose: bool verbose: bool
selector: DeploymentSelector selector: DeploymentSelector
@@ -1347,19 +1480,66 @@ def _resolve_or_create(
) )
if found is not None: if found is not None:
return found.id, step return found.id, step
created, step = _create_deployment( try:
ctx.client, created, step = _create_deployment(
step, ctx.client,
name=ctx.selector.name if isinstance(ctx.selector, ByName) else None, step,
agent=asdict(ctx.selector) if isinstance(ctx.selector, ByAgent) else None, name=ctx.selector.name if isinstance(ctx.selector, ByName) else None,
source=source, agent=asdict(ctx.selector) if isinstance(ctx.selector, ByAgent) else None,
source_config={"deployment_type": ctx.deployment_type}, source=source,
source_revision_config={}, source_config={"deployment_type": ctx.deployment_type},
secrets=ctx.secrets, source_revision_config={},
) secrets=ctx.secrets,
)
except HostBackendError as err:
if _needs_a_listener(err):
raise ListenerRequiredError(
"The image has to come from a registry you manage, so re-run with "
"--push-to <registry>/<repository>."
) from None
raise
return created.id, step return created.id, step
class ListenerRequiredError(click.UsageError):
def __init__(self, remedy: str) -> None:
super().__init__(
"This workspace deploys through a listener in your own cluster. "
f"{remedy}\nLearn about listeners: {_LISTENERS_DOCS_URL}"
)
def _needs_a_listener(err: HostBackendError) -> bool:
return err.status_code == 400 and _LISTENER_REQUIRED_MARKER in (
err.detail or err.message
)
def _requested_listener(client: HostBackendClient, listener_id: str) -> Listener:
try:
resource = _call_host_backend_with_optional_tenant(
client, lambda c: c.get_listener(listener_id)
)
except HostBackendError as err:
if err.status_code not in _LISTENER_NOT_FOUND_STATUSES:
raise
available = _available_listeners(client)
if not available:
raise click.UsageError(_NO_LISTENERS) from None
raise click.UsageError(
f"Listener {listener_id} was not found in this workspace. "
f"Available listeners:\n{_describe_listeners(available)}"
) from None
return Listener.from_resource(resource)
def _available_listeners(client: HostBackendClient) -> tuple[Listener, ...]:
resources = _call_host_backend_with_optional_tenant(
client, lambda c: c.list_listeners()
)
return tuple(Listener.from_resource(resource) for resource in resources)
def _ensure_customer_registry_source(existing: ExistingDeployment) -> None: def _ensure_customer_registry_source(existing: ExistingDeployment) -> None:
if existing.source != _CUSTOMER_REGISTRY_SOURCE: if existing.source != _CUSTOMER_REGISTRY_SOURCE:
raise click.UsageError( raise click.UsageError(
@@ -1422,6 +1602,7 @@ class RemoteBuildSource:
class CustomerRegistrySource: class CustomerRegistrySource:
reference: ImageReference reference: ImageReference
prebuilt_image: str | None prebuilt_image: str | None
requested_placement: RequestedPlacement
def run(self, ctx: DeployContext) -> DeployOutcome: def run(self, ctx: DeployContext) -> DeployOutcome:
if isinstance(ctx.selector, ById): if isinstance(ctx.selector, ById):
@@ -1443,6 +1624,7 @@ class CustomerRegistrySource:
self, ctx: DeployContext, existing: ExistingDeployment, step: int self, ctx: DeployContext, existing: ExistingDeployment, step: int
) -> DeployOutcome: ) -> DeployOutcome:
_ensure_customer_registry_source(existing) _ensure_customer_registry_source(existing)
self.requested_placement.ensure_not_requested(existing.id)
image_uri, step = self._publish(ctx, step) image_uri, step = self._publish(ctx, step)
_log_deploy_step(step, f"Updating deployment {existing.id}") _log_deploy_step(step, f"Updating deployment {existing.id}")
updated = ctx.client.update_deployment( updated = ctx.client.update_deployment(
@@ -1456,7 +1638,25 @@ class CustomerRegistrySource:
existing.id, _image_revision_result(updated, "Deployment updated") existing.id, _image_revision_result(updated, "Deployment updated")
) )
def _resolve_placement(self, ctx: DeployContext) -> Placement:
requested = self.requested_placement
if requested.listener_id is not None:
return requested.on(_requested_listener(ctx.client, requested.listener_id))
if not (ctx.endpoints.is_cloud or requested.requested):
return Unplaced()
return requested.among(_available_listeners(ctx.client))
def _announce(self, placement: Placement) -> None:
if isinstance(placement, OnListener):
_get_emitter().info(
placement.summary,
listener_id=placement.listener_id,
k8s_namespace=placement.k8s_namespace,
)
def _create(self, ctx: DeployContext, name: str | None, step: int) -> DeployOutcome: def _create(self, ctx: DeployContext, name: str | None, step: int) -> DeployOutcome:
placement = self._resolve_placement(ctx)
self._announce(placement)
image_uri, step = self._publish(ctx, step) image_uri, step = self._publish(ctx, step)
try: try:
created, _ = _create_deployment( created, _ = _create_deployment(
@@ -1467,13 +1667,19 @@ class CustomerRegistrySource:
if isinstance(ctx.selector, ByAgent) if isinstance(ctx.selector, ByAgent)
else None, else None,
source=_CUSTOMER_REGISTRY_SOURCE, source=_CUSTOMER_REGISTRY_SOURCE,
source_config={"resource_spec": _OPERATOR_DEFAULT_RESOURCE_SPEC}, source_config={
"resource_spec": _OPERATOR_DEFAULT_RESOURCE_SPEC,
**placement.source_config(),
},
source_revision_config={"image_uri": image_uri}, source_revision_config={"image_uri": image_uri},
secrets=ctx.secrets, secrets=ctx.secrets,
) )
except HostBackendError as err: except HostBackendError as err:
if err.status_code == 400 and _LISTENER_REQUIRED_MARKER in err.message: if _needs_a_listener(err):
raise click.ClickException(_HYBRID_LISTENER_GUIDANCE) from None raise ListenerRequiredError(
"Re-run with --listener-id and --k8s-namespace.\n"
f"{err.detail or err.message}"
) from None
raise raise
return DeployOutcome( return DeployOutcome(
created.id, _image_revision_result(created.resource, "Deployment created") created.id, _image_revision_result(created.resource, "Deployment created")
@@ -1534,14 +1740,31 @@ def _select_source(
image_name: str | None, image_name: str | None,
tag: str | None, tag: str | None,
remote_build_flag: bool | None, remote_build_flag: bool | None,
placement: RequestedPlacement,
selector: DeploymentSelector,
) -> DeploymentSource: ) -> DeploymentSource:
if push_to is None and placement.requested:
raise click.UsageError(
"--listener-id and --k8s-namespace only apply when creating a "
"deployment with --push-to."
)
if placement.requested and isinstance(selector, ById):
raise click.UsageError(
"Listener and namespace are fixed when a deployment is created, so "
"they cannot be set for an existing --deployment-id. Drop them, or "
"create a new deployment with --name."
)
if push_to is not None: if push_to is not None:
if remote_build_flag is True: if remote_build_flag is True:
raise click.UsageError("--push-to cannot be combined with --remote.") raise click.UsageError("--push-to cannot be combined with --remote.")
reference = _push_reference(push_to, tag) reference = _push_reference(push_to, tag)
if image is None: if image is None:
_require_local_docker() _require_local_docker()
return CustomerRegistrySource(reference, prebuilt_image=image) return CustomerRegistrySource(
reference=reference,
prebuilt_image=image,
requested_placement=placement,
)
if image and remote_build_flag is True: if image and remote_build_flag is True:
raise click.UsageError("--image cannot be combined with --remote builds.") raise click.UsageError("--image cannot be combined with --remote builds.")
use_remote_build, local_build_error = _resolve_build_mode( use_remote_build, local_build_error = _resolve_build_mode(
@@ -1647,9 +1870,7 @@ def _call_host_backend_with_optional_tenant(
prompted_for_tenant = True prompted_for_tenant = True
continue continue
if err.status_code == 403 and "not enabled" in err.message.lower(): if err.status_code == 403 and "not enabled" in err.message.lower():
smith_base = ControlPlaneEndpoints.from_control_plane_url( smith_base = client.endpoints.dashboard_url
client.base_url
).dashboard_url
raise HostBackendError( raise HostBackendError(
"LangSmith Deployment is not enabled for this organization. " "LangSmith Deployment is not enabled for this organization. "
f"Enable it at {smith_base}/host/deployments" f"Enable it at {smith_base}/host/deployments"
@@ -1690,11 +1911,17 @@ OPT_HOST_URL = click.option(
) )
OPT_AGENT_ID = click.option( OPT_AGENT_ID = click.option(
"--agent-id", help="Logical agent ID (requires agent mode enabled for the tenant)." "--agent-id",
envvar="LANGSMITH_AGENT_ID",
show_envvar=True,
help="Logical agent ID (requires agent mode enabled for the tenant).",
) )
OPT_AGENT_ENVIRONMENT = click.option( OPT_AGENT_ENVIRONMENT = click.option(
"--environment", "--agent-environment",
"environment",
envvar="LANGSMITH_AGENT_ENVIRONMENT",
show_envvar=True,
type=click.Choice(["development", "staging", "production"]), type=click.Choice(["development", "staging", "production"]),
help="Agent environment (requires agent mode enabled for the tenant).", help="Agent environment (requires agent mode enabled for the tenant).",
) )
@@ -1848,6 +2075,21 @@ def _deploy_base_options(
"Give the tag here or with --tag (default: latest)." "Give the tag here or with --tag (default: latest)."
), ),
), ),
click.option(
"--listener-id",
help=(
"Listener that will run the deployment, for workspaces that "
"deploy through a listener in your own cluster. Only used when "
"creating a deployment with --push-to."
),
),
click.option(
"--k8s-namespace",
help=(
"Kubernetes namespace the listener deploys into. Only used when "
"creating a deployment with --push-to."
),
),
click.option( click.option(
"--config", "--config",
"-c", "-c",
@@ -1958,6 +2200,8 @@ def _deploy_cmd(
image_name: str | None, image_name: str | None,
image: str | None, image: str | None,
push_to: str | None, push_to: str | None,
listener_id: str | None,
k8s_namespace: str | None,
tag: str | None, tag: str | None,
base_image: str | None, base_image: str | None,
install_command: str | None, install_command: str | None,
@@ -1982,13 +2226,14 @@ def _deploy_cmd(
validate_deploy_commands(install_command, build_command) validate_deploy_commands(install_command, build_command)
agent = None agent = None
if agent_id is not None or environment is not None: if agent_id is not None or environment is not None:
em.note("Note: --agent-id and --agent-environment flags are in private beta")
if not agent_id or not agent_id.strip() or not environment: if not agent_id or not agent_id.strip() or not environment:
raise click.UsageError( raise click.UsageError(
"--agent-id and --environment are required together." "--agent-id and --agent-environment are required together."
) )
if name is not None or deployment_id is not None: if name is not None or deployment_id is not None:
raise click.UsageError( raise click.UsageError(
"--agent-id and --environment cannot be combined with --name or --deployment-id." "--agent-id and --agent-environment cannot be combined with --name or --deployment-id."
) )
agent = {"agent_id": agent_id, "environment": environment} agent = {"agent_id": agent_id, "environment": environment}
if not config.exists(): if not config.exists():
@@ -2024,12 +2269,15 @@ def _deploy_cmd(
secrets = _secrets_from_env(_env_without_deployment_name(env_vars)) secrets = _secrets_from_env(_env_without_deployment_name(env_vars))
selector = ByAgent(**agent) if agent else deployment_selector(deployment_id, name)
source = _select_source( source = _select_source(
push_to=push_to, push_to=push_to,
image=image, image=image,
image_name=image_name, image_name=image_name,
tag=tag, tag=tag,
remote_build_flag=remote_build_flag, remote_build_flag=remote_build_flag,
placement=RequestedPlacement(listener_id, k8s_namespace),
selector=selector,
) )
client = _create_host_backend_client(host_url, api_key, env_vars=env_vars) client = _create_host_backend_client(host_url, api_key, env_vars=env_vars)
@@ -2042,6 +2290,7 @@ def _deploy_cmd(
outcome = source.run( outcome = source.run(
DeployContext( DeployContext(
client=client, client=client,
endpoints=client.endpoints,
spec=BuildSpec( spec=BuildSpec(
config=config, config=config,
config_json=config_json, config_json=config_json,
@@ -2053,9 +2302,7 @@ def _deploy_cmd(
build_command=build_command, build_command=build_command,
), ),
verbose=verbose, verbose=verbose,
selector=ByAgent(**agent) selector=selector,
if agent
else deployment_selector(deployment_id, name),
deployment_type=deployment_type, deployment_type=deployment_type,
secrets=secrets, secrets=secrets,
tracked_packages=tracked_packages, tracked_packages=tracked_packages,
@@ -2064,7 +2311,7 @@ def _deploy_cmd(
dep_status_url = _emit_deployment_status_url( dep_status_url = _emit_deployment_status_url(
outcome.build_result.updated, outcome.build_result.updated,
outcome.deployment_id, outcome.deployment_id,
client.base_url, client.endpoints,
) )
if no_wait: if no_wait:
@@ -2138,6 +2385,11 @@ def deploy_list(
agent_id: str | None, agent_id: str | None,
environment: str | None, environment: str | None,
) -> None: ) -> None:
if agent_id is not None or environment is not None:
click.secho(
"Note: --agent-id and --agent-environment flags are in private beta",
fg="yellow",
)
if agent_id is not None and not agent_id.strip(): if agent_id is not None and not agent_id.strip():
raise click.UsageError("--agent-id must not be empty.") raise click.UsageError("--agent-id must not be empty.")
filters = {} filters = {}
@@ -2146,16 +2398,10 @@ def deploy_list(
if environment is not None: if environment is not None:
filters["agent_environment"] = environment filters["agent_environment"] = environment
client = _create_host_backend_client(host_url, api_key) client = _create_host_backend_client(host_url, api_key)
response = _call_host_backend_with_optional_tenant( deployments = _call_host_backend_with_optional_tenant(
client, client,
lambda c: c.list_deployments(name_contains=name_contains, **filters), lambda c: c.list_deployments(name_contains=name_contains, **filters),
) )
resources = response.get("resources") if isinstance(response, dict) else None
deployments = (
[item for item in resources if isinstance(item, dict)]
if isinstance(resources, list)
else []
)
if not deployments: if not deployments:
click.echo("No deployments found.") click.echo("No deployments found.")
return return
@@ -2195,16 +2441,10 @@ def deploy_revisions_list(
api_key: str | None, host_url: str | None, limit: int, deployment_id: str api_key: str | None, host_url: str | None, limit: int, deployment_id: str
) -> None: ) -> None:
client = _create_host_backend_client(host_url, api_key) client = _create_host_backend_client(host_url, api_key)
response = _call_host_backend_with_optional_tenant( revisions = _call_host_backend_with_optional_tenant(
client, client,
lambda c: c.list_revisions(deployment_id, limit=limit), lambda c: c.list_revisions(deployment_id, limit=limit),
) )
resources = response.get("resources") if isinstance(response, dict) else None
revisions = (
[item for item in resources if isinstance(item, dict)]
if isinstance(resources, list)
else []
)
if not revisions: if not revisions:
click.echo(f"No revisions found for deployment {deployment_id}.") click.echo(f"No revisions found for deployment {deployment_id}.")
return return
@@ -2354,17 +2594,12 @@ def deploy_logs(
dep_id = found.id dep_id = found.id
if log_type == "build" and not revision_id: if log_type == "build" and not revision_id:
revisions_resp = client.list_revisions(dep_id, limit=1) revisions = client.list_revisions(dep_id, limit=1)
resources = ( if not revisions:
revisions_resp.get("resources", [])
if isinstance(revisions_resp, dict)
else []
)
if not resources:
raise click.ClickException( raise click.ClickException(
"No revisions found for this deployment. Cannot fetch build logs." "No revisions found for this deployment. Cannot fetch build logs."
) )
revision_id = str(resources[0]["id"]) revision_id = str(revisions[0]["id"])
click.secho(f"Using latest revision: {revision_id}", fg="cyan") click.secho(f"Using latest revision: {revision_id}", fg="cyan")
payload: dict = {"limit": limit, "order": "desc"} payload: dict = {"limit": limit, "order": "desc"}
+72 -19
View File
@@ -18,6 +18,7 @@ CLOUD_DASHBOARD_HOST = "smith.langchain.com"
CONTROL_PLANE_PATH = "/api-host" CONTROL_PLANE_PATH = "/api-host"
LANGSMITH_API_PATHS = ("/api/v1", "/api") LANGSMITH_API_PATHS = ("/api/v1", "/api")
LOCAL_HOSTNAMES = ("localhost", "127.0.0.1") LOCAL_HOSTNAMES = ("localhost", "127.0.0.1")
MAX_PAGE_SIZE = 100
SourceName = Literal["internal_docker", "internal_source", "external_docker"] SourceName = Literal["internal_docker", "internal_source", "external_docker"]
@@ -36,6 +37,13 @@ class ControlPlaneEndpoints:
return cls.from_langsmith_endpoint(langsmith_endpoint) return cls.from_langsmith_endpoint(langsmith_endpoint)
return cls(CLOUD_CONTROL_PLANE_URL, CLOUD_DASHBOARD_URL) return cls(CLOUD_CONTROL_PLANE_URL, CLOUD_DASHBOARD_URL)
@property
def is_cloud(self) -> bool:
hostname = urlparse(self.control_plane_url).hostname or ""
return hostname == CLOUD_CONTROL_PLANE_HOST or hostname.endswith(
f".{CLOUD_CONTROL_PLANE_HOST}"
)
@classmethod @classmethod
def from_control_plane_url(cls, url: str) -> ControlPlaneEndpoints: def from_control_plane_url(cls, url: str) -> ControlPlaneEndpoints:
control_plane_url = url.rstrip("/") control_plane_url = url.rstrip("/")
@@ -83,12 +91,36 @@ def _without_api_path(path: str) -> str:
return path return path
def _resources(payload: object) -> list[dict[str, Any]]:
if not isinstance(payload, dict):
return []
resources = payload.get("resources")
if not isinstance(resources, list):
return []
return [item for item in resources if isinstance(item, dict)]
class HostBackendError(click.ClickException): class HostBackendError(click.ClickException):
"""Raised when the host backend returns an error response.""" """Raised when the host backend returns an error response."""
def __init__(self, message: str, status_code: int | None = None): def __init__(
self,
message: str,
status_code: int | None = None,
detail: str | None = None,
):
super().__init__(message) super().__init__(message)
self.status_code = status_code self.status_code = status_code
self.detail = detail
def _error_detail(response: httpx.Response) -> str | None:
try:
body = response.json()
except ValueError:
return None
detail = body.get("detail") if isinstance(body, dict) else None
return detail if isinstance(detail, str) else None
class HostBackendClient: class HostBackendClient:
@@ -110,7 +142,8 @@ class HostBackendClient:
} }
if tenant_id: if tenant_id:
headers["X-Tenant-ID"] = tenant_id headers["X-Tenant-ID"] = tenant_id
self._base_url = base_url.rstrip("/") self._endpoints = ControlPlaneEndpoints.from_control_plane_url(base_url)
self._base_url = self._endpoints.control_plane_url
self._client = httpx.Client( self._client = httpx.Client(
base_url=self._base_url, base_url=self._base_url,
headers=headers, headers=headers,
@@ -122,6 +155,10 @@ class HostBackendClient:
def base_url(self) -> str: def base_url(self) -> str:
return self._base_url return self._base_url
@property
def endpoints(self) -> ControlPlaneEndpoints:
return self._endpoints
def set_tenant(self, tenant_id: str) -> None: def set_tenant(self, tenant_id: str) -> None:
self._client.headers["X-Tenant-ID"] = tenant_id self._client.headers["X-Tenant-ID"] = tenant_id
@@ -136,10 +173,12 @@ class HostBackendClient:
resp = self._client.request(method, path, json=payload, params=params) resp = self._client.request(method, path, json=payload, params=params)
resp.raise_for_status() resp.raise_for_status()
except httpx.HTTPStatusError as err: except httpx.HTTPStatusError as err:
detail = err.response.text or str(err.response.status_code) detail = _error_detail(err.response)
reason = detail or err.response.text or str(err.response.status_code)
raise HostBackendError( raise HostBackendError(
f"{method} {path} failed with status {err.response.status_code}: {detail}", f"{method} {path} failed with status {err.response.status_code}: {reason}",
status_code=err.response.status_code, status_code=err.response.status_code,
detail=detail,
) from None ) from None
except httpx.TransportError as err: except httpx.TransportError as err:
raise HostBackendError(str(err)) from None raise HostBackendError(str(err)) from None
@@ -178,20 +217,29 @@ class HostBackendClient:
def list_deployments( def list_deployments(
self, self,
name_contains: str = "",
*, *,
name: str | None = None,
name_contains: str | None = None,
limit: int | None = None,
agent_id: str | None = None, agent_id: str | None = None,
agent_environment: str | None = None, agent_environment: str | None = None,
) -> dict[str, Any]: ) -> list[dict[str, Any]]:
params = {"name_contains": name_contains} given = (
if agent_id is not None: ("name", name),
params["agent_id"] = agent_id ("name_contains", name_contains),
if agent_environment is not None: ("limit", limit),
params["agent_environment"] = agent_environment ("agent_id", agent_id),
return self._request( ("agent_environment", agent_environment),
"GET", )
"/v2/deployments", params = {key: value for key, value in given if value is not None}
params=params, return _resources(self._request("GET", "/v2/deployments", params=params))
def get_listener(self, listener_id: str) -> dict[str, Any]:
return self._request("GET", f"/v2/listeners/{listener_id}")
def list_listeners(self) -> list[dict[str, Any]]:
return _resources(
self._request("GET", "/v2/listeners", params={"limit": MAX_PAGE_SIZE})
) )
def get_deployment(self, deployment_id: str) -> dict[str, Any]: def get_deployment(self, deployment_id: str) -> dict[str, Any]:
@@ -266,10 +314,15 @@ class HostBackendClient:
payload["secrets"] = secrets payload["secrets"] = secrets
return self._request("PATCH", f"/v2/deployments/{deployment_id}", payload) return self._request("PATCH", f"/v2/deployments/{deployment_id}", payload)
def list_revisions(self, deployment_id: str, limit: int = 1) -> dict[str, Any]: def list_revisions(
return self._request( self, deployment_id: str, limit: int = 1
"GET", ) -> list[dict[str, Any]]:
f"/v2/deployments/{deployment_id}/revisions?limit={limit}", return _resources(
self._request(
"GET",
f"/v2/deployments/{deployment_id}/revisions",
params={"limit": limit},
)
) )
def get_revision(self, deployment_id: str, revision_id: str) -> dict[str, Any]: def get_revision(self, deployment_id: str, revision_id: str) -> dict[str, Any]:
+5 -1
View File
@@ -650,7 +650,8 @@ class Config(TypedDict, total=False):
pip_config_file: str | None pip_config_file: str | None
"""Optional. Path to a pip config file (e.g., "/etc/pip.conf" or "pip.ini") for controlling """Optional. Path to a pip config file (e.g., "/etc/pip.conf" or "pip.ini") for controlling
package installation (custom indices, credentials, etc.). package installation (custom indices, timeouts, etc.). The file is copied into the
generated image, so it must not contain credentials or other secrets.
Only relevant if Python dependencies are installed via pip. If omitted, default pip settings are used. Only relevant if Python dependencies are installed via pip. If omitted, default pip settings are used.
""" """
@@ -689,6 +690,9 @@ class Config(TypedDict, total=False):
- "." or "./src" if you have a local Python package - "." or "./src" if you have a local Python package
- str (aka "anthropic") for a PyPI package - str (aka "anthropic") for a PyPI package
- "git+https://github.com/org/repo.git@main" for a Git-based package - "git+https://github.com/org/repo.git@main" for a Git-based package
Git HTTP URLs must not contain userinfo such as a username or token. For private
dependencies, provide short-lived credentials through the build environment's
secret-backed Git credential helper.
Defaults to an empty list, meaning no additional packages installed beyond your base environment. Defaults to an empty list, meaning no additional packages installed beyond your base environment.
This field is not supported when `source.kind` is `uv`. This field is not supported when `source.kind` is `uv`.
+10
View File
@@ -880,6 +880,7 @@ def python_config_to_docker_uv_lock(
_get_node_pm_install_cmd, _get_node_pm_install_cmd,
_get_pip_cleanup_lines, _get_pip_cleanup_lines,
_image_supports_uv, _image_supports_uv,
_validate_git_http_url_userinfo_files,
docker_tag, docker_tag,
) )
@@ -890,11 +891,20 @@ def python_config_to_docker_uv_lock(
) )
config_root = config_path.parent.resolve() config_root = config_path.parent.resolve()
source_root = config["source"].get("root", ".")
project_root = (config_root / source_root).resolve()
_validate_git_http_url_userinfo_files(
[project_root / "pyproject.toml", project_root / "uv.lock"]
)
install_cmd = "uv pip install --system" install_cmd = "uv pip install --system"
_, global_reqs_pip_install, pip_config_file_str = _build_python_install_commands( _, global_reqs_pip_install, pip_config_file_str = _build_python_install_commands(
config, install_cmd config, install_cmd
) )
plan = _plan_uv_lock_workspace(config_path, config) plan = _plan_uv_lock_workspace(config_path, config)
_validate_git_http_url_userinfo_files(
package.pyproject_path for package in plan.install_order
)
_update_uv_lock_graph_paths(config_path, config, plan) _update_uv_lock_graph_paths(config_path, config, plan)
for section, key in [ for section, key in [
+2 -2
View File
@@ -28,7 +28,7 @@
"type": "null" "type": "null"
} }
], ],
"description": "Optional. Path to a pip config file (e.g., \"/etc/pip.conf\" or \"pip.ini\") for controlling\npackage installation (custom indices, credentials, etc.).\n\nOnly relevant if Python dependencies are installed via pip. If omitted, default pip settings are used.\n" "description": "Optional. Path to a pip config file (e.g., \"/etc/pip.conf\" or \"pip.ini\") for controlling\npackage installation (custom indices, timeouts, etc.). The file is copied into the\ngenerated image, so it must not contain credentials or other secrets.\n\nOnly relevant if Python dependencies are installed via pip. If omitted, default pip settings are used.\n"
}, },
"_INTERNAL_docker_tag": { "_INTERNAL_docker_tag": {
"anyOf": [ "anyOf": [
@@ -270,7 +270,7 @@
"type": "null" "type": "null"
} }
], ],
"description": "Optional. Path to a pip config file (e.g., \"/etc/pip.conf\" or \"pip.ini\") for controlling\npackage installation (custom indices, credentials, etc.).\n\nOnly relevant if Python dependencies are installed via pip. If omitted, default pip settings are used.\n" "description": "Optional. Path to a pip config file (e.g., \"/etc/pip.conf\" or \"pip.ini\") for controlling\npackage installation (custom indices, timeouts, etc.). The file is copied into the\ngenerated image, so it must not contain credentials or other secrets.\n\nOnly relevant if Python dependencies are installed via pip. If omitted, default pip settings are used.\n"
}, },
"_INTERNAL_docker_tag": { "_INTERNAL_docker_tag": {
"anyOf": [ "anyOf": [
+2 -2
View File
@@ -28,7 +28,7 @@
"type": "null" "type": "null"
} }
], ],
"description": "Optional. Path to a pip config file (e.g., \"/etc/pip.conf\" or \"pip.ini\") for controlling\npackage installation (custom indices, credentials, etc.).\n\nOnly relevant if Python dependencies are installed via pip. If omitted, default pip settings are used.\n" "description": "Optional. Path to a pip config file (e.g., \"/etc/pip.conf\" or \"pip.ini\") for controlling\npackage installation (custom indices, timeouts, etc.). The file is copied into the\ngenerated image, so it must not contain credentials or other secrets.\n\nOnly relevant if Python dependencies are installed via pip. If omitted, default pip settings are used.\n"
}, },
"_INTERNAL_docker_tag": { "_INTERNAL_docker_tag": {
"anyOf": [ "anyOf": [
@@ -270,7 +270,7 @@
"type": "null" "type": "null"
} }
], ],
"description": "Optional. Path to a pip config file (e.g., \"/etc/pip.conf\" or \"pip.ini\") for controlling\npackage installation (custom indices, credentials, etc.).\n\nOnly relevant if Python dependencies are installed via pip. If omitted, default pip settings are used.\n" "description": "Optional. Path to a pip config file (e.g., \"/etc/pip.conf\" or \"pip.ini\") for controlling\npackage installation (custom indices, timeouts, etc.). The file is copied into the\ngenerated image, so it must not contain credentials or other secrets.\n\nOnly relevant if Python dependencies are installed via pip. If omitted, default pip settings are used.\n"
}, },
"_INTERNAL_docker_tag": { "_INTERNAL_docker_tag": {
"anyOf": [ "anyOf": [
+27 -31
View File
@@ -382,20 +382,18 @@ def test_deploy_list_command(monkeypatch) -> None:
def list_deployments(self, name_contains: str = ""): def list_deployments(self, name_contains: str = ""):
captured["name_contains"] = name_contains captured["name_contains"] = name_contains
return { return [
"resources": [ {
{ "id": "dep-123",
"id": "dep-123", "name": "alpha",
"name": "alpha", "source_config": {"custom_url": "https://alpha.example.com"},
"source_config": {"custom_url": "https://alpha.example.com"}, },
}, {
{ "id": "dep-456",
"id": "dep-456", "name": "beta",
"name": "beta", "source_config": {"custom_url": "https://beta.example.com"},
"source_config": {"custom_url": "https://beta.example.com"}, },
}, ]
]
}
monkeypatch.setattr(deploy_module, "HostBackendClient", FakeClient) monkeypatch.setattr(deploy_module, "HostBackendClient", FakeClient)
@@ -435,7 +433,7 @@ def test_deploy_list_command_no_results(monkeypatch) -> None:
pass pass
def list_deployments(self, name_contains: str = ""): def list_deployments(self, name_contains: str = ""):
return {"resources": []} return []
monkeypatch.setattr(deploy_module, "HostBackendClient", FakeClient) monkeypatch.setattr(deploy_module, "HostBackendClient", FakeClient)
@@ -468,20 +466,18 @@ def test_deploy_revisions_list_command(monkeypatch) -> None:
def list_revisions(self, deployment_id: str, limit: int = 1): def list_revisions(self, deployment_id: str, limit: int = 1):
captured["deployment_id"] = deployment_id captured["deployment_id"] = deployment_id
captured["limit"] = str(limit) captured["limit"] = str(limit)
return { return [
"resources": [ {
{ "id": "rev-123",
"id": "rev-123", "status": "CREATING",
"status": "CREATING", "created_at": "2023-11-07T05:31:56Z",
"created_at": "2023-11-07T05:31:56Z", },
}, {
{ "id": "rev-456",
"id": "rev-456", "status": "DEPLOYED",
"status": "DEPLOYED", "created_at": "2023-11-08T10:00:00Z",
"created_at": "2023-11-08T10:00:00Z", },
}, ]
]
}
monkeypatch.setattr(deploy_module, "HostBackendClient", FakeClient) monkeypatch.setattr(deploy_module, "HostBackendClient", FakeClient)
@@ -522,7 +518,7 @@ def test_deploy_revisions_list_command_no_results(monkeypatch) -> None:
pass pass
def list_revisions(self, deployment_id: str, limit: int = 1): def list_revisions(self, deployment_id: str, limit: int = 1):
return {"resources": []} return []
monkeypatch.setattr(deploy_module, "HostBackendClient", FakeClient) monkeypatch.setattr(deploy_module, "HostBackendClient", FakeClient)
@@ -555,7 +551,7 @@ def test_deploy_revisions_list_command_with_explicit_limit(monkeypatch) -> None:
def list_revisions(self, deployment_id: str, limit: int = 1): def list_revisions(self, deployment_id: str, limit: int = 1):
captured["deployment_id"] = deployment_id captured["deployment_id"] = deployment_id
captured["limit"] = str(limit) captured["limit"] = str(limit)
return {"resources": []} return []
monkeypatch.setattr(deploy_module, "HostBackendClient", FakeClient) monkeypatch.setattr(deploy_module, "HostBackendClient", FakeClient)
@@ -1,5 +1,6 @@
import asyncio import asyncio
import json import json
import uuid
from collections.abc import Callable, Iterator from collections.abc import Callable, Iterator
from contextlib import contextmanager from contextlib import contextmanager
from dataclasses import dataclass, field from dataclasses import dataclass, field
@@ -17,6 +18,7 @@ from langgraph_cli.host_backend import HostBackendClient
from langgraph_cli.image_reference import ImageReference from langgraph_cli.image_reference import ImageReference
CONTROL_PLANE_URL = "https://control-plane.example.com" CONTROL_PLANE_URL = "https://control-plane.example.com"
CLOUD_CONTROL_PLANE_URL = "https://api.host.langchain.com"
REGISTRY_URL = "https://registry.example.com/team" REGISTRY_URL = "https://registry.example.com/team"
PUSH_TOKEN = "push-token" PUSH_TOKEN = "push-token"
PUSHED_IMAGE = "registry.example.com/team/my-app:latest" PUSHED_IMAGE = "registry.example.com/team/my-app:latest"
@@ -24,10 +26,25 @@ PUSHED_DIGEST = "registry.example.com/team/my-app@sha256:abc123"
PUSH_REPOSITORY = "registry.example.com/team/agent" PUSH_REPOSITORY = "registry.example.com/team/agent"
EXTERNAL_IMAGE = f"{PUSH_REPOSITORY}:latest" EXTERNAL_IMAGE = f"{PUSH_REPOSITORY}:latest"
EXTERNAL_DIGEST = f"{PUSH_REPOSITORY}@sha256:abc123" EXTERNAL_DIGEST = f"{PUSH_REPOSITORY}@sha256:abc123"
LISTENER_REQUIRED = ( LISTENER_ID = "11111111-1111-4111-8111-111111111111"
"Source configuration error: 'source_config.listener_id' is required for " OTHER_LISTENER_ID = "22222222-2222-4222-8222-222222222222"
"workspace with available listener IDs: ['listener-1']" PAGE_TWO_LISTENER_ID = "33333333-3333-4333-8333-333333333333"
) UNKNOWN_LISTENER_ID = "99999999-9999-4999-8999-999999999999"
LISTENER = {
"id": LISTENER_ID,
"compute_id": "prod-cluster",
"compute_config": {"k8s_namespaces": ["agents"]},
}
OTHER_LISTENER = {
"id": OTHER_LISTENER_ID,
"compute_id": "other-cluster",
"compute_config": {"k8s_namespaces": ["agents"]},
}
TWO_NAMESPACE_LISTENER = {
"id": LISTENER_ID,
"compute_id": "prod-cluster",
"compute_config": {"k8s_namespaces": ["agents", "agents-staging"]},
}
CREATED_ID = "dep-created" CREATED_ID = "dep-created"
TRACKED_PACKAGES = ["langgraph:1.0.0"] TRACKED_PACKAGES = ["langgraph:1.0.0"]
SIGNED_UPLOAD_URL = "https://storage.example.com/signed" SIGNED_UPLOAD_URL = "https://storage.example.com/signed"
@@ -38,7 +55,12 @@ DIGESTS_FORMAT = "{{json .RepoDigests}}"
NOT_A_CLI_DEPLOYMENT = ( NOT_A_CLI_DEPLOYMENT = (
"push token is only available for 'internal_docker' source deployments" "push token is only available for 'internal_docker' source deployments"
) )
LISTENER_REQUIRED = (
"Source configuration error: 'source_config.listener_id' is required "
f"for workspace with available listener IDs: ['{LISTENER_ID}']"
)
LIST_DEPLOYMENTS = "GET /v2/deployments" LIST_DEPLOYMENTS = "GET /v2/deployments"
LIST_LISTENERS = "GET /v2/listeners"
CREATE_DEPLOYMENT = "POST /v2/deployments" CREATE_DEPLOYMENT = "POST /v2/deployments"
@@ -58,12 +80,22 @@ def _get(deployment_id: str) -> str:
return f"GET /v2/deployments/{deployment_id}" return f"GET /v2/deployments/{deployment_id}"
def _looks_like_a_uuid(value: str) -> bool:
try:
uuid.UUID(value)
except ValueError:
return False
return True
@dataclass @dataclass
class ControlPlaneDouble: class ControlPlaneDouble:
timeline: list[str] timeline: list[str]
existing_deployments: list[dict] = field(default_factory=list) existing_deployments: list[dict] = field(default_factory=list)
push_token_status: int = 200 push_token_status: int = 200
create_error: str | None = None create_error: str | None = None
listeners: list[dict] = field(default_factory=list)
listeners_by_id: dict[str, dict] = field(default_factory=dict)
bodies: dict[str, dict] = field(default_factory=dict) bodies: dict[str, dict] = field(default_factory=dict)
def handle(self, request: httpx.Request) -> httpx.Response: def handle(self, request: httpx.Request) -> httpx.Response:
@@ -71,14 +103,45 @@ class ControlPlaneDouble:
self.timeline.append(route) self.timeline.append(route)
if request.content: if request.content:
self.bodies[route] = json.loads(request.content) self.bodies[route] = json.loads(request.content)
return self._respond(request.method, request.url.path) return self._respond(request)
def _respond(self, method: str, path: str) -> httpx.Response: def _respond(self, request: httpx.Request) -> httpx.Response:
method, path = request.method, request.url.path
if (method, path) == ("GET", "/v2/listeners"):
return httpx.Response(200, json={"resources": self.listeners})
if method == "GET" and path.startswith("/v2/listeners/"):
listener_id = path.rsplit("/", 1)[-1]
if not _looks_like_a_uuid(listener_id):
return httpx.Response(
422,
json={
"detail": [
{"type": "uuid_parsing", "loc": ["path", "listener_id"]}
]
},
)
known = {listener["id"]: listener for listener in self.listeners}
known.update(self.listeners_by_id)
if listener_id not in known:
return httpx.Response(
404, json={"detail": f"Listener ID {listener_id} not found."}
)
return httpx.Response(200, json=known[listener_id])
if (method, path) == ("GET", "/v2/deployments"): if (method, path) == ("GET", "/v2/deployments"):
return httpx.Response(200, json={"resources": self.existing_deployments}) name = request.url.params.get("name")
return httpx.Response(
200,
json={
"resources": [
deployment
for deployment in self.existing_deployments
if name is None or deployment.get("name") == name
]
},
)
if (method, path) == ("POST", "/v2/deployments"): if (method, path) == ("POST", "/v2/deployments"):
if self.create_error is not None: if self.create_error is not None:
return httpx.Response(400, text=self.create_error) return httpx.Response(400, json={"detail": self.create_error})
return httpx.Response(201, json={"id": CREATED_ID, "tenant_id": "tenant-1"}) return httpx.Response(201, json={"id": CREATED_ID, "tenant_id": "tenant-1"})
if path.endswith("/push-token"): if path.endswith("/push-token"):
if self.push_token_status != 200: if self.push_token_status != 200:
@@ -199,7 +262,7 @@ class DeployProject:
timeline: list[str] timeline: list[str]
uploads: list[tuple[str, str, int]] uploads: list[tuple[str, str, int]]
def run(self, *args: str) -> Result: def run(self, *args: str, host_url: str = CONTROL_PLANE_URL) -> Result:
return CliRunner().invoke( return CliRunner().invoke(
cli, cli,
[ [
@@ -207,7 +270,7 @@ class DeployProject:
"--api-key", "--api-key",
"test-key", "test-key",
"--host-url", "--host-url",
CONTROL_PLANE_URL, host_url,
"--name", "--name",
"my-app", "my-app",
"--no-input", "--no-input",
@@ -611,18 +674,6 @@ def test_push_to_rejects_a_non_external_deployment_before_any_docker_work(
assert deploy_project.docker.verbs() == [] assert deploy_project.docker.verbs() == []
def test_push_to_explains_the_listener_requirement_of_hybrid_workspaces(
deploy_project: DeployProject,
) -> None:
deploy_project.control_plane.create_error = LISTENER_REQUIRED
result = deploy_project.run("--push-to", PUSH_REPOSITORY)
assert result.exit_code != 0
assert "listener" in result.output
assert "--deployment-id" in result.output
def test_push_to_with_deployment_id_fetches_the_deployment_once( def test_push_to_with_deployment_id_fetches_the_deployment_once(
deploy_project: DeployProject, deploy_project: DeployProject,
) -> None: ) -> None:
@@ -652,3 +703,414 @@ def test_invalid_tag_fails_before_any_control_plane_call(
assert result.exit_code != 0 assert result.exit_code != 0
assert "Image tag may only contain" in result.output assert "Image tag may only contain" in result.output
assert deploy_project.timeline == [] assert deploy_project.timeline == []
def test_push_to_places_a_new_deployment_on_the_only_listener(
deploy_project: DeployProject,
) -> None:
deploy_project.control_plane.listeners = [LISTENER]
result = deploy_project.run(
"--push-to", PUSH_REPOSITORY, host_url=CLOUD_CONTROL_PLANE_URL
)
assert result.exit_code == 0, result.output
assert deploy_project.timeline == [
LIST_DEPLOYMENTS,
LIST_LISTENERS,
"docker build",
"docker push",
"docker inspect-digest",
CREATE_DEPLOYMENT,
]
assert deploy_project.control_plane.bodies[CREATE_DEPLOYMENT]["source_config"] == {
"resource_spec": {},
"listener_id": LISTENER_ID,
"listener_config": {"k8s_namespace": "agents"},
}
assert f"Deploying through listener {LISTENER_ID} in namespace agents" in (
result.output
)
def test_push_to_places_a_new_deployment_on_the_chosen_listener(
deploy_project: DeployProject,
) -> None:
deploy_project.control_plane.listeners = [LISTENER, OTHER_LISTENER]
result = deploy_project.run(
"--push-to",
PUSH_REPOSITORY,
"--listener-id",
OTHER_LISTENER_ID,
"--k8s-namespace",
"agents",
host_url=CLOUD_CONTROL_PLANE_URL,
)
assert result.exit_code == 0, result.output
assert deploy_project.control_plane.bodies[CREATE_DEPLOYMENT]["source_config"] == {
"resource_spec": {},
"listener_id": OTHER_LISTENER_ID,
"listener_config": {"k8s_namespace": "agents"},
}
@pytest.mark.parametrize(
("listeners", "args", "message"),
[
pytest.param(
[LISTENER, OTHER_LISTENER], (), "--listener-id", id="two_listeners"
),
pytest.param(
[TWO_NAMESPACE_LISTENER], (), "--k8s-namespace", id="two_namespaces"
),
pytest.param(
[LISTENER],
("--k8s-namespace", "nope"),
"does not serve namespace",
id="unknown_namespace",
),
],
)
def test_push_to_refuses_an_unresolved_placement_before_any_docker_work(
deploy_project: DeployProject, listeners, args, message
) -> None:
deploy_project.control_plane.listeners = listeners
result = deploy_project.run(
"--push-to", PUSH_REPOSITORY, *args, host_url=CLOUD_CONTROL_PLANE_URL
)
assert result.exit_code != 0
assert message in result.output
assert deploy_project.docker.verbs() == []
assert CREATE_DEPLOYMENT not in deploy_project.timeline
def test_self_hosted_control_plane_keeps_its_default_placement(
deploy_project: DeployProject,
) -> None:
deploy_project.control_plane.listeners = [LISTENER]
result = deploy_project.run("--push-to", PUSH_REPOSITORY)
assert result.exit_code == 0, result.output
assert deploy_project.control_plane.bodies[CREATE_DEPLOYMENT]["source_config"] == {
"resource_spec": {}
}
def test_self_hosted_control_plane_places_when_asked(
deploy_project: DeployProject,
) -> None:
deploy_project.control_plane.listeners = [LISTENER]
result = deploy_project.run(
"--push-to", PUSH_REPOSITORY, "--listener-id", LISTENER_ID
)
assert result.exit_code == 0, result.output
assert deploy_project.control_plane.bodies[CREATE_DEPLOYMENT]["source_config"] == {
"resource_spec": {},
"listener_id": LISTENER_ID,
"listener_config": {"k8s_namespace": "agents"},
}
def test_updating_a_deployment_never_looks_up_listeners(
deploy_project: DeployProject,
) -> None:
deploy_project.control_plane.listeners = [LISTENER]
deploy_project.control_plane.existing_deployments = [
{"id": "dep-ext", "name": "my-app", "source": "external_docker"}
]
result = deploy_project.run(
"--push-to", PUSH_REPOSITORY, host_url=CLOUD_CONTROL_PLANE_URL
)
assert result.exit_code == 0, result.output
assert LIST_LISTENERS not in deploy_project.timeline
def test_listener_flags_are_refused_for_a_deployment_id_without_any_call(
deploy_project: DeployProject,
) -> None:
result = deploy_project.run(
"--push-to",
PUSH_REPOSITORY,
"--deployment-id",
"dep-ext",
"--k8s-namespace",
"agents",
host_url=CLOUD_CONTROL_PLANE_URL,
)
assert result.exit_code != 0
assert "fixed when a deployment is created" in result.output
assert deploy_project.timeline == []
def test_listener_flags_are_refused_on_an_existing_deployment(
deploy_project: DeployProject,
) -> None:
deploy_project.control_plane.listeners = [LISTENER]
deploy_project.control_plane.existing_deployments = [
{"id": "dep-ext", "name": "my-app", "source": "external_docker"}
]
result = deploy_project.run(
"--push-to",
PUSH_REPOSITORY,
"--listener-id",
LISTENER_ID,
host_url=CLOUD_CONTROL_PLANE_URL,
)
assert result.exit_code != 0
assert "fixed when a deployment is created" in result.output
assert deploy_project.docker.verbs() == []
def test_a_deployment_without_a_listener_announces_nothing(
deploy_project: DeployProject,
) -> None:
result = deploy_project.run("--push-to", PUSH_REPOSITORY)
assert result.exit_code == 0, result.output
assert "listener" not in result.output
def test_a_self_hosted_create_without_flags_never_looks_up_listeners(
deploy_project: DeployProject,
) -> None:
deploy_project.control_plane.listeners = [LISTENER]
result = deploy_project.run("--push-to", PUSH_REPOSITORY)
assert result.exit_code == 0, result.output
assert LIST_LISTENERS not in deploy_project.timeline
def test_a_control_plane_that_demands_a_listener_names_the_flags(
deploy_project: DeployProject,
) -> None:
deploy_project.control_plane.create_error = LISTENER_REQUIRED
result = deploy_project.run("--push-to", PUSH_REPOSITORY)
assert result.exit_code != 0
assert "--listener-id" in result.output
assert "--k8s-namespace" in result.output
assert LISTENER_ID in result.output
assert "{" not in result.output
assert "POST /v2/deployments failed" not in result.output
def test_listener_flags_without_push_to_make_no_call_at_all(
deploy_project: DeployProject,
) -> None:
result = deploy_project.run("--listener-id", LISTENER_ID)
assert result.exit_code != 0
assert "--push-to" in result.output
assert deploy_project.timeline == []
def test_a_truncated_listener_page_says_so(deploy_project: DeployProject) -> None:
deploy_project.control_plane.listeners = [
{
"id": str(uuid.UUID(int=index)),
"compute_id": "cluster",
"compute_config": {"k8s_namespaces": ["agents"]},
}
for index in range(100)
]
result = deploy_project.run(
"--push-to", PUSH_REPOSITORY, host_url=CLOUD_CONTROL_PLANE_URL
)
assert result.exit_code != 0
assert "first 100" in result.output
def test_a_managed_build_in_a_listener_workspace_points_at_push_to(
deploy_project: DeployProject,
) -> None:
deploy_project.control_plane.create_error = LISTENER_REQUIRED
result = deploy_project.run("--no-remote")
assert result.exit_code != 0
assert "--push-to" in result.output
assert deploy_project.docker.verbs() == []
@pytest.mark.parametrize(
"args",
[
pytest.param(("--no-remote",), id="managed_build"),
pytest.param(("--push-to", PUSH_REPOSITORY), id="push_to"),
],
)
def test_a_listener_requirement_links_the_listener_docs(
deploy_project: DeployProject, args: tuple[str, ...]
) -> None:
deploy_project.control_plane.create_error = LISTENER_REQUIRED
result = deploy_project.run(*args)
assert result.exit_code != 0
assert "https://docs.langchain.com/langsmith/control-plane#listeners" in (
result.output
)
def test_a_managed_control_plane_without_listeners_creates_as_before(
deploy_project: DeployProject,
) -> None:
result = deploy_project.run(
"--push-to", PUSH_REPOSITORY, host_url=CLOUD_CONTROL_PLANE_URL
)
assert result.exit_code == 0, result.output
assert deploy_project.control_plane.bodies[CREATE_DEPLOYMENT]["source_config"] == {
"resource_spec": {}
}
assert deploy_project.timeline.count(LIST_LISTENERS) == 1
def test_a_listener_without_an_id_is_reported_rather_than_ignored(
deploy_project: DeployProject,
) -> None:
deploy_project.control_plane.listeners = [
{"compute_id": "broken", "compute_config": {"k8s_namespaces": ["agents"]}},
LISTENER,
]
result = deploy_project.run(
"--push-to", PUSH_REPOSITORY, host_url=CLOUD_CONTROL_PLANE_URL
)
assert result.exit_code != 0
assert "without an id" in result.output
assert deploy_project.docker.verbs() == []
def _listener_route(listener_id: str) -> str:
return f"GET /v2/listeners/{listener_id}"
def test_an_explicit_listener_is_fetched_by_id_not_searched(
deploy_project: DeployProject,
) -> None:
deploy_project.control_plane.listeners = [LISTENER, OTHER_LISTENER]
result = deploy_project.run(
"--push-to",
PUSH_REPOSITORY,
"--listener-id",
OTHER_LISTENER_ID,
host_url=CLOUD_CONTROL_PLANE_URL,
)
assert result.exit_code == 0, result.output
assert _listener_route(OTHER_LISTENER_ID) in deploy_project.timeline
assert LIST_LISTENERS not in deploy_project.timeline
assert deploy_project.control_plane.bodies[CREATE_DEPLOYMENT]["source_config"] == {
"resource_spec": {},
"listener_id": OTHER_LISTENER_ID,
"listener_config": {"k8s_namespace": "agents"},
}
def test_an_explicit_listener_beyond_the_first_page_still_works(
deploy_project: DeployProject,
) -> None:
deploy_project.control_plane.listeners = [
{
"id": str(uuid.UUID(int=index)),
"compute_id": "cluster",
"compute_config": {"k8s_namespaces": ["agents"]},
}
for index in range(100)
]
deploy_project.control_plane.listeners_by_id = {
PAGE_TWO_LISTENER_ID: {
"id": PAGE_TWO_LISTENER_ID,
"compute_id": "far-cluster",
"compute_config": {"k8s_namespaces": ["agents"]},
}
}
result = deploy_project.run(
"--push-to",
PUSH_REPOSITORY,
"--listener-id",
PAGE_TWO_LISTENER_ID,
host_url=CLOUD_CONTROL_PLANE_URL,
)
assert result.exit_code == 0, result.output
assert deploy_project.control_plane.bodies[CREATE_DEPLOYMENT]["source_config"] == {
"resource_spec": {},
"listener_id": PAGE_TWO_LISTENER_ID,
"listener_config": {"k8s_namespace": "agents"},
}
def test_an_unknown_listener_names_the_ones_that_exist(
deploy_project: DeployProject,
) -> None:
deploy_project.control_plane.listeners = [LISTENER]
result = deploy_project.run(
"--push-to",
PUSH_REPOSITORY,
"--listener-id",
UNKNOWN_LISTENER_ID,
host_url=CLOUD_CONTROL_PLANE_URL,
)
assert result.exit_code != 0
assert "was not found" in result.output
assert LISTENER_ID in result.output
assert "prod-cluster" in result.output
assert deploy_project.docker.verbs() == []
def test_an_explicit_listener_in_a_workspace_without_any_is_refused(
deploy_project: DeployProject,
) -> None:
result = deploy_project.run(
"--push-to",
PUSH_REPOSITORY,
"--listener-id",
LISTENER_ID,
host_url=CLOUD_CONTROL_PLANE_URL,
)
assert result.exit_code != 0
assert "no listeners" in result.output
assert deploy_project.docker.verbs() == []
def test_a_listener_id_that_is_not_an_identifier_still_names_the_real_ones(
deploy_project: DeployProject,
) -> None:
deploy_project.control_plane.listeners = [LISTENER]
result = deploy_project.run(
"--push-to",
PUSH_REPOSITORY,
"--listener-id",
"not-a-listener",
host_url=CLOUD_CONTROL_PLANE_URL,
)
assert result.exit_code != 0
assert "was not found" in result.output
assert LISTENER_ID in result.output
assert "uuid_parsing" not in result.output
+237
View File
@@ -255,6 +255,243 @@ def test_validate_config():
) )
@pytest.mark.parametrize(
"dependency",
[
"git+https://user:secret-token@github.com/org/private.git@main",
"private-package @ git+http://token@github.com/org/private.git",
"git+HTTPS://user%40example.com:secret%2Ftoken@github.com/org/private.git",
"git+https://${GIT_TOKEN}@github.com/org/private.git",
],
)
def test_validate_config_rejects_git_http_url_userinfo(dependency: str):
with pytest.raises(click.UsageError) as exc_info:
validate_config(
{
"python_version": "3.11",
"dependencies": [dependency],
"graphs": {"agent": "./agent.py:graph"},
}
)
message = str(exc_info.value)
assert "must not contain credentials or other URL userinfo" in message
assert "secret-token" not in message
assert "secret%2Ftoken" not in message
def test_validate_config_file_reports_source_for_git_http_url_userinfo(
tmp_path: pathlib.Path,
):
config_path = tmp_path / "langgraph.json"
config_path.write_text(
json.dumps(
{
"python_version": "3.11",
"dependencies": ["git+https://secret-token@github.com/org/private.git"],
"graphs": {"agent": "./agent.py:graph"},
}
)
)
with pytest.raises(click.UsageError) as exc_info:
validate_config_file(config_path)
message = str(exc_info.value)
assert "secret-token" not in message
assert f"Found in: {config_path.resolve()}" in message
@pytest.mark.parametrize(
"manifest", ["package.json", "package-lock.json", "yarn.lock", "pnpm-lock.yaml"]
)
def test_config_to_docker_rejects_git_http_url_userinfo_in_node_files(
tmp_path: pathlib.Path, manifest: str
):
config_path = tmp_path / "langgraph.json"
config_path.write_text("{}\n")
(tmp_path / "agent.js").write_text("export const graph = {};\n")
(tmp_path / "package.json").write_text('{"name":"agent"}\n')
(tmp_path / manifest).write_text(
'"priv": "git+https://user:secret-token@github.com/org/private.git"\n'
)
config = validate_config(
{
"node_version": "20",
"graphs": {"agent": "./agent.js:graph"},
}
)
with pytest.raises(click.UsageError) as exc_info:
config_to_docker(
config_path,
config,
base_image="langchain/langgraphjs-api",
)
message = str(exc_info.value)
assert "must not contain credentials or other URL userinfo" in message
assert "secret-token" not in message
assert f"Found in: {(tmp_path / manifest).resolve()}" in message
def test_config_to_docker_allows_node_git_urls_without_http_userinfo(
tmp_path: pathlib.Path,
):
config_path = tmp_path / "langgraph.json"
config_path.write_text("{}\n")
(tmp_path / "agent.js").write_text("export const graph = {};\n")
(tmp_path / "package.json").write_text(
'{"dependencies":{"public":"git+https://github.com/org/public.git"}}\n'
)
config = validate_config(
{
"node_version": "20",
"graphs": {"agent": "./agent.js:graph"},
}
)
docker, _ = config_to_docker(
config_path,
config,
base_image="langchain/langgraphjs-api",
)
assert f"ADD . /deps/{tmp_path.name}" in docker
def test_config_to_docker_rejects_git_http_url_userinfo_in_node_workspace(
tmp_path: pathlib.Path,
):
config_root = tmp_path / "apps" / "agent"
config_root.mkdir(parents=True)
config_path = config_root / "langgraph.json"
config_path.write_text("{}\n")
(config_root / "agent.js").write_text("export const graph = {};\n")
(config_root / "package.json").write_text(
'{"dependencies":{"priv":"git+https://secret-token@github.com/org/private.git"}}\n'
)
(tmp_path / "package.json").write_text('{"name":"workspace"}\n')
config = validate_config(
{
"node_version": "20",
"graphs": {"agent": "./agent.js:graph"},
}
)
with pytest.raises(click.UsageError) as exc_info:
config_to_docker(
config_path,
config,
base_image="langchain/langgraphjs-api",
build_context=str(tmp_path),
)
message = str(exc_info.value)
assert "secret-token" not in message
assert f"Found in: {(config_root / 'package.json').resolve()}" in message
@pytest.mark.parametrize(
"dependency",
[
"git+https://github.com/org/public.git@main",
"private-package @ git+https://github.com/org/private.git@main",
"git+ssh://git@github.com/org/private.git@main",
],
)
def test_validate_config_allows_git_urls_without_http_userinfo(dependency: str):
config = validate_config(
{
"python_version": "3.11",
"dependencies": [dependency],
"graphs": {"agent": "./agent.py:graph"},
}
)
assert config["dependencies"] == [dependency]
def test_config_to_docker_rejects_git_http_url_userinfo_in_requirements(
tmp_path: pathlib.Path,
):
config_path = tmp_path / "langgraph.json"
config_path.write_text("{}\n")
(tmp_path / "agent.py").write_text("graph = object()\n")
(tmp_path / "requirements.txt").write_text(
"private @ git+https://secret-token@github.com/org/private.git\n"
)
config = validate_config(
{
"python_version": "3.11",
"dependencies": ["."],
"graphs": {"agent": "./agent.py:graph"},
}
)
with pytest.raises(click.UsageError) as exc_info:
config_to_docker(
config_path,
config,
base_image="langchain/langgraph-api:0.2.47",
)
message = str(exc_info.value)
assert "must not contain credentials or other URL userinfo" in message
assert "secret-token" not in message
assert f"Found in: {(tmp_path / 'requirements.txt').resolve()}" in message
@pytest.mark.parametrize("manifest", ["pyproject.toml", "uv.lock"])
def test_config_to_docker_rejects_git_http_url_userinfo_in_uv_files(
tmp_path: pathlib.Path, manifest: str
):
config_path = tmp_path / "langgraph.json"
config_path.write_text("{}\n")
(tmp_path / "src").mkdir()
(tmp_path / "src" / "agent.py").write_text("graph = object()\n")
pyproject = textwrap.dedent(
"""
[project]
name = "agent"
version = "0.1.0"
dependencies = ["private"]
[tool.uv.sources]
private = { git = "https://github.com/org/private.git" }
"""
).strip()
uv_lock = "# uv lock file\n"
if manifest == "pyproject.toml":
pyproject = pyproject.replace(
"https://github.com", "https://secret-token@github.com"
)
else:
uv_lock += (
'source = { git = "https://secret-token@github.com/org/private.git" }\n'
)
(tmp_path / "pyproject.toml").write_text(pyproject + "\n")
(tmp_path / "uv.lock").write_text(uv_lock)
config = validate_config(
{
"python_version": "3.11",
"graphs": {"agent": "./src/agent.py:graph"},
"source": {"kind": "uv"},
}
)
with pytest.raises(click.UsageError) as exc_info:
config_to_docker(
config_path,
config,
base_image="langchain/langgraph-api:0.2.47",
)
message = str(exc_info.value)
assert "must not contain credentials or other URL userinfo" in message
assert "secret-token" not in message
def test_validate_config_image_distro(): def test_validate_config_image_distro():
"""Test validation of image_distro field.""" """Test validation of image_distro field."""
# Valid image_distro values should work # Valid image_distro values should work
@@ -58,7 +58,7 @@ AGENT_ARGS = [
"deploy", "deploy",
"--agent-id", "--agent-id",
"customer-support", "customer-support",
"--environment", "--agent-environment",
"staging", "staging",
"--remote", "--remote",
"--no-wait", "--no-wait",
@@ -72,9 +72,9 @@ def test_agent_create(deployment_api, tmp_path, monkeypatch):
result = CliRunner().invoke(cli, AGENT_ARGS) result = CliRunner().invoke(cli, AGENT_ARGS)
assert result.exit_code == 0, result.output assert result.exit_code == 0, result.output
assert dict(requests[0].url.params) == { assert dict(requests[0].url.params) == {
"name_contains": "",
"agent_id": "customer-support", "agent_id": "customer-support",
"agent_environment": "staging", "agent_environment": "staging",
"limit": "100",
} }
payload = json.loads(requests[1].content) payload = json.loads(requests[1].content)
assert payload["agent"] == { assert payload["agent"] == {
@@ -103,3 +103,17 @@ def test_agent_rejects_explicit_name(deployment_api, monkeypatch):
assert result.exit_code == 2 assert result.exit_code == 2
assert "cannot be combined" in result.output assert "cannot be combined" in result.output
assert not requests assert not requests
def test_agent_lookup_refuses_a_control_plane_that_ignores_the_filter(deployment_api):
state, requests, _ = deployment_api
state["resources"] = [
{"id": "someone-elses", "is_preview": False},
{"id": "another", "is_preview": False},
]
result = CliRunner().invoke(cli, AGENT_ARGS)
assert result.exit_code != 0
assert "does not filter deployments by agent" in result.output
assert len(requests) == 1
@@ -13,10 +13,17 @@ import pytest
import langgraph_cli.deploy as deploy_mod import langgraph_cli.deploy as deploy_mod
from langgraph_cli.deploy import ( from langgraph_cli.deploy import (
ById,
ByName,
CustomerRegistrySource, CustomerRegistrySource,
DockerBuildCommand, DockerBuildCommand,
ExistingDeployment,
Listener,
ManagedRegistrySource, ManagedRegistrySource,
OnListener,
RemoteBuildSource, RemoteBuildSource,
RequestedPlacement,
Unplaced,
_call_host_backend_with_optional_tenant, _call_host_backend_with_optional_tenant,
_create_host_backend_client, _create_host_backend_client,
_docker_config_for_token, _docker_config_for_token,
@@ -27,6 +34,7 @@ from langgraph_cli.deploy import (
_resolve_pushed_image_digest, _resolve_pushed_image_digest,
_select_source, _select_source,
_validate_prebuilt_image, _validate_prebuilt_image,
find_deployment_by_name,
normalize_image_tag, normalize_image_tag,
normalize_name, normalize_name,
) )
@@ -280,11 +288,13 @@ class TestCallHostBackendWithOptionalTenant:
return c return c
def test_success_passes_through(self): def test_success_passes_through(self):
client = self._make_client(lambda req: httpx.Response(200, json={"ok": True})) client = self._make_client(
lambda req: httpx.Response(200, json={"resources": [{"id": "dep-1"}]})
)
result = _call_host_backend_with_optional_tenant( result = _call_host_backend_with_optional_tenant(
client, lambda c: c.list_deployments() client, lambda c: c.list_deployments()
) )
assert result == {"ok": True} assert result == [{"id": "dep-1"}]
def test_403_not_enabled_gives_actionable_error(self): def test_403_not_enabled_gives_actionable_error(self):
detail = ( detail = (
@@ -607,6 +617,8 @@ class TestSelectSource:
"image_name": None, "image_name": None,
"tag": None, "tag": None,
"remote_build_flag": None, "remote_build_flag": None,
"placement": RequestedPlacement(),
"selector": ByName("my-app"),
} }
REPOSITORY = "registry.example.com/app" REPOSITORY = "registry.example.com/app"
@@ -617,7 +629,9 @@ class TestSelectSource:
{"push_to": REPOSITORY}, {"push_to": REPOSITORY},
True, True,
CustomerRegistrySource( CustomerRegistrySource(
ImageReference(REPOSITORY, "latest"), prebuilt_image=None reference=ImageReference(REPOSITORY, "latest"),
prebuilt_image=None,
requested_placement=RequestedPlacement(),
), ),
id="push_to_selects_the_external_source_with_the_default_tag", id="push_to_selects_the_external_source_with_the_default_tag",
), ),
@@ -625,7 +639,9 @@ class TestSelectSource:
{"push_to": f"{REPOSITORY}:v2"}, {"push_to": f"{REPOSITORY}:v2"},
True, True,
CustomerRegistrySource( CustomerRegistrySource(
ImageReference(REPOSITORY, "v2"), prebuilt_image=None reference=ImageReference(REPOSITORY, "v2"),
prebuilt_image=None,
requested_placement=RequestedPlacement(),
), ),
id="push_to_keeps_a_tag_given_in_the_reference", id="push_to_keeps_a_tag_given_in_the_reference",
), ),
@@ -633,7 +649,9 @@ class TestSelectSource:
{"push_to": REPOSITORY, "tag": "v3"}, {"push_to": REPOSITORY, "tag": "v3"},
True, True,
CustomerRegistrySource( CustomerRegistrySource(
ImageReference(REPOSITORY, "v3"), prebuilt_image=None reference=ImageReference(REPOSITORY, "v3"),
prebuilt_image=None,
requested_placement=RequestedPlacement(),
), ),
id="tag_flag_composes_with_push_to", id="tag_flag_composes_with_push_to",
), ),
@@ -641,10 +659,25 @@ class TestSelectSource:
{"push_to": REPOSITORY, "image": "app:dev"}, {"push_to": REPOSITORY, "image": "app:dev"},
False, False,
CustomerRegistrySource( CustomerRegistrySource(
ImageReference(REPOSITORY, "latest"), prebuilt_image="app:dev" reference=ImageReference(REPOSITORY, "latest"),
prebuilt_image="app:dev",
requested_placement=RequestedPlacement(),
), ),
id="prebuilt_image_is_retagged_for_push_to_without_docker_checks", id="prebuilt_image_is_retagged_for_push_to_without_docker_checks",
), ),
pytest.param(
{
"push_to": REPOSITORY,
"placement": RequestedPlacement("listener-1", "agents"),
},
True,
CustomerRegistrySource(
reference=ImageReference(REPOSITORY, "latest"),
prebuilt_image=None,
requested_placement=RequestedPlacement("listener-1", "agents"),
),
id="push_to_carries_the_requested_placement",
),
pytest.param( pytest.param(
{"remote_build_flag": True}, {"remote_build_flag": True},
True, True,
@@ -720,6 +753,16 @@ class TestSelectSource:
"--image cannot be combined with --remote builds.", "--image cannot be combined with --remote builds.",
id="image_with_remote", id="image_with_remote",
), ),
pytest.param(
{"placement": RequestedPlacement(listener_id="listener-1")},
"only apply when creating a deployment with --push-to",
id="listener_without_push_to",
),
pytest.param(
{"placement": RequestedPlacement(k8s_namespace="agents")},
"only apply when creating a deployment with --push-to",
id="namespace_without_push_to",
),
], ],
) )
def test_conflicting_flags_are_rejected(self, monkeypatch, flags, message): def test_conflicting_flags_are_rejected(self, monkeypatch, flags, message):
@@ -890,3 +933,289 @@ class TestResolvePushedImageDigest:
frame_locals = captured["coro"].cr_frame.f_locals frame_locals = captured["coro"].cr_frame.f_locals
assert "--config" not in frame_locals["args"] assert "--config" not in frame_locals["args"]
captured["coro"].close() captured["coro"].close()
class TestListener:
@pytest.mark.parametrize(
("resource", "expected"),
[
pytest.param(
{
"id": "listener-1",
"compute_id": "prod-cluster",
"compute_config": {"k8s_namespaces": ["agents", "agents-staging"]},
},
Listener("listener-1", "prod-cluster", ("agents", "agents-staging")),
id="reads_id_cluster_and_namespaces",
),
pytest.param(
{"id": "listener-1", "compute_id": "c", "compute_config": {}},
Listener("listener-1", "c", ()),
id="missing_namespaces",
),
pytest.param(
{"id": "listener-1", "compute_id": "c", "compute_config": None},
Listener("listener-1", "c", ()),
id="null_compute_config",
),
pytest.param(
{"id": "listener-1"},
Listener("listener-1", "", ()),
id="only_an_id",
),
],
)
def test_from_resource_reads_the_control_plane_shape(self, resource, expected):
assert Listener.from_resource(resource) == expected
ONE_NAMESPACE = Listener("listener-1", "prod-cluster", ("agents",))
TWO_NAMESPACES = Listener("listener-2", "multi-cluster", ("agents", "agents-staging"))
NO_NAMESPACE = Listener("listener-3", "broken-cluster", ())
class TestRequestedPlacement:
@pytest.mark.parametrize(
("request_", "listeners", "expected"),
[
pytest.param(
RequestedPlacement(), (), Unplaced(), id="no_listeners_no_request"
),
pytest.param(
RequestedPlacement(),
(ONE_NAMESPACE,),
OnListener("listener-1", "agents"),
id="uses_the_only_possible_answer",
),
pytest.param(
RequestedPlacement(k8s_namespace="agents-staging"),
(TWO_NAMESPACES,),
OnListener("listener-2", "agents-staging"),
id="namespace_alone_picks_the_only_listener",
),
],
)
def test_resolves_to_a_placement(self, request_, listeners, expected):
assert request_.among(listeners) == expected
@pytest.mark.parametrize(
("request_", "listeners", "message"),
[
pytest.param(
RequestedPlacement(listener_id="listener-1"),
(),
"no listeners",
id="workspace_has_no_listeners",
),
pytest.param(
RequestedPlacement(),
(ONE_NAMESPACE, TWO_NAMESPACES),
"--listener-id",
id="several_listeners_need_a_choice",
),
pytest.param(
RequestedPlacement(k8s_namespace="agents"),
(ONE_NAMESPACE, TWO_NAMESPACES),
"--listener-id",
id="namespace_alone_is_ambiguous_with_several_listeners",
),
pytest.param(
RequestedPlacement(k8s_namespace="agents"),
(),
"no listeners",
id="namespace_without_any_listener",
),
pytest.param(
RequestedPlacement(),
(TWO_NAMESPACES,),
"--k8s-namespace",
id="several_namespaces_need_a_choice",
),
],
)
def test_refuses_and_names_the_choices(self, request_, listeners, message):
with pytest.raises(click.UsageError, match=message):
request_.among(listeners)
def test_the_error_lists_every_listener_with_its_cluster_and_namespaces(self):
with pytest.raises(click.UsageError) as error:
RequestedPlacement().among((ONE_NAMESPACE, TWO_NAMESPACES))
assert "listener-1" in error.value.message
assert "prod-cluster" in error.value.message
assert "agents-staging" in error.value.message
@pytest.mark.parametrize(
("placement", "expected"),
[
pytest.param(Unplaced(), {}, id="unplaced_adds_nothing"),
pytest.param(
OnListener("listener-1", "agents"),
{
"listener_id": "listener-1",
"listener_config": {"k8s_namespace": "agents"},
},
id="placed_carries_listener_and_namespace",
),
],
)
def test_source_config_matches_the_control_plane_shape(self, placement, expected):
assert placement.source_config() == expected
def test_finding_a_deployment_by_name_narrows_the_search_for_every_server_version():
seen: dict = {}
def handler(req: httpx.Request) -> httpx.Response:
seen["params"] = dict(req.url.params)
return httpx.Response(
200,
json={"resources": [{"id": "dep-1", "name": "agent", "source": "github"}]},
)
client = HostBackendClient(
"https://api.example.com", "key", transport=httpx.MockTransport(handler)
)
found = find_deployment_by_name(client, "agent")
assert seen["params"] == {
"name": "agent",
"name_contains": "agent",
"limit": "100",
}
assert found == ExistingDeployment("dep-1", "github")
def test_a_server_that_ignores_the_exact_name_filter_never_matches_another_deployment():
client = HostBackendClient(
"https://api.example.com",
"key",
transport=httpx.MockTransport(
lambda req: httpx.Response(
200,
json={
"resources": [
{
"id": "dep-other",
"name": "another-teams-agent",
"source": "external_docker",
}
]
},
)
),
)
assert find_deployment_by_name(client, "brand-new-agent") is None
def test_a_full_page_without_a_match_refuses_to_claim_the_name_is_free():
page = [
{"id": f"dep-{index}", "name": f"other-agent-{index}"} for index in range(100)
]
client = HostBackendClient(
"https://api.example.com",
"key",
transport=httpx.MockTransport(
lambda req: httpx.Response(200, json={"resources": page})
),
)
with pytest.raises(click.ClickException, match="--deployment-id"):
find_deployment_by_name(client, "brand-new-agent")
def test_a_partial_page_without_a_match_means_the_name_is_free():
client = HostBackendClient(
"https://api.example.com",
"key",
transport=httpx.MockTransport(
lambda req: httpx.Response(
200, json={"resources": [{"id": "dep-1", "name": "other"}]}
)
),
)
assert find_deployment_by_name(client, "brand-new-agent") is None
@pytest.mark.parametrize(
"resource",
[
pytest.param({"compute_id": "c"}, id="no_id"),
pytest.param({"id": ""}, id="empty_id"),
],
)
def test_a_listener_without_an_id_is_refused(resource):
with pytest.raises(HostBackendError, match="without an id"):
Listener.from_resource(resource)
def test_a_deployment_id_with_listener_flags_is_refused_without_probing_docker(
monkeypatch,
):
def explode() -> tuple[bool, str | None]:
raise AssertionError("docker must not be probed for an argv-only conflict")
monkeypatch.setattr(deploy_mod, "can_build_locally", explode)
with pytest.raises(click.UsageError, match="--deployment-id"):
_select_source(
push_to="registry.example.com/app",
image=None,
image_name=None,
tag=None,
remote_build_flag=None,
placement=RequestedPlacement(listener_id="listener-1"),
selector=ById("dep-1"),
)
class TestPlacementOnAKnownListener:
@pytest.mark.parametrize(
("request_", "listener", "expected"),
[
pytest.param(
RequestedPlacement(listener_id="listener-1"),
ONE_NAMESPACE,
OnListener("listener-1", "agents"),
id="the_only_namespace_is_used",
),
pytest.param(
RequestedPlacement(listener_id="listener-2", k8s_namespace="agents"),
TWO_NAMESPACES,
OnListener("listener-2", "agents"),
id="the_chosen_namespace_is_used",
),
],
)
def test_places_on_the_listener(self, request_, listener, expected):
assert request_.on(listener) == expected
@pytest.mark.parametrize(
("request_", "listener", "message"),
[
pytest.param(
RequestedPlacement(listener_id="listener-2"),
TWO_NAMESPACES,
"--k8s-namespace",
id="several_namespaces_need_a_choice",
),
pytest.param(
RequestedPlacement(listener_id="listener-2", k8s_namespace="nope"),
TWO_NAMESPACES,
"does not serve namespace",
id="unknown_namespace",
),
pytest.param(
RequestedPlacement(listener_id="listener-3"),
NO_NAMESPACE,
"serves no namespaces",
id="listener_without_namespaces",
),
],
)
def test_refuses_and_names_the_namespaces(self, request_, listener, message):
with pytest.raises(click.UsageError, match=message):
request_.on(listener)
+142 -14
View File
@@ -79,19 +79,6 @@ def test_request_transport_error_raises():
c._request("GET", "/test") c._request("GET", "/test")
def test_list_deployments_sends_query_params():
def handler(req: httpx.Request) -> httpx.Response:
assert req.url.path == "/v2/deployments"
assert req.url.params["name_contains"] == "my app"
return httpx.Response(200, json={"ok": True})
c = HostBackendClient(
"https://api.example.com", "test-key", transport=httpx.MockTransport(handler)
)
result = c.list_deployments("my app")
assert result == {"ok": True}
def _capturing_client(captured: dict) -> HostBackendClient: def _capturing_client(captured: dict) -> HostBackendClient:
def handler(req: httpx.Request) -> httpx.Response: def handler(req: httpx.Request) -> httpx.Response:
captured["body"] = req.read() captured["body"] = req.read()
@@ -421,7 +408,7 @@ def test_injected_transport_receives_requests_under_the_prefixed_base_url():
transport=httpx.MockTransport(handler), transport=httpx.MockTransport(handler),
) )
assert c.list_revisions("dep-1", limit=2) == {"ok": True} assert c.list_revisions("dep-1", limit=2) == []
assert seen == { assert seen == {
"url": "https://smith.example.com/api-host/v2/deployments/dep-1/revisions?limit=2", "url": "https://smith.example.com/api-host/v2/deployments/dep-1/revisions?limit=2",
"api_key": "key", "api_key": "key",
@@ -546,3 +533,144 @@ def test_control_plane_endpoints_resolve(host_url, langsmith_endpoint, expected)
endpoints = ControlPlaneEndpoints.resolve(host_url, langsmith_endpoint) endpoints = ControlPlaneEndpoints.resolve(host_url, langsmith_endpoint)
assert (endpoints.control_plane_url, endpoints.dashboard_url) == expected assert (endpoints.control_plane_url, endpoints.dashboard_url) == expected
@pytest.mark.parametrize(
("payload", "expected"),
[
pytest.param(
{"resources": [{"id": "a"}, {"id": "b"}]},
[{"id": "a"}, {"id": "b"}],
id="list_returns_the_resources",
),
pytest.param({"resources": []}, [], id="empty_list"),
pytest.param({}, [], id="missing_key"),
pytest.param({"resources": None}, [], id="null_resources"),
pytest.param(
{"resources": ["nope", {"id": "a"}]}, [{"id": "a"}], id="skips_non_objects"
),
pytest.param([], [], id="unexpected_envelope"),
],
)
def test_list_endpoints_return_resource_objects(payload, expected):
def handler(req: httpx.Request) -> httpx.Response:
return httpx.Response(200, json=payload)
c = HostBackendClient(
"https://api.example.com", "key", transport=httpx.MockTransport(handler)
)
assert c.list_deployments() == expected
def test_list_listeners_asks_for_a_full_page():
seen: dict = {}
def handler(req: httpx.Request) -> httpx.Response:
seen["url"] = str(req.url)
return httpx.Response(200, json={"resources": [{"id": "listener-1"}]})
c = HostBackendClient(
"https://api.example.com", "key", transport=httpx.MockTransport(handler)
)
assert c.list_listeners() == [{"id": "listener-1"}]
assert seen["url"] == "https://api.example.com/v2/listeners?limit=100"
@pytest.mark.parametrize(
("control_plane_url", "expected"),
[
pytest.param("https://api.host.langchain.com", True, id="cloud"),
pytest.param("https://eu.api.host.langchain.com", True, id="cloud_region"),
pytest.param("https://dev.api.host.langchain.com", True, id="cloud_dev"),
pytest.param("https://smith.example.com/api-host", False, id="self_hosted"),
pytest.param(
"https://corp.example.com/langsmith/api-host",
False,
id="self_hosted_prefix",
),
pytest.param("http://localhost:8080/api-host", False, id="local"),
pytest.param(
"https://evil-api.host.langchain.com", False, id="lookalike_needs_a_dot"
),
],
)
def test_is_cloud_recognises_the_managed_control_plane(control_plane_url, expected):
endpoints = ControlPlaneEndpoints.from_control_plane_url(control_plane_url)
assert endpoints.is_cloud is expected
@pytest.mark.parametrize(
("call", "expected_params"),
[
pytest.param(
lambda c: c.list_deployments(name="agent"),
{"name": "agent"},
id="exact_name_filters_server_side",
),
pytest.param(
lambda c: c.list_deployments(name_contains="age"),
{"name_contains": "age"},
id="substring_search_keeps_its_own_parameter",
),
pytest.param(
lambda c: c.list_deployments(),
{},
id="no_filter_sends_no_parameters",
),
pytest.param(
lambda c: c.list_deployments(
name="agent", name_contains="agent", limit=100
),
{"name": "agent", "name_contains": "agent", "limit": "100"},
id="both_filters_travel_together_for_older_servers",
),
],
)
def test_list_deployments_sends_one_name_filter(call, expected_params):
seen: dict = {}
def handler(req: httpx.Request) -> httpx.Response:
seen.update(dict(req.url.params))
return httpx.Response(200, json={"resources": []})
call(
HostBackendClient(
"https://api.example.com", "key", transport=httpx.MockTransport(handler)
)
)
assert seen == expected_params
@pytest.mark.parametrize(
("body", "expected"),
[
pytest.param(
{"detail": "Source configuration error: bad listener"},
"Source configuration error: bad listener",
id="fastapi_detail_is_unwrapped",
),
pytest.param(
{"detail": {"loc": ["body"], "msg": "nope"}},
None,
id="a_structured_detail_is_left_alone",
),
pytest.param({"other": "shape"}, None, id="an_unknown_shape_is_left_alone"),
],
)
def test_error_detail_is_readable(body, expected):
c = HostBackendClient(
"https://api.example.com",
"key",
transport=httpx.MockTransport(lambda req: httpx.Response(400, json=body)),
)
with pytest.raises(HostBackendError) as error:
c.get_deployment("dep-1")
assert error.value.detail == expected
if expected is not None:
assert error.value.message.endswith(expected)
+14 -4
View File
@@ -1950,7 +1950,7 @@ class Pregel(
run_tasks: list[PregelTaskWrites] = [] run_tasks: list[PregelTaskWrites] = []
run_task_ids: list[str] = [] run_task_ids: list[str] = []
for as_node, values, provided_task_id in valid_updates: for i, (as_node, values, provided_task_id) in enumerate(valid_updates):
# create task to run all writers of the chosen node # create task to run all writers of the chosen node
writers = self.nodes[as_node].flat_writers writers = self.nodes[as_node].flat_writers
if not writers: if not writers:
@@ -1964,7 +1964,7 @@ class Pregel(
task_id = provided_task_id or ( task_id = provided_task_id or (
prepared_task_ids.popleft() prepared_task_ids.popleft()
if prepared_task_ids if prepared_task_ids
else str(uuid5(UUID(checkpoint["id"]), INTERRUPT)) else _update_task_id(checkpoint["id"], i)
) )
run_tasks.append(task) run_tasks.append(task)
run_task_ids.append(task_id) run_task_ids.append(task_id)
@@ -2410,7 +2410,7 @@ class Pregel(
run_tasks: list[PregelTaskWrites] = [] run_tasks: list[PregelTaskWrites] = []
run_task_ids: list[str] = [] run_task_ids: list[str] = []
for as_node, values, provided_task_id in valid_updates: for i, (as_node, values, provided_task_id) in enumerate(valid_updates):
# create task to run all writers of the chosen node # create task to run all writers of the chosen node
writers = self.nodes[as_node].flat_writers writers = self.nodes[as_node].flat_writers
if not writers: if not writers:
@@ -2424,7 +2424,7 @@ class Pregel(
task_id = provided_task_id or ( task_id = provided_task_id or (
prepared_task_ids.popleft() prepared_task_ids.popleft()
if prepared_task_ids if prepared_task_ids
else str(uuid5(UUID(checkpoint["id"]), INTERRUPT)) else _update_task_id(checkpoint["id"], i)
) )
run_tasks.append(task) run_tasks.append(task)
run_task_ids.append(task_id) run_task_ids.append(task_id)
@@ -4172,6 +4172,16 @@ class Pregel(
await self.cache.aclear(namespaces) await self.cache.aclear(namespaces)
def _update_task_id(checkpoint_id: str, i: int) -> str:
"""Task id for the `i`th update of a superstep that has no task to reuse.
Savers keep one write per `(task_id, idx)`, so updates sharing an id lose
all but the first one's writes, which a `DeltaChannel` replays from. The
first update keeps the id a lone update has always had.
"""
return str(uuid5(UUID(checkpoint_id), INTERRUPT if i == 0 else f"{INTERRUPT}:{i}"))
def _trigger_to_nodes(nodes: dict[str, PregelNode]) -> Mapping[str, Sequence[str]]: def _trigger_to_nodes(nodes: dict[str, PregelNode]) -> Mapping[str, Sequence[str]]:
"""Index from a trigger to nodes that depend on it.""" """Index from a trigger to nodes that depend on it."""
trigger_to_nodes: defaultdict[str, list[str]] = defaultdict(list) trigger_to_nodes: defaultdict[str, list[str]] = defaultdict(list)
@@ -20,6 +20,7 @@ from typing import Annotated, Any
import pytest import pytest
from langchain_core.messages import HumanMessage from langchain_core.messages import HumanMessage
from langgraph.checkpoint.base import BaseCheckpointSaver
from langgraph.checkpoint.memory import InMemorySaver from langgraph.checkpoint.memory import InMemorySaver
from langgraph.checkpoint.serde.types import _DeltaSnapshot from langgraph.checkpoint.serde.types import _DeltaSnapshot
from typing_extensions import TypedDict from typing_extensions import TypedDict
@@ -27,16 +28,17 @@ from typing_extensions import TypedDict
from langgraph.channels.delta import DeltaChannel from langgraph.channels.delta import DeltaChannel
from langgraph.graph import START, StateGraph from langgraph.graph import START, StateGraph
from langgraph.graph.message import _messages_delta_reducer from langgraph.graph.message import _messages_delta_reducer
from langgraph.types import StateUpdate from langgraph.types import StateSnapshot, StateUpdate
pytestmark = pytest.mark.anyio pytestmark = pytest.mark.anyio
def _build_graph( def _build_graph(
checkpointer: InMemorySaver, checkpointer: BaseCheckpointSaver,
*, *,
two_nodes: bool = False, two_nodes: bool = False,
snapshot_frequency: int = 1000, snapshot_frequency: int = 1000,
interrupt_before: list[str] | None = None,
) -> Any: ) -> Any:
"""Compile a minimal DeltaChannel-backed `messages` graph. """Compile a minimal DeltaChannel-backed `messages` graph.
@@ -63,7 +65,7 @@ def _build_graph(
builder.set_finish_point("assistant") builder.set_finish_point("assistant")
else: else:
builder.set_finish_point("model") builder.set_finish_point("model")
return builder.compile(checkpointer=checkpointer) return builder.compile(checkpointer=checkpointer, interrupt_before=interrupt_before)
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
@@ -273,10 +275,6 @@ def test_bulk_update_state_multi_task_per_superstep_delta_channel() -> None:
that each call `put_writes`. Guards the regression where moving that each call `put_writes`. Guards the regression where moving
`put_writes` outside the per-task loop would persist only the last `put_writes` outside the per-task loop would persist only the last
task's writes. task's writes.
Explicit `task_id`s are required to disambiguate writes belonging to
different `StateUpdate`s targeting the same node — otherwise both share
the deterministic interrupt-derived id and collide in the saver.
""" """
saver = InMemorySaver() saver = InMemorySaver()
@@ -310,6 +308,92 @@ def test_bulk_update_state_multi_task_per_superstep_delta_channel() -> None:
assert sorted(ids) == ["m1", "m2"] assert sorted(ids) == ["m1", "m2"]
def _update(content: str, as_node: str) -> StateUpdate:
return StateUpdate(
values={"messages": [HumanMessage(content=content, id=content)]},
as_node=as_node,
)
def _contents(state: StateSnapshot) -> list[str]:
return [m.content for m in state.values["messages"]]
def test_bulk_update_state_keeps_every_update_without_task_ids(
sync_checkpointer: BaseCheckpointSaver,
) -> None:
graph = _build_graph(sync_checkpointer, two_nodes=True)
config = {"configurable": {"thread_id": "bulk-no-task-ids"}}
graph.invoke({"messages": [HumanMessage(content="hi", id="hi")]}, config)
graph.bulk_update_state(
config,
[
[
_update("first", "model"),
_update("second", "model"),
_update("third", "assistant"),
]
],
)
contents = _contents(graph.get_state(config))
assert sorted(contents) == ["first", "hi", "second", "third"], (
f"every update's writes must persist; got {contents}"
)
async def test_abulk_update_state_keeps_every_update_without_task_ids(
async_checkpointer: BaseCheckpointSaver,
) -> None:
graph = _build_graph(async_checkpointer, two_nodes=True)
config = {"configurable": {"thread_id": "bulk-no-task-ids"}}
await graph.ainvoke({"messages": [HumanMessage(content="hi", id="hi")]}, config)
await graph.abulk_update_state(
config,
[
[
_update("first", "model"),
_update("second", "model"),
_update("third", "assistant"),
]
],
)
contents = _contents(await graph.aget_state(config))
assert sorted(contents) == ["first", "hi", "second", "third"], (
f"every update's writes must persist; got {contents}"
)
def test_bulk_update_state_keeps_every_update_next_to_a_pending_task(
sync_checkpointer: BaseCheckpointSaver,
) -> None:
graph = _build_graph(
sync_checkpointer, two_nodes=True, interrupt_before=["assistant"]
)
config = {"configurable": {"thread_id": "bulk-pending-task"}}
graph.invoke({"messages": [HumanMessage(content="hi", id="hi")]}, config)
assert graph.get_state(config).next == ("assistant",)
graph.bulk_update_state(
config,
[
[
_update("first", "assistant"),
_update("second", "model"),
_update("third", "model"),
]
],
)
contents = _contents(graph.get_state(config))
assert sorted(contents) == ["first", "hi", "second", "third"], (
f"every update's writes must persist; got {contents}"
)
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
# Public-API observation of fresh-thread checkpoint shape # Public-API observation of fresh-thread checkpoint shape
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------