diff --git a/interfaces/admin/static/chat.js b/interfaces/admin/static/chat.js index d49c6528..4a328673 100644 --- a/interfaces/admin/static/chat.js +++ b/interfaces/admin/static/chat.js @@ -2649,10 +2649,6 @@ async function renderChatUi() { const payload = msg.payload || {}; onEvent({ event: String(msg.event || ""), payload }); const d = msg.event === "chat" ? { phase: String(payload.state || "") } : {}; - if (msg.event === "session.message") { - const role = String(((payload.message || {}).role) || "").toLowerCase(); - if (role === "assistant") return doneMeta; - } const phase = String(d.phase || ""); if (phase === "final" || phase === "error" || phase === "aborted") { return doneMeta; @@ -2783,7 +2779,13 @@ async function renderChatUi() { let last = null; for (let i = msgs.length - 1; i >= 0; i--) { const m = msgs[i]; - if (String((m && m.role) || "").toLowerCase() === "assistant") { + if (String((m && m.role) || "").toLowerCase() !== "assistant") continue; + // Recovery should only accept visible assistant body, not intermediate + // reasoning/tool-call events; otherwise we may terminate on a partial line. + const et = String((m && m.event_type) || "").trim().toLowerCase(); + if (et && et !== "assistant_text" && et !== "assistant") continue; + if (!String((m && m.content) || "").trim()) continue; + { last = m; break; } @@ -3055,6 +3057,10 @@ ${autoLimit ? `
auto-added claus const abortController = new AbortController(); currentStreamAbortController = abortController; currentAbortMeta = { sessionId: String(activeId || ""), runId: "" }; + // End-of-turn reload gating: avoid clearing/repainting messagesEl after we already + // finalized the stream bubble in-place (prevents end "flash"). + let turnFinalized = false; + let turnStreamedEnough = false; try { const transport = new OclawWsChatTransport({ tokenProvider: () => token, @@ -3064,6 +3070,7 @@ ${autoLimit ? `
auto-added claus // Immediate assistant placeholder bubble with dynamic status. ensureStreamBubble(); _startDynamicStreamStatus(); + let wsLastActivityAt = Date.now(); const wsSendPromise = transport.sendSessionSend({ sessionId: activeId, text: userText, @@ -3078,6 +3085,7 @@ ${autoLimit ? `
auto-added claus memoryMode: String(localStorage.getItem(CHAT_MEMORY_MODE_KEY) || "default"), signal: abortController.signal, onEvent: async (frame) => { + wsLastActivityAt = Date.now(); const eventName = String((frame && frame.event) || ""); const payload = (frame && frame.payload) || {}; if (eventName === "agent.event") { @@ -3164,6 +3172,8 @@ ${autoLimit ? `
auto-added claus streamTextBuffer = ""; chatRunId = null; const streamedEnough = hasRealStreamText || chatStreamSegments.length > 0 || sawWsChatEvent; + turnFinalized = true; + turnStreamedEnough = streamedEnough; chatStreamSegments = []; // Avoid end-of-turn flash: only reload history when stream had no usable content. if (!streamedEnough) { @@ -3180,6 +3190,8 @@ ${autoLimit ? `
auto-added claus streamDisplayShown = streamDisplayTarget; const ok = await appendFinalAssistant(payload.message, chatStream); if (!ok) _markStreamTerminal("error", "aborted"); + turnFinalized = true; + turnStreamedEnough = hasRealStreamText || chatStreamSegments.length > 0 || sawWsChatEvent; chatStream = ""; streamTextBuffer = ""; chatRunId = null; @@ -3190,6 +3202,8 @@ ${autoLimit ? `
auto-added claus sawWsTerminalEvent = true; _stopDynamicStreamStatus(); _markStreamTerminal("error", String(payload.errorMessage || "chat error")); + turnFinalized = true; + turnStreamedEnough = hasRealStreamText || chatStreamSegments.length > 0 || sawWsChatEvent; chatStream = ""; streamTextBuffer = ""; chatRunId = null; @@ -3198,12 +3212,25 @@ ${autoLimit ? `
auto-added claus } }, }); - doneMeta = await Promise.race([ - wsSendPromise, - new Promise((_, reject) => { - setTimeout(() => reject(new Error(`ws_send_timeout:${WS_CHAT_SEND_TIMEOUT_MS}`)), WS_CHAT_SEND_TIMEOUT_MS); - }), - ]); + let wsWatchdog = 0; + const wsInactivityTimeoutPromise = new Promise((_, reject) => { + wsWatchdog = setInterval(() => { + if (Date.now() - wsLastActivityAt > WS_CHAT_SEND_TIMEOUT_MS) { + if (wsWatchdog) { + clearInterval(wsWatchdog); + wsWatchdog = 0; + } + reject(new Error(`ws_send_timeout:${WS_CHAT_SEND_TIMEOUT_MS}`)); + } + }, 1000); + }); + const wsSendObserved = wsSendPromise.finally(() => { + if (wsWatchdog) { + clearInterval(wsWatchdog); + wsWatchdog = 0; + } + }); + doneMeta = await Promise.race([wsSendObserved, wsInactivityTimeoutPromise]); if (doneMeta && typeof doneMeta === "object") doneMeta.__transport = "ws"; if (doneMeta && typeof doneMeta === "object") { const startToRunning = @@ -3231,11 +3258,17 @@ ${autoLimit ? `
auto-added claus emsg.includes("closed") || emsg.includes("timeout"); if (wsLikeFailure) { - const wsLikelyAlreadyProducedReply = sawWsTerminalEvent || hasRealStreamText || sawWsChatEvent; - if (wsLikelyAlreadyProducedReply) { - setTimeout(() => { - loadMessagesForActive().catch(() => {}); - }, 350); + // Only short-circuit when a terminal event was already received. + // If we only saw partial deltas and WS drops, we must attempt recovery. + if (sawWsTerminalEvent) { + // If the stream bubble was already finalized with usable content, do NOT + // repaint messagesEl (prevents end-of-turn flash). Only reload when we + // have nothing usable and need to recover from persisted history. + if (!turnFinalized || !turnStreamedEnough) { + setTimeout(() => { + loadMessagesForActive().catch(() => {}); + }, 350); + } return doneMeta; } // WS timeout may happen while backend still computes and persists final reply. diff --git a/interfaces/ws/turn_runner.py b/interfaces/ws/turn_runner.py index 40d44b47..b26af3b9 100644 --- a/interfaces/ws/turn_runner.py +++ b/interfaces/ws/turn_runner.py @@ -260,18 +260,102 @@ async def run_agent_turn_via_bridge( on_tool_ui=on_tool_ui, ) + async def _heartbeat_during_turn(stop_evt: asyncio.Event) -> None: + # Keep WS active while model/tool pipeline is still preparing first token. + # Some proxies/load balancers close idle WS in 20-30s without downstream frames. + while not stop_evt.is_set(): + try: + await asyncio.wait_for(stop_evt.wait(), timeout=4.0) + break + except asyncio.TimeoutError: + rid = run_id_holder.get("run_id") or "" + if not rid: + continue + with conn._abort_lock: + if rid in conn._aborted_run_ids: + continue + try: + await conn.emit_agent_event( + run_id=rid, + stream="lifecycle", + data={"phase": "running", "event": "heartbeat"}, + ) + except Exception: + # Best-effort keepalive; never fail the turn on heartbeat errors. + pass + + hb_stop = asyncio.Event() + hb_task: asyncio.Task[Any] | None = asyncio.create_task(_heartbeat_during_turn(hb_stop)) try: result = await asyncio.to_thread(_run_turn_sync) except Exception as exc: + hb_stop.set() + if hb_task is not None: + try: + await hb_task + except Exception: + pass + err_text = str(exc or "agent_failed") + user_facing_error = ( + "本轮执行失败:工具不可用或执行异常。" + f"\n\n错误信息:{err_text}\n\n请重试,或改用其它可用工具。" + ) + fail_msg: dict[str, Any] = { + "role": "assistant", + "content": user_facing_error, + "timestamp": now_ms(), + } + try: + store.add_message( + session_id=str(session_id), + role="assistant", + content=str(user_facing_error), + turn_uuid=str(run_id_holder.get("run_id") or "") or None, + event_type="assistant_text", + ) + except Exception: + pass if send_response: - await conn.send_res(req_id, ok=False, error=error_shape("UNAVAILABLE", str(exc or "agent_failed"))) + await conn.send_res( + req_id, + ok=True, + payload={ + "runId": run_id_holder.get("run_id") or "", + "acceptedAt": now_ms(), + "mode": "sync_direct", + "taskId": "", + "traceId": "", + "reply": user_facing_error, + "selectedSpecialist": str(p.get("specialist") or "generalist"), + "interactionMode": str(p.get("interaction_mode") or "comprehensive"), + "dispatchReason": "execution_failed", + "managerSelectedSpecialist": str(p.get("specialist") or "generalist"), + "requestedSpecialist": str(p.get("specialist") or "generalist"), + "dynamicAgentUsed": False, + "dynamicAgentName": "", + "relayPointerCount": 0, + "relayEnvelopePresent": False, + "relayEnvelopePointerCount": 0, + "relayTtlTurnCount": int(marker_turn_count), + "relayTtlSessionCount": int(marker_session_count), + "relayTtlKeepCount": int(marker_keep_count), + "status": "failed", + "error": err_text, + }, + ) await conn.emit_agent_event( run_id=run_id_holder.get("run_id") or "", stream="lifecycle", # Always emit a terminal lifecycle event even on failure. data={"phase": "end", "status": "error", "error": str(exc or "agent_failed")}, ) - await conn.emit_chat_event(run_id=run_id_holder.get("run_id") or "", state="error", error=str(exc or "agent_failed")) + await conn.emit_chat_event( + run_id=run_id_holder.get("run_id") or "", + state="final", + reply=user_facing_error, + message=fail_msg, + session_key=str(session_id), + ) try: await conn.send_event( "session.marker", @@ -287,8 +371,18 @@ async def run_agent_turn_via_bridge( ) except Exception: pass + try: + await conn.send_event("session.message", {"sessionKey": str(session_id), "message": fail_msg}) + except Exception: + pass return finally: + hb_stop.set() + if hb_task is not None: + try: + await hb_task + except Exception: + pass _rid0 = str(run_id_holder.get("run_id") or "") if _rid0: with conn._abort_lock: @@ -296,6 +390,25 @@ async def run_agent_turn_via_bridge( conn._aborted_run_ids.discard(_rid0) rid = str(getattr(result, "run_id", "") or run_id_holder.get("run_id") or "") + run_status = "success" + run_last_error_code = "" + run_stop_reason = "" + run_error_detail = "" + try: + rr = store.oclaw_run_get(run_id=rid, tenant_id=tenant_id or None) if rid else None + if rr is not None: + run_status = str(getattr(rr, "status", "") or "success").strip().lower() or "success" + payload = getattr(rr, "payload", {}) or {} + if isinstance(payload, dict): + run_last_error_code = str(payload.get("last_error_code") or "").strip() + run_stop_reason = str(payload.get("stop_reason") or "").strip() + if rid and run_status == "failed": + attempts = store.oclaw_attempt_list(run_id=rid, limit=30) + if attempts: + latest = attempts[-1] if isinstance(attempts[-1], dict) else {} + run_error_detail = str((latest or {}).get("reason") or "").strip() + except Exception: + pass ttft_payload: dict[str, Any] | None = None try: ft = first_token_ms_holder.get("ms") @@ -362,6 +475,10 @@ async def run_agent_turn_via_bridge( "relayTtlSessionCount": int(getattr(result, "relay_ttl_session_count", 0) or 0), "relayTtlKeepCount": int(getattr(result, "relay_ttl_keep_count", 0) or 0), "ttft": ttft_payload if isinstance(ttft_payload, dict) else None, + "status": run_status, + "lastErrorCode": run_last_error_code, + "stopReason": run_stop_reason, + "errorDetail": run_error_detail, }, ) await conn.emit_agent_event( @@ -369,7 +486,7 @@ async def run_agent_turn_via_bridge( stream="lifecycle", data={ "phase": "end", - "status": "ok", + "status": "error" if run_status == "failed" else "ok", "reply": str(getattr(result, "reply_text", "") or ""), "elapsedMs": int(getattr(result, "elapsed_ms", 0) or 0), "mode": str(getattr(result, "mode", "sync_direct") or "sync_direct"), @@ -382,26 +499,49 @@ async def run_agent_turn_via_bridge( with buf_lock: final_text = "".join(token_chunks) final_msg: dict[str, Any] = {"role": "assistant", "content": final_text, "timestamp": now_ms()} - try: - persisted = store.get_messages(session_id=session_id, limit=12) - for m in reversed(list(persisted or [])): - if str(getattr(m, "role", "") or "").lower() != "assistant": - continue - content = str(getattr(m, "content", "") or "") - if not content.strip(): - continue - final_text = content - final_msg = { - "id": int(getattr(m, "id", 0) or 0), - "role": "assistant", - "content": content, - "timestamp": str(getattr(m, "timestamp", "") or ""), - "tool_calls": getattr(m, "tool_calls", None), - "attachments": getattr(m, "attachments", None), - } - break - except Exception: - pass + if run_status == "failed" and not str(final_text or "").strip(): + err_code = run_last_error_code or "unknown_error" + stop_reason = run_stop_reason or "failed" + detail_line = f"\n详细原因:{run_error_detail}" if run_error_detail else "" + final_text = ( + "本轮执行失败,已提前结束。" + f"\n\n错误代码:{err_code}" + f"\n停止原因:{stop_reason}" + f"{detail_line}" + "\n\n请重试,或减少输入复杂度后再试。" + ) + final_msg = {"role": "assistant", "content": final_text, "timestamp": now_ms()} + try: + store.add_message( + session_id=str(session_id), + role="assistant", + content=str(final_text), + turn_uuid=str(rid or "") or None, + event_type="assistant_text", + ) + except Exception: + pass + elif run_status != "failed": + try: + persisted = store.get_messages(session_id=session_id, limit=12) + for m in reversed(list(persisted or [])): + if str(getattr(m, "role", "") or "").lower() != "assistant": + continue + content = str(getattr(m, "content", "") or "") + if not content.strip(): + continue + final_text = content + final_msg = { + "id": int(getattr(m, "id", 0) or 0), + "role": "assistant", + "content": content, + "timestamp": str(getattr(m, "timestamp", "") or ""), + "tool_calls": getattr(m, "tool_calls", None), + "attachments": getattr(m, "attachments", None), + } + break + except Exception: + pass await conn.emit_chat_event(run_id=rid, state="final", reply=str(final_text or ""), message=final_msg, session_key=str(session_id)) try: diff --git a/platform/llm/transports/openai_chat_completions.py b/platform/llm/transports/openai_chat_completions.py index f3f85e97..cf5b3e2d 100644 --- a/platform/llm/transports/openai_chat_completions.py +++ b/platform/llm/transports/openai_chat_completions.py @@ -24,6 +24,12 @@ def _model_id_suggests_gemini(model: str | None) -> bool: return "gemini" in (model or "").lower() +def _is_minimax_compat(model: str | None, base_url: str | None) -> bool: + m = str(model or "").strip().lower() + b = str(base_url or "").strip().lower() + return ("minimax" in m) or ("minimax" in b) + + def _find_thought_signature_in_obj(o: Any) -> str | None: if isinstance(o, dict): for k in ("thought_signature", "thoughtSignature"): @@ -126,12 +132,34 @@ class OpenAIChatModel(ChatModel): except Exception as exc: logger.warning("replay_policy apply failed (%s); continuing without rewrite", exc) + minimax_compat = _is_minimax_compat(self.model, self.base_url) use_tools = bool(tools) cleaned_msgs: list[dict[str, Any]] = [] for m in norm_msgs or []: if not isinstance(m, dict): continue - if str(m.get("role") or "") == "tool": + role = str(m.get("role") or "") + if minimax_compat and role == "assistant" and isinstance(m.get("tool_calls"), list): + # MiniMax OpenAI-compat can reject long-history tool_use/tool_result replays. + # Keep text semantics, but remove historical wire-level tool_calls. + mm = dict(m) + mm.pop("tool_calls", None) + cleaned_msgs.append(mm) + continue + if minimax_compat and role == "tool": + # Downgrade history tool_result rows to plain assistant text context, + # avoiding strict tool_result sequence validation on replay. + tcid = str(m.get("tool_call_id") or m.get("call_id") or "").strip() + tname = str(m.get("name") or "").strip() + raw = str(m.get("content") or "") + prefix = "[tool_result replay]" + if tname: + prefix += f" name={tname}" + if tcid: + prefix += f" id={tcid}" + cleaned_msgs.append({"role": "assistant", "content": f"{prefix}\n{raw}".strip()}) + continue + if role == "tool": mm = dict(m) for k in ("tool_call_id", "call_id"): if k in mm: diff --git a/runtime/agent_core_attempt.py b/runtime/agent_core_attempt.py index b9feb771..636d8699 100644 --- a/runtime/agent_core_attempt.py +++ b/runtime/agent_core_attempt.py @@ -28,6 +28,7 @@ ALL_ATTEMPT_ERROR_CODES = ( "control_interrupted", "auth_invalid_credentials", "input_invalid_request", + "tool_replay_protocol_mismatch", "context_overflow", "tool_loop_guard", "tool_execution_failed", @@ -83,6 +84,8 @@ def _classify_attempt_error(exc: Exception) -> tuple[str, str, bool]: return ("control_interrupted", raw[:500], False) if "api_key" in low or "invalid api key" in low or "unauthorized" in low or "401" in low or "forbidden" in low or "403" in low: return ("auth_invalid_credentials", raw[:500], False) + if "invalid tool_result sequence" in low or "unexpected tool_use_id" in low: + return ("tool_replay_protocol_mismatch", raw[:500], True) if "invalid_request" in low or "bad_request" in low or "400" in low: return ("input_invalid_request", raw[:500], False) if "context_length" in low or "token limit" in low or "max context" in low: diff --git a/runtime/agent_core_run.py b/runtime/agent_core_run.py index 6e92be3d..ee2ced58 100644 --- a/runtime/agent_core_run.py +++ b/runtime/agent_core_run.py @@ -61,6 +61,7 @@ DEFAULT_RETRYABLE_ERROR_CODES = ( "provider_rate_limited", "provider_temporary_error", "provider_unavailable", + "tool_replay_protocol_mismatch", "context_overflow", "tool_execution_failed", )