From 331056cdc8d93bcb117cc97999af93539176b8ee Mon Sep 17 00:00:00 2001 From: Xi Zhang <106144707+X-iZhang@users.noreply.github.com> Date: Tue, 19 May 2026 13:35:36 +0200 Subject: [PATCH] =?UTF-8?q?feat(middleware):=20upgrade=20deepagents=200.5.?= =?UTF-8?q?7=20=E2=86=92=200.6.2=20(#231)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * feat(middleware): add CodeInterpreterMiddleware with project-specific configuration chore(config): increase checkpoint retention limit for runaway conversations fix(tests): update database schema references from 'blob' to 'value' chore(deps): update deepagents dependency to include quickjs support * feat(deepagents): update to version 0.6.1 and add optional dependencies for quickjs * feat(sessions): improve error handling for message deltas and update Overwrite type check * Enhance PruningCheckpointer with DeltaChannel Awareness - Introduced a new pruning strategy in `_prune_after_put` to preserve the `_DeltaSnapshot` chain during checkpoint pruning. - Implemented methods to fetch recent checkpoint IDs and walk to snapshot ancestors, ensuring that necessary checkpoints are retained. - Updated SQL queries to handle checkpoint and write deletions more efficiently. - Added comprehensive tests for DeltaChannel-aware pruning, ensuring that the pruning logic correctly handles various checkpoint scenarios, including those with and without snapshot seeds. - Refactored `_load_checkpoint_messages` to utilize the new saver interface, improving message reconstruction from checkpoints. * feat(tests): add migration sweep test to preserve snapshot ancestor * feat(sessions): enhance checkpoint retrieval to prevent transcript leakage in multi-agent scenarios * feat(middleware): enhance CodeInterpreterMiddleware with configurable timeout and result character limit feat(config): add CodeInterpreterMiddleware tuning parameters to EvoScientistConfig feat(sessions): implement inline message delta reducer for improved message handling * feat(dependencies): update deepagents version to 0.6.2 in pyproject.toml and uv.lock --- EvoScientist/EvoScientist.py | 8 + EvoScientist/config/settings.py | 20 +- EvoScientist/middleware/__init__.py | 2 + EvoScientist/middleware/code_interpreter.py | 84 +++ EvoScientist/sessions.py | 621 ++++++++++++++----- pyproject.toml | 2 +- tests/test_sessions.py | 637 +++++++++++++++++++- uv.lock | 90 ++- 8 files changed, 1270 insertions(+), 194 deletions(-) create mode 100644 EvoScientist/middleware/code_interpreter.py diff --git a/EvoScientist/EvoScientist.py b/EvoScientist/EvoScientist.py index 6467d1f..19afb25 100644 --- a/EvoScientist/EvoScientist.py +++ b/EvoScientist/EvoScientist.py @@ -468,6 +468,7 @@ def _get_default_middleware(*, for_async_subagent: bool = False): ContextOverflowMapperMiddleware, ModelFallbackMiddleware, ToolErrorHandlerMiddleware, + create_code_interpreter_middleware, create_context_editing_middleware, create_memory_middleware, create_tool_selector_middleware, @@ -497,6 +498,13 @@ def _get_default_middleware(*, for_async_subagent: bool = False): from .middleware.ask_user import AskUserMiddleware mw.insert(0, AskUserMiddleware()) + + mw.append( + create_code_interpreter_middleware( + timeout=cfg.code_interpreter_timeout, + max_result_chars=cfg.code_interpreter_max_result_chars, + ) + ) return mw diff --git a/EvoScientist/config/settings.py b/EvoScientist/config/settings.py index 725d27c..2b11111 100644 --- a/EvoScientist/config/settings.py +++ b/EvoScientist/config/settings.py @@ -259,8 +259,22 @@ class EvoScientistConfig: # Agent features enable_ask_user: bool = True # Enable ask_user tool for agent-initiated questions + # CodeInterpreterMiddleware (PTC — Parallel Tool Calls) tuning + # The PTC allowlist itself is hardcoded in + # ``EvoScientist/middleware/code_interpreter.py`` as a load-bearing safety + # decision (excludes ``execute`` so PTC can't bypass HITL approval, + # excludes ``write_file``/``edit_file`` because batched writes have no + # benefit). Only the resource budget knobs are user-tunable. + code_interpreter_timeout: float = 60.0 # seconds per JS eval + code_interpreter_max_result_chars: int = 10000 # truncate large JSON results + # Checkpoint pruning (sessions.db retention per (thread_id, checkpoint_ns)) - checkpoint_keep_per_thread: int = 10 + # Safety net for runaway conversations. Under DeltaChannel (deepagents 0.6+) + # normal usage produces linear growth, so this default is set well above + # any realistic conversation length (~180-450 turns of dialogue) while + # still capping legacy bloat at upgrade time. 0 disables ongoing pruning + # entirely; the one-time legacy migration sweep still runs. + checkpoint_keep_per_thread: int = 1000 # DM access control policy dm_policy: str = "allowlist" @@ -358,6 +372,8 @@ def _coerce_value(value: Any, field_type: Any) -> Any: return bool(value) if field_type == "int" or field_type is int: return int(value) + if field_type == "float" or field_type is float: + return float(value) return str(value) @@ -454,6 +470,8 @@ _ENV_MAPPINGS = { "checkpoint_keep_per_thread": "EVOSCIENTIST_CHECKPOINT_KEEP_PER_THREAD", "enable_async_subagents": "EVOSCIENTIST_ENABLE_ASYNC_SUBAGENTS", "langgraph_dev_port": "EVOSCIENTIST_LANGGRAPH_DEV_PORT", + "code_interpreter_timeout": "EVOSCIENTIST_CODE_INTERPRETER_TIMEOUT", + "code_interpreter_max_result_chars": "EVOSCIENTIST_CODE_INTERPRETER_MAX_RESULT_CHARS", "langgraph_dev_file_persistence": "EVOSCIENTIST_LANGGRAPH_DEV_FILE_PERSISTENCE", "langgraph_dev_jobs_per_worker": "EVOSCIENTIST_LANGGRAPH_DEV_JOBS_PER_WORKER", "recursion_limit": "EVOSCIENTIST_RECURSION_LIMIT", diff --git a/EvoScientist/middleware/__init__.py b/EvoScientist/middleware/__init__.py index f725650..92ad71e 100644 --- a/EvoScientist/middleware/__init__.py +++ b/EvoScientist/middleware/__init__.py @@ -11,6 +11,7 @@ from .ask_user import ( Choice, Question, ) +from .code_interpreter import create_code_interpreter_middleware from .configurable_model import ConfigurableModelMiddleware from .context_editing import ( compute_context_editing_trigger, @@ -42,6 +43,7 @@ __all__ = [ "Question", "ToolErrorHandlerMiddleware", "compute_context_editing_trigger", + "create_code_interpreter_middleware", "create_context_editing_middleware", "create_memory_middleware", "create_tool_selector_middleware", diff --git a/EvoScientist/middleware/code_interpreter.py b/EvoScientist/middleware/code_interpreter.py new file mode 100644 index 0000000..bc2acfe --- /dev/null +++ b/EvoScientist/middleware/code_interpreter.py @@ -0,0 +1,84 @@ +"""CodeInterpreterMiddleware configuration for EvoScientist. + +Wraps ``langchain-quickjs``'s ``CodeInterpreterMiddleware`` with project-specific +defaults: a PTC allowlist scoped to read-only, batch-friendly tools relevant to +the scientific research workflow (search, sub-agent dispatch, file inspection), +a longer per-eval timeout suitable for LLM-authored algorithms, a larger result +budget for returning structured JSON, and a user-facing tool name that LLMs +recognize from ChatGPT Code Interpreter training data. + +Excluded from PTC by design: + - ``execute`` (shell) — would bypass ``HumanInTheLoopMiddleware`` approval + - ``write_file`` / ``edit_file`` — side-effectful, no batch benefit + - ``think_tool`` — reflection is not batchable + - ``tavily_search`` — only mounted on the ``research-agent`` sub-agent, + not on the main agent; main agent reaches search via ``task`` dispatch + - MCP tools — dynamic at runtime; add manually if a specific server needs PTC + +Usage:: + + from EvoScientist.middleware import create_code_interpreter_middleware + + middleware = create_code_interpreter_middleware( + timeout=60.0, max_result_chars=10000 + ) +""" + +from __future__ import annotations + +from langchain_quickjs import CodeInterpreterMiddleware + +# Defaults match the historical hardcoded values. Callers (the agent +# builder in ``EvoScientist.py``) pass the resolved ``EvoScientistConfig`` +# values; tests / ad-hoc callers can omit and get sensible defaults. +_DEFAULT_TIMEOUT_SECONDS: float = 60.0 +_DEFAULT_MAX_RESULT_CHARS: int = 10000 + +# Read-only, batchable tools that benefit from being callable inside JS. +# Multi-agent orchestration is the killer use case: ``Promise.all`` over +# ``start_async_task`` fans out experiments / writing / data-analysis in +# parallel without each dispatch costing a separate LLM round-trip. Names +# that don't exist at runtime (e.g. async tools when langgraph dev isn't +# reachable) are silently skipped by ``filter_tools_for_ptc``. +_DEFAULT_PTC_ALLOWLIST: list[str] = [ + # Sub-agent dispatch — sync (deepagents) + async (langgraph dev) + "task", + "start_async_task", + "check_async_task", + "update_async_task", + "cancel_async_task", + "list_async_tasks", + # Workspace inspection (read-only, batchable) + "read_file", + "grep", + "glob", + "ls", +] + + +def create_code_interpreter_middleware( + *, + timeout: float = _DEFAULT_TIMEOUT_SECONDS, + max_result_chars: int = _DEFAULT_MAX_RESULT_CHARS, +) -> CodeInterpreterMiddleware: + """Build a project-tuned CodeInterpreterMiddleware instance. + + Args: + timeout: Per-eval timeout in seconds. Defaults to 60s — long enough + for LLM-authored algorithms that touch async sub-agent dispatch + (``start_async_task`` + ``check_async_task`` polling). + max_result_chars: Maximum characters of JS eval output passed back + to the LLM. Defaults to 10k — fits structured JSON aggregations + of file reads / sub-agent results without truncating useful + payloads. Larger values trade tokens for completeness. + + Returns: + Configured ``CodeInterpreterMiddleware`` ready to append to an agent's + middleware stack. + """ + return CodeInterpreterMiddleware( + ptc=_DEFAULT_PTC_ALLOWLIST, + timeout=timeout, + max_result_chars=max_result_chars, + tool_name="code_interpreter", + ) diff --git a/EvoScientist/sessions.py b/EvoScientist/sessions.py index d4cfa97..6259d36 100644 --- a/EvoScientist/sessions.py +++ b/EvoScientist/sessions.py @@ -27,8 +27,16 @@ from pathlib import Path from typing import Any import aiosqlite +from langchain_core.messages import ( + AnyMessage, + BaseMessage, + RemoveMessage, + convert_to_messages, +) from langgraph.checkpoint.serde.jsonplus import JsonPlusSerializer from langgraph.checkpoint.sqlite.aio import AsyncSqliteSaver +from langgraph.graph.message import REMOVE_ALL_MESSAGES +from langgraph.types import Overwrite _logger = logging.getLogger(__name__) @@ -102,7 +110,10 @@ def generate_thread_id() -> str: # Default kept when the caller cannot resolve config (tests, unit-init paths). # Production callers use ``EvoScientistConfig.checkpoint_keep_per_thread``. -_DEFAULT_KEEP_PER_NS = 10 +# Kept in sync with the dataclass default so config-failure fallbacks +# don't silently regress to the pre-DeltaChannel aggressive value (which +# could prune away ``_DeltaSnapshot`` seeds and break message replay). +_DEFAULT_KEEP_PER_NS = 1000 class PruningCheckpointer(AsyncSqliteSaver): @@ -192,85 +203,211 @@ class PruningCheckpointer(AsyncSqliteSaver): _logger.warning("checkpoint pruning failed: %s", exc, exc_info=True) return result + # Safety cap on the snapshot walk. Upstream default + # ``snapshot_frequency`` is 1000; ``DELTA_MAX_SUPERSTEPS_SINCE_SNAPSHOT`` + # is 5000. 10000 is a generous ceiling that catches pathological data + # (cycles, malformed parents) without raising a hard limit on normal + # operation. + _MAX_SNAPSHOT_WALK_STEPS = 10000 + async def _prune_after_put(self, thread_id: str, checkpoint_ns: str) -> None: - """Run the two DELETE queries (writes first, then checkpoints). + """Prune old checkpoints with DeltaChannel awareness. + + Naively keeping the N most-recent rows can sever the + ``_DeltaSnapshot`` chain that ``messages`` reconstruction + depends on — the surviving "latest" checkpoint is rarely a + snapshot point itself, so delta channels silently reconstruct + as empty (upstream ``BaseCheckpointSaver.prune`` spells out the + same failure mode). + + After selecting the N most-recent anchor ids, walk back from + the OLDEST anchor's parent via ``parent_checkpoint_id`` until + hitting an ancestor whose ``channel_values["messages"]`` is a + seed (``_DeltaSnapshot`` blob or plain list — both detected via + ``_unwrap_messages_seed``). All visited ancestors are preserved + alongside the anchor set. Anchors form a contiguous head, so a + single walk from the oldest one covers all of them. Restricted to rows whose ``metadata.agent_name == AGENT_NAME``. ``json_extract(metadata, '$.agent_name') = ?`` evaluates to NULL - (and so fails the predicate) for any row whose metadata does not - carry an ``agent_name`` key — by design, those rows belong to - third-party LangGraph users and must never be pruned by us. + (and so fails the predicate) for any row whose metadata lacks an + ``agent_name`` key — by design, those rows belong to third-party + LangGraph users and must never be pruned by us. - Runs through the saver's connection and lock for atomicity with - concurrent ``aput()`` calls on the same thread. + Walk + DELETEs held under ``self.lock`` for atomicity with + concurrent ``aput()`` on the same thread. """ keep = self._keep_per_ns agent = AGENT_NAME - # writes first — we look up which checkpoint ids will be deleted, then - # delete the writes pointing at them. Doing checkpoints first would - # leave orphan writes whose `checkpoint_id` we can no longer resolve. - del_writes = ( - "DELETE FROM writes " + async with self.lock: + # ``writes`` table is checked inside ``_delete_outside`` so a + # legacy DB that only has ``checkpoints`` still gets pruned + # (writes DELETE silently skipped; checkpoints DELETE runs). + # The migration sweep depends on this — it walks legacy DBs + # that often pre-date the ``writes`` table entirely. + anchor_ids = await self._fetch_recent_checkpoint_ids( + thread_id, checkpoint_ns, agent, keep + ) + if len(anchor_ids) < keep: + return # nothing to prune yet + + extra_preserve = await self._walk_to_snapshot_ancestor( + thread_id, checkpoint_ns, anchor_ids[-1] + ) + kept = set(anchor_ids) | extra_preserve + + await self._delete_outside(thread_id, checkpoint_ns, agent, kept) + await self.conn.commit() + + async def _fetch_recent_checkpoint_ids( + self, + thread_id: str, + checkpoint_ns: str, + agent: str, + limit: int, + ) -> list[str]: + """Return the ``limit`` most-recent checkpoint ids (newest first).""" + query = ( + "SELECT checkpoint_id FROM checkpoints " "WHERE thread_id = ? AND checkpoint_ns = ? " - " AND checkpoint_id IN (" - " SELECT checkpoint_id FROM checkpoints " - " WHERE thread_id = ? AND checkpoint_ns = ? " - " AND json_extract(metadata, '$.agent_name') = ? " - " AND checkpoint_id NOT IN (" - " SELECT checkpoint_id FROM checkpoints " - " WHERE thread_id = ? AND checkpoint_ns = ? " - " AND json_extract(metadata, '$.agent_name') = ? " - " ORDER BY checkpoint_id DESC LIMIT ?" - " )" - " )" + " AND json_extract(metadata, '$.agent_name') = ? " + "ORDER BY checkpoint_id DESC LIMIT ?" ) + async with self.conn.execute( + query, (thread_id, checkpoint_ns, agent, limit) + ) as cur: + rows = await cur.fetchall() + return [r[0] for r in rows] + + async def _walk_to_snapshot_ancestor( + self, + thread_id: str, + checkpoint_ns: str, + oldest_anchor_id: str, + ) -> set[str]: + """Walk parent chain until hitting a ``messages`` seed. + + Returns the set of ancestor ids to preserve (inclusive of the + snapshot ancestor). On chain-break or deserialization failure, + returns what was visited so far — the safe side is over-preserve. + """ + extra: set[str] = set() + cursor = await self._fetch_parent_checkpoint_id( + thread_id, checkpoint_ns, oldest_anchor_id + ) + steps = 0 + while cursor is not None and steps < self._MAX_SNAPSHOT_WALK_STEPS: + steps += 1 + blob = await self._fetch_checkpoint_blob(thread_id, checkpoint_ns, cursor) + if blob is None: + break # chain broken (legacy DB); preserve what we have + extra.add(cursor) + try: + ck = self.serde.loads_typed(blob) + except Exception as exc: + _logger.warning( + "Failed to deserialize checkpoint %s while walking to " + "snapshot for thread %s: %s", + cursor, + thread_id, + exc, + ) + break # safe-side: preserve everything visited so far + cv = ck.get("channel_values") or {} + if _unwrap_messages_seed(cv.get("messages")) is not None: + break # found seed; this ancestor anchors reconstruction + cursor = await self._fetch_parent_checkpoint_id( + thread_id, checkpoint_ns, cursor + ) + return extra + + async def _fetch_parent_checkpoint_id( + self, thread_id: str, checkpoint_ns: str, checkpoint_id: str + ) -> str | None: + query = ( + "SELECT parent_checkpoint_id FROM checkpoints " + "WHERE thread_id = ? AND checkpoint_ns = ? AND checkpoint_id = ?" + ) + async with self.conn.execute( + query, (thread_id, checkpoint_ns, checkpoint_id) + ) as cur: + row = await cur.fetchone() + return row[0] if row and row[0] else None + + async def _fetch_checkpoint_blob( + self, thread_id: str, checkpoint_ns: str, checkpoint_id: str + ) -> tuple[str, bytes] | None: + query = ( + "SELECT type, checkpoint FROM checkpoints " + "WHERE thread_id = ? AND checkpoint_ns = ? AND checkpoint_id = ?" + ) + async with self.conn.execute( + query, (thread_id, checkpoint_ns, checkpoint_id) + ) as cur: + row = await cur.fetchone() + if not row or not row[0] or not row[1]: + return None + return (row[0], row[1]) + + async def _delete_outside( + self, + thread_id: str, + checkpoint_ns: str, + agent: str, + kept_ids: set[str], + ) -> None: + """DELETE rows whose ``checkpoint_id`` is NOT in ``kept_ids``. + + Writes deleted first to preserve referential ordering — if we + dropped checkpoints first, surviving writes' ``checkpoint_id`` + would become orphans. + + Empty ``kept_ids`` is a no-op rather than "delete everything" — + a defensive check; the caller always passes anchor_ids which is + non-empty by construction (already checked ``len >= keep`` in + the caller). + """ + if not kept_ids: + return + kept_list = list(kept_ids) + placeholders = ",".join("?" * len(kept_list)) + + # Writes DELETE only runs if the ``writes`` table exists. Legacy + # DBs from pre-DeltaChannel builds may have only ``checkpoints`` — + # we still want to prune those, just skipping the writes step. + if await _table_exists(self.conn, "writes"): + del_writes = ( + "DELETE FROM writes " + "WHERE thread_id = ? AND checkpoint_ns = ? " + " AND checkpoint_id IN (" + " SELECT checkpoint_id FROM checkpoints " + " WHERE thread_id = ? AND checkpoint_ns = ? " + " AND json_extract(metadata, '$.agent_name') = ? " + f" AND checkpoint_id NOT IN ({placeholders})" + " )" + ) + await self.conn.execute( + del_writes, + ( + thread_id, + checkpoint_ns, + thread_id, + checkpoint_ns, + agent, + *kept_list, + ), + ) + del_checkpoints = ( "DELETE FROM checkpoints " "WHERE thread_id = ? AND checkpoint_ns = ? " " AND json_extract(metadata, '$.agent_name') = ? " - " AND checkpoint_id NOT IN (" - " SELECT checkpoint_id FROM checkpoints " - " WHERE thread_id = ? AND checkpoint_ns = ? " - " AND json_extract(metadata, '$.agent_name') = ? " - " ORDER BY checkpoint_id DESC LIMIT ?" - " )" + f" AND checkpoint_id NOT IN ({placeholders})" + ) + await self.conn.execute( + del_checkpoints, + (thread_id, checkpoint_ns, agent, *kept_list), ) - async with self.lock: - # Skip the writes DELETE on a legacy DB that only has the - # checkpoints table — without this guard, a missing ``writes`` - # would raise sqlite3.OperationalError, the outer try/except - # in ``aput`` would log+swallow it, and pruning would silently - # stop working on the very databases this feature is meant to - # clean up. ``AsyncSqliteSaver.setup()`` normally creates both - # tables, but inherited DBs from older builds may lag. - if await _table_exists(self.conn, "writes"): - await self.conn.execute( - del_writes, - ( - thread_id, - checkpoint_ns, - thread_id, - checkpoint_ns, - agent, - thread_id, - checkpoint_ns, - agent, - keep, - ), - ) - await self.conn.execute( - del_checkpoints, - ( - thread_id, - checkpoint_ns, - agent, - thread_id, - checkpoint_ns, - agent, - keep, - ), - ) - await self.conn.commit() # --------------------------------------------------------------------------- @@ -371,49 +508,256 @@ async def _table_exists(conn: aiosqlite.Connection, table: str) -> bool: return await cur.fetchone() is not None +def _reduce_messages_delta( + state: list[AnyMessage], writes: list[Any] +) -> list[AnyMessage]: + """Inline copy of deepagents' ``_messages_delta_reducer``. + + The upstream reducer lives in ``deepagents._messages_reducer`` (a + private module) and itself adapts langgraph's experimental + ``_messages_delta_reducer`` (PR #7729). Both surfaces are + pre-stable — langgraph marks DeltaChannel as Beta, and the deepagents + file's leading underscore signals it's not part of the public API. + + We copy the implementation here so a future upstream rename or + semantic shift doesn't silently break thread reconstruction. + Behavior MUST stay equivalent: dedups by message ``id``, tombstones + via ``RemoveMessage``, resets on ``REMOVE_ALL_MESSAGES``. ID-less + messages are appended without ID assignment — checkpointers + serialize pending writes before ``update()`` runs, so IDs assigned + inside the reducer never reach stored writes and would differ on + replay, defeating deduplication. + + Raw dict / string / tuple inputs are coerced to typed ``BaseMessage`` + so HTTP-driven graphs (and persisted blobs that round-tripped + through JSON) reconstruct correctly without a separate coercion + step. + """ + flat: list[Any] = [] + for w in writes: + if isinstance(w, list): + flat.extend(w) + else: + flat.append(w) + # Steady-state writes from this module already typed; only raw input + # (deserialized blobs, dict shorthands) needs ``convert_to_messages``. + state_msgs: list[AnyMessage] = ( + state + if state and isinstance(state[0], BaseMessage) + else convert_to_messages(state) # type: ignore[arg-type] + ) + msgs: list[AnyMessage] = convert_to_messages(flat) # type: ignore[assignment] + + # ``REMOVE_ALL_MESSAGES`` resets everything; honor the last sentinel + # in the batch — discard prior state plus every write before it. + remove_all_idx: int | None = 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 :] + + 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: + mid = msg.id + if mid is None: + result.append(msg) + elif isinstance(msg, RemoveMessage): + if mid in index: + result[index[mid]] = None + del index[mid] + elif mid in index: + result[index[mid]] = msg + else: + index[mid] = len(result) + result.append(msg) + return [m for m in result if m is not None] + + async def _load_checkpoint_messages( - conn: aiosqlite.Connection, + saver: AsyncSqliteSaver, thread_id: str, - serde: JsonPlusSerializer, ) -> list: """Load messages from the most recent checkpoint for *thread_id*. + Delegates the ``messages`` DeltaChannel walk to upstream + ``BaseCheckpointSaver.aget_delta_channel_history`` — it finds the + nearest ancestor whose ``channel_values["messages"]`` carries a seed + (``_DeltaSnapshot`` blob or plain list) and returns that plus every + on-path pending write oldest→newest. The local + ``_reduce_messages_delta`` (inline copy of deepagents' reducer; see + that function's docstring) is then applied in a single batched call, + preserving dedup-by-id, ``RemoveMessage`` tombstones, and + ``REMOVE_ALL_MESSAGES`` reset semantics. + + ``Overwrite`` (``langgraph.types.Overwrite``) is not a message-like + and the reducer doesn't recognize it — split the batch at each + occurrence, replacing accumulated state with the wrapped value before + resuming reducer application. + + ``_summarization_event`` doesn't ride on the ``messages`` channel, + so it's fetched separately from the latest checkpoint's + ``channel_values``. + Returns a list of LangChain message objects, or an empty list on failure. """ - channel_values = await _load_checkpoint_channel_values(conn, thread_id, serde) - messages = channel_values.get("messages", []) - if not isinstance(messages, list): - return [] - event = channel_values.get("_summarization_event") - return _apply_summarization_event( - messages, event if isinstance(event, dict) else None + # Pre-resolve the latest EvoScientist checkpoint_id with an + # ``agent_name`` filter, then pin it into the config so + # ``aget_tuple`` fetches THAT specific row. Without the pin, + # ``aget_tuple`` returns the latest by ``checkpoint_id`` alone — in + # a multi-agent DB where a third-party tool shares the same + # ``(thread_id, checkpoint_ns)`` and happens to have a higher id, + # we'd leak that agent's transcript into our /resume. The ancestor + # walk via ``parent_checkpoint_id`` chain is unambiguous (specific + # ids), so pinning the head is sufficient — the rest of the chain + # follows EvoScientist's parent links. + head_query = ( + "SELECT checkpoint_id FROM checkpoints " + "WHERE thread_id = ? AND checkpoint_ns = '' " + " AND json_extract(metadata, '$.agent_name') = ? " + "ORDER BY checkpoint_id DESC LIMIT 1" ) + async with saver.conn.execute(head_query, (thread_id, AGENT_NAME)) as cur: + head_row = await cur.fetchone() + if head_row is None: + return [] + config = { + "configurable": { + "thread_id": thread_id, + "checkpoint_ns": "", + "checkpoint_id": head_row[0], + } + } + target = await saver.aget_tuple(config) + if target is None: + return [] + + # ``aget_delta_channel_history`` walks from ``target.parent_config`` — + # it deliberately excludes the target itself (its caller is the + # runtime preparing to apply a NEW delta on top). For /resume we want + # the state AT the latest checkpoint, so check the target's own + # ``channel_values`` first, falling back to the ancestor walk only + # when no seed is materialized locally. + target_cv = target.checkpoint.get("channel_values") or {} + target_seed = _unwrap_messages_seed(target_cv.get("messages")) + if target_seed is not None: + accumulated: list = target_seed + writes: list = [w for w in (target.pending_writes or []) if w[1] == "messages"] + else: + history = await saver.aget_delta_channel_history( + config=config, channels=["messages"] + ) + entry = history.get("messages", {}) + accumulated = _unwrap_messages_seed(entry.get("seed")) or [] + writes = list(entry.get("writes", [])) + # ``target.pending_writes`` are the deltas recorded at this step; + # apply them on top of whatever the ancestor walk reconstructed. + writes.extend(w for w in (target.pending_writes or []) if w[1] == "messages") + + # Batched reducer: collect contiguous message-like writes and flush + # in one call. Overwrite splits the batch because it resets state + # rather than appending. + batch: list = [] + + def _flush() -> None: + nonlocal accumulated + if not batch: + return + try: + accumulated = _reduce_messages_delta(accumulated, batch) + except Exception as exc: + _logger.warning( + "Failed to apply %d messages deltas for thread %s: %s", + len(batch), + thread_id, + exc, + ) + batch.clear() + + for _task_id, _channel, delta in writes: + # ``Overwrite`` wraps a value with "replace this channel" + # semantics — flush any pending writes first, then reset state + # to the wrapped value. + if isinstance(delta, Overwrite): + _flush() + inner = getattr(delta, "value", None) + accumulated = ( + list(inner) + if isinstance(inner, list) + else ([inner] if inner is not None else []) + ) + else: + batch.append(delta) + _flush() + + # ``_summarization_event`` rides on its own channel; pick it off the + # target's ``channel_values`` we already deserialized above. + event = target_cv.get("_summarization_event") + summarization_event = event if isinstance(event, dict) else None + + if not accumulated: + await _log_orphan_warning_if_pruned(saver.conn, thread_id) + + if not isinstance(accumulated, list): + return [] + return _apply_summarization_event(accumulated, summarization_event) -async def _load_checkpoint_channel_values( - conn: aiosqlite.Connection, - thread_id: str, - serde: JsonPlusSerializer, -) -> dict: - """Load channel_values from the most recent checkpoint for *thread_id*.""" +async def _log_orphan_warning_if_pruned( + conn: aiosqlite.Connection, thread_id: str +) -> None: + """Emit a WARNING if *thread_id* has orphan ``parent_checkpoint_id`` refs. + + Empty reconstructed history + a broken parent chain is the signature + of a pre-fix DB where the old ``keep_per_ns=10`` default pruned early + writes before the snapshot frequency materialized a ``_DeltaSnapshot`` + seed. The DB has no path to recover those messages — log so /resume + showing a stub history isn't silent. + """ query = """ - SELECT type, checkpoint - FROM checkpoints - WHERE thread_id = ? - AND json_extract(metadata, '$.agent_name') = ? - ORDER BY checkpoint_id DESC + SELECT 1 FROM checkpoints c1 + WHERE c1.thread_id = ? AND c1.checkpoint_ns = '' + AND c1.parent_checkpoint_id IS NOT NULL + AND NOT EXISTS ( + SELECT 1 FROM checkpoints c2 + WHERE c2.thread_id = c1.thread_id + AND c2.checkpoint_ns = '' + AND c2.checkpoint_id = c1.parent_checkpoint_id + ) LIMIT 1 """ - async with conn.execute(query, (thread_id, AGENT_NAME)) as cur: - row = await cur.fetchone() - if not row or not row[0] or not row[1]: - return {} - try: - data = serde.loads_typed((row[0], row[1])) - channel_values = data.get("channel_values", {}) - return channel_values if isinstance(channel_values, dict) else {} - except (ValueError, TypeError, KeyError): - return {} + async with conn.execute(query, (thread_id,)) as cur: + if await cur.fetchone(): + _logger.warning( + "Thread %s has orphan checkpoints (pre-fix pruning) " + "and no surviving DeltaChannel snapshot; reconstructed " + "history is empty. Early messages cannot be recovered.", + thread_id, + ) + + +def _unwrap_messages_seed(value: object) -> list | None: + """Coerce a ``channel_values["messages"]`` snapshot seed into a plain list. + + LangGraph 1.2 stores snapshot blobs as ``_DeltaSnapshot(value=[...])`` + (a ``NamedTuple``, NOT a list subclass), so a bare ``isinstance(v, list)`` + check silently ignores the seed and reconstruction starts from whatever + writes survived pruning. Pre-migration / non-DeltaChannel checkpoints + still store a plain list. Returns ``None`` when no usable seed is + present (caller leaves accumulated state untouched). + """ + if value is None: + return None + if isinstance(value, list): + return list(value) + inner = getattr(value, "value", None) + if isinstance(inner, list): + return list(inner) + return None def _apply_summarization_event(messages: list, event: dict | None) -> list: @@ -532,9 +876,11 @@ async def list_threads( ] if (include_message_count or include_preview) and threads: + # Share one saver across all threads so ``setup()`` runs once. serde = JsonPlusSerializer() + saver = AsyncSqliteSaver(conn, serde=serde) for t in threads: - msgs = await _load_checkpoint_messages(conn, t["thread_id"], serde) + msgs = await _load_checkpoint_messages(saver, t["thread_id"]) if include_message_count: t["message_count"] = len(msgs) if include_preview: @@ -677,6 +1023,11 @@ async def get_thread_messages(thread_id: str) -> list: Only returns messages for EvoScientist threads. Returns an empty list if the thread has no checkpoints. + + Reconstructs the full message history by walking the checkpoint chain + and applying pending writes — required under deepagents 0.6 + ``DeltaChannel`` where messages live in the ``writes`` table rather + than the latest checkpoint's ``channel_values``. """ db_path = str(get_db_path()) async with aiosqlite.connect(db_path, timeout=30.0) as conn: @@ -692,10 +1043,8 @@ async def get_thread_messages(thread_id: str) -> list: if not await cur.fetchone(): return [] serde = JsonPlusSerializer() - channel_values = await _load_checkpoint_channel_values(conn, thread_id, serde) - messages = channel_values.get("messages", []) - event = channel_values.get("_summarization_event") - return _apply_summarization_event(messages, event) + saver = AsyncSqliteSaver(conn, serde=serde) + return await _load_checkpoint_messages(saver, thread_id) # --------------------------------------------------------------------------- @@ -850,9 +1199,6 @@ async def _run_migration_sweep( return 0 if await _get_user_version(conn) >= _MIGRATION_VERSION: return 0 - # ``writes`` is optional on legacy DBs: skip the writes DELETE - # if the table is absent rather than aborting the whole sweep. - has_writes = await _table_exists(conn, "writes") async with conn.execute( "SELECT DISTINCT thread_id, checkpoint_ns FROM checkpoints " @@ -861,62 +1207,23 @@ async def _run_migration_sweep( ) as cur: pairs = await cur.fetchall() - del_writes = ( - "DELETE FROM writes " - "WHERE thread_id = ? AND checkpoint_ns = ? " - " AND checkpoint_id IN (" - " SELECT checkpoint_id FROM checkpoints " - " WHERE thread_id = ? AND checkpoint_ns = ? " - " AND json_extract(metadata, '$.agent_name') = ? " - " AND checkpoint_id NOT IN (" - " SELECT checkpoint_id FROM checkpoints " - " WHERE thread_id = ? AND checkpoint_ns = ? " - " AND json_extract(metadata, '$.agent_name') = ? " - " ORDER BY checkpoint_id DESC LIMIT ?" - " )" - " )" - ) - del_checkpoints = ( - "DELETE FROM checkpoints " - "WHERE thread_id = ? AND checkpoint_ns = ? " - " AND json_extract(metadata, '$.agent_name') = ? " - " AND checkpoint_id NOT IN (" - " SELECT checkpoint_id FROM checkpoints " - " WHERE thread_id = ? AND checkpoint_ns = ? " - " AND json_extract(metadata, '$.agent_name') = ? " - " ORDER BY checkpoint_id DESC LIMIT ?" - " )" - ) + # Reuse the DeltaChannel-aware prune logic from PruningCheckpointer + # instead of running naive keep_latest SQL: legacy DBs almost always + # have threads where the latest N checkpoints sit ABOVE a + # ``_DeltaSnapshot`` ancestor, and the naive form would sever the + # snapshot chain — exactly the failure mode the steady-state Fix + # already prevents. Sharing one saver across all pairs means + # ``setup()`` and the in-class lock are constructed once. + # + # We invoke ``_prune_after_put`` directly (not ``aput``) — the + # sweep is a bulk cleanup, not a checkpoint write. As a result + # ``saver._aput_lock`` (the outer put+prune pair lock) is + # intentionally unused here; only the inner ``self.lock`` that + # ``_prune_after_put`` itself acquires runs. + saver = PruningCheckpointer(conn, keep_per_ns=keep) for thread_id, checkpoint_ns in pairs: ns = checkpoint_ns or "" - if has_writes: - await conn.execute( - del_writes, - ( - thread_id, - ns, - thread_id, - ns, - AGENT_NAME, - thread_id, - ns, - AGENT_NAME, - keep, - ), - ) - await conn.execute( - del_checkpoints, - ( - thread_id, - ns, - AGENT_NAME, - thread_id, - ns, - AGENT_NAME, - keep, - ), - ) - await conn.commit() + await saver._prune_after_put(str(thread_id), ns) pairs_pruned += 1 if progress_cb is not None: try: diff --git a/pyproject.toml b/pyproject.toml index fb81dae..fef0df7 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -16,7 +16,7 @@ classifiers = [ "Programming Language :: Python :: 3", ] dependencies = [ - "deepagents>=0.5.7", + "deepagents[quickjs]~=0.6.2", "langchain>=1.2", "langchain-anthropic>=1.4", "langchain-openai>=1.2", diff --git a/tests/test_sessions.py b/tests/test_sessions.py index 3979fd5..11f2179 100644 --- a/tests/test_sessions.py +++ b/tests/test_sessions.py @@ -145,12 +145,18 @@ class TestThreadFunctions(unittest.TestCase): idx INTEGER NOT NULL, channel TEXT NOT NULL, type TEXT, - blob BLOB, + value BLOB, PRIMARY KEY (thread_id, checkpoint_ns, checkpoint_id, task_id, idx) ) """) - # Insert test checkpoints + # Insert test checkpoints. ``type`` + ``checkpoint`` are + # populated with a serialized empty-state blob so upstream + # ``aget_tuple`` (used by message reconstruction) can + # deserialize them — production checkpoints always have + # these set; bare-metadata rows are a test fiction. + serde = JsonPlusSerializer() + empty_ck_type, empty_ck_blob = serde.dumps_typed({"channel_values": {}}) for i, tid in enumerate(["abc12345", "abc12399", "def00001"]): meta = json.dumps( { @@ -161,8 +167,8 @@ class TestThreadFunctions(unittest.TestCase): } ) await conn.execute( - "INSERT INTO checkpoints (thread_id, checkpoint_ns, checkpoint_id, metadata) VALUES (?, '', ?, ?)", - (tid, f"cp_{i}", meta), + "INSERT INTO checkpoints (thread_id, checkpoint_ns, checkpoint_id, type, checkpoint, metadata) VALUES (?, '', ?, ?, ?, ?)", + (tid, f"cp_{i}", empty_ck_type, empty_ck_blob, meta), ) # Insert a non-EvoScientist checkpoint (should be filtered) @@ -173,8 +179,8 @@ class TestThreadFunctions(unittest.TestCase): } ) await conn.execute( - "INSERT INTO checkpoints (thread_id, checkpoint_ns, checkpoint_id, metadata) VALUES (?, '', ?, ?)", - ("zzz99999", "cp_other", other_meta), + "INSERT INTO checkpoints (thread_id, checkpoint_ns, checkpoint_id, type, checkpoint, metadata) VALUES (?, '', ?, ?, ?, ?)", + ("zzz99999", "cp_other", empty_ck_type, empty_ck_blob, other_meta), ) await conn.commit() @@ -364,6 +370,254 @@ class TestThreadFunctions(unittest.TestCase): finally: _run(_cleanup()) + def test_get_thread_messages_reconstructs_multi_delta_chain(self): + """3-checkpoint chain with ``_DeltaSnapshot`` seed + pending writes. + + Exercises the upstream ``aget_delta_channel_history`` walk: the + latest checkpoint has no materialized seed, so the walk must climb + back through an intermediate delta-only ancestor to a snapshot + further back, then accumulate writes oldest→newest on top. + """ + from langgraph.checkpoint.serde.types import _DeltaSnapshot + + async def _insert(): + import aiosqlite + + serde = JsonPlusSerializer() + + seed_messages = [ + HumanMessage(content="m1", id="m1"), + AIMessage(content="m2", id="m2"), + ] + cp1_type, cp1_blob = serde.dumps_typed( + {"channel_values": {"messages": _DeltaSnapshot(value=seed_messages)}} + ) + cp_empty_type, cp_empty_blob = serde.dumps_typed({"channel_values": {}}) + + w2_type, w2_blob = serde.dumps_typed([HumanMessage(content="m3", id="m3")]) + w3_type, w3_blob = serde.dumps_typed([AIMessage(content="m4", id="m4")]) + + meta = json.dumps( + { + "agent_name": AGENT_NAME, + "updated_at": "2025-01-26T10:00:00+00:00", + } + ) + + async with aiosqlite.connect(self._db_path) as conn: + for cid, parent, ck_type, ck_blob in [ + ("cp_chain_1", None, cp1_type, cp1_blob), + ("cp_chain_2", "cp_chain_1", cp_empty_type, cp_empty_blob), + ("cp_chain_3", "cp_chain_2", cp_empty_type, cp_empty_blob), + ]: + await conn.execute( + "INSERT INTO checkpoints (thread_id, checkpoint_ns, checkpoint_id, parent_checkpoint_id, type, checkpoint, metadata) " + "VALUES (?, '', ?, ?, ?, ?, ?)", + ("chain12345", cid, parent, ck_type, ck_blob, meta), + ) + for cid, wtype, wblob in [ + ("cp_chain_2", w2_type, w2_blob), + ("cp_chain_3", w3_type, w3_blob), + ]: + await conn.execute( + "INSERT INTO writes (thread_id, checkpoint_ns, checkpoint_id, task_id, idx, channel, type, value) " + "VALUES (?, '', ?, ?, ?, ?, ?, ?)", + ("chain12345", cid, "task0", 0, "messages", wtype, wblob), + ) + await conn.commit() + + async def _cleanup(): + import aiosqlite + + async with aiosqlite.connect(self._db_path) as conn: + await conn.execute( + "DELETE FROM checkpoints WHERE thread_id = ?", + ("chain12345",), + ) + await conn.execute( + "DELETE FROM writes WHERE thread_id = ?", + ("chain12345",), + ) + await conn.commit() + + async def _assert_walk_branch_active(): + """Confirm cp_chain_3 has no materialized ``messages`` seed. + + Without this guard, a future refactor that ends up writing a + ``_DeltaSnapshot`` at every checkpoint would silently move + this test onto the hybrid's "target seed" branch — the + reconstruction would still match, but the ancestor walk + being tested here would never run. The assertion locks in + which branch the test exercises. + """ + import aiosqlite + + async with aiosqlite.connect(self._db_path) as conn: + async with conn.execute( + "SELECT type, checkpoint FROM checkpoints " + "WHERE thread_id = ? AND checkpoint_id = ?", + ("chain12345", "cp_chain_3"), + ) as cur: + row = await cur.fetchone() + assert row is not None + ck = JsonPlusSerializer().loads_typed((row[0], row[1])) + assert "messages" not in (ck.get("channel_values") or {}), ( + "cp_chain_3 must NOT carry a messages seed — this test " + "exercises the ancestor-walk branch, not the target-seed " + "shortcut." + ) + + _run(_insert()) + try: + _run(_assert_walk_branch_active()) + messages = _run(get_thread_messages("chain12345")) + assert [m.content for m in messages] == ["m1", "m2", "m3", "m4"] + assert isinstance(messages[0], HumanMessage) + assert isinstance(messages[1], AIMessage) + assert isinstance(messages[2], HumanMessage) + assert isinstance(messages[3], AIMessage) + finally: + _run(_cleanup()) + + def test_get_thread_messages_handles_overwrite_bare_message(self): + """``Overwrite(value=)`` wraps to a single-element list. + + The ``Overwrite`` reset branch in ``_load_checkpoint_messages`` + has three sub-cases: list value (most common), ``None`` (clears + state), and a bare ``BaseMessage`` (rare but valid). The + last case has no other test coverage — this guards against a + refactor that silently drops the ``[inner]`` wrapping fallback. + """ + from langgraph.checkpoint.serde.types import _DeltaSnapshot + from langgraph.types import Overwrite + + async def _insert(): + import aiosqlite + + serde = JsonPlusSerializer() + seed_messages = [ + HumanMessage(content="m1", id="m1"), + AIMessage(content="m2", id="m2"), + ] + ck_type, ck_blob = serde.dumps_typed( + {"channel_values": {"messages": _DeltaSnapshot(value=seed_messages)}} + ) + # Bare message (NOT wrapped in a list) — the rare third case + # the Overwrite branch handles. + ow = Overwrite(value=HumanMessage(content="replaced", id="repl")) + w_type, w_blob = serde.dumps_typed(ow) + meta = json.dumps({"agent_name": AGENT_NAME}) + async with aiosqlite.connect(self._db_path) as conn: + await conn.execute( + "INSERT INTO checkpoints (thread_id, checkpoint_ns, checkpoint_id, " + "parent_checkpoint_id, type, checkpoint, metadata) " + "VALUES (?, '', ?, NULL, ?, ?, ?)", + ("ow_bare01", "cp_bare", ck_type, ck_blob, meta), + ) + await conn.execute( + "INSERT INTO writes (thread_id, checkpoint_ns, checkpoint_id, task_id, idx, channel, type, value) " + "VALUES (?, '', ?, 'task0', 0, 'messages', ?, ?)", + ("ow_bare01", "cp_bare", w_type, w_blob), + ) + await conn.commit() + + async def _cleanup(): + import aiosqlite + + async with aiosqlite.connect(self._db_path) as conn: + await conn.execute( + "DELETE FROM checkpoints WHERE thread_id = ?", ("ow_bare01",) + ) + await conn.execute( + "DELETE FROM writes WHERE thread_id = ?", ("ow_bare01",) + ) + await conn.commit() + + _run(_insert()) + try: + messages = _run(get_thread_messages("ow_bare01")) + # Overwrite replaced the seed completely; bare message wrapped + # in a 1-element list. + assert len(messages) == 1 + assert isinstance(messages[0], HumanMessage) + assert messages[0].content == "replaced" + assert messages[0].id == "repl" + finally: + _run(_cleanup()) + + def test_get_thread_messages_ignores_colliding_other_agent(self): + """Multi-agent DB with thread_id collision: must surface only ours. + + Without the agent_name filter on the head-checkpoint lookup, + ``saver.aget_tuple()`` returns the latest by ``checkpoint_id`` + alone — so if a third-party agent's checkpoint for the same + ``thread_id`` happens to have a higher id (lexicographically), + we'd leak its transcript into ``/resume``. Pinning the head to + the latest EvoScientist-agent row prevents that. + """ + from langgraph.checkpoint.serde.types import _DeltaSnapshot + + async def _insert(): + import aiosqlite + + serde = JsonPlusSerializer() + evo_messages = [ + HumanMessage(content="ours_1", id="o1"), + AIMessage(content="ours_2", id="o2"), + ] + other_messages = [ + HumanMessage(content="theirs_1", id="t1"), + AIMessage(content="theirs_2", id="t2"), + HumanMessage(content="theirs_3", id="t3"), + ] + evo_type, evo_blob = serde.dumps_typed( + {"channel_values": {"messages": _DeltaSnapshot(value=evo_messages)}} + ) + other_type, other_blob = serde.dumps_typed( + {"channel_values": {"messages": _DeltaSnapshot(value=other_messages)}} + ) + evo_meta = json.dumps({"agent_name": AGENT_NAME}) + other_meta = json.dumps({"agent_name": "ThirdPartyAgent"}) + + async with aiosqlite.connect(self._db_path) as conn: + # EvoScientist's checkpoint id is LEXICOGRAPHICALLY + # SMALLER than the third-party agent's, so a naive + # "latest by checkpoint_id" lookup would pick the wrong + # one. + await conn.execute( + "INSERT INTO checkpoints (thread_id, checkpoint_ns, checkpoint_id, " + "parent_checkpoint_id, type, checkpoint, metadata) " + "VALUES (?, '', 'aaa_evo', NULL, ?, ?, ?)", + ("collide01", evo_type, evo_blob, evo_meta), + ) + await conn.execute( + "INSERT INTO checkpoints (thread_id, checkpoint_ns, checkpoint_id, " + "parent_checkpoint_id, type, checkpoint, metadata) " + "VALUES (?, '', 'zzz_other', NULL, ?, ?, ?)", + ("collide01", other_type, other_blob, other_meta), + ) + await conn.commit() + + async def _cleanup(): + import aiosqlite + + async with aiosqlite.connect(self._db_path) as conn: + await conn.execute( + "DELETE FROM checkpoints WHERE thread_id = ?", ("collide01",) + ) + await conn.commit() + + _run(_insert()) + try: + messages = _run(get_thread_messages("collide01")) + assert [m.content for m in messages] == ["ours_1", "ours_2"] + # Defense-in-depth: explicitly forbid leakage of the other + # agent's content. + for msg in messages: + assert not msg.content.startswith("theirs_") + finally: + _run(_cleanup()) + # -- Agent isolation: OtherAgent data should never be visible -- def test_thread_exists_ignores_other_agent(self): @@ -403,7 +657,7 @@ class TestThreadFunctions(unittest.TestCase): (shared_tid, "cp_evo_shared", evo_meta), ) await conn.execute( - "INSERT INTO writes (thread_id, checkpoint_ns, checkpoint_id, task_id, idx, channel, type, blob) " + "INSERT INTO writes (thread_id, checkpoint_ns, checkpoint_id, task_id, idx, channel, type, value) " "VALUES (?, '', ?, 't1', 0, 'ch', 'str', X'AA')", (shared_tid, "cp_evo_shared"), ) @@ -420,7 +674,7 @@ class TestThreadFunctions(unittest.TestCase): (shared_tid, "cp_other_shared", other_meta), ) await conn.execute( - "INSERT INTO writes (thread_id, checkpoint_ns, checkpoint_id, task_id, idx, channel, type, blob) " + "INSERT INTO writes (thread_id, checkpoint_ns, checkpoint_id, task_id, idx, channel, type, value) " "VALUES (?, '', ?, 't2', 0, 'ch', 'str', X'BB')", (shared_tid, "cp_other_shared"), ) @@ -869,6 +1123,264 @@ class TestPruningCheckpointer(unittest.TestCase): assert survivors == [ids[4], ids[3]] +class TestPruningCheckpointerDeltaChannel(unittest.TestCase): + """Tests for DeltaChannel-aware pruning. + + The naive ``keep_latest`` pruner can sever the ``_DeltaSnapshot`` + chain that ``messages`` reconstruction depends on. These tests + exercise the walk-to-snapshot-ancestor extension that preserves the + chain head between each kept anchor and the nearest snapshot. + + Checkpoints are inserted directly via SQL with explicit + ``parent_checkpoint_id`` to give precise control over the chain + structure (the chained-``aput`` pattern in + ``TestPruningCheckpointer`` doesn't let us pick which checkpoint + materializes a snapshot). + """ + + def setUp(self): + self._tmpdir = tempfile.mkdtemp() + self._db_path = os.path.join(self._tmpdir, "delta_prune.db") + + def tearDown(self): + try: + os.unlink(self._db_path) + os.rmdir(self._tmpdir) + except OSError: + pass + + @staticmethod + async def _insert_chain(conn, thread_id, specs): + """Insert a chain of checkpoints in oldest→newest order. + + ``specs`` is a list of ``(cid, parent_cid, msgs_seed)`` tuples: + - ``cid``: checkpoint_id + - ``parent_cid``: parent_checkpoint_id (``None`` for chain root) + - ``msgs_seed``: value stored at ``channel_values["messages"]`` + — pass ``None`` for delta-only (no seed), a ``list`` for + plain seed, or a ``_DeltaSnapshot`` for wrapped seed. + """ + serde = JsonPlusSerializer() + meta = json.dumps({"agent_name": AGENT_NAME}) + for cid, parent_cid, msgs in specs: + cv = {"messages": msgs} if msgs is not None else {} + ck_type, ck_blob = serde.dumps_typed({"channel_values": cv}) + await conn.execute( + "INSERT INTO checkpoints (thread_id, checkpoint_ns, checkpoint_id, " + "parent_checkpoint_id, type, checkpoint, metadata) " + "VALUES (?, '', ?, ?, ?, ?, ?)", + (thread_id, cid, parent_cid, ck_type, ck_blob, meta), + ) + await conn.commit() + + @staticmethod + async def _surviving_ids(conn, thread_id): + async with conn.execute( + "SELECT checkpoint_id FROM checkpoints WHERE thread_id = ? " + "ORDER BY checkpoint_id ASC", + (thread_id,), + ) as cur: + return [r[0] for r in await cur.fetchall()] + + def test_preserves_snapshot_ancestor(self): + """Snapshot lives outside the anchor window → walk reaches and stops.""" + from langgraph.checkpoint.serde.types import _DeltaSnapshot + + from EvoScientist.sessions import PruningCheckpointer + + tid = "tdsa_001" + + async def _go(): + import aiosqlite + + async with aiosqlite.connect(self._db_path) as conn: + saver = PruningCheckpointer(conn, keep_per_ns=5) + await saver.setup() + # 10 checkpoints; snapshot at cp_003 (outside the anchor + # window of cp_006..cp_010). Walk from cp_005 → cp_004 → + # cp_003 (seed found, stop). Survivors: cp_003..cp_010 + # (8). Pruned: cp_001, cp_002. + specs = [ + ("cp_001", None, None), + ("cp_002", "cp_001", None), + ("cp_003", "cp_002", _DeltaSnapshot(value=[])), + *[(f"cp_{i:03d}", f"cp_{i - 1:03d}", None) for i in range(4, 11)], + ] + await self._insert_chain(conn, tid, specs) + await saver._prune_after_put(tid, "") + await conn.commit() + return await self._surviving_ids(conn, tid) + + survivors = _run(_go()) + assert survivors == [f"cp_{i:03d}" for i in range(3, 11)] + + def test_preserves_full_chain_when_no_snapshot(self): + """No snapshot anywhere → walk reaches root, preserves everything.""" + from EvoScientist.sessions import PruningCheckpointer + + tid = "tdsa_002" + + async def _go(): + import aiosqlite + + async with aiosqlite.connect(self._db_path) as conn: + saver = PruningCheckpointer(conn, keep_per_ns=5) + await saver.setup() + # 10 delta-only checkpoints, no seed anywhere. Walk + # exhausts to root (cp_001's parent is None → break). + # All 10 must survive — the alternative is silent + # truncation, which is the bug Fix #B prevents. + specs = [("cp_001", None, None)] + [ + (f"cp_{i:03d}", f"cp_{i - 1:03d}", None) for i in range(2, 11) + ] + await self._insert_chain(conn, tid, specs) + await saver._prune_after_put(tid, "") + await conn.commit() + return await self._surviving_ids(conn, tid) + + survivors = _run(_go()) + assert survivors == [f"cp_{i:03d}" for i in range(1, 11)] + + def test_plain_list_seed_also_terminates_walk(self): + """Pre-DeltaChannel format (plain list in channel_values) also counts as seed.""" + from EvoScientist.sessions import PruningCheckpointer + + tid = "tdsa_003" + + async def _go(): + import aiosqlite + + async with aiosqlite.connect(self._db_path) as conn: + saver = PruningCheckpointer(conn, keep_per_ns=3) + await saver.setup() + # 6 checkpoints; cp_002 has plain-list seed (legacy + # format). Walk from cp_003 → cp_002 (seed) → stop. + # Survivors: cp_002..cp_006 (5). Pruned: cp_001. + specs = [ + ("cp_001", None, None), + ("cp_002", "cp_001", []), # plain list seed + ("cp_003", "cp_002", None), + ("cp_004", "cp_003", None), + ("cp_005", "cp_004", None), + ("cp_006", "cp_005", None), + ] + await self._insert_chain(conn, tid, specs) + await saver._prune_after_put(tid, "") + await conn.commit() + return await self._surviving_ids(conn, tid) + + survivors = _run(_go()) + assert survivors == ["cp_002", "cp_003", "cp_004", "cp_005", "cp_006"] + + def test_chain_break_stops_walk_cleanly(self): + """Missing ancestor row breaks the chain; walk stops without raising.""" + from EvoScientist.sessions import PruningCheckpointer + + tid = "tdsa_004" + + async def _go(): + import aiosqlite + + async with aiosqlite.connect(self._db_path) as conn: + saver = PruningCheckpointer(conn, keep_per_ns=2) + await saver.setup() + # Insert cp_001..cp_005, then DELETE cp_002 to break + # the chain. Walk from cp_003 → tries cp_002 → + # _fetch_checkpoint_blob returns None → break with + # nothing added (because we add the cursor's id only + # AFTER fetching its blob succeeds). + specs = [ + ("cp_001", None, None), + ("cp_002", "cp_001", None), + ("cp_003", "cp_002", None), + ("cp_004", "cp_003", None), + ("cp_005", "cp_004", None), + ] + await self._insert_chain(conn, tid, specs) + await conn.execute( + "DELETE FROM checkpoints WHERE thread_id = ? AND checkpoint_id = ?", + (tid, "cp_002"), + ) + await conn.commit() + await saver._prune_after_put(tid, "") + await conn.commit() + return await self._surviving_ids(conn, tid) + + survivors = _run(_go()) + # anchors = cp_004, cp_005. Walk visits cp_003 (preserved), + # then cp_002 → None → break. cp_001 pruned. cp_002 already + # absent. Survivors: cp_003, cp_004, cp_005. + assert survivors == ["cp_003", "cp_004", "cp_005"] + + def test_deserialization_failure_safe_side_over_preserves(self): + """Corrupt blob mid-walk: pruner preserves what it visited so far.""" + from EvoScientist.sessions import PruningCheckpointer + + tid = "tdsa_005" + + async def _go(): + import aiosqlite + + async with aiosqlite.connect(self._db_path) as conn: + saver = PruningCheckpointer(conn, keep_per_ns=2) + await saver.setup() + # cp_001..cp_005, all delta-only. Then overwrite cp_003 + # with a corrupt blob. Walk from cp_003: fetch blob + # succeeds (returns garbage bytes), add cp_003 to + # extra_preserve, deserialize FAILS → break with cp_003 + # already preserved. + specs = [ + ("cp_001", None, None), + ("cp_002", "cp_001", None), + ("cp_003", "cp_002", None), + ("cp_004", "cp_003", None), + ("cp_005", "cp_004", None), + ] + await self._insert_chain(conn, tid, specs) + await conn.execute( + "UPDATE checkpoints SET type = ?, checkpoint = ? " + "WHERE thread_id = ? AND checkpoint_id = ?", + ("garbage_type", b"not a real blob", tid, "cp_003"), + ) + await conn.commit() + await saver._prune_after_put(tid, "") + await conn.commit() + return await self._surviving_ids(conn, tid) + + survivors = _run(_go()) + # anchors = cp_004, cp_005. Walk visits cp_003 (added to + # extra_preserve before deserialize fails). cp_001, cp_002 + # pruned. Survivors: cp_003, cp_004, cp_005. + assert survivors == ["cp_003", "cp_004", "cp_005"] + + def test_anchor_count_below_keep_is_noop(self): + """When checkpoint count < keep_per_ns, prune returns early without DELETE.""" + from EvoScientist.sessions import PruningCheckpointer + + tid = "tdsa_006" + + async def _go(): + import aiosqlite + + async with aiosqlite.connect(self._db_path) as conn: + saver = PruningCheckpointer(conn, keep_per_ns=5) + await saver.setup() + # Only 3 checkpoints; keep=5. anchor_ids has 3 items, + # 3 < 5, prune returns early — all survive untouched. + specs = [ + ("cp_001", None, None), + ("cp_002", "cp_001", None), + ("cp_003", "cp_002", None), + ] + await self._insert_chain(conn, tid, specs) + await saver._prune_after_put(tid, "") + await conn.commit() + return await self._surviving_ids(conn, tid) + + survivors = _run(_go()) + assert survivors == ["cp_001", "cp_002", "cp_003"] + + class TestMigrationSweep(unittest.TestCase): """Tests for the legacy-bloat migration sweep.""" @@ -933,7 +1445,7 @@ class TestMigrationSweep(unittest.TestCase): idx INTEGER NOT NULL, channel TEXT NOT NULL, type TEXT, - blob BLOB, + value BLOB, PRIMARY KEY (thread_id, checkpoint_ns, checkpoint_id, task_id, idx) ) """ @@ -1114,6 +1626,107 @@ class TestMigrationSweep(unittest.TestCase): assert _run(_second()) + def test_sweep_preserves_snapshot_ancestor(self): + """Migration sweep must apply the same DeltaChannel walk as steady-state. + + Without this, legacy users upgrading to PR #231 would hit a + one-shot silent truncation: the bloat sweep would naive-prune + the ``_DeltaSnapshot`` seed out of long threads, then the + ``user_version`` marker locks the sweep so it never re-runs — + leaving permanently empty ``/resume`` history. + + Mirrors ``TestPruningCheckpointerDeltaChannel.test_preserves_ + snapshot_ancestor`` but drives via ``_run_migration_sweep``. + """ + from langgraph.checkpoint.serde.types import _DeltaSnapshot + + from EvoScientist.sessions import _run_migration_sweep + + tid = "tsweep_delta" + + async def _seed(): + import aiosqlite + + serde = JsonPlusSerializer() + snapshot_type, snapshot_blob = serde.dumps_typed( + {"channel_values": {"messages": _DeltaSnapshot(value=[])}} + ) + empty_type, empty_blob = serde.dumps_typed({"channel_values": {}}) + meta = json.dumps({"agent_name": AGENT_NAME}) + + async with aiosqlite.connect(self._db_path) as conn: + await conn.execute( + """ + CREATE TABLE IF NOT EXISTS checkpoints ( + thread_id TEXT NOT NULL, + checkpoint_ns TEXT NOT NULL DEFAULT '', + checkpoint_id TEXT NOT NULL, + parent_checkpoint_id TEXT, + type TEXT, + checkpoint BLOB, + metadata TEXT NOT NULL DEFAULT '{}', + PRIMARY KEY (thread_id, checkpoint_ns, checkpoint_id) + ) + """ + ) + await conn.execute( + """ + CREATE TABLE IF NOT EXISTS writes ( + thread_id TEXT NOT NULL, + checkpoint_ns TEXT NOT NULL DEFAULT '', + checkpoint_id TEXT NOT NULL, + task_id TEXT NOT NULL, + idx INTEGER NOT NULL, + channel TEXT NOT NULL, + type TEXT, + value BLOB, + PRIMARY KEY (thread_id, checkpoint_ns, checkpoint_id, task_id, idx) + ) + """ + ) + # cp_001..cp_002 delta-only, cp_003 carries the snapshot, + # cp_004..cp_010 delta-only. Anchor window with keep=5 is + # cp_006..cp_010; walk from cp_005 backward hits cp_003 + # (seed) → stop. Survivors: cp_003..cp_010. + for i in range(1, 11): + cid = f"cp_{i:03d}" + parent = f"cp_{i - 1:03d}" if i > 1 else None + if i == 3: + ct, cb = snapshot_type, snapshot_blob + else: + ct, cb = empty_type, empty_blob + await conn.execute( + "INSERT INTO checkpoints (thread_id, checkpoint_ns, " + "checkpoint_id, parent_checkpoint_id, type, checkpoint, metadata) " + "VALUES (?, '', ?, ?, ?, ?, ?)", + (tid, cid, parent, ct, cb, meta), + ) + await conn.commit() + + _run(_seed()) + pairs = _run(_run_migration_sweep(keep=5)) + assert pairs == 1 + + async def _survivors(): + import aiosqlite + + async with aiosqlite.connect(self._db_path) as conn: + async with conn.execute( + "SELECT checkpoint_id FROM checkpoints WHERE thread_id = ? " + "ORDER BY checkpoint_id ASC", + (tid,), + ) as cur: + return [r[0] for r in await cur.fetchall()] + + survivors = _run(_survivors()) + # cp_001, cp_002 pruned. cp_003 (snapshot) + walk-through (cp_004, + # cp_005) + anchors (cp_006..cp_010) survive. + assert survivors == [f"cp_{i:03d}" for i in range(3, 11)] + # Explicit absence of the pruned ids — guards against a future + # refactor that accidentally returns an empty survivors list. + assert "cp_001" not in survivors + assert "cp_002" not in survivors + class TestDbStats(unittest.TestCase): """Tests for the read-only ``db_stats`` diagnostic helper.""" @@ -1167,7 +1780,7 @@ class TestDbStats(unittest.TestCase): idx INTEGER NOT NULL, channel TEXT NOT NULL, type TEXT, - blob BLOB, + value BLOB, PRIMARY KEY (thread_id, checkpoint_ns, checkpoint_id, task_id, idx) ) """ @@ -1197,7 +1810,7 @@ class TestDbStats(unittest.TestCase): # (counted by db_stats via the JOIN to checkpoints). for i in range(4): await conn.execute( - "INSERT INTO writes (thread_id, checkpoint_ns, checkpoint_id, task_id, idx, channel, type, blob) " + "INSERT INTO writes (thread_id, checkpoint_ns, checkpoint_id, task_id, idx, channel, type, value) " "VALUES ('evo01', '', 'ce01_0', 't1', ?, 'ch', 'str', X'AA')", (i,), ) @@ -1206,7 +1819,7 @@ class TestDbStats(unittest.TestCase): # checkpoints and filters by agent_name). for i in range(2): await conn.execute( - "INSERT INTO writes (thread_id, checkpoint_ns, checkpoint_id, task_id, idx, channel, type, blob) " + "INSERT INTO writes (thread_id, checkpoint_ns, checkpoint_id, task_id, idx, channel, type, value) " "VALUES ('oth01', '', 'co01_0', 't2', ?, 'ch', 'str', X'BB')", (i,), ) diff --git a/uv.lock b/uv.lock index 199a97f..b6f6d0e 100644 --- a/uv.lock +++ b/uv.lock @@ -794,7 +794,7 @@ wheels = [ [[package]] name = "deepagents" -version = "0.5.7" +version = "0.6.2" source = { registry = "https://pypi.org/simple" } dependencies = [ { name = "langchain" }, @@ -804,9 +804,14 @@ dependencies = [ { name = "langsmith" }, { name = "wcmatch" }, ] -sdist = { url = "https://files.pythonhosted.org/packages/5e/81/7c2285e29816a0f1c33211cda31770ddca3b6bb6baf8d645e4d85f133e7d/deepagents-0.5.7.tar.gz", hash = "sha256:8a9e28f2d2c48b5eb1f659573c95829ac31db563cb223d561b0d1dddce133187", size = 165168, upload-time = "2026-05-05T15:54:23.992Z" } +sdist = { url = "https://files.pythonhosted.org/packages/98/3f/eade375f843b05abdae6c667f341300f13b810b3d69e7f1ef9a093d17bae/deepagents-0.6.2.tar.gz", hash = "sha256:81bc6f69c6f54bd0b54ce929e8299003a604747404e7f9d208a66842c5fbf1f3", size = 179651, upload-time = "2026-05-18T23:57:50.072Z" } wheels = [ - { url = "https://files.pythonhosted.org/packages/fd/bb/c1c68ad0ca96cdddd601e7fca0f274be202411378af6da76a35f0e4a1b66/deepagents-0.5.7-py3-none-any.whl", hash = "sha256:2783f87cfeade4c6fbbb86b427e35bf9d12c6cc29d1f0f2b978cd2529082bb66", size = 187645, upload-time = "2026-05-05T15:54:22.764Z" }, + { url = "https://files.pythonhosted.org/packages/8d/cf/a760686e1bb347287bcb9f741134a088980112a7633d3b2017fb13a8adb0/deepagents-0.6.2-py3-none-any.whl", hash = "sha256:202c42c06b5a7abb762dedba6a57073253096fa28effa6e156db280e10027322", size = 204754, upload-time = "2026-05-18T23:57:48.665Z" }, +] + +[package.optional-dependencies] +quickjs = [ + { name = "langchain-quickjs" }, ] [[package]] @@ -876,7 +881,7 @@ name = "evoscientist" version = "0.1.0" source = { editable = "." } dependencies = [ - { name = "deepagents" }, + { name = "deepagents", extra = ["quickjs"] }, { name = "filelock" }, { name = "httpx" }, { name = "langchain" }, @@ -976,7 +981,7 @@ requires-dist = [ { name = "certifi", marker = "extra == 'wechat'", specifier = ">=2024.0" }, { name = "cryptography", marker = "extra == 'all-channels'", specifier = ">=41.0" }, { name = "cryptography", marker = "extra == 'qq'", specifier = ">=41.0" }, - { name = "deepagents", specifier = ">=0.5.7" }, + { name = "deepagents", extras = ["quickjs"], specifier = "~=0.6.2" }, { name = "discord-py", marker = "extra == 'all-channels'", specifier = ">=2.3" }, { name = "discord-py", marker = "extra == 'discord'", specifier = ">=2.3" }, { name = "faster-whisper", marker = "extra == 'stt'", specifier = ">=1.0" }, @@ -1939,16 +1944,16 @@ wheels = [ [[package]] name = "langchain" -version = "1.2.17" +version = "1.3.1" source = { registry = "https://pypi.org/simple" } dependencies = [ { name = "langchain-core" }, { name = "langgraph" }, { name = "pydantic" }, ] -sdist = { url = "https://files.pythonhosted.org/packages/46/35/322d13339acb61d7a733d03a73a9ade968c64ac0eb982f497d24e22a998f/langchain-1.2.17.tar.gz", hash = "sha256:c30b578c0eebbde8bec9247dbbbae1a791128557b99b65c8be1e007040975d09", size = 577779, upload-time = "2026-04-30T20:25:34.626Z" } +sdist = { url = "https://files.pythonhosted.org/packages/11/e5/6350e77a9e2764eaafcb2d581cbf0b800f53c6bc98fdf5ebc85f3a931ded/langchain-1.3.1.tar.gz", hash = "sha256:bc283c220233230f48b8e50ab1fbf1b688bcb206d933fa448d40a9b143177f62", size = 581329, upload-time = "2026-05-15T18:14:55.368Z" } wheels = [ - { url = "https://files.pythonhosted.org/packages/d1/cf/b183dba8667f7b6d1be546fb8089a3bc3bc12b514f551f5317ae03815770/langchain-1.2.17-py3-none-any.whl", hash = "sha256:ff881cdfbe90e0b6afac42eea7999657c282cc73db059c910d803f4e9f8ff305", size = 113131, upload-time = "2026-04-30T20:25:32.895Z" }, + { url = "https://files.pythonhosted.org/packages/78/11/3d7ed10b535413a07ed5e15682abcb77f3c4204ac49586977a495f9b24e6/langchain-1.3.1-py3-none-any.whl", hash = "sha256:154e9c30c90b391eba4315296f6bf6b6fac6b058ddea4cc771a10470968fe36f", size = 114345, upload-time = "2026-05-15T18:14:53.984Z" }, ] [[package]] @@ -1967,7 +1972,7 @@ wheels = [ [[package]] name = "langchain-core" -version = "1.3.3" +version = "1.4.0" source = { registry = "https://pypi.org/simple" } dependencies = [ { name = "jsonpatch" }, @@ -1980,9 +1985,9 @@ dependencies = [ { name = "typing-extensions" }, { name = "uuid-utils" }, ] -sdist = { url = "https://files.pythonhosted.org/packages/d3/ae/8b74458fc3850ec3d150eb9f45e857db129dafa801fb5cf173dfc9f8bbf3/langchain_core-1.3.3.tar.gz", hash = "sha256:fa510a5db8efdc0c6ff41c0939fb5c00a0183c11f6b84233e892e3227ff69182", size = 915041, upload-time = "2026-05-05T19:02:36.612Z" } +sdist = { url = "https://files.pythonhosted.org/packages/59/de/679a53472c25860837e32c0442c962fa86e95317a36460e2c9d5c91b17c2/langchain_core-1.4.0.tar.gz", hash = "sha256:1dc341eed802ed9c117c0df3923c991e5e9e226571e5725c194eeb5bd93d1a7f", size = 920260, upload-time = "2026-05-11T18:42:35.919Z" } wheels = [ - { url = "https://files.pythonhosted.org/packages/1f/01/4771b7ab2af1d1aba5b710bd8f13d9225c609425214b357590a17b01be77/langchain_core-1.3.3-py3-none-any.whl", hash = "sha256:18aae8506f37da7f74398492279a7d6efcee4f8e23c4c41c7af080eeb7ef7bd1", size = 543857, upload-time = "2026-05-05T19:02:34.52Z" }, + { url = "https://files.pythonhosted.org/packages/0f/1a/86c38c27b81913a1c6c12448cab55defb5a1097c7dc9a4cea83f55477a2d/langchain_core-1.4.0-py3-none-any.whl", hash = "sha256:23cbbdb46e38ddd1dd5247e6167e96013eae74bea4c5949c550809970a9e565c", size = 548120, upload-time = "2026-05-11T18:42:33.992Z" }, ] [[package]] @@ -2081,9 +2086,25 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/c2/e9/06c47ecb2aff08f83dfa30058da3bf86be64862c19569043ed5331bbeecd/langchain_protocol-0.0.14-py3-none-any.whl", hash = "sha256:ffc35089779bd8ca217015180cef5e660fc3b074efdaa0f2e95df73583f1a047", size = 6984, upload-time = "2026-04-29T16:40:17.841Z" }, ] +[[package]] +name = "langchain-quickjs" +version = "0.1.2" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "deepagents" }, + { name = "langchain" }, + { name = "langchain-core" }, + { name = "langgraph" }, + { name = "quickjs-rs" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/d2/f9/0298d3da290498684cad75d595543b7d538b193e1a0efa5a565b511720b1/langchain_quickjs-0.1.2.tar.gz", hash = "sha256:508ee976852a4ca40f7826433989307f7f3d1ae8ec444ba9e6f0d9789ead5ca4", size = 208337, upload-time = "2026-05-11T19:52:37.138Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/0b/40/a9cb0479268474ebcd03df1cafa5be454569855a6e92fc85130f9f3a40dc/langchain_quickjs-0.1.2-py3-none-any.whl", hash = "sha256:78221685eb70a030b57dc6fc74a50ec5c8f730766f8c8521fd9f5ca1b1f8b745", size = 37149, upload-time = "2026-05-11T19:52:35.648Z" }, +] + [[package]] name = "langgraph" -version = "1.1.10" +version = "1.2.0" source = { registry = "https://pypi.org/simple" } dependencies = [ { name = "langchain-core" }, @@ -2093,9 +2114,9 @@ dependencies = [ { name = "pydantic" }, { name = "xxhash" }, ] -sdist = { url = "https://files.pythonhosted.org/packages/9a/b3/7dec224369c7938eb3227ff69542a0d0f517862a0d27945b8c395f2a781f/langgraph-1.1.10.tar.gz", hash = "sha256:3115beb58203283c98d8752a90c034f3432177d2979a1fe205f76e5f1b744500", size = 560685, upload-time = "2026-04-27T17:19:10.426Z" } +sdist = { url = "https://files.pythonhosted.org/packages/58/61/d5d25e783035aa307d289b37e082258a6061c0fb4caa4a284f3bf1e87169/langgraph-1.2.0.tar.gz", hash = "sha256:4a9baaf62afc5d5f63144a50095140a34b9aa9b7cea695d25326d564775348e7", size = 690248, upload-time = "2026-05-12T03:46:39.164Z" } wheels = [ - { url = "https://files.pythonhosted.org/packages/80/07/057dc1aa7991115fca53f1fa6573a7cc0dd296c05360c672cc67fdb6245b/langgraph-1.1.10-py3-none-any.whl", hash = "sha256:8a4f163f72f4401648d0c11b48ee906947d938ba8cf1f474540fe591534f0d17", size = 173750, upload-time = "2026-04-27T17:19:09.073Z" }, + { url = "https://files.pythonhosted.org/packages/f6/e8/e3304ac0015c2bdb04ad9785e4ed65c788855ce7857ce6104dd2f5d322db/langgraph-1.2.0-py3-none-any.whl", hash = "sha256:03fd5895a8d4b70db1ff63ebc3bacead29dd20cd794a8b1a483e7ec9018f7a65", size = 234262, upload-time = "2026-05-12T03:46:37.971Z" }, ] [[package]] @@ -2141,15 +2162,15 @@ wheels = [ [[package]] name = "langgraph-checkpoint" -version = "4.0.3" +version = "4.1.0" source = { registry = "https://pypi.org/simple" } dependencies = [ { name = "langchain-core" }, { name = "ormsgpack" }, ] -sdist = { url = "https://files.pythonhosted.org/packages/7c/e1/885e49cdafceb4c74dae4573bc5dd6054c6c640382ee73104532f33dca46/langgraph_checkpoint-4.0.3.tar.gz", hash = "sha256:a7b5e2ca18fb79b55edf19396d4ee446f8a53dcb7a4ec62ce6f1c7e00bb5af7f", size = 174009, upload-time = "2026-04-27T14:34:02.777Z" } +sdist = { url = "https://files.pythonhosted.org/packages/02/b4/6005c5dd88ad484fe6235d4c43a0d2cee7e91b08ad85a180985c2662df87/langgraph_checkpoint-4.1.0.tar.gz", hash = "sha256:e5bb304e30fc1363ac8fcb5f7dee5ca2185d77fe475b0d01de2c5f91324c2c21", size = 181942, upload-time = "2026-05-12T03:33:49.888Z" } wheels = [ - { url = "https://files.pythonhosted.org/packages/19/ee/ecd3fa2e893746dde3b768daca2a4935208bc77d09445437ccfffb4a8c9b/langgraph_checkpoint-4.0.3-py3-none-any.whl", hash = "sha256:b91b765712a2311a5b198760f714b7ab9b376d01c047ed78d9b9a3e80df802a3", size = 51682, upload-time = "2026-04-27T14:34:01.51Z" }, + { url = "https://files.pythonhosted.org/packages/93/74/d3be2b41955e20ccd624dba5f6fe9d38dcee385ba470a6e13ed86732fc86/langgraph_checkpoint-4.1.0-py3-none-any.whl", hash = "sha256:8bc2a0466a20c38b865ce6671b42093fd5c041133f32351cae4222e0eeaf7fb5", size = 56047, upload-time = "2026-05-12T03:33:48.548Z" }, ] [[package]] @@ -2190,15 +2211,15 @@ inmem = [ [[package]] name = "langgraph-prebuilt" -version = "1.0.13" +version = "1.1.0" source = { registry = "https://pypi.org/simple" } dependencies = [ { name = "langchain-core" }, { name = "langgraph-checkpoint" }, ] -sdist = { url = "https://files.pythonhosted.org/packages/b5/a4/f8ac75fa7c503103f0cf7680944e28bbaaef74c19a8d163d7346869cc369/langgraph_prebuilt-1.0.13.tar.gz", hash = "sha256:ad219782a80e1718e7e7794de49e0ae307111d45cbcffab9a52725a66a609456", size = 172913, upload-time = "2026-04-30T01:48:15.742Z" } +sdist = { url = "https://files.pythonhosted.org/packages/29/66/ed9b93f56bc17ef22d551892f0ac2b225a97fe0fcf23a511b857f70d590b/langgraph_prebuilt-1.1.0.tar.gz", hash = "sha256:3c579cf6eed2d17f9c157c2d0fcaddcd8688524e7022d3b22b37a3bf4589d528", size = 178833, upload-time = "2026-05-12T03:37:49.332Z" } wheels = [ - { url = "https://files.pythonhosted.org/packages/69/ef/5ada0bef4013ef5ae53a0ca1de5736517f1076a54d313f156ca545ec65d5/langgraph_prebuilt-1.0.13-py3-none-any.whl", hash = "sha256:7055e9fad41fbd3593800aed0aea0a6e974b17f33ed51b80d3d3a031212dd7c0", size = 37214, upload-time = "2026-04-30T01:48:14.507Z" }, + { url = "https://files.pythonhosted.org/packages/e9/43/3fe1a700b8490ed02679cdbbc8c915eb23a092faf496c9c1118abcd10be3/langgraph_prebuilt-1.1.0-py3-none-any.whl", hash = "sha256:51e311747d755b751d5c6b39b0c1446124d3a7643d2515017e6714b323508fc9", size = 41043, upload-time = "2026-05-12T03:37:48.007Z" }, ] [[package]] @@ -2234,7 +2255,7 @@ wheels = [ [[package]] name = "langsmith" -version = "0.8.0" +version = "0.8.5" source = { registry = "https://pypi.org/simple" } dependencies = [ { name = "httpx" }, @@ -2247,9 +2268,9 @@ dependencies = [ { name = "xxhash" }, { name = "zstandard" }, ] -sdist = { url = "https://files.pythonhosted.org/packages/a8/64/95f1f013531395f4e8ed73caeee780f65c7c58fe028cb543f8937b45611b/langsmith-0.8.0.tar.gz", hash = "sha256:59fe5b2a56bbbe14a08aa76691f84b49e8675dd21e11b57d80c6db8c08bac2e3", size = 4432996, upload-time = "2026-04-30T22:13:07.341Z" } +sdist = { url = "https://files.pythonhosted.org/packages/17/eb/8883d1158c743d0aac350f09df7880714d27283497e8c80bb9fe3480f165/langsmith-0.8.5.tar.gz", hash = "sha256:3615243d99c12f4047f13042bdc05a373dce232d106a6511b3ca7b48c5af1c2c", size = 4462348, upload-time = "2026-05-15T21:31:41.093Z" } wheels = [ - { url = "https://files.pythonhosted.org/packages/f3/e1/a4be2e696c9473bb53298df398237da5674704d781d4b748ed35aeef592a/langsmith-0.8.0-py3-none-any.whl", hash = "sha256:12cc4bc5622b835a6d841964d6034df3617bdb912dae0c1381fd0a68a9b3a3ef", size = 393268, upload-time = "2026-04-30T22:13:05.56Z" }, + { url = "https://files.pythonhosted.org/packages/23/85/968c88a63e32a59b3e5c68afd2fe114ce0708a125db0be1a85efc25fb2ea/langsmith-0.8.5-py3-none-any.whl", hash = "sha256:efc779f9d450dcaf9d97bc8894f4926276509d6e730e05289af9a64debce06ae", size = 399564, upload-time = "2026-05-15T21:31:39.046Z" }, ] [package.optional-dependencies] @@ -3601,6 +3622,29 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/3c/26/1062c7ec1b053db9e499b4d2d5bc231743201b74051c973dadeac80a8f43/questionary-2.1.1-py3-none-any.whl", hash = "sha256:a51af13f345f1cdea62347589fbb6df3b290306ab8930713bfae4d475a7d4a59", size = 36753, upload-time = "2025-08-28T19:00:19.56Z" }, ] +[[package]] +name = "quickjs-rs" +version = "0.1.2" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/96/59/4c596144ee2dfe49024cd3abf32a97ba62c34cd9f0f72e6fe23021a73181/quickjs_rs-0.1.2.tar.gz", hash = "sha256:95c42dcb40f067ae3b95eb3f79836cc8bcab62d4fd4cb276970ca2a211e687c7", size = 155522, upload-time = "2026-05-05T04:32:48.315Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/f5/cd/0375e8324d581e3c8c54c167359643063cbd6a824a73d0d5d37ad84ecad0/quickjs_rs-0.1.2-cp311-cp311-macosx_10_12_x86_64.whl", hash = "sha256:4a03488c73626a129072eae966bcadede77dcc79cfe36f2ea3b0f1ae351bae75", size = 1358448, upload-time = "2026-05-05T04:32:20.562Z" }, + { url = "https://files.pythonhosted.org/packages/18/c3/f5f98bb3f7e4d8df209c470592d700b54d3f26b91c5b3203133d5940e0c4/quickjs_rs-0.1.2-cp311-cp311-macosx_11_0_arm64.whl", hash = "sha256:bf0941f4d6903ac4ce24d3734f71c4391aff403d34774e66405f25045a538ad5", size = 1271913, upload-time = "2026-05-05T04:32:22.821Z" }, + { url = "https://files.pythonhosted.org/packages/a9/18/4452836b1d94a3e9f4a9294c06ae4c4e6a0e8ec6e4842cf95ac7cc51fec1/quickjs_rs-0.1.2-cp311-cp311-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:71f22f01b871800978b8bc23aeeb9edc4c708ada2552da1bfe15227af1717e11", size = 1296910, upload-time = "2026-05-05T04:32:24.759Z" }, + { url = "https://files.pythonhosted.org/packages/c4/eb/591cd04c38cbc0770a0609ed87fdb1121e38af64b78b507d67d4ad31f7cc/quickjs_rs-0.1.2-cp311-cp311-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:efb68efe12b600e18dc7b8709dfb3bd61acf3246683713757101bec6db2dd5ad", size = 1395102, upload-time = "2026-05-05T04:32:26.838Z" }, + { url = "https://files.pythonhosted.org/packages/1f/d1/4d2cb369278c0b6333b68224d35a37161ba1408dc42588c7033756a62416/quickjs_rs-0.1.2-cp311-cp311-win_amd64.whl", hash = "sha256:93c1ac55e811f582c1b7affd92274ded945e77009010d4a028bfa2d6577b8322", size = 1304595, upload-time = "2026-05-05T04:32:28.686Z" }, + { url = "https://files.pythonhosted.org/packages/67/01/e50258535da0f38167d1df6e8370027a1629f3553ecebe188a501eaee826/quickjs_rs-0.1.2-cp312-cp312-macosx_10_12_x86_64.whl", hash = "sha256:543ae2a983f677f7c83457fbea6fb33abf938215b512439340c90768c15b0b96", size = 1358211, upload-time = "2026-05-05T04:32:30.344Z" }, + { url = "https://files.pythonhosted.org/packages/51/32/6f2c18a39388b1f45fecbc8639053b15b5a7d8d5345dba1886e7ae120dc7/quickjs_rs-0.1.2-cp312-cp312-macosx_11_0_arm64.whl", hash = "sha256:68bb46c3909bd40768da20f0ffbdd1e0da70702e0638711f5e2bb31d778dd92a", size = 1272808, upload-time = "2026-05-05T04:32:32.176Z" }, + { url = "https://files.pythonhosted.org/packages/d0/f2/77f04f177df196674e926eb7a019127cea5551fa8174d621300536888be9/quickjs_rs-0.1.2-cp312-cp312-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:7b3d5bde9f2c12bfce1702f9d1bc2ccfa14f3993211e3d9ed8eceaf7069d432d", size = 1295774, upload-time = "2026-05-05T04:32:33.987Z" }, + { url = "https://files.pythonhosted.org/packages/e2/cb/ee09b2d9797d248ca9c492e88ded527798693bb013f5d921c9fb07229f8f/quickjs_rs-0.1.2-cp312-cp312-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:a7dc9d446b00c24437ae22ba0acefc238abfbdca223716bd8c2ce4223b836d14", size = 1394148, upload-time = "2026-05-05T04:32:35.662Z" }, + { url = "https://files.pythonhosted.org/packages/00/a7/72c50e1f52ac928902657f19ca0f549676e6be84252e5d3fa286cecc6c01/quickjs_rs-0.1.2-cp312-cp312-win_amd64.whl", hash = "sha256:b4229db014881f56828c8ffd120b6be822f2886d570ab1bea99cfe70da83f401", size = 1302547, upload-time = "2026-05-05T04:32:37.29Z" }, + { url = "https://files.pythonhosted.org/packages/81/e2/6497237126abea9f058e72904d406f31d0bde387796f67b13cdb7b86184c/quickjs_rs-0.1.2-cp313-cp313-macosx_10_12_x86_64.whl", hash = "sha256:673a5c9824905754f67f52e7d05a23ca8094828b44cd389dd63976caa109420f", size = 1358790, upload-time = "2026-05-05T04:32:39.492Z" }, + { url = "https://files.pythonhosted.org/packages/5c/55/f5fc4bf2c237a8570ce57ee8f9ddd76bc34b8c4e0013d408bc3a697dea0e/quickjs_rs-0.1.2-cp313-cp313-macosx_11_0_arm64.whl", hash = "sha256:2af1fa9cca2009e8ab69d39d5b4579b59ec69607fb2cf4394cccfab14bd89caa", size = 1273427, upload-time = "2026-05-05T04:32:41.254Z" }, + { url = "https://files.pythonhosted.org/packages/63/bf/88760ad6331a08d1749678d50127d3d7ba7309a815fb4c325d4d6ffa740e/quickjs_rs-0.1.2-cp313-cp313-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:806eec8c3ba699931372e34743b3894795eb382f2ee70321a25795228524840f", size = 1295904, upload-time = "2026-05-05T04:32:43.012Z" }, + { url = "https://files.pythonhosted.org/packages/5e/b9/7f0aef4047a7a2469a9945cdb493b5e8937249b036c8b53699c84eed9e0b/quickjs_rs-0.1.2-cp313-cp313-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:fed6aa295b80f5ebc7b077af37895a70234c566961735ca9c803a3fe8be511ec", size = 1394390, upload-time = "2026-05-05T04:32:45.052Z" }, + { url = "https://files.pythonhosted.org/packages/0c/db/9e3bfba55953fe5af5f39f1144b033fe572ae55a5e1bc68a4408f636cb54/quickjs_rs-0.1.2-cp313-cp313-win_amd64.whl", hash = "sha256:606a7b4da741ceb51e4f79723ad366f9272f69bca4f438b985d1b54c56c8fcc0", size = 1302778, upload-time = "2026-05-05T04:32:46.761Z" }, +] + [[package]] name = "referencing" version = "0.37.0"