2e5db60dc5
The runtime (llm/runtime.py, stream/stop.py) is a separate development line whose stop adapter drives LangGraph internals. Upstream v0.3.0 bumps its dependencies for currency, and two exact-version assertions in that adapter turned the bump into a silent regression: stop ownership was refused, so cancellation/continuation runs never reached a terminal state. Resolved without touching the adapter's logic: - stream/stop.py: claim ownership by *capability* instead of an exact version string. The internals the adapter swaps (_graph_aiter / _pump_cond / _exhausted / _aborting / _anext_task / _mux) and the SQLite saver's connection lock are present and identical in langgraph 1.2.6 and 1.2.11, and langgraph-checkpoint-sqlite 3.1.1 exposes the same barrier as 3.0.3. A new patch release can no longer disable stop ownership by being newer; a release that really drops the internals still fails closed with CHECKPOINT_STOP_ADAPTER_UNSUPPORTED. - EvoScientist.py: supply TodoListMiddleware only when deepagents' own default chain lacks it. deepagents 0.7 dropped it (upstream adds one back); 0.6.x still ships it, and a second instance collides by name in langchain's create_agent. - backends.py: fall back to a shape-compatible DeleteResult when deepagents has no delete support, so upstream v0.3.0's delete refusals import and run on either line. - tests/test_backends.py: gate the delete-behaviour tests on the framework actually providing backend deletion instead of asserting a specific stack. Verified: 4216 passed / 33 skipped / 29 failed / 16 errors — every remaining failure is pre-existing on the untouched pre-merge tree except two (a google-stream cleanup-order assertion and one webui launcher test).
146 lines
5.4 KiB
Python
146 lines
5.4 KiB
Python
"""Instance-local stop adapter for LangGraph 1.2.x, SQLite saver 3.0.x/3.1.x.
|
|
|
|
No detached writers or foreign update_state/raw SQL writers are supported.
|
|
aiosqlite 0.22.1's FIFO commit is the connection barrier, not its saver lock.
|
|
|
|
Ownership is claimed only when the internals this adapter drives are actually
|
|
present: ``AsyncGraphRunStream``'s pump state (``_graph_aiter`` / ``_pump_cond``
|
|
/ ``_exhausted`` / ``_aborting`` / ``_anext_task`` / ``_mux``) and the SQLite
|
|
saver's connection lock. Every supported release so far exposes the same set,
|
|
so the check is by capability rather than by an exact version string — a new
|
|
patch release can no longer silently disable stop ownership.
|
|
"""
|
|
import asyncio
|
|
from importlib.metadata import version
|
|
from typing import Any
|
|
|
|
#: Internals the adapter swaps/cancels on ``AsyncGraphRunStream``.
|
|
_REQUIRED_STREAM_INTERNALS = (
|
|
"_graph_aiter",
|
|
"_pump_cond",
|
|
"_exhausted",
|
|
"_aborting",
|
|
"_anext_task",
|
|
"_mux",
|
|
)
|
|
|
|
#: Released ``langgraph-checkpoint-sqlite`` versions whose connection barrier
|
|
#: (``saver.lock`` + the aiosqlite FIFO commit) this adapter has been verified
|
|
#: against.
|
|
_SUPPORTED_CHECKPOINT_SQLITE = ("3.0.3", "3.1.1")
|
|
|
|
|
|
class StreamStop:
|
|
def __init__(self):
|
|
self.graph: Any = None
|
|
self.stream = None
|
|
self.pulls = set()
|
|
self.exits = set()
|
|
self.errors = []
|
|
self.task = None
|
|
self.unsupported = False
|
|
self.exit_tasks_observed = 0
|
|
self.unconfirmed_exit = asyncio.Event()
|
|
|
|
def attach(self, stream):
|
|
from langgraph.stream.run_stream import AsyncGraphRunStream
|
|
if (
|
|
version("langgraph").split(".")[:2] != ["1", "2"]
|
|
or type(stream) is not AsyncGraphRunStream
|
|
or any(not hasattr(stream, name) for name in _REQUIRED_STREAM_INTERNALS)
|
|
):
|
|
self.unsupported = True
|
|
raise RuntimeError("CHECKPOINT_STOP_ADAPTER_UNSUPPORTED")
|
|
self.stream = stream
|
|
self.graph = stream._graph_aiter
|
|
stream._graph_aiter = self
|
|
|
|
def __aiter__(self):
|
|
return self
|
|
|
|
def _record(self, exc):
|
|
self.exits.update(arg for arg in exc.args if isinstance(arg, asyncio.Future))
|
|
|
|
async def __anext__(self):
|
|
task = asyncio.current_task()
|
|
self.pulls.add(task)
|
|
try:
|
|
return await self.graph.__anext__()
|
|
except BaseException as exc:
|
|
self._record(exc)
|
|
raise
|
|
finally:
|
|
self.pulls.discard(task)
|
|
|
|
async def _retry_close(self, close):
|
|
while True:
|
|
try:
|
|
await close()
|
|
return
|
|
except BaseException as exc:
|
|
self._record(exc)
|
|
self.errors.append(exc)
|
|
await self._settle_exits()
|
|
await asyncio.sleep(0.05)
|
|
|
|
async def _settle_exits(self):
|
|
while self.exits:
|
|
tasks = tuple(self.exits)
|
|
self.exit_tasks_observed += len(tasks)
|
|
results = await asyncio.gather(*tasks, return_exceptions=True)
|
|
self.exits.difference_update(tasks)
|
|
for result in results:
|
|
if isinstance(result, BaseException):
|
|
self._record(result)
|
|
self.errors.append(result)
|
|
# A failed/cancelled stack exit is not executor settlement.
|
|
# There is no supported recovery handle for that case.
|
|
await self.unconfirmed_exit.wait()
|
|
|
|
def start(self, producers):
|
|
if self.task is None:
|
|
self.task = asyncio.create_task(self._close(producers), name="evo-v3-stop")
|
|
return self.task
|
|
|
|
async def _close(self, producers):
|
|
if self.unsupported:
|
|
await self.unconfirmed_exit.wait()
|
|
stream = self.stream
|
|
if stream is None:
|
|
return
|
|
# Freeze pumping before cancellation can discard the dependency's pull.
|
|
async with stream._pump_cond:
|
|
stream._exhausted = True
|
|
stream._aborting = True
|
|
pulls = set(self.pulls)
|
|
if stream._anext_task is not None:
|
|
pulls.add(stream._anext_task)
|
|
stream._pump_cond.notify_all()
|
|
for task in pulls:
|
|
if not task.done():
|
|
task.cancel()
|
|
results = await asyncio.gather(*pulls, return_exceptions=True)
|
|
for result in results:
|
|
if isinstance(result, BaseException):
|
|
self._record(result)
|
|
await self._settle_exits()
|
|
await self._retry_close(self.graph.aclose)
|
|
await self._settle_exits()
|
|
await self._retry_close(stream._mux.aclose)
|
|
for task in producers:
|
|
if not task.done():
|
|
task.cancel()
|
|
await asyncio.gather(*producers, return_exceptions=True)
|
|
|
|
|
|
async def sqlite_barrier(saver):
|
|
from langgraph.checkpoint.sqlite.aio import AsyncSqliteSaver
|
|
import aiosqlite
|
|
if (type(saver) is not AsyncSqliteSaver or type(saver.conn) is not aiosqlite.Connection
|
|
or version("langgraph-checkpoint-sqlite") not in _SUPPORTED_CHECKPOINT_SQLITE
|
|
or version("aiosqlite") != "0.22.1"):
|
|
raise RuntimeError("CHECKPOINT_STOP_ADAPTER_UNSUPPORTED")
|
|
async with saver.lock:
|
|
# Worker executes even cancelled queued futures. Await a later commit
|
|
# after graph/executor/repair writers settle, covering those operations.
|
|
await saver.conn.commit() |