From 76f62d8faca70a699c465e558f8c048c5d9b821e Mon Sep 17 00:00:00 2001 From: oliver Date: Tue, 11 Aug 2026 11:13:41 +0800 Subject: [PATCH] Cut ops turn idle loops and WhatsApp progress spam. Align tool validation with playbook recipes, add turn checklist/idle guard, and cap interim WA status ticks so multi-hop CLI turns stop flooding the group. Co-authored-by: Cursor --- docs/ENVIRONMENT_VARIABLES.md | 33 +++ docs/ENVIRONMENT_VARIABLES_CHANGELOG.md | 60 +++++ .../application/gateway/ops_short_intent.py | 11 +- .../application/gateway/whatsapp_progress.py | 72 ++++-- runtime/chat/tool_runtime.py | 8 + runtime/chat/turn_idle_guard.py | 209 ++++++++++++++++ runtime/direct_loop.py | 170 ++++++++++++- runtime/gateway.py | 8 + runtime/tools/playbook_contracts.py | 229 ++++++++++++++++++ runtime/tools/tool_validation.py | 17 +- .../test_playbook_contracts_and_idle_guard.py | 130 ++++++++++ tests/test_whatsapp_turn_progress.py | 40 ++- 12 files changed, 944 insertions(+), 43 deletions(-) create mode 100644 runtime/chat/turn_idle_guard.py create mode 100644 runtime/tools/playbook_contracts.py create mode 100644 tests/test_playbook_contracts_and_idle_guard.py diff --git a/docs/ENVIRONMENT_VARIABLES.md b/docs/ENVIRONMENT_VARIABLES.md index 4b2f1943..67ee1104 100644 --- a/docs/ENVIRONMENT_VARIABLES.md +++ b/docs/ENVIRONMENT_VARIABLES.md @@ -31,6 +31,22 @@ - 作用:单轮工具循环轮次上限(一轮 = 模型思考一次并执行一批工具) - 生效:`runtime/gateway.py`, `runtime/direct_loop.py`, `runtime/worker.py` +- `AIA_TURN_IDLE_GUARD` + - 默认:`1`(开启) + - 作用:turn 内 idle guard:短指令只叙述不调工具时 nudge 一次;连续无进展 / 参数校验失败过多时 early finalize,减少空转 + - 取值:`0` / `false` / `no` / `off` 关闭 + - 生效:`runtime/chat/turn_idle_guard.py`, `runtime/direct_loop.py` + +- `AIA_TURN_IDLE_MAX_ROUNDS` + - 默认:`2` + - 作用:连续「有工具调用但 0 成功」轮次达到该值且本 turn 尚无成功工具时,触发 early finalize + - 生效:`runtime/chat/turn_idle_guard.py` + +- `AIA_TURN_IDLE_MAX_SCHEMA_FAILS` + - 默认:`3` + - 作用:本 turn `tool_invalid_arguments` / schema_validation 累计次数达到该值且尚无成功工具时,触发 early finalize + - 生效:`runtime/chat/turn_idle_guard.py` + - `AIA_TURN_MAX_CONTEXT_MESSAGES` - 默认:`80` - 作用:上下文消息上限 @@ -547,6 +563,23 @@ - 作用:每用户最近事件回放缓冲上限(用于 `connect.params.lastSeq` 断线补偿) - 生效:`oclaw/interfaces/ws/common.py`, `oclaw/interfaces/ws/runtime_impl.py`, `oclaw/interfaces/ws/runtime_helpers.py` +## WhatsApp 中间进度 + +- `OCLAW_WHATSAPP_TURN_PROGRESS` + - 默认:开启(未设置或非 `0/false/no/off`) + - 作用:长 turn 是否向 WhatsApp 发中间进度(终稿不受影响) + - 生效:`runtime/application/gateway/whatsapp_progress.py`, `inbound_service.py` + +- `OCLAW_WHATSAPP_PROGRESS_MIN_INTERVAL_SEC` + - 默认:`45` + - 作用:两条中间进度的最小间隔(秒) + - 生效:`runtime/application/gateway/whatsapp_progress.py` + +- `OCLAW_WHATSAPP_PROGRESS_MAX_PER_TURN` + - 默认:`2` + - 作用:单次入站 turn 中间进度硬上限;同名 long tool(如多跳 `execManagedNe`)每 turn 只播报一次;不再发 mid-turn `tools done / composing` + - 生效:`runtime/application/gateway/whatsapp_progress.py` + ## WeCom 长连接 - `AIA_WECOM_LONGCONN_WORKERS` diff --git a/docs/ENVIRONMENT_VARIABLES_CHANGELOG.md b/docs/ENVIRONMENT_VARIABLES_CHANGELOG.md index 23096453..2e5d1dfb 100644 --- a/docs/ENVIRONMENT_VARIABLES_CHANGELOG.md +++ b/docs/ENVIRONMENT_VARIABLES_CHANGELOG.md @@ -18,6 +18,66 @@ --- +## 2026-08-11 / WhatsApp progress noise reduction + +### Added +- `OCLAW_WHATSAPP_PROGRESS_MAX_PER_TURN` + - 默认值:`2` + - 用途:单次入站 turn 的中间进度消息硬上限(终稿回复不计) + - 影响模块:`runtime/application/gateway/whatsapp_progress.py` + +### Changed +- `OCLAW_WHATSAPP_PROGRESS_MIN_INTERVAL_SEC` + - 变更前:默认 `12` + - 变更后:默认 `45` + - 影响:中间进度更稀疏;仍可用环境变量覆盖 + - 是否需要重启:是 +- WhatsApp progress 行为:不再转发 mid-turn `tools done / composing`;同一 long tool(如重复 `execManagedNe`)每 turn 只播报一次 + +### Deprecated +- (无) + +### Removed +- (无) + +### Migration Checklist +- [ ] 重启 gateway +- [ ] 现场长 CLI 路径查询验证:中间进度 ≤2 条,终稿仍正常 +- [ ] 若需完全关闭进度:`OCLAW_WHATSAPP_TURN_PROGRESS=0` + +--- + +## 2026-08-11 / turn idle guard + playbook contracts + +### Added +- `AIA_TURN_IDLE_GUARD` + - 默认值:`1` + - 用途:开启 turn 内 idle guard(短指令 narration nudge / 无进展 early finalize) + - 影响模块:`runtime/chat/turn_idle_guard.py`, `runtime/direct_loop.py` +- `AIA_TURN_IDLE_MAX_ROUNDS` + - 默认值:`2` + - 用途:连续无成功工具轮次阈值 + - 影响模块:`runtime/chat/turn_idle_guard.py` +- `AIA_TURN_IDLE_MAX_SCHEMA_FAILS` + - 默认值:`3` + - 用途:schema 校验失败累计阈值 + - 影响模块:`runtime/chat/turn_idle_guard.py` + +### Changed +- (无) + +### Deprecated +- (无) + +### Removed +- (无) + +### Migration Checklist +- [ ] 重启 gateway / worker 后生效 +- [ ] 若需关闭:设 `AIA_TURN_IDLE_GUARD=0` + +--- + ## 模板 ```md diff --git a/runtime/application/gateway/ops_short_intent.py b/runtime/application/gateway/ops_short_intent.py index 26ef6fab..3c228319 100644 --- a/runtime/application/gateway/ops_short_intent.py +++ b/runtime/application/gateway/ops_short_intent.py @@ -162,7 +162,16 @@ def build_ops_short_intent_hint(*, intent: str, lang: str = "en") -> str: if not pair: return "" en, zh = pair - return zh if str(lang or "").strip().lower().startswith("zh") else en + base = zh if str(lang or "").strip().lower().startswith("zh") else en + try: + from runtime.tools.playbook_contracts import build_turn_checklist + + checklist = build_turn_checklist(intent=key, lang=lang) + if checklist: + return f"{base}\n{checklist}".strip() + except Exception: + pass + return base def maybe_ops_short_intent_system_hint(*, text: str, lang: str = "en") -> str: diff --git a/runtime/application/gateway/whatsapp_progress.py b/runtime/application/gateway/whatsapp_progress.py index 6e31c0e1..9d2a9da6 100644 --- a/runtime/application/gateway/whatsapp_progress.py +++ b/runtime/application/gateway/whatsapp_progress.py @@ -1,4 +1,8 @@ -"""Rate-limited WhatsApp interim progress during long channel turns.""" +"""Rate-limited WhatsApp interim progress during long channel turns. + +Field groups hate spam: keep at most a few ticks, never alternate +"Running CLI" ↔ "composing", and announce each long tool at most once. +""" from __future__ import annotations @@ -52,11 +56,21 @@ def whatsapp_turn_progress_enabled() -> bool: def progress_min_interval_sec() -> float: + """Default 45s — field turns often run multi-hop CLI for minutes.""" raw = str(os.environ.get("OCLAW_WHATSAPP_PROGRESS_MIN_INTERVAL_SEC") or "").strip() try: - return max(3.0, min(float(raw), 120.0)) + return max(5.0, min(float(raw), 300.0)) except Exception: - return 12.0 + return 45.0 + + +def progress_max_per_turn() -> int: + """Hard cap on interim WA ticks per inbound turn (final reply is separate).""" + raw = str(os.environ.get("OCLAW_WHATSAPP_PROGRESS_MAX_PER_TURN") or "").strip() + try: + return max(0, min(int(raw), 20)) + except Exception: + return 2 def normalize_tool_key(name: str) -> str: @@ -84,7 +98,7 @@ def humanize_long_tool(*, tool_name: str, lang: str = "en") -> str | None: def should_forward_progress_text(text: str) -> bool: - """Filter noisy think/retry ticks; keep meaningful wait signals.""" + """Filter noisy think/retry/composing ticks; keep rare meaningful waits only.""" t = str(text or "").strip() if not t: return False @@ -95,13 +109,12 @@ def should_forward_progress_text(text: str) -> bool: return False if "retry-empty" in low or "retry-native-tool-calls" in low: return False - m = _TOOLS_DONE_RE.search(t) - if m: - try: - return int(m.group(1)) >= 8000 - except Exception: - return False - # Specialist / other explicit progress lines + if "idle-guard" in low: + return False + # Mid-turn "tools done / composing" is the main WA spam pattern when CLI loops. + if _TOOLS_DONE_RE.search(t) or "composing" in low or "整理回复" in t: + return False + # Specialist / other explicit progress lines (rare) if low.startswith("oclaw:"): return True return len(t) >= 8 @@ -110,11 +123,11 @@ def should_forward_progress_text(text: str) -> bool: def humanize_progress_text(*, text: str, lang: str = "en") -> str: t = str(text or "").strip() is_zh = str(lang or "").strip().lower().startswith("zh") - m = _TOOLS_DONE_RE.search(t) - if m: + # tools-done is filtered by should_forward; keep a soft fallback if callers bypass. + if _TOOLS_DONE_RE.search(t): if is_zh: - return "工具已完成,正在整理回复…" - return "Tools finished; composing the reply…" + return "仍在处理,请稍候…" + return "Still working, please wait…" if t.lower().startswith("oclaw:"): body = t.split(":", 1)[-1].strip() if is_zh: @@ -150,6 +163,7 @@ class WhatsappTurnProgressPublisher: is_group: bool = False, inbound: Any = None, min_interval_sec: float | None = None, + max_per_turn: int | None = None, enabled: bool | None = None, clock: Callable[[], float] | None = None, ) -> None: @@ -158,12 +172,14 @@ class WhatsappTurnProgressPublisher: self._is_group = bool(is_group) self._inbound = inbound self._min_interval = float(min_interval_sec if min_interval_sec is not None else progress_min_interval_sec()) + self._max_per_turn = int(max_per_turn if max_per_turn is not None else progress_max_per_turn()) self._enabled = bool(whatsapp_turn_progress_enabled() if enabled is None else enabled) self._clock = clock or time.monotonic self._lock = threading.Lock() self._last_sent_at = 0.0 self._last_text = "" self._sent_count = 0 + self._announced_tools: set[str] = set() @property def sent_count(self) -> int: @@ -183,10 +199,18 @@ class WhatsappTurnProgressPublisher: if str(event or "").strip() != "tool_use_call": return pl = payload if isinstance(payload, dict) else {} - label = humanize_long_tool(tool_name=str(pl.get("tool_name") or ""), lang=self._lang) + tool_name = str(pl.get("tool_name") or "") + key = normalize_tool_key(tool_name) + label = humanize_long_tool(tool_name=tool_name, lang=self._lang) if not label: return - self._maybe_send(label) + with self._lock: + # Same long tool (e.g. repeated execManagedNe hops) → one WA tick only. + if key and key in self._announced_tools: + return + if self._maybe_send(label) and key: + with self._lock: + self._announced_tools.add(key) def _reply_metadata(self) -> dict[str, Any] | None: if not self._is_group or self._inbound is None: @@ -196,23 +220,26 @@ class WhatsappTurnProgressPublisher: except Exception: return None - def _maybe_send(self, text: str) -> None: + def _maybe_send(self, text: str) -> bool: msg = str(text or "").strip() if not msg: - return + return False with self._lock: + if self._max_per_turn >= 0 and self._sent_count >= self._max_per_turn: + return False now = float(self._clock()) if msg == self._last_text and self._sent_count > 0: - return + return False if self._sent_count > 0 and (now - self._last_sent_at) < self._min_interval: - return + return False try: self._enqueue(msg, self._reply_metadata()) except Exception: - return + return False self._last_sent_at = now self._last_text = msg self._sent_count += 1 + return True __all__ = [ @@ -221,6 +248,7 @@ __all__ = [ "humanize_long_tool", "humanize_progress_text", "normalize_tool_key", + "progress_max_per_turn", "progress_min_interval_sec", "should_forward_progress_text", "whatsapp_turn_progress_enabled", diff --git a/runtime/chat/tool_runtime.py b/runtime/chat/tool_runtime.py index 9faf47f8..c0a5d818 100644 --- a/runtime/chat/tool_runtime.py +++ b/runtime/chat/tool_runtime.py @@ -829,10 +829,18 @@ class ToolExecutor: ok, v_err = validate_tool_arguments(tool.parameters, tool_args) if not ok: + intent = None + try: + md = ctx.inbound_metadata if isinstance(ctx.inbound_metadata, dict) else {} + intent = str(md.get("ops_short_intent") or "").strip() or None + except Exception: + intent = None return format_invalid_arguments_error( tool.parameters or {}, str(v_err or "invalid"), lang=str(ctx.lang or "zh"), + tool_name=str(tc.name or tool.name or ""), + intent=intent, ), int((time.perf_counter() - t0) * 1000) try: diff --git a/runtime/chat/turn_idle_guard.py b/runtime/chat/turn_idle_guard.py new file mode 100644 index 00000000..459b3e54 --- /dev/null +++ b/runtime/chat/turn_idle_guard.py @@ -0,0 +1,209 @@ +"""Turn-local idle / checklist guard to cut narration-only and no-progress loops.""" + +from __future__ import annotations + +import os +from dataclasses import dataclass, field +from typing import Any, Literal + + +IdleAction = Literal["continue", "nudge", "early_finalize"] + + +def _env_int(name: str, default: int, *, min_v: int = 1, max_v: int = 20) -> int: + raw = str(os.getenv(name) or "").strip() + if not raw: + return default + try: + n = int(raw) + except Exception: + return default + return max(min_v, min(int(n), max_v)) + + +def idle_guard_enabled() -> bool: + raw = str(os.getenv("AIA_TURN_IDLE_GUARD") or "1").strip().lower() + return raw not in ("0", "false", "no", "off") + + +@dataclass +class RoundStats: + round_idx: int + had_tool_calls: bool + ok_count: int = 0 + fail_count: int = 0 + schema_fail_count: int = 0 + retry_guard_count: int = 0 + tool_names: list[str] = field(default_factory=list) + + +@dataclass +class IdleGuardDecision: + action: IdleAction + reason: str + nudge_text: str = "" + + +@dataclass +class TurnIdleTracker: + """Tracks per-turn progress and decides nudge / early finalize.""" + + lang: str = "en" + short_intent: str | None = None + max_idle_rounds: int = field(default_factory=lambda: _env_int("AIA_TURN_IDLE_MAX_ROUNDS", 2)) + max_schema_fails: int = field(default_factory=lambda: _env_int("AIA_TURN_IDLE_MAX_SCHEMA_FAILS", 3)) + rounds: list[RoundStats] = field(default_factory=list) + nudged: bool = False + total_ok: int = 0 + total_schema_fails: int = 0 + + def record_from_traces( + self, + *, + round_idx: int, + had_tool_calls: bool, + round_traces: list[dict[str, Any]], + results_by_id: dict[str, tuple[dict[str, Any], int]] | None = None, + ) -> RoundStats: + ok = 0 + fail = 0 + schema = 0 + guard = 0 + names: list[str] = [] + # Prefer live results when available (richer failure_class). + payloads: list[dict[str, Any]] = [] + if results_by_id: + for _cid, (payload, _dur) in results_by_id.items(): + if isinstance(payload, dict): + payloads.append(payload) + for tr in round_traces or []: + name = str((tr or {}).get("name") or "").strip() + if name: + names.append(name) + if not payloads: + if (tr or {}).get("ok") is True: + ok += 1 + elif (tr or {}).get("ok") is False: + fail += 1 + for payload in payloads: + if payload.get("ok") is True: + ok += 1 + continue + fail += 1 + klass = str(payload.get("failure_class") or payload.get("error_class") or "").lower() + code = str(payload.get("error_code") or "").lower() + if klass == "schema_validation" or code in {"tool_invalid_arguments", "invalid_arguments"}: + schema += 1 + if klass == "retry_guard" or code in { + "identical_retry_blocked", + "retry_forbidden_blocked", + "tool_loop_guard", + }: + guard += 1 + stats = RoundStats( + round_idx=int(round_idx), + had_tool_calls=bool(had_tool_calls), + ok_count=int(ok), + fail_count=int(fail), + schema_fail_count=int(schema), + retry_guard_count=int(guard), + tool_names=names, + ) + self.rounds.append(stats) + self.total_ok += int(ok) + self.total_schema_fails += int(schema) + return stats + + def decide_after_assistant_no_tools(self) -> IdleGuardDecision: + """Model returned text without tool_calls.""" + if not idle_guard_enabled(): + return IdleGuardDecision(action="continue", reason="disabled") + # Short-intent recipes require an evidence tool before answering. + intent = str(self.short_intent or "").strip() + if intent and intent != "continue" and self.total_ok == 0 and not self.nudged: + self.nudged = True + return IdleGuardDecision( + action="nudge", + reason="short_intent_narration_only", + nudge_text=self._nudge_text(reason="narration_only"), + ) + return IdleGuardDecision(action="continue", reason="allow_text_reply") + + def decide_after_tools(self, stats: RoundStats) -> IdleGuardDecision: + if not idle_guard_enabled(): + return IdleGuardDecision(action="continue", reason="disabled") + + if self.total_schema_fails >= int(self.max_schema_fails) and self.total_ok == 0: + return IdleGuardDecision( + action="early_finalize", + reason="schema_fail_budget", + nudge_text=self._nudge_text(reason="schema_fail_budget"), + ) + + # Count trailing idle rounds: tools ran but zero successes. + idle_streak = 0 + for r in reversed(self.rounds): + if r.had_tool_calls and r.ok_count == 0: + idle_streak += 1 + continue + break + + if idle_streak >= int(self.max_idle_rounds) and self.total_ok == 0: + return IdleGuardDecision( + action="early_finalize", + reason="idle_no_progress", + nudge_text=self._nudge_text(reason="idle_no_progress"), + ) + + # One soft nudge when a round is all retry_guard / schema with no success yet. + if ( + stats.had_tool_calls + and stats.ok_count == 0 + and (stats.schema_fail_count + stats.retry_guard_count) > 0 + and self.total_ok == 0 + and not self.nudged + ): + self.nudged = True + return IdleGuardDecision( + action="nudge", + reason="round_no_progress", + nudge_text=self._nudge_text(reason="round_no_progress"), + ) + + return IdleGuardDecision(action="continue", reason="progress_or_recoverable") + + def _nudge_text(self, *, reason: str) -> str: + is_zh = str(self.lang or "").strip().lower().startswith("zh") + intent = str(self.short_intent or "").strip() + from runtime.tools.playbook_contracts import short_intent_first_step + + step = short_intent_first_step(intent) + if is_zh: + lines = ["[Idle guard] 本轮未取得有效工具证据,禁止空转。"] + if step: + tool, example = step + lines.append(f"立即调用:{tool} 参数示例 {example}") + elif reason == "schema_fail_budget": + lines.append("参数多次不合法:按 tool 返回的 example/playbook 修正后只重试一次,然后直接作答。") + else: + lines.append("改参数或换 fallback 工具;若仍无证据则基于已知信息直接答复用户。") + return "\n".join(lines) + lines = ["[Idle guard] No usable tool evidence yet — stop spinning."] + if step: + tool, example = step + lines.append(f"Call now: {tool} with example args {example}") + elif reason == "schema_fail_budget": + lines.append( + "Repeated invalid arguments: fix once using the tool example/playbook, then answer." + ) + else: + lines.append("Change args or switch fallback tools; if still blocked, answer from what you have.") + return "\n".join(lines) + + +__all__ = [ + "IdleGuardDecision", + "RoundStats", + "TurnIdleTracker", + "idle_guard_enabled", +] diff --git a/runtime/direct_loop.py b/runtime/direct_loop.py index b1ea9697..192f09f6 100644 --- a/runtime/direct_loop.py +++ b/runtime/direct_loop.py @@ -1255,8 +1255,48 @@ def run_oclaw_direct_loop( user_facing_hints: list[str] = [] final_text = "" hit_tool_round_limit = False + early_finalize_reason = "" workspace_lane_role = str(skill_binding_role or wire_policy_role or "generalist").strip().lower() or "generalist" + # Turn checklist / idle guard (P0: cut narration-only + no-progress loops). + short_intent = "" + try: + md0 = inbound_metadata if isinstance(inbound_metadata, dict) else {} + short_intent = str(md0.get("ops_short_intent") or "").strip() + except Exception: + short_intent = "" + if not short_intent: + try: + from runtime.application.gateway.ops_short_intent import detect_ops_short_intent + + short_intent = str(detect_ops_short_intent(str(user_text or "")) or "").strip() + except Exception: + short_intent = "" + turn_system_suffix = "" + if short_intent: + try: + from runtime.tools.playbook_contracts import build_turn_checklist + + turn_system_suffix = build_turn_checklist( + intent=short_intent, + lang=lang, + goal=str(user_text or "")[:160], + ) + except Exception: + turn_system_suffix = "" + from runtime.chat.turn_idle_guard import TurnIdleTracker + + idle_tracker = TurnIdleTracker(lang=str(lang or "en"), short_intent=short_intent or None) + + def _effective_system_prompt() -> str: + base = str(system_prompt or "") + extra = str(turn_system_suffix or "").strip() + if not extra: + return base + if extra in base: + return base + return f"{base}\n\n{extra}".strip() + base_url = str(getattr(model, "base_url", "") or "") model_id = str(getattr(model, "model", "") or "") allow_dsml_text_tools = dsml_text_tools_enabled(base_url=base_url, model_id=model_id) @@ -1275,7 +1315,7 @@ def run_oclaw_direct_loop( store=store, session_id=session_id, max_messages=max_messages, - system_prompt=system_prompt, + system_prompt=_effective_system_prompt(), model=model, lang=lang, memory_context=memory_context, @@ -1373,6 +1413,36 @@ def run_oclaw_direct_loop( ) final_text = step.assistant_text if not step.llm_tool_calls: + idle_decision = idle_tracker.decide_after_assistant_no_tools() + if idle_decision.action == "nudge" and idle_decision.nudge_text: + turn_system_suffix = ( + f"{turn_system_suffix}\n\n{idle_decision.nudge_text}".strip() + if turn_system_suffix + else idle_decision.nudge_text + ) + try: + if trace_id: + _emit_direct_loop_trace( + store=store, + session_id=session_id, + trace_id=trace_id, + parent_span_id=parent_span_id, + event_type="turn_idle_guard", + payload={ + "action": idle_decision.action, + "reason": idle_decision.reason, + "round": int(round_idx + 1), + "short_intent": short_intent or "", + }, + run_id=run_id, + attempt_no=attempt_no, + lang=lang, + ) + except Exception: + pass + if on_progress: + on_progress("oclaw: idle-guard nudge…") + continue break if round_idx == (max_rounds - 1): # Reached tool-round cap with pending tool calls. Execute this batch, then @@ -1420,17 +1490,18 @@ def run_oclaw_direct_loop( inbound_metadata=inbound_metadata, ) + round_traces: list[dict[str, Any]] = [] for tc in step.llm_tool_calls: result, dur = results_by_id.get(str(getattr(tc, "id", "") or ""), ({}, 0)) - tool_traces.append( - { - "name": str(getattr(tc, "name", "") or ""), - "tool_call_id": str(getattr(tc, "id", "") or ""), - "ok": bool((result or {}).get("ok")) if isinstance(result, dict) else None, - "duration_ms": int(dur), - "round": int(round_idx + 1), - } - ) + tr = { + "name": str(getattr(tc, "name", "") or ""), + "tool_call_id": str(getattr(tc, "id", "") or ""), + "ok": bool((result or {}).get("ok")) if isinstance(result, dict) else None, + "duration_ms": int(dur), + "round": int(round_idx + 1), + } + round_traces.append(tr) + tool_traces.append(tr) if isinstance(result, dict): uh = str(result.get("user_facing_hint") or "").strip() if uh: @@ -1439,7 +1510,80 @@ def run_oclaw_direct_loop( if on_progress: on_progress(f"oclaw: tools done ({elapsed_ms}ms)") - need_finalize = hit_tool_round_limit or (bool(tool_traces) and not str(final_text or "").strip()) + idle_stats = idle_tracker.record_from_traces( + round_idx=int(round_idx + 1), + had_tool_calls=True, + round_traces=round_traces, + results_by_id=results_by_id, + ) + idle_decision = idle_tracker.decide_after_tools(idle_stats) + if idle_decision.action == "nudge" and idle_decision.nudge_text: + turn_system_suffix = ( + f"{turn_system_suffix}\n\n{idle_decision.nudge_text}".strip() + if turn_system_suffix + else idle_decision.nudge_text + ) + try: + if trace_id: + _emit_direct_loop_trace( + store=store, + session_id=session_id, + trace_id=trace_id, + parent_span_id=parent_span_id, + event_type="turn_idle_guard", + payload={ + "action": idle_decision.action, + "reason": idle_decision.reason, + "round": int(round_idx + 1), + "short_intent": short_intent or "", + "ok_count": int(idle_stats.ok_count), + "schema_fail_count": int(idle_stats.schema_fail_count), + }, + run_id=run_id, + attempt_no=attempt_no, + lang=lang, + ) + except Exception: + pass + if on_progress: + on_progress("oclaw: idle-guard nudge…") + elif idle_decision.action == "early_finalize": + early_finalize_reason = str(idle_decision.reason or "idle") + if idle_decision.nudge_text: + turn_system_suffix = ( + f"{turn_system_suffix}\n\n{idle_decision.nudge_text}".strip() + if turn_system_suffix + else idle_decision.nudge_text + ) + try: + if trace_id: + _emit_direct_loop_trace( + store=store, + session_id=session_id, + trace_id=trace_id, + parent_span_id=parent_span_id, + event_type="turn_idle_guard", + payload={ + "action": idle_decision.action, + "reason": idle_decision.reason, + "round": int(round_idx + 1), + "short_intent": short_intent or "", + }, + run_id=run_id, + attempt_no=attempt_no, + lang=lang, + ) + except Exception: + pass + if on_progress: + on_progress("oclaw: idle-guard finalize…") + break + + need_finalize = ( + hit_tool_round_limit + or bool(early_finalize_reason) + or (bool(tool_traces) and not str(final_text or "").strip()) + ) if need_finalize: _check_stop(should_stop) if on_progress: @@ -1448,10 +1592,10 @@ def run_oclaw_direct_loop( finalize_suffix = build_finalize_system_suffix( lang=lang, - hit_tool_round_limit=hit_tool_round_limit, + hit_tool_round_limit=hit_tool_round_limit or bool(early_finalize_reason), user_facing_hints=user_facing_hints, ) - finalize_system = str(system_prompt or "") + finalize_system = _effective_system_prompt() if finalize_suffix: finalize_system = f"{finalize_system}\n\n{finalize_suffix}".strip() msgs = _build_model_context( diff --git a/runtime/gateway.py b/runtime/gateway.py index da1c8ce6..049825d3 100644 --- a/runtime/gateway.py +++ b/runtime/gateway.py @@ -870,6 +870,14 @@ class OclawGateway: short_intent = detect_ops_short_intent( str(msg.text or md_intent.get("raw_inbound_text") or "") ) + if short_intent: + # Stamp for tool validation playbook examples + turn idle guard. + try: + if not isinstance(msg.metadata, dict): + msg.metadata = {} + msg.metadata["ops_short_intent"] = str(short_intent) + except Exception: + pass if ops_short_intent_should_filter_tools(short_intent): before_n = len(tools.list()) filtered_specs = filter_tool_specs_for_ops_short_intent( diff --git a/runtime/tools/playbook_contracts.py b/runtime/tools/playbook_contracts.py new file mode 100644 index 00000000..15528680 --- /dev/null +++ b/runtime/tools/playbook_contracts.py @@ -0,0 +1,229 @@ +"""Canonical ops playbook tool contracts. + +Keeps skill recipes (ops-netx-*-playbook) and runtime JSON schemas aligned by +providing executable examples used in: + +- invalid-argument error payloads (self-correct without blind retry) +- short-intent turn checklists +- schema↔playbook regression tests +""" + +from __future__ import annotations + +from typing import Any + +# Bare MCP tool names (without mcp__netx__) and always-on expert tools. +_PLAYBOOK_EXAMPLES: dict[str, dict[str, dict[str, Any]]] = { + "ume_alarm_xlsx_report": { + "fiber_cut": {"mode": "fiber_cut", "deliverable": True}, + "offline": {"mode": "offline", "deliverable": True}, + "alarm_tally": {"mode": "aggregate_by_host", "severity": "critical", "deliverable": True}, + "excel_export": {"mode": "list", "deliverable": True}, + "license": {"mode": "list", "keyword": "license", "deliverable": True}, + "congestion": {"mode": "list", "keyword": "bandwidth", "deliverable": True}, + "default": {"mode": "list", "deliverable": True}, + }, + "write_xlsx": { + "excel_export": { + "sheets": [{"name": "Sheet1", "headers": ["host_name", "count"], "rows": [["NE-1", 1]]}], + "deliverable": True, + "name": "alarms.xlsx", + }, + "default": { + "sheets": [{"name": "Sheet1", "headers": ["col"], "rows": [["val"]]}], + "deliverable": True, + "name": "report.xlsx", + }, + }, + "aggregateumealarms": { + "alarm_tally": {"severity": "critical", "top_ne": 20}, + "default": {"severity": "critical", "top_ne": 20}, + }, + "queryumealarms": { + "default": {"host_name": "", "page_size": 50}, + }, + "queryumealarmsraw": { + "congestion": {"keyword": "bandwidth", "field_preset": "evidence", "page_size": 50}, + "license": {"keyword": "license", "field_preset": "evidence", "page_size": 50}, + "default": {"keyword": "", "field_preset": "evidence", "page_size": 50}, + }, + "runumediagnostics": { + "default": {}, + }, + "listmanagedne": { + "default": {"keyword": "", "connect_status": "pass"}, + }, + "getmanagedne": { + "default": {"ne_id": ""}, + }, + "execmanagedne": { + "default": { + "ume_ne_ids": ["", ""], + "commands": ["show version"], + "read_timeout_sec": 60, + }, + }, + "listclitargets": { + "default": {"source": "ume", "keyword": ""}, + }, + "findtopologypaths": { + "default": { + "from_ume_ne_id": "", + "to_ume_ne_id": "", + "detail": "summary", + }, + }, +} + +# Short-intent → preferred first tool + example key (playbook recipe). +_SHORT_INTENT_FIRST_STEP: dict[str, tuple[str, str]] = { + "fiber_cut": ("ume_alarm_xlsx_report", "fiber_cut"), + "offline": ("ume_alarm_xlsx_report", "offline"), + "alarm_tally": ("ume_alarm_xlsx_report", "alarm_tally"), + "excel_export": ("ume_alarm_xlsx_report", "excel_export"), + "license": ("ume_alarm_xlsx_report", "license"), + "congestion": ("ume_alarm_xlsx_report", "congestion"), +} + + +def canonical_tool_key(tool_name: str) -> str: + raw = str(tool_name or "").strip() + if not raw: + return "" + if "__" in raw: + raw = raw.rsplit("__", 1)[-1] + # Legacy snake_case netx_* → last segment style already handled by rsplit. + if raw.lower().startswith("netx_"): + raw = raw[5:] + return raw.strip().lower() + + +def playbook_examples_for_tool(tool_name: str) -> dict[str, dict[str, Any]]: + return dict(_PLAYBOOK_EXAMPLES.get(canonical_tool_key(tool_name)) or {}) + + +def playbook_example_for_tool( + tool_name: str, + *, + intent: str | None = None, +) -> dict[str, Any] | None: + examples = playbook_examples_for_tool(tool_name) + if not examples: + return None + key = str(intent or "").strip() + if key and key in examples: + return dict(examples[key]) + if "default" in examples: + return dict(examples["default"]) + # First recipe as fallback. + first = next(iter(examples.values()), None) + return dict(first) if isinstance(first, dict) else None + + +def short_intent_first_step(intent: str | None) -> tuple[str, dict[str, Any]] | None: + key = str(intent or "").strip() + if not key: + return None + pair = _SHORT_INTENT_FIRST_STEP.get(key) + if not pair: + return None + tool, example_key = pair + example = playbook_example_for_tool(tool, intent=example_key) or {} + return tool, example + + +def build_turn_checklist( + *, + intent: str | None = None, + lang: str = "en", + goal: str | None = None, +) -> str: + """Compact in-turn checklist injected into system prompt (ops short intents).""" + step = short_intent_first_step(intent) + is_zh = str(lang or "").strip().lower().startswith("zh") + lines: list[str] = [] + if is_zh: + lines.append("[本轮 checklist — 先工具后结论,禁止只叙述不调用]") + else: + lines.append("[Turn checklist — call tools first; do not narrate-only]") + goal_s = str(goal or "").strip() + if goal_s: + lines.append(f"- goal: {goal_s[:160]}") + if step: + tool, example = step + lines.append(f"- step1: {tool}({_fmt_args(example)})") + if is_zh: + lines.append("- 完成后用 Result/Evidence 短答;勿翻页或开无关 playbook") + else: + lines.append("- then reply with Result/Evidence; no pagination / unrelated playbooks") + elif is_zh: + lines.append("- 需要证据时立刻调用工具;失败时改参数或换 fallback,禁止相同参数盲重试") + else: + lines.append("- If evidence is required, call a tool now; on failure change args or switch tools") + return "\n".join(lines) + + +def enrich_invalid_arguments_with_playbook( + payload: dict[str, Any], + *, + tool_name: str, + intent: str | None = None, +) -> dict[str, Any]: + """Attach playbook-aligned example when schema validation fails.""" + if not isinstance(payload, dict): + return payload + out = dict(payload) + example = playbook_example_for_tool(tool_name, intent=intent) + if example: + # Prefer playbook recipe over generic schema-derived example. + out["example"] = example + out["playbook_example"] = True + prior = str(out.get("hint") or "").strip() + tip = f"Playbook recipe: {_fmt_args(example)}" + if tip not in prior: + out["hint"] = f"{prior} {tip}".strip() if prior else tip + out["tool"] = str(tool_name or "").strip() or out.get("tool") + return out + + +def schema_playbook_mismatches( + tool_name: str, + parameters: dict[str, Any] | None, +) -> list[str]: + """Return human-readable mismatches between playbook examples and JSON schema.""" + from runtime.tools.tool_validation import validate_tool_arguments + + examples = playbook_examples_for_tool(tool_name) + if not examples or not isinstance(parameters, dict) or not parameters: + return [] + issues: list[str] = [] + for label, args in examples.items(): + ok, err = validate_tool_arguments(parameters, dict(args)) + if not ok: + issues.append(f"{canonical_tool_key(tool_name)}/{label}: {err}") + return issues + + +def _fmt_args(args: dict[str, Any]) -> str: + parts: list[str] = [] + for k, v in (args or {}).items(): + if isinstance(v, bool): + parts.append(f"{k}={'true' if v else 'false'}") + elif isinstance(v, (int, float)): + parts.append(f"{k}={v}") + elif isinstance(v, str): + parts.append(f'{k}="{v}"') + else: + parts.append(f"{k}=…") + return ", ".join(parts) + + +__all__ = [ + "build_turn_checklist", + "canonical_tool_key", + "enrich_invalid_arguments_with_playbook", + "playbook_example_for_tool", + "playbook_examples_for_tool", + "schema_playbook_mismatches", + "short_intent_first_step", +] diff --git a/runtime/tools/tool_validation.py b/runtime/tools/tool_validation.py index 57be959c..4ed16771 100644 --- a/runtime/tools/tool_validation.py +++ b/runtime/tools/tool_validation.py @@ -58,6 +58,8 @@ def format_invalid_arguments_error( message: str, *, lang: str = "en", + tool_name: str | None = None, + intent: str | None = None, ) -> dict[str, Any]: """Rich invalid-arg payload so the model can self-correct without blind retries.""" props = parameters.get("properties") if isinstance(parameters.get("properties"), dict) else {} @@ -69,7 +71,7 @@ def format_invalid_arguments_error( else: err = f"Invalid arguments: {message}" hint = "Fix arguments to match the schema example; do not retry with the same payload." - return { + out: dict[str, Any] = { "ok": False, "error_code": "tool_invalid_arguments", "error": err, @@ -78,7 +80,20 @@ def format_invalid_arguments_error( "properties": sorted(str(k) for k in props.keys()), "example": example, "hint": hint, + "failure_class": "schema_validation", } + if tool_name: + try: + from runtime.tools.playbook_contracts import enrich_invalid_arguments_with_playbook + + out = enrich_invalid_arguments_with_playbook( + out, + tool_name=str(tool_name), + intent=intent, + ) + except Exception: + out["tool"] = str(tool_name) + return out def validate_tool_arguments(parameters: dict[str, Any], arguments: dict[str, Any]) -> tuple[bool, str | None]: diff --git a/tests/test_playbook_contracts_and_idle_guard.py b/tests/test_playbook_contracts_and_idle_guard.py new file mode 100644 index 00000000..9f19c2dc --- /dev/null +++ b/tests/test_playbook_contracts_and_idle_guard.py @@ -0,0 +1,130 @@ +"""Playbook contracts ↔ tool schema alignment + idle guard.""" + +from __future__ import annotations + +from runtime.chat.turn_idle_guard import TurnIdleTracker +from runtime.tools.experts.network_ops.ume_alarm_xlsx_report import ume_alarm_xlsx_report_tool +from runtime.tools.playbook_contracts import ( + build_turn_checklist, + enrich_invalid_arguments_with_playbook, + playbook_example_for_tool, + schema_playbook_mismatches, + short_intent_first_step, +) +from runtime.tools.public.write_xlsx_tool import write_xlsx_tool +from runtime.tools.tool_validation import format_invalid_arguments_error + + +def test_ume_alarm_xlsx_playbook_examples_match_schema() -> None: + tool = ume_alarm_xlsx_report_tool() + issues = schema_playbook_mismatches(tool.name, tool.parameters) + assert issues == [], issues + + +def test_write_xlsx_playbook_examples_match_schema() -> None: + tool = write_xlsx_tool() + issues = schema_playbook_mismatches(tool.name, tool.parameters) + assert issues == [], issues + + +def test_short_intent_first_step_fiber_cut() -> None: + step = short_intent_first_step("fiber_cut") + assert step is not None + tool, example = step + assert tool == "ume_alarm_xlsx_report" + assert example.get("mode") == "fiber_cut" + assert example.get("deliverable") is True + + +def test_checklist_contains_step1_tool() -> None: + text = build_turn_checklist(intent="offline", lang="en", goal="offline NE list") + assert "ume_alarm_xlsx_report" in text + assert "mode=\"offline\"" in text or "offline" in text + assert "Turn checklist" in text + + +def test_invalid_args_enriched_with_playbook_example() -> None: + tool = ume_alarm_xlsx_report_tool() + raw = format_invalid_arguments_error( + tool.parameters or {}, + "'bogus' is not one of ['list', 'aggregate_by_host', 'fiber_cut', 'offline']", + lang="en", + tool_name="ume_alarm_xlsx_report", + intent="fiber_cut", + ) + assert raw.get("playbook_example") is True + assert raw.get("example", {}).get("mode") == "fiber_cut" + assert "Playbook recipe" in str(raw.get("hint") or "") + + +def test_enrich_invalid_arguments_mcp_alias() -> None: + base = { + "ok": False, + "error_code": "tool_invalid_arguments", + "example": {"severity": "x"}, + "hint": "fix it", + } + out = enrich_invalid_arguments_with_playbook( + base, + tool_name="mcp__netx__aggregateUmeAlarms", + intent="alarm_tally", + ) + assert out["example"].get("top_ne") == 20 + assert out.get("playbook_example") is True + + +def test_idle_guard_nudges_short_intent_narration() -> None: + tracker = TurnIdleTracker(lang="en", short_intent="fiber_cut") + d1 = tracker.decide_after_assistant_no_tools() + assert d1.action == "nudge" + assert "Idle guard" in d1.nudge_text + assert "ume_alarm_xlsx_report" in d1.nudge_text + d2 = tracker.decide_after_assistant_no_tools() + assert d2.action == "continue" + + +def test_idle_guard_early_finalize_on_idle_streak() -> None: + tracker = TurnIdleTracker(lang="en", short_intent="alarm_tally", max_idle_rounds=2) + # Round 1: all failures + s1 = tracker.record_from_traces( + round_idx=1, + had_tool_calls=True, + round_traces=[{"name": "ume_alarm_xlsx_report", "ok": False}], + results_by_id={ + "1": ( + { + "ok": False, + "error_code": "tool_invalid_arguments", + "failure_class": "schema_validation", + }, + 10, + ) + }, + ) + d1 = tracker.decide_after_tools(s1) + assert d1.action in {"nudge", "continue", "early_finalize"} + # Round 2: still no success + s2 = tracker.record_from_traces( + round_idx=2, + had_tool_calls=True, + round_traces=[{"name": "ume_alarm_xlsx_report", "ok": False}], + results_by_id={ + "2": ( + { + "ok": False, + "error_code": "identical_retry_blocked", + "failure_class": "retry_guard", + }, + 5, + ) + }, + ) + d2 = tracker.decide_after_tools(s2) + assert d2.action == "early_finalize" + assert d2.reason in {"idle_no_progress", "schema_fail_budget"} + + +def test_playbook_example_for_mcp_namespaced_tool() -> None: + ex = playbook_example_for_tool("mcp__netx__execManagedNe") + assert ex is not None + assert "commands" in ex diff --git a/tests/test_whatsapp_turn_progress.py b/tests/test_whatsapp_turn_progress.py index a16a5762..4f9237aa 100644 --- a/tests/test_whatsapp_turn_progress.py +++ b/tests/test_whatsapp_turn_progress.py @@ -38,12 +38,14 @@ class WhatsappProgressHelpersTests(unittest.TestCase): self.assertFalse(should_forward_progress_text("oclaw: running…")) self.assertFalse(should_forward_progress_text("oclaw: think (1)…")) self.assertFalse(should_forward_progress_text("oclaw: tools done (120ms)")) - self.assertTrue(should_forward_progress_text("oclaw: tools done (12000ms)")) + # Mid-turn composing spam must not reach WhatsApp (was the main noise pattern). + self.assertFalse(should_forward_progress_text("oclaw: tools done (12000ms)")) + self.assertFalse(should_forward_progress_text("oclaw: idle-guard nudge…")) self.assertTrue(should_forward_progress_text("oclaw: image specialist (legacy multimodal HTTP)…")) def test_humanize_defaults_to_english(self) -> None: self.assertIn("CLI", str(humanize_long_tool(tool_name="mcp__netx__execManagedNe") or "")) - self.assertIn("composing", humanize_progress_text(text="oclaw: tools done (9000ms)").lower()) + self.assertIn("wait", humanize_progress_text(text="oclaw: tools done (9000ms)").lower()) class WhatsappTurnProgressPublisherTests(unittest.TestCase): @@ -59,6 +61,7 @@ class WhatsappTurnProgressPublisherTests(unittest.TestCase): lang="en", is_group=False, min_interval_sec=10.0, + max_per_turn=5, enabled=True, clock=lambda: clock["t"], ) @@ -71,18 +74,43 @@ class WhatsappTurnProgressPublisherTests(unittest.TestCase): clock["t"] = 105.0 pub.on_tool_ui("tool_use_call", {"tool_name": "mcp__netx__queryUmeAlarms"}) - self.assertEqual(len(sent), 1) # throttled + self.assertEqual(len(sent), 1) # throttled by interval clock["t"] = 111.0 pub.on_tool_ui("tool_use_call", {"tool_name": "mcp__netx__queryUmeAlarms"}) self.assertEqual(len(sent), 2) + # tools-done / composing never forwarded pub.on_progress("oclaw: tools done (15000ms)") self.assertEqual(len(sent), 2) - clock["t"] = 122.0 + clock["t"] = 160.0 pub.on_progress("oclaw: tools done (15000ms)") - self.assertEqual(len(sent), 3) - self.assertIn("composing", sent[2][0].lower()) + self.assertEqual(len(sent), 2) + + def test_same_tool_announced_once_and_max_cap(self) -> None: + sent: list[str] = [] + clock = {"t": 0.0} + pub = WhatsappTurnProgressPublisher( + enqueue=lambda t, _m: sent.append(t), + lang="en", + min_interval_sec=1.0, + max_per_turn=2, + enabled=True, + clock=lambda: clock["t"], + ) + pub.on_tool_ui("tool_use_call", {"tool_name": "mcp__netx__listCliTargets"}) + clock["t"] = 10.0 + # Repeated execManagedNe hops must not spam. + for _ in range(6): + clock["t"] += 60.0 + pub.on_tool_ui("tool_use_call", {"tool_name": "mcp__netx__execManagedNe"}) + self.assertEqual(len(sent), 2) + self.assertIn("CLI targets", sent[0]) + self.assertIn("CLI", sent[1]) + # Cap reached: another distinct tool is dropped. + clock["t"] += 60.0 + pub.on_tool_ui("tool_use_call", {"tool_name": "mcp__netx__findTopologyPaths"}) + self.assertEqual(len(sent), 2) def test_group_metadata_mentions_without_quote(self) -> None: sent: list[tuple[str, dict[str, Any] | None]] = []