diff --git a/runtime/application/gateway/channel_turn_gate.py b/runtime/application/gateway/channel_turn_gate.py index 7307611d..6283b8dd 100644 --- a/runtime/application/gateway/channel_turn_gate.py +++ b/runtime/application/gateway/channel_turn_gate.py @@ -30,16 +30,30 @@ def merge_channel_pending_jobs(jobs: list[dict[str, Any]]) -> dict[str, Any]: lang = str(rows[-1].get("lang") or "").strip().lower() if lang.startswith("zh"): - header = "处理上一问期间又收到多条跟进,请一并回答:" + header = "处理上一问期间又收到多条跟进,请按发言人分别回答:" else: - header = "Several follow-up questions arrived while the previous request was still running. Please answer them together:" + header = ( + "Several follow-up questions arrived while the previous request was still running. " + "Answer each item for its own sender:" + ) parts: list[str] = [] attachments: list[dict[str, Any]] = [] for idx, job in enumerate(rows, start=1): text = str(job.get("user_text") or "").strip() + md = job.get("msg_metadata") if isinstance(job.get("msg_metadata"), dict) else {} + sender = str( + md.get("group_sender_label") + or md.get("push_name") + or md.get("group_sender_id") + or md.get("external_user_id") + or "" + ).strip() if text: - parts.append(f"{idx}) {text}") + if sender and not text.lstrip().startswith(("[Sender:", "[发言:")): + parts.append(f"{idx}) [{sender}] {text}") + else: + parts.append(f"{idx}) {text}") for att in job.get("attachments") or []: if isinstance(att, dict): attachments.append(att) diff --git a/runtime/application/gateway/inbound_service.py b/runtime/application/gateway/inbound_service.py index 0342d57c..519b9c50 100644 --- a/runtime/application/gateway/inbound_service.py +++ b/runtime/application/gateway/inbound_service.py @@ -1376,6 +1376,7 @@ def process_inbound_payload(payload: dict[str, Any]) -> dict[str, Any]: tenant_id=str(tenant_id or ""), account_id=wa_acct, quoted_ctx=quoted_ctx, + lang=lang, ) if not user_text and gw_attachments: user_text = ( @@ -1388,6 +1389,17 @@ def process_inbound_payload(payload: dict[str, Any]) -> dict[str, Any]: from runtime.types import StandardMessage from runtime.application.gateway.channel_turn_gate import get_channel_turn_gate + raw_meta = ( + meta_for_inbound.get("raw") + if isinstance(meta_for_inbound.get("raw"), dict) + else {} + ) + sender_label = str( + raw_meta.get("pushName") + or meta_for_inbound.get("push_name") + or inbound.external_user_id + or "" + ).strip() msg_metadata = { "tenant_id": tenant_id, "user_id": user_id, @@ -1400,6 +1412,9 @@ def process_inbound_payload(payload: dict[str, Any]) -> dict[str, Any]: "external_user_id": inbound.external_user_id, "external_chat_id": inbound.external_chat_id, "group_sender_id": inbound.external_user_id, + "group_sender_label": sender_label, + "push_name": sender_label, + "group_session_scope": group_policy.session_scope, "mentioned_jids": mention_jids_for_ctx, "mention_names": mention_names_for_ctx, "raw_inbound_text": raw_user_text, @@ -1436,12 +1451,20 @@ def process_inbound_payload(payload: dict[str, Any]) -> dict[str, Any]: busy_meta = build_whatsapp_group_progress_metadata( inbound=inbound ) - busy_text = ( - "还在处理上一条请求,这条会排队合并处理。" - if str(lang or "").startswith("zh") - else "Still working on your previous request; " - "I'll merge this follow-up next." - ) + if str(lang or "").startswith("zh"): + busy_text = ( + "本会话还有请求在处理,你这条已排队,完成后会继续回答。" + if bool(getattr(inbound, "is_group", False)) + else "还在处理上一条请求,这条会排队合并处理。" + ) + else: + busy_text = ( + "Still working on an earlier request in this chat; " + "your message is queued and will be answered next." + if bool(getattr(inbound, "is_group", False)) + else "Still working on your previous request; " + "I'll merge this follow-up next." + ) _enqueue_whatsapp_inbound_reply( store, inbound=inbound, diff --git a/runtime/gateway.py b/runtime/gateway.py index 6be61542..079f14a8 100644 --- a/runtime/gateway.py +++ b/runtime/gateway.py @@ -484,6 +484,15 @@ class OclawGateway: return str(maybe_ops_short_intent_system_hint(text=str(msg.text or ""), lang=lang) or "").strip() + @staticmethod + def _group_focus_system_hint(msg: StandardMessage, lang: str) -> str: + md = msg.metadata if isinstance(msg.metadata, dict) else {} + if not bool(md.get("is_group")): + return "" + from runtime.orchestration.group_ingest import build_group_focus_instruction + + return str(build_group_focus_instruction(lang=lang) or "").strip() + @staticmethod def _tabular_query_system_hint(lang: str) -> str: limits = OclawGateway._tabular_limits_from_config() @@ -1180,6 +1189,9 @@ class OclawGateway: ch_hint = self._channel_file_delivery_system_hint(lang) if ch_hint: sys_prompt = f"{sys_prompt}\n\n{ch_hint}".strip() + focus_hint = self._group_focus_system_hint(msg, lang) + if focus_hint: + sys_prompt = f"{sys_prompt}\n\n{focus_hint}".strip() if str(manager_specialist or requested_specialist or "").strip().lower() == "ops": intent_hint = self._ops_short_intent_system_hint(msg, lang) if intent_hint: diff --git a/runtime/orchestration/group_ingest.py b/runtime/orchestration/group_ingest.py index 82c5f550..5d22da38 100644 --- a/runtime/orchestration/group_ingest.py +++ b/runtime/orchestration/group_ingest.py @@ -415,13 +415,20 @@ def should_process_group_inbound( return not require_mention -def build_group_sender_context(*, metadata: dict[str, Any] | None, external_user_id: str) -> str: +def build_group_sender_context( + *, + metadata: dict[str, Any] | None, + external_user_id: str, + lang: str = "en", +) -> str: meta = metadata if isinstance(metadata, dict) else {} raw = meta.get("raw") if isinstance(meta.get("raw"), dict) else {} push_name = str(raw.get("pushName") or meta.get("push_name") or meta.get("display_name") or "").strip() sender = str(external_user_id or "").strip() label = push_name or sender or "unknown" - return f"[发言: {label}]" + if str(lang or "").strip().lower().startswith("zh"): + return f"[发言: {label}]" + return f"[Sender: {label}]" def _mention_tokens_for_jids(jids: list[str]) -> list[str]: @@ -565,6 +572,7 @@ def prepare_group_user_text_for_model( tenant_id: str = "", account_id: str = "", quoted_ctx: str = "", + lang: str = "en", ) -> str: body = str(text or "").strip() body = strip_bot_mentions_from_text(text=body, bot_jid=bot_jid, metadata=metadata) @@ -582,10 +590,11 @@ def prepare_group_user_text_for_model( account_id=account_id, ) prefix_parts: list[str] = [] - if normalize_group_session_scope(session_scope) == "chat": - prefix_parts.append( - build_group_sender_context(metadata=metadata, external_user_id=external_user_id) - ) + # Always tag the speaker in groups (shared or per-user) so history/quotes cannot 串话. + _ = session_scope # retained for call-site compatibility / future policy forks + prefix_parts.append( + build_group_sender_context(metadata=metadata, external_user_id=external_user_id, lang=lang) + ) quote = str(quoted_ctx or "").strip() if quote: prefix_parts.append(quote) @@ -597,10 +606,15 @@ def prepare_group_user_text_for_model( def build_group_focus_instruction(*, lang: str = "en") -> str: if str(lang or "").strip().lower().startswith("zh"): - return "[群聊规则:只回答当前发言人的问题;除非本条消息明确引用或承接前文,否则不要默认继承其他群成员的上下文。]" + return ( + "[群聊规则:只回答当前发言人(见 [发言: …])的问题;" + "除非本条消息明确引用或承接前文,否则不要默认继承其他群成员的上下文;" + "若排队合并了多条跟进,按条分别回应并 @ 对应发言人。]" + ) return ( - "[Group chat rule: answer only the current sender's request. " - "Do not assume context from other members unless this message explicitly quotes or references it.]" + "[Group chat rule: answer only the current sender (see [Sender: …]). " + "Do not assume context from other members unless this message explicitly quotes or references it. " + "If several follow-ups were merged while busy, answer each item and @ the matching sender.]" ) diff --git a/runtime/tools/tool_error_hints.py b/runtime/tools/tool_error_hints.py index d8600742..ad1088b9 100644 --- a/runtime/tools/tool_error_hints.py +++ b/runtime/tools/tool_error_hints.py @@ -77,19 +77,29 @@ def enrich_mcp_scope_error(result: dict[str, Any]) -> dict[str, Any]: out["error"] = f"insufficient_scope:{scope}" if scope else "insufficient_scope" if scope == "sql:query": out["hint"] = ( - "Current netx token lacks sql:query. Prefer aggregateUmeAlarms / queryUmeAlarmsRaw / " - "ume_alarm_xlsx_report; ask an admin to grant sql:query only if SQL is required." + "Current netx token lacks sql:query. Do not retry sqlQueryUme. " + "Prefer aggregateUmeAlarms / queryUmeAlarmsRaw / ume_alarm_xlsx_report; " + "ask an admin to grant sql:query only if SQL is truly required." ) out["fallback_tools"] = [ "mcp__netx__aggregateUmeAlarms", "mcp__netx__queryUmeAlarmsRaw", "ume_alarm_xlsx_report", ] + out["user_facing_hint"] = ( + "SQL query is not enabled for this bot token. " + "I will use alarm aggregate/report tools instead, or an admin can grant sql:query." + ) else: out["hint"] = ( f"Current netx token lacks scope {scope or '(unknown)'}. " - "Ask an admin to grant it on the netx API token, or use tools that do not need this scope." + "Do not retry the same tool; ask an admin to grant it, or use tools that do not need this scope." ) + out["user_facing_hint"] = ( + f"Permission missing ({scope or 'scope'}). An admin needs to grant this on the netx API token." + ) + out["failure_class"] = "auth" + out["retry_forbidden"] = True return out diff --git a/tests/test_channel_turn_gate.py b/tests/test_channel_turn_gate.py index 22024ec5..bad06850 100644 --- a/tests/test_channel_turn_gate.py +++ b/tests/test_channel_turn_gate.py @@ -48,9 +48,31 @@ def test_merge_channel_pending_jobs_zh_header() -> None: {"user_text": "还有告警吗?", "lang": "zh"}, ] ) - assert "一并回答" in str(merged.get("user_text") or "") - assert "1) 光功率?" in str(merged.get("user_text") or "") - assert "2) 还有告警吗?" in str(merged.get("user_text") or "") + text = str(merged.get("user_text") or "") + assert "分别回答" in text + assert "1) 光功率?" in text + assert "2) 还有告警吗?" in text + + +def test_merge_channel_pending_jobs_keeps_sender_labels() -> None: + merged = merge_channel_pending_jobs( + [ + { + "user_text": "LOS count?", + "lang": "en", + "msg_metadata": {"group_sender_label": "Alice"}, + }, + { + "user_text": "offline NEs?", + "lang": "en", + "msg_metadata": {"group_sender_label": "Bob"}, + }, + ] + ) + text = str(merged.get("user_text") or "") + assert "[Alice]" in text + assert "[Bob]" in text + assert "own sender" in text.lower() or "Answer each" in text def test_isolated_sessions_and_reset() -> None: diff --git a/tests/test_group_ingest.py b/tests/test_group_ingest.py index 54ffabb4..92f05a26 100644 --- a/tests/test_group_ingest.py +++ b/tests/test_group_ingest.py @@ -359,8 +359,14 @@ def test_build_group_sender_context() -> None: metadata={"raw": {"pushName": "Alice"}}, external_user_id="111@s.whatsapp.net", ) - assert ctx == "[发言: Alice]" + assert ctx == "[Sender: Alice]" assert "111@s.whatsapp.net" not in ctx + zh = build_group_sender_context( + metadata={"raw": {"pushName": "Alice"}}, + external_user_id="111@s.whatsapp.net", + lang="zh", + ) + assert zh == "[发言: Alice]" def test_strip_bot_mentions_from_text() -> None: @@ -385,16 +391,17 @@ def test_normalize_mentioned_users_in_text_uses_nickname() -> None: def test_prepare_group_user_text_for_model_user_in_chat() -> None: out = prepare_group_user_text_for_model( text="@162788605444170 每三分钟提醒@200846277140511 喝水", - metadata={"bot_lid": "162788605444170@lid", "bot_push_name": "oliver"}, + metadata={"bot_lid": "162788605444170@lid", "bot_push_name": "oliver", "raw": {"pushName": "Bob"}}, mentions=["162788605444170@lid", "200846277140511@lid"], bot_jid="8618142387786@s.whatsapp.net", session_scope="user_in_chat", external_user_id="8618142387786@s.whatsapp.net", filtered_mention_jids=["200846277140511@lid"], mention_names=["吴华"], + lang="zh", ) - assert out == "每三分钟提醒@吴华 喝水" - assert "[发言:" not in out + assert out.startswith("[发言: Bob]") + assert "每三分钟提醒@吴华 喝水" in out def test_prepare_group_user_text_for_model_shared_chat_prefix() -> None: @@ -405,8 +412,9 @@ def test_prepare_group_user_text_for_model_shared_chat_prefix() -> None: bot_jid="999@s.whatsapp.net", session_scope="chat", external_user_id="111@s.whatsapp.net", + lang="en", ) - assert out.startswith("[发言: Alice]") + assert out.startswith("[Sender: Alice]") assert out.endswith("帮忙") assert "@999" not in out @@ -415,6 +423,7 @@ def test_build_group_focus_instruction() -> None: en = build_group_focus_instruction() zh = build_group_focus_instruction(lang="zh") assert "current sender" in en + assert "[Sender:" in en or "Sender" in en assert "群聊规则" in zh @@ -742,8 +751,8 @@ def test_inbound_group_mention_uses_per_user_session_and_sender_prefix( assert len(session_ids) == 2 assert session_ids[0] != session_ids[1] assert "[群成员:" not in captured["text"] - assert "[发言:" not in captured["text"] - assert captured["text"] == "again" + assert captured["text"].startswith("[Sender: Bob]") + assert captured["text"].endswith("again") assert "群聊规则" not in captured["text"] sid = store.get_or_create_channel_session_v2( diff --git a/tests/test_oclaw_gateway_trace.py b/tests/test_oclaw_gateway_trace.py index 54f40d70..c2c553ec 100644 --- a/tests/test_oclaw_gateway_trace.py +++ b/tests/test_oclaw_gateway_trace.py @@ -954,3 +954,31 @@ def test_channel_file_delivery_hint_goes_to_system_not_user_message() -> None: metadata={}, ) assert not OclawGateway._is_channel_delivery_channel(msg_admin) + + +def test_group_focus_system_hint_only_for_groups() -> None: + from runtime.types import StandardMessage + + group_msg = StandardMessage( + session_id="s1", + tenant_id="t1", + user_id="u1", + role="member", + channel="whatsapp", + text="@bot alarms", + attachments=[], + metadata={"is_group": True}, + ) + dm = StandardMessage( + session_id="s1", + tenant_id="t1", + user_id="u1", + role="member", + channel="whatsapp", + text="alarms", + attachments=[], + metadata={"is_group": False}, + ) + hint = OclawGateway._group_focus_system_hint(group_msg, "en") + assert "current sender" in hint + assert OclawGateway._group_focus_system_hint(dm, "en") == "" diff --git a/tests/test_tool_error_hints.py b/tests/test_tool_error_hints.py index e37af5a3..463048b7 100644 --- a/tests/test_tool_error_hints.py +++ b/tests/test_tool_error_hints.py @@ -33,6 +33,10 @@ def test_enrich_mcp_scope_sql() -> None: raw = {"ok": False, "error_code": "mcp_rpc_error_-32001", "error": "insufficient_scope:sql:query"} out = enrich_mcp_scope_error(raw) assert out["error_code"] == "insufficient_scope" + assert out.get("failure_class") == "auth" + assert out.get("retry_forbidden") is True + assert "aggregateUmeAlarms" in str(out.get("hint") or "") + assert "SQL query is not enabled" in str(out.get("user_facing_hint") or "") assert out["required_scope"] == "sql:query" assert "fallback_tools" in out assert "ume_alarm_xlsx_report" in out["fallback_tools"]