diff --git a/examples/draft-revise-loop.py b/examples/draft-revise-loop.py index 28fea5e36..eb90e2872 100644 --- a/examples/draft-revise-loop.py +++ b/examples/draft-revise-loop.py @@ -2,8 +2,8 @@ from __future__ import annotations from langchain.chat_models.openai import ChatOpenAI from langchain.output_parsers.openai_functions import JsonOutputFunctionsParser -from langchain.prompts import SystemMessagePromptTemplate -from langchain.schema.output_parser import StrOutputParser +from langchain_core.output_parsers import StrOutputParser +from langchain_core.prompts import SystemMessagePromptTemplate from permchain import Channel, Pregel diff --git a/examples/rag.py b/examples/rag.py index 14a5a5c50..392288de2 100644 --- a/examples/rag.py +++ b/examples/rag.py @@ -1,8 +1,8 @@ from langchain.chat_models import ChatOpenAI from langchain.embeddings import OpenAIEmbeddings -from langchain.prompts import PromptTemplate -from langchain.schema.messages import AIMessage, AnyMessage, FunctionMessage from langchain.vectorstores import FAISS +from langchain_core.messages import AIMessage, AnyMessage, FunctionMessage +from langchain_core.prompts import PromptTemplate from permchain import Channel, Pregel from permchain.channels import Topic diff --git a/examples/recursive-web-loader.py b/examples/recursive-web-loader.py index 5d4b8fae1..346b890c2 100644 --- a/examples/recursive-web-loader.py +++ b/examples/recursive-web-loader.py @@ -2,9 +2,9 @@ from contextlib import asynccontextmanager, contextmanager from typing import AsyncGenerator, Callable, FrozenSet, Generator, Optional, TypedDict import httpx -from langchain.schema import Document -from langchain.schema.runnable import RunnableLambda, RunnablePassthrough -from langchain.utils.html import extract_sub_links +from langchain_core.documents import Document +from langchain_core.runnables import RunnableLambda, RunnablePassthrough +from langchain_core.utils.html import extract_sub_links from permchain import Channel, Pregel from permchain.channels.context import Context diff --git a/permchain/channels/base.py b/permchain/channels/base.py index bd53e37b4..ab951e42e 100644 --- a/permchain/channels/base.py +++ b/permchain/channels/base.py @@ -14,11 +14,11 @@ from typing import ( from typing_extensions import Self -from permchain.constants import CHECKPOINT_KEY_TS, CHECKPOINT_KEY_VERSION +from permchain.checkpoint.base import Checkpoint Value = TypeVar("Value") Update = TypeVar("Update") -Checkpoint = TypeVar("Checkpoint") +C = TypeVar("C") class EmptyChannelError(Exception): @@ -34,7 +34,7 @@ class InvalidUpdateError(Exception): pass -class BaseChannel(Generic[Value, Update, Checkpoint], ABC): +class BaseChannel(Generic[Value, Update, C], ABC): @property @abstractmethod def ValueType(self) -> Any: @@ -47,14 +47,12 @@ class BaseChannel(Generic[Value, Update, Checkpoint], ABC): @contextmanager @abstractmethod - def empty( - self, checkpoint: Optional[Checkpoint] = None - ) -> Generator[Self, None, None]: + def empty(self, checkpoint: Optional[C] = None) -> Generator[Self, None, None]: """Return a new identical channel, optionally initialized from a checkpoint.""" @asynccontextmanager async def aempty( - self, checkpoint: Optional[Checkpoint] = None + self, checkpoint: Optional[C] = None ) -> AsyncGenerator[Self, None]: """Return a new identical channel, optionally initialized from a checkpoint.""" with self.empty(checkpoint) as value: @@ -74,7 +72,7 @@ class BaseChannel(Generic[Value, Update, Checkpoint], ABC): Raises EmptyChannelError if the channel is empty (never updated yet).""" @abstractmethod - def checkpoint(self) -> Checkpoint | None: + def checkpoint(self) -> C | None: """Return a string representation of the channel's current state. Raises EmptyChannelError if the channel is empty (never updated yet), @@ -84,12 +82,13 @@ class BaseChannel(Generic[Value, Update, Checkpoint], ABC): @contextmanager def ChannelsManager( channels: Mapping[str, BaseChannel], - checkpoint: Optional[Mapping[str, Any]], + checkpoint: Checkpoint, ) -> Generator[Mapping[str, BaseChannel], None, None]: """Manage channels for the lifetime of a Pregel invocation (multiple steps).""" # TODO use https://docs.python.org/3/library/contextlib.html#contextlib.ExitStack - checkpoint = checkpoint or {} - empty = {k: v.empty(checkpoint.get(k)) for k, v in channels.items()} + empty = { + k: v.empty(checkpoint["channel_values"].get(k)) for k, v in channels.items() + } try: yield {k: v.__enter__() for k, v in empty.items()} finally: @@ -100,11 +99,12 @@ def ChannelsManager( @asynccontextmanager async def AsyncChannelsManager( channels: Mapping[str, BaseChannel], - checkpoint: Optional[Mapping[str, Any]], + checkpoint: Checkpoint, ) -> AsyncGenerator[Mapping[str, BaseChannel], None]: """Manage channels for the lifetime of a Pregel invocation (multiple steps).""" - checkpoint = checkpoint or {} - empty = {k: v.aempty(checkpoint.get(k)) for k, v in channels.items()} + empty = { + k: v.aempty(checkpoint["channel_values"].get(k)) for k, v in channels.items() + } try: yield {k: await v.__aenter__() for k, v in empty.items()} finally: @@ -112,15 +112,31 @@ async def AsyncChannelsManager( await v.__aexit__(None, None, None) -def create_checkpoint(channels: Mapping[str, BaseChannel]) -> Mapping[str, Any]: +def create_checkpoint( + checkpoint: Checkpoint, channels: Mapping[str, BaseChannel] +) -> Checkpoint: """Create a checkpoint for the given channels.""" - checkpoint = { - CHECKPOINT_KEY_VERSION: 1, - CHECKPOINT_KEY_TS: datetime.now(timezone.utc).isoformat(), - } + checkpoint = Checkpoint( + v=1, + ts=datetime.now(timezone.utc).isoformat(), + channel_values=checkpoint["channel_values"], + channel_versions=checkpoint["channel_versions"], + versions_seen=checkpoint["versions_seen"], + ) for k, v in channels.items(): try: - checkpoint[k] = v.checkpoint() + checkpoint["channel_values"][k] = v.checkpoint() except EmptyChannelError: pass return checkpoint + + +def channel_values(channels: Mapping[str, BaseChannel]) -> dict[str, Any]: + """Return a dictionary of channel values.""" + values: dict[str, Any] = {} + for k, v in channels.items(): + try: + values[k] = v.get() + except EmptyChannelError: + pass + return values diff --git a/permchain/checkpoint/base.py b/permchain/checkpoint/base.py index 7c4b96384..b43cea377 100644 --- a/permchain/checkpoint/base.py +++ b/permchain/checkpoint/base.py @@ -1,14 +1,35 @@ import asyncio from abc import ABC, abstractmethod -from typing import Any, Mapping +from collections import defaultdict +from datetime import datetime, timezone +from typing import Any, TypedDict -from langchain.load.serializable import Serializable -from langchain.schema.runnable import RunnableConfig -from langchain.schema.runnable.utils import ConfigurableFieldSpec +from langchain_core.load.serializable import Serializable +from langchain_core.pydantic_v1 import Field +from langchain_core.runnables import RunnableConfig +from langchain_core.runnables.utils import ConfigurableFieldSpec from permchain.utils import StrEnum +class Checkpoint(TypedDict): + v: int + ts: str + channel_values: dict[str, Any] + channel_versions: defaultdict[str, int] + versions_seen: defaultdict[str, defaultdict[str, int]] + + +def empty_checkpoint() -> Checkpoint: + return Checkpoint( + v=1, + ts=datetime.now(timezone.utc).isoformat(), + channel_values={}, + channel_versions=defaultdict(int), + versions_seen=defaultdict(lambda: defaultdict(int)), + ) + + class CheckpointAt(StrEnum): END_OF_STEP = "end_of_step" END_OF_RUN = "end_of_run" @@ -22,17 +43,23 @@ class BaseCheckpointAdapter(Serializable, ABC): return [] @abstractmethod - def get(self, config: RunnableConfig) -> Mapping[str, Any] | None: + def get(self, config: RunnableConfig) -> Checkpoint | None: ... @abstractmethod - def put(self, config: RunnableConfig, checkpoint: Mapping[str, Any]) -> None: + def put(self, config: RunnableConfig, checkpoint: Checkpoint) -> None: ... - async def aget(self, config: RunnableConfig) -> Mapping[str, Any] | None: + async def aget(self, config: RunnableConfig) -> Checkpoint | None: return await asyncio.get_running_loop().run_in_executor(None, self.get, config) - async def aput(self, config: RunnableConfig, checkpoint: Mapping[str, Any]) -> None: + async def aput(self, config: RunnableConfig, checkpoint: Checkpoint) -> None: return await asyncio.get_running_loop().run_in_executor( None, self.put, config, checkpoint ) + + +class CheckpointView(Serializable): + values: dict[str, Any] = Field(frozen=True) + + step: int diff --git a/permchain/checkpoint/memory.py b/permchain/checkpoint/memory.py index e036c877b..7b150203f 100644 --- a/permchain/checkpoint/memory.py +++ b/permchain/checkpoint/memory.py @@ -1,14 +1,12 @@ -from typing import Any, Dict, Mapping +from langchain_core.pydantic_v1 import Field +from langchain_core.runnables import RunnableConfig +from langchain_core.runnables.utils import ConfigurableFieldSpec -from langchain.pydantic_v1 import Field -from langchain.schema.runnable import RunnableConfig -from langchain.schema.runnable.utils import ConfigurableFieldSpec - -from permchain.checkpoint.base import BaseCheckpointAdapter +from permchain.checkpoint.base import BaseCheckpointAdapter, Checkpoint class MemoryCheckpoint(BaseCheckpointAdapter): - storage: Dict[str, Mapping[str, Any]] = Field(default_factory=dict) + storage: dict[str, Checkpoint] = Field(default_factory=dict) @property def config_specs(self) -> list[ConfigurableFieldSpec]: @@ -23,8 +21,8 @@ class MemoryCheckpoint(BaseCheckpointAdapter): ), ] - def get(self, config: RunnableConfig) -> Mapping[str, Any] | None: + def get(self, config: RunnableConfig) -> Checkpoint | None: return self.storage.get(config["configurable"]["thread_id"], None) - def put(self, config: RunnableConfig, checkpoint: Mapping[str, Any]) -> None: + def put(self, config: RunnableConfig, checkpoint: Checkpoint) -> None: return self.storage.update({config["configurable"]["thread_id"]: checkpoint}) diff --git a/permchain/constants.py b/permchain/constants.py index bb96e6361..6e5f9c37f 100644 --- a/permchain/constants.py +++ b/permchain/constants.py @@ -1,4 +1,2 @@ CONFIG_KEY_SEND = "__pregel_send" CONFIG_KEY_READ = "__pregel_read" -CHECKPOINT_KEY_VERSION = "__pregel_version" -CHECKPOINT_KEY_TS = "__pregel_ts" diff --git a/permchain/pregel/__init__.py b/permchain/pregel/__init__.py index 3edcc2f14..eff58980b 100644 --- a/permchain/pregel/__init__.py +++ b/permchain/pregel/__init__.py @@ -19,23 +19,23 @@ from typing import ( overload, ) -from langchain.callbacks.manager import ( +from langchain_core.callbacks.manager import ( AsyncCallbackManagerForChainRun, CallbackManagerForChainRun, ) -from langchain.globals import get_debug -from langchain.pydantic_v1 import BaseModel, Field, create_model, root_validator -from langchain.schema.runnable import ( +from langchain_core.globals import get_debug +from langchain_core.pydantic_v1 import BaseModel, Field, create_model, root_validator +from langchain_core.runnables import ( Runnable, RunnableSerializable, ) -from langchain.schema.runnable.base import Input, Output, coerce_to_runnable -from langchain.schema.runnable.config import ( +from langchain_core.runnables.base import Input, Output, coerce_to_runnable +from langchain_core.runnables.config import ( RunnableConfig, get_executor_for_config, patch_config, ) -from langchain.schema.runnable.utils import ( +from langchain_core.runnables.utils import ( ConfigurableFieldSpec, get_unique_config_specs, ) @@ -45,9 +45,17 @@ from permchain.channels.base import ( BaseChannel, ChannelsManager, EmptyChannelError, + channel_values, create_checkpoint, ) -from permchain.checkpoint.base import BaseCheckpointAdapter, CheckpointAt +from permchain.channels.last_value import LastValue +from permchain.checkpoint.base import ( + BaseCheckpointAdapter, + Checkpoint, + CheckpointAt, + CheckpointView, + empty_checkpoint, +) from permchain.constants import CONFIG_KEY_READ, CONFIG_KEY_SEND from permchain.pregel.debug import print_checkpoint, print_step_start from permchain.pregel.io import map_input, map_output @@ -84,7 +92,10 @@ class Channel: @classmethod def subscribe_to( - cls, channels: str | Sequence[str], key: Optional[str] = None + cls, + channels: str | Sequence[str], + key: Optional[str] = None, + when: Callable[[Any], bool] | None = None, ) -> ChannelInvoke: """Runs process.invoke() each time channels are updated, with a dict of the channel values as input.""" @@ -100,6 +111,7 @@ class Channel: else {chan: chan for chan in channels}, ), triggers=[channels] if isinstance(channels, str) else channels, + when=when, ) @classmethod @@ -135,7 +147,7 @@ class Pregel(RunnableSerializable[dict[str, Any] | Any, dict[str, Any] | Any]): debug: bool = Field(default_factory=get_debug) - checkpoint: Optional[BaseCheckpointAdapter] = None + saver: Optional[BaseCheckpointAdapter] = None class Config: arbitrary_types_allowed = True @@ -151,7 +163,7 @@ class Pregel(RunnableSerializable[dict[str, Any] | Any, dict[str, Any] | Any]): def config_specs(self) -> list[ConfigurableFieldSpec]: return get_unique_config_specs( [spec for chain in self.chains.values() for spec in chain.config_specs] - + (self.checkpoint.config_specs if self.checkpoint is not None else []) + + (self.saver.config_specs if self.saver is not None else []) ) @property @@ -194,27 +206,27 @@ class Pregel(RunnableSerializable[dict[str, Any] | Any, dict[str, Any] | Any]): input: Iterator[dict[str, Any] | Any], run_manager: CallbackManagerForChainRun, config: RunnableConfig, - ) -> Iterator[dict[str, Any] | Any]: + ) -> Iterator[tuple[dict[str, Any] | Any, CheckpointView]]: if config["recursion_limit"] < 1: raise ValueError("recursion_limit must be at least 1") + # copy chains to ignore mutations during execution processes = {**self.chains} - checkpoint = ( - self.checkpoint.get(config) if self.checkpoint is not None else None - ) + # get checkpoint from saver, or create an empty one + checkpoint = self.saver.get(config) if self.saver else None + checkpoint = checkpoint or empty_checkpoint() + # create channels from checkpoint with ChannelsManager( self.channels, checkpoint ) as channels, get_executor_for_config(config) as executor: - next_tasks = _apply_writes_and_prepare_next_tasks( - processes, + # map inputs to channel updates + _apply_writes( + checkpoint, channels, deque(w for c in input for w in map_input(self.input, c)), config, 0, ) - if not next_tasks: - return - read = partial(_read_channel, channels) # Similarly to Bulk Synchronous Parallel / Pregel model @@ -223,6 +235,12 @@ class Pregel(RunnableSerializable[dict[str, Any] | Any, dict[str, Any] | Any]): # channels are guaranteed to be immutable for the duration of the step, # with channel updates applied only at the transition between steps for step in range(config["recursion_limit"]): + next_tasks = _prepare_next_tasks(checkpoint, processes, channels) + + # if no more tasks, we're done + if not next_tasks: + break + if self.debug: print_step_start(step, next_tasks) @@ -255,63 +273,55 @@ class Pregel(RunnableSerializable[dict[str, Any] | Any, dict[str, Any] | Any]): # interrupt on failure or timeout _interrupt_or_proceed(done, inflight, step) - # apply writes to channels, decide on next step - - next_tasks = _apply_writes_and_prepare_next_tasks( - processes, channels, pending_writes, config, step + 1 - ) + # apply writes to channels + _apply_writes(checkpoint, channels, pending_writes, config, step + 1) if self.debug: print_checkpoint(step, channels) - # if any write to output channels in this step, yield current value - for output in map_output(self.output, pending_writes, channels): - yield output + # yield current value and checkpoint view + view = CheckpointView( + values=channel_values(channels), + step=step + 1, + ) + yield map_output(self.output, pending_writes, channels), view + # if view was updated, apply writes to channels + _apply_writes_from_view(checkpoint, channels, view) # save end of step checkpoint - if ( - self.checkpoint is not None - and self.checkpoint.at == CheckpointAt.END_OF_STEP - ): - checkpoint = create_checkpoint(channels) - self.checkpoint.put(config, checkpoint) - - # if no more tasks, we're done - if not next_tasks: - break + if self.saver is not None and self.saver.at == CheckpointAt.END_OF_STEP: + checkpoint = create_checkpoint(checkpoint, channels) + self.saver.put(config, checkpoint) # save end of run checkpoint - if ( - self.checkpoint is not None - and self.checkpoint.at == CheckpointAt.END_OF_RUN - ): - checkpoint = create_checkpoint(channels) - self.checkpoint.put(config, checkpoint) + if self.saver is not None and self.saver.at == CheckpointAt.END_OF_RUN: + checkpoint = create_checkpoint(checkpoint, channels) + self.saver.put(config, checkpoint) async def _atransform( self, input: AsyncIterator[dict[str, Any] | Any], run_manager: AsyncCallbackManagerForChainRun, config: RunnableConfig, - ) -> AsyncIterator[dict[str, Any] | Any]: + ) -> AsyncIterator[tuple[dict[str, Any] | Any, CheckpointView]]: if config["recursion_limit"] < 1: raise ValueError("recursion_limit must be at least 1") + # copy chains to ignore mutations during execution processes = {**self.chains} - checkpoint = ( - await self.checkpoint.aget(config) if self.checkpoint is not None else None - ) + # get checkpoint from saver, or create an empty one + checkpoint = await self.saver.aget(config) if self.saver else None + checkpoint = checkpoint or empty_checkpoint() + # create channels from checkpoint async with AsyncChannelsManager(self.channels, checkpoint) as channels: - next_tasks = _apply_writes_and_prepare_next_tasks( - processes, + # map inputs to channel updates + _apply_writes( + checkpoint, channels, deque([w async for c in input for w in map_input(self.input, c)]), config, 0, ) - if not next_tasks: - return - read = partial(_read_channel, channels) # Similarly to Bulk Synchronous Parallel / Pregel model @@ -320,6 +330,12 @@ class Pregel(RunnableSerializable[dict[str, Any] | Any, dict[str, Any] | Any]): # channels are guaranteed to be immutable for the duration of the step, # channel updates being applied only at the transition between steps for step in range(config["recursion_limit"]): + next_tasks = _prepare_next_tasks(checkpoint, processes, channels) + + # if no more tasks, we're done + if not next_tasks: + break + if self.debug: print_step_start(step, next_tasks) @@ -355,37 +371,30 @@ class Pregel(RunnableSerializable[dict[str, Any] | Any, dict[str, Any] | Any]): # interrupt on failure or timeout _interrupt_or_proceed(done, inflight, step) - # apply writes to channels, decide on next step - next_tasks = _apply_writes_and_prepare_next_tasks( - processes, channels, pending_writes, config, step + 1 - ) + # apply writes to channels + _apply_writes(checkpoint, channels, pending_writes, config, step + 1) if self.debug: print_checkpoint(step, channels) - # if any write to output channels in this step, yield current value - for output in map_output(self.output, pending_writes, channels): - yield output + # yield current value and checkpoint view + view = CheckpointView( + values=channel_values(channels), + step=step + 1, + ) + yield map_output(self.output, pending_writes, channels), view + # if view was updated, apply writes to channels + _apply_writes_from_view(checkpoint, channels, view) # save end of step checkpoint - if ( - self.checkpoint is not None - and self.checkpoint.at == CheckpointAt.END_OF_STEP - ): - checkpoint = create_checkpoint(channels) - await self.checkpoint.aput(config, checkpoint) - - # if no more tasks, we're done - if not next_tasks: - break + if self.saver is not None and self.saver.at == CheckpointAt.END_OF_STEP: + checkpoint = create_checkpoint(checkpoint, channels) + await self.saver.aput(config, checkpoint) # save end of run checkpoint - if ( - self.checkpoint is not None - and self.checkpoint.at == CheckpointAt.END_OF_RUN - ): - checkpoint = create_checkpoint(channels) - await self.checkpoint.aput(config, checkpoint) + if self.saver is not None and self.saver.at == CheckpointAt.END_OF_RUN: + checkpoint = create_checkpoint(checkpoint, channels) + await self.saver.aput(config, checkpoint) def invoke( self, @@ -412,9 +421,22 @@ class Pregel(RunnableSerializable[dict[str, Any] | Any, dict[str, Any] | Any]): config: RunnableConfig | None = None, **kwargs: Any | None, ) -> Iterator[dict[str, Any] | Any]: - return self._transform_stream_with_config( + for output, _ in self._transform_stream_with_config( input, self._transform, config, **kwargs - ) + ): + if output is not None: + yield output + + def step( + self, + input: dict[str, Any] | Any, + config: RunnableConfig | None = None, + **kwargs: Any, + ) -> Iterator[tuple[dict[str, Any] | Any, CheckpointView]]: + for tup in self._transform_stream_with_config( + iter([input]), self._transform, config, **kwargs + ): + yield cast(tuple[dict[str, Any] | Any, CheckpointView], tup) async def ainvoke( self, @@ -445,10 +467,25 @@ class Pregel(RunnableSerializable[dict[str, Any] | Any, dict[str, Any] | Any]): config: RunnableConfig | None = None, **kwargs: Any | None, ) -> AsyncIterator[dict[str, Any] | Any]: - async for chunk in self._atransform_stream_with_config( + async for output, _ in self._atransform_stream_with_config( input, self._atransform, config, **kwargs ): - yield chunk + if output is not None: + yield output + + async def astep( + self, + input: dict[str, Any] | Any, + config: RunnableConfig | None = None, + **kwargs: Any, + ) -> AsyncIterator[tuple[dict[str, Any] | Any, CheckpointView]]: + async def input_stream() -> AsyncIterator[dict[str, Any] | Any]: + yield input + + async for tup in self._atransform_stream_with_config( + input_stream(), self._atransform, config, **kwargs + ): + yield cast(tuple[dict[str, Any] | Any, CheckpointView], tup) def _interrupt_or_proceed( @@ -484,13 +521,13 @@ def _read_channel( return None -def _apply_writes_and_prepare_next_tasks( - processes: Mapping[str, ChannelInvoke | ChannelBatch], +def _apply_writes( + checkpoint: Checkpoint, channels: Mapping[str, BaseChannel], pending_writes: Sequence[tuple[str, Any]], config: RunnableConfig, for_step: int, -) -> list[tuple[Runnable, Any, str]]: +) -> None: pending_writes_by_channel: dict[str, list[Any]] = defaultdict(list) # Group writes by channel for chan, val in pending_writes: @@ -508,6 +545,7 @@ def _apply_writes_and_prepare_next_tasks( for chan, vals in pending_writes_by_channel.items(): if chan in channels: channels[chan].update(vals) + checkpoint["channel_versions"][chan] += 1 updated_channels.add(chan) else: logger.warning(f"Skipping write for channel {chan} which has no readers") @@ -516,16 +554,43 @@ def _apply_writes_and_prepare_next_tasks( if chan not in updated_channels: channels[chan].update([]) + +def _apply_writes_from_view( + checkpoint: Checkpoint, + channels: Mapping[str, BaseChannel], + view: CheckpointView, +) -> None: + for chan, value in view.values.items(): + if value == channels[chan].get(): + continue + + assert isinstance(channels[chan], LastValue), ( + f"Can't modify channel {chan} of type " + f"{channels[chan].__class__.__name__}" + ) + checkpoint["channel_versions"][chan] += 1 + channels[chan].update([view.values[chan]]) + + +def _prepare_next_tasks( + checkpoint: Checkpoint, + processes: Mapping[str, ChannelInvoke | ChannelBatch], + channels: Mapping[str, BaseChannel], +) -> list[tuple[Runnable, Any, str]]: tasks: list[tuple[Runnable, Any, str]] = [] # Check if any processes should be run in next step # If so, prepare the values to be passed to them for name, proc in processes.items(): + seen = checkpoint["versions_seen"][name] if isinstance(proc, ChannelInvoke): # If any of the channels read by this process were updated - if any(chan in updated_channels for chan in proc.triggers): + if any( + checkpoint["channel_versions"][chan] > seen[chan] + for chan in proc.triggers + ): # If all channels subscribed by this process have been initialized try: - val = { + val: Any = { k: _read_channel( channels, chan, catch=chan not in proc.triggers ) @@ -539,10 +604,20 @@ def _apply_writes_and_prepare_next_tasks( if list(proc.channels.keys()) == [None]: val = val[None] - tasks.append((proc, val, name)) + # update seen versions + seen.update( + { + chan: checkpoint["channel_versions"][chan] + for chan in proc.triggers + } + ) + + # skip if condition is not met + if proc.when is None or proc.when(val): + tasks.append((proc, val, name)) elif isinstance(proc, ChannelBatch): # If the channel read by this process was updated - if proc.channel in updated_channels: + if checkpoint["channel_versions"][proc.channel] > seen[proc.channel]: # Here we don't catch EmptyChannelError because the channel # must be intialized if the previous `if` condition is true val = channels[proc.channel].get() @@ -550,5 +625,6 @@ def _apply_writes_and_prepare_next_tasks( val = [{proc.key: v} for v in val] tasks.append((proc, val, name)) + seen[proc.channel] = checkpoint["channel_versions"][proc.channel] return tasks diff --git a/permchain/pregel/debug.py b/permchain/pregel/debug.py index 613749858..c1e3cc2be 100644 --- a/permchain/pregel/debug.py +++ b/permchain/pregel/debug.py @@ -1,8 +1,8 @@ from pprint import pformat from typing import Any, Iterator, Mapping -from langchain.schema.runnable import Runnable -from langchain.utils.input import get_bolded_text, get_colored_text +from langchain_core.runnables import Runnable +from langchain_core.utils.input import get_bolded_text, get_colored_text from permchain.channels.base import BaseChannel, EmptyChannelError diff --git a/permchain/pregel/io.py b/permchain/pregel/io.py index 460554877..d49193b49 100644 --- a/permchain/pregel/io.py +++ b/permchain/pregel/io.py @@ -5,10 +5,12 @@ from permchain.pregel.log import logger def map_input( - input_channels: str | Sequence[str], chunk: dict[str, Any] | Any + input_channels: str | Sequence[str], chunk: dict[str, Any] | Any | None ) -> Iterator[tuple[str, Any]]: """Map input chunk to a sequence of pending writes in the form (channel, value).""" - if isinstance(input_channels, str): + if chunk is None: + return + elif isinstance(input_channels, str): yield (input_channels, chunk) else: if not isinstance(chunk, dict): @@ -24,11 +26,12 @@ def map_output( output_channels: str | Sequence[str], pending_writes: Sequence[tuple[str, Any]], channels: Mapping[str, BaseChannel], -) -> Iterator[dict[str, Any] | Any]: +) -> dict[str, Any] | Any | None: """Map pending writes (a sequence of tuples (channel, value)) to output chunk.""" if isinstance(output_channels, str): if any(chan == output_channels for chan, _ in pending_writes): - yield channels[output_channels].get() + return channels[output_channels].get() else: if updated := {c for c, _ in pending_writes if c in output_channels}: - yield {chan: channels[chan].get() for chan in updated} + return {chan: channels[chan].get() for chan in updated} + return None diff --git a/permchain/pregel/read.py b/permchain/pregel/read.py index b8691885a..1c9b1dc13 100644 --- a/permchain/pregel/read.py +++ b/permchain/pregel/read.py @@ -2,7 +2,7 @@ from __future__ import annotations from typing import Any, Callable, List, Mapping, Optional, Sequence -from langchain.pydantic_v1 import Field +from langchain_core.pydantic_v1 import Field from langchain_core.runnables import ( Runnable, RunnableConfig, @@ -70,7 +70,7 @@ class ChannelInvoke(RunnableBindingBase): triggers: List[str] = Field(default_factory=list) - skip: Optional[Callable[[Any], bool]] = None + when: Optional[Callable[[Any], bool]] = None bound: Runnable[Any, Any] = Field(default=default_bound) @@ -80,6 +80,7 @@ class ChannelInvoke(RunnableBindingBase): self, channels: Mapping[None, str] | Mapping[str, str], triggers: Sequence[str], + when: Optional[Callable[[Any], bool]] = None, *, bound: Optional[Runnable[Any, Any]] = None, kwargs: Optional[Mapping[str, Any]] = None, @@ -89,6 +90,7 @@ class ChannelInvoke(RunnableBindingBase): super().__init__( channels=channels, triggers=triggers, + when=when, bound=bound or default_bound, kwargs=kwargs or {}, config=config, @@ -108,6 +110,7 @@ class ChannelInvoke(RunnableBindingBase): **{chan: chan for chan in channels}, }, triggers=self.triggers, + when=self.when, bound=self.bound, kwargs=self.kwargs, config=self.config, @@ -123,6 +126,7 @@ class ChannelInvoke(RunnableBindingBase): return ChannelInvoke( channels=self.channels, triggers=self.triggers, + when=self.when, bound=coerce_to_runnable(other), kwargs=self.kwargs, config=self.config, @@ -131,6 +135,7 @@ class ChannelInvoke(RunnableBindingBase): return ChannelInvoke( channels=self.channels, triggers=self.triggers, + when=self.when, # delegate to __or__ in self.bound bound=self.bound | other, kwargs=self.kwargs, diff --git a/permchain/pregel/validate.py b/permchain/pregel/validate.py index a0483f4cb..5c0dd6718 100644 --- a/permchain/pregel/validate.py +++ b/permchain/pregel/validate.py @@ -2,15 +2,9 @@ from typing import Any, Mapping, Sequence from permchain.channels.base import BaseChannel from permchain.channels.last_value import LastValue -from permchain.constants import CHECKPOINT_KEY_TS, CHECKPOINT_KEY_VERSION from permchain.pregel.read import ChannelBatch, ChannelInvoke from permchain.pregel.reserved import ReservedChannels -FORBIDDEN_CHANNEL_NAMES = { - CHECKPOINT_KEY_TS, - CHECKPOINT_KEY_VERSION, -} - def validate_chains_channels( chains: Mapping[str, ChannelInvoke | ChannelBatch], @@ -55,10 +49,6 @@ def validate_chains_channels( if chan not in channels: channels[chan] = LastValue(Any) # type: ignore[arg-type] - for name in FORBIDDEN_CHANNEL_NAMES: - if name in channels: - raise ValueError(f"Channel name {name} is reserved") - for chan in ReservedChannels: if chan not in channels: channels[chan] = LastValue(Any) # type: ignore[arg-type] diff --git a/permchain/pregel/write.py b/permchain/pregel/write.py index 5174096f2..3c21d39c8 100644 --- a/permchain/pregel/write.py +++ b/permchain/pregel/write.py @@ -2,12 +2,12 @@ from __future__ import annotations from typing import Any, Callable, Sequence -from langchain.schema.runnable import ( +from langchain_core.runnables import ( Runnable, RunnableConfig, RunnablePassthrough, ) -from langchain.schema.runnable.utils import ConfigurableFieldSpec +from langchain_core.runnables.utils import ConfigurableFieldSpec from permchain.constants import CONFIG_KEY_SEND diff --git a/poetry.lock b/poetry.lock index 60bbf53bc..9f474956c 100644 --- a/poetry.lock +++ b/poetry.lock @@ -826,72 +826,73 @@ files = [ [[package]] name = "greenlet" -version = "3.0.1" +version = "3.0.3" description = "Lightweight in-process concurrent programming" optional = false python-versions = ">=3.7" files = [ - {file = "greenlet-3.0.1-cp310-cp310-macosx_10_9_universal2.whl", hash = "sha256:f89e21afe925fcfa655965ca8ea10f24773a1791400989ff32f467badfe4a064"}, - {file = "greenlet-3.0.1-cp310-cp310-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:28e89e232c7593d33cac35425b58950789962011cc274aa43ef8865f2e11f46d"}, - {file = "greenlet-3.0.1-cp310-cp310-manylinux_2_17_ppc64le.manylinux2014_ppc64le.whl", hash = "sha256:b8ba29306c5de7717b5761b9ea74f9c72b9e2b834e24aa984da99cbfc70157fd"}, - {file = "greenlet-3.0.1-cp310-cp310-manylinux_2_17_s390x.manylinux2014_s390x.whl", hash = "sha256:19bbdf1cce0346ef7341705d71e2ecf6f41a35c311137f29b8a2dc2341374565"}, - {file = "greenlet-3.0.1-cp310-cp310-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:599daf06ea59bfedbec564b1692b0166a0045f32b6f0933b0dd4df59a854caf2"}, - {file = "greenlet-3.0.1-cp310-cp310-manylinux_2_24_x86_64.manylinux_2_28_x86_64.whl", hash = "sha256:b641161c302efbb860ae6b081f406839a8b7d5573f20a455539823802c655f63"}, - {file = "greenlet-3.0.1-cp310-cp310-musllinux_1_1_aarch64.whl", hash = "sha256:d57e20ba591727da0c230ab2c3f200ac9d6d333860d85348816e1dca4cc4792e"}, - {file = "greenlet-3.0.1-cp310-cp310-musllinux_1_1_x86_64.whl", hash = "sha256:5805e71e5b570d490938d55552f5a9e10f477c19400c38bf1d5190d760691846"}, - {file = "greenlet-3.0.1-cp310-cp310-win_amd64.whl", hash = "sha256:52e93b28db27ae7d208748f45d2db8a7b6a380e0d703f099c949d0f0d80b70e9"}, - {file = "greenlet-3.0.1-cp311-cp311-macosx_10_9_universal2.whl", hash = "sha256:f7bfb769f7efa0eefcd039dd19d843a4fbfbac52f1878b1da2ed5793ec9b1a65"}, - {file = "greenlet-3.0.1-cp311-cp311-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:91e6c7db42638dc45cf2e13c73be16bf83179f7859b07cfc139518941320be96"}, - {file = "greenlet-3.0.1-cp311-cp311-manylinux_2_17_ppc64le.manylinux2014_ppc64le.whl", hash = "sha256:1757936efea16e3f03db20efd0cd50a1c86b06734f9f7338a90c4ba85ec2ad5a"}, - {file = "greenlet-3.0.1-cp311-cp311-manylinux_2_17_s390x.manylinux2014_s390x.whl", hash = "sha256:19075157a10055759066854a973b3d1325d964d498a805bb68a1f9af4aaef8ec"}, - {file = "greenlet-3.0.1-cp311-cp311-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:e9d21aaa84557d64209af04ff48e0ad5e28c5cca67ce43444e939579d085da72"}, - {file = "greenlet-3.0.1-cp311-cp311-manylinux_2_24_x86_64.manylinux_2_28_x86_64.whl", hash = "sha256:2847e5d7beedb8d614186962c3d774d40d3374d580d2cbdab7f184580a39d234"}, - {file = "greenlet-3.0.1-cp311-cp311-musllinux_1_1_aarch64.whl", hash = "sha256:97e7ac860d64e2dcba5c5944cfc8fa9ea185cd84061c623536154d5a89237884"}, - {file = "greenlet-3.0.1-cp311-cp311-musllinux_1_1_x86_64.whl", hash = "sha256:b2c02d2ad98116e914d4f3155ffc905fd0c025d901ead3f6ed07385e19122c94"}, - {file = "greenlet-3.0.1-cp311-cp311-win_amd64.whl", hash = "sha256:22f79120a24aeeae2b4471c711dcf4f8c736a2bb2fabad2a67ac9a55ea72523c"}, - {file = "greenlet-3.0.1-cp312-cp312-macosx_10_9_universal2.whl", hash = "sha256:100f78a29707ca1525ea47388cec8a049405147719f47ebf3895e7509c6446aa"}, - {file = "greenlet-3.0.1-cp312-cp312-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:60d5772e8195f4e9ebf74046a9121bbb90090f6550f81d8956a05387ba139353"}, - {file = "greenlet-3.0.1-cp312-cp312-manylinux_2_17_ppc64le.manylinux2014_ppc64le.whl", hash = "sha256:daa7197b43c707462f06d2c693ffdbb5991cbb8b80b5b984007de431493a319c"}, - {file = "greenlet-3.0.1-cp312-cp312-manylinux_2_17_s390x.manylinux2014_s390x.whl", hash = "sha256:ea6b8aa9e08eea388c5f7a276fabb1d4b6b9d6e4ceb12cc477c3d352001768a9"}, - {file = "greenlet-3.0.1-cp312-cp312-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:8d11ebbd679e927593978aa44c10fc2092bc454b7d13fdc958d3e9d508aba7d0"}, - {file = "greenlet-3.0.1-cp312-cp312-manylinux_2_24_x86_64.manylinux_2_28_x86_64.whl", hash = "sha256:dbd4c177afb8a8d9ba348d925b0b67246147af806f0b104af4d24f144d461cd5"}, - {file = "greenlet-3.0.1-cp312-cp312-musllinux_1_1_aarch64.whl", hash = "sha256:20107edf7c2c3644c67c12205dc60b1bb11d26b2610b276f97d666110d1b511d"}, - {file = "greenlet-3.0.1-cp312-cp312-musllinux_1_1_x86_64.whl", hash = "sha256:8bef097455dea90ffe855286926ae02d8faa335ed8e4067326257cb571fc1445"}, - {file = "greenlet-3.0.1-cp312-cp312-win_amd64.whl", hash = "sha256:b2d3337dcfaa99698aa2377c81c9ca72fcd89c07e7eb62ece3f23a3fe89b2ce4"}, - {file = "greenlet-3.0.1-cp37-cp37m-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:80ac992f25d10aaebe1ee15df45ca0d7571d0f70b645c08ec68733fb7a020206"}, - {file = "greenlet-3.0.1-cp37-cp37m-manylinux_2_17_ppc64le.manylinux2014_ppc64le.whl", hash = "sha256:337322096d92808f76ad26061a8f5fccb22b0809bea39212cd6c406f6a7060d2"}, - {file = "greenlet-3.0.1-cp37-cp37m-manylinux_2_17_s390x.manylinux2014_s390x.whl", hash = "sha256:b9934adbd0f6e476f0ecff3c94626529f344f57b38c9a541f87098710b18af0a"}, - {file = "greenlet-3.0.1-cp37-cp37m-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:dc4d815b794fd8868c4d67602692c21bf5293a75e4b607bb92a11e821e2b859a"}, - {file = "greenlet-3.0.1-cp37-cp37m-manylinux_2_24_x86_64.manylinux_2_28_x86_64.whl", hash = "sha256:41bdeeb552d814bcd7fb52172b304898a35818107cc8778b5101423c9017b3de"}, - {file = "greenlet-3.0.1-cp37-cp37m-musllinux_1_1_aarch64.whl", hash = "sha256:6e6061bf1e9565c29002e3c601cf68569c450be7fc3f7336671af7ddb4657166"}, - {file = "greenlet-3.0.1-cp37-cp37m-musllinux_1_1_x86_64.whl", hash = "sha256:fa24255ae3c0ab67e613556375a4341af04a084bd58764731972bcbc8baeba36"}, - {file = "greenlet-3.0.1-cp37-cp37m-win32.whl", hash = "sha256:b489c36d1327868d207002391f662a1d163bdc8daf10ab2e5f6e41b9b96de3b1"}, - {file = "greenlet-3.0.1-cp37-cp37m-win_amd64.whl", hash = "sha256:f33f3258aae89da191c6ebaa3bc517c6c4cbc9b9f689e5d8452f7aedbb913fa8"}, - {file = "greenlet-3.0.1-cp38-cp38-macosx_11_0_universal2.whl", hash = "sha256:d2905ce1df400360463c772b55d8e2518d0e488a87cdea13dd2c71dcb2a1fa16"}, - {file = "greenlet-3.0.1-cp38-cp38-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:0a02d259510b3630f330c86557331a3b0e0c79dac3d166e449a39363beaae174"}, - {file = "greenlet-3.0.1-cp38-cp38-manylinux_2_17_ppc64le.manylinux2014_ppc64le.whl", hash = "sha256:55d62807f1c5a1682075c62436702aaba941daa316e9161e4b6ccebbbf38bda3"}, - {file = "greenlet-3.0.1-cp38-cp38-manylinux_2_17_s390x.manylinux2014_s390x.whl", hash = "sha256:3fcc780ae8edbb1d050d920ab44790201f027d59fdbd21362340a85c79066a74"}, - {file = "greenlet-3.0.1-cp38-cp38-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:4eddd98afc726f8aee1948858aed9e6feeb1758889dfd869072d4465973f6bfd"}, - {file = "greenlet-3.0.1-cp38-cp38-manylinux_2_24_x86_64.manylinux_2_28_x86_64.whl", hash = "sha256:eabe7090db68c981fca689299c2d116400b553f4b713266b130cfc9e2aa9c5a9"}, - {file = "greenlet-3.0.1-cp38-cp38-musllinux_1_1_aarch64.whl", hash = "sha256:f2f6d303f3dee132b322a14cd8765287b8f86cdc10d2cb6a6fae234ea488888e"}, - {file = "greenlet-3.0.1-cp38-cp38-musllinux_1_1_x86_64.whl", hash = "sha256:d923ff276f1c1f9680d32832f8d6c040fe9306cbfb5d161b0911e9634be9ef0a"}, - {file = "greenlet-3.0.1-cp38-cp38-win32.whl", hash = "sha256:0b6f9f8ca7093fd4433472fd99b5650f8a26dcd8ba410e14094c1e44cd3ceddd"}, - {file = "greenlet-3.0.1-cp38-cp38-win_amd64.whl", hash = "sha256:990066bff27c4fcf3b69382b86f4c99b3652bab2a7e685d968cd4d0cfc6f67c6"}, - {file = "greenlet-3.0.1-cp39-cp39-macosx_10_9_universal2.whl", hash = "sha256:ce85c43ae54845272f6f9cd8320d034d7a946e9773c693b27d620edec825e376"}, - {file = "greenlet-3.0.1-cp39-cp39-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:89ee2e967bd7ff85d84a2de09df10e021c9b38c7d91dead95b406ed6350c6997"}, - {file = "greenlet-3.0.1-cp39-cp39-manylinux_2_17_ppc64le.manylinux2014_ppc64le.whl", hash = "sha256:87c8ceb0cf8a5a51b8008b643844b7f4a8264a2c13fcbcd8a8316161725383fe"}, - {file = "greenlet-3.0.1-cp39-cp39-manylinux_2_17_s390x.manylinux2014_s390x.whl", hash = "sha256:d6a8c9d4f8692917a3dc7eb25a6fb337bff86909febe2f793ec1928cd97bedfc"}, - {file = "greenlet-3.0.1-cp39-cp39-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:9fbc5b8f3dfe24784cee8ce0be3da2d8a79e46a276593db6868382d9c50d97b1"}, - {file = "greenlet-3.0.1-cp39-cp39-manylinux_2_24_x86_64.manylinux_2_28_x86_64.whl", hash = "sha256:85d2b77e7c9382f004b41d9c72c85537fac834fb141b0296942d52bf03fe4a3d"}, - {file = "greenlet-3.0.1-cp39-cp39-musllinux_1_1_aarch64.whl", hash = "sha256:696d8e7d82398e810f2b3622b24e87906763b6ebfd90e361e88eb85b0e554dc8"}, - {file = "greenlet-3.0.1-cp39-cp39-musllinux_1_1_x86_64.whl", hash = "sha256:329c5a2e5a0ee942f2992c5e3ff40be03e75f745f48847f118a3cfece7a28546"}, - {file = "greenlet-3.0.1-cp39-cp39-win32.whl", hash = "sha256:cf868e08690cb89360eebc73ba4be7fb461cfbc6168dd88e2fbbe6f31812cd57"}, - {file = "greenlet-3.0.1-cp39-cp39-win_amd64.whl", hash = "sha256:ac4a39d1abae48184d420aa8e5e63efd1b75c8444dd95daa3e03f6c6310e9619"}, - {file = "greenlet-3.0.1.tar.gz", hash = "sha256:816bd9488a94cba78d93e1abb58000e8266fa9cc2aa9ccdd6eb0696acb24005b"}, + {file = "greenlet-3.0.3-cp310-cp310-macosx_11_0_universal2.whl", hash = "sha256:9da2bd29ed9e4f15955dd1595ad7bc9320308a3b766ef7f837e23ad4b4aac31a"}, + {file = "greenlet-3.0.3-cp310-cp310-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:d353cadd6083fdb056bb46ed07e4340b0869c305c8ca54ef9da3421acbdf6881"}, + {file = "greenlet-3.0.3-cp310-cp310-manylinux_2_17_ppc64le.manylinux2014_ppc64le.whl", hash = "sha256:dca1e2f3ca00b84a396bc1bce13dd21f680f035314d2379c4160c98153b2059b"}, + {file = "greenlet-3.0.3-cp310-cp310-manylinux_2_17_s390x.manylinux2014_s390x.whl", hash = "sha256:3ed7fb269f15dc662787f4119ec300ad0702fa1b19d2135a37c2c4de6fadfd4a"}, + {file = "greenlet-3.0.3-cp310-cp310-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:dd4f49ae60e10adbc94b45c0b5e6a179acc1736cf7a90160b404076ee283cf83"}, + {file = "greenlet-3.0.3-cp310-cp310-manylinux_2_24_x86_64.manylinux_2_28_x86_64.whl", hash = "sha256:73a411ef564e0e097dbe7e866bb2dda0f027e072b04da387282b02c308807405"}, + {file = "greenlet-3.0.3-cp310-cp310-musllinux_1_1_aarch64.whl", hash = "sha256:7f362975f2d179f9e26928c5b517524e89dd48530a0202570d55ad6ca5d8a56f"}, + {file = "greenlet-3.0.3-cp310-cp310-musllinux_1_1_x86_64.whl", hash = "sha256:649dde7de1a5eceb258f9cb00bdf50e978c9db1b996964cd80703614c86495eb"}, + {file = "greenlet-3.0.3-cp310-cp310-win_amd64.whl", hash = "sha256:68834da854554926fbedd38c76e60c4a2e3198c6fbed520b106a8986445caaf9"}, + {file = "greenlet-3.0.3-cp311-cp311-macosx_11_0_universal2.whl", hash = "sha256:b1b5667cced97081bf57b8fa1d6bfca67814b0afd38208d52538316e9422fc61"}, + {file = "greenlet-3.0.3-cp311-cp311-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:52f59dd9c96ad2fc0d5724107444f76eb20aaccb675bf825df6435acb7703559"}, + {file = "greenlet-3.0.3-cp311-cp311-manylinux_2_17_ppc64le.manylinux2014_ppc64le.whl", hash = "sha256:afaff6cf5200befd5cec055b07d1c0a5a06c040fe5ad148abcd11ba6ab9b114e"}, + {file = "greenlet-3.0.3-cp311-cp311-manylinux_2_17_s390x.manylinux2014_s390x.whl", hash = "sha256:fe754d231288e1e64323cfad462fcee8f0288654c10bdf4f603a39ed923bef33"}, + {file = "greenlet-3.0.3-cp311-cp311-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:2797aa5aedac23af156bbb5a6aa2cd3427ada2972c828244eb7d1b9255846379"}, + {file = "greenlet-3.0.3-cp311-cp311-manylinux_2_24_x86_64.manylinux_2_28_x86_64.whl", hash = "sha256:b7f009caad047246ed379e1c4dbcb8b020f0a390667ea74d2387be2998f58a22"}, + {file = "greenlet-3.0.3-cp311-cp311-musllinux_1_1_aarch64.whl", hash = "sha256:c5e1536de2aad7bf62e27baf79225d0d64360d4168cf2e6becb91baf1ed074f3"}, + {file = "greenlet-3.0.3-cp311-cp311-musllinux_1_1_x86_64.whl", hash = "sha256:894393ce10ceac937e56ec00bb71c4c2f8209ad516e96033e4b3b1de270e200d"}, + {file = "greenlet-3.0.3-cp311-cp311-win_amd64.whl", hash = "sha256:1ea188d4f49089fc6fb283845ab18a2518d279c7cd9da1065d7a84e991748728"}, + {file = "greenlet-3.0.3-cp312-cp312-macosx_11_0_universal2.whl", hash = "sha256:70fb482fdf2c707765ab5f0b6655e9cfcf3780d8d87355a063547b41177599be"}, + {file = "greenlet-3.0.3-cp312-cp312-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:d4d1ac74f5c0c0524e4a24335350edad7e5f03b9532da7ea4d3c54d527784f2e"}, + {file = "greenlet-3.0.3-cp312-cp312-manylinux_2_17_ppc64le.manylinux2014_ppc64le.whl", hash = "sha256:149e94a2dd82d19838fe4b2259f1b6b9957d5ba1b25640d2380bea9c5df37676"}, + {file = "greenlet-3.0.3-cp312-cp312-manylinux_2_17_s390x.manylinux2014_s390x.whl", hash = "sha256:15d79dd26056573940fcb8c7413d84118086f2ec1a8acdfa854631084393efcc"}, + {file = "greenlet-3.0.3-cp312-cp312-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:881b7db1ebff4ba09aaaeae6aa491daeb226c8150fc20e836ad00041bcb11230"}, + {file = "greenlet-3.0.3-cp312-cp312-manylinux_2_24_x86_64.manylinux_2_28_x86_64.whl", hash = "sha256:fcd2469d6a2cf298f198f0487e0a5b1a47a42ca0fa4dfd1b6862c999f018ebbf"}, + {file = "greenlet-3.0.3-cp312-cp312-musllinux_1_1_aarch64.whl", hash = "sha256:1f672519db1796ca0d8753f9e78ec02355e862d0998193038c7073045899f305"}, + {file = "greenlet-3.0.3-cp312-cp312-musllinux_1_1_x86_64.whl", hash = "sha256:2516a9957eed41dd8f1ec0c604f1cdc86758b587d964668b5b196a9db5bfcde6"}, + {file = "greenlet-3.0.3-cp312-cp312-win_amd64.whl", hash = "sha256:bba5387a6975598857d86de9eac14210a49d554a77eb8261cc68b7d082f78ce2"}, + {file = "greenlet-3.0.3-cp37-cp37m-macosx_11_0_universal2.whl", hash = "sha256:5b51e85cb5ceda94e79d019ed36b35386e8c37d22f07d6a751cb659b180d5274"}, + {file = "greenlet-3.0.3-cp37-cp37m-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:daf3cb43b7cf2ba96d614252ce1684c1bccee6b2183a01328c98d36fcd7d5cb0"}, + {file = "greenlet-3.0.3-cp37-cp37m-manylinux_2_17_ppc64le.manylinux2014_ppc64le.whl", hash = "sha256:99bf650dc5d69546e076f413a87481ee1d2d09aaaaaca058c9251b6d8c14783f"}, + {file = "greenlet-3.0.3-cp37-cp37m-manylinux_2_17_s390x.manylinux2014_s390x.whl", hash = "sha256:2dd6e660effd852586b6a8478a1d244b8dc90ab5b1321751d2ea15deb49ed414"}, + {file = "greenlet-3.0.3-cp37-cp37m-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:e3391d1e16e2a5a1507d83e4a8b100f4ee626e8eca43cf2cadb543de69827c4c"}, + {file = "greenlet-3.0.3-cp37-cp37m-manylinux_2_24_x86_64.manylinux_2_28_x86_64.whl", hash = "sha256:e1f145462f1fa6e4a4ae3c0f782e580ce44d57c8f2c7aae1b6fa88c0b2efdb41"}, + {file = "greenlet-3.0.3-cp37-cp37m-musllinux_1_1_aarch64.whl", hash = "sha256:1a7191e42732df52cb5f39d3527217e7ab73cae2cb3694d241e18f53d84ea9a7"}, + {file = "greenlet-3.0.3-cp37-cp37m-musllinux_1_1_x86_64.whl", hash = "sha256:0448abc479fab28b00cb472d278828b3ccca164531daab4e970a0458786055d6"}, + {file = "greenlet-3.0.3-cp37-cp37m-win32.whl", hash = "sha256:b542be2440edc2d48547b5923c408cbe0fc94afb9f18741faa6ae970dbcb9b6d"}, + {file = "greenlet-3.0.3-cp37-cp37m-win_amd64.whl", hash = "sha256:01bc7ea167cf943b4c802068e178bbf70ae2e8c080467070d01bfa02f337ee67"}, + {file = "greenlet-3.0.3-cp38-cp38-macosx_11_0_universal2.whl", hash = "sha256:1996cb9306c8595335bb157d133daf5cf9f693ef413e7673cb07e3e5871379ca"}, + {file = "greenlet-3.0.3-cp38-cp38-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:3ddc0f794e6ad661e321caa8d2f0a55ce01213c74722587256fb6566049a8b04"}, + {file = "greenlet-3.0.3-cp38-cp38-manylinux_2_17_ppc64le.manylinux2014_ppc64le.whl", hash = "sha256:c9db1c18f0eaad2f804728c67d6c610778456e3e1cc4ab4bbd5eeb8e6053c6fc"}, + {file = "greenlet-3.0.3-cp38-cp38-manylinux_2_17_s390x.manylinux2014_s390x.whl", hash = "sha256:7170375bcc99f1a2fbd9c306f5be8764eaf3ac6b5cb968862cad4c7057756506"}, + {file = "greenlet-3.0.3-cp38-cp38-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:6b66c9c1e7ccabad3a7d037b2bcb740122a7b17a53734b7d72a344ce39882a1b"}, + {file = "greenlet-3.0.3-cp38-cp38-manylinux_2_24_x86_64.manylinux_2_28_x86_64.whl", hash = "sha256:098d86f528c855ead3479afe84b49242e174ed262456c342d70fc7f972bc13c4"}, + {file = "greenlet-3.0.3-cp38-cp38-musllinux_1_1_aarch64.whl", hash = "sha256:81bb9c6d52e8321f09c3d165b2a78c680506d9af285bfccbad9fb7ad5a5da3e5"}, + {file = "greenlet-3.0.3-cp38-cp38-musllinux_1_1_x86_64.whl", hash = "sha256:fd096eb7ffef17c456cfa587523c5f92321ae02427ff955bebe9e3c63bc9f0da"}, + {file = "greenlet-3.0.3-cp38-cp38-win32.whl", hash = "sha256:d46677c85c5ba00a9cb6f7a00b2bfa6f812192d2c9f7d9c4f6a55b60216712f3"}, + {file = "greenlet-3.0.3-cp38-cp38-win_amd64.whl", hash = "sha256:419b386f84949bf0e7c73e6032e3457b82a787c1ab4a0e43732898a761cc9dbf"}, + {file = "greenlet-3.0.3-cp39-cp39-macosx_11_0_universal2.whl", hash = "sha256:da70d4d51c8b306bb7a031d5cff6cc25ad253affe89b70352af5f1cb68e74b53"}, + {file = "greenlet-3.0.3-cp39-cp39-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:086152f8fbc5955df88382e8a75984e2bb1c892ad2e3c80a2508954e52295257"}, + {file = "greenlet-3.0.3-cp39-cp39-manylinux_2_17_ppc64le.manylinux2014_ppc64le.whl", hash = "sha256:d73a9fe764d77f87f8ec26a0c85144d6a951a6c438dfe50487df5595c6373eac"}, + {file = "greenlet-3.0.3-cp39-cp39-manylinux_2_17_s390x.manylinux2014_s390x.whl", hash = "sha256:b7dcbe92cc99f08c8dd11f930de4d99ef756c3591a5377d1d9cd7dd5e896da71"}, + {file = "greenlet-3.0.3-cp39-cp39-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:1551a8195c0d4a68fac7a4325efac0d541b48def35feb49d803674ac32582f61"}, + {file = "greenlet-3.0.3-cp39-cp39-manylinux_2_24_x86_64.manylinux_2_28_x86_64.whl", hash = "sha256:64d7675ad83578e3fc149b617a444fab8efdafc9385471f868eb5ff83e446b8b"}, + {file = "greenlet-3.0.3-cp39-cp39-musllinux_1_1_aarch64.whl", hash = "sha256:b37eef18ea55f2ffd8f00ff8fe7c8d3818abd3e25fb73fae2ca3b672e333a7a6"}, + {file = "greenlet-3.0.3-cp39-cp39-musllinux_1_1_x86_64.whl", hash = "sha256:77457465d89b8263bca14759d7c1684df840b6811b2499838cc5b040a8b5b113"}, + {file = "greenlet-3.0.3-cp39-cp39-win32.whl", hash = "sha256:57e8974f23e47dac22b83436bdcf23080ade568ce77df33159e019d161ce1d1e"}, + {file = "greenlet-3.0.3-cp39-cp39-win_amd64.whl", hash = "sha256:c5ee858cfe08f34712f548c3c363e807e7186f03ad7a5039ebadb29e8c6be067"}, + {file = "greenlet-3.0.3.tar.gz", hash = "sha256:43374442353259554ce33599da8b692d5aa96f8976d567d4badf263371fbe491"}, ] [package.extras] -docs = ["Sphinx"] +docs = ["Sphinx", "furo"] test = ["objgraph", "psutil"] [[package]] @@ -3070,60 +3071,60 @@ files = [ [[package]] name = "sqlalchemy" -version = "2.0.23" +version = "2.0.24" description = "Database Abstraction Library" optional = false python-versions = ">=3.7" files = [ - {file = "SQLAlchemy-2.0.23-cp310-cp310-macosx_10_9_x86_64.whl", hash = "sha256:638c2c0b6b4661a4fd264f6fb804eccd392745c5887f9317feb64bb7cb03b3ea"}, - {file = "SQLAlchemy-2.0.23-cp310-cp310-macosx_11_0_arm64.whl", hash = "sha256:e3b5036aa326dc2df50cba3c958e29b291a80f604b1afa4c8ce73e78e1c9f01d"}, - {file = "SQLAlchemy-2.0.23-cp310-cp310-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:787af80107fb691934a01889ca8f82a44adedbf5ef3d6ad7d0f0b9ac557e0c34"}, - {file = "SQLAlchemy-2.0.23-cp310-cp310-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:c14eba45983d2f48f7546bb32b47937ee2cafae353646295f0e99f35b14286ab"}, - {file = "SQLAlchemy-2.0.23-cp310-cp310-musllinux_1_1_aarch64.whl", hash = "sha256:0666031df46b9badba9bed00092a1ffa3aa063a5e68fa244acd9f08070e936d3"}, - {file = "SQLAlchemy-2.0.23-cp310-cp310-musllinux_1_1_x86_64.whl", hash = "sha256:89a01238fcb9a8af118eaad3ffcc5dedaacbd429dc6fdc43fe430d3a941ff965"}, - {file = "SQLAlchemy-2.0.23-cp310-cp310-win32.whl", hash = "sha256:cabafc7837b6cec61c0e1e5c6d14ef250b675fa9c3060ed8a7e38653bd732ff8"}, - {file = "SQLAlchemy-2.0.23-cp310-cp310-win_amd64.whl", hash = "sha256:87a3d6b53c39cd173990de2f5f4b83431d534a74f0e2f88bd16eabb5667e65c6"}, - {file = "SQLAlchemy-2.0.23-cp311-cp311-macosx_10_9_x86_64.whl", hash = "sha256:d5578e6863eeb998980c212a39106ea139bdc0b3f73291b96e27c929c90cd8e1"}, - {file = "SQLAlchemy-2.0.23-cp311-cp311-macosx_11_0_arm64.whl", hash = "sha256:62d9e964870ea5ade4bc870ac4004c456efe75fb50404c03c5fd61f8bc669a72"}, - {file = "SQLAlchemy-2.0.23-cp311-cp311-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:c80c38bd2ea35b97cbf7c21aeb129dcbebbf344ee01a7141016ab7b851464f8e"}, - {file = "SQLAlchemy-2.0.23-cp311-cp311-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:75eefe09e98043cff2fb8af9796e20747ae870c903dc61d41b0c2e55128f958d"}, - {file = "SQLAlchemy-2.0.23-cp311-cp311-musllinux_1_1_aarch64.whl", hash = "sha256:bd45a5b6c68357578263d74daab6ff9439517f87da63442d244f9f23df56138d"}, - {file = "SQLAlchemy-2.0.23-cp311-cp311-musllinux_1_1_x86_64.whl", hash = "sha256:a86cb7063e2c9fb8e774f77fbf8475516d270a3e989da55fa05d08089d77f8c4"}, - {file = "SQLAlchemy-2.0.23-cp311-cp311-win32.whl", hash = "sha256:b41f5d65b54cdf4934ecede2f41b9c60c9f785620416e8e6c48349ab18643855"}, - {file = "SQLAlchemy-2.0.23-cp311-cp311-win_amd64.whl", hash = "sha256:9ca922f305d67605668e93991aaf2c12239c78207bca3b891cd51a4515c72e22"}, - {file = "SQLAlchemy-2.0.23-cp312-cp312-macosx_10_9_x86_64.whl", hash = "sha256:d0f7fb0c7527c41fa6fcae2be537ac137f636a41b4c5a4c58914541e2f436b45"}, - {file = "SQLAlchemy-2.0.23-cp312-cp312-macosx_11_0_arm64.whl", hash = "sha256:7c424983ab447dab126c39d3ce3be5bee95700783204a72549c3dceffe0fc8f4"}, - {file = "SQLAlchemy-2.0.23-cp312-cp312-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:f508ba8f89e0a5ecdfd3761f82dda2a3d7b678a626967608f4273e0dba8f07ac"}, - {file = "SQLAlchemy-2.0.23-cp312-cp312-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:6463aa765cf02b9247e38b35853923edbf2f6fd1963df88706bc1d02410a5577"}, - {file = "SQLAlchemy-2.0.23-cp312-cp312-musllinux_1_1_aarch64.whl", hash = "sha256:e599a51acf3cc4d31d1a0cf248d8f8d863b6386d2b6782c5074427ebb7803bda"}, - {file = "SQLAlchemy-2.0.23-cp312-cp312-musllinux_1_1_x86_64.whl", hash = "sha256:fd54601ef9cc455a0c61e5245f690c8a3ad67ddb03d3b91c361d076def0b4c60"}, - {file = "SQLAlchemy-2.0.23-cp312-cp312-win32.whl", hash = "sha256:42d0b0290a8fb0165ea2c2781ae66e95cca6e27a2fbe1016ff8db3112ac1e846"}, - {file = "SQLAlchemy-2.0.23-cp312-cp312-win_amd64.whl", hash = "sha256:227135ef1e48165f37590b8bfc44ed7ff4c074bf04dc8d6f8e7f1c14a94aa6ca"}, - {file = "SQLAlchemy-2.0.23-cp37-cp37m-macosx_10_9_x86_64.whl", hash = "sha256:14aebfe28b99f24f8a4c1346c48bc3d63705b1f919a24c27471136d2f219f02d"}, - {file = "SQLAlchemy-2.0.23-cp37-cp37m-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:3e983fa42164577d073778d06d2cc5d020322425a509a08119bdcee70ad856bf"}, - {file = "SQLAlchemy-2.0.23-cp37-cp37m-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:7e0dc9031baa46ad0dd5a269cb7a92a73284d1309228be1d5935dac8fb3cae24"}, - {file = "SQLAlchemy-2.0.23-cp37-cp37m-musllinux_1_1_aarch64.whl", hash = "sha256:5f94aeb99f43729960638e7468d4688f6efccb837a858b34574e01143cf11f89"}, - {file = "SQLAlchemy-2.0.23-cp37-cp37m-musllinux_1_1_x86_64.whl", hash = "sha256:63bfc3acc970776036f6d1d0e65faa7473be9f3135d37a463c5eba5efcdb24c8"}, - {file = "SQLAlchemy-2.0.23-cp37-cp37m-win32.whl", hash = "sha256:f48ed89dd11c3c586f45e9eec1e437b355b3b6f6884ea4a4c3111a3358fd0c18"}, - {file = "SQLAlchemy-2.0.23-cp37-cp37m-win_amd64.whl", hash = "sha256:1e018aba8363adb0599e745af245306cb8c46b9ad0a6fc0a86745b6ff7d940fc"}, - {file = "SQLAlchemy-2.0.23-cp38-cp38-macosx_10_9_x86_64.whl", hash = "sha256:64ac935a90bc479fee77f9463f298943b0e60005fe5de2aa654d9cdef46c54df"}, - {file = "SQLAlchemy-2.0.23-cp38-cp38-macosx_11_0_arm64.whl", hash = "sha256:c4722f3bc3c1c2fcc3702dbe0016ba31148dd6efcd2a2fd33c1b4897c6a19693"}, - {file = "SQLAlchemy-2.0.23-cp38-cp38-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:4af79c06825e2836de21439cb2a6ce22b2ca129bad74f359bddd173f39582bf5"}, - {file = "SQLAlchemy-2.0.23-cp38-cp38-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:683ef58ca8eea4747737a1c35c11372ffeb84578d3aab8f3e10b1d13d66f2bc4"}, - {file = "SQLAlchemy-2.0.23-cp38-cp38-musllinux_1_1_aarch64.whl", hash = "sha256:d4041ad05b35f1f4da481f6b811b4af2f29e83af253bf37c3c4582b2c68934ab"}, - {file = "SQLAlchemy-2.0.23-cp38-cp38-musllinux_1_1_x86_64.whl", hash = "sha256:aeb397de65a0a62f14c257f36a726945a7f7bb60253462e8602d9b97b5cbe204"}, - {file = "SQLAlchemy-2.0.23-cp38-cp38-win32.whl", hash = "sha256:42ede90148b73fe4ab4a089f3126b2cfae8cfefc955c8174d697bb46210c8306"}, - {file = "SQLAlchemy-2.0.23-cp38-cp38-win_amd64.whl", hash = "sha256:964971b52daab357d2c0875825e36584d58f536e920f2968df8d581054eada4b"}, - {file = "SQLAlchemy-2.0.23-cp39-cp39-macosx_10_9_x86_64.whl", hash = "sha256:616fe7bcff0a05098f64b4478b78ec2dfa03225c23734d83d6c169eb41a93e55"}, - {file = "SQLAlchemy-2.0.23-cp39-cp39-macosx_11_0_arm64.whl", hash = "sha256:0e680527245895aba86afbd5bef6c316831c02aa988d1aad83c47ffe92655e74"}, - {file = "SQLAlchemy-2.0.23-cp39-cp39-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:9585b646ffb048c0250acc7dad92536591ffe35dba624bb8fd9b471e25212a35"}, - {file = "SQLAlchemy-2.0.23-cp39-cp39-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:4895a63e2c271ffc7a81ea424b94060f7b3b03b4ea0cd58ab5bb676ed02f4221"}, - {file = "SQLAlchemy-2.0.23-cp39-cp39-musllinux_1_1_aarch64.whl", hash = "sha256:cc1d21576f958c42d9aec68eba5c1a7d715e5fc07825a629015fe8e3b0657fb0"}, - {file = "SQLAlchemy-2.0.23-cp39-cp39-musllinux_1_1_x86_64.whl", hash = "sha256:967c0b71156f793e6662dd839da54f884631755275ed71f1539c95bbada9aaab"}, - {file = "SQLAlchemy-2.0.23-cp39-cp39-win32.whl", hash = "sha256:0a8c6aa506893e25a04233bc721c6b6cf844bafd7250535abb56cb6cc1368884"}, - {file = "SQLAlchemy-2.0.23-cp39-cp39-win_amd64.whl", hash = "sha256:f3420d00d2cb42432c1d0e44540ae83185ccbbc67a6054dcc8ab5387add6620b"}, - {file = "SQLAlchemy-2.0.23-py3-none-any.whl", hash = "sha256:31952bbc527d633b9479f5f81e8b9dfada00b91d6baba021a869095f1a97006d"}, - {file = "SQLAlchemy-2.0.23.tar.gz", hash = "sha256:c1bda93cbbe4aa2aa0aa8655c5aeda505cd219ff3e8da91d1d329e143e4aff69"}, + {file = "SQLAlchemy-2.0.24-cp310-cp310-macosx_10_9_x86_64.whl", hash = "sha256:5f801d85ba4753d4ed97181d003e5d3fa330ac7c4587d131f61d7f968f416862"}, + {file = "SQLAlchemy-2.0.24-cp310-cp310-macosx_11_0_arm64.whl", hash = "sha256:b35c35e3923ade1e7ac44e150dec29f5863513246c8bf85e2d7d313e3832bcfb"}, + {file = "SQLAlchemy-2.0.24-cp310-cp310-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:1d9b3fd5eca3c0b137a5e0e468e24ca544ed8ca4783e0e55341b7ed2807518ee"}, + {file = "SQLAlchemy-2.0.24-cp310-cp310-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:7a6209e689d0ff206c40032b6418e3cfcfc5af044b3f66e381d7f1ae301544b4"}, + {file = "SQLAlchemy-2.0.24-cp310-cp310-musllinux_1_1_aarch64.whl", hash = "sha256:37e89d965b52e8b20571b5d44f26e2124b26ab63758bf1b7598a0e38fb2c4005"}, + {file = "SQLAlchemy-2.0.24-cp310-cp310-musllinux_1_1_x86_64.whl", hash = "sha256:c6910eb4ea90c0889f363965cd3c8c45a620ad27b526a7899f0054f6c1b9219e"}, + {file = "SQLAlchemy-2.0.24-cp310-cp310-win32.whl", hash = "sha256:d8e7e8a150e7b548e7ecd6ebb9211c37265991bf2504297d9454e01b58530fc6"}, + {file = "SQLAlchemy-2.0.24-cp310-cp310-win_amd64.whl", hash = "sha256:396f05c552f7fa30a129497c41bef5b4d1423f9af8fe4df0c3dcd38f3e3b9a14"}, + {file = "SQLAlchemy-2.0.24-cp311-cp311-macosx_10_9_x86_64.whl", hash = "sha256:adbd67dac4ebf54587198b63cd30c29fd7eafa8c0cab58893d9419414f8efe4b"}, + {file = "SQLAlchemy-2.0.24-cp311-cp311-macosx_11_0_arm64.whl", hash = "sha256:a0f611b431b84f55779cbb7157257d87b4a2876b067c77c4f36b15e44ced65e2"}, + {file = "SQLAlchemy-2.0.24-cp311-cp311-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:56a0e90a959e18ac5f18c80d0cad9e90cb09322764f536e8a637426afb1cae2f"}, + {file = "SQLAlchemy-2.0.24-cp311-cp311-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:6db686a1d9f183c639f7e06a2656af25d4ed438eda581de135d15569f16ace33"}, + {file = "SQLAlchemy-2.0.24-cp311-cp311-musllinux_1_1_aarch64.whl", hash = "sha256:f0cc0b486a56dff72dddae6b6bfa7ff201b0eeac29d4bc6f0e9725dc3c360d71"}, + {file = "SQLAlchemy-2.0.24-cp311-cp311-musllinux_1_1_x86_64.whl", hash = "sha256:4a1d4856861ba9e73bac05030cec5852eabfa9ef4af8e56c19d92de80d46fc34"}, + {file = "SQLAlchemy-2.0.24-cp311-cp311-win32.whl", hash = "sha256:a3c2753bf4f48b7a6024e5e8a394af49b1b12c817d75d06942cae03d14ff87b3"}, + {file = "SQLAlchemy-2.0.24-cp311-cp311-win_amd64.whl", hash = "sha256:38732884eabc64982a09a846bacf085596ff2371e4e41d20c0734f7e50525d01"}, + {file = "SQLAlchemy-2.0.24-cp312-cp312-macosx_10_9_x86_64.whl", hash = "sha256:9f992e0f916201731993eab8502912878f02287d9f765ef843677ff118d0e0b1"}, + {file = "SQLAlchemy-2.0.24-cp312-cp312-macosx_11_0_arm64.whl", hash = "sha256:2587e108463cc2e5b45a896b2e7cc8659a517038026922a758bde009271aed11"}, + {file = "SQLAlchemy-2.0.24-cp312-cp312-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:0bb7cedcddffca98c40bb0becd3423e293d1fef442b869da40843d751785beb3"}, + {file = "SQLAlchemy-2.0.24-cp312-cp312-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:83fa6df0e035689df89ff77a46bf8738696785d3156c2c61494acdcddc75c69d"}, + {file = "SQLAlchemy-2.0.24-cp312-cp312-musllinux_1_1_aarch64.whl", hash = "sha256:cc889fda484d54d0b31feec409406267616536d048a450fc46943e152700bb79"}, + {file = "SQLAlchemy-2.0.24-cp312-cp312-musllinux_1_1_x86_64.whl", hash = "sha256:57ef6f2cb8b09a042d0dbeaa46a30f2df5dd1e1eb889ba258b0d5d7d6011b81c"}, + {file = "SQLAlchemy-2.0.24-cp312-cp312-win32.whl", hash = "sha256:ea490564435b5b204d8154f0e18387b499ea3cedc1e6af3b3a2ab18291d85aa7"}, + {file = "SQLAlchemy-2.0.24-cp312-cp312-win_amd64.whl", hash = "sha256:ccfd336f96d4c9bbab0309f2a565bf15c468c2d8b2d277a32f89c5940f71fcf9"}, + {file = "SQLAlchemy-2.0.24-cp37-cp37m-macosx_10_9_x86_64.whl", hash = "sha256:9aaaaa846b10dfbe1bda71079d0e31a7e2cebedda9409fa7dba3dfed1ae803e8"}, + {file = "SQLAlchemy-2.0.24-cp37-cp37m-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:95bae3d38f8808d79072da25d5e5a6095f36fe1f9d6c614dd72c59ca8397c7c0"}, + {file = "SQLAlchemy-2.0.24-cp37-cp37m-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:a04191a7c8d77e63f6fc1e8336d6c6e93176c0c010833e74410e647f0284f5a1"}, + {file = "SQLAlchemy-2.0.24-cp37-cp37m-musllinux_1_1_aarch64.whl", hash = "sha256:acc58b7c2e40235712d857fdfc8f2bda9608f4a850d8d9ac0dd1fc80939ca6ac"}, + {file = "SQLAlchemy-2.0.24-cp37-cp37m-musllinux_1_1_x86_64.whl", hash = "sha256:00d76fe5d7cdb5d84d625ce002ce29fefba0bfd98e212ae66793fed30af73931"}, + {file = "SQLAlchemy-2.0.24-cp37-cp37m-win32.whl", hash = "sha256:29e51f848f843bbd75d74ae64ab1ab06302cb1dccd4549d1f5afe6b4a946edb2"}, + {file = "SQLAlchemy-2.0.24-cp37-cp37m-win_amd64.whl", hash = "sha256:e9d036e343a604db3f5a6c33354018a84a1d3f6dcae3673358b404286204798c"}, + {file = "SQLAlchemy-2.0.24-cp38-cp38-macosx_10_9_x86_64.whl", hash = "sha256:9bafaa05b19dc07fa191c1966c5e852af516840b0d7b46b7c3303faf1a349bc9"}, + {file = "SQLAlchemy-2.0.24-cp38-cp38-macosx_11_0_arm64.whl", hash = "sha256:e69290b921b7833c04206f233d6814c60bee1d135b09f5ae5d39229de9b46cd4"}, + {file = "SQLAlchemy-2.0.24-cp38-cp38-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:e8398593ccc4440ce6dffcc4f47d9b2d72b9fe7112ac12ea4a44e7d4de364db1"}, + {file = "SQLAlchemy-2.0.24-cp38-cp38-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:f073321a79c81e1a009218a21089f61d87ee5fa3c9563f6be94f8b41ff181812"}, + {file = "SQLAlchemy-2.0.24-cp38-cp38-musllinux_1_1_aarch64.whl", hash = "sha256:9036ebfd934813990c5b9f71f297e77ed4963720db7d7ceec5a3fdb7cd2ef6ce"}, + {file = "SQLAlchemy-2.0.24-cp38-cp38-musllinux_1_1_x86_64.whl", hash = "sha256:fcf84fe93397a0f67733aa2a38ed4eab9fc6348189fc950e656e1ea198f45668"}, + {file = "SQLAlchemy-2.0.24-cp38-cp38-win32.whl", hash = "sha256:6f5e75de91c754365c098ac08c13fdb267577ce954fa239dd49228b573ca88d7"}, + {file = "SQLAlchemy-2.0.24-cp38-cp38-win_amd64.whl", hash = "sha256:9f29c7f0f4b42337ec5a779e166946a9f86d7d56d827e771b69ecbdf426124ac"}, + {file = "SQLAlchemy-2.0.24-cp39-cp39-macosx_10_9_x86_64.whl", hash = "sha256:07cc423892f2ceda9ae1daa28c0355757f362ecc7505b1ab1a3d5d8dc1c44ac6"}, + {file = "SQLAlchemy-2.0.24-cp39-cp39-macosx_11_0_arm64.whl", hash = "sha256:2a479aa1ab199178ff1956b09ca8a0693e70f9c762875d69292d37049ffd0d8f"}, + {file = "SQLAlchemy-2.0.24-cp39-cp39-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:9b8d0e8578e7f853f45f4512b5c920f6a546cd4bed44137460b2a56534644205"}, + {file = "SQLAlchemy-2.0.24-cp39-cp39-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:e17e7e27af178d31b436dda6a596703b02a89ba74a15e2980c35ecd9909eea3a"}, + {file = "SQLAlchemy-2.0.24-cp39-cp39-musllinux_1_1_aarch64.whl", hash = "sha256:1ca7903d5e7db791a355b579c690684fac6304478b68efdc7f2ebdcfe770d8d7"}, + {file = "SQLAlchemy-2.0.24-cp39-cp39-musllinux_1_1_x86_64.whl", hash = "sha256:db09e424d7bb89b6215a184ca93b4f29d7f00ea261b787918a1af74143b98c06"}, + {file = "SQLAlchemy-2.0.24-cp39-cp39-win32.whl", hash = "sha256:a5cd7d30e47f87b21362beeb3e86f1b5886e7d9b0294b230dde3d3f4a1591375"}, + {file = "SQLAlchemy-2.0.24-cp39-cp39-win_amd64.whl", hash = "sha256:7ae5d44517fe81079ce75cf10f96978284a6db2642c5932a69c82dbae09f009a"}, + {file = "SQLAlchemy-2.0.24-py3-none-any.whl", hash = "sha256:8f358f5cfce04417b6ff738748ca4806fe3d3ae8040fb4e6a0c9a6973ccf9b6e"}, + {file = "SQLAlchemy-2.0.24.tar.gz", hash = "sha256:6db97656fd3fe3f7e5b077f12fa6adb5feb6e0b567a3e99f47ecf5f7ea0a09e3"}, ] [package.dependencies] @@ -3133,7 +3134,7 @@ typing-extensions = ">=4.2.0" [package.extras] aiomysql = ["aiomysql (>=0.2.0)", "greenlet (!=0.4.17)"] aioodbc = ["aioodbc", "greenlet (!=0.4.17)"] -aiosqlite = ["aiosqlite", "greenlet (!=0.4.17)", "typing-extensions (!=3.10.0.1)"] +aiosqlite = ["aiosqlite", "greenlet (!=0.4.17)", "typing_extensions (!=3.10.0.1)"] asyncio = ["greenlet (!=0.4.17)"] asyncmy = ["asyncmy (>=0.2.3,!=0.2.4,!=0.2.6)", "greenlet (!=0.4.17)"] mariadb-connector = ["mariadb (>=1.0.1,!=1.1.2,!=1.1.5)"] @@ -3143,7 +3144,7 @@ mssql-pyodbc = ["pyodbc"] mypy = ["mypy (>=0.910)"] mysql = ["mysqlclient (>=1.4.0)"] mysql-connector = ["mysql-connector-python"] -oracle = ["cx-oracle (>=8)"] +oracle = ["cx_oracle (>=8)"] oracle-oracledb = ["oracledb (>=1.0.1)"] postgresql = ["psycopg2 (>=2.7)"] postgresql-asyncpg = ["asyncpg", "greenlet (!=0.4.17)"] @@ -3153,7 +3154,7 @@ postgresql-psycopg2binary = ["psycopg2-binary"] postgresql-psycopg2cffi = ["psycopg2cffi"] postgresql-psycopgbinary = ["psycopg[binary] (>=3.0.7)"] pymysql = ["pymysql"] -sqlcipher = ["sqlcipher3-binary"] +sqlcipher = ["sqlcipher3_binary"] [[package]] name = "stack-data" @@ -3598,4 +3599,4 @@ testing = ["big-O", "jaraco.functools", "jaraco.itertools", "more-itertools", "p [metadata] lock-version = "2.0" python-versions = ">=3.8.1,<4.0" -content-hash = "31c0918244f47b18e634bb487f32918ff5fde3ecc84dcd416b955b1077ffe663" +content-hash = "2777c34c9c1b5c056206649d287e08e55671a647721374fb7c3b6edc711e796e" diff --git a/pyproject.toml b/pyproject.toml index 43d8ba149..40536ecf9 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -9,7 +9,7 @@ repository = "https://www.github.com/langchain-ai/permchain" [tool.poetry.dependencies] python = ">=3.8.1,<4.0" -langchain = "^0.0.352" +langchain-core = "^0.1.3" [tool.poetry.group.test.dependencies] @@ -37,6 +37,7 @@ optional = true [tool.poetry.group.dev.dependencies] jupyter = "^1.0.0" openai = "^0.27.8" +langchain = "^0.0.352" [tool.ruff] select = [ "E", "F", "I" ] diff --git a/tests/test_pregel.py b/tests/test_pregel.py index 4d5c32f74..aad297018 100644 --- a/tests/test_pregel.py +++ b/tests/test_pregel.py @@ -5,7 +5,7 @@ from contextlib import contextmanager from typing import Generator import pytest -from langchain.schema.runnable import RunnablePassthrough +from langchain_core.runnables import RunnablePassthrough from pytest_mock import MockerFixture from permchain import Channel, Pregel @@ -139,17 +139,51 @@ def test_invoke_single_process_in_dict_out_dict(mocker: MockerFixture) -> None: def test_invoke_two_processes_in_out(mocker: MockerFixture) -> None: add_one = mocker.Mock(side_effect=lambda x: x + 1) chain_one = Channel.subscribe_to("input") | add_one | Channel.write_to("inbox") - chain_two = ( - Channel.subscribe_to_each("inbox") | add_one | Channel.write_to("output") - ) + chain_two = Channel.subscribe_to("inbox") | add_one | Channel.write_to("output") app = Pregel( chains={"chain_one": chain_one, "chain_two": chain_two}, - channels={"inbox": Topic(int)}, ) assert app.invoke(2) == 4 + for output, view in app.step(2): + if view.step == 1: + assert view.values == { + "inbox": 3, + "input": 2, + "is_last_step": False, + } + assert output is None + elif view.step == 2: + assert view.values == { + "output": 4, + "inbox": 3, + "input": 2, + "is_last_step": False, + } + assert output == 4 + + for output, view in app.step(2): + if view.step == 1: + assert view.values == { + "inbox": 3, + "input": 2, + "is_last_step": False, + } + assert output is None + # modify inbox value + view.values["inbox"] = 5 + elif view.step == 2: + assert view.values == { + "output": 6, + "inbox": 5, + "input": 2, + "is_last_step": False, + } + # output is different now + assert output == 6 + def test_invoke_two_processes_in_dict_out(mocker: MockerFixture) -> None: add_one = mocker.Mock(side_effect=lambda x: x + 1) @@ -286,34 +320,34 @@ def test_invoke_checkpoint(mocker: MockerFixture) -> None: app = Pregel( chains={"chain_one": chain_one}, channels={"total": BinaryOperatorAggregate(int, operator.add)}, - checkpoint=memory, + saver=memory, ) # total starts out as 0, so output is 0+2=2 assert app.invoke(2, {"configurable": {"thread_id": "1"}}) == 2 checkpoint = memory.get({"configurable": {"thread_id": "1"}}) assert checkpoint is not None - assert checkpoint.get("total") == 2 + assert checkpoint["channel_values"].get("total") == 2 # total is now 2, so output is 2+3=5 assert app.invoke(3, {"configurable": {"thread_id": "1"}}) == 5 checkpoint = memory.get({"configurable": {"thread_id": "1"}}) assert checkpoint is not None - assert checkpoint.get("total") == 7 + assert checkpoint["channel_values"].get("total") == 7 # total is now 2+5=7, so output would be 7+4=11, but raises ValueError with pytest.raises(ValueError): app.invoke(4, {"configurable": {"thread_id": "1"}}) # checkpoint is not updated checkpoint = memory.get({"configurable": {"thread_id": "1"}}) assert checkpoint is not None - assert checkpoint.get("total") == 7 + assert checkpoint["channel_values"].get("total") == 7 # on a new thread, total starts out as 0, so output is 0+5=5 assert app.invoke(5, {"configurable": {"thread_id": "2"}}) == 5 checkpoint = memory.get({"configurable": {"thread_id": "1"}}) assert checkpoint is not None - assert checkpoint.get("total") == 7 + assert checkpoint["channel_values"].get("total") == 7 checkpoint = memory.get({"configurable": {"thread_id": "2"}}) assert checkpoint is not None - assert checkpoint.get("total") == 5 + assert checkpoint["channel_values"].get("total") == 5 def test_invoke_two_processes_two_in_join_two_out(mocker: MockerFixture) -> None: diff --git a/tests/test_pregel_async.py b/tests/test_pregel_async.py index 8a59fa4a0..b32dc17d2 100644 --- a/tests/test_pregel_async.py +++ b/tests/test_pregel_async.py @@ -4,7 +4,7 @@ from contextlib import asynccontextmanager, contextmanager from typing import Any, AsyncGenerator, AsyncIterator, Generator import pytest -from langchain.schema.runnable import RunnablePassthrough +from langchain_core.runnables import RunnablePassthrough from pytest_mock import MockerFixture from permchain import Channel, Pregel @@ -144,17 +144,51 @@ async def test_invoke_single_process_in_dict_out_dict(mocker: MockerFixture) -> async def test_invoke_two_processes_in_out(mocker: MockerFixture) -> None: add_one = mocker.Mock(side_effect=lambda x: x + 1) chain_one = Channel.subscribe_to("input") | add_one | Channel.write_to("inbox") - chain_two = ( - Channel.subscribe_to_each("inbox") | add_one | Channel.write_to("output") - ) + chain_two = Channel.subscribe_to("inbox") | add_one | Channel.write_to("output") app = Pregel( chains={"chain_one": chain_one, "chain_two": chain_two}, - channels={"inbox": Topic(int)}, ) assert await app.ainvoke(2) == 4 + async for output, view in app.astep(2): + if view.step == 1: + assert view.values == { + "inbox": 3, + "input": 2, + "is_last_step": False, + } + assert output is None + elif view.step == 2: + assert view.values == { + "output": 4, + "inbox": 3, + "input": 2, + "is_last_step": False, + } + assert output == 4 + + async for output, view in app.astep(2): + if view.step == 1: + assert view.values == { + "inbox": 3, + "input": 2, + "is_last_step": False, + } + assert output is None + # modify inbox value + view.values["inbox"] = 5 + elif view.step == 2: + assert view.values == { + "output": 6, + "inbox": 5, + "input": 2, + "is_last_step": False, + } + # output is different now + assert output == 6 + async def test_invoke_two_processes_in_dict_out(mocker: MockerFixture) -> None: add_one = mocker.Mock(side_effect=lambda x: x + 1) @@ -299,34 +333,34 @@ async def test_invoke_checkpoint(mocker: MockerFixture) -> None: app = Pregel( chains={"chain_one": chain_one}, channels={"total": BinaryOperatorAggregate(int, operator.add)}, - checkpoint=memory, + saver=memory, ) # total starts out as 0, so output is 0+2=2 assert await app.ainvoke(2, {"configurable": {"thread_id": "1"}}) == 2 checkpoint = await memory.aget({"configurable": {"thread_id": "1"}}) assert checkpoint is not None - assert checkpoint.get("total") == 2 + assert checkpoint["channel_values"].get("total") == 2 # total is now 2, so output is 2+3=5 assert await app.ainvoke(3, {"configurable": {"thread_id": "1"}}) == 5 checkpoint = await memory.aget({"configurable": {"thread_id": "1"}}) assert checkpoint is not None - assert checkpoint.get("total") == 7 + assert checkpoint["channel_values"].get("total") == 7 # total is now 2+5=7, so output would be 7+4=11, but raises ValueError with pytest.raises(ValueError): await app.ainvoke(4, {"configurable": {"thread_id": "1"}}) # checkpoint is not updated checkpoint = await memory.aget({"configurable": {"thread_id": "1"}}) assert checkpoint is not None - assert checkpoint.get("total") == 7 + assert checkpoint["channel_values"].get("total") == 7 # on a new thread, total starts out as 0, so output is 0+5=5 assert await app.ainvoke(5, {"configurable": {"thread_id": "2"}}) == 5 checkpoint = await memory.aget({"configurable": {"thread_id": "1"}}) assert checkpoint is not None - assert checkpoint.get("total") == 7 + assert checkpoint["channel_values"].get("total") == 7 checkpoint = await memory.aget({"configurable": {"thread_id": "2"}}) assert checkpoint is not None - assert checkpoint.get("total") == 5 + assert checkpoint["channel_values"].get("total") == 5 async def test_invoke_two_processes_two_in_join_two_out(mocker: MockerFixture) -> None: