From 560e746ceae75e32f88f1e75117284ccdbaa4a8f Mon Sep 17 00:00:00 2001 From: jinliyl <6469360+jinliyl@users.noreply.github.com> Date: Tue, 16 Sep 2025 16:49:56 +0800 Subject: [PATCH] reformat code, support flowllm 0.1.9 (#26) * reformat code, support flowllm 0.1.9 * Update README.md add FLOW_USE_FRAMEWORK=true * Update README_ZH.md add FLOW_USE_FRAMEWORK=true * Update index.md add FLOW_USE_FRAMEWORK=true * add env FLOW_APP_NAME=ReMe --- README.md | 4 +-- README_ZH.md | 4 +-- docs/index.md | 4 +-- docs/sop_memory/making_sop_memories.md | 2 +- example.env | 7 ++-- pyproject.toml | 4 +-- reme_ai/__init__.py | 2 +- reme_ai/app.py | 10 +++--- reme_ai/config/default.yaml | 34 +------------------ reme_ai/react/simple_react_op.py | 2 +- reme_ai/retrieve/personal/extract_time_op.py | 4 +-- reme_ai/retrieve/personal/fuse_rerank_op.py | 4 +-- reme_ai/retrieve/personal/print_memory_op.py | 4 +-- reme_ai/retrieve/personal/read_message_op.py | 4 +-- .../retrieve/personal/retrieve_memory_op.py | 4 +-- reme_ai/retrieve/personal/semantic_rank_op.py | 4 +-- reme_ai/retrieve/personal/set_query_op.py | 4 +-- reme_ai/retrieve/task/build_query_op.py | 4 +-- reme_ai/retrieve/task/merge_memory_op.py | 4 +-- reme_ai/retrieve/task/rerank_memory_op.py | 4 +-- reme_ai/retrieve/task/rewrite_memory_op.py | 4 +-- ...y => agentscope_runtime_memory_service.py} | 6 +++- reme_ai/service/personal_memory_service.py | 31 +++++++---------- reme_ai/service/task_memory_service.py | 32 +++++++---------- reme_ai/summary/personal/contra_repeat_op.py | 6 ++-- .../summary/personal/get_observation_op.py | 6 ++-- .../personal/get_observation_with_time_op.py | 6 ++-- .../personal/get_reflection_subject_op.py | 4 +-- reme_ai/summary/personal/info_filter_op.py | 6 ++-- .../summary/personal/load_today_memory_op.py | 4 +-- .../summary/personal/long_contra_repeat_op.py | 6 ++-- reme_ai/summary/personal/update_insight_op.py | 4 +-- .../summary/task/comparative_extraction_op.py | 4 +-- reme_ai/summary/task/failure_extraction_op.py | 4 +-- .../summary/task/memory_deduplication_op.py | 4 +-- reme_ai/summary/task/memory_validation_op.py | 4 +-- .../task/simple_comparative_summary_op.py | 4 +-- reme_ai/summary/task/simple_summary_op.py | 4 +-- reme_ai/summary/task/success_extraction_op.py | 4 +-- .../summary/task/trajectory_preprocess_op.py | 4 +-- .../task/trajectory_segmentation_op.py | 4 +-- reme_ai/vector_store/delete_memory_op.py | 4 +-- .../vector_store/recall_vector_store_op.py | 4 +-- reme_ai/vector_store/update_memory_freq_op.py | 4 +-- .../vector_store/update_memory_utility_op.py | 4 +-- .../vector_store/update_vector_store_op.py | 4 +-- .../vector_store/vector_store_action_op.py | 4 +-- 47 files changed, 118 insertions(+), 170 deletions(-) rename reme_ai/service/{base_memory_service.py => agentscope_runtime_memory_service.py} (93%) diff --git a/README.md b/README.md index 3306fcda..486b9027 100644 --- a/README.md +++ b/README.md @@ -100,11 +100,9 @@ pip install . Copy `example.env` to .env and modify the corresponding parameters: ```bash -# Required: LLM API Configuration +FLOW_APP_NAME=ReMe FLOW_LLM_API_KEY=sk-xxxx FLOW_LLM_BASE_URL=https://xxxx/v1 - -# Required: Embedding Model Configuration FLOW_EMBEDDING_API_KEY=sk-xxxx FLOW_EMBEDDING_BASE_URL=https://xxxx/v1 diff --git a/README_ZH.md b/README_ZH.md index 824d5db4..b7ea5355 100644 --- a/README_ZH.md +++ b/README_ZH.md @@ -89,11 +89,9 @@ pip install . 复制 `example.env` 为 .env并修改其中对应参数: ```bash -# 必需:LLM API配置 +FLOW_APP_NAME=ReMe FLOW_LLM_API_KEY=sk-xxxx FLOW_LLM_BASE_URL=https://xxxx/v1 - -# 必需:嵌入模型配置 FLOW_EMBEDDING_API_KEY=sk-xxxx FLOW_EMBEDDING_BASE_URL=https://xxxx/v1 diff --git a/docs/index.md b/docs/index.md index c0ca34c3..1e23ba22 100644 --- a/docs/index.md +++ b/docs/index.md @@ -105,11 +105,9 @@ pip install . Copy `example.env` to .env and modify the corresponding parameters: ```bash -# Required: LLM API Configuration +FLOW_APP_NAME=ReMe FLOW_LLM_API_KEY=sk-xxxx FLOW_LLM_BASE_URL=https://xxxx/v1 - -# Required: Embedding Model Configuration FLOW_EMBEDDING_API_KEY=sk-xxxx FLOW_EMBEDDING_BASE_URL=https://xxxx/v1 diff --git a/docs/sop_memory/making_sop_memories.md b/docs/sop_memory/making_sop_memories.md index 0970160a..f6d97d6f 100644 --- a/docs/sop_memory/making_sop_memories.md +++ b/docs/sop_memory/making_sop_memories.md @@ -21,7 +21,7 @@ framework. Each operation (Op) needs to define the following core attributes: ```python -class BaseOp: +class BaseAsyncToolOp: description: str # Description of the operation input_schema: Dict[str, ParamAttr] # Input parameter schema definition output_schema: Dict[str, ParamAttr] # Output parameter schema definition diff --git a/example.env b/example.env index d655b58b..589365fe 100644 --- a/example.env +++ b/example.env @@ -1,9 +1,8 @@ +FLOW_APP_NAME=ReMe FLOW_EMBEDDING_API_KEY=sk-xxxx FLOW_EMBEDDING_BASE_URL=https://xxxx/v1 - FLOW_LLM_API_KEY=sk-xxxx FLOW_LLM_BASE_URL=https://xxxx/v1 -FLOW_ES_HOSTS=http://0.0.0.0:9200 - -FLOW_USE_FRAMEWORK=true +# FLOW_ES_HOSTS=http://0.0.0.0:9200 +# FLOW_DASHSCOPE_API_KEY=sk-xxx \ No newline at end of file diff --git a/pyproject.toml b/pyproject.toml index 14828cbe..4fd5660a 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta" [project] name = "reme_ai" -version = "0.1.8" +version = "0.1.9" description = "Remember me" authors = [ { name = "jinli.yl", email = "jinli.yl@alibaba-inc.com" }, @@ -24,7 +24,7 @@ classifiers = [ keywords = ["llm", "memory", "experience", "memoryscope", "ai", "mcp", "http"] dependencies = [ - "flowllm>=0.1.8", + "flowllm>=0.1.9", ] [tool.setuptools.packages.find] diff --git a/reme_ai/__init__.py b/reme_ai/__init__.py index 8c2a11e4..50195757 100644 --- a/reme_ai/__init__.py +++ b/reme_ai/__init__.py @@ -3,4 +3,4 @@ from reme_ai import retrieve from reme_ai import summary from reme_ai import vector_store -__version__ = "0.1.8" +__version__ = "0.1.9" diff --git a/reme_ai/app.py b/reme_ai/app.py index 15a7762e..f4931b25 100644 --- a/reme_ai/app.py +++ b/reme_ai/app.py @@ -4,18 +4,16 @@ import warnings warnings.filterwarnings("ignore", category=DeprecationWarning, module="websockets") warnings.filterwarnings("ignore", category=DeprecationWarning, module="uvicorn") -from flowllm.service.base_service import BaseService +from flowllm.app import FlowLLMApp from reme_ai.config.config_parser import ConfigParser def main(): - with BaseService.get_service(*sys.argv[1:], parser=ConfigParser) as service: - service(logo="ReMe") - + with FlowLLMApp(args=sys.argv[1:], parser=ConfigParser) as app: + app.run_service() if __name__ == "__main__": main() -# python -m build -# twine upload dist/* +# python -m build && twine upload dist/* diff --git a/reme_ai/config/default.yaml b/reme_ai/config/default.yaml index 4debc881..7984694a 100644 --- a/reme_ai/config/default.yaml +++ b/reme_ai/config/default.yaml @@ -1,7 +1,5 @@ backend: http -language: "" -thread_pool_max_workers: 32 -ray_max_workers: 1 +thread_pool_max_workers: 64 mcp: transport: sse @@ -17,9 +15,6 @@ http: flow: retrieve_task_memory: flow_content: build_query_op >> recall_vector_store_op >> rerank_memory_op >> rewrite_memory_op - stream: false - use_async: true - service_type: http+mcp description: "Retrieves the most relevant top-k memory experiences from historical data based on the current query to enhance task-solving capabilities" input_schema: query: @@ -29,9 +24,6 @@ flow: summary_task_memory: flow_content: trajectory_preprocess_op >> (success_extraction_op|failure_extraction_op|comparative_extraction_op) >> memory_validation_op >> update_vector_store_op - stream: false - use_async: true - service_type: http+mcp description: "Summarizes conversation trajectories or messages into structured memory representations for long-term storage" input_schema: trajectories: @@ -41,9 +33,6 @@ flow: retrieve_personal_memory: flow_content: set_query_op >> (extract_time_op | (retrieve_memory_op >> semantic_rank_op)) >> fuse_rerank_op - stream: false - use_async: true - service_type: http+mcp description: "Retrieves the most relevant personal memories from historical data based on the query to enhance response quality" input_schema: query: @@ -53,9 +42,6 @@ flow: summary_personal_memory: flow_content: info_filter_op >> (get_observation_op | get_observation_with_time_op | load_today_memory_op) >> contra_repeat_op >> update_vector_store_op - stream: false - use_async: true - service_type: http+mcp description: "Consolidates user observations and memories by filtering information and removing redundancies for efficient storage" input_schema: trajectories: @@ -65,9 +51,6 @@ flow: retrieve_task_memory_simple: flow_content: build_query_op >> recall_vector_store_op >> merge_memory_op - stream: false - use_async: true - service_type: http+mcp description: "Retrieves the most relevant top-k memory experiences from historical data based on the current query with simplified processing" input_schema: query: @@ -77,9 +60,6 @@ flow: summary_task_memory_simple: flow_content: simple_summary_op >> update_vector_store_op - stream: false - use_async: true - service_type: http+mcp description: "Summarizes conversation trajectories or messages into memories using a simplified approach" input_schema: trajectories: @@ -89,9 +69,6 @@ flow: vector_store: flow_content: vector_store_action_op - stream: false - use_async: true - service_type: http+mcp description: "Directly operates on the vector store with various management actions" input_schema: action: @@ -102,9 +79,6 @@ flow: record_task_memory: flow_content: update_memory_freq_op >> update_memory_utility_op >> update_vector_store_op - stream: false - use_async: true - service_type: http+mcp description: "Update the freq & utility attributes of retrieved task memories" input_schema: workspace_id: @@ -122,9 +96,6 @@ flow: delete_task_memory: flow_content: delete_memory_op >> update_vector_store_op - stream: false - use_async: true - service_type: http+mcp description: "Delete task memories when utility/freq < utility_threshold and freq >= freq_threshold" input_schema: workspace_id: @@ -142,9 +113,6 @@ flow: react: flow_content: simple_react_op - stream: false - use_async: true - service_type: http+mcp description: "React to the current task with an agent" input_schema: query: diff --git a/reme_ai/react/simple_react_op.py b/reme_ai/react/simple_react_op.py index 7dccafae..4929949f 100644 --- a/reme_ai/react/simple_react_op.py +++ b/reme_ai/react/simple_react_op.py @@ -2,7 +2,7 @@ import asyncio from flowllm import C from flowllm.context.flow_context import FlowContext -from flowllm.op.llm.react_llm_op import ReactLLMOp +from flowllm.op.gallery.react_llm_op import ReactLLMOp @C.register_op() diff --git a/reme_ai/retrieve/personal/extract_time_op.py b/reme_ai/retrieve/personal/extract_time_op.py index 26ec7b2b..70266b87 100644 --- a/reme_ai/retrieve/personal/extract_time_op.py +++ b/reme_ai/retrieve/personal/extract_time_op.py @@ -1,7 +1,7 @@ import re from typing import Dict -from flowllm import C, BaseLLMOp +from flowllm import C, BaseAsyncOp from flowllm.enumeration.role import Role from flowllm.schema.message import Message from loguru import logger @@ -12,7 +12,7 @@ from reme_ai.utils.datetime_handler import DatetimeHandler @C.register_op() -class ExtractTimeOp(BaseLLMOp): +class ExtractTimeOp(BaseAsyncOp): file_path: str = __file__ EXTRACT_TIME_PATTERN = r"-\s*(\S+)[::]\s*(\S+)" diff --git a/reme_ai/retrieve/personal/fuse_rerank_op.py b/reme_ai/retrieve/personal/fuse_rerank_op.py index 369b4fc9..28e3dc1b 100644 --- a/reme_ai/retrieve/personal/fuse_rerank_op.py +++ b/reme_ai/retrieve/personal/fuse_rerank_op.py @@ -1,6 +1,6 @@ from typing import Dict, List -from flowllm import C, BaseLLMOp +from flowllm import C, BaseAsyncOp from loguru import logger from reme_ai.constants.common_constants import EXTRACT_TIME_DICT @@ -8,7 +8,7 @@ from reme_ai.schema.memory import BaseMemory @C.register_op() -class FuseRerankOp(BaseLLMOp): +class FuseRerankOp(BaseAsyncOp): """ Reranks the memory nodes by scores, types, and temporal relevance. Formats the top-K reranked nodes to print. """ diff --git a/reme_ai/retrieve/personal/print_memory_op.py b/reme_ai/retrieve/personal/print_memory_op.py index ff3e47ec..f7296f2b 100644 --- a/reme_ai/retrieve/personal/print_memory_op.py +++ b/reme_ai/retrieve/personal/print_memory_op.py @@ -1,13 +1,13 @@ from typing import List -from flowllm import C, BaseOp +from flowllm import C, BaseAsyncOp from loguru import logger from reme_ai.schema.memory import BaseMemory @C.register_op() -class PrintMemoryOp(BaseOp): +class PrintMemoryOp(BaseAsyncOp): """ Formats the memories to print. """ diff --git a/reme_ai/retrieve/personal/read_message_op.py b/reme_ai/retrieve/personal/read_message_op.py index 0bb06095..307e9511 100644 --- a/reme_ai/retrieve/personal/read_message_op.py +++ b/reme_ai/retrieve/personal/read_message_op.py @@ -1,12 +1,12 @@ from typing import List -from flowllm import C, BaseOp +from flowllm import C, BaseAsyncOp from flowllm.schema.message import Message from loguru import logger @C.register_op() -class ReadMessageOp(BaseOp): +class ReadMessageOp(BaseAsyncOp): """ Fetches unmemorized chat messages. """ diff --git a/reme_ai/retrieve/personal/retrieve_memory_op.py b/reme_ai/retrieve/personal/retrieve_memory_op.py index 8ee6fb17..a02a125e 100644 --- a/reme_ai/retrieve/personal/retrieve_memory_op.py +++ b/reme_ai/retrieve/personal/retrieve_memory_op.py @@ -1,6 +1,6 @@ from typing import List -from flowllm import C, BaseLLMOp +from flowllm import C, BaseAsyncOp from flowllm.schema.vector_node import VectorNode from loguru import logger @@ -8,7 +8,7 @@ from reme_ai.schema.memory import BaseMemory, vector_node_to_memory @C.register_op() -class RetrieveMemoryOp(BaseLLMOp): +class RetrieveMemoryOp(BaseAsyncOp): """ Retrieves memories based on specified criteria such as status, type, and timestamp. Processes these memories concurrently, sorts them by similarity, and logs the activity, diff --git a/reme_ai/retrieve/personal/semantic_rank_op.py b/reme_ai/retrieve/personal/semantic_rank_op.py index 0d7ffc20..2a267d4d 100644 --- a/reme_ai/retrieve/personal/semantic_rank_op.py +++ b/reme_ai/retrieve/personal/semantic_rank_op.py @@ -2,7 +2,7 @@ import json import re from typing import List -from flowllm import C, BaseLLMOp +from flowllm import C, BaseAsyncOp from loguru import logger from reme_ai.schema import Message, Role @@ -10,7 +10,7 @@ from reme_ai.schema.memory import BaseMemory @C.register_op() -class SemanticRankOp(BaseLLMOp): +class SemanticRankOp(BaseAsyncOp): """ The SemanticRankOp class processes queries by retrieving memory nodes, removing duplicates, ranking them based on semantic relevance using a model, diff --git a/reme_ai/retrieve/personal/set_query_op.py b/reme_ai/retrieve/personal/set_query_op.py index 5144f357..aa9617dc 100644 --- a/reme_ai/retrieve/personal/set_query_op.py +++ b/reme_ai/retrieve/personal/set_query_op.py @@ -1,14 +1,14 @@ import datetime from typing import Tuple -from flowllm import C, BaseOp +from flowllm import C, BaseAsyncOp from loguru import logger from reme_ai.constants.common_constants import QUERY_WITH_TS @C.register_op() -class SetQueryOp(BaseOp): +class SetQueryOp(BaseAsyncOp): """ The `SetQueryOp` class is responsible for setting a query and its associated timestamp into the context, utilizing either provided parameters or details from the context. diff --git a/reme_ai/retrieve/task/build_query_op.py b/reme_ai/retrieve/task/build_query_op.py index 80e11b92..e2463b3d 100644 --- a/reme_ai/retrieve/task/build_query_op.py +++ b/reme_ai/retrieve/task/build_query_op.py @@ -1,4 +1,4 @@ -from flowllm import C, BaseLLMOp +from flowllm import C, BaseAsyncOp from flowllm.utils.llm_utils import merge_messages_content from loguru import logger @@ -6,7 +6,7 @@ from reme_ai.schema import Message, Role @C.register_op() -class BuildQueryOp(BaseLLMOp): +class BuildQueryOp(BaseAsyncOp): file_path: str = __file__ async def async_execute(self): diff --git a/reme_ai/retrieve/task/merge_memory_op.py b/reme_ai/retrieve/task/merge_memory_op.py index ab1e8e88..6d67d203 100644 --- a/reme_ai/retrieve/task/merge_memory_op.py +++ b/reme_ai/retrieve/task/merge_memory_op.py @@ -1,13 +1,13 @@ from typing import List -from flowllm import C, BaseOp +from flowllm import C, BaseAsyncOp from loguru import logger from reme_ai.schema.memory import BaseMemory @C.register_op() -class MergeMemoryOp(BaseOp): +class MergeMemoryOp(BaseAsyncOp): async def async_execute(self): memory_list: List[BaseMemory] = self.context.response.metadata["memory_list"] diff --git a/reme_ai/retrieve/task/rerank_memory_op.py b/reme_ai/retrieve/task/rerank_memory_op.py index fec119a8..298300b6 100644 --- a/reme_ai/retrieve/task/rerank_memory_op.py +++ b/reme_ai/retrieve/task/rerank_memory_op.py @@ -2,7 +2,7 @@ import json import re from typing import List -from flowllm import C, BaseLLMOp +from flowllm import C, BaseAsyncOp from flowllm.enumeration.role import Role from flowllm.schema.message import Message from loguru import logger @@ -11,7 +11,7 @@ from reme_ai.schema.memory import BaseMemory @C.register_op() -class RerankMemoryOp(BaseLLMOp): +class RerankMemoryOp(BaseAsyncOp): """ Rerank and filter recalled experiences using LLM and score-based filtering """ diff --git a/reme_ai/retrieve/task/rewrite_memory_op.py b/reme_ai/retrieve/task/rewrite_memory_op.py index eae0b0a1..37ad179b 100644 --- a/reme_ai/retrieve/task/rewrite_memory_op.py +++ b/reme_ai/retrieve/task/rewrite_memory_op.py @@ -2,7 +2,7 @@ import json import re from typing import List -from flowllm import C, BaseLLMOp +from flowllm import C, BaseAsyncOp from flowllm.enumeration.role import Role from flowllm.schema.message import Message from loguru import logger @@ -11,7 +11,7 @@ from reme_ai.schema.memory import BaseMemory @C.register_op() -class RewriteMemoryOp(BaseLLMOp): +class RewriteMemoryOp(BaseAsyncOp): """ Generate and rewrite context messages from reranked experiences """ diff --git a/reme_ai/service/base_memory_service.py b/reme_ai/service/agentscope_runtime_memory_service.py similarity index 93% rename from reme_ai/service/base_memory_service.py rename to reme_ai/service/agentscope_runtime_memory_service.py index 208b43b8..e4c21713 100644 --- a/reme_ai/service/base_memory_service.py +++ b/reme_ai/service/agentscope_runtime_memory_service.py @@ -1,12 +1,16 @@ from abc import abstractmethod, ABC from typing import Optional, Dict, Any +from flowllm import FlowLLMApp from pydantic import Field +from reme_ai.config.config_parser import ConfigParser -class BaseMemoryService(ABC): + +class AgentscopeRuntimeMemoryService(ABC): def __init__(self): + self.app = FlowLLMApp(parser=ConfigParser, load_default_config=True) self.session_id_dict: dict = {} def add_session_memory_id(self, session_id: str, memory_id): diff --git a/reme_ai/service/personal_memory_service.py b/reme_ai/service/personal_memory_service.py index 60a85a27..30e8bc95 100644 --- a/reme_ai/service/personal_memory_service.py +++ b/reme_ai/service/personal_memory_service.py @@ -1,31 +1,26 @@ import asyncio from typing import Optional, Dict, Any, List -from flowllm import C -from flowllm.flow import BaseToolFlow from flowllm.schema.flow_response import FlowResponse from loguru import logger from pydantic import Field, BaseModel -from reme_ai.config.config_parser import ConfigParser from reme_ai.schema.memory import PersonalMemory -from reme_ai.service.base_memory_service import BaseMemoryService +from reme_ai.service.agentscope_runtime_memory_service import AgentscopeRuntimeMemoryService -class PersonalMemoryService(BaseMemoryService): +class PersonalMemoryService(AgentscopeRuntimeMemoryService): - async def start(self) -> None: - C.set_service_config(parser=ConfigParser, config_name="config=default").init_by_service_config() + async def start(self): + return await self.app.async_start() async def stop(self) -> None: - C.stop_by_service_config() + return await self.app.async_stop() async def health(self) -> bool: return True async def add_memory(self, user_id: str, messages: list, session_id: Optional[str] = None) -> None: - summary_flow: BaseToolFlow = C.flow_dict["summary_personal_memory"] - new_messages: List[dict] = [] for message in messages: if isinstance(message, dict): @@ -42,7 +37,7 @@ class PersonalMemoryService(BaseMemoryService): ] } - result: FlowResponse = await summary_flow(**kwargs) + result: FlowResponse = await self.app.async_execute_flow(name="summary_personal_memory", **kwargs) memory_list: List[PersonalMemory] = result.metadata.get("memory_list", []) for memory in memory_list: memory_id = memory.memory_id @@ -54,9 +49,6 @@ class PersonalMemoryService(BaseMemoryService): "such as top_k, score etc.", default=None, )) -> list: - - retrieve_flow: BaseToolFlow = C.flow_dict["retrieve_personal_memory"] - new_messages: List[dict] = [] for message in messages: if isinstance(message, dict): @@ -75,7 +67,7 @@ class PersonalMemoryService(BaseMemoryService): "top_k": filters.get("top_k", 1) if filters else 1 } - result: FlowResponse = await retrieve_flow(**kwargs) + result: FlowResponse = await self.app.async_execute_flow(name="retrieve_personal_memory", **kwargs) logger.info(f"[personal_memory_service] user_id={user_id} search result: {result.model_dump_json()}") return [result.answer] @@ -85,8 +77,7 @@ class PersonalMemoryService(BaseMemoryService): "such as top_k, score etc.", default=None, )) -> list: - vector_store_flow: BaseToolFlow = C.flow_dict["vector_store"] - result = await vector_store_flow(workspace_id=user_id, action="list") + result = await self.app.async_execute_flow(name="vector_store", workspace_id=user_id, action="list") logger.info(f"[personal_memory_service] list_memory result: {result}") result = result.metadata["action_result"] @@ -99,8 +90,10 @@ class PersonalMemoryService(BaseMemoryService): if not delete_ids: return - vector_store_flow: BaseToolFlow = C.flow_dict["vector_store"] - result = await vector_store_flow(workspace_id=user_id, action="delete_ids", memory_ids=delete_ids) + result = await self.app.async_execute_flow(name="vector_store", + workspace_id=user_id, + action="delete_ids", + memory_ids=delete_ids) result = result.metadata["action_result"] logger.info(f"[personal_memory_service] delete memory result={result}") diff --git a/reme_ai/service/task_memory_service.py b/reme_ai/service/task_memory_service.py index 32499d05..319334de 100644 --- a/reme_ai/service/task_memory_service.py +++ b/reme_ai/service/task_memory_service.py @@ -1,31 +1,26 @@ import asyncio from typing import Optional, Dict, Any, List -from flowllm import C -from flowllm.flow import BaseToolFlow from flowllm.schema.flow_response import FlowResponse from loguru import logger from pydantic import Field, BaseModel -from reme_ai.config.config_parser import ConfigParser from reme_ai.schema.memory import TaskMemory -from reme_ai.service.base_memory_service import BaseMemoryService +from reme_ai.service.agentscope_runtime_memory_service import AgentscopeRuntimeMemoryService -class TaskMemoryService(BaseMemoryService): +class TaskMemoryService(AgentscopeRuntimeMemoryService): - async def start(self) -> None: - C.set_service_config(parser=ConfigParser, config_name="config=default").init_by_service_config() + async def start(self): + return await self.app.async_start() async def stop(self) -> None: - C.stop_by_service_config() + return await self.app.async_stop() async def health(self) -> bool: return True async def add_memory(self, user_id: str, messages: list, session_id: Optional[str] = None) -> None: - summary_flow: BaseToolFlow = C.flow_dict["summary_task_memory"] - new_messages: List[dict] = [] for message in messages: if isinstance(message, dict): @@ -42,7 +37,7 @@ class TaskMemoryService(BaseMemoryService): ] } - result: FlowResponse = await summary_flow(**kwargs) + result: FlowResponse = await self.app.async_execute_flow(name="summary_task_memory", **kwargs) memory_list: List[TaskMemory] = result.metadata.get("memory_list", []) for memory in memory_list: memory_id = memory.memory_id @@ -54,9 +49,6 @@ class TaskMemoryService(BaseMemoryService): "such as top_k, score etc.", default=None, )) -> list: - - retrieve_flow: BaseToolFlow = C.flow_dict["retrieve_task_memory"] - new_messages: List[dict] = [] for message in messages: if isinstance(message, dict): @@ -72,7 +64,7 @@ class TaskMemoryService(BaseMemoryService): "top_k": filters.get("top_k", 1) if filters else 1 } - result: FlowResponse = await retrieve_flow(**kwargs) + result: FlowResponse = await self.app.async_execute_flow(name="retrieve_task_memory", **kwargs) logger.info(f"[task_memory_service] user_id={user_id} add result: {result.model_dump_json()}") return [result.answer] @@ -82,11 +74,9 @@ class TaskMemoryService(BaseMemoryService): "such as top_k, score etc.", default=None, )) -> list: - vector_store_flow: BaseToolFlow = C.flow_dict["vector_store"] - result = await vector_store_flow(workspace_id=user_id, action="list") + result = await self.app.async_execute_flow(name="vector_store", workspace_id=user_id, action="list") print("list_memory result:", result) - result = result.metadata["action_result"] for i, line in enumerate(result): logger.info(f"[task_memory_service] list memory.{i}={line}") @@ -97,8 +87,10 @@ class TaskMemoryService(BaseMemoryService): if not delete_ids: return - vector_store_flow: BaseToolFlow = C.flow_dict["vector_store"] - result = await vector_store_flow(workspace_id=user_id, action="delete_ids", memory_ids=delete_ids) + result = await self.app.async_execute_flow(name="vector_store", + workspace_id=user_id, + action="delete_ids", + memory_ids=delete_ids) result = result.metadata["action_result"] logger.info(f"[task_memory_service] delete memory result={result}") diff --git a/reme_ai/summary/personal/contra_repeat_op.py b/reme_ai/summary/personal/contra_repeat_op.py index 9330e45c..fc85f436 100644 --- a/reme_ai/summary/personal/contra_repeat_op.py +++ b/reme_ai/summary/personal/contra_repeat_op.py @@ -2,7 +2,7 @@ import json import re from typing import List, Tuple -from flowllm import C, BaseLLMOp +from flowllm import C, BaseAsyncOp from flowllm.enumeration.role import Role from flowllm.schema.message import Message from loguru import logger @@ -11,10 +11,10 @@ from reme_ai.schema.memory import BaseMemory @C.register_op() -class ContraRepeatOp(BaseLLMOp): +class ContraRepeatOp(BaseAsyncOp): """ The `ContraRepeatOp` class specializes in processing memory nodes to identify and handle - contradictory and repetitive information. It extends the base functionality of `BaseLLMOp`. + contradictory and repetitive information. It extends the base functionality of `BaseAsyncOp`. Responsibilities: - Collects observation memories from context. diff --git a/reme_ai/summary/personal/get_observation_op.py b/reme_ai/summary/personal/get_observation_op.py index 6c7abd1a..14ead179 100644 --- a/reme_ai/summary/personal/get_observation_op.py +++ b/reme_ai/summary/personal/get_observation_op.py @@ -1,7 +1,7 @@ import re from typing import List -from flowllm import C, BaseLLMOp +from flowllm import C, BaseAsyncOp from flowllm.schema.message import Message from loguru import logger @@ -10,9 +10,9 @@ from reme_ai.utils.datetime_handler import DatetimeHandler @C.register_op() -class GetObservationOp(BaseLLMOp): +class GetObservationOp(BaseAsyncOp): """ - A specialized operation class to generate observations from chat messages using BaseLLMOp. + A specialized operation class to generate observations from chat messages using BaseAsyncOp. """ file_path: str = __file__ diff --git a/reme_ai/summary/personal/get_observation_with_time_op.py b/reme_ai/summary/personal/get_observation_with_time_op.py index cd3056b3..b19dfa06 100644 --- a/reme_ai/summary/personal/get_observation_with_time_op.py +++ b/reme_ai/summary/personal/get_observation_with_time_op.py @@ -1,7 +1,7 @@ import re from typing import List -from flowllm import C, BaseLLMOp +from flowllm import C, BaseAsyncOp from flowllm.schema.message import Message from loguru import logger @@ -10,9 +10,9 @@ from reme_ai.utils.datetime_handler import DatetimeHandler @C.register_op() -class GetObservationWithTimeOp(BaseLLMOp): +class GetObservationWithTimeOp(BaseAsyncOp): """ - A specialized operation class to extract observations with time information from chat messages using BaseLLMOp. + A specialized operation class to extract observations with time information from chat messages using BaseAsyncOp. """ file_path: str = __file__ diff --git a/reme_ai/summary/personal/get_reflection_subject_op.py b/reme_ai/summary/personal/get_reflection_subject_op.py index 328963de..a88a9e7d 100644 --- a/reme_ai/summary/personal/get_reflection_subject_op.py +++ b/reme_ai/summary/personal/get_reflection_subject_op.py @@ -1,6 +1,6 @@ from typing import List -from flowllm import C, BaseLLMOp +from flowllm import C, BaseAsyncOp from flowllm.schema.message import Message from loguru import logger @@ -8,7 +8,7 @@ from reme_ai.schema.memory import BaseMemory, PersonalMemory @C.register_op() -class GetReflectionSubjectOp(BaseLLMOp): +class GetReflectionSubjectOp(BaseAsyncOp): """ A specialized operation class responsible for retrieving unreflected memory nodes, generating reflection prompts with current insights, invoking an LLM for fresh insights, diff --git a/reme_ai/summary/personal/info_filter_op.py b/reme_ai/summary/personal/info_filter_op.py index 6bc7d186..0e741852 100644 --- a/reme_ai/summary/personal/info_filter_op.py +++ b/reme_ai/summary/personal/info_filter_op.py @@ -1,7 +1,7 @@ import re from typing import List -from flowllm import C, BaseLLMOp +from flowllm import C, BaseAsyncOp from flowllm.schema.message import Message, Trajectory from loguru import logger @@ -9,9 +9,9 @@ from reme_ai.schema.memory import PersonalMemory @C.register_op() -class InfoFilterOp(BaseLLMOp): +class InfoFilterOp(BaseAsyncOp): """ - A specialized operation class to filter messages based on information content scores using BaseLLMOp. + A specialized operation class to filter messages based on information content scores using BaseAsyncOp. This filters chat messages by retaining only those that include significant information about the user. """ file_path: str = __file__ diff --git a/reme_ai/summary/personal/load_today_memory_op.py b/reme_ai/summary/personal/load_today_memory_op.py index 3cc48f78..989269d7 100644 --- a/reme_ai/summary/personal/load_today_memory_op.py +++ b/reme_ai/summary/personal/load_today_memory_op.py @@ -1,6 +1,6 @@ from typing import List -from flowllm import C, BaseLLMOp +from flowllm import C, BaseAsyncOp from flowllm.schema.vector_node import VectorNode from loguru import logger @@ -9,7 +9,7 @@ from reme_ai.utils.datetime_handler import DatetimeHandler @C.register_op() -class LoadTodayMemoryOp(BaseLLMOp): +class LoadTodayMemoryOp(BaseAsyncOp): """ Operation to load today's memories from vector store for deduplication. Focuses specifically on retrieving and deduplicating memories from the current date. diff --git a/reme_ai/summary/personal/long_contra_repeat_op.py b/reme_ai/summary/personal/long_contra_repeat_op.py index bf77837d..bb16ede1 100644 --- a/reme_ai/summary/personal/long_contra_repeat_op.py +++ b/reme_ai/summary/personal/long_contra_repeat_op.py @@ -1,7 +1,7 @@ import re from typing import List -from flowllm import C, BaseLLMOp +from flowllm import C, BaseAsyncOp from flowllm.enumeration.role import Role from flowllm.schema.message import Message from loguru import logger @@ -10,10 +10,10 @@ from reme_ai.schema.memory import BaseMemory, PersonalMemory @C.register_op() -class LongContraRepeatOp(BaseLLMOp): +class LongContraRepeatOp(BaseAsyncOp): """ Manages and updates memory entries within a conversation scope by identifying - and handling contradictions or redundancies. It extends BaseLLMOp to provide + and handling contradictions or redundancies. It extends BaseAsyncOp to provide specialized functionality for long conversations with potential contradictory or repetitive statements. """ diff --git a/reme_ai/summary/personal/update_insight_op.py b/reme_ai/summary/personal/update_insight_op.py index b20e2882..17740561 100644 --- a/reme_ai/summary/personal/update_insight_op.py +++ b/reme_ai/summary/personal/update_insight_op.py @@ -1,7 +1,7 @@ import re from typing import List -from flowllm import C, BaseLLMOp +from flowllm import C, BaseAsyncOp from flowllm.schema.message import Message from loguru import logger @@ -9,7 +9,7 @@ from reme_ai.schema.memory import PersonalMemory @C.register_op() -class UpdateInsightOp(BaseLLMOp): +class UpdateInsightOp(BaseAsyncOp): """ This class is responsible for updating insight value in a memory system. It filters insight nodes based on their association with observed nodes, utilizes a ranking model to prioritize them, diff --git a/reme_ai/summary/task/comparative_extraction_op.py b/reme_ai/summary/task/comparative_extraction_op.py index 9823c944..71e72a7d 100644 --- a/reme_ai/summary/task/comparative_extraction_op.py +++ b/reme_ai/summary/task/comparative_extraction_op.py @@ -1,6 +1,6 @@ from typing import List, Tuple, Optional -from flowllm import C, BaseLLMOp +from flowllm import C, BaseAsyncOp from flowllm.enumeration.role import Role from flowllm.schema.message import Message as FlowMessage from loguru import logger @@ -11,7 +11,7 @@ from reme_ai.utils.op_utils import merge_messages_content, parse_json_experience @C.register_op() -class ComparativeExtractionOp(BaseLLMOp): +class ComparativeExtractionOp(BaseAsyncOp): file_path: str = __file__ async def async_execute(self): diff --git a/reme_ai/summary/task/failure_extraction_op.py b/reme_ai/summary/task/failure_extraction_op.py index 1b5ed7dc..2a7b67c9 100644 --- a/reme_ai/summary/task/failure_extraction_op.py +++ b/reme_ai/summary/task/failure_extraction_op.py @@ -1,6 +1,6 @@ from typing import List -from flowllm import C, BaseLLMOp +from flowllm import C, BaseAsyncOp from flowllm.enumeration.role import Role from flowllm.schema.message import Message as FlowMessage from loguru import logger @@ -11,7 +11,7 @@ from reme_ai.utils.op_utils import merge_messages_content, parse_json_experience @C.register_op() -class FailureExtractionOp(BaseLLMOp): +class FailureExtractionOp(BaseAsyncOp): file_path: str = __file__ async def async_execute(self): diff --git a/reme_ai/summary/task/memory_deduplication_op.py b/reme_ai/summary/task/memory_deduplication_op.py index be884eda..e5b86a32 100644 --- a/reme_ai/summary/task/memory_deduplication_op.py +++ b/reme_ai/summary/task/memory_deduplication_op.py @@ -1,13 +1,13 @@ from typing import List -from flowllm import C, BaseOp +from flowllm import C, BaseAsyncOp from loguru import logger from reme_ai.schema.memory import BaseMemory @C.register_op() -class MemoryDeduplicationOp(BaseOp): +class MemoryDeduplicationOp(BaseAsyncOp): file_path: str = __file__ async def async_execute(self): diff --git a/reme_ai/summary/task/memory_validation_op.py b/reme_ai/summary/task/memory_validation_op.py index 3e6a208f..485c66e1 100644 --- a/reme_ai/summary/task/memory_validation_op.py +++ b/reme_ai/summary/task/memory_validation_op.py @@ -2,7 +2,7 @@ import json import re from typing import List, Dict, Any -from flowllm import C, BaseLLMOp +from flowllm import C, BaseAsyncOp from flowllm.enumeration.role import Role from flowllm.schema.message import Message as FlowMessage from loguru import logger @@ -12,7 +12,7 @@ from reme_ai.schema.memory import BaseMemory @C.register_op() -class MemoryValidationOp(BaseLLMOp): +class MemoryValidationOp(BaseAsyncOp): file_path: str = __file__ async def async_execute(self): diff --git a/reme_ai/summary/task/simple_comparative_summary_op.py b/reme_ai/summary/task/simple_comparative_summary_op.py index 04f83bed..b2d212b5 100644 --- a/reme_ai/summary/task/simple_comparative_summary_op.py +++ b/reme_ai/summary/task/simple_comparative_summary_op.py @@ -1,7 +1,7 @@ import json from typing import List, Dict -from flowllm import C, BaseLLMOp +from flowllm import C, BaseAsyncOp from flowllm.enumeration.role import Role from flowllm.schema.message import Message as FlowMessage from loguru import logger @@ -12,7 +12,7 @@ from reme_ai.utils.op_utils import merge_messages_content @C.register_op() -class SimpleComparativeSummaryOp(BaseLLMOp): +class SimpleComparativeSummaryOp(BaseAsyncOp): file_path: str = __file__ async def compare_summary_trajectory(self, trajectory_a: Trajectory, trajectory_b: Trajectory) -> List[BaseMemory]: diff --git a/reme_ai/summary/task/simple_summary_op.py b/reme_ai/summary/task/simple_summary_op.py index 6e9838e2..65e66d80 100644 --- a/reme_ai/summary/task/simple_summary_op.py +++ b/reme_ai/summary/task/simple_summary_op.py @@ -1,7 +1,7 @@ import json from typing import List -from flowllm import C, BaseLLMOp +from flowllm import C, BaseAsyncOp from flowllm.enumeration.role import Role from flowllm.schema.message import Message as FlowMessage from loguru import logger @@ -12,7 +12,7 @@ from reme_ai.utils.op_utils import merge_messages_content @C.register_op() -class SimpleSummaryOp(BaseLLMOp): +class SimpleSummaryOp(BaseAsyncOp): file_path: str = __file__ async def summary_trajectory(self, trajectory: Trajectory) -> List[BaseMemory]: diff --git a/reme_ai/summary/task/success_extraction_op.py b/reme_ai/summary/task/success_extraction_op.py index 8db77d06..72bed938 100644 --- a/reme_ai/summary/task/success_extraction_op.py +++ b/reme_ai/summary/task/success_extraction_op.py @@ -1,6 +1,6 @@ from typing import List -from flowllm import C, BaseLLMOp +from flowllm import C, BaseAsyncOp from flowllm.enumeration.role import Role from flowllm.schema.message import Message as FlowMessage from loguru import logger @@ -11,7 +11,7 @@ from reme_ai.utils.op_utils import merge_messages_content, parse_json_experience @C.register_op() -class SuccessExtractionOp(BaseLLMOp): +class SuccessExtractionOp(BaseAsyncOp): file_path: str = __file__ async def async_execute(self): diff --git a/reme_ai/summary/task/trajectory_preprocess_op.py b/reme_ai/summary/task/trajectory_preprocess_op.py index 97dc847e..7d8357ca 100644 --- a/reme_ai/summary/task/trajectory_preprocess_op.py +++ b/reme_ai/summary/task/trajectory_preprocess_op.py @@ -1,13 +1,13 @@ from typing import List, Dict -from flowllm import C, BaseOp +from flowllm import C, BaseAsyncOp from loguru import logger from reme_ai.schema import Trajectory @C.register_op() -class TrajectoryPreprocessOp(BaseOp): +class TrajectoryPreprocessOp(BaseAsyncOp): file_path: str = __file__ async def async_execute(self): diff --git a/reme_ai/summary/task/trajectory_segmentation_op.py b/reme_ai/summary/task/trajectory_segmentation_op.py index 647913e1..4e0f0b54 100644 --- a/reme_ai/summary/task/trajectory_segmentation_op.py +++ b/reme_ai/summary/task/trajectory_segmentation_op.py @@ -2,7 +2,7 @@ import json import re from typing import List -from flowllm import C, BaseLLMOp +from flowllm import C, BaseAsyncOp from flowllm.enumeration.role import Role from flowllm.schema.message import Message as FlowMessage from loguru import logger @@ -11,7 +11,7 @@ from reme_ai.schema import Message, Trajectory @C.register_op() -class TrajectorySegmentationOp(BaseLLMOp): +class TrajectorySegmentationOp(BaseAsyncOp): file_path: str = __file__ async def async_execute(self): diff --git a/reme_ai/vector_store/delete_memory_op.py b/reme_ai/vector_store/delete_memory_op.py index df7e6484..893dcb8c 100644 --- a/reme_ai/vector_store/delete_memory_op.py +++ b/reme_ai/vector_store/delete_memory_op.py @@ -1,11 +1,11 @@ from typing import Iterable -from flowllm import C, BaseLLMOp +from flowllm import C, BaseAsyncOp from flowllm.schema.vector_node import VectorNode @C.register_op() -class DeleteMemoryOp(BaseLLMOp): +class DeleteMemoryOp(BaseAsyncOp): file_path: str = __file__ async def async_execute(self): diff --git a/reme_ai/vector_store/recall_vector_store_op.py b/reme_ai/vector_store/recall_vector_store_op.py index 5ba526d4..354bddfe 100644 --- a/reme_ai/vector_store/recall_vector_store_op.py +++ b/reme_ai/vector_store/recall_vector_store_op.py @@ -1,6 +1,6 @@ from typing import List -from flowllm import C, BaseLLMOp +from flowllm import C, BaseAsyncOp from flowllm.schema.vector_node import VectorNode from loguru import logger @@ -8,7 +8,7 @@ from reme_ai.schema.memory import BaseMemory, vector_node_to_memory @C.register_op() -class RecallVectorStoreOp(BaseLLMOp): +class RecallVectorStoreOp(BaseAsyncOp): async def async_execute(self): recall_key: str = self.op_params.get("recall_key", "query") diff --git a/reme_ai/vector_store/update_memory_freq_op.py b/reme_ai/vector_store/update_memory_freq_op.py index bc2dff15..f3a88fdc 100644 --- a/reme_ai/vector_store/update_memory_freq_op.py +++ b/reme_ai/vector_store/update_memory_freq_op.py @@ -1,13 +1,13 @@ from typing import List -from flowllm import C, BaseOp +from flowllm import C, BaseAsyncOp from loguru import logger from reme_ai.schema.memory import BaseMemory, dict_to_memory @C.register_op() -class UpdateMemoryFreqOp(BaseOp): +class UpdateMemoryFreqOp(BaseAsyncOp): file_path: str = __file__ async def async_execute(self): diff --git a/reme_ai/vector_store/update_memory_utility_op.py b/reme_ai/vector_store/update_memory_utility_op.py index 999cacbf..08278ede 100644 --- a/reme_ai/vector_store/update_memory_utility_op.py +++ b/reme_ai/vector_store/update_memory_utility_op.py @@ -1,13 +1,13 @@ from typing import List -from flowllm import C, BaseOp +from flowllm import C, BaseAsyncOp from loguru import logger from reme_ai.schema.memory import BaseMemory @C.register_op() -class UpdateMemoryUtilityOp(BaseOp): +class UpdateMemoryUtilityOp(BaseAsyncOp): file_path: str = __file__ async def async_execute(self): diff --git a/reme_ai/vector_store/update_vector_store_op.py b/reme_ai/vector_store/update_vector_store_op.py index 1fd2487c..9e156554 100644 --- a/reme_ai/vector_store/update_vector_store_op.py +++ b/reme_ai/vector_store/update_vector_store_op.py @@ -1,7 +1,7 @@ import json from typing import List -from flowllm import C, BaseLLMOp +from flowllm import C, BaseAsyncOp from flowllm.schema.vector_node import VectorNode from loguru import logger @@ -9,7 +9,7 @@ from reme_ai.schema.memory import BaseMemory @C.register_op() -class UpdateVectorStoreOp(BaseLLMOp): +class UpdateVectorStoreOp(BaseAsyncOp): async def async_execute(self): workspace_id: str = self.context.workspace_id diff --git a/reme_ai/vector_store/vector_store_action_op.py b/reme_ai/vector_store/vector_store_action_op.py index efd28ca3..ae9182d9 100644 --- a/reme_ai/vector_store/vector_store_action_op.py +++ b/reme_ai/vector_store/vector_store_action_op.py @@ -1,11 +1,11 @@ -from flowllm import C, BaseLLMOp +from flowllm import C, BaseAsyncOp from flowllm.schema.vector_node import VectorNode from reme_ai.schema.memory import vector_node_to_memory, dict_to_memory, BaseMemory @C.register_op() -class VectorStoreActionOp(BaseLLMOp): +class VectorStoreActionOp(BaseAsyncOp): async def async_execute(self): workspace_id: str = self.context.workspace_id