diff --git a/runtime/application/gateway/inbound_service.py b/runtime/application/gateway/inbound_service.py index 59fa3f1e..5b91f034 100644 --- a/runtime/application/gateway/inbound_service.py +++ b/runtime/application/gateway/inbound_service.py @@ -1193,6 +1193,7 @@ def process_inbound_payload(payload: dict[str, Any]) -> dict[str, Any]: target={"is_group": bool(inbound.is_group), "mention_all": bool(mention_all)}, ) d = pe.decide_action(ctx=act) + channel_is_wa = str(inbound.channel or "").strip().lower() == "whatsapp" if d.needs_confirmation: token_key = f"confirm_token:{session_id}" token = (store.get_setting(token_key) or "").strip() @@ -1200,12 +1201,22 @@ def process_inbound_payload(payload: dict[str, Any]) -> dict[str, Any]: token = pe.new_confirmation_token() store.set_setting(token_key, token) if not has_explicit_confirmation_token(inbound.text, token): - reply = f"该动作需要确认。请回复 `confirm {token}` 或包含 `[confirm:{token}]`。" + if channel_is_wa: + reply = ( + f"This action needs confirmation. " + f"Reply `confirm {token}` or include `[confirm:{token}]`." + ) + else: + reply = f"该动作需要确认。请回复 `confirm {token}` 或包含 `[confirm:{token}]`。" else: reply = f"[assistant] ok scope={scope} session={session_id[:8]} (confirmed)" else: if not _role_can_write(role, inbound.text): - reply = "你的角色暂无写入权限。请联系管理员提升权限。" + reply = ( + "Your role cannot write yet. Ask an administrator to raise your permissions." + if channel_is_wa + else "你的角色暂无写入权限。请联系管理员提升权限。" + ) else: cmd_reply = _handle_productivity_commands( text=inbound.text, @@ -1410,6 +1421,39 @@ def process_inbound_payload(payload: dict[str, Any]) -> dict[str, Any]: turn_handle = gate.try_begin(str(session_id), turn_job) if turn_handle is None: # Another turn is active; this message will be merged after it. + if ( + str(inbound.channel or "").strip().lower() == "whatsapp" + and _whatsapp_inbound_queue_delivery_enabled() + and gate.pending_count(str(session_id)) == 1 + ): + try: + busy_meta = None + if bool(getattr(inbound, "is_group", False)): + from runtime.application.gateway.whatsapp_progress import ( + build_whatsapp_group_progress_metadata, + ) + + 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." + ) + _enqueue_whatsapp_inbound_reply( + store, + inbound=inbound, + account_id=account_id, + tenant_id=str(tenant_id or ""), + reply_text=busy_text, + reply_attachments=None, + reply_metadata=busy_meta, + kind="inbound_progress", + ) + except Exception: + pass return { "ok": True, "replies": [], diff --git a/runtime/application/gateway/ops_short_intent.py b/runtime/application/gateway/ops_short_intent.py new file mode 100644 index 00000000..35129ff3 --- /dev/null +++ b/runtime/application/gateway/ops_short_intent.py @@ -0,0 +1,99 @@ +"""Detect ops short intents for WhatsApp/field English recipes.""" + +from __future__ import annotations + +import re +from typing import Any + +_BOT_MENTION_RE = re.compile(r"@\S+") + +# intent -> (en hint, zh hint) +_HINTS: dict[str, tuple[str, str]] = { + "fiber_cut": ( + "[Ops short-intent: fiber/LOS. Prefer ume_alarm_xlsx_report(mode=fiber_cut) in ≤3 tool calls; " + "do not paginate or re-list CLI targets first.]", + "[短指令:断纤/LOS。优先 ume_alarm_xlsx_report(mode=fiber_cut),≤3 次工具;勿先翻页或反复 listCliTargets。]", + ), + "offline": ( + "[Ops short-intent: offline NE. Prefer ume_alarm_xlsx_report(mode=offline) in ≤3 tool calls.]", + "[短指令:离线网元。优先 ume_alarm_xlsx_report(mode=offline),≤3 次工具。]", + ), + "alarm_tally": ( + "[Ops short-intent: alarm tally/top. Prefer aggregateUmeAlarms or " + "ume_alarm_xlsx_report(mode=aggregate_by_host); ≤3 tool calls.]", + "[短指令:告警统计/Top。优先 aggregateUmeAlarms 或 ume_alarm_xlsx_report(mode=aggregate_by_host);≤3 次工具。]", + ), + "excel_export": ( + "[Ops short-intent: Excel export. Prefer ume_alarm_xlsx_report or write_xlsx(deliverable=true); " + "do not build xlsx via run_command.]", + "[短指令:导出 Excel。优先 ume_alarm_xlsx_report 或 write_xlsx(deliverable=true);禁止 run_command 造表。]", + ), + "continue": ( + "[Ops short-intent: continue/confirm. Resume the unfinished prior task immediately; " + "do not re-ask confirmation or restart the query from scratch.]", + "[短指令:继续/确认。立即承接上一未完成任务;勿再问确认或重开查询。]", + ), +} + + +def normalize_ops_user_text(text: str) -> str: + s = str(text or "").strip() + if not s: + return "" + s = _BOT_MENTION_RE.sub(" ", s) + return re.sub(r"\s+", " ", s).strip().lower() + + +def detect_ops_short_intent(text: str) -> str | None: + """Return a recipe key for ultra-short field asks, else None.""" + t = normalize_ops_user_text(text) + if not t or len(t) > 160: + return None + + if re.fullmatch(r"(yes|y|ok|okay|confirm|continue|please continue|go ahead|可以|确认|继续|好的|行)", t): + return "continue" + if any(k in t for k in ("fiber", "los", "cable cut", "optical", "断纤", "光缆", "光路")): + return "fiber_cut" + if any(k in t for k in ("offline", "board offline", "ne communication", "离线", "单板离线", "通信中断")): + return "offline" + if any(k in t for k in ("excel", "xlsx", "spreadsheet", "export", "send me the table", "表格", "导出")): + return "excel_export" + if any( + k in t + for k in ( + "alarm", + "critical", + "tally", + "top ne", + "how many", + "告警", + "统计", + "多少", + ) + ): + return "alarm_tally" + return None + + +def build_ops_short_intent_hint(*, intent: str, lang: str = "en") -> str: + key = str(intent or "").strip() + pair = _HINTS.get(key) + if not pair: + return "" + en, zh = pair + return zh if str(lang or "").strip().lower().startswith("zh") else en + + +def maybe_ops_short_intent_system_hint(*, text: str, lang: str = "en") -> str: + intent = detect_ops_short_intent(text) + if not intent: + return "" + return build_ops_short_intent_hint(intent=intent, lang=lang) + + +__all__ = [ + "build_ops_short_intent_hint", + "detect_ops_short_intent", + "maybe_ops_short_intent_system_hint", + "normalize_ops_user_text", +] diff --git a/runtime/gateway.py b/runtime/gateway.py index 13ba5ed0..6be61542 100644 --- a/runtime/gateway.py +++ b/runtime/gateway.py @@ -478,6 +478,12 @@ class OclawGateway: return str(build_channel_file_delivery_instruction(lang=lang) or "").strip() + @staticmethod + def _ops_short_intent_system_hint(msg: StandardMessage, lang: str) -> str: + from runtime.application.gateway.ops_short_intent import maybe_ops_short_intent_system_hint + + return str(maybe_ops_short_intent_system_hint(text=str(msg.text or ""), lang=lang) or "").strip() + @staticmethod def _tabular_query_system_hint(lang: str) -> str: limits = OclawGateway._tabular_limits_from_config() @@ -1174,6 +1180,10 @@ class OclawGateway: ch_hint = self._channel_file_delivery_system_hint(lang) if ch_hint: sys_prompt = f"{sys_prompt}\n\n{ch_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: + sys_prompt = f"{sys_prompt}\n\n{intent_hint}".strip() if self._has_tabular_ref_attachments(msg): sys_prompt = f"{sys_prompt}\n\n{self._tabular_query_system_hint(lang)}".strip() if self._has_text_ref_attachments(msg): diff --git a/runtime/tools/mcp/adapter.py b/runtime/tools/mcp/adapter.py index 4ca96e6e..80eb952e 100644 --- a/runtime/tools/mcp/adapter.py +++ b/runtime/tools/mcp/adapter.py @@ -139,14 +139,9 @@ class _McpBoundTool: "instead of listing again." ) if tool_name == "execManagedNe" and res.get("ok") is False: - err = str(res.get("error") or res.get("error_code") or "") - low = err.lower() - if "timeout" in low or res.get("error_code") == "tool_timeout_or_failed": - res = dict(res) - res["hint"] = ( - "CLI timed out. Raise read_timeout_sec (60–120), reduce commands, " - "or reuse prior listCliTargets ids — do not blind-retry identical calls." - ) + from runtime.tools.tool_error_hints import enrich_exec_managed_ne_error + + res = enrich_exec_managed_ne_error(res) return res return ToolSpec( diff --git a/runtime/tools/tool_error_hints.py b/runtime/tools/tool_error_hints.py index 9784b89c..4ed442e2 100644 --- a/runtime/tools/tool_error_hints.py +++ b/runtime/tools/tool_error_hints.py @@ -93,7 +93,114 @@ def enrich_mcp_scope_error(result: dict[str, Any]) -> dict[str, Any]: return out +def _unwrap_nested_error_blob(raw: Any) -> tuple[str, str]: + """Best-effort unwrap double-encoded MCP/netx error payloads.""" + import json + + text = str(raw or "").strip() + code = "" + if not text: + return "", "" + cur: Any = text + for _ in range(3): + if isinstance(cur, dict): + code = str(cur.get("error_code") or cur.get("code") or code or "") + nested = cur.get("error") + if nested is None and cur.get("data") is not None: + nested = cur.get("data") + if isinstance(nested, (dict, list)): + cur = nested + continue + if nested is not None: + text = str(nested) + cur = nested + if isinstance(cur, str) and cur.strip().startswith("{"): + try: + cur = json.loads(cur) + continue + except Exception: + break + break + if isinstance(cur, str) and cur.strip().startswith("{"): + try: + cur = json.loads(cur) + continue + except Exception: + text = cur + break + if isinstance(cur, str): + text = cur + break + if isinstance(cur, dict): + text = str(cur.get("error") or cur.get("message") or text) + code = str(cur.get("error_code") or cur.get("code") or code or "") + return str(text or "").strip(), str(code or "").strip() + + +def enrich_exec_managed_ne_error(result: dict[str, Any]) -> dict[str, Any]: + """Classify execManagedNe failures so agents stop blind-retrying.""" + if not isinstance(result, dict) or result.get("ok") is not False: + return result + out = dict(result) + raw_err = out.get("error") + raw_code = str(out.get("error_code") or "") + unwrapped, nested_code = _unwrap_nested_error_blob(raw_err) + if unwrapped and unwrapped != str(raw_err or "").strip(): + out["error_detail"] = unwrapped + code = (nested_code or raw_code or "").strip() + blob = f"{unwrapped} {code} {raw_err}".lower() + + error_class = "exec_failed" + hint = ( + "CLI failed. Check ne_id/ume_ne_id, avoid identical blind retries, " + "and prefer batching show commands in one execManagedNe call." + ) + if "timeout" in blob or code in {"tool_timeout_or_failed", "read_timeout", "deadline_exceeded"}: + error_class = "timeout" + hint = ( + "CLI timed out. Raise read_timeout_sec (60–120), reduce commands, " + "or reuse prior listCliTargets ids — do not blind-retry identical calls." + ) + elif any( + x in blob + for x in ( + "unreachable", + "connection refused", + "no route", + "timed out connecting", + "host unreachable", + "network is unreachable", + "connect_failed", + "ssh_connect", + ) + ): + error_class = "unreachable" + hint = ( + "Device unreachable / connect failed. Do not spam retries; report the NE as unreachable " + "and try another target or verify UME→CLI credentials/jump host." + ) + elif any(x in blob for x in ("auth", "permission denied", "login failed", "authentication", "password")): + error_class = "auth" + hint = ( + "CLI authentication failed. Do not retry the same credentials; " + "fix UME→CLI / managed-NE credentials instead." + ) + elif any(x in blob for x in ("command", "syntax", "invalid input", "ambiguous command", "%error")): + error_class = "command_error" + hint = ( + "Command rejected by the device. Fix the CLI syntax or vendor dialect; " + "do not retry the identical command string." + ) + + out["error_class"] = error_class + if code and not out.get("error_code"): + out["error_code"] = code + out["hint"] = hint + return out + + __all__ = [ + "enrich_exec_managed_ne_error", "enrich_mcp_scope_error", "format_unregistered_tool_error", "suggest_tool_names", diff --git a/tests/test_ops_short_intent_and_exec_hints.py b/tests/test_ops_short_intent_and_exec_hints.py new file mode 100644 index 00000000..b05ab46d --- /dev/null +++ b/tests/test_ops_short_intent_and_exec_hints.py @@ -0,0 +1,41 @@ +from __future__ import annotations + +from runtime.application.gateway.ops_short_intent import ( + detect_ops_short_intent, + maybe_ops_short_intent_system_hint, +) +from runtime.tools.tool_error_hints import enrich_exec_managed_ne_error + + +def test_detect_ops_short_intent_english_field() -> None: + assert detect_ops_short_intent("@bot fiber cut sites") == "fiber_cut" + assert detect_ops_short_intent("offline NE list") == "offline" + assert detect_ops_short_intent("please continue") == "continue" + assert detect_ops_short_intent("YES") == "continue" + assert detect_ops_short_intent("export excel") == "excel_export" + assert detect_ops_short_intent("critical top alarms") == "alarm_tally" + assert detect_ops_short_intent("hello there how are you doing today with something else") is None + + +def test_ops_short_intent_hint_english_default() -> None: + hint = maybe_ops_short_intent_system_hint(text="LOS on these sites", lang="en") + assert "fiber" in hint.lower() or "LOS" in hint + assert "ume_alarm_xlsx_report" in hint + assert "断纤" not in hint + + +def test_enrich_exec_timeout() -> None: + out = enrich_exec_managed_ne_error( + {"ok": False, "error_code": "tool_timeout_or_failed", "error": "timeout"} + ) + assert out["error_class"] == "timeout" + assert "read_timeout_sec" in out["hint"] + + +def test_enrich_exec_unreachable_nested_json() -> None: + nested = '{"ok": false, "error": "ssh_connect failed: host unreachable"}' + out = enrich_exec_managed_ne_error( + {"ok": False, "error_code": "mcp_tool_call_failed", "error": nested} + ) + assert out["error_class"] == "unreachable" + assert "unreachable" in out["hint"].lower() diff --git a/tests/test_tool_error_hints.py b/tests/test_tool_error_hints.py index c1060422..e37af5a3 100644 --- a/tests/test_tool_error_hints.py +++ b/tests/test_tool_error_hints.py @@ -36,3 +36,10 @@ def test_enrich_mcp_scope_sql() -> None: assert out["required_scope"] == "sql:query" assert "fallback_tools" in out assert "ume_alarm_xlsx_report" in out["fallback_tools"] + + +def test_enrich_exec_auth() -> None: + from runtime.tools.tool_error_hints import enrich_exec_managed_ne_error + + out = enrich_exec_managed_ne_error({"ok": False, "error": "authentication failed: Permission denied"}) + assert out["error_class"] == "auth" diff --git a/tests/test_whatsapp_inbound_queue_cancel.py b/tests/test_whatsapp_inbound_queue_cancel.py index 3cad1fa4..998308b1 100644 --- a/tests/test_whatsapp_inbound_queue_cancel.py +++ b/tests/test_whatsapp_inbound_queue_cancel.py @@ -235,10 +235,19 @@ class WhatsappInboundSerialQueueTests(unittest.TestCase): pending = self.store.list_pending_channel_outbound_messages( channel="whatsapp", account_id="wa-default", limit=10 ) - self.assertEqual(len(pending), 2) + self.assertEqual(len(pending), 3) texts = " ".join(str(p.get("text") or "") for p in pending) self.assertIn("ans:1", texts) self.assertIn("ans:2", texts) + self.assertIn("Still working", texts) + kinds = [] + for p in pending: + try: + kinds.append(json.loads(str(p.get("source") or "{}")).get("kind")) + except Exception: + kinds.append(None) + self.assertIn("inbound_progress", kinds) + self.assertEqual(kinds.count("inbound_reply"), 2) if __name__ == "__main__":