From ff379b91cf46929918e94bf9780c3a7e4ec008f2 Mon Sep 17 00:00:00 2001 From: Xubin Ren <52506698+Re-bin@users.noreply.github.com> Date: Sun, 26 Jul 2026 17:15:41 +0800 Subject: [PATCH] fix(agent): preserve merged runtime context markers --- nanobot/agent/runner.py | 52 ++++++++++++- nanobot/runtime_context.py | 64 +++++++++++++++ tests/agent/test_runner_injections.py | 107 +++++++++++++++++++++++--- 3 files changed, 208 insertions(+), 15 deletions(-) diff --git a/nanobot/agent/runner.py b/nanobot/agent/runner.py index 060c0cef..9071a34a 100644 --- a/nanobot/agent/runner.py +++ b/nanobot/agent/runner.py @@ -19,6 +19,11 @@ from nanobot.agent.context_governance import ( from nanobot.agent.hook import AgentHook, AgentHookContext, AgentRunHookContext from nanobot.agent.tools.registry import ToolRegistry, is_tool_error_result from nanobot.providers.base import LLMProvider, LLMResponse, ToolCallRequest +from nanobot.runtime_context import ( + RUNTIME_CONTEXT_MESSAGE_META, + detach_runtime_context, + reattach_runtime_context, +) from nanobot.session.history_visibility import is_hidden_history_message from nanobot.utils.helpers import ( IncrementalThinkExtractor, @@ -137,10 +142,51 @@ class AgentRunner: and not is_hidden_history_message(messages[-1]) ): merged = dict(messages[-1]) - merged["content"] = cls._merge_message_content( - merged.get("content"), - injection.get("content"), + left_meta = merged.get("_meta") + right_meta = injection.get("_meta") + left_marker = ( + left_meta.get(RUNTIME_CONTEXT_MESSAGE_META) + if isinstance(left_meta, dict) + else None ) + right_marker = ( + right_meta.get(RUNTIME_CONTEXT_MESSAGE_META) + if isinstance(right_meta, dict) + else None + ) + detached_left = ( + detach_runtime_context(merged.get("content"), left_marker) + if isinstance(left_marker, dict) + else (merged.get("content"), [], []) + ) + detached_right = ( + detach_runtime_context(injection.get("content"), right_marker) + if isinstance(right_marker, dict) + else (injection.get("content"), [], []) + ) + if detached_left is not None and detached_right is not None: + left_content, left_sources, left_blocks = detached_left + right_content, right_sources, right_blocks = detached_right + merged_content = cls._merge_message_content(left_content, right_content) + context_blocks = [*left_blocks, *right_blocks] + if context_blocks: + merged_content, marker = reattach_runtime_context( + merged_content, + [*left_sources, *right_sources], + context_blocks, + ) + internal_meta = dict(left_meta) if isinstance(left_meta, dict) else {} + if isinstance(right_meta, dict): + for key, value in right_meta.items(): + internal_meta.setdefault(key, value) + internal_meta[RUNTIME_CONTEXT_MESSAGE_META] = marker + merged["_meta"] = internal_meta + merged["content"] = merged_content + else: + merged["content"] = cls._merge_message_content( + merged.get("content"), + injection.get("content"), + ) messages[-1] = merged continue messages.append(injection) diff --git a/nanobot/runtime_context.py b/nanobot/runtime_context.py index bb45e68e..29d9f6c0 100644 --- a/nanobot/runtime_context.py +++ b/nanobot/runtime_context.py @@ -139,6 +139,70 @@ def append_runtime_context( } +def detach_runtime_context( + content: Any, + marker: Mapping[str, Any], +) -> tuple[Any, list[str], list[dict[str, Any]]] | None: + """Detach one validated runtime-context suffix for safe message merging.""" + if marker.get("version") != 1: + return None + raw_sources = marker.get("sources") + sources = [ + source + for source in raw_sources + if isinstance(source, str) and source + ] if isinstance(raw_sources, list) else [] + + suffix = marker.get("suffix") + if isinstance(content, str) and isinstance(suffix, str) and suffix: + if content == suffix: + clean_content = "" + elif content.endswith("\n\n" + suffix): + clean_content = content[: -(len(suffix) + 2)] + else: + return None + return clean_content, sources, [{"type": "text", "text": suffix}] + + expected = marker.get("blocks") + if isinstance(content, list) and isinstance(expected, list) and expected: + count = len(expected) + if content[-count:] != expected: + return None + return content[:-count], sources, deepcopy(expected) + return None + + +def reattach_runtime_context( + content: Any, + sources: Sequence[str], + blocks: Sequence[Mapping[str, Any]], +) -> tuple[Any, dict[str, Any]]: + """Append detached runtime-context blocks after visible messages are merged.""" + context_blocks = [deepcopy(dict(block)) for block in blocks] + if isinstance(content, str) and all( + block.get("type") == "text" and isinstance(block.get("text"), str) + for block in context_blocks + ): + suffix = "\n\n".join(block["text"] for block in context_blocks) + merged = f"{content}\n\n{suffix}" if content else suffix + return merged, { + "version": 1, + "sources": list(sources), + "suffix": suffix, + } + + visible_blocks = ( + [*content] + if isinstance(content, list) + else ([] if content is None else [{"type": "text", "text": str(content)}]) + ) + return [*visible_blocks, *context_blocks], { + "version": 1, + "sources": list(sources), + "blocks": context_blocks, + } + + def public_history_message(message: Mapping[str, Any]) -> dict[str, Any]: """Return a user-visible copy with trusted runtime context removed exactly.""" cleaned = deepcopy(dict(message)) diff --git a/tests/agent/test_runner_injections.py b/tests/agent/test_runner_injections.py index 9270ba1e..9b96b4a8 100644 --- a/tests/agent/test_runner_injections.py +++ b/tests/agent/test_runner_injections.py @@ -527,6 +527,17 @@ async def test_pending_injection_resolves_its_own_runtime_context(tmp_path): "thread_id": "topic-7", }, )) + await pending_queue.put(InboundMessage( + channel="telegram", + sender_id="user-c", + chat_id="group-1", + content="another follow-up", + metadata={ + "message_id": "message-3", + "sender_name": "Carol", + "thread_id": "topic-7", + }, + )) _, _, all_messages, _, _ = await loop._run_agent_loop( [{"role": "user", "content": "initial message from user A"}], @@ -538,28 +549,48 @@ async def test_pending_injection_resolves_its_own_runtime_context(tmp_path): pending_queue=pending_queue, ) - assert seen_contexts == [( - "telegram", - "group-1", - "user-b", - "message-2", - session.key, - "follow-up from the second speaker", - "Bob", - "topic-7", - )] + assert seen_contexts == [ + ( + "telegram", + "group-1", + "user-b", + "message-2", + session.key, + "follow-up from the second speaker", + "Bob", + "topic-7", + ), + ( + "telegram", + "group-1", + "user-c", + "message-3", + session.key, + "another follow-up", + "Carol", + "topic-7", + ), + ] injected = [message for message in all_messages if message.get("role") == "user"][-1] assert "follow-up from the second speaker" in str(injected["content"]) model_messages = provider.chat_with_retry.await_args_list[-1].kwargs["messages"] assert "telegram | group-1 | user-b | message-2" in str(model_messages) assert "Bob | topic-7" in str(model_messages) - assert injected["_meta"][RUNTIME_CONTEXT_MESSAGE_META]["sources"] == ["identity"] + assert "telegram | group-1 | user-c | message-3" in str(model_messages) + assert "Carol | topic-7" in str(model_messages) + assert injected["_meta"][RUNTIME_CONTEXT_MESSAGE_META]["sources"] == [ + "identity", + "identity", + ] loop._save_turn(session, all_messages, skip=1) persisted = [message for message in session.messages if message.get("role") == "user"][-1] assert "telegram | group-1 | user-b | message-2" in str(persisted["content"]) - assert public_history_message(persisted)["content"] == "follow-up from the second speaker" + assert "telegram | group-1 | user-c | message-3" in str(persisted["content"]) + assert public_history_message(persisted)["content"] == ( + "follow-up from the second speaker\n\nanother follow-up" + ) @pytest.mark.asyncio @@ -688,6 +719,58 @@ async def test_runner_merges_multiple_injected_user_messages_without_losing_medi ) +def test_runner_merge_preserves_runtime_markers_with_media() -> None: + from nanobot.agent.runner import AgentRunner + from nanobot.runtime_context import ( + RUNTIME_CONTEXT_HISTORY_META, + RUNTIME_CONTEXT_MESSAGE_META, + RuntimeContextBlock, + append_runtime_context, + public_history_message, + ) + + first_visible = [ + {"type": "text", "text": "first"}, + {"type": "image_url", "image_url": {"url": "data:image/png;base64,AA=="}}, + ] + first_content, first_marker = append_runtime_context( + first_visible, + [RuntimeContextBlock(source="first", content="private first")], + ) + second_content, second_marker = append_runtime_context( + "second", + [RuntimeContextBlock(source="second", content="private second")], + ) + messages: list[dict] = [] + + AgentRunner._append_injected_messages(messages, [ + { + "role": "user", + "content": first_content, + "_meta": {RUNTIME_CONTEXT_MESSAGE_META: first_marker}, + }, + { + "role": "user", + "content": second_content, + "_meta": {RUNTIME_CONTEXT_MESSAGE_META: second_marker}, + }, + ]) + + assert len(messages) == 1 + merged = messages[0] + assert "private first" in str(merged["content"]) + assert "private second" in str(merged["content"]) + persisted = { + "role": "user", + "content": merged["content"], + RUNTIME_CONTEXT_HISTORY_META: merged["_meta"][RUNTIME_CONTEXT_MESSAGE_META], + } + assert public_history_message(persisted)["content"] == [ + *first_visible, + {"type": "text", "text": "second"}, + ] + + @pytest.mark.asyncio async def test_injection_cycles_capped_at_max(): """Injection cycles should be capped at _MAX_INJECTION_CYCLES."""