mirror of
https://github.com/hansjone/oclaw.git
synced 2026-10-09 04:40:45 +08:00
Steer WhatsApp ops short intents and clarify busy/CLI failures.
Inject English recipe hints for fiber/offline/excel/continue, ack the first queued follow-up, localize gate copy, and classify execManagedNe errors so agents stop blind retries. Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
parent
47fd83597e
commit
f1d025ed60
8 changed files with 323 additions and 11 deletions
|
|
@ -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": [],
|
||||
|
|
|
|||
99
runtime/application/gateway/ops_short_intent.py
Normal file
99
runtime/application/gateway/ops_short_intent.py
Normal file
|
|
@ -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",
|
||||
]
|
||||
|
|
@ -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):
|
||||
|
|
|
|||
|
|
@ -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(
|
||||
|
|
|
|||
|
|
@ -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",
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue