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
This commit is contained in:
jinliyl 2025-09-16 16:49:56 +08:00 committed by GitHub
parent 00b81ebb7e
commit 560e746cea
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
47 changed files with 118 additions and 170 deletions

View file

@ -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

View file

@ -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

View file

@ -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

View file

@ -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

View file

@ -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

View file

@ -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]

View file

@ -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"

View file

@ -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/*

View file

@ -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:

View file

@ -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()

View file

@ -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+)"

View file

@ -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.
"""

View file

@ -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.
"""

View file

@ -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.
"""

View file

@ -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,

View file

@ -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,

View file

@ -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.

View file

@ -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):

View file

@ -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"]

View file

@ -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
"""

View file

@ -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
"""

View file

@ -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):

View file

@ -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}")

View file

@ -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}")

View file

@ -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.

View file

@ -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__

View file

@ -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__

View file

@ -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,

View file

@ -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__

View file

@ -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.

View file

@ -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.
"""

View file

@ -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,

View file

@ -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):

View file

@ -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):

View file

@ -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):

View file

@ -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):

View file

@ -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]:

View file

@ -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]:

View file

@ -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):

View file

@ -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):

View file

@ -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):

View file

@ -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):

View file

@ -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")

View file

@ -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):

View file

@ -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):

View file

@ -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

View file

@ -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