Compare commits

..
Author SHA1 Message Date
Quanzheng LongandCursor 8da59aba37 chore: remove unnecessary missing-typed-dict-key suppression
Co-authored-by: Cursor <cursoragent@cursor.com>
2026-05-08 09:49:57 -07:00
Quanzheng LongandCursor 44805588b6 docs: clarify get_delta_channel_keepset is a basic reference implementation
Make it clear that custom backends may override with a more efficient
version but the default works correctly for any saver.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-05-07 17:10:20 -07:00
Quanzheng LongandCursor 68f8893847 test(conformance): add migration case for delta_channel_history
Test that a pre-delta plain value in channel_values[ch] (from a thread
that used BinaryOperatorAggregate before switching to DeltaChannel) is
correctly treated as the seed and terminates the walk.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-05-07 17:09:10 -07:00
Quanzheng Long 4cf12f9500 rm 2026-05-07 12:13:35 -07:00
Quanzheng LongandCursor a8ab1b1638 fix: resolve ruff import sorting in conformance wrapper tests
Co-authored-by: Cursor <cursoragent@cursor.com>
2026-05-07 11:26:40 -07:00
Quanzheng Long c78270c40f simp 2026-05-07 11:16:52 -07:00
Quanzheng LongandCursor 9dc68d6986 test: add delta-channel conformance wrappers for InMemory and SQLite savers
Runs the three new delta-channel conformance capabilities
(delta_channel_history, delta_channel_keepset, delta_channel_reconstruction)
against InMemorySaver and AsyncSqliteSaver. Guarded by importorskip so
they're skipped gracefully when checkpoint-conformance isn't installed.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-05-07 11:15:34 -07:00
Quanzheng Long d6c29f8157 simp 2026-05-07 11:14:51 -07:00
Quanzheng LongandCursor 4b3839af0f fix(conformance): lazy-import _DeltaSnapshot to avoid CI collection failure
Move _DeltaSnapshot imports from module level to function bodies so the
conformance test files can be collected even when the installed
langgraph-checkpoint version doesn't yet export the symbol.

Also fix get_delta_channel_keepset to check the target checkpoint's own
channel_values before walking the parent chain — if the target itself
has a snapshot, no walk is needed.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-05-07 11:07:39 -07:00
Quanzheng Long 569f2d2d14 done 2026-05-07 10:58:13 -07:00
Quanzheng Long 0d32281d5d dc 2026-05-07 10:17:24 -07:00
69f2d3a430 chore(langgraph): re-implement exit mode for delta channel (#7730)
## Summary

Replaces `durability="exit"`'s blanket force-snapshot of every
`DeltaChannel` with proper write persistence that honors per-channel
`snapshot_frequency`, plus closes two latent bugs the force-snapshot was
masking.

Before: every exit-mode run wrote a full `_DeltaSnapshot` blob for every
delta channel, even when the channel had zero updates this run and was
nowhere near its `snapshot_frequency`. After: the same count-based
decision used by `durability="sync"`/`"async"` applies — channels at or
above `snapshot_frequency` snapshot; channels below it persist their
accumulated writes via a lazy "stub" anchor; untouched channels write
nothing.

## What changed

**Core redesign** (`pregel/_loop.py`, `pregel/_checkpoint.py`)

- Drop `force_delta_snapshot` from `create_checkpoint` and
`_should_snapshot_delta`.
- Add `decide_delta_snapshots(channels, counts)` pure helper used by
both `create_checkpoint` and the new exit-mode peek-ahead path.
- Add `_exit_delta_writes` accumulator: every delta-channel write
produced during a `durability="exit"` run (input writes from `_first` +
per-superstep writes captured before `pending_writes.clear()` in
`after_tick`) is collected into this list.
- Add `_put_exit_delta_writes` (sync + async): runs from
`_suppress_interrupt` BEFORE `_put_checkpoint(exiting=True)`. Filters
out channels that will snapshot, then persists remaining writes to
`checkpoint_writes` under an anchor parent. The anchor is the existing
saved parent on resumed runs, or a lazily-created empty stub on first
runs.
- Visibility ordering: stub put goes onto `_put_checkpoint_fut` (becomes
the next put's `prev`); exit-write futures go onto `_delta_write_futs`.
The existing `_checkpointer_put_after_previous` already drains both
before calling `saver.put`, so `final_checkpoint` is structurally
guaranteed to land last — readers never see a partial view.

**Latent bugs fixed (previously masked by force-snapshot)**

- **Sync drain race**: `SyncPregelLoop` now initializes
`_delta_write_futs = []` in `__enter__` and drains it in sync
`_checkpointer_put_after_previous` before `put`, mirroring the async
version. Without this, a multi-worker `BackgroundExecutor` could publish
a checkpoint before the writes that produced it.
- **Count double-bump in exit mode**: in `_put_checkpoint`,
`delta_updates_since_snapshot` was being incremented twice for the last
superstep — once by the intermediate `after_tick` call, once by
`_suppress_interrupt`. Force-snapshot used to reset all counts to 0 so
this never persisted; without it, snapshots would fire one superstep
early after every exit-mode run. Fixed by gating the count-bump behind
`not exiting`.

**Pre-existing input-durability gap**

- In the plain (non-Command) input path of `_first`, delta-channel input
writes are now persisted via `put_writes` (mirroring the Command path),
so sub-frequency inputs survive a `get_state` on resumed runs in
`sync`/`async` durability. Note: first-run `sync`/`async` still has the
same gap (writes orphan on the synthetic-empty parent id). That's
flagged as a follow-up — out of scope for this PR.

## Test plan

- Existing `tests/test_pregel.py` and `tests/test_pregel_async.py` pass
unchanged.
- Existing `tests/test_channels.py` (29 tests) and
`tests/test_delta_channel_migration.py` pass unchanged.
- New `tests/test_exit_delta_persistence.py` (11 tests) covers:
- **Write-path**: zero-write exit (no stub), all-snapshot first run (no
stub), sub-freq first run (single shared stub), sub-freq resumed run
(anchor on saved parent), sync-vs-exit count parity, mixed
snapshot/non-snapshot channels, snapshot fires at frequency.
- **Read-path**: K-run replay chain reads correctly across
stub→saved-parent transition; metadata `delta_updates_since_snapshot`
round-trips correctly; mixed sync/exit durability alternation produces
correct final state; snapshot+tail-deltas combination reads correctly.
- `make format && make lint && make test` in `libs/langgraph/`.

---------

Co-authored-by: Cursor <cursoragent@cursor.com>
Co-authored-by: Sydney Runkle <sydneymarierunkle@gmail.com>
2026-05-07 09:47:01 -07:00
Asamu DavidandGitHub 95b41d058f release: bump cli version (#7734)
## Summary

This release adds the following to the`langgraph deploy` command:
- json event output 
- non-interactive mode
2026-05-07 17:29:13 +01:00
dependabot[bot]GitHubdependabot[bot] <49699333+dependabot[bot]@users.noreply.github.com>John KennedyClaude Opus 4.7
e49c093f48 chore(deps): bump ty from 0.0.23 to 0.0.33 in /libs/sdk-py (#7666)
Bumps the minor-and-patch group with 1 update in the /libs/sdk-py
directory: [ty](https://github.com/astral-sh/ty).

Updates `ty` from 0.0.23 to 0.0.33
<details>
<summary>Release notes</summary>
<p><em>Sourced from <a
href="https://github.com/astral-sh/ty/releases">ty's
releases</a>.</em></p>
<blockquote>
<h2>0.0.33</h2>
<h2>Release Notes</h2>
<p>Released on 2026-04-28.</p>
<h3>Notable changes</h3>
<ul>
<li>
<p>ty now prefers the declared type of an annotated assignment in more
situations (<a
href="https://redirect.github.com/astral-sh/ruff/pull/24802">#24802</a>).
Consider this example:</p>
<pre lang="py"><code>from some_library import untyped_function
<p>threshold: int | None = 0
result: str = untyped_function()
</code></pre></p>
<p>ty previously favored the <em>inferred</em> type of the right hand
side expression when <code>threshold</code> and <code>result</code> were
used. This is useful for <code>threshold</code>, as it allows something
like <code>threshold += 1</code> to work without an error: we know that
<code>threshold</code> could later become <code>None</code>, but
<em>right now</em>, we see that it is an <code>int</code>. However, for
<code>result</code>, the inferred type is <code>Unknown</code>. This is
<em>not</em> a useful type and it can lead to false negatives. Starting
with this release, ty will therefore prefer
the declared type <em>if the inferred and declared types are mutually
assignable</em>. In the above example, <code>threshold</code> will still
be inferred as <code>int</code> (or rather <code>Literal[1]</code>), but
<code>result</code> will now be inferred as <code>str</code>. If you
previously added <code>cast</code>s to work around this behavior, you
should be able to remove them after upgrading.</p>
</li>
</ul>
<h3>Bug fixes</h3>
<ul>
<li>Fix reporting of annotation-only locals as unused (<a
href="https://redirect.github.com/astral-sh/ruff/pull/24811">#24811</a>)</li>
<li>Fix project and workspace selection (<a
href="https://redirect.github.com/astral-sh/ruff/pull/24824">#24824</a>)</li>
<li>Fix go-to definition for generic classes (<a
href="https://redirect.github.com/astral-sh/ruff/pull/24714">#24714</a>)</li>
<li>Fix receiver coloring for aliased decorators (<a
href="https://redirect.github.com/astral-sh/ruff/pull/24884">#24884</a>)</li>
</ul>
<h3>LSP server</h3>
<ul>
<li>Add support for go-to definition in literal enum member inlay hints
(<a
href="https://redirect.github.com/astral-sh/ruff/pull/24792">#24792</a>)</li>
<li>Add support for &quot;baking&quot; keyword argument inlay hints into
the source code (<a
href="https://redirect.github.com/astral-sh/ruff/pull/24667">#24667</a>)</li>
<li>Don't allow inlay hint edits when introducing a non global scope
symbol (<a
href="https://redirect.github.com/astral-sh/ruff/pull/24797">#24797</a>)</li>
<li>Omit semantic highlighting for unresolved symbols (<a
href="https://redirect.github.com/astral-sh/ruff/pull/24718">#24718</a>)</li>
</ul>
<h3>Core type checking</h3>
<ul>
<li>Support narrowing with aliased conditional expressions (<a
href="https://redirect.github.com/astral-sh/ruff/pull/24302">#24302</a>)</li>
<li>Model short-circuiting control flow in Boolean expressions (<a
href="https://redirect.github.com/astral-sh/ruff/pull/24458">#24458</a>)</li>
<li>Handle <code>finally</code> blocks where all
<code>try</code>/<code>except</code> blocks are terminal (<a
href="https://redirect.github.com/astral-sh/ruff/pull/24882">#24882</a>)</li>
<li>Detect invalid <code>ClassVar</code> vs instance-attribute overrides
(<a
href="https://redirect.github.com/astral-sh/ruff/pull/24767">#24767</a>)</li>
<li>Emit diagnostic for invalid uses of <code>Unpack[...]</code> (<a
href="https://redirect.github.com/astral-sh/ruff/pull/24868">#24868</a>)</li>
<li>Infer lambda parameter types with <code>Callable</code> type context
(<a
href="https://redirect.github.com/astral-sh/ruff/pull/24317">#24317</a>)</li>
<li>Support <code>**</code> unpacking of <code>TypedDict</code> in
dict-literal assignments (<a
href="https://redirect.github.com/astral-sh/ruff/pull/24703">#24703</a>)</li>
<li>Support <code>Unpack[TypedDict]</code> in <code>**kwargs</code>
signatures (<a
href="https://redirect.github.com/astral-sh/ruff/pull/24653">#24653</a>)</li>
<li>Treat <code>[*xs]</code> as an irrefutable pattern when matching on
<code>Sequence</code> (<a
href="https://redirect.github.com/astral-sh/ruff/pull/24787">#24787</a>)</li>
<li>Improve generics solving for unions in invariant positions (<a
href="https://redirect.github.com/astral-sh/ruff/pull/24698">#24698</a>)</li>
<li>Improve generics solving for unions when matching against protocols
(<a
href="https://redirect.github.com/astral-sh/ruff/pull/24837">#24837</a>)</li>
</ul>
<h3>Diagnostics</h3>
<!-- raw HTML omitted -->
</blockquote>
<p>... (truncated)</p>
</details>
<details>
<summary>Changelog</summary>
<p><em>Sourced from <a
href="https://github.com/astral-sh/ty/blob/main/CHANGELOG.md">ty's
changelog</a>.</em></p>
<blockquote>
<h2>0.0.33</h2>
<p>Released on 2026-04-28.</p>
<h3>Notable changes</h3>
<ul>
<li>
<p>ty now prefers the declared type of an annotated assignment in more
situations (<a
href="https://redirect.github.com/astral-sh/ruff/pull/24802">#24802</a>).
Consider this example:</p>
<pre lang="py"><code>from some_library import untyped_function
<p>threshold: int | None = 0
result: str = untyped_function()
</code></pre></p>
<p>ty previously favored the <em>inferred</em> type of the right hand
side expression when <code>threshold</code> and <code>result</code> were
used. This is useful for <code>threshold</code>, as it allows something
like <code>threshold += 1</code> to work without an error: we know that
<code>threshold</code> could later become <code>None</code>, but
<em>right now</em>, we see that it is an <code>int</code>. However, for
<code>result</code>, the inferred type is <code>Unknown</code>. This is
<em>not</em> a useful type and it can lead to false negatives. Starting
with this release, ty will therefore prefer
the declared type <em>if the inferred and declared types are mutually
assignable</em>. In the above example, <code>threshold</code> will still
be inferred as <code>int</code> (or rather <code>Literal[1]</code>), but
<code>result</code> will now be inferred as <code>str</code>. If you
previously added <code>cast</code>s to work around this behavior, you
should be able to remove them after upgrading.</p>
</li>
</ul>
<h3>Bug fixes</h3>
<ul>
<li>Fix reporting of annotation-only locals as unused (<a
href="https://redirect.github.com/astral-sh/ruff/pull/24811">#24811</a>)</li>
<li>Fix project and workspace selection (<a
href="https://redirect.github.com/astral-sh/ruff/pull/24824">#24824</a>)</li>
<li>Fix go-to definition for generic classes (<a
href="https://redirect.github.com/astral-sh/ruff/pull/24714">#24714</a>)</li>
<li>Fix receiver coloring for aliased decorators (<a
href="https://redirect.github.com/astral-sh/ruff/pull/24884">#24884</a>)</li>
</ul>
<h3>LSP server</h3>
<ul>
<li>Add support for go-to definition in literal enum member inlay hints
(<a
href="https://redirect.github.com/astral-sh/ruff/pull/24792">#24792</a>)</li>
<li>Add support for &quot;baking&quot; keyword argument inlay hints into
the source code (<a
href="https://redirect.github.com/astral-sh/ruff/pull/24667">#24667</a>)</li>
<li>Don't allow inlay hint edits when introducing a non global scope
symbol (<a
href="https://redirect.github.com/astral-sh/ruff/pull/24797">#24797</a>)</li>
<li>Omit semantic highlighting for unresolved symbols (<a
href="https://redirect.github.com/astral-sh/ruff/pull/24718">#24718</a>)</li>
</ul>
<h3>Core type checking</h3>
<ul>
<li>Support narrowing with aliased conditional expressions (<a
href="https://redirect.github.com/astral-sh/ruff/pull/24302">#24302</a>)</li>
<li>Model short-circuiting control flow in Boolean expressions (<a
href="https://redirect.github.com/astral-sh/ruff/pull/24458">#24458</a>)</li>
<li>Handle <code>finally</code> blocks where all
<code>try</code>/<code>except</code> blocks are terminal (<a
href="https://redirect.github.com/astral-sh/ruff/pull/24882">#24882</a>)</li>
<li>Detect invalid <code>ClassVar</code> vs instance-attribute overrides
(<a
href="https://redirect.github.com/astral-sh/ruff/pull/24767">#24767</a>)</li>
<li>Emit diagnostic for invalid uses of <code>Unpack[...]</code> (<a
href="https://redirect.github.com/astral-sh/ruff/pull/24868">#24868</a>)</li>
<li>Infer lambda parameter types with <code>Callable</code> type context
(<a
href="https://redirect.github.com/astral-sh/ruff/pull/24317">#24317</a>)</li>
<li>Support <code>**</code> unpacking of <code>TypedDict</code> in
dict-literal assignments (<a
href="https://redirect.github.com/astral-sh/ruff/pull/24703">#24703</a>)</li>
<li>Support <code>Unpack[TypedDict]</code> in <code>**kwargs</code>
signatures (<a
href="https://redirect.github.com/astral-sh/ruff/pull/24653">#24653</a>)</li>
<li>Treat <code>[*xs]</code> as an irrefutable pattern when matching on
<code>Sequence</code> (<a
href="https://redirect.github.com/astral-sh/ruff/pull/24787">#24787</a>)</li>
<li>Improve generics solving for unions in invariant positions (<a
href="https://redirect.github.com/astral-sh/ruff/pull/24698">#24698</a>)</li>
<li>Improve generics solving for unions when matching against protocols
(<a
href="https://redirect.github.com/astral-sh/ruff/pull/24837">#24837</a>)</li>
</ul>
<h3>Diagnostics</h3>
<ul>
<li>Add error context to <code>invalid-return-type</code> diagnostics,
<code>invalid-yield</code> diagnostics, attribute assignment diagnostics
(<a
href="https://redirect.github.com/astral-sh/ruff/pull/24770">#24770</a>,
<a
href="https://redirect.github.com/astral-sh/ruff/pull/24771">#24771</a>)</li>
</ul>
<!-- raw HTML omitted -->
</blockquote>
<p>... (truncated)</p>
</details>
<details>
<summary>Commits</summary>
<ul>
<li><a
href="https://github.com/astral-sh/ty/commit/c512d8425418a2170e92aa7fbbd70952d4e04118"><code>c512d84</code></a>
Bump version to 0.0.33 (<a
href="https://redirect.github.com/astral-sh/ty/issues/3368">#3368</a>)</li>
<li><a
href="https://github.com/astral-sh/ty/commit/4cd7b334b90eba09042700e7b654044c2d6bcd15"><code>4cd7b33</code></a>
Upgrade Depot runners from macOS 14 to 15 (<a
href="https://redirect.github.com/astral-sh/ty/issues/3363">#3363</a>)</li>
<li><a
href="https://github.com/astral-sh/ty/commit/c78b8324515bf15a662f61e1d74fd85484d67e8c"><code>c78b832</code></a>
Update rui314/setup-mold digest to 9c9c13b (<a
href="https://redirect.github.com/astral-sh/ty/issues/3342">#3342</a>)</li>
<li><a
href="https://github.com/astral-sh/ty/commit/dea338134aef96d741a48886e414f622dbf10426"><code>dea3381</code></a>
Update actions/cache action to v5.0.5 (<a
href="https://redirect.github.com/astral-sh/ty/issues/3343">#3343</a>)</li>
<li><a
href="https://github.com/astral-sh/ty/commit/d451af477bd5517e132609c3b016799d697f4182"><code>d451af4</code></a>
update typing-features and faqs (<a
href="https://redirect.github.com/astral-sh/ty/issues/3335">#3335</a>)</li>
<li><a
href="https://github.com/astral-sh/ty/commit/052d70bc1a7457f50c5acbafe8de7a846f7d5e77"><code>052d70b</code></a>
Update prek dependencies (<a
href="https://redirect.github.com/astral-sh/ty/issues/3344">#3344</a>)</li>
<li><a
href="https://github.com/astral-sh/ty/commit/66b5e878163ce4ab4810a45a2679a98011932b8b"><code>66b5e87</code></a>
Update astral-sh/setup-uv action to v8.1.0 (<a
href="https://redirect.github.com/astral-sh/ty/issues/3345">#3345</a>)</li>
<li><a
href="https://github.com/astral-sh/ty/commit/7ec6712a6f02d0255de4c7f0066d5bb23987a9bb"><code>7ec6712</code></a>
Add a 'Diagnostics improvements' section to the changelogs (<a
href="https://redirect.github.com/astral-sh/ty/issues/3309">#3309</a>)</li>
<li><a
href="https://github.com/astral-sh/ty/commit/978dfdb38dfb568943f73779d769a965c1d5a397"><code>978dfdb</code></a>
Add version metadata publishing to the release process (<a
href="https://redirect.github.com/astral-sh/ty/issues/3292">#3292</a>)</li>
<li><a
href="https://github.com/astral-sh/ty/commit/4d1e1fc57ca8bfdcbcee513ba92135d2932eb279"><code>4d1e1fc</code></a>
Bump version to 0.0.32 (<a
href="https://redirect.github.com/astral-sh/ty/issues/3302">#3302</a>)</li>
<li>Additional commits viewable in <a
href="https://github.com/astral-sh/ty/compare/0.0.23...0.0.33">compare
view</a></li>
</ul>
</details>
<br />

---------

Signed-off-by: dependabot[bot] <support@github.com>
Co-authored-by: dependabot[bot] <49699333+dependabot[bot]@users.noreply.github.com>
Co-authored-by: John Kennedy <65985482+jkennedyvz@users.noreply.github.com>
Co-authored-by: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
2026-05-07 08:21:49 -04:00
Sydney RunkleandGitHub 9032a3f90a docs(checkpoint): mark DeltaChannel and delta-history APIs as beta (#7732)
## Summary

Adds a Beta admonition to the `DeltaChannel` surface area so users can
distinguish stable from in-progress contracts. Single-sourced on the
base class; subclass overrides in `checkpoint-postgres` /
`checkpoint-sqlite` / memory inherit the marker via their existing
references back to the base.

Touched:
- `DeltaChannel` (`libs/langgraph/langgraph/channels/delta.py`)
- `BaseCheckpointSaver.get_delta_channel_history` and
`aget_delta_channel_history` (`libs/checkpoint/.../base/__init__.py`)
- `DeltaChannelHistory` TypedDict
- `CheckpointMetadata.delta_updates_since_snapshot`

## Why docstring admonitions, not `@beta`

The base `get_delta_channel_history` methods are designed to be
overridden by savers. Wrapping with `@beta` (from `langchain_core._api`)
would emit warnings whenever a subclass called
`super().get_delta_channel_history(...)` or whenever the default
ancestor walk fired. Docstring-only keeps the signal advisory and
noise-free. `_DeltaSnapshot` is already leading-underscore-private, so
it implicitly signals "internal."

## Versioning note

This intentionally stays a **minor** bump for the checkpoint releases
(4.0 → 4.1, 3.0 → 3.1):

- The new methods are strictly additive — defaults provided on the base,
no signatures changed, no removed APIs. Third-party savers keep working
without overriding anything.
- The beta marker and the version bump do orthogonal jobs: semver
answers "is this a breaking change?" (no), the marker answers "is this
contract stable?" (no).
- Bumping major now would consume the lever you want available for when
the delta contract actually changes in a breaking way.

## Test plan

- [ ] Docstring-only — no behavioral change
- [x] `make format && make lint` clean in `libs/checkpoint` and
`libs/langgraph`
2026-05-07 06:44:54 -04:00
Sydney RunkleandGitHub 1a989f22bb fix(checkpoint-postgres): add column aliases to seed-blob branch of delta stage-2 UNION ALL (#7728)
## Summary

- Ports the langgraph API checkpointer fix to the OSS Postgres
checkpoint.
- The `'b'` (seed-blob) branch of `_build_delta_stage2_sql`'s `UNION
ALL` was missing column aliases (`AS _kind`, `AS checkpoint_id`, `AS
task_id`, `AS idx`).
- In PostgreSQL, `UNION ALL` column names are taken from the **first**
`SELECT`. When `channels_with_chain` is empty but `channels_with_seed`
is not (e.g. very first run of a DeltaChannel graph, or
`snapshot_frequency=1`), the `'b'` branch becomes the first `SELECT`, so
columns got default names (`text`, `null`) instead of the expected
aliases — causing `KeyError: '_kind'` in
`_build_delta_channels_writes_history`.

## Test plan

- [ ] Existing `checkpoint-postgres` delta channel tests pass (`make
test` in `libs/checkpoint-postgres`)
- [ ] Manually verified fix matches the correction suggested in the
upstream API PR review comment

🤖 Generated with [Claude Code](https://claude.com/claude-code)
2026-05-06 14:36:30 -04:00
36 changed files with 2310 additions and 737 deletions
+3
View File
@@ -63,6 +63,9 @@ The suite tests **base** capabilities (required) and **extended** capabilities (
| `delete_for_runs` | no | `adelete_for_runs` |
| `copy_thread` | no | `acopy_thread` |
| `prune` | no | `aprune` |
| `delta_channel_history` | no | `aget_delta_channel_history` |
| `delta_channel_keepset` | no | `aget_delta_channel_keepset` |
| `delta_channel_reconstruction` | no | `aput` |
Extended capabilities are detected by checking whether the method is overridden from `BaseCheckpointSaver`. If not overridden, those tests are skipped.
@@ -23,6 +23,9 @@ class Capability(str, Enum):
DELETE_FOR_RUNS = "delete_for_runs"
COPY_THREAD = "copy_thread"
PRUNE = "prune"
DELTA_CHANNEL_HISTORY = "delta_channel_history"
DELTA_CHANNEL_KEEPSET = "delta_channel_keepset"
DELTA_CHANNEL_RECONSTRUCTION = "delta_channel_reconstruction"
# Capabilities that every checkpointer must support.
@@ -42,6 +45,9 @@ EXTENDED_CAPABILITIES = frozenset(
Capability.DELETE_FOR_RUNS,
Capability.COPY_THREAD,
Capability.PRUNE,
Capability.DELTA_CHANNEL_HISTORY,
Capability.DELTA_CHANNEL_KEEPSET,
Capability.DELTA_CHANNEL_RECONSTRUCTION,
}
)
@@ -57,6 +63,9 @@ _CAPABILITY_METHOD_MAP: dict[Capability, str] = {
Capability.DELETE_FOR_RUNS: "adelete_for_runs",
Capability.COPY_THREAD: "acopy_thread",
Capability.PRUNE: "aprune",
Capability.DELTA_CHANNEL_HISTORY: "aget_delta_channel_history",
Capability.DELTA_CHANNEL_KEEPSET: "aget_tuple",
Capability.DELTA_CHANNEL_RECONSTRUCTION: "aput",
}
@@ -9,6 +9,15 @@ from langgraph.checkpoint.conformance.spec.test_delete_for_runs import (
from langgraph.checkpoint.conformance.spec.test_delete_thread import (
run_delete_thread_tests,
)
from langgraph.checkpoint.conformance.spec.test_delta_channel_history import (
run_delta_channel_history_tests,
)
from langgraph.checkpoint.conformance.spec.test_delta_channel_keepset import (
run_delta_channel_keepset_tests,
)
from langgraph.checkpoint.conformance.spec.test_delta_channel_reconstruction import (
run_delta_channel_reconstruction_tests,
)
from langgraph.checkpoint.conformance.spec.test_get_tuple import run_get_tuple_tests
from langgraph.checkpoint.conformance.spec.test_list import run_list_tests
from langgraph.checkpoint.conformance.spec.test_prune import run_prune_tests
@@ -24,4 +33,7 @@ __all__ = [
"run_delete_for_runs_tests",
"run_copy_thread_tests",
"run_prune_tests",
"run_delta_channel_history_tests",
"run_delta_channel_keepset_tests",
"run_delta_channel_reconstruction_tests",
]
@@ -0,0 +1,99 @@
"""Shared fixtures for delta-channel conformance tests.
Builds a parent chain with `_DeltaSnapshot` blobs at known positions via
direct `aput` / `aput_writes` calls. No langgraph or Pregel dependency.
"""
from __future__ import annotations
from collections.abc import Sequence
from typing import Any
from uuid import uuid4
from langchain_core.runnables import RunnableConfig
from langgraph.checkpoint.base import BaseCheckpointSaver, Checkpoint
from langgraph.checkpoint.base.id import uuid6
from langgraph.checkpoint.conformance.test_utils import generate_metadata
async def build_delta_chain(
saver: BaseCheckpointSaver,
*,
thread_id: str | None = None,
checkpoint_ns: str = "",
channel: str = "messages",
snapshots_at_steps: Sequence[int] = (0,),
total_steps: int = 6,
write_value_fn: Any | None = None,
) -> list[RunnableConfig]:
"""Build a parent chain with `_DeltaSnapshot` at known positions.
Args:
saver: Checkpointer instance.
thread_id: Defaults to a random UUID.
checkpoint_ns: Namespace (default root).
channel: Channel name used for snapshots and writes.
snapshots_at_steps: Steps at which a `_DeltaSnapshot` blob is stored
in `channel_values[channel]`. Step 0 is the oldest checkpoint.
total_steps: Number of checkpoints in the chain.
write_value_fn: Callable(step) -> write value. Defaults to step index.
Returns:
List of stored configs (oldest first), one per step.
"""
if write_value_fn is None:
def write_value_fn(step: int) -> Any:
return step
from langgraph.checkpoint.serde.types import _DeltaSnapshot
thread_id = thread_id or str(uuid4())
snapshot_set = set(snapshots_at_steps)
stored: list[RunnableConfig] = []
parent_cfg: RunnableConfig | None = None
for step in range(total_steps):
config: RunnableConfig = {
"configurable": {
"thread_id": thread_id,
"checkpoint_ns": checkpoint_ns,
}
}
if parent_cfg:
config["configurable"]["checkpoint_id"] = parent_cfg["configurable"][
"checkpoint_id"
]
channel_values: dict[str, Any] = {}
channel_versions: dict[str, int] = {}
if step in snapshot_set:
channel_values[channel] = _DeltaSnapshot(
write_value_fn(step),
)
channel_versions[channel] = step + 1
cp = Checkpoint(
v=1,
id=str(uuid6(clock_seq=-1)),
ts="",
channel_values=channel_values,
channel_versions=channel_versions,
versions_seen={},
updated_channels=None,
)
new_versions = dict(channel_versions)
parent_cfg = await saver.aput(
config, cp, generate_metadata(step=step), new_versions
)
stored.append(parent_cfg)
# Write a pending write for non-snapshot steps so the walk has
# something to collect.
if step not in snapshot_set:
await saver.aput_writes(
parent_cfg, [(channel, write_value_fn(step))], str(uuid4())
)
return stored
@@ -0,0 +1,247 @@
"""DELTA_CHANNEL_HISTORY capability tests — aget_delta_channel_history contract."""
from __future__ import annotations
import traceback
from collections.abc import Callable
from uuid import uuid4
from langgraph.checkpoint.base import BaseCheckpointSaver
from langgraph.checkpoint.conformance.spec._delta_fixtures import build_delta_chain
async def test_history_returns_writes_oldest_first(
saver: BaseCheckpointSaver,
) -> None:
"""Writes are returned oldest-to-newest."""
tid = str(uuid4())
# 5 steps: snapshot at 0, writes at 1,2,3,4.
# Head is step 4. Walk starts at step 3 (parent of head).
# Collects writes from steps 1,2,3 (between snapshot at 0 and head's parent).
configs = await build_delta_chain(
saver, thread_id=tid, channel="ch", snapshots_at_steps=[0], total_steps=5
)
head = configs[-1]
result = await saver.aget_delta_channel_history(config=head, channels=["ch"])
writes = result["ch"]["writes"]
values = [w[2] for w in writes]
assert values == [1, 2, 3], f"Expected [1,2,3], got {values}"
async def test_history_seed_is_nearest_snapshot(
saver: BaseCheckpointSaver,
) -> None:
"""Seed is the value from the nearest ancestor with channel_values populated."""
tid = str(uuid4())
# 6 steps: snapshots at 0 and 3, writes at 1,2,4,5.
# Head is step 5. Walk from step 4 backward stops at step 3 (snapshot).
# Collects writes from step 4 only (between step 3 and head's parent step 4).
configs = await build_delta_chain(
saver,
thread_id=tid,
channel="ch",
snapshots_at_steps=[0, 3],
total_steps=6,
)
head = configs[-1]
result = await saver.aget_delta_channel_history(config=head, channels=["ch"])
assert "seed" in result["ch"], "Expected seed from snapshot at step 3"
seed = result["ch"]["seed"]
from langgraph.checkpoint.serde.types import _DeltaSnapshot
actual_value = seed.value if isinstance(seed, _DeltaSnapshot) else seed
assert actual_value == 3, f"Expected seed value 3 (step 3), got {actual_value}"
writes = result["ch"]["writes"]
values = [w[2] for w in writes]
assert values == [4], f"Expected [4], got {values}"
async def test_history_excludes_target_pending_writes(
saver: BaseCheckpointSaver,
) -> None:
"""Target's own pending_writes are NOT included in the history."""
tid = str(uuid4())
configs = await build_delta_chain(
saver, thread_id=tid, channel="ch", snapshots_at_steps=[0], total_steps=3
)
head = configs[-1]
# Add writes directly to the head checkpoint
await saver.aput_writes(head, [("ch", "extra")], str(uuid4()))
result = await saver.aget_delta_channel_history(config=head, channels=["ch"])
writes = result["ch"]["writes"]
values = [w[2] for w in writes]
assert "extra" not in values, f"Target's writes should be excluded, got {values}"
async def test_history_multi_channel(
saver: BaseCheckpointSaver,
) -> None:
"""Multiple channels have independent walk termination."""
tid = str(uuid4())
configs: list = []
parent_cfg = None
from langgraph.checkpoint.base import Checkpoint
from langgraph.checkpoint.base.id import uuid6
from langgraph.checkpoint.serde.types import _DeltaSnapshot
from langgraph.checkpoint.conformance.test_utils import generate_metadata
for step in range(5):
config = {"configurable": {"thread_id": tid, "checkpoint_ns": ""}}
if parent_cfg:
config["configurable"]["checkpoint_id"] = parent_cfg["configurable"][
"checkpoint_id"
]
cv: dict = {}
cvs: dict = {}
if step == 1:
cv["a"] = _DeltaSnapshot("snap_a")
cvs["a"] = step + 1
if step == 3:
cv["b"] = _DeltaSnapshot("snap_b")
cvs["b"] = step + 1
cp = Checkpoint(
v=1,
id=str(uuid6(clock_seq=-1)),
ts="",
channel_values=cv,
channel_versions=cvs,
versions_seen={},
updated_channels=None,
)
parent_cfg = await saver.aput(config, cp, generate_metadata(step=step), cvs)
configs.append(parent_cfg)
await saver.aput_writes(parent_cfg, [("a", step), ("b", step)], str(uuid4()))
head = configs[-1]
result = await saver.aget_delta_channel_history(config=head, channels=["a", "b"])
a_writes = [w[2] for w in result["a"]["writes"]]
b_writes = [w[2] for w in result["b"]["writes"]]
assert a_writes == [1, 2, 3], f"Expected a writes [1,2,3], got {a_writes}"
assert b_writes == [3], f"Expected b writes [3], got {b_writes}"
async def test_history_empty_channels_returns_empty(
saver: BaseCheckpointSaver,
) -> None:
"""Empty channels list returns empty mapping."""
tid = str(uuid4())
configs = await build_delta_chain(
saver, thread_id=tid, channel="ch", snapshots_at_steps=[0], total_steps=3
)
result = await saver.aget_delta_channel_history(config=configs[-1], channels=[])
assert result == {}
async def test_history_walk_to_root_no_seed(
saver: BaseCheckpointSaver,
) -> None:
"""Walk reaches root without finding seed — no 'seed' key in result."""
tid = str(uuid4())
configs = await build_delta_chain(
saver,
thread_id=tid,
channel="ch",
snapshots_at_steps=[],
total_steps=4,
)
head = configs[-1]
result = await saver.aget_delta_channel_history(config=head, channels=["ch"])
assert "seed" not in result["ch"], f"Expected no seed, got {result['ch']}"
async def test_history_migration_plain_value_as_seed(
saver: BaseCheckpointSaver,
) -> None:
"""Pre-delta plain value in channel_values acts as seed (migration case).
When a thread was originally using a regular channel (BinaryOperatorAggregate)
and later switches to DeltaChannel, the old checkpoint has a plain value in
channel_values[ch] (not a _DeltaSnapshot). The walk should treat it as the
seed and terminate there.
"""
from langgraph.checkpoint.base import Checkpoint
from langgraph.checkpoint.base.id import uuid6
from langgraph.checkpoint.conformance.test_utils import generate_metadata
tid = str(uuid4())
configs: list = []
parent_cfg = None
for step in range(4):
config = {"configurable": {"thread_id": tid, "checkpoint_ns": ""}}
if parent_cfg:
config["configurable"]["checkpoint_id"] = parent_cfg["configurable"][
"checkpoint_id"
]
cv: dict = {}
cvs: dict = {}
# Step 1: plain value (migration case — old checkpoint before delta)
if step == 1:
cv["ch"] = [10, 20, 30]
cvs["ch"] = step + 1
cp = Checkpoint(
v=1,
id=str(uuid6(clock_seq=-1)),
ts="",
channel_values=cv,
channel_versions=cvs,
versions_seen={},
updated_channels=None,
)
parent_cfg = await saver.aput(config, cp, generate_metadata(step=step), cvs)
configs.append(parent_cfg)
if step != 1:
await saver.aput_writes(parent_cfg, [("ch", step)], str(uuid4()))
head = configs[-1]
result = await saver.aget_delta_channel_history(config=head, channels=["ch"])
# Seed should be the plain value from step 1
assert "seed" in result["ch"], "Expected seed from migration plain value at step 1"
seed = result["ch"]["seed"]
assert seed == [10, 20, 30], f"Expected plain value [10,20,30], got {seed}"
# Writes should be from step 2 only (between seed at step 1 and head's parent step 2)
writes = result["ch"]["writes"]
values = [w[2] for w in writes]
assert values == [2], f"Expected [2], got {values}"
ALL_DELTA_CHANNEL_HISTORY_TESTS = [
test_history_returns_writes_oldest_first,
test_history_seed_is_nearest_snapshot,
test_history_excludes_target_pending_writes,
test_history_multi_channel,
test_history_empty_channels_returns_empty,
test_history_walk_to_root_no_seed,
test_history_migration_plain_value_as_seed,
]
async def run_delta_channel_history_tests(
saver: BaseCheckpointSaver,
on_test_result: Callable[[str, str, bool, str | None], None] | None = None,
) -> tuple[int, int, list[str]]:
"""Run all delta_channel_history tests. Returns (passed, failed, failure_names)."""
passed = 0
failed = 0
failures: list[str] = []
for test_fn in ALL_DELTA_CHANNEL_HISTORY_TESTS:
try:
await test_fn(saver)
passed += 1
if on_test_result:
on_test_result("delta_channel_history", test_fn.__name__, True, None)
except Exception:
failed += 1
msg = f"{test_fn.__name__}: {traceback.format_exc()}"
failures.append(msg)
if on_test_result:
on_test_result(
"delta_channel_history",
test_fn.__name__,
False,
traceback.format_exc(),
)
return passed, failed, failures
@@ -0,0 +1,172 @@
"""DELTA_CHANNEL_KEEPSET capability tests — aget_delta_channel_keepset contract."""
from __future__ import annotations
import traceback
from collections.abc import Callable
from uuid import uuid4
from langgraph.checkpoint.base import BaseCheckpointSaver
from langgraph.checkpoint.conformance.spec._delta_fixtures import build_delta_chain
async def test_keepset_empty_channels_returns_target_only(
saver: BaseCheckpointSaver,
) -> None:
"""Empty channels → keep-set is just {target_id}."""
tid = str(uuid4())
configs = await build_delta_chain(
saver, thread_id=tid, channel="ch", snapshots_at_steps=[0], total_steps=4
)
head = configs[-1]
keep = await saver.aget_delta_channel_keepset(config=head, channels=[])
head_id = head["configurable"]["checkpoint_id"]
assert keep == {head_id}, f"Expected only target, got {keep}"
async def test_keepset_snapshot_at_target(
saver: BaseCheckpointSaver,
) -> None:
"""When target itself has a snapshot, keep-set is just {target_id}."""
tid = str(uuid4())
configs = await build_delta_chain(
saver,
thread_id=tid,
channel="ch",
snapshots_at_steps=[0, 3],
total_steps=4,
)
head = configs[3]
keep = await saver.aget_delta_channel_keepset(config=head, channels=["ch"])
head_id = head["configurable"]["checkpoint_id"]
assert keep == {head_id}, f"Snapshot at target should yield only target, got {keep}"
async def test_keepset_snapshot_n_back(
saver: BaseCheckpointSaver,
) -> None:
"""Snapshot N steps back → target + intermediates + snapshot ancestor."""
tid = str(uuid4())
configs = await build_delta_chain(
saver, thread_id=tid, channel="ch", snapshots_at_steps=[0], total_steps=5
)
head = configs[-1]
keep = await saver.aget_delta_channel_keepset(config=head, channels=["ch"])
expected_ids = {c["configurable"]["checkpoint_id"] for c in configs}
assert keep == expected_ids, f"Expected all ancestors, got {keep}"
async def test_keepset_multi_channel_union(
saver: BaseCheckpointSaver,
) -> None:
"""Multi-channel keep-set is the union (max chain per channel)."""
tid = str(uuid4())
from langgraph.checkpoint.base import Checkpoint
from langgraph.checkpoint.base.id import uuid6
from langgraph.checkpoint.serde.types import _DeltaSnapshot
from langgraph.checkpoint.conformance.test_utils import generate_metadata
configs: list = []
parent_cfg = None
for step in range(6):
config = {"configurable": {"thread_id": tid, "checkpoint_ns": ""}}
if parent_cfg:
config["configurable"]["checkpoint_id"] = parent_cfg["configurable"][
"checkpoint_id"
]
cv: dict = {}
cvs: dict = {}
# Channel "a" has snapshot at step 3 (recent)
if step == 3:
cv["a"] = _DeltaSnapshot("snap_a")
cvs["a"] = step + 1
# Channel "b" has snapshot at step 1 (further back)
if step == 1:
cv["b"] = _DeltaSnapshot("snap_b")
cvs["b"] = step + 1
cp = Checkpoint(
v=1,
id=str(uuid6(clock_seq=-1)),
ts="",
channel_values=cv,
channel_versions=cvs,
versions_seen={},
updated_channels=None,
)
parent_cfg = await saver.aput(config, cp, generate_metadata(step=step), cvs)
configs.append(parent_cfg)
head = configs[-1]
keep = await saver.aget_delta_channel_keepset(config=head, channels=["a", "b"])
# Union: b needs back to step 1, so steps 1..5 (all except step 0) plus head
expected_ids = {c["configurable"]["checkpoint_id"] for c in configs[1:]}
expected_ids.add(head["configurable"]["checkpoint_id"])
assert keep == expected_ids, f"Expected union, got {keep} vs {expected_ids}"
async def test_keepset_walk_to_root(
saver: BaseCheckpointSaver,
) -> None:
"""No snapshot anywhere → entire chain to root is in keep-set."""
tid = str(uuid4())
configs = await build_delta_chain(
saver, thread_id=tid, channel="ch", snapshots_at_steps=[], total_steps=4
)
head = configs[-1]
keep = await saver.aget_delta_channel_keepset(config=head, channels=["ch"])
all_ids = {c["configurable"]["checkpoint_id"] for c in configs}
assert keep == all_ids, f"Expected full chain, got {keep}"
async def test_keepset_deterministic(
saver: BaseCheckpointSaver,
) -> None:
"""Same inputs return identical sets."""
tid = str(uuid4())
configs = await build_delta_chain(
saver, thread_id=tid, channel="ch", snapshots_at_steps=[0], total_steps=5
)
head = configs[-1]
keep1 = await saver.aget_delta_channel_keepset(config=head, channels=["ch"])
keep2 = await saver.aget_delta_channel_keepset(config=head, channels=["ch"])
assert keep1 == keep2
ALL_DELTA_CHANNEL_KEEPSET_TESTS = [
test_keepset_empty_channels_returns_target_only,
test_keepset_snapshot_at_target,
test_keepset_snapshot_n_back,
test_keepset_multi_channel_union,
test_keepset_walk_to_root,
test_keepset_deterministic,
]
async def run_delta_channel_keepset_tests(
saver: BaseCheckpointSaver,
on_test_result: Callable[[str, str, bool, str | None], None] | None = None,
) -> tuple[int, int, list[str]]:
"""Run all delta_channel_keepset tests. Returns (passed, failed, failure_names)."""
passed = 0
failed = 0
failures: list[str] = []
for test_fn in ALL_DELTA_CHANNEL_KEEPSET_TESTS:
try:
await test_fn(saver)
passed += 1
if on_test_result:
on_test_result("delta_channel_keepset", test_fn.__name__, True, None)
except Exception:
failed += 1
msg = f"{test_fn.__name__}: {traceback.format_exc()}"
failures.append(msg)
if on_test_result:
on_test_result(
"delta_channel_keepset",
test_fn.__name__,
False,
traceback.format_exc(),
)
return passed, failed, failures
@@ -0,0 +1,164 @@
"""DELTA_CHANNEL_RECONSTRUCTION capability tests — end-to-end round-trip.
Exercises: aput + aput_writes + aget_delta_channel_history + reconstruction.
This catches the most common silent-corruption mode: failing to round-trip
`_DeltaSnapshot` blobs through serialization.
NOTE: This test does NOT import from `langgraph` (which is not a dependency
of checkpoint-conformance). Instead it inlines a minimal reconstruction
equivalent: seed + fold writes through a simple list-append reducer.
"""
from __future__ import annotations
import traceback
from collections.abc import Callable
from typing import Any
from uuid import uuid4
from langgraph.checkpoint.base import BaseCheckpointSaver
from langgraph.checkpoint.conformance.spec._delta_fixtures import build_delta_chain
def _reconstruct(seed: Any, writes: list) -> list:
"""Minimal DeltaChannel reconstruction: list-append reducer.
Mirrors DeltaChannel.from_checkpoint(seed) + replay_writes(writes).
"""
from langgraph.checkpoint.serde.types import _DeltaSnapshot
if seed is None:
base: list = []
elif isinstance(seed, _DeltaSnapshot):
base = list(seed.value)
else:
base = list(seed)
for _task_id, _ch, value in writes:
base = base + value
return base
async def test_reconstruction_basic(
saver: BaseCheckpointSaver,
) -> None:
"""Reconstruct delta channel value from history matches expected."""
tid = str(uuid4())
# 5 steps: snapshot at 0 (value=[0]), writes at 1,2,3,4.
# Head = step 4. Walk from parent (step 3) collects writes 1,2,3.
configs = await build_delta_chain(
saver,
thread_id=tid,
channel="msgs",
snapshots_at_steps=[0],
total_steps=5,
write_value_fn=lambda step: [step],
)
head = configs[-1]
result = await saver.aget_delta_channel_history(config=head, channels=["msgs"])
history = result["msgs"]
seed = history.get("seed")
reconstructed = _reconstruct(seed, history["writes"])
# seed=[0] from step 0, writes from steps 1,2,3
expected = [0] + [1] + [2] + [3]
assert reconstructed == expected, (
f"Reconstructed {reconstructed} != expected {expected}"
)
async def test_reconstruction_mid_chain_snapshot(
saver: BaseCheckpointSaver,
) -> None:
"""Reconstruction works when snapshot is mid-chain."""
tid = str(uuid4())
# 6 steps: snapshots at 0 and 3, writes at 1,2,4,5.
# Head = step 5. Walk from step 4 stops at step 3 (snapshot).
# Collects writes from step 4.
configs = await build_delta_chain(
saver,
thread_id=tid,
channel="msgs",
snapshots_at_steps=[0, 3],
total_steps=6,
write_value_fn=lambda step: [step],
)
head = configs[-1]
result = await saver.aget_delta_channel_history(config=head, channels=["msgs"])
history = result["msgs"]
seed = history.get("seed")
reconstructed = _reconstruct(seed, history["writes"])
# Snapshot at step 3 = [3], write from step 4
expected = [3] + [4]
assert reconstructed == expected, (
f"Reconstructed {reconstructed} != expected {expected}"
)
async def test_reconstruction_no_snapshot(
saver: BaseCheckpointSaver,
) -> None:
"""Reconstruction from root (no snapshot) gives all writes accumulated."""
tid = str(uuid4())
# 4 steps: no snapshot, writes at 0,1,2,3.
# Head = step 3. Walk from step 2 collects writes 0,1,2.
configs = await build_delta_chain(
saver,
thread_id=tid,
channel="msgs",
snapshots_at_steps=[],
total_steps=4,
write_value_fn=lambda step: [step],
)
head = configs[-1]
result = await saver.aget_delta_channel_history(config=head, channels=["msgs"])
history = result["msgs"]
seed = history.get("seed")
reconstructed = _reconstruct(seed, history["writes"])
# No seed → start empty, writes from steps 0,1,2
expected = [0] + [1] + [2]
assert reconstructed == expected, (
f"Reconstructed {reconstructed} != expected {expected}"
)
ALL_DELTA_CHANNEL_RECONSTRUCTION_TESTS = [
test_reconstruction_basic,
test_reconstruction_mid_chain_snapshot,
test_reconstruction_no_snapshot,
]
async def run_delta_channel_reconstruction_tests(
saver: BaseCheckpointSaver,
on_test_result: Callable[[str, str, bool, str | None], None] | None = None,
) -> tuple[int, int, list[str]]:
"""Run all reconstruction tests. Returns (passed, failed, failure_names)."""
passed = 0
failed = 0
failures: list[str] = []
for test_fn in ALL_DELTA_CHANNEL_RECONSTRUCTION_TESTS:
try:
await test_fn(saver)
passed += 1
if on_test_result:
on_test_result(
"delta_channel_reconstruction", test_fn.__name__, True, None
)
except Exception:
failed += 1
msg = f"{test_fn.__name__}: {traceback.format_exc()}"
failures.append(msg)
if on_test_result:
on_test_result(
"delta_channel_reconstruction",
test_fn.__name__,
False,
traceback.format_exc(),
)
return passed, failed, failures
@@ -19,6 +19,15 @@ from langgraph.checkpoint.conformance.spec.test_delete_for_runs import (
from langgraph.checkpoint.conformance.spec.test_delete_thread import (
run_delete_thread_tests,
)
from langgraph.checkpoint.conformance.spec.test_delta_channel_history import (
run_delta_channel_history_tests,
)
from langgraph.checkpoint.conformance.spec.test_delta_channel_keepset import (
run_delta_channel_keepset_tests,
)
from langgraph.checkpoint.conformance.spec.test_delta_channel_reconstruction import (
run_delta_channel_reconstruction_tests,
)
from langgraph.checkpoint.conformance.spec.test_get_tuple import run_get_tuple_tests
from langgraph.checkpoint.conformance.spec.test_list import run_list_tests
from langgraph.checkpoint.conformance.spec.test_prune import run_prune_tests
@@ -35,6 +44,9 @@ _RUNNERS = {
Capability.DELETE_FOR_RUNS: run_delete_for_runs_tests,
Capability.COPY_THREAD: run_copy_thread_tests,
Capability.PRUNE: run_prune_tests,
Capability.DELTA_CHANNEL_HISTORY: run_delta_channel_history_tests,
Capability.DELTA_CHANNEL_KEEPSET: run_delta_channel_keepset_tests,
Capability.DELTA_CHANNEL_RECONSTRUCTION: run_delta_channel_reconstruction_tests,
}
@@ -43,7 +43,11 @@ asyncio_mode = "auto"
# The extended methods (acopy_thread, adelete_for_runs, aprune) are checked
# at runtime via capability detection and may not exist on the installed
# base class. Dict literal inference is also overly strict for RunnableConfig.
# Delta-channel tests import from `langgraph` (not a declared dep of this
# package — at test time it is installed alongside); private `_DeltaSnapshot`
# imports are intentional (beta surface).
unresolved-attribute = "ignore"
unresolved-import = "ignore"
invalid-argument-type = "ignore"
invalid-return-type = "ignore"
@@ -58,6 +62,9 @@ lint.select = [
lint.ignore = ["E501", "B008"]
target-version = "py310"
[tool.uv.sources]
langgraph-checkpoint = {path = "../checkpoint", editable = true}
[[tool.uv.index]]
name = "testpypi"
url = "https://test.pypi.org/simple/"
+605 -504
View File
File diff suppressed because it is too large Load Diff
@@ -279,8 +279,8 @@ def _build_delta_stage2_sql(
)
for _ in channels_with_seed:
branches.append(
"SELECT 'b'::text, NULL, channel, "
"type, blob, NULL, NULL, version "
"SELECT 'b'::text AS _kind, NULL::text AS checkpoint_id, channel, "
"type, blob, NULL::text AS task_id, NULL::int AS idx, version "
"FROM checkpoint_blobs "
"WHERE thread_id = %s AND checkpoint_ns = %s AND channel = %s "
"AND version = %s"
@@ -0,0 +1,37 @@
"""Run delta-channel conformance capabilities against AsyncSqliteSaver."""
from __future__ import annotations
import pytest
pytest.importorskip(
"langgraph.checkpoint.conformance",
reason="langgraph-checkpoint-conformance not installed",
)
pytest.importorskip("aiosqlite", reason="aiosqlite not installed")
@pytest.mark.asyncio
async def test_delta_channel_conformance():
from langgraph.checkpoint.conformance import validate
from langgraph.checkpoint.conformance.initializer import checkpointer_test
from langgraph.checkpoint.sqlite.aio import AsyncSqliteSaver
@checkpointer_test(name="AsyncSqliteSaver")
async def sqlite_saver():
async with AsyncSqliteSaver.from_conn_string(":memory:") as saver:
yield saver
report = await validate(
sqlite_saver,
capabilities={
"delta_channel_history",
"delta_channel_keepset",
"delta_channel_reconstruction",
},
)
for cap, result in report.results.items():
if result.passed is False:
details = "\n".join(result.failures or [])
pytest.fail(f"Capability {cap} failed:\n{details}")
@@ -63,6 +63,11 @@ class CheckpointMetadata(TypedDict, total=False):
delta_updates_since_snapshot: dict[str, int]
"""Per-channel update count since the last `_DeltaSnapshot` was written.
!!! warning "Beta"
This metadata field backs `DeltaChannel` (beta). The key name and
contents may change while the delta-channel design stabilizes.
Maps channel name → number of supersteps that wrote to this channel
since its last snapshot blob. Used by `pregel.create_checkpoint` to
decide when to write the next snapshot (when the count reaches the
@@ -135,6 +140,11 @@ class CheckpointTuple(NamedTuple):
class DeltaChannelHistory(TypedDict):
"""Per-channel result entry from `BaseCheckpointSaver.get_delta_channel_history`.
!!! warning "Beta"
Part of the `DeltaChannel` support surface; in beta. Field names and
semantics may change.
Storage-level view of what one channel contributed across the ancestor
chain of a target checkpoint:
@@ -317,6 +327,14 @@ class BaseCheckpointSaver(Generic[V]):
Args:
run_ids: The run IDs whose checkpoints should be deleted.
!!! warning "DeltaChannel"
Deleting a run that produced ancestor `checkpoint_writes` — or
the only `_DeltaSnapshot` blob — for a still-live thread will
break reconstruction of any `DeltaChannel` whose history
depended on those rows. See the `DeltaChannel` note on `prune`
for safe-recovery strategies.
"""
raise NotImplementedError
@@ -330,6 +348,17 @@ class BaseCheckpointSaver(Generic[V]):
Args:
source_thread_id: The thread ID to copy from.
target_thread_id: The thread ID to copy to.
!!! warning "DeltaChannel"
Implementations must copy the **complete** parent chain (all
ancestor checkpoints and their `checkpoint_writes`) — copying
only the head checkpoint will leave the target thread with
`DeltaChannel` state that cannot be reconstructed (no path back
to a `_DeltaSnapshot` ancestor). Equivalently, the copy must
include enough ancestors that every `DeltaChannel`-backed key
has either a `_DeltaSnapshot` in `channel_values` somewhere in
the chain, or a complete write history back to the chain root.
"""
raise NotImplementedError
@@ -345,6 +374,34 @@ class BaseCheckpointSaver(Generic[V]):
thread_ids: The thread IDs to prune.
strategy: The pruning strategy. `"keep_latest"` retains only the most
recent checkpoint per namespace. `"delete"` removes all checkpoints.
!!! warning "DeltaChannel"
Custom implementations must be `DeltaChannel`-aware. `DeltaChannel`
stores only a sentinel in `channel_values` for non-snapshot steps;
reconstruction walks the parent chain via
`get_delta_channel_history`, accumulating rows from
`checkpoint_writes` until it reaches an ancestor whose
`channel_values` contains a `_DeltaSnapshot` blob (written every
`snapshot_frequency` updates).
A naive `"keep_latest"` that drops intermediate checkpoints and
their writes can sever that chain: the surviving "latest"
checkpoint is rarely a snapshot point itself, so its delta
channels would silently reconstruct as empty (no error raised —
`get_delta_channel_history` simply returns no `seed`). Safe
options when the graph uses `DeltaChannel`:
* Walk back from each kept checkpoint and preserve every
ancestor (plus its `checkpoint_writes`) up to the nearest one
whose `channel_values` already contains a `_DeltaSnapshot` for
every `DeltaChannel`-backed key.
* Force a fresh snapshot on the kept checkpoint before deleting
ancestors — rewrite `channel_values[k] = _DeltaSnapshot(value)`
for each delta channel `k` (resolving `value` via the existing
ancestor walk first), then prune.
* Skip pruning threads whose graph uses `DeltaChannel` until one
of the above is implemented.
"""
raise NotImplementedError
@@ -461,6 +518,13 @@ class BaseCheckpointSaver(Generic[V]):
Args:
run_ids: The run IDs whose checkpoints should be deleted.
!!! warning "DeltaChannel"
See `delete_for_runs` — deleting rows a still-live thread's
`DeltaChannel` reconstruction depends on (writes between the
head and its nearest `_DeltaSnapshot` ancestor) will silently
corrupt that channel's state.
"""
raise NotImplementedError
@@ -474,6 +538,13 @@ class BaseCheckpointSaver(Generic[V]):
Args:
source_thread_id: The thread ID to copy from.
target_thread_id: The thread ID to copy to.
!!! warning "DeltaChannel"
See `copy_thread` — the copy must carry the complete parent
chain (or at least back to a `_DeltaSnapshot` ancestor for every
`DeltaChannel`) so the target thread can reconstruct delta
state.
"""
raise NotImplementedError
@@ -489,6 +560,13 @@ class BaseCheckpointSaver(Generic[V]):
thread_ids: The thread IDs to prune.
strategy: The pruning strategy. `"keep_latest"` retains only the most
recent checkpoint per namespace. `"delete"` removes all checkpoints.
!!! warning "DeltaChannel"
See `prune` for the full `DeltaChannel` caveat. In short:
`"keep_latest"` must not drop ancestor checkpoints / writes that
sit between the kept checkpoint and the nearest `_DeltaSnapshot`
ancestor, or delta channels will silently reconstruct as empty.
"""
raise NotImplementedError
@@ -497,6 +575,14 @@ class BaseCheckpointSaver(Generic[V]):
) -> Mapping[str, DeltaChannelHistory]:
"""Walk the parent chain returning per-channel writes + seed.
!!! warning "Beta"
This method is part of the `DeltaChannel` support surface and is
in beta. The signature, return shape (`DeltaChannelHistory`), and
interaction with `_DeltaSnapshot` blobs may change. Override at
your own risk; the default implementation will continue to work
against the public `BaseCheckpointSaver` contract.
For each requested channel, walks ancestors of the checkpoint
identified by `config` (following `parent_config`) and accumulates
`pending_writes` for that channel. The walk terminates per-channel
@@ -556,7 +642,13 @@ class BaseCheckpointSaver(Generic[V]):
async def aget_delta_channel_history(
self, *, config: RunnableConfig, channels: Sequence[str]
) -> Mapping[str, DeltaChannelHistory]:
"""Async version of `get_delta_channel_history`."""
"""Async version of `get_delta_channel_history`.
!!! warning "Beta"
This method is part of the `DeltaChannel` support surface and is
in beta. See `get_delta_channel_history` for caveats.
"""
if not channels:
return {}
collected_by_ch: dict[str, list[PendingWrite]] = {c: [] for c in channels}
@@ -588,6 +680,123 @@ class BaseCheckpointSaver(Generic[V]):
result[ch] = entry
return result
def get_delta_channel_keepset(
self,
*,
config: RunnableConfig,
channels: Sequence[str],
) -> set[str]:
"""Return ancestor checkpoint_ids that must survive deletion.
!!! warning "Beta"
This method is part of the `DeltaChannel` support surface and is
in beta. The signature may change while the delta-channel design
stabilizes.
Walks the parent chain from `config` backward, collecting visited
checkpoint_ids (inclusive of the target), and terminates per-channel
when that channel has a populated `channel_values[ch]` (a
`_DeltaSnapshot` blob or a pre-migration plain value). The returned
set is the minimum keep-set: every checkpoint_id whose removal would
break reconstruction of the listed channels at `config`.
Pass `channels=[]` to return just `{config.checkpoint_id}` — useful
for graphs that don't use `DeltaChannel`.
Compose this into custom `prune` / `delete_for_runs` / `copy_thread`::
keep = saver.get_delta_channel_keepset(
config=head_config, channels=delta_channels,
)
delete_rows_not_in(keep)
Note:
The default implementation here uses repeated `get_tuple` calls
to walk the parent chain. This is a basic reference implementation
suitable for low-frequency maintenance operations (prune, etc.).
Custom checkpointer backends may override with a more efficient
version tailored to their data model (e.g. a single SQL query
with a recursive CTE), but it is not required — the default works
correctly for any saver that implements `get_tuple`.
Args:
config: Configuration identifying the target checkpoint.
channels: Channel names whose delta history must be preserved.
Empty sequence means only the target checkpoint_id is kept.
Returns:
Set of checkpoint_ids that must not be deleted.
"""
target_tuple = self.get_tuple(config)
if target_tuple is None:
return set()
target_id = target_tuple.config["configurable"]["checkpoint_id"]
keep: set[str] = {target_id}
if not channels:
return keep
remaining: set[str] = set(channels)
for ch in list(remaining):
if ch in target_tuple.checkpoint["channel_values"]:
remaining.discard(ch)
if not remaining:
return keep
cursor_config: RunnableConfig | None = target_tuple.parent_config
while cursor_config is not None and remaining:
tup = self.get_tuple(cursor_config)
if tup is None:
break
cid = tup.config["configurable"]["checkpoint_id"]
keep.add(cid)
for ch in list(remaining):
if ch in tup.checkpoint["channel_values"]:
remaining.discard(ch)
if not remaining:
break
cursor_config = tup.parent_config
return keep
async def aget_delta_channel_keepset(
self,
*,
config: RunnableConfig,
channels: Sequence[str],
) -> set[str]:
"""Async version of `get_delta_channel_keepset`.
!!! warning "Beta"
This method is part of the `DeltaChannel` support surface and is
in beta. See `get_delta_channel_keepset` for full documentation.
"""
target_tuple = await self.aget_tuple(config)
if target_tuple is None:
return set()
target_id = target_tuple.config["configurable"]["checkpoint_id"]
keep: set[str] = {target_id}
if not channels:
return keep
remaining: set[str] = set(channels)
for ch in list(remaining):
if ch in target_tuple.checkpoint["channel_values"]:
remaining.discard(ch)
if not remaining:
return keep
cursor_config: RunnableConfig | None = target_tuple.parent_config
while cursor_config is not None and remaining:
tup = await self.aget_tuple(cursor_config)
if tup is None:
break
cid = tup.config["configurable"]["checkpoint_id"]
keep.add(cid)
for ch in list(remaining):
if ch in tup.checkpoint["channel_values"]:
remaining.discard(ch)
if not remaining:
break
cursor_config = tup.parent_config
return keep
def get_next_version(self, current: V | None, channel: None) -> V:
"""Generate the next version ID for a channel.
@@ -0,0 +1,35 @@
"""Run delta-channel conformance capabilities against InMemorySaver."""
from __future__ import annotations
import pytest
conformance = pytest.importorskip(
"langgraph.checkpoint.conformance",
reason="langgraph-checkpoint-conformance not installed",
)
@pytest.mark.asyncio
async def test_delta_channel_conformance():
from langgraph.checkpoint.conformance import validate
from langgraph.checkpoint.conformance.initializer import checkpointer_test
from langgraph.checkpoint.memory import InMemorySaver
@checkpointer_test(name="InMemorySaver")
async def mem_saver():
yield InMemorySaver()
report = await validate(
mem_saver,
capabilities={
"delta_channel_history",
"delta_channel_keepset",
"delta_channel_reconstruction",
},
)
for cap, result in report.results.items():
if result.passed is False:
details = "\n".join(result.failures or [])
pytest.fail(f"Capability {cap} failed:\n{details}")
+1 -1
View File
@@ -1 +1 @@
__version__ = "0.4.24"
__version__ = "0.4.25"
@@ -26,6 +26,15 @@ class DeltaChannel(Generic[Value], BaseChannel[Any, Any, Any]):
"""Reducer channel that stores only a sentinel in checkpoint blobs and
reconstructs state by replaying ancestor writes through the reducer.
!!! warning "Beta"
`DeltaChannel` is in beta. The API and on-disk representation may
change in future releases. Threads written with `DeltaChannel` today
are expected to remain readable, but the surrounding contract
(`BaseCheckpointSaver.get_delta_channel_history`, the
`_DeltaSnapshot` blob shape, the `delta_updates_since_snapshot`
metadata field) is not yet stable.
The reducer receives the current accumulated value and a batch of writes
in one call: `reducer(state, [write1, write2, ...]) -> new_state`.
+16 -41
View File
@@ -249,15 +249,17 @@ def _messages_delta_reducer(
) -> list[AnyMessage]:
"""**Experimental.** Batch reducer for use with `DeltaChannel`.
Provides full `add_messages` parity: dedup by ID, `RemoveMessage`
tombstoning, `REMOVE_ALL_MESSAGES` reset, `BaseMessageChunk` coercion,
and UUID assignment for ID-less messages — all in a single batched pass.
Processes all writes in one pass — dedup by ID, `RemoveMessage`
tombstoning — without calling `add_messages`.
This reducer is batching-invariant, as required by `DeltaChannel`:
`reducer(reducer(state, xs), ys) == reducer(state, xs + ys)`.
Raw dict / string / tuple inputs are coerced to typed `BaseMessage`
objects so that HTTP-driven graphs work without a separate coercion step.
objects so that HTTP-driven graphs work without a separate coercion
step. This is not full `add_messages` parity — `REMOVE_ALL_MESSAGES`,
unknown-id `RemoveMessage` errors, missing-id UUID assignment, and
`BaseMessageChunk` conversion are not handled here.
Example::
@@ -278,51 +280,24 @@ def _messages_delta_reducer(
flat.extend(w)
else:
flat.append(w)
# Steady state: the reducer's own output is already typed BaseMessages
# (never chunks), so skip convert_to_messages on the fast path.
# Steady state: the reducer's own output is already typed, so skip
# `convert_to_messages` on state when the first element is a BaseMessage.
# Only raw input (initial dicts, deserialized blobs) hits the slow path.
if state and isinstance(state[0], BaseMessage):
state_msgs = state
else:
state_msgs = cast(
"list[AnyMessage]",
[
message_chunk_to_message(cast(BaseMessageChunk, m))
for m in convert_to_messages(state)
],
)
# Coerce chunks to full messages — streaming nodes can emit BaseMessageChunk.
msgs = cast(
"list[AnyMessage]",
[
message_chunk_to_message(cast(BaseMessageChunk, m))
for m in convert_to_messages(flat)
],
)
state_msgs = cast("list[AnyMessage]", convert_to_messages(state))
msgs = cast("list[AnyMessage]", convert_to_messages(flat))
# REMOVE_ALL_MESSAGES resets everything; find the last sentinel and
# discard all state plus all writes before it.
remove_all_idx = None
for idx, m in enumerate(msgs):
if isinstance(m, RemoveMessage) and m.id == REMOVE_ALL_MESSAGES:
remove_all_idx = idx
if remove_all_idx is not None:
state_msgs = []
msgs = msgs[remove_all_idx + 1 :]
# Build index and assign missing IDs in one pass (parity with add_messages
# so that eviction and RemoveMessage tombstoning work on ID-less messages).
index: dict[str, int] = {}
for i, m in enumerate(state_msgs):
if m.id is None:
m.id = str(uuid.uuid4())
index[m.id] = i
index: dict[str, int] = {
m.id: i for i, m in enumerate(state_msgs) if m.id is not None
}
result: list[AnyMessage | None] = list(state_msgs)
for msg in msgs:
if msg.id is None:
msg.id = str(uuid.uuid4())
mid = msg.id
if isinstance(msg, RemoveMessage):
if mid is None:
result.append(msg)
elif isinstance(msg, RemoveMessage):
if mid in index:
result[index[mid]] = None
del index[mid]
+36 -58
View File
@@ -34,28 +34,23 @@ def empty_checkpoint() -> Checkpoint:
)
def _should_snapshot_delta(
name: str,
ch: DeltaChannel,
updates_since_snapshot: Mapping[str, int],
*,
force: bool,
) -> bool:
"""Decide whether `ch` should write a `_DeltaSnapshot` this step.
def delta_channels_to_snapshot(
channels: Mapping[str, BaseChannel],
counts: Mapping[str, int],
) -> set[str]:
"""Return the set of DeltaChannel names that should snapshot now.
Triggers:
* `force` — always snapshot (used by `durability="exit"`).
* Update-count: this channel has accumulated at least
`snapshot_frequency` updates since its last snapshot. The count
is supplied by the caller via `updates_since_snapshot[name]` and
is reset to `0` whenever a snapshot fires.
Version-format-independent: works for `int`, `float`, and `str`
versioning schemes alike.
A channel snapshots when its accumulated update count (since the last
snapshot) reaches or exceeds `snapshot_frequency`. This is a pure
predicate — no mutation.
"""
if force:
return True
return updates_since_snapshot.get(name, 0) >= ch.snapshot_frequency
return {
name
for name, ch in channels.items()
if isinstance(ch, DeltaChannel)
and ch.is_available()
and counts.get(name, 0) >= ch.snapshot_frequency
}
def create_checkpoint(
@@ -66,34 +61,19 @@ def create_checkpoint(
id: str | None = None,
updated_channels: set[str] | None = None,
get_next_version: GetNextVersion | None = None,
force_delta_snapshot: bool = False,
updates_since_snapshot: Mapping[str, int] | None = None,
new_updates_since_snapshot: dict[str, int] | None = None,
channels_to_snapshot: set[str] | None = None,
) -> Checkpoint:
"""Create a checkpoint for the given channels.
"""Build a new Checkpoint from the previous one and live channel state.
For each `DeltaChannel`, a `_DeltaSnapshot(value)` blob is written into
`channel_values[k]` when this channel has accumulated at least
`snapshot_frequency` updates since its last snapshot (counter supplied
via `updates_since_snapshot`). Otherwise the channel is omitted from
`channel_values`; its `channel_versions` entry still bumps so that the
saver tracks the channel and the ancestor walk can replay writes.
Snapshots are eager: even if the channel had no write this step, a
version bump is forced (via `get_next_version`) so `put()` includes
the channel in `new_versions` and stores the blob.
`force_delta_snapshot` ignores the cadence and always snapshots —
used by `durability="exit"` where intermediate writes are not stored
as ancestor `checkpoint_writes`.
If `new_updates_since_snapshot` is provided, the function resets the
counter to `0` for any channel that snapshotted this step. Counters
for channels that did not snapshot are left untouched (the caller is
responsible for incrementing them based on `updated_channels`).
For each name in `channels_to_snapshot`, a `_DeltaSnapshot(value)` blob
is written into `channel_values[k]`. Other delta channels are omitted
from `channel_values` — the ancestor walk reconstructs their state
from `checkpoint_writes`. Callers compute the set via
`delta_channels_to_snapshot(channels, counts)`; defaults to empty
(no snapshots) when not provided.
"""
ts = datetime.now(timezone.utc).isoformat()
counts = updates_since_snapshot or {}
channels_to_snapshot = channels_to_snapshot or set()
if channels is None:
values = checkpoint["channel_values"]
channel_versions = checkpoint["channel_versions"]
@@ -104,25 +84,23 @@ def create_checkpoint(
if k not in channel_versions:
continue
ch = channels[k]
if (
isinstance(ch, DeltaChannel)
and ch.is_available()
and _should_snapshot_delta(
k,
ch,
counts,
force=force_delta_snapshot,
)
):
# Eager snapshot: bump version if not already written this step
# so put() includes this channel in new_versions and stores blob.
if k in channels_to_snapshot:
# In exit mode, the snapshot decision is deferred to exit
# time (intermediate steps have do_checkpoint=False). The
# channel's count may have reached snapshot_frequency over
# several supersteps, but the LAST superstep may not have
# written to this channel. In that case apply_writes()
# (in _algo.py) didn't bump this channel's version, so
# saver.put() wouldn't include it in new_versions and
# the snapshot blob would be silently dropped. The manual
# bump below closes the gap. In sync/async durability this
# branch is effectively dead code (the step that pushes
# the count to freq always writes the channel).
if get_next_version is not None and (
updated_channels is None or k not in updated_channels
):
channel_versions[k] = get_next_version(channel_versions[k], None)
values[k] = _DeltaSnapshot(ch.get())
if new_updates_since_snapshot is not None:
new_updates_since_snapshot[k] = 0
else:
v = ch.checkpoint()
if v is not MISSING:
+216 -21
View File
@@ -100,6 +100,7 @@ from langgraph.pregel._checkpoint import (
channels_from_checkpoint,
copy_checkpoint,
create_checkpoint,
delta_channels_to_snapshot,
empty_checkpoint,
)
from langgraph.pregel._executor import (
@@ -194,8 +195,40 @@ class PregelLoop:
_migrate_checkpoint: Callable[[Checkpoint], None] | None
submit: Submit
channels: Mapping[str, BaseChannel]
# Only set on AsyncPregelLoop; sync loops keep this as None.
# Futures from `checkpointer.put_writes` calls that produced delta-channel
# writes. `_checkpointer_put_after_previous` drains this list (swap to a
# local `futs` then reset to `[]` and wait/gather) before putting the
# next checkpoint, so a checkpoint never becomes durable before the
# writes that produced it. Initialised to `[]` in both sync and async
# `__enter__`; stays `None` only when no checkpointer.
_delta_write_futs: list[Any] | None = None
# Exit-mode accumulator: every delta-channel write produced during this
# run (input writes from `_first` + per-superstep writes captured in
# `after_tick`). At exit, `_put_exit_delta_writes` filters out channels
# that will snapshot, then persists the rest under an anchor parent.
# `None` when not in exit mode (so the capture sites are no-ops).
# Each tuple is `(step, task_id, channel, value)` — `step` drives the
# synthetic step-prefixed task_id used to preserve chronological order
# under the saver's `ORDER BY task_id, idx` sorting.
_exit_delta_writes: list[tuple[int, str, str, Any]] | None = None
# The checkpoint_config that points at the parent loaded at `__enter__`
# (or the synthetic-empty checkpoint, on first run). We capture it
# eagerly because every `_put_checkpoint` advances `self.checkpoint_config`
# to the newly-saved checkpoint's id — by exit time the original parent
# config would otherwise be lost. `_put_exit_delta_writes` uses this:
# on resumed runs as the anchor for exit delta writes; on first runs
# to derive the lazy stub's config (its `checkpoint_id` is the
# synthetic-empty id we want the stub persisted under).
_initial_checkpoint_config: RunnableConfig
# True iff the saver actually returned a tuple at `__enter__`. False
# on the first-ever run for a thread (no parent persisted yet).
# `_put_exit_delta_writes` uses this to decide between anchoring on
# the existing parent (True) or creating a lazy stub (False).
_has_persisted_parent: bool = False
managed: ManagedValueMapping
checkpoint: Checkpoint
checkpoint_id_saved: str
@@ -637,6 +670,11 @@ class PregelLoop:
self._emit(
"values", map_output_values, self.output_keys, writes, self.channels
)
# capture delta-channel writes for exit-mode accumulator before clearing
if self._exit_delta_writes is not None:
for tid, ch, v in self.checkpoint_pending_writes:
if isinstance(self.specs.get(ch), DeltaChannel):
self._exit_delta_writes.append((self.step, tid, ch, v))
# clear pending writes
self.checkpoint_pending_writes.clear()
# only replay (re-execute) done tasks on the first tick
@@ -854,6 +892,27 @@ class PregelLoop:
self.checkpointer_get_next_version,
self.trigger_to_nodes,
)
# Input writes go through `apply_writes` directly (above) — they
# never enter `checkpoint_pending_writes`, so the after_tick
# capture site does not see them. In exit mode, capture them
# here so `_exit_delta_writes` includes the input's delta writes
# alongside per-superstep writes; otherwise the input would be
# lost on read (it's not in final_checkpoint.channel_values for
# sub-freq channels, and walks ignore target.pending_writes).
if self._exit_delta_writes is not None:
for c, v in input_writes:
if isinstance(self.specs.get(c), DeltaChannel):
self._exit_delta_writes.append((self.step, NULL_TASK_ID, c, v))
# Persist delta-channel input writes so sub-freq inputs are
# recoverable via ancestor walk (mirrors the Command input path).
if self.durability != "exit":
delta_input = [
(c, v)
for c, v in input_writes
if isinstance(self.specs.get(c), DeltaChannel)
]
if delta_input:
self.put_writes(NULL_TASK_ID, delta_input)
# save input checkpoint
self.updated_channels = updated_channels
self._put_checkpoint({"source": "input"})
@@ -905,36 +964,60 @@ class PregelLoop:
return updated_channels
def _put_checkpoint(self, metadata: CheckpointMetadata) -> None:
# assign step and parents
# `is` (object identity) — not `==`. Three of four call sites pass a
# fresh dict ({"source":"input"|"loop"|"fork"}); only
# `_suppress_interrupt`(will rename to _on_loop_exit soon)
# at exit reuses the existing `self.checkpoint_metadata` instance. So
# `metadata is self.checkpoint_metadata` is True only on the exit call,
# which is what we use to gate exit-only behaviour (skip count-bump,
# don't replace metadata). Could be replaced by an explicit
# `exiting: bool = False` parameter; left as-is to match the existing
# idiom in this file.
# TODO: replace with an explicit `exiting: bool = False` parameter.
exiting = metadata is self.checkpoint_metadata
if exiting and self.checkpoint["id"] == self.checkpoint_id_saved:
# checkpoint already saved
return
# Carry per-delta-channel update bookkeeping forward across
# supersteps. Capture from the OLD metadata before potentially
# replacing it with a fresh dict that wouldn't contain it. Then
# increment for any delta channel updated this step (so the count
# reflects "supersteps that wrote to this channel since last
# snapshot"). create_checkpoint will reset entries to 0 for any
# channel that fires a snapshot this step.
prev_counts = dict(
self.checkpoint_metadata.get("delta_updates_since_snapshot", {}) or {}
)
new_counts = dict(prev_counts)
if self.updated_channels:
for ch_name in self.updated_channels:
ch_obj = self.channels.get(ch_name)
if isinstance(ch_obj, DeltaChannel):
new_counts[ch_name] = new_counts.get(ch_name, 0) + 1
# Per-delta-channel update bookkeeping.
#
# `_put_checkpoint` is called once per superstep with a fresh
# metadata dict (source="input"|"loop"|"fork") — those are the
# intermediate calls that bump the count by +1 for each delta
# channel touched that step. In exit mode,
# `_suppress_interrupt`(will rename to _on_loop_exit soon)
# additionally calls `_put_checkpoint(self.checkpoint_metadata)` AT
# EXIT to commit the final checkpoint — this runs *after* the last
# intermediate call already counted the last superstep. So the
# exit call must NOT bump again or it would double-count the last
# superstep. (Sync/async durability does not call `_put_checkpoint`
# at exit, so the issue only surfaces in exit mode. force_delta_snapshot
# used to mask this latent bug by resetting every count to 0.)
if not exiting:
prev_counts = dict(
self.checkpoint_metadata.get("delta_updates_since_snapshot", {}) or {}
)
new_counts = dict(prev_counts)
if self.updated_channels:
for ch_name in self.updated_channels:
if isinstance(self.channels.get(ch_name), DeltaChannel):
new_counts[ch_name] = new_counts.get(ch_name, 0) + 1
metadata["step"] = self.step
metadata["parents"] = self.config[CONF].get(CONFIG_KEY_CHECKPOINT_MAP, {})
self.checkpoint_metadata = metadata
else:
new_counts = dict(
self.checkpoint_metadata.get("delta_updates_since_snapshot", {}) or {}
)
# do checkpoint?
do_checkpoint = self._checkpointer_put_after_previous is not None and (
exiting or self.durability != "exit"
)
# create new checkpoint
channels_to_snapshot = (
delta_channels_to_snapshot(self.channels, new_counts)
if do_checkpoint
else set()
)
self.checkpoint = create_checkpoint(
self.checkpoint,
self.channels if do_checkpoint else None,
@@ -944,10 +1027,10 @@ class PregelLoop:
get_next_version=self.checkpointer_get_next_version
if do_checkpoint
else None,
force_delta_snapshot=exiting and self.durability == "exit",
updates_since_snapshot=new_counts,
new_updates_since_snapshot=new_counts,
channels_to_snapshot=channels_to_snapshot,
)
for k in channels_to_snapshot:
new_counts[k] = 0
if new_counts:
self.checkpoint_metadata["delta_updates_since_snapshot"] = new_counts
elif "delta_updates_since_snapshot" in self.checkpoint_metadata:
@@ -1010,6 +1093,97 @@ class PregelLoop:
# increment step
self.step += 1
def _put_exit_delta_writes(self) -> None:
"""Stage stub + accumulated delta writes so final_checkpoint's put
waits on them (visibility invariant: both must be durable before
final_checkpoint becomes visible to readers).
Stub is created lazily — only when no persisted parent exists AND at
least one delta channel has writes that won't be snapshotted.
"""
if (
not self._exit_delta_writes
or self.checkpointer is None
or self._checkpointer_put_after_previous is None
or self.checkpointer_put_writes is None
):
return
counts = self.checkpoint_metadata.get("delta_updates_since_snapshot", {}) or {}
channels_to_snapshot = delta_channels_to_snapshot(self.channels, counts)
pending = [
(step, tid, ch, v)
for (step, tid, ch, v) in self._exit_delta_writes
if ch not in channels_to_snapshot
]
if not pending:
return
if self._has_persisted_parent:
# _initial_checkpoint_config's checkpoint_id is the saved parent's
# id (saver returned a real tuple at __enter__).
anchor_config = self._initial_checkpoint_config
else:
stub_cp = empty_checkpoint()
stub_cp["id"] = self.checkpoint_id_saved
stub_cp["ts"] = datetime.now(timezone.utc).isoformat()
# Stub has no parent (checkpoint_id=None in config).
stub_put_config = patch_configurable(
self._initial_checkpoint_config,
{CONFIG_KEY_CHECKPOINT_ID: None},
)
# Anchor config for put_writes: checkpoint_id = stub's id.
anchor_config = patch_configurable(
self._initial_checkpoint_config,
{CONFIG_KEY_CHECKPOINT_ID: stub_cp["id"]},
)
self._put_checkpoint_fut = self.submit(
self._checkpointer_put_after_previous,
getattr(self, "_put_checkpoint_fut", None),
stub_put_config,
stub_cp,
{"step": -2},
{},
)
# Set checkpoint_config so final_checkpoint's _put_checkpoint
# sees the stub as its parent.
self.checkpoint_config = anchor_config
# Step-prefixed synthetic task_id preserves chronological superstep
# order under the saver's ORDER BY task_id, idx sorting.
grouped: dict[tuple[int, str], list[tuple[str, Any]]] = {}
for step, tid, ch, v in pending:
grouped.setdefault((step, tid), []).append((ch, v))
anchor_write_config = patch_configurable(
anchor_config,
{
CONFIG_KEY_CHECKPOINT_NS: self.config[CONF].get(
CONFIG_KEY_CHECKPOINT_NS, ""
),
CONFIG_KEY_CHECKPOINT_ID: anchor_config[CONF][CONFIG_KEY_CHECKPOINT_ID],
},
)
for (step, tid), entries in grouped.items():
synth_tid = f"{step:08d}-{tid}"
if self.checkpointer_put_writes_accepts_task_path:
fut = self.submit(
self.checkpointer_put_writes,
anchor_write_config,
entries,
synth_tid,
"",
)
else:
fut = self.submit(
self.checkpointer_put_writes,
anchor_write_config,
entries,
synth_tid,
)
if self._delta_write_futs is not None:
self._delta_write_futs.append(fut)
def _suppress_interrupt(
self,
exc_type: type[BaseException] | None,
@@ -1025,6 +1199,7 @@ class PregelLoop:
# or a nested graph with checkpointer=True
or all(NS_END not in part for part in self.checkpoint_ns)
):
self._put_exit_delta_writes()
self._put_checkpoint(self.checkpoint_metadata)
self._put_pending_writes()
# suppress interrupt
@@ -1230,6 +1405,9 @@ class SyncPregelLoop(PregelLoop, AbstractContextManager):
metadata: CheckpointMetadata,
new_versions: ChannelVersions,
) -> RunnableConfig:
if self._delta_write_futs:
futs, self._delta_write_futs = self._delta_write_futs, []
concurrent.futures.wait(futs)
try:
if prev is not None:
prev.result()
@@ -1347,6 +1525,10 @@ class SyncPregelLoop(PregelLoop, AbstractContextManager):
# graph/thread. Returns None on first invocation.
saved = self.checkpointer.get_tuple(self.checkpoint_config)
# Capture before the synthetic-empty fallback below overwrites `saved`.
# `_put_exit_delta_writes` uses this on first run (no persisted parent)
# to lazy-create a stub instead of anchoring delta writes on a parent.
self._has_persisted_parent = saved is not None
if saved is None:
saved = CheckpointTuple(
self.checkpoint_config, empty_checkpoint(), {"step": -2}, None, []
@@ -1362,6 +1544,7 @@ class SyncPregelLoop(PregelLoop, AbstractContextManager):
**saved.config.get(CONF, {}),
},
}
self._initial_checkpoint_config = self.checkpoint_config
self.prev_checkpoint_config = saved.parent_config
self.checkpoint_id_saved = saved.checkpoint["id"]
self.checkpoint = saved.checkpoint
@@ -1371,6 +1554,10 @@ class SyncPregelLoop(PregelLoop, AbstractContextManager):
if saved.pending_writes is not None
else []
)
self._delta_write_futs = []
self._exit_delta_writes = (
[] if self.durability == "exit" and self.checkpointer is not None else None
)
self.submit = self.stack.enter_context(BackgroundExecutor(self.config))
self.channels, self.managed = channels_from_checkpoint(
self.specs,
@@ -1596,6 +1783,10 @@ class AsyncPregelLoop(PregelLoop, AbstractAsyncContextManager):
# graph/thread. Returns None on first invocation.
saved = await self.checkpointer.aget_tuple(self.checkpoint_config)
# Capture before the synthetic-empty fallback below overwrites `saved`.
# `_put_exit_delta_writes` uses this on first run (no persisted parent)
# to lazy-create a stub instead of anchoring delta writes on a parent.
self._has_persisted_parent = saved is not None
if saved is None:
saved = CheckpointTuple(
self.checkpoint_config, empty_checkpoint(), {"step": -2}, None, []
@@ -1611,6 +1802,7 @@ class AsyncPregelLoop(PregelLoop, AbstractAsyncContextManager):
**saved.config.get(CONF, {}),
},
}
self._initial_checkpoint_config = self.checkpoint_config
self.prev_checkpoint_config = saved.parent_config
self.checkpoint_id_saved = saved.checkpoint["id"]
self.checkpoint = saved.checkpoint
@@ -1621,6 +1813,9 @@ class AsyncPregelLoop(PregelLoop, AbstractAsyncContextManager):
else []
)
self._delta_write_futs = []
self._exit_delta_writes = (
[] if self.durability == "exit" and self.checkpointer is not None else None
)
self.submit = await self.stack.enter_async_context(
AsyncBackgroundExecutor(self.config)
)
+2 -62
View File
@@ -3,12 +3,7 @@ from collections.abc import Sequence
from typing import Annotated
import pytest
from langchain_core.messages import (
AIMessage,
AIMessageChunk,
HumanMessage,
RemoveMessage,
)
from langchain_core.messages import AIMessage, HumanMessage, RemoveMessage
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.checkpoint.serde.types import _DeltaSnapshot
from typing_extensions import NotRequired, TypedDict
@@ -21,7 +16,7 @@ from langgraph.channels.topic import Topic
from langgraph.channels.untracked_value import UntrackedValue
from langgraph.errors import EmptyChannelError, InvalidUpdateError
from langgraph.graph import START, StateGraph
from langgraph.graph.message import REMOVE_ALL_MESSAGES, _messages_delta_reducer
from langgraph.graph.message import _messages_delta_reducer
from langgraph.graph.state import _get_channel
from langgraph.types import Overwrite
@@ -289,61 +284,6 @@ def test_messages_delta_reducer_tuple_write_is_one_message() -> None:
assert result[0].content == "hi"
def test_messages_delta_reducer_assigns_uuid_to_id_less_messages() -> None:
"""Messages without IDs get UUIDs assigned, matching add_messages behavior.
Without UUID assignment, RemoveMessage tombstoning fails on messages that
were created without explicit IDs.
"""
m1 = HumanMessage(content="hi")
m2 = AIMessage(content="hello")
assert m1.id is None
assert m2.id is None
result = _messages_delta_reducer([], [[m1, m2]])
assert len(result) == 2
assert result[0].id is not None
assert result[1].id is not None
# RemoveMessage tombstoning must work on the now-assigned IDs.
result2 = _messages_delta_reducer(result, [RemoveMessage(id=result[1].id)])
assert len(result2) == 1
assert result2[0].content == "hi"
def test_messages_delta_reducer_remove_all_messages() -> None:
"""REMOVE_ALL_MESSAGES sentinel clears all state and preceding writes."""
state = [HumanMessage(content="old", id="h1"), AIMessage(content="prior", id="a1")]
# Sentinel mid-batch: everything before it (including state) is discarded.
result = _messages_delta_reducer(
state,
[
[
RemoveMessage(id=REMOVE_ALL_MESSAGES),
HumanMessage(content="fresh", id="h2"),
]
],
)
assert len(result) == 1
assert result[0].content == "fresh"
# Batching-invariant: split across two calls must equal one combined call.
step1 = _messages_delta_reducer(state, [[RemoveMessage(id=REMOVE_ALL_MESSAGES)]])
step2 = _messages_delta_reducer(step1, [[HumanMessage(content="fresh", id="h2")]])
assert step2 == result
def test_messages_delta_reducer_coerces_message_chunks() -> None:
"""BaseMessageChunk writes are coerced to full messages."""
chunk = AIMessageChunk(content="hello", id="a1")
result = _messages_delta_reducer([], [[chunk]])
assert len(result) == 1
assert not isinstance(result[0], AIMessageChunk)
assert result[0].content == "hello"
assert result[0].id == "a1"
def test_delta_channel_checkpoint_returns_missing() -> None:
"""checkpoint() always returns MISSING regardless of state.
@@ -0,0 +1,365 @@
"""Tests for exit-mode delta channel persistence redesign.
Validates that `durability="exit"` correctly persists delta-channel writes
using count-based snapshot decisions (rather than force-snapshotting every
channel), lazy stub creation when no parent exists, and proper read-path
reconstruction via ancestor walks.
"""
from typing import Annotated, Any
import pytest
from langchain_core.messages import AIMessage, HumanMessage
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.checkpoint.serde.types import _DeltaSnapshot
from typing_extensions import TypedDict
from langgraph.channels.delta import DeltaChannel
from langgraph.graph import START, StateGraph
from langgraph.graph.message import _messages_delta_reducer
pytestmark = pytest.mark.anyio
def _build_graph(
checkpointer: InMemorySaver,
*,
freq: int = 1000,
) -> Any:
channel = DeltaChannel(_messages_delta_reducer, snapshot_frequency=freq)
# Functional TypedDict form: class form can't reference `channel` (a
# local variable) inside Annotated due to forward-ref evaluation rules.
State = TypedDict("State", {"messages": Annotated[list, channel]}) # type: ignore[call-overload] # noqa: UP013
def respond(state: dict) -> dict:
i = len(state["messages"])
return {"messages": [AIMessage(content=f"reply-{i}", id=f"ai{i}")]}
builder = StateGraph(State)
builder.add_node("respond", respond)
builder.add_edge(START, "respond")
return builder.compile(checkpointer=checkpointer)
# ---------------------------------------------------------------------------
# 8a. Write-path / structural tests
# ---------------------------------------------------------------------------
async def test_exit_first_run_no_delta_writes() -> None:
"""Graph with delta channel invoked with input that doesn't touch it.
Only one checkpoint row, no stub."""
State = TypedDict( # noqa: UP013
"State",
{
"messages": Annotated[list, DeltaChannel(_messages_delta_reducer)],
"value": str,
},
) # type: ignore[call-overload]
def noop(state: dict) -> dict:
return {"value": "done"}
saver = InMemorySaver()
builder = StateGraph(State)
builder.add_node("noop", noop)
builder.add_edge(START, "noop")
graph = builder.compile(checkpointer=saver)
config = {"configurable": {"thread_id": "no-delta-writes"}}
graph.invoke({"value": "start"}, config, durability="exit")
checkpoints = list(saver.list(config))
assert len(checkpoints) == 1
stubs = [t for t in checkpoints if t.metadata.get("step") == -2]
assert len(stubs) == 0
async def test_exit_first_run_all_snapshot() -> None:
"""snapshot_frequency=1 forces every channel to snapshot.
No stub needed; final_checkpoint has _DeltaSnapshot."""
saver = InMemorySaver()
graph = _build_graph(saver, freq=1)
config = {"configurable": {"thread_id": "all-snapshot"}}
result = graph.invoke(
{"messages": [HumanMessage(content="hi", id="h1")]},
config,
durability="exit",
)
assert len(result["messages"]) == 2
checkpoints = list(saver.list(config))
stubs = [t for t in checkpoints if t.metadata.get("step") == -2]
assert len(stubs) == 0
head = saver.get_tuple(config)
assert head is not None
assert isinstance(head.checkpoint["channel_values"].get("messages"), _DeltaSnapshot)
state = graph.get_state(config)
assert [m.content for m in state.values["messages"]] == ["hi", "reply-1"]
async def test_exit_first_run_sub_freq_with_writes() -> None:
"""First run with default snapshot_frequency (1000), writes below threshold.
A stub is created; writes are anchored under it; get_state reconstructs."""
saver = InMemorySaver()
graph = _build_graph(saver)
config = {"configurable": {"thread_id": "sub-freq-first"}}
result = graph.invoke(
{"messages": [HumanMessage(content="hello", id="h1")]},
config,
durability="exit",
)
assert [m.content for m in result["messages"]] == ["hello", "reply-1"]
checkpoints = list(saver.list(config))
stubs = [t for t in checkpoints if t.metadata.get("step") == -2]
assert len(stubs) == 1, f"Expected 1 stub, got {len(stubs)}"
head = saver.get_tuple(config)
assert head is not None
assert "messages" not in head.checkpoint["channel_values"]
assert "messages" in head.checkpoint["channel_versions"]
state = graph.get_state(config)
assert [m.content for m in state.values["messages"]] == ["hello", "reply-1"]
async def test_exit_resumed_run_sub_freq() -> None:
"""Two consecutive exit runs. Second run anchors on the first's
final_checkpoint (no new stub). Ordering preserved."""
saver = InMemorySaver()
graph = _build_graph(saver)
config = {"configurable": {"thread_id": "resumed-sub-freq"}}
graph.invoke(
{"messages": [HumanMessage(content="msg1", id="h1")]},
config,
durability="exit",
)
graph.invoke(
{"messages": [HumanMessage(content="msg2", id="h2")]},
config,
durability="exit",
)
checkpoints = list(saver.list(config))
stubs = [t for t in checkpoints if t.metadata.get("step") == -2]
assert len(stubs) == 1
state = graph.get_state(config)
contents = [m.content for m in state.values["messages"]]
assert len(contents) == 4
assert contents[0] == "msg1"
assert contents[2] == "msg2"
assert contents[0:4:2] == ["msg1", "msg2"]
async def test_exit_count_parity_sync_vs_exit() -> None:
"""Sync and exit durability produce the same delta_updates_since_snapshot
after an equivalent run."""
for durability in ("sync", "exit"):
saver = InMemorySaver()
graph = _build_graph(saver)
config = {"configurable": {"thread_id": f"parity-{durability}"}}
graph.invoke(
{"messages": [HumanMessage(content="hi", id="h1")]},
config,
durability=durability,
)
head = saver.get_tuple(config)
assert head is not None
counts = head.metadata.get("delta_updates_since_snapshot", {})
assert counts.get("messages") == 2, (
f"durability={durability}: expected count=2, got {counts}"
)
async def test_exit_snapshot_fires_at_frequency() -> None:
"""With snapshot_frequency=3, after 3 exit runs (each incrementing count
by 2: input + superstep), the 2nd run hits count=4>=3, triggering snapshot.
After that run, count resets to 0 and channel_values has _DeltaSnapshot."""
saver = InMemorySaver()
graph = _build_graph(saver, freq=3)
config = {"configurable": {"thread_id": "snapshot-at-freq"}}
graph.invoke(
{"messages": [HumanMessage(content="m1", id="h1")]},
config,
durability="exit",
)
head = saver.get_tuple(config)
assert head is not None
count1 = head.metadata.get("delta_updates_since_snapshot", {}).get("messages", 0)
assert count1 == 2
graph.invoke(
{"messages": [HumanMessage(content="m2", id="h2")]},
config,
durability="exit",
)
head = saver.get_tuple(config)
assert head is not None
count2 = head.metadata.get("delta_updates_since_snapshot", {}).get("messages", 0)
assert count2 == 0, f"Expected reset to 0 after snapshot, got {count2}"
assert isinstance(head.checkpoint["channel_values"].get("messages"), _DeltaSnapshot)
async def test_exit_mixed_snapshot_and_non_snapshot() -> None:
"""One delta channel at freq=1 (always snapshot) and one at freq=1000
(never snapshot within this test). Verify correct behavior for both."""
fast_ch = DeltaChannel(_messages_delta_reducer, snapshot_frequency=1)
slow_ch = DeltaChannel(_messages_delta_reducer, snapshot_frequency=1000)
State = TypedDict( # noqa: UP013
"State",
{"fast": Annotated[list, fast_ch], "slow": Annotated[list, slow_ch]},
) # type: ignore[call-overload]
def respond(state: dict) -> dict:
return {
"fast": [AIMessage(content="fast-reply", id="f1")],
"slow": [AIMessage(content="slow-reply", id="s1")],
}
saver = InMemorySaver()
builder = StateGraph(State)
builder.add_node("respond", respond)
builder.add_edge(START, "respond")
graph = builder.compile(checkpointer=saver)
config = {"configurable": {"thread_id": "mixed-freq"}}
graph.invoke(
{
"fast": [HumanMessage(content="fast-in", id="fi")],
"slow": [HumanMessage(content="slow-in", id="si")],
},
config,
durability="exit",
)
head = saver.get_tuple(config)
assert head is not None
assert isinstance(head.checkpoint["channel_values"].get("fast"), _DeltaSnapshot)
assert "slow" not in head.checkpoint["channel_values"]
state = graph.get_state(config)
assert [m.content for m in state.values["fast"]] == ["fast-in", "fast-reply"]
assert [m.content for m in state.values["slow"]] == ["slow-in", "slow-reply"]
# ---------------------------------------------------------------------------
# 8b. Read-path tests
# ---------------------------------------------------------------------------
async def test_exit_multi_run_replay_chain() -> None:
"""K=4 consecutive exit runs, each adding a message. After each run,
get_state returns all messages in chronological order."""
saver = InMemorySaver()
graph = _build_graph(saver)
config = {"configurable": {"thread_id": "replay-chain"}}
for i in range(4):
graph.invoke(
{"messages": [HumanMessage(content=f"user-{i}", id=f"h{i}")]},
config,
durability="exit",
)
state = graph.get_state(config)
contents = [m.content for m in state.values["messages"]]
user_msgs = [c for c in contents if c.startswith("user-")]
assert user_msgs == [f"user-{j}" for j in range(i + 1)], (
f"After run {i}: user messages out of order: {user_msgs}"
)
assert len(contents) == (i + 1) * 2
async def test_exit_metadata_round_trip() -> None:
"""K=5 consecutive exit runs with snapshot_frequency=5. Verify metadata
delta_updates_since_snapshot increments correctly across runs."""
freq = 5
saver = InMemorySaver()
graph = _build_graph(saver, freq=freq)
config = {"configurable": {"thread_id": "metadata-rt"}}
for i in range(1, 6):
graph.invoke(
{"messages": [HumanMessage(content=f"m{i}", id=f"h{i}")]},
config,
durability="exit",
)
head = saver.get_tuple(config)
assert head is not None
count = head.metadata.get("delta_updates_since_snapshot", {}).get("messages", 0)
cumulative = i * 2
if cumulative >= freq:
assert count == 0 or count == cumulative % freq or count < freq, (
f"After run {i}: count={count} should have reset or be partial"
)
else:
assert count == cumulative, (
f"After run {i}: expected {cumulative}, got {count}"
)
async def test_exit_mixed_durability_round_trip() -> None:
"""Alternate sync and exit durability; verify counts stay monotonic
and state accumulates correctly."""
saver = InMemorySaver()
graph = _build_graph(saver)
config = {"configurable": {"thread_id": "mixed-durability"}}
for i, dur in enumerate(["sync", "exit", "sync", "exit"]):
graph.invoke(
{"messages": [HumanMessage(content=f"msg-{i}", id=f"h{i}")]},
config,
durability=dur,
)
state = graph.get_state(config)
contents = [m.content for m in state.values["messages"]]
user_msgs = [c for c in contents if c.startswith("msg-")]
assert user_msgs == [f"msg-{j}" for j in range(i + 1)], (
f"After run {i} (durability={dur}): {user_msgs}"
)
assert len(contents) == (i + 1) * 2
async def test_exit_snapshot_then_tail_deltas() -> None:
"""Run 1 forces snapshot (freq=1). Run 2 at freq=1000 adds more writes
that don't snapshot. Reading after run 2 must combine the snapshot seed
with the tail deltas."""
saver = InMemorySaver()
graph1 = _build_graph(saver, freq=1)
config = {"configurable": {"thread_id": "snapshot-then-tail"}}
graph1.invoke(
{"messages": [HumanMessage(content="seed-msg", id="h1")]},
config,
durability="exit",
)
head = saver.get_tuple(config)
assert head is not None
assert isinstance(head.checkpoint["channel_values"].get("messages"), _DeltaSnapshot)
graph2 = _build_graph(saver, freq=1000)
graph2.invoke(
{"messages": [HumanMessage(content="tail-msg", id="h2")]},
config,
durability="exit",
)
state = graph2.get_state(config)
contents = [m.content for m in state.values["messages"]]
assert "seed-msg" in contents
assert "tail-msg" in contents
assert contents.index("seed-msg") < contents.index("tail-msg")
+2 -2
View File
@@ -1849,14 +1849,14 @@ dev = [
{ name = "pytest-watch" },
{ name = "ruff", specifier = "==0.15.12" },
{ name = "starlette" },
{ name = "ty", specifier = "==0.0.23" },
{ name = "ty", specifier = "==0.0.33" },
]
lint = [
{ name = "codespell" },
{ name = "mypy", specifier = "==1.20.2" },
{ name = "ruff", specifier = "==0.15.12" },
{ name = "starlette" },
{ name = "ty", specifier = "==0.0.23" },
{ name = "ty", specifier = "==0.0.33" },
]
test = [
{ name = "pytest" },
+2 -2
View File
@@ -618,14 +618,14 @@ dev = [
{ name = "pytest-watch" },
{ name = "ruff", specifier = "==0.15.12" },
{ name = "starlette" },
{ name = "ty", specifier = "==0.0.23" },
{ name = "ty", specifier = "==0.0.33" },
]
lint = [
{ name = "codespell" },
{ name = "mypy", specifier = "==1.20.2" },
{ name = "ruff", specifier = "==0.15.12" },
{ name = "starlette" },
{ name = "ty", specifier = "==0.0.23" },
{ name = "ty", specifier = "==0.0.33" },
]
test = [
{ name = "pytest" },
+3 -3
View File
@@ -110,7 +110,7 @@ def get_client(
if url is None:
url = "http://api"
if os.environ.get("__LANGGRAPH_DEFER_LOOPBACK_TRANSPORT") == "true":
transport = get_asgi_transport()(app=None, root_path="/noauth") # type: ignore[invalid-argument-type]
transport = get_asgi_transport()(app=None, root_path="/noauth") # ty: ignore[invalid-argument-type]
_registered_transports.append(transport)
else:
try:
@@ -122,7 +122,7 @@ def get_client(
"Failed to connect to in-process LangGraph server. Deferring configuration.",
exc_info=True,
)
transport = get_asgi_transport()(app=None, root_path="/noauth") # type: ignore[invalid-argument-type]
transport = get_asgi_transport()(app=None, root_path="/noauth") # ty: ignore[invalid-argument-type]
_registered_transports.append(transport)
if transport is None:
@@ -131,7 +131,7 @@ def get_client(
base_url=url,
transport=transport,
timeout=(
httpx.Timeout(timeout) # type: ignore[arg-type]
httpx.Timeout(timeout) # ty: ignore[invalid-argument-type]
if timeout is not None
else httpx.Timeout(connect=5, read=300, write=300, pool=5)
),
+1 -1
View File
@@ -49,7 +49,7 @@ async def _wrap_stream_v2(
async for part in raw:
v2 = _sse_to_v2_dict(part.event, part.data)
if v2 is not None:
yield v2
yield v2 # ty: ignore[invalid-yield]
class RunsClient:
@@ -144,8 +144,9 @@ def _resolve_timezone(tz: str | tzinfo | ZoneInfo | None) -> str | None:
return tz
if isinstance(tz, tzinfo):
# ZoneInfo objects have a .key attribute with the IANA name
if hasattr(tz, "key"):
return tz.key # type: ignore[union-attr]
key = getattr(tz, "key", None)
if isinstance(key, str):
return key
# Fall back to tzname for fixed-offset timezones like datetime.timezone.utc
name = tz.tzname(None)
if name is not None:
@@ -209,7 +210,7 @@ def configure_loopback_transports(app: Any) -> None:
@functools.lru_cache(maxsize=1)
def get_asgi_transport() -> type[httpx.ASGITransport]:
try:
from langgraph_api import asgi_transport # type: ignore[unresolved-import]
from langgraph_api import asgi_transport # ty: ignore[unresolved-import]
return asgi_transport.ASGITransport
except ImportError:
+1 -1
View File
@@ -77,7 +77,7 @@ def get_sync_client(
base_url=url,
transport=transport,
timeout=(
httpx.Timeout(timeout) # type: ignore[arg-type]
httpx.Timeout(timeout) # ty: ignore[invalid-argument-type]
if timeout is not None
else httpx.Timeout(connect=5, read=300, write=300, pool=5)
),
+1 -1
View File
@@ -49,7 +49,7 @@ def _wrap_stream_v2_sync(
for part in raw:
v2 = _sse_to_v2_dict(part.event, part.data)
if v2 is not None:
yield v2
yield v2 # ty: ignore[invalid-yield]
class SyncRunsClient:
+8 -5
View File
@@ -16,10 +16,10 @@ T = TypeVar("T")
CacheStatus = Literal["miss", "fresh", "stale", "expired"]
try:
from langgraph_api.cache import ( # type: ignore[unresolved-import]
from langgraph_api.cache import ( # ty: ignore[unresolved-import]
cache_get as _cache_get,
)
from langgraph_api.cache import ( # type: ignore[unresolved-import]
from langgraph_api.cache import ( # ty: ignore[unresolved-import]
cache_set as _cache_set,
)
except ImportError:
@@ -28,8 +28,8 @@ except ImportError:
try:
from langgraph_api.cache import SWRResult # type: ignore[unresolved-import]
from langgraph_api.cache import swr as _api_swr # type: ignore[unresolved-import]
from langgraph_api.cache import SWRResult # ty: ignore[unresolved-import]
from langgraph_api.cache import swr as _api_swr # ty: ignore[unresolved-import]
except ImportError:
_api_swr = None
@@ -40,7 +40,10 @@ except ImportError:
value: T
status: CacheStatus
async def mutate(self, value: T = ...) -> T: # type: ignore[assignment]
async def mutate(
self,
value: T = ..., # ty: ignore[invalid-parameter-default]
) -> T: # ty: ignore[empty-body]
"""Update or revalidate the cached value."""
...
+1 -1
View File
@@ -37,7 +37,7 @@ class APIError(httpx.HTTPStatusError, LangGraphError):
req = response_or_request
response = None
httpx.HTTPStatusError.__init__(self, message, request=req, response=response) # type: ignore[arg-type]
httpx.HTTPStatusError.__init__(self, message, request=req, response=response) # ty: ignore[invalid-argument-type]
LangGraphError.__init__(self, message)
self.request = req
+1 -1
View File
@@ -156,7 +156,7 @@ class _ExecutionRuntime(_ServerRuntimeBase[ContextT], Generic[ContextT]):
This API is in beta and may change in future releases.
"""
context: ContextT = field(default=None) # type: ignore[assignment]
context: ContextT = field(default=None) # ty: ignore[invalid-assignment]
"""The graph run context, typed by the graph's `context_schema`.
Only available during `threads.create_run`.
+3 -3
View File
@@ -55,7 +55,7 @@ class BytesLineDecoder:
# Include any existing buffer in the first portion of the
# splitlines result.
self.buffer.extend(lines[0])
lines = cast(list[BytesLike], [self.buffer, *lines[1:]])
lines = [self.buffer, *lines[1:]]
self.buffer = bytearray()
if not trailing_newline:
@@ -69,7 +69,7 @@ class BytesLineDecoder:
if not self.buffer and not self.trailing_cr:
return []
lines = [self.buffer]
lines: list[BytesLike] = [self.buffer]
self.buffer = bytearray()
self.trailing_cr = False
return lines
@@ -102,7 +102,7 @@ class SSEDecoder:
sse = StreamPart(
event=self._event,
data=orjson.loads(self._data) if self._data else None, # type: ignore[invalid-argument-type]
data=orjson.loads(self._data) if self._data else None, # ty: ignore[invalid-argument-type]
id=self.last_event_id,
)
+1 -1
View File
@@ -33,7 +33,7 @@ lint = [
"ruff==0.15.12",
"codespell",
"mypy==1.20.2",
"ty==0.0.23",
"ty==0.0.33",
"starlette",
]
dev = [
+2 -2
View File
@@ -388,7 +388,7 @@ async def test_async_stream_v2_client_side_conversion() -> None:
event="values", data={"messages": [{"role": "user", "content": "hi"}]}
)
yield StreamPart(event="updates|sub:abc", data={"node": {"out": 1}})
yield StreamPart(event="end", data=None) # type: ignore[arg-type]
yield StreamPart(event="end", data=None) # ty: ignore[invalid-argument-type]
parts: list[StreamPartV2] = [part async for part in _wrap_stream_v2(mock_stream())]
assert len(parts) == 3
@@ -420,7 +420,7 @@ def test_sync_stream_v2_client_side_conversion() -> None:
def mock_stream() -> Any:
yield StreamPart(event="metadata", data={"run_id": "r1"})
yield StreamPart(event="values", data={"state": "full"})
yield StreamPart(event="end", data=None) # type: ignore[arg-type]
yield StreamPart(event="end", data=None) # ty: ignore[invalid-argument-type]
parts: list[StreamPartV2] = list(_wrap_stream_v2_sync(mock_stream()))
assert len(parts) == 2
+1 -1
View File
@@ -67,6 +67,6 @@ class TestHandlerValidation:
with pytest.raises(TypeError, match="must accept exactly 2 parameters"):
@encryption.encrypt.blob # type: ignore[arg-type]
@encryption.encrypt.blob # ty: ignore[invalid-argument-type]
async def wrong_params(ctx):
return ctx
+20 -20
View File
@@ -533,14 +533,14 @@ dev = [
{ name = "pytest-watch" },
{ name = "ruff", specifier = "==0.15.12" },
{ name = "starlette" },
{ name = "ty", specifier = "==0.0.23" },
{ name = "ty", specifier = "==0.0.33" },
]
lint = [
{ name = "codespell" },
{ name = "mypy", specifier = "==1.20.2" },
{ name = "ruff", specifier = "==0.15.12" },
{ name = "starlette" },
{ name = "ty", specifier = "==0.0.23" },
{ name = "ty", specifier = "==0.0.33" },
]
test = [
{ name = "pytest" },
@@ -1275,26 +1275,26 @@ wheels = [
[[package]]
name = "ty"
version = "0.0.23"
version = "0.0.33"
source = { registry = "https://pypi.org/simple" }
sdist = { url = "https://files.pythonhosted.org/packages/75/ba/d3c998ff4cf6b5d75b39356db55fe1b7caceecc522b9586174e6a5dee6f7/ty-0.0.23.tar.gz", hash = "sha256:5fb05db58f202af366f80ef70f806e48f5237807fe424ec787c9f289e3f3a4ef", size = 5341461, upload-time = "2026-03-13T12:34:23.125Z" }
sdist = { url = "https://files.pythonhosted.org/packages/84/44/9478c50c266826c1bf30d1692e589755bffa8f1c0a3eb7af8a346c255991/ty-0.0.33.tar.gz", hash = "sha256:46d63bda07403322cb6c28ccfdd5536be916e13df725c29f7ccd0a21f06bd9e8", size = 5559373, upload-time = "2026-04-28T10:45:13.18Z" }
wheels = [
{ url = "https://files.pythonhosted.org/packages/f4/21/aab32603dfdfacd4819e52fa8c6074e7bd578218a5142729452fc6a62db6/ty-0.0.23-py3-none-linux_armv6l.whl", hash = "sha256:e810eef1a5f1cfc0731a58af8d2f334906a96835829767aed00026f1334a8dd7", size = 10329096, upload-time = "2026-03-13T12:34:09.432Z" },
{ url = "https://files.pythonhosted.org/packages/9f/a9/dd3287a82dce3df546ec560296208d4905dcf06346b6e18c2f3c63523bd1/ty-0.0.23-py3-none-macosx_10_12_x86_64.whl", hash = "sha256:e43d36bd89a151ddcad01acaeff7dcc507cb73ff164c1878d2d11549d39a061c", size = 10156631, upload-time = "2026-03-13T12:34:53.122Z" },
{ url = "https://files.pythonhosted.org/packages/0f/01/3f25909b02fac29bb0a62b2251f8d62e65d697781ffa4cf6b47a4c075c85/ty-0.0.23-py3-none-macosx_11_0_arm64.whl", hash = "sha256:bd6a340969577b4645f231572c4e46012acba2d10d4c0c6570fe1ab74e76ae00", size = 9653211, upload-time = "2026-03-13T12:34:15.049Z" },
{ url = "https://files.pythonhosted.org/packages/d5/60/bfc0479572a6f4b90501c869635faf8d84c8c68ffc5dd87d04f049affabc/ty-0.0.23-py3-none-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:341441783e626eeb7b1ec2160432956aed5734932ab2d1c26f94d0c98b229937", size = 10156143, upload-time = "2026-03-13T12:34:34.468Z" },
{ url = "https://files.pythonhosted.org/packages/3a/81/8a93e923535a340f54bea20ff196f6b2787782b2f2f399bd191c4bc132d6/ty-0.0.23-py3-none-manylinux_2_17_armv7l.manylinux2014_armv7l.whl", hash = "sha256:8ce1dc66c26d4167e2c78d12fa870ef5a7ec9cc344d2baaa6243297cfa88bd52", size = 10136632, upload-time = "2026-03-13T12:34:28.832Z" },
{ url = "https://files.pythonhosted.org/packages/da/cb/2ac81c850c58acc9f976814404d28389c9c1c939676e32287b9cff61381e/ty-0.0.23-py3-none-manylinux_2_17_i686.manylinux2014_i686.whl", hash = "sha256:bae1e7a294bf8528836f7617dc5c360ea2dddb63789fc9471ae6753534adca05", size = 10655025, upload-time = "2026-03-13T12:34:37.105Z" },
{ url = "https://files.pythonhosted.org/packages/b5/9b/bac771774c198c318ae699fc013d8cd99ed9caf993f661fba11238759244/ty-0.0.23-py3-none-manylinux_2_17_ppc64le.manylinux2014_ppc64le.whl", hash = "sha256:d2b162768764d9dc177c83fb497a51532bb67cbebe57b8fa0f2668436bf53f3c", size = 11230107, upload-time = "2026-03-13T12:34:20.751Z" },
{ url = "https://files.pythonhosted.org/packages/14/09/7644fb0e297265e18243f878aca343593323b9bb19ed5278dcbc63781be0/ty-0.0.23-py3-none-manylinux_2_17_s390x.manylinux2014_s390x.whl", hash = "sha256:d28384e48ca03b34e4e2beee0e230c39bbfb68994bb44927fec61ef3642900da", size = 10934177, upload-time = "2026-03-13T12:34:17.904Z" },
{ url = "https://files.pythonhosted.org/packages/18/14/69a25a0cad493fb6a947302471b579a03516a3b00e7bece77fdc6b4afb9b/ty-0.0.23-py3-none-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:559d9a299df793cb7a7902caed5eda8a720ff69164c31c979673e928f02251ee", size = 10752487, upload-time = "2026-03-13T12:34:31.785Z" },
{ url = "https://files.pythonhosted.org/packages/9d/2a/42fc3cbccf95af0a62308ebed67e084798ab7a85ef073c9986ef18032743/ty-0.0.23-py3-none-musllinux_1_2_aarch64.whl", hash = "sha256:32a7b8a14a98e1d20a9d8d2af23637ed7efdb297ac1fa2450b8e465d05b94482", size = 10133007, upload-time = "2026-03-13T12:34:42.838Z" },
{ url = "https://files.pythonhosted.org/packages/e1/69/307833f1b52fa3670e0a1d496e43ef7df556ecde838192d3fcb9b35e360d/ty-0.0.23-py3-none-musllinux_1_2_armv7l.whl", hash = "sha256:6f803b9b9cca87af793467973b9abdd4b83e6b96d9b5e749d662cff7ead70b6d", size = 10169698, upload-time = "2026-03-13T12:34:12.351Z" },
{ url = "https://files.pythonhosted.org/packages/89/ae/5dd379ec22d0b1cba410d7af31c366fcedff191d5b867145913a64889f66/ty-0.0.23-py3-none-musllinux_1_2_i686.whl", hash = "sha256:4a0bf086ec8e2197b7ea7ebfcf4be36cb6a52b235f8be61647ef1b2d99d6ffd3", size = 10346080, upload-time = "2026-03-13T12:34:40.012Z" },
{ url = "https://files.pythonhosted.org/packages/98/c7/dfc83203d37998620bba9c4873a080c8850a784a8a46f56f8163c5b4e320/ty-0.0.23-py3-none-musllinux_1_2_x86_64.whl", hash = "sha256:252539c3fcd7aeb9b8d5c14e2040682c3e1d7ff640906d63fd2c4ce35865a4ba", size = 10848162, upload-time = "2026-03-13T12:34:45.421Z" },
{ url = "https://files.pythonhosted.org/packages/89/08/05481511cfbcc1fd834b6c67aaae090cb609a079189ddf2032139ccfc490/ty-0.0.23-py3-none-win32.whl", hash = "sha256:51b591d19eef23bbc3807aef77d38fa1f003c354e1da908aa80ea2dca0993f77", size = 9748283, upload-time = "2026-03-13T12:34:50.607Z" },
{ url = "https://files.pythonhosted.org/packages/31/2e/eaed4ff5c85e857a02415084c394e02c30476b65e158eec1938fdaa9a205/ty-0.0.23-py3-none-win_amd64.whl", hash = "sha256:1e137e955f05c501cfbb81dd2190c8fb7d01ec037c7e287024129c722a83c9ad", size = 10698355, upload-time = "2026-03-13T12:34:26.134Z" },
{ url = "https://files.pythonhosted.org/packages/91/29/b32cb7b4c7d56b9ed50117f8ad6e45834aec293e4cb14749daab4e9236d5/ty-0.0.23-py3-none-win_arm64.whl", hash = "sha256:a0399bd13fd2cd6683fd0a2d59b9355155d46546d8203e152c556ddbdeb20842", size = 10155890, upload-time = "2026-03-13T12:34:48.082Z" },
{ url = "https://files.pythonhosted.org/packages/e9/24/e287388c63a19191be26b32ff4dbd06029834068150ebe2532939bc4c851/ty-0.0.33-py3-none-linux_armv6l.whl", hash = "sha256:94d0a9d2234261a8911396d59e506b5923fe0971dbda43b9dcea287936887fcc", size = 11021308, upload-time = "2026-04-28T10:45:43.34Z" },
{ url = "https://files.pythonhosted.org/packages/00/ca/ba1eed819895bd239fba8ee35dfcd5fcb266c203b0914a17a59579096bb5/ty-0.0.33-py3-none-macosx_10_12_x86_64.whl", hash = "sha256:e4a2b5ba078f90de342f56b5f7979bb77c9b9b1d8625a041352ffc6ee93c4073", size = 10777272, upload-time = "2026-04-28T10:45:32.905Z" },
{ url = "https://files.pythonhosted.org/packages/25/a8/c3131d37b44b3fea1d6654a1c929a0cd0873822f77a90482b8ec28f6fbbd/ty-0.0.33-py3-none-macosx_11_0_arm64.whl", hash = "sha256:84ff5707825e9af9668d2bcf66975f93e520a63b524ab494e3a8265735be2563", size = 10201078, upload-time = "2026-04-28T10:45:23.374Z" },
{ url = "https://files.pythonhosted.org/packages/7b/db/d8e37ff0045810cc65e1ff36aa0da0a2253c05659787ac987df8a16c7897/ty-0.0.33-py3-none-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:e375285736f57886868e7af0b11c7b0ec5b6543fa15e7ad2a714fed9f077d4e0", size = 10732347, upload-time = "2026-04-28T10:45:21.444Z" },
{ url = "https://files.pythonhosted.org/packages/e0/1a/20e83a412506a918e4684fc67b567cf7cc13b105470b3428cb23c3d5aa13/ty-0.0.33-py3-none-manylinux_2_17_armv7l.manylinux2014_armv7l.whl", hash = "sha256:5680f6350c3b4e46b8bff6d7bb132366ea239463d6cad4892725d06046e65464", size = 10808238, upload-time = "2026-04-28T10:45:38.565Z" },
{ url = "https://files.pythonhosted.org/packages/5d/4b/d0a39f4464dc6cb4cc2c159473ce216bd1846bfb684c0323a3cb36dce5c6/ty-0.0.33-py3-none-manylinux_2_17_i686.manylinux2014_i686.whl", hash = "sha256:c5535538bad8d0f7e62bcdff02197cdb30e41451d80b35d27e17d128f2e1dc5d", size = 11288348, upload-time = "2026-04-28T10:45:08.419Z" },
{ url = "https://files.pythonhosted.org/packages/35/7e/f1745e0f9583363d7a83d9a4990fc244f76ecc30840ddad83dc16a33c52d/ty-0.0.33-py3-none-manylinux_2_17_ppc64le.manylinux2014_ppc64le.whl", hash = "sha256:da196c42bbbc069e1e21e3e52107c061aa9660352dae57a41930690b56e2c02d", size = 11789907, upload-time = "2026-04-28T10:45:19.064Z" },
{ url = "https://files.pythonhosted.org/packages/a5/71/25f39f46a12d662859d45bc648555d0661044eb43db6b5648c9947487da9/ty-0.0.33-py3-none-manylinux_2_17_s390x.manylinux2014_s390x.whl", hash = "sha256:9281672921ef6d4460e03146b5e6c18cb1a3e3a3b8a1a88f6f33226d05a469b7", size = 11500774, upload-time = "2026-04-28T10:45:48.012Z" },
{ url = "https://files.pythonhosted.org/packages/94/ec/136959ecbb7c71cb90537f5aea441c73f4ab24612868a6ecdc9d7444d32d/ty-0.0.33-py3-none-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:82c1b8f303f82da64e878108e764be3ecbcd7c9903ac0a7f7031614ed00b97ab", size = 11360314, upload-time = "2026-04-28T10:45:05.402Z" },
{ url = "https://files.pythonhosted.org/packages/cf/95/32809575c222f00beed498cb728e9290a0f5009f930025381bb7253b2206/ty-0.0.33-py3-none-musllinux_1_2_aarch64.whl", hash = "sha256:efe3af412c9ff67bce5fa37d0a2b0d8555c24072b145a5bac6c79637f1c83abe", size = 10707785, upload-time = "2026-04-28T10:45:10.836Z" },
{ url = "https://files.pythonhosted.org/packages/13/89/c8e9531f7aa4a093359e15fa32c8e1277fbbe90d16894d7c6032d29f4b34/ty-0.0.33-py3-none-musllinux_1_2_armv7l.whl", hash = "sha256:aeec29c91ea768601747da546c3efc20b72c2fb1bd52bcc786a5c6eeff51d27b", size = 10834987, upload-time = "2026-04-28T10:45:40.738Z" },
{ url = "https://files.pythonhosted.org/packages/31/16/9835fbcf5338af1a1917bd28fdb8a7193c210b83f243aa286fa9f79cb3ad/ty-0.0.33-py3-none-musllinux_1_2_i686.whl", hash = "sha256:a535977c52bbb5f7e96b8b70a6ad375ad077f4a9ff2492508ea3816a2b403819", size = 10968968, upload-time = "2026-04-28T10:45:30.26Z" },
{ url = "https://files.pythonhosted.org/packages/36/69/64c76aabc1bc70c7f24b686cd93c3407f8ea430905e395f59bf9603ef571/ty-0.0.33-py3-none-musllinux_1_2_x86_64.whl", hash = "sha256:1d732facf39fcb221ba279d469c5040d37883e964f123b1563888efd34818180", size = 11458077, upload-time = "2026-04-28T10:45:45.971Z" },
{ url = "https://files.pythonhosted.org/packages/91/84/fae27b0c4718776a298690d31ca4cc1995f2e3e1c63a7b59e84c41498e9a/ty-0.0.33-py3-none-win32.whl", hash = "sha256:d90960b574428dc252f85e8598ec5fcb7f619794196b2fc95a90da075ed4681c", size = 10345364, upload-time = "2026-04-28T10:45:16.836Z" },
{ url = "https://files.pythonhosted.org/packages/3c/a0/a2938b23ae3e1a09a2d7c189e2ac5f7113676bae4e0e23948b568e18e5f8/ty-0.0.33-py3-none-win_amd64.whl", hash = "sha256:c1c3aec62c44de610c6e95f0a4e97ac3dbc07934bfdbf1fd90d758c9ff72f48e", size = 11342470, upload-time = "2026-04-28T10:45:26.455Z" },
{ url = "https://files.pythonhosted.org/packages/ab/62/7fb948aace38d2f6329261bb33c035a8484549c74f1db28649c7a4c6fed9/ty-0.0.33-py3-none-win_arm64.whl", hash = "sha256:0d44f99ba1b441e55e2aa301b2ac0a21112784931b46a5f66f4ea9efe5620d97", size = 10742673, upload-time = "2026-04-28T10:45:35.555Z" },
]
[[package]]