From 8787adc3cce07f92eb743ec3e4bf63dba3ec72be Mon Sep 17 00:00:00 2001 From: oliver Date: Sat, 30 May 2026 22:53:49 +0800 Subject: [PATCH] fix(chat): repair async worker and video specialist routing Fix trace_id/run_id initialization so async turns persist messages. Pass skill_binding_role through the worker and route video/image experts synchronously for the DashScope legacy lane. Co-authored-by: Cursor --- runtime/router.py | 10 ++++++++++ runtime/worker.py | 30 +++++++++++++++++++++++++----- tests/test_oclaw_router.py | 18 ++++++++++++++++++ 3 files changed, 53 insertions(+), 5 deletions(-) diff --git a/runtime/router.py b/runtime/router.py index 285e1a69..e79d278a 100644 --- a/runtime/router.py +++ b/runtime/router.py @@ -111,6 +111,16 @@ def decide_route(msg: StandardMessage, *, store: Any | None = None, model: Any | skill_count = int(md.get("skills_total") or 0) except Exception: skill_count = 0 + # Image/video legacy lanes run synchronously in-process (DashScope HTTP + poll). + # Do not queue them as async_task for long prompts with attachments. + if requested_specialist in ("video", "image"): + return RouterDecision( + mode="sync_direct", + reason=f"{requested_specialist}_expert_legacy_lane", + skill_signal=f"skills={int(skill_count)}", + interaction_mode=interaction_mode, + requested_specialist=requested_specialist, + ) mode = _router_mode_from_store(store) if mode == "llm_json": d = _decide_llm_json(msg, model=model) diff --git a/runtime/worker.py b/runtime/worker.py index 287c3c01..9614cfae 100644 --- a/runtime/worker.py +++ b/runtime/worker.py @@ -10,7 +10,7 @@ from runtime.agents.factory import build_gateway_executor from runtime.agent_core_run import AgentCoreRunInput, run_agent_core from runtime.memory_stage import build_memory_context from runtime.relay_pointer import build_acp_relay_result, validate_relay_share_envelope -from runtime.types import StandardMessage +from runtime.types import StandardMessage, normalize_interaction_mode, normalize_requested_specialist from runtime.chat.model_path_audit import ensure_no_tool_or_embedded_image_payload from runtime.session_auto_title import ( AUTO_TITLE_SYSTEM_PROMPT_EN, @@ -171,8 +171,8 @@ def _worker_loop(*, store: Any, worker_id: str, poll_interval_s: float) -> None: except Exception: payload = {} - trace_id = str(payload.get("trace_id") or "") - run_id = str(payload.get("run_id") or "").strip() or None + trace_id = str(payload.get("trace_id") or "") + run_id = str(payload.get("run_id") or "").strip() or None session_id = str(task.session_id or "") try: if trace_id: @@ -239,7 +239,24 @@ def _worker_loop(*, store: Any, worker_id: str, poll_interval_s: float) -> None: user_id = str(payload.get("user_id") or "") viewer_username = str(payload.get("viewer_username") or "") model_profile_id = str(payload.get("model_profile_id") or "") or None - selected_specialist = str(payload.get("selected_specialist") or "") or str(metadata.get("selected_specialist") or "") + interaction_mode = normalize_interaction_mode( + str(payload.get("interaction_mode") or metadata.get("interaction_mode") or "") + ) + requested_specialist = normalize_requested_specialist( + str(payload.get("requested_specialist") or metadata.get("selected_specialist") or "") + ) + manager_specialist = normalize_requested_specialist( + str( + payload.get("manager_selected_specialist") + or payload.get("selected_specialist") + or requested_specialist + or "" + ) + ) + skill_binding_role = str(manager_specialist or "generalist") + wire_policy_role = ( + "manager" if interaction_mode == "comprehensive" else str(requested_specialist or skill_binding_role) + ) memory_ctx = build_memory_context( store=store, session_id=session_id, @@ -250,7 +267,7 @@ def _worker_loop(*, store: Any, worker_id: str, poll_interval_s: float) -> None: executor = build_gateway_executor( store, lang=lang, - specialist=selected_specialist or "generalist", + specialist=skill_binding_role, profile_id=model_profile_id, viewer_user_id=user_id or None, viewer_username=viewer_username or None, @@ -293,6 +310,7 @@ def _worker_loop(*, store: Any, worker_id: str, poll_interval_s: float) -> None: store=store, data=AgentCoreRunInput( msg=msg, + persisted_user_text=str(user_text or ""), lang=lang, system_prompt=system_prompt, model=executor.model, @@ -307,6 +325,8 @@ def _worker_loop(*, store: Any, worker_id: str, poll_interval_s: float) -> None: memory_context=memory_ctx, oclaw_task_id=str(task.id), oclaw_worker_id=worker_id, + skill_binding_role=skill_binding_role, + wire_policy_role=wire_policy_role, ), ) base_result = { diff --git a/tests/test_oclaw_router.py b/tests/test_oclaw_router.py index 2a958d98..bc7216a9 100644 --- a/tests/test_oclaw_router.py +++ b/tests/test_oclaw_router.py @@ -114,6 +114,24 @@ def test_requested_specialist_normalization_defaults_to_generalist() -> None: assert normalize_requested_specialist("unknown") == "generalist" +def test_router_video_expert_sync_despite_long_attachments() -> None: + long_text = "x" * 150 + msg = StandardMessage( + session_id="s1", + tenant_id="t1", + user_id="u1", + role="user", + channel="admin_chat", + text=long_text, + attachments=[{"type": "image_ref", "attachment_id": "a" * 64}], + metadata={"interaction_mode": "expert", "selected_specialist": "video"}, + ) + d = decide_route(msg) + assert d.mode == "sync_direct" + assert d.reason == "video_expert_legacy_lane" + assert d.requested_specialist == "video" + + def test_router_carries_interaction_mode_and_requested_specialist() -> None: msg = StandardMessage( session_id="s1",