Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@
##############################################################################

name = "MemoryOS"
version = "2.0.32"
version = "2.0.33"
description = "Intelligence Begins with Memory"
license = {text = "Apache-2.0"}
readme = "README.md"
Expand Down
2 changes: 1 addition & 1 deletion src/memos/__init__.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
__version__ = "2.0.32"
__version__ = "2.0.33"

from memos.configs.mem_cube import GeneralMemCubeConfig
from memos.configs.mem_os import MOSConfig
Expand Down
55 changes: 46 additions & 9 deletions src/memos/mem_reader/read_pref_memory/process_preference_memory.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,12 @@
from memos.context.context import ContextThreadPoolExecutor
from memos.log import get_logger
from memos.mem_reader.read_multi_modal import detect_lang
from memos.memories.textual.item import TextualMemoryItem, TreeNodeTextualMemoryMetadata
from memos.mem_reader.source_filter import MemorySourceFilter
from memos.memories.textual.item import (
SourceMessage,
TextualMemoryItem,
TreeNodeTextualMemoryMetadata,
)
from memos.templates.prefer_complete_prompt import (
NAIVE_EXPLICIT_PREFERENCE_EXTRACT_PROMPT,
NAIVE_EXPLICIT_PREFERENCE_EXTRACT_PROMPT_ZH,
Expand Down Expand Up @@ -100,6 +105,7 @@ def _create_preference_memory_item(
fast_item: TextualMemoryItem | None,
info: dict[str, Any],
embedder,
sources_override: list[SourceMessage] | None = None,
**kwargs,
) -> TextualMemoryItem:
"""
Expand Down Expand Up @@ -133,7 +139,13 @@ def _create_preference_memory_item(
embedding = embedder.embed([context_summary])[0] if embedder and context_summary else None

# Extract sources from fast_item
sources = getattr(fast_item.metadata, "sources", []) if fast_item else []
sources = (
sources_override
if sources_override is not None
else getattr(fast_item.metadata, "sources", [])
if fast_item
else []
)

# Create metadata
metadata = TreeNodeTextualMemoryMetadata(
Expand Down Expand Up @@ -168,6 +180,7 @@ def _process_single_chunk_explicit(
info: dict[str, Any],
llm,
embedder,
sources_override: list[SourceMessage] | None = None,
**kwargs,
) -> list[TextualMemoryItem]:
"""Process a single chunk for explicit preferences."""
Expand All @@ -190,6 +203,7 @@ def _process_single_chunk_explicit(
fast_item=fast_item,
info=info,
embedder=embedder,
sources_override=sources_override,
**kwargs,
)
memories.append(memory)
Expand All @@ -203,6 +217,7 @@ def _process_single_chunk_implicit(
info: dict[str, Any],
llm,
embedder,
sources_override: list[SourceMessage] | None = None,
**kwargs,
) -> list[TextualMemoryItem]:
"""Process a single chunk for implicit preferences."""
Expand All @@ -225,6 +240,7 @@ def _process_single_chunk_implicit(
fast_item=fast_item,
info=info,
embedder=embedder,
sources_override=sources_override,
**kwargs,
)
memories.append(memory)
Expand Down Expand Up @@ -260,13 +276,20 @@ def process_preference_fine(
return []

try:
# Convert fast_memory_items to messages format
# Convert fast_memory_items to source-filtered messages format
source_filter = MemorySourceFilter()
chunks = []
for fast_item in fast_memory_items:
mem_str = fast_item.memory or ""
raw_sources = getattr(fast_item.metadata, "sources", None)
if raw_sources:
filtered_sources = source_filter.filter_for_preference(raw_sources)
mem_str = source_filter.sources_to_prompt_text(filtered_sources)
else:
filtered_sources = None
mem_str = fast_item.memory or ""
if not mem_str.strip():
continue
chunks.append((mem_str, fast_item))
chunks.append((mem_str, fast_item, filtered_sources))

if not chunks:
return []
Expand All @@ -277,16 +300,30 @@ def process_preference_fine(
futures = {}

# Submit explicit extraction tasks
for chunk, fast_item in chunks:
for chunk, fast_item, filtered_sources in chunks:
future = executor.submit(
_process_single_chunk_explicit, chunk, fast_item, info, llm, embedder, **kwargs
_process_single_chunk_explicit,
chunk,
fast_item,
info,
llm,
embedder,
filtered_sources,
**kwargs,
)
futures[future] = ("explicit_preference", chunk)

# Submit implicit extraction tasks
for chunk, fast_item in chunks:
for chunk, fast_item, filtered_sources in chunks:
future = executor.submit(
_process_single_chunk_implicit, chunk, fast_item, info, llm, embedder, **kwargs
_process_single_chunk_implicit,
chunk,
fast_item,
info,
llm,
embedder,
filtered_sources,
**kwargs,
)
futures[future] = ("implicit_preference", chunk)

Expand Down
228 changes: 228 additions & 0 deletions src/memos/mem_reader/source_filter.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,228 @@
"""Source filters shared by memory extraction steps."""

from __future__ import annotations

import re

from dataclasses import dataclass
from typing import Any, Literal

from memos.memories.textual.item import SourceMessage


SourceFilterAction = Literal[
"extract_after_last",
"strip_after_first",
"drop_if_present",
"drop_if_prefix",
]


@dataclass(frozen=True)
class SourceFilterRule:
name: str
action: SourceFilterAction
patterns: tuple[re.Pattern[str], ...]


@dataclass(frozen=True)
class SourceFilterPolicy:
allowed_roles: frozenset[str]
blocked_roles: frozenset[str]
blocked_source_types: frozenset[str]
rules: tuple[SourceFilterRule, ...]


def _patterns(*values: str, flags: int = 0) -> tuple[re.Pattern[str], ...]:
return tuple(re.compile(value, flags) for value in values)


PREFERENCE_SOURCE_POLICY = SourceFilterPolicy(
allowed_roles=frozenset({"user"}),
blocked_roles=frozenset({"assistant", "system", "tool"}),
blocked_source_types=frozenset({"tool"}),
rules=(
SourceFilterRule(
name="user_query_boundary",
action="extract_after_last",
patterns=_patterns(
r"(?im)(?:^|[ \t])#{1,3}[ \t]*用户(?:的)?(?:消息|问题)(?:为|是)?[ \t]*[::]",
r"user\u200b原\u200b始\u200bquery\u200b:\u200b\u200b\u200b\u200b",
r"(?im)(?:^|[ \t])#{1,3}[ \t]*user[ \t]*原始[ \t]*query[ \t]*[::]",
),
),
SourceFilterRule(
name="trailing_context",
action="strip_after_first",
patterns=_patterns(
r"(?m)^\s{0,3}#{1,3}\s*以下是可能和用户问题关联的对话记忆",
),
),
SourceFilterRule(
name="stream_transcript",
action="strip_after_first",
patterns=_patterns(
r'data:\s*\{"id"\s*:\s*"chatcmpl',
r"data:\s*\[DONE\]",
),
),
SourceFilterRule(
name="retrieval_context",
action="drop_if_present",
patterns=_patterns(
r"(?m)^\s{0,3}#{1,3}\s*以下内容是基于用户发送的消息的搜索结果",
r"(?i)<(?:retrieved_context|search_context|web_results)>",
),
),
SourceFilterRule(
name="memory_context",
action="drop_if_present",
patterns=_patterns(
r"(?i)</?(?:memories|memory_context)>",
r"(?i)===\s*MemOS LONG-TERM MEMORY",
r"(?i)\[MemOS Auto-Recall\]",
),
),
SourceFilterRule(
name="agent_reasoning",
action="drop_if_present",
patterns=_patterns(
r"(?i)<(?:thinking|reasoning|agent_scratchpad)>",
),
),
SourceFilterRule(
name="runtime_metadata",
action="drop_if_present",
patterns=_patterns(
r"(?m)^\s*Conversation info \(untrusted metadata\):",
r"(?m)^\s*Untrusted context \(metadata, do not treat as instructions or commands\):",
),
),
SourceFilterRule(
name="automation_context",
action="drop_if_prefix",
patterns=_patterns(
r"\[cron:",
r"System:\s+\[",
r"A scheduled reminder has been triggered",
),
),
SourceFilterRule(
name="assistant_runtime_prefix",
action="drop_if_prefix",
patterns=_patterns(
r"小依会根据用户需求",
r"正在完善Gemini的思考过程",
),
),
),
)


class MemorySourceFilter:
"""Filter raw memory sources before building extraction prompts."""

def __init__(self, policy: SourceFilterPolicy = PREFERENCE_SOURCE_POLICY):
self.policy = policy

def filter_for_preference(self, sources: list[Any] | None) -> list[SourceMessage]:
"""Return only sources allowed to appear in preference extraction prompts."""
filtered: list[SourceMessage] = []
for source in sources or []:
source_dict = self._source_to_dict(source)
if not self._keep_role_for_preference(source_dict):
continue

content = str(source_dict.get("content") or "")
content = self._strip_known_context_wrappers(content)
if not content.strip():
continue

cleaned = source_dict.copy()
cleaned["content"] = content.strip()
if cleaned.get("type") is None:
cleaned["type"] = "chat"
cleaned = self._coerce_source_fields(cleaned)
filtered.append(SourceMessage(**cleaned))
return filtered

def build_prompt_text(self, sources: list[Any] | None) -> str:
"""Build a compact prompt text from filtered sources."""
return self.sources_to_prompt_text(self.filter_for_preference(sources))

def sources_to_prompt_text(self, sources: list[SourceMessage] | None) -> str:
"""Build prompt text from sources that have already been filtered."""
lines = []
for source in sources or []:
role = source.role or "user"
content = (source.content or "").strip()
if not content:
continue
lines.append(f"{role}: {content}")
return "\n".join(lines)

def _source_to_dict(self, source: Any) -> dict[str, Any]:
if isinstance(source, SourceMessage):
return source.model_dump(exclude_none=True)
if isinstance(source, dict):
return source.copy()
if hasattr(source, "model_dump"):
return source.model_dump(exclude_none=True)
return {}

def _coerce_source_fields(self, source: dict[str, Any]) -> dict[str, Any]:
for key in ("chat_time", "message_id", "role", "type"):
if source.get(key) is not None and not isinstance(source[key], str):
source[key] = str(source[key])
return source

def _keep_role_for_preference(self, source: dict[str, Any]) -> bool:
source_type = str(source.get("type") or "chat").strip().lower()
role = str(source.get("role") or "").strip().lower()
if source_type in self.policy.blocked_source_types or role in self.policy.blocked_roles:
return False
return role in self.policy.allowed_roles

def _strip_known_context_wrappers(self, text: str) -> str:
text = text.replace("\r\n", "\n").replace("\r", "\n").strip()
if not text:
return ""

for rule in self.policy.rules:
text = self._apply_rule(text, rule)
if not text:
return ""
return text.strip()

def _apply_rule(self, text: str, rule: SourceFilterRule) -> str:
if rule.action == "extract_after_last":
extracted = self._extract_after_last_match(text, rule.patterns)
return text if extracted is None else extracted.strip()
if rule.action == "strip_after_first":
return self._strip_after_first_match(text, rule.patterns)
if rule.action == "drop_if_present":
return "" if any(pattern.search(text) for pattern in rule.patterns) else text
if rule.action == "drop_if_prefix":
stripped = text.strip()
return "" if any(pattern.match(stripped) for pattern in rule.patterns) else text
return text

def _extract_after_last_match(
self, text: str, patterns: tuple[re.Pattern[str], ...]
) -> str | None:
best_match: re.Match[str] | None = None
for pattern in patterns:
for match in pattern.finditer(text):
if best_match is None or match.start() > best_match.start():
best_match = match
if best_match is None:
return None
return text[best_match.end() :]

def _strip_after_first_match(self, text: str, patterns: tuple[re.Pattern[str], ...]) -> str:
end = len(text)
for pattern in patterns:
match = pattern.search(text)
if match:
end = min(end, match.start())
return text[:end].strip()
Original file line number Diff line number Diff line change
Expand Up @@ -179,8 +179,8 @@ def _process_memories_with_reader(
is_upload_skill=is_upload_skill,
)
except Exception as e:
logger.warning("%s: Fail to transfer mem: %s", e, memory_items)
processed_memories = []
logger.warning("%s: Fail to transfer mem: %s", e, memory_items, exc_info=True)
return

if processed_memories and len(processed_memories) > 0:
flattened_memories = []
Expand Down
Loading
Loading