mirror of
https://github.com/hansjone/oclaw.git
synced 2026-10-11 03:00:45 +08:00
feat: plan agent v2, global chat mode, and user-mode v2 flag sync
- Add runtime/plan_agent_v2 package and shims; gateway/direct_loop/WS wiring - Admin chat: interaction mode and specialist only in user menu; session API stores memory_mode and execution_mode only - POST /admin/api/chat/user-mode mirrors plan_agent_version to AIA_EXPERT_PLAN_AGENT_V2_ENABLED (v2 to 1, v1 to 0) - Composer cleanup (hidden mode select, no reasoning toggle in meta bar); tests and _local/system.env.example Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
parent
ae44cbcad5
commit
37522f492a
49 changed files with 4146 additions and 238 deletions
|
|
@ -35,6 +35,8 @@ _DIRECT_LOOP_OC_STAGE: dict[str, str] = {
|
|||
"tool_pairing_guard": "tool_pairing_guard",
|
||||
}
|
||||
_THINK_BLOCK_RE = re.compile(r"<(think|redacted_thinking)>\s*(.*?)\s*</\1>\s*", flags=re.IGNORECASE | re.DOTALL)
|
||||
_DSML_INVOKE_NAME_RE = re.compile(r"invoke\s+name\s*=\s*['\"]([^'\"\s>]+)['\"]", flags=re.IGNORECASE)
|
||||
_JSON_TOOL_NAME_RE = re.compile(r"['\"]name['\"]\s*:\s*['\"]([^'\"\s]{1,120})['\"]", flags=re.IGNORECASE)
|
||||
_TOOL_WIRE_CACHE_LOCK = threading.Lock()
|
||||
_TOOL_WIRE_CACHE: dict[str, tuple[float, list[dict[str, Any]]]] = {}
|
||||
_TOOL_WIRE_CACHE_TTL_SEC = 300.0
|
||||
|
|
@ -54,6 +56,16 @@ def _safe_int(raw: Any, default: int, *, min_value: int = 1, max_value: int = 2_
|
|||
return min(value, max_value)
|
||||
|
||||
|
||||
def _safe_nonneg_int(raw: Any, default: int, *, max_value: int = 2_000_000) -> int:
|
||||
try:
|
||||
value = int(raw)
|
||||
except Exception:
|
||||
return max(0, int(default))
|
||||
if value < 0:
|
||||
return max(0, int(default))
|
||||
return min(value, max_value)
|
||||
|
||||
|
||||
def _oclaw_config_path() -> Path:
|
||||
raw = str(os.getenv("AIA_OCLAW_CONFIG_PATH") or "").strip()
|
||||
if raw:
|
||||
|
|
@ -820,6 +832,163 @@ def _tool_names_for_trace(tools: list[dict[str, Any]]) -> list[str]:
|
|||
return out
|
||||
|
||||
|
||||
def _chat_with_empty_body_retry(
|
||||
*,
|
||||
model: Any,
|
||||
msgs: list[dict[str, Any]],
|
||||
llm_tools: list[dict[str, Any]],
|
||||
on_token: Optional[Callable[[str], None]],
|
||||
on_progress: Optional[Callable[[str], None]],
|
||||
progress_label: str = "oclaw: think",
|
||||
) -> Any:
|
||||
# Empty assistant body can occur transiently at upstream gateways.
|
||||
# Retry until non-empty (bounded by retry count and total timeout).
|
||||
retry_max = _safe_nonneg_int(os.getenv("AIA_EMPTY_ASSISTANT_RETRY_MAX"), 1, max_value=3)
|
||||
retry_delay_ms = _safe_nonneg_int(os.getenv("AIA_EMPTY_ASSISTANT_RETRY_DELAY_MS"), 1200, max_value=15_000)
|
||||
retry_total_timeout_ms = _safe_nonneg_int(os.getenv("AIA_EMPTY_ASSISTANT_RETRY_TOTAL_TIMEOUT_MS"), 30_000, max_value=300_000)
|
||||
started = time.perf_counter()
|
||||
retries_done = 0
|
||||
resp = model.chat(msgs, llm_tools, on_token=on_token)
|
||||
while True:
|
||||
content = str(getattr(resp, "content", "") or "")
|
||||
tool_calls = list(getattr(resp, "tool_calls", []) or [])
|
||||
textual_tool_intent = (not tool_calls) and bool(_extract_textual_tool_intent_names(content))
|
||||
if (content.strip() or tool_calls) and not textual_tool_intent:
|
||||
return resp
|
||||
elapsed_ms = int((time.perf_counter() - started) * 1000.0)
|
||||
if retries_done >= retry_max or elapsed_ms >= retry_total_timeout_ms:
|
||||
return resp
|
||||
if textual_tool_intent:
|
||||
if on_progress:
|
||||
on_progress(f"{progress_label} retry-native-tool-calls ({retries_done + 1}/{retry_max})…")
|
||||
repair_msgs = list(msgs) + [
|
||||
{
|
||||
"role": "system",
|
||||
"content": (
|
||||
"Do not output textual tool intent/templates (DSML/XML/JSON). "
|
||||
"If a tool is needed, return native tool_calls only."
|
||||
),
|
||||
}
|
||||
]
|
||||
retries_done += 1
|
||||
resp = model.chat(repair_msgs, llm_tools, on_token=on_token)
|
||||
continue
|
||||
if on_progress:
|
||||
on_progress(f"{progress_label} retry-empty ({retries_done + 1}/{retry_max})…")
|
||||
if retry_delay_ms > 0:
|
||||
time.sleep(float(retry_delay_ms) / 1000.0)
|
||||
retries_done += 1
|
||||
resp = model.chat(msgs, llm_tools, on_token=on_token)
|
||||
|
||||
|
||||
def _extract_dsml_invoke_names(text: str) -> list[str]:
|
||||
raw = str(text or "")
|
||||
if not raw:
|
||||
return []
|
||||
out: list[str] = []
|
||||
seen: set[str] = set()
|
||||
for m in _DSML_INVOKE_NAME_RE.finditer(raw):
|
||||
nm = str(m.group(1) or "").strip()
|
||||
if not nm or nm in seen:
|
||||
continue
|
||||
seen.add(nm)
|
||||
out.append(nm)
|
||||
if len(out) >= 8:
|
||||
break
|
||||
return out
|
||||
|
||||
|
||||
def _extract_textual_tool_intent_names(text: str) -> list[str]:
|
||||
raw = str(text or "")
|
||||
if not raw:
|
||||
return []
|
||||
lower = raw.lower()
|
||||
marker_hit = ("tool_calls" in lower) or ("invoke name" in lower) or ("parameter name" in lower)
|
||||
if not marker_hit:
|
||||
return []
|
||||
out: list[str] = []
|
||||
seen: set[str] = set()
|
||||
for nm in _extract_dsml_invoke_names(raw):
|
||||
key = str(nm or "").strip()
|
||||
if key and key not in seen:
|
||||
seen.add(key)
|
||||
out.append(key)
|
||||
if len(out) < 8:
|
||||
for m in _JSON_TOOL_NAME_RE.finditer(raw):
|
||||
nm = str(m.group(1) or "").strip()
|
||||
if not nm or nm in seen:
|
||||
continue
|
||||
seen.add(nm)
|
||||
out.append(nm)
|
||||
if len(out) >= 8:
|
||||
break
|
||||
if out:
|
||||
return out
|
||||
return ["unknown_tool"]
|
||||
|
||||
|
||||
def _persist_dsml_protocol_mismatch_step(
|
||||
*,
|
||||
store: Any,
|
||||
session_id: str,
|
||||
turn_uuid: str,
|
||||
assistant_text: str,
|
||||
invoke_names: list[str],
|
||||
) -> _LoopStepResult:
|
||||
names = [str(x or "").strip() for x in (invoke_names or []) if str(x or "").strip()]
|
||||
if not names:
|
||||
names = ["unknown_tool"]
|
||||
stored_tool_calls: list[dict[str, Any]] = []
|
||||
for nm in names:
|
||||
stored_tool_calls.append(
|
||||
{
|
||||
"id": f"call_dsml_{uuid.uuid4().hex}",
|
||||
"name": nm,
|
||||
"arguments": {},
|
||||
"thought_signature": None,
|
||||
}
|
||||
)
|
||||
assistant_row = store.add_message(
|
||||
session_id=session_id,
|
||||
role="assistant",
|
||||
content="",
|
||||
tool_calls=stored_tool_calls,
|
||||
turn_uuid=turn_uuid,
|
||||
event_type="tool_call",
|
||||
event_payload={
|
||||
"protocol_mismatch": "textual_tool_intent",
|
||||
"raw_excerpt": str(assistant_text or "")[:2000],
|
||||
},
|
||||
)
|
||||
for tc in stored_tool_calls:
|
||||
tcid = str(tc.get("id") or "").strip()
|
||||
tname = str(tc.get("name") or "").strip() or "unknown_tool"
|
||||
tool_result = {
|
||||
"ok": False,
|
||||
"error_code": "model_protocol_mismatch_dsml",
|
||||
"error": "model_returned_textual_tool_intent_instead_of_native_tool_calls",
|
||||
"detail": {"tool_name": tname},
|
||||
}
|
||||
store.add_message(
|
||||
session_id=session_id,
|
||||
role="tool",
|
||||
content=_json_dumps_safe(tool_result),
|
||||
tool_calls={
|
||||
"tool_call_id": tcid,
|
||||
"name": tname,
|
||||
"assistant_message_id": int(getattr(assistant_row, "id", 0) or 0),
|
||||
},
|
||||
turn_uuid=turn_uuid,
|
||||
event_type="tool_result",
|
||||
event_payload={"tool_name": tname, "protocol_mismatch": "textual_tool_intent"},
|
||||
)
|
||||
return _LoopStepResult(
|
||||
assistant_text="",
|
||||
llm_tool_calls=[],
|
||||
assistant_msg_id=int(getattr(assistant_row, "id", 0) or 0),
|
||||
)
|
||||
|
||||
|
||||
def _persist_assistant_step(
|
||||
*,
|
||||
store: Any,
|
||||
|
|
@ -843,10 +1012,8 @@ def _persist_assistant_step(
|
|||
|
||||
reasoning_chunks, assistant_body = _split_reasoning_and_body(assistant_text, explicit_reasoning=reasoning_text)
|
||||
reasoning_full = "\n".join([str(x or "").strip() for x in reasoning_chunks if str(x or "").strip()]).strip()
|
||||
if not str(assistant_body or "").strip() and not stored_tool_calls:
|
||||
# Provider/model can occasionally return an empty body; persist a visible stub
|
||||
# so UI doesn't look "stuck" and operators can diagnose from history.
|
||||
assistant_body = "(空响应)模型返回了空内容,请重试一次;若持续出现,请检查模型网关/上游返回。"
|
||||
# Keep empty body as-is when model returns nothing and there are no tool calls.
|
||||
# The UI should treat this as an invisible intermediate/final empty response.
|
||||
if not thinking_mode_enabled:
|
||||
for idx, chunk in enumerate(reasoning_chunks):
|
||||
store.add_message(
|
||||
|
|
@ -1022,20 +1189,37 @@ def run_oclaw_direct_loop(
|
|||
lang=lang,
|
||||
wire_policy_role=wire_policy_role,
|
||||
)
|
||||
resp = model.chat(msgs, llm_tools, on_token=on_token)
|
||||
resp = _chat_with_empty_body_retry(
|
||||
model=model,
|
||||
msgs=msgs,
|
||||
llm_tools=llm_tools,
|
||||
on_token=on_token,
|
||||
on_progress=on_progress,
|
||||
progress_label="oclaw: think",
|
||||
)
|
||||
assistant_text = str(getattr(resp, "content", "") or "")
|
||||
reasoning_text = str(getattr(resp, "reasoning_content", "") or "")
|
||||
llm_tool_calls = list(getattr(resp, "tool_calls", []) or [])
|
||||
textual_tool_intent_names = _extract_textual_tool_intent_names(assistant_text) if not llm_tool_calls else []
|
||||
|
||||
step = _persist_assistant_step(
|
||||
store=store,
|
||||
session_id=session_id,
|
||||
turn_uuid=turn_uuid,
|
||||
assistant_text=assistant_text,
|
||||
reasoning_text=reasoning_text,
|
||||
llm_tool_calls=llm_tool_calls,
|
||||
thinking_mode_enabled=bool(getattr(model, "thinking_mode_enabled", False)),
|
||||
)
|
||||
if textual_tool_intent_names:
|
||||
step = _persist_dsml_protocol_mismatch_step(
|
||||
store=store,
|
||||
session_id=session_id,
|
||||
turn_uuid=turn_uuid,
|
||||
assistant_text=assistant_text,
|
||||
invoke_names=textual_tool_intent_names,
|
||||
)
|
||||
else:
|
||||
step = _persist_assistant_step(
|
||||
store=store,
|
||||
session_id=session_id,
|
||||
turn_uuid=turn_uuid,
|
||||
assistant_text=assistant_text,
|
||||
reasoning_text=reasoning_text,
|
||||
llm_tool_calls=llm_tool_calls,
|
||||
thinking_mode_enabled=bool(getattr(model, "thinking_mode_enabled", False)),
|
||||
)
|
||||
final_text = step.assistant_text
|
||||
if not step.llm_tool_calls:
|
||||
break
|
||||
|
|
@ -1109,7 +1293,14 @@ def run_oclaw_direct_loop(
|
|||
active_turn_uuid=turn_uuid,
|
||||
)
|
||||
# Final pass forbids extra tool calls; model must synthesize answer.
|
||||
resp = model.chat(msgs, [], on_token=on_token)
|
||||
resp = _chat_with_empty_body_retry(
|
||||
model=model,
|
||||
msgs=msgs,
|
||||
llm_tools=[],
|
||||
on_token=on_token,
|
||||
on_progress=on_progress,
|
||||
progress_label="oclaw: finalize",
|
||||
)
|
||||
step = _persist_assistant_step(
|
||||
store=store,
|
||||
session_id=session_id,
|
||||
|
|
|
|||
|
|
@ -36,6 +36,7 @@ from oclaw.runtime.worker import ensure_worker_started
|
|||
from oclaw.runtime.orchestration.trace import new_span_id, new_trace_id
|
||||
from oclaw.runtime.chat.tool_runtime import compact_turn_tool_messages_for_storage
|
||||
from oclaw.runtime.chat.model_path_audit import ensure_no_tool_or_embedded_image_payload
|
||||
from oclaw.runtime.tools.base import ToolRegistry
|
||||
from oclaw.runtime.tools.local_sdk import local_adapter_startup_self_check
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
|
@ -951,6 +952,92 @@ class OclawGateway:
|
|||
except Exception:
|
||||
selected_executor = executor
|
||||
|
||||
system_prompt_override = ""
|
||||
tools_override = None
|
||||
if interaction_mode == "expert":
|
||||
from oclaw.runtime.plan_agent_v2.gateway_adapter import evaluate_gateway_expert_turn_shadow
|
||||
from oclaw.runtime.plan_agent_v2.tool_specs import DEFAULT_SESSION_KEY, materialize_plan_mode_v2_tools
|
||||
|
||||
execution_mode = str(base_metadata.get("execution_mode") or "agent").strip().lower()
|
||||
if execution_mode not in {"agent", "plan"}:
|
||||
execution_mode = "agent"
|
||||
try:
|
||||
self.store.set_setting(DEFAULT_SESSION_KEY, str(msg.session_id or ""))
|
||||
except Exception:
|
||||
pass
|
||||
# Respect store setting AIA_EXPERT_PLAN_AGENT_V2_ENABLED (default off); do not force cutover.
|
||||
shadow = evaluate_gateway_expert_turn_shadow(
|
||||
store=self.store,
|
||||
msg=msg,
|
||||
lang=lang,
|
||||
interaction_mode=interaction_mode,
|
||||
requested_specialist=requested_specialist,
|
||||
execution_mode=execution_mode,
|
||||
base_system_prompt=str(getattr(selected_executor, "system_prompt", "") or ""),
|
||||
force_flag=False,
|
||||
trace_id=trace_id,
|
||||
parent_span_id=None,
|
||||
)
|
||||
if shadow.used_v2 and shadow.decision is not None:
|
||||
action = str(shadow.decision.action or "")
|
||||
if action in {"enter_plan", "stay_plan"}:
|
||||
elapsed_ms = int((time.perf_counter() - t0) * 1000)
|
||||
_trace_local(
|
||||
event_type="response_sent",
|
||||
payload={"ok": True, "elapsed_ms": elapsed_ms, "mode": "sync_direct", "plan_action": action},
|
||||
started_at=t0,
|
||||
)
|
||||
_flush_trace_rows()
|
||||
return OclawGatewayResult(
|
||||
run_id=rid,
|
||||
reply_text=str(shadow.decision.reply_text or ""),
|
||||
trace_id=trace_id,
|
||||
elapsed_ms=elapsed_ms,
|
||||
mode="sync_direct",
|
||||
selected_specialist=requested_specialist,
|
||||
interaction_mode=interaction_mode,
|
||||
dispatch_reason=f"plan_agent_v2:{action}",
|
||||
manager_selected_specialist=requested_specialist,
|
||||
requested_specialist=requested_specialist,
|
||||
dynamic_agent_used=False,
|
||||
dynamic_agent_name="",
|
||||
relay_pointer_count=int(relay_stats.get("relay_pointer_count") or 0),
|
||||
relay_envelope_present=bool(relay_stats.get("relay_envelope_present")),
|
||||
relay_envelope_pointer_count=int(relay_stats.get("relay_envelope_pointer_count") or 0),
|
||||
relay_ttl_turn_count=int(ttl_stats.get("turn") or 0),
|
||||
relay_ttl_session_count=int(ttl_stats.get("session") or 0),
|
||||
relay_ttl_keep_count=int(ttl_stats.get("keep") or 0),
|
||||
)
|
||||
if action == "run_agent":
|
||||
system_prompt_override = str(shadow.decision.system_prompt_override or "")
|
||||
exec_tools = getattr(selected_executor, "tools", None)
|
||||
if isinstance(exec_tools, ToolRegistry):
|
||||
merged = ToolRegistry(exec_tools.list() + materialize_plan_mode_v2_tools(store=self.store))
|
||||
tools_override = merged
|
||||
_trace_local(
|
||||
event_type="plan_mode_tools_augmented",
|
||||
payload={"base_count": len(exec_tools.list()), "merged_count": len(merged.list())},
|
||||
started_at=t0,
|
||||
)
|
||||
try:
|
||||
plan_mode = str((shadow.decision.plan_state or {}).get("mode") or "").strip().lower()
|
||||
except Exception:
|
||||
plan_mode = ""
|
||||
if plan_mode == "plan":
|
||||
from oclaw.runtime.plan_agent_v2.tool_policy import filter_tools_for_mode
|
||||
|
||||
if isinstance(tools_override, ToolRegistry):
|
||||
filtered = filter_tools_for_mode(registry=tools_override, mode="plan")
|
||||
tools_override = ToolRegistry(filtered)
|
||||
_trace_local(
|
||||
event_type="plan_mode_tools_filtered",
|
||||
payload={
|
||||
"before_count": len(merged.list()) if isinstance(exec_tools, ToolRegistry) else len(filtered),
|
||||
"after_count": len(filtered),
|
||||
},
|
||||
started_at=t0,
|
||||
)
|
||||
|
||||
route_mode = "sync_direct"
|
||||
route_msg = StandardMessage(
|
||||
session_id=msg.session_id,
|
||||
|
|
@ -1064,10 +1151,10 @@ class OclawGateway:
|
|||
)
|
||||
try:
|
||||
model = getattr(selected_executor, "model", None)
|
||||
tools = getattr(selected_executor, "tools", None)
|
||||
tools = tools_override if tools_override is not None else getattr(selected_executor, "tools", None)
|
||||
if model is None or tools is None:
|
||||
raise RuntimeError("executor missing model/tools")
|
||||
sys_prompt = str(getattr(selected_executor, "system_prompt", "") or "")
|
||||
sys_prompt = system_prompt_override or str(getattr(selected_executor, "system_prompt", "") or "")
|
||||
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):
|
||||
|
|
|
|||
42
runtime/plan_agent_v2/__init__.py
Normal file
42
runtime/plan_agent_v2/__init__.py
Normal file
|
|
@ -0,0 +1,42 @@
|
|||
from .adapter import PlanAgentV2Decision, evaluate_for_expert_mode
|
||||
from .compat import build_shadow_gateway_result, legacy_gateway_result_keys
|
||||
from .gateway_adapter import GatewayPlanV2AdapterOutput, evaluate_gateway_expert_turn_shadow
|
||||
from .manager import PlanModeManagerV2
|
||||
from .models import PLAN_MODE_NORMAL, PLAN_MODE_PLAN, PlanAgentStateV2
|
||||
from .prompt_injector import build_plan_mode_prefix, inject_plan_context
|
||||
from .state_store import PlanAgentStateStoreV2
|
||||
from .switch import should_route_to_v2, v2_feature_enabled
|
||||
from .tool_policy import filter_tools_for_mode, plan_mode_allowed_tool_names
|
||||
from .tool_specs import (
|
||||
enter_plan_mode_v2_tool,
|
||||
exit_plan_mode_v2_tool,
|
||||
is_plan_mode_v2_active,
|
||||
materialize_plan_mode_v2_tools,
|
||||
)
|
||||
from .trace import emit_plan_agent_v2_trace
|
||||
|
||||
__all__ = [
|
||||
"PLAN_MODE_NORMAL",
|
||||
"PLAN_MODE_PLAN",
|
||||
"PlanAgentStateV2",
|
||||
"PlanAgentStateStoreV2",
|
||||
"PlanModeManagerV2",
|
||||
"build_plan_mode_prefix",
|
||||
"inject_plan_context",
|
||||
"filter_tools_for_mode",
|
||||
"plan_mode_allowed_tool_names",
|
||||
"enter_plan_mode_v2_tool",
|
||||
"exit_plan_mode_v2_tool",
|
||||
"materialize_plan_mode_v2_tools",
|
||||
"is_plan_mode_v2_active",
|
||||
"v2_feature_enabled",
|
||||
"should_route_to_v2",
|
||||
"PlanAgentV2Decision",
|
||||
"evaluate_for_expert_mode",
|
||||
"GatewayPlanV2AdapterOutput",
|
||||
"evaluate_gateway_expert_turn_shadow",
|
||||
"emit_plan_agent_v2_trace",
|
||||
"legacy_gateway_result_keys",
|
||||
"build_shadow_gateway_result",
|
||||
]
|
||||
|
||||
281
runtime/plan_agent_v2/adapter.py
Normal file
281
runtime/plan_agent_v2/adapter.py
Normal file
|
|
@ -0,0 +1,281 @@
|
|||
from __future__ import annotations
|
||||
|
||||
from dataclasses import dataclass
|
||||
from typing import Any
|
||||
|
||||
from .manager import PlanModeManagerV2
|
||||
from .models import PLAN_MODE_PLAN
|
||||
from .prompt_injector import build_plan_mode_prefix, inject_plan_context
|
||||
from .trace import emit_plan_agent_v2_trace
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class PlanAgentV2Decision:
|
||||
action: str # enter_plan | stay_plan | run_agent
|
||||
reply_text: str
|
||||
plan_state: dict[str, Any]
|
||||
system_prompt_override: str = ""
|
||||
|
||||
|
||||
def _is_confirm_text(text: str) -> bool:
|
||||
t = str(text or "").strip().lower()
|
||||
return t in {"确认", "确认计划", "同意", "通过", "approve", "approved", "confirm", "yes"}
|
||||
|
||||
|
||||
def _normalize_user_text(text: str) -> str:
|
||||
return " ".join(str(text or "").strip().lower().split())
|
||||
|
||||
|
||||
def _is_low_signal_continue(text_norm: str) -> bool:
|
||||
t = str(text_norm or "").strip().lower()
|
||||
return t in {
|
||||
"继续",
|
||||
"继续啊",
|
||||
"继续吧",
|
||||
"可以",
|
||||
"好的",
|
||||
"好",
|
||||
"ok",
|
||||
"okay",
|
||||
"go on",
|
||||
"continue",
|
||||
}
|
||||
|
||||
|
||||
def _confirm_strategy(store: Any) -> str:
|
||||
try:
|
||||
raw = str(store.get_setting("AIA_EXPERT_PLAN_CONFIRM_STRATEGY") or "").strip().lower()
|
||||
except Exception:
|
||||
raw = ""
|
||||
if raw in {"auto", "strict", "off"}:
|
||||
return raw
|
||||
return "strict"
|
||||
|
||||
|
||||
def _last_user_text_norm_from_history(*, store: Any, session_id: str) -> str:
|
||||
"""Most recent persisted user message (current turn is usually not persisted yet)."""
|
||||
try:
|
||||
msgs = store.get_messages(session_id=session_id, limit=120)
|
||||
except Exception:
|
||||
return ""
|
||||
for m in reversed(msgs):
|
||||
if str(getattr(m, "role", "") or "").strip().lower() == "user":
|
||||
return _normalize_user_text(str(getattr(m, "content", "") or ""))
|
||||
return ""
|
||||
|
||||
|
||||
def _agent_conversation_stall_suffix(*, lang: str) -> str:
|
||||
is_en = str(lang or "").startswith("en")
|
||||
if is_en:
|
||||
return (
|
||||
"\n\n[Conversation stall guard — agent mode]\n"
|
||||
"The user's latest message matches their previous user message in this session.\n"
|
||||
"- Do not repeat your last assistant reply or restate \"I will now…\" boilerplate.\n"
|
||||
"- Make substantive progress: execute the next concrete tool step, produce new actionable output, "
|
||||
"or ask exactly one specific blocking question.\n"
|
||||
)
|
||||
return (
|
||||
"\n\n【对话停滞防护 · agent 模式】\n"
|
||||
"检测到用户本条输入与上一轮用户输入相同(会话已持久化部分)。\n"
|
||||
"- 禁止复述上一轮助手回复或重复「接下来我将…」式独白。\n"
|
||||
"- 必须给出实质进展:执行具体工具步骤、写出新的可执行结果,或只提一个关键追问。\n"
|
||||
)
|
||||
|
||||
|
||||
def evaluate_for_expert_mode(
|
||||
*,
|
||||
store: Any,
|
||||
session_id: str,
|
||||
lang: str,
|
||||
requested_specialist: str,
|
||||
user_text: str,
|
||||
execution_mode: str = "agent",
|
||||
base_system_prompt: str,
|
||||
trace_id: str | None = None,
|
||||
parent_span_id: str | None = None,
|
||||
) -> PlanAgentV2Decision:
|
||||
mgr = PlanModeManagerV2(store=store)
|
||||
st = mgr.load_state(session_id=session_id)
|
||||
txt = str(user_text or "").strip()
|
||||
txt_norm = _normalize_user_text(txt)
|
||||
exec_mode = str(execution_mode or "").strip().lower()
|
||||
if exec_mode not in {"agent", "plan"}:
|
||||
exec_mode = "plan"
|
||||
confirm_strategy = _confirm_strategy(store)
|
||||
|
||||
if exec_mode == "agent" and st.mode != PLAN_MODE_PLAN:
|
||||
emit_plan_agent_v2_trace(
|
||||
store=store,
|
||||
session_id=session_id,
|
||||
trace_id=trace_id,
|
||||
parent_span_id=parent_span_id,
|
||||
event_type="plan_mode_bypassed",
|
||||
payload={"requested_mode": "agent", "plan_mode_state": str(st.mode or "")},
|
||||
)
|
||||
last_user_norm = _last_user_text_norm_from_history(store=store, session_id=session_id)
|
||||
stall = bool(txt_norm and last_user_norm and txt_norm == last_user_norm)
|
||||
override = ""
|
||||
if stall:
|
||||
emit_plan_agent_v2_trace(
|
||||
store=store,
|
||||
session_id=session_id,
|
||||
trace_id=trace_id,
|
||||
parent_span_id=parent_span_id,
|
||||
event_type="agent_mode_conversation_stall",
|
||||
payload={"reason": "repeated_user_message"},
|
||||
)
|
||||
base = str(base_system_prompt or "").strip()
|
||||
suffix = _agent_conversation_stall_suffix(lang=lang).strip()
|
||||
override = f"{base}\n\n{suffix}".strip()
|
||||
return PlanAgentV2Decision(
|
||||
action="run_agent",
|
||||
reply_text="",
|
||||
plan_state=st.to_dict(),
|
||||
system_prompt_override=override,
|
||||
)
|
||||
|
||||
if st.mode != PLAN_MODE_PLAN:
|
||||
entered = mgr.enter(session_id=session_id, owner_specialist=requested_specialist, force_new_plan=False)
|
||||
emit_plan_agent_v2_trace(
|
||||
store=store,
|
||||
session_id=session_id,
|
||||
trace_id=trace_id,
|
||||
parent_span_id=parent_span_id,
|
||||
event_type="plan_mode_entered",
|
||||
payload={"owner_specialist": entered.owner_specialist, "plan_id": entered.plan_id},
|
||||
)
|
||||
prefix = build_plan_mode_prefix(state=entered, lang=lang)
|
||||
return PlanAgentV2Decision(
|
||||
action="run_agent",
|
||||
reply_text="",
|
||||
plan_state=entered.to_dict(),
|
||||
system_prompt_override=f"{prefix}\n\n{str(base_system_prompt or '').strip()}".strip(),
|
||||
)
|
||||
|
||||
st = mgr.refresh_plan_content(session_id=session_id)
|
||||
st = mgr.update_loop_guard(session_id=session_id, user_text_norm=txt_norm)
|
||||
if _is_confirm_text(txt):
|
||||
if exec_mode != "agent" and confirm_strategy == "strict":
|
||||
emit_plan_agent_v2_trace(
|
||||
store=store,
|
||||
session_id=session_id,
|
||||
trace_id=trace_id,
|
||||
parent_span_id=parent_span_id,
|
||||
event_type="plan_mode_confirm_blocked",
|
||||
payload={
|
||||
"reason": "execution_mode_not_agent",
|
||||
"requested_mode": exec_mode,
|
||||
"confirm_strategy": confirm_strategy,
|
||||
},
|
||||
)
|
||||
blocked_reply = (
|
||||
"Plan is ready. Please switch to agent mode, then confirm to execute."
|
||||
if str(lang or "").startswith("en")
|
||||
else "计划已就绪。请先切换到 agent 模式,再回复“确认”开始执行。"
|
||||
)
|
||||
return PlanAgentV2Decision(
|
||||
action="stay_plan",
|
||||
reply_text=blocked_reply,
|
||||
plan_state=st.to_dict(),
|
||||
system_prompt_override="",
|
||||
)
|
||||
if exec_mode != "agent" and confirm_strategy == "auto":
|
||||
emit_plan_agent_v2_trace(
|
||||
store=store,
|
||||
session_id=session_id,
|
||||
trace_id=trace_id,
|
||||
parent_span_id=parent_span_id,
|
||||
event_type="plan_mode_confirm_auto_switched",
|
||||
payload={"from_mode": exec_mode, "to_mode": "agent", "confirm_strategy": confirm_strategy},
|
||||
)
|
||||
confirmed = mgr.confirm(session_id=session_id)
|
||||
emit_plan_agent_v2_trace(
|
||||
store=store,
|
||||
session_id=session_id,
|
||||
trace_id=trace_id,
|
||||
parent_span_id=parent_span_id,
|
||||
event_type="plan_mode_confirmed",
|
||||
payload={
|
||||
"plan_id": confirmed.plan_id,
|
||||
"plan_confirmed": bool(confirmed.plan_confirmed),
|
||||
"confirm_strategy": confirm_strategy,
|
||||
},
|
||||
)
|
||||
next_system = inject_plan_context(base_system=base_system_prompt, state=confirmed, lang=lang)
|
||||
reply = mgr.build_approved_execution_message(state=confirmed, lang=lang)
|
||||
return PlanAgentV2Decision(
|
||||
action="run_agent",
|
||||
reply_text=reply,
|
||||
plan_state=confirmed.to_dict(),
|
||||
system_prompt_override=next_system,
|
||||
)
|
||||
|
||||
if _is_low_signal_continue(txt_norm):
|
||||
emit_plan_agent_v2_trace(
|
||||
store=store,
|
||||
session_id=session_id,
|
||||
trace_id=trace_id,
|
||||
parent_span_id=parent_span_id,
|
||||
event_type="plan_mode_active",
|
||||
payload={"plan_id": st.plan_id, "plan_path": st.plan_path, "loop_guard": "low_signal_continue"},
|
||||
)
|
||||
low_signal_reply = (
|
||||
"Plan mode detected a low-information continuation. "
|
||||
"Please provide concrete plan adjustments, or switch to agent mode and reply 'confirm' to execute."
|
||||
if str(lang or "").startswith("en")
|
||||
else "检测到低信息续写(如“继续/可以”)。请给出具体计划修改点,或切换到 agent 模式后回复“确认”直接执行。"
|
||||
)
|
||||
return PlanAgentV2Decision(
|
||||
action="stay_plan",
|
||||
reply_text=low_signal_reply,
|
||||
plan_state=st.to_dict(),
|
||||
system_prompt_override="",
|
||||
)
|
||||
|
||||
if int(st.plan_loop_count or 0) >= 2:
|
||||
emit_plan_agent_v2_trace(
|
||||
store=store,
|
||||
session_id=session_id,
|
||||
trace_id=trace_id,
|
||||
parent_span_id=parent_span_id,
|
||||
event_type="plan_mode_active",
|
||||
payload={"plan_id": st.plan_id, "plan_path": st.plan_path, "loop_guard": "hard_block"},
|
||||
)
|
||||
anti_loop_reply = (
|
||||
"I am in plan mode. I will only output a concise executable plan. "
|
||||
"If you want me to execute, switch to agent mode and reply 'confirm'."
|
||||
if str(lang or "").startswith("en")
|
||||
else "当前为 plan 模式,我只输出可执行计划。若要开始执行,请切换到 agent 模式并回复“确认”。"
|
||||
)
|
||||
return PlanAgentV2Decision(
|
||||
action="stay_plan",
|
||||
reply_text=anti_loop_reply,
|
||||
plan_state=st.to_dict(),
|
||||
system_prompt_override="",
|
||||
)
|
||||
|
||||
prefix = build_plan_mode_prefix(state=st, lang=lang)
|
||||
anti_loop_suffix = (
|
||||
"\n\n[Anti-loop guard]\n"
|
||||
"- Do not repeat the previous response.\n"
|
||||
"- If user asks similarly, refine with more concrete steps, checks, and fallback.\n"
|
||||
"- Keep output as plan only; do not pretend execution is complete."
|
||||
)
|
||||
emit_plan_agent_v2_trace(
|
||||
store=store,
|
||||
session_id=session_id,
|
||||
trace_id=trace_id,
|
||||
parent_span_id=parent_span_id,
|
||||
event_type="plan_mode_active",
|
||||
payload={"plan_id": st.plan_id, "plan_path": st.plan_path, "loop_count": int(st.plan_loop_count or 0)},
|
||||
)
|
||||
return PlanAgentV2Decision(
|
||||
action="run_agent",
|
||||
reply_text="",
|
||||
plan_state=st.to_dict(),
|
||||
system_prompt_override=f"{prefix}{anti_loop_suffix}\n\n{str(base_system_prompt or '').strip()}".strip(),
|
||||
)
|
||||
|
||||
|
||||
__all__ = ["PlanAgentV2Decision", "evaluate_for_expert_mode"]
|
||||
|
||||
45
runtime/plan_agent_v2/compat.py
Normal file
45
runtime/plan_agent_v2/compat.py
Normal file
|
|
@ -0,0 +1,45 @@
|
|||
from __future__ import annotations
|
||||
|
||||
from typing import Any
|
||||
|
||||
from .adapter import PlanAgentV2Decision
|
||||
from oclaw.runtime.gateway import OclawGatewayResult
|
||||
|
||||
|
||||
def legacy_gateway_result_keys() -> set[str]:
|
||||
return set(OclawGatewayResult.__dataclass_fields__.keys())
|
||||
|
||||
|
||||
def build_shadow_gateway_result(
|
||||
*,
|
||||
decision: PlanAgentV2Decision,
|
||||
run_id: str,
|
||||
trace_id: str,
|
||||
elapsed_ms: int,
|
||||
requested_specialist: str,
|
||||
) -> dict[str, Any]:
|
||||
return {
|
||||
"run_id": str(run_id),
|
||||
"reply_text": str(decision.reply_text or ""),
|
||||
"trace_id": str(trace_id),
|
||||
"elapsed_ms": int(elapsed_ms),
|
||||
"mode": "sync_direct",
|
||||
"task_id": None,
|
||||
"selected_specialist": str((decision.plan_state or {}).get("owner_specialist") or requested_specialist or "generalist"),
|
||||
"interaction_mode": "expert",
|
||||
"dispatch_reason": f"plan_agent_v2:{decision.action}",
|
||||
"manager_selected_specialist": str((decision.plan_state or {}).get("owner_specialist") or requested_specialist or "generalist"),
|
||||
"requested_specialist": str(requested_specialist or "generalist"),
|
||||
"dynamic_agent_used": False,
|
||||
"dynamic_agent_name": "",
|
||||
"relay_pointer_count": 0,
|
||||
"relay_envelope_present": False,
|
||||
"relay_envelope_pointer_count": 0,
|
||||
"relay_ttl_turn_count": 0,
|
||||
"relay_ttl_session_count": 0,
|
||||
"relay_ttl_keep_count": 0,
|
||||
}
|
||||
|
||||
|
||||
__all__ = ["build_shadow_gateway_result", "legacy_gateway_result_keys"]
|
||||
|
||||
54
runtime/plan_agent_v2/gateway_adapter.py
Normal file
54
runtime/plan_agent_v2/gateway_adapter.py
Normal file
|
|
@ -0,0 +1,54 @@
|
|||
from __future__ import annotations
|
||||
|
||||
from dataclasses import dataclass
|
||||
from typing import Any
|
||||
|
||||
from .adapter import PlanAgentV2Decision, evaluate_for_expert_mode
|
||||
from .switch import should_route_to_v2
|
||||
from oclaw.runtime.types import StandardMessage
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class GatewayPlanV2AdapterOutput:
|
||||
used_v2: bool
|
||||
decision: PlanAgentV2Decision | None
|
||||
|
||||
|
||||
def evaluate_gateway_expert_turn_shadow(
|
||||
*,
|
||||
store: Any,
|
||||
msg: StandardMessage,
|
||||
lang: str,
|
||||
interaction_mode: str,
|
||||
requested_specialist: str,
|
||||
execution_mode: str = "",
|
||||
base_system_prompt: str,
|
||||
force_flag: bool = False,
|
||||
trace_id: str | None = None,
|
||||
parent_span_id: str | None = None,
|
||||
) -> GatewayPlanV2AdapterOutput:
|
||||
if not should_route_to_v2(store=store, interaction_mode=interaction_mode, force_flag=force_flag):
|
||||
return GatewayPlanV2AdapterOutput(used_v2=False, decision=None)
|
||||
meta = msg.metadata if isinstance(msg.metadata, dict) else {}
|
||||
if "plan_agent_version" in meta:
|
||||
if str(meta.get("plan_agent_version") or "").strip().lower() != "v2":
|
||||
return GatewayPlanV2AdapterOutput(used_v2=False, decision=None)
|
||||
eff_mode = str(execution_mode or "").strip().lower()
|
||||
if eff_mode not in {"agent", "plan"}:
|
||||
eff_mode = "plan" if force_flag else "agent"
|
||||
dec = evaluate_for_expert_mode(
|
||||
store=store,
|
||||
session_id=str(msg.session_id or ""),
|
||||
lang=lang,
|
||||
requested_specialist=requested_specialist,
|
||||
user_text=str(msg.text or ""),
|
||||
execution_mode=eff_mode,
|
||||
base_system_prompt=base_system_prompt,
|
||||
trace_id=trace_id,
|
||||
parent_span_id=parent_span_id,
|
||||
)
|
||||
return GatewayPlanV2AdapterOutput(used_v2=True, decision=dec)
|
||||
|
||||
|
||||
__all__ = ["GatewayPlanV2AdapterOutput", "evaluate_gateway_expert_turn_shadow"]
|
||||
|
||||
185
runtime/plan_agent_v2/manager.py
Normal file
185
runtime/plan_agent_v2/manager.py
Normal file
|
|
@ -0,0 +1,185 @@
|
|||
from __future__ import annotations
|
||||
|
||||
import time
|
||||
import uuid
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
from .models import PLAN_MODE_NORMAL, PLAN_MODE_PLAN, PlanAgentStateV2
|
||||
from .state_store import PlanAgentStateStoreV2
|
||||
|
||||
|
||||
def _default_plan_dir() -> Path:
|
||||
return Path(__file__).resolve().parents[2] / "data" / "plans"
|
||||
|
||||
|
||||
def _resolve_plan_dir(store: Any) -> Path:
|
||||
raw = str(store.get_setting("AIA_EXPERT_PLAN_FILE_DIR") or "").strip()
|
||||
if raw:
|
||||
p = Path(raw)
|
||||
return p if p.is_absolute() else (Path(__file__).resolve().parents[2] / p)
|
||||
return _default_plan_dir()
|
||||
|
||||
|
||||
def _plan_template() -> str:
|
||||
return (
|
||||
"# Plan\n\n"
|
||||
"## Goal\n"
|
||||
"- \n\n"
|
||||
"## Scope\n"
|
||||
"- \n\n"
|
||||
"## Steps\n"
|
||||
"1. \n"
|
||||
"2. \n"
|
||||
"3. \n\n"
|
||||
"## Risks\n"
|
||||
"- \n\n"
|
||||
"## Acceptance\n"
|
||||
"- \n"
|
||||
)
|
||||
|
||||
|
||||
class PlanModeManagerV2:
|
||||
def __init__(self, *, store: Any):
|
||||
self._store = store
|
||||
self._state_store = PlanAgentStateStoreV2(store)
|
||||
|
||||
def load_state(self, *, session_id: str) -> PlanAgentStateV2:
|
||||
return self._state_store.load(session_id=session_id)
|
||||
|
||||
def enter(
|
||||
self,
|
||||
*,
|
||||
session_id: str,
|
||||
owner_specialist: str,
|
||||
force_new_plan: bool = False,
|
||||
) -> PlanAgentStateV2:
|
||||
prev = self._state_store.load(session_id=session_id)
|
||||
if prev.mode == PLAN_MODE_PLAN and not force_new_plan:
|
||||
return prev
|
||||
sid = str(session_id or "").strip()
|
||||
if not sid:
|
||||
return prev
|
||||
plan_id = uuid.uuid4().hex
|
||||
plan_root = _resolve_plan_dir(self._store) / sid
|
||||
plan_root.mkdir(parents=True, exist_ok=True)
|
||||
plan_path = plan_root / f"{plan_id}.md"
|
||||
plan_content = _plan_template()
|
||||
plan_path.write_text(plan_content, encoding="utf-8")
|
||||
now_ms = int(time.time() * 1000)
|
||||
next_state = PlanAgentStateV2(
|
||||
mode=PLAN_MODE_PLAN,
|
||||
owner_specialist=str(owner_specialist or "generalist").strip().lower() or "generalist",
|
||||
plan_id=plan_id,
|
||||
plan_path=str(plan_path),
|
||||
plan_content=plan_content,
|
||||
plan_confirmed=False,
|
||||
entered_at_ms=now_ms,
|
||||
updated_at_ms=now_ms,
|
||||
last_user_text_norm="",
|
||||
plan_loop_count=0,
|
||||
)
|
||||
return self._state_store.save(session_id=sid, state=next_state)
|
||||
|
||||
def refresh_plan_content(self, *, session_id: str) -> PlanAgentStateV2:
|
||||
st = self._state_store.load(session_id=session_id)
|
||||
p = Path(str(st.plan_path or "").strip())
|
||||
if not p.exists() or not p.is_file():
|
||||
return st
|
||||
content = p.read_text(encoding="utf-8", errors="replace")
|
||||
return self._state_store.save(
|
||||
session_id=session_id,
|
||||
state=PlanAgentStateV2(
|
||||
mode=st.mode,
|
||||
owner_specialist=st.owner_specialist,
|
||||
plan_id=st.plan_id,
|
||||
plan_path=st.plan_path,
|
||||
plan_content=content,
|
||||
plan_confirmed=st.plan_confirmed,
|
||||
entered_at_ms=st.entered_at_ms,
|
||||
updated_at_ms=st.updated_at_ms,
|
||||
last_user_text_norm=st.last_user_text_norm,
|
||||
plan_loop_count=st.plan_loop_count,
|
||||
),
|
||||
)
|
||||
|
||||
def update_loop_guard(self, *, session_id: str, user_text_norm: str) -> PlanAgentStateV2:
|
||||
st = self._state_store.load(session_id=session_id)
|
||||
nxt_count = int(st.plan_loop_count or 0) + 1 if user_text_norm and user_text_norm == st.last_user_text_norm else 0
|
||||
return self._state_store.save(
|
||||
session_id=session_id,
|
||||
state=PlanAgentStateV2(
|
||||
mode=st.mode,
|
||||
owner_specialist=st.owner_specialist,
|
||||
plan_id=st.plan_id,
|
||||
plan_path=st.plan_path,
|
||||
plan_content=st.plan_content,
|
||||
plan_confirmed=st.plan_confirmed,
|
||||
entered_at_ms=st.entered_at_ms,
|
||||
updated_at_ms=st.updated_at_ms,
|
||||
last_user_text_norm=user_text_norm,
|
||||
plan_loop_count=nxt_count,
|
||||
),
|
||||
)
|
||||
|
||||
def confirm(self, *, session_id: str) -> PlanAgentStateV2:
|
||||
st = self.refresh_plan_content(session_id=session_id)
|
||||
return self._state_store.save(
|
||||
session_id=session_id,
|
||||
state=PlanAgentStateV2(
|
||||
mode=PLAN_MODE_NORMAL,
|
||||
owner_specialist=st.owner_specialist,
|
||||
plan_id=st.plan_id,
|
||||
plan_path=st.plan_path,
|
||||
plan_content=st.plan_content,
|
||||
plan_confirmed=True,
|
||||
entered_at_ms=st.entered_at_ms,
|
||||
updated_at_ms=st.updated_at_ms,
|
||||
last_user_text_norm="",
|
||||
plan_loop_count=0,
|
||||
),
|
||||
)
|
||||
|
||||
def build_approved_execution_message(self, *, state: PlanAgentStateV2, lang: str) -> str:
|
||||
is_en = str(lang or "").startswith("en")
|
||||
plan_path = str(state.plan_path or "").strip() or "unknown"
|
||||
plan_content = str(state.plan_content or "").strip()
|
||||
if plan_content:
|
||||
if is_en:
|
||||
return (
|
||||
"User has approved your plan. You can now start implementation.\n\n"
|
||||
f"Plan file: {plan_path}\n\n"
|
||||
f"## Approved Plan\n{plan_content}"
|
||||
)
|
||||
return (
|
||||
"用户已确认计划,你可以开始执行实现。\n\n"
|
||||
f"计划文件:{plan_path}\n\n"
|
||||
f"## 已确认计划\n{plan_content}"
|
||||
)
|
||||
return (
|
||||
f"Plan approved. You can now start implementation. Plan file: {plan_path}"
|
||||
if is_en
|
||||
else f"计划已确认,你可以开始执行实现。计划文件:{plan_path}"
|
||||
)
|
||||
|
||||
def exit_without_confirm(self, *, session_id: str) -> PlanAgentStateV2:
|
||||
st = self._state_store.load(session_id=session_id)
|
||||
return self._state_store.save(
|
||||
session_id=session_id,
|
||||
state=PlanAgentStateV2(
|
||||
mode=PLAN_MODE_NORMAL,
|
||||
owner_specialist=st.owner_specialist,
|
||||
plan_id=st.plan_id,
|
||||
plan_path=st.plan_path,
|
||||
plan_content=st.plan_content,
|
||||
plan_confirmed=False,
|
||||
entered_at_ms=st.entered_at_ms,
|
||||
updated_at_ms=st.updated_at_ms,
|
||||
last_user_text_norm="",
|
||||
plan_loop_count=0,
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
__all__ = ["PlanModeManagerV2"]
|
||||
|
||||
49
runtime/plan_agent_v2/models.py
Normal file
49
runtime/plan_agent_v2/models.py
Normal file
|
|
@ -0,0 +1,49 @@
|
|||
from __future__ import annotations
|
||||
|
||||
from dataclasses import asdict, dataclass
|
||||
from typing import Any
|
||||
|
||||
|
||||
PLAN_MODE_NORMAL = "normal"
|
||||
PLAN_MODE_PLAN = "plan"
|
||||
_VALID_MODES = {PLAN_MODE_NORMAL, PLAN_MODE_PLAN}
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class PlanAgentStateV2:
|
||||
mode: str = PLAN_MODE_NORMAL
|
||||
owner_specialist: str = "generalist"
|
||||
plan_id: str = ""
|
||||
plan_path: str = ""
|
||||
plan_content: str = ""
|
||||
plan_confirmed: bool = False
|
||||
entered_at_ms: int = 0
|
||||
updated_at_ms: int = 0
|
||||
last_user_text_norm: str = ""
|
||||
plan_loop_count: int = 0
|
||||
|
||||
def to_dict(self) -> dict[str, Any]:
|
||||
return asdict(self)
|
||||
|
||||
@staticmethod
|
||||
def from_dict(raw: dict[str, Any] | None) -> "PlanAgentStateV2":
|
||||
obj = raw if isinstance(raw, dict) else {}
|
||||
mode = str(obj.get("mode") or PLAN_MODE_NORMAL).strip().lower()
|
||||
if mode not in _VALID_MODES:
|
||||
mode = PLAN_MODE_NORMAL
|
||||
return PlanAgentStateV2(
|
||||
mode=mode,
|
||||
owner_specialist=str(obj.get("owner_specialist") or "generalist").strip().lower() or "generalist",
|
||||
plan_id=str(obj.get("plan_id") or "").strip(),
|
||||
plan_path=str(obj.get("plan_path") or "").strip(),
|
||||
plan_content=str(obj.get("plan_content") or ""),
|
||||
plan_confirmed=bool(obj.get("plan_confirmed")),
|
||||
entered_at_ms=int(obj.get("entered_at_ms") or 0),
|
||||
updated_at_ms=int(obj.get("updated_at_ms") or 0),
|
||||
last_user_text_norm=str(obj.get("last_user_text_norm") or "").strip().lower(),
|
||||
plan_loop_count=int(obj.get("plan_loop_count") or 0),
|
||||
)
|
||||
|
||||
|
||||
__all__ = ["PLAN_MODE_NORMAL", "PLAN_MODE_PLAN", "PlanAgentStateV2"]
|
||||
|
||||
120
runtime/plan_agent_v2/prompt_injector.py
Normal file
120
runtime/plan_agent_v2/prompt_injector.py
Normal file
|
|
@ -0,0 +1,120 @@
|
|||
from __future__ import annotations
|
||||
|
||||
from pathlib import Path
|
||||
|
||||
from .models import PLAN_MODE_PLAN, PlanAgentStateV2
|
||||
|
||||
|
||||
def _plan_file_info(state: PlanAgentStateV2) -> str:
|
||||
plan_path = str(state.plan_path or "").strip()
|
||||
if not plan_path:
|
||||
return "No plan file path is available yet."
|
||||
p = Path(plan_path)
|
||||
if p.exists():
|
||||
return (
|
||||
f"A plan file already exists at {plan_path}. "
|
||||
"You can read it and make incremental edits."
|
||||
)
|
||||
return (
|
||||
f"No plan file exists yet. You should create your plan at {plan_path}."
|
||||
)
|
||||
|
||||
|
||||
def build_plan_mode_prefix(*, state: PlanAgentStateV2, lang: str) -> str:
|
||||
if state.mode != PLAN_MODE_PLAN:
|
||||
return ""
|
||||
is_en = str(lang or "").startswith("en")
|
||||
file_info = _plan_file_info(state)
|
||||
if is_en:
|
||||
return (
|
||||
"Plan mode is active. The user does not want execution yet.\n"
|
||||
"You MUST NOT make real project edits, run non-readonly tools, or claim implementation is done.\n"
|
||||
"## Execution discipline (critical)\n"
|
||||
"- Do NOT narrate as if you will run scripts, migrate files, or touch disk *in this turn*. "
|
||||
"Phrases like \"I'll write the script and run it\", \"starting migration now\", or "
|
||||
"\"let me execute\" mislead the user—refuse that pattern.\n"
|
||||
"- If the user needs real execution, say explicitly: switch to **agent mode** in the UI, "
|
||||
"then confirm; you cannot perform execution while plan mode is active.\n"
|
||||
"- Do NOT repeat the same \"next I will…\" monologue across turns. On vague follow-ups, "
|
||||
"edit the plan file or ask **one** concrete question—do not restate boilerplate.\n"
|
||||
"- Do not ask the user to \"approve the plan\" in chat when the product expects mode switch + "
|
||||
"confirm; instead tell them the handoff: agent mode → confirm to execute.\n\n"
|
||||
"## Plan File Info\n"
|
||||
f"{file_info}\n"
|
||||
"Only the plan file is allowed to be edited while in plan mode.\n\n"
|
||||
"## Plan Workflow\n"
|
||||
"### Phase 1: Initial Understanding\n"
|
||||
"- Understand the request and inspect relevant codepaths.\n"
|
||||
"- Reuse existing functions/utilities/patterns when possible.\n\n"
|
||||
"### Phase 2: Design\n"
|
||||
"- Propose a concrete implementation strategy with trade-offs.\n\n"
|
||||
"### Phase 3: Review\n"
|
||||
"- Validate alignment with user intent and constraints.\n"
|
||||
"- Clarify unresolved requirements only when necessary.\n\n"
|
||||
"### Phase 4: Final Plan\n"
|
||||
"- Output sections: Context, Changes, Critical files, Verification.\n"
|
||||
"- Prefer one recommended approach over listing many alternatives.\n\n"
|
||||
"### Phase 5: Execution Handoff\n"
|
||||
"- Ask user to switch to agent mode and confirm before execution.\n"
|
||||
"- Do not execute while still in plan mode.\n\n"
|
||||
"## Plan mode tools (lifecycle)\n"
|
||||
"- Built-in tools `enter_plan_mode_v2` and `exit_plan_mode_v2` mirror cc-mini-style "
|
||||
"Enter/Exit plan mode: they bind or release plan-mode state for this session.\n"
|
||||
"- Prefer updating the plan file in place while staying in plan mode; only call "
|
||||
"`enter_plan_mode_v2` with `force_new_plan: true` when the user explicitly wants a new plan document.\n"
|
||||
"- When the written plan is ready for review, either keep plan mode and summarize next steps for the user, "
|
||||
"or call `exit_plan_mode_v2`. Use `confirm: true` only when the user has explicitly approved executing "
|
||||
"this plan; use `confirm: false` to leave plan mode without marking the plan approved for execution.\n"
|
||||
"- If the user sends low-content prompts such as 'continue' or 'ok', do not repeat long boilerplate; "
|
||||
"revise the plan file or ask one concrete clarification."
|
||||
)
|
||||
return (
|
||||
"当前处于 plan 模式,用户暂不要求执行。\n"
|
||||
"你必须不做真实项目改动、不调用非只读工具,也不要声称已经实现完成。\n"
|
||||
"## 执行纪律(必须遵守)\n"
|
||||
"- 禁止用「我现在写脚本并执行」「开始迁移/复制」「让我跑一下」等表述,假装本回合会动磁盘或执行命令。\n"
|
||||
"- 若用户需要真实执行,必须明确说明:请在界面切换到 **agent 模式**,再按产品流程确认;"
|
||||
"在 plan 模式下你无法代为执行。\n"
|
||||
"- 禁止多轮重复同一套「接下来我将……」的独白;用户只说「继续/好的」时,应小幅改计划文件或只提一个具体问题,"
|
||||
"不要复读长模板。\n"
|
||||
"- 不要用闲聊式「你同意这个计划吗?」代替产品要求的 **切 agent + 确认**;应提示用户按界面切换到 agent 模式后再确认执行。\n\n"
|
||||
"## 计划文件信息\n"
|
||||
f"{file_info}\n"
|
||||
"在 plan 模式下,只允许围绕计划文件进行编辑。\n\n"
|
||||
"## 计划工作流\n"
|
||||
"### 阶段1:理解问题\n"
|
||||
"- 先理解需求并检查相关代码路径。\n"
|
||||
"- 优先复用现有函数、工具和既有模式。\n\n"
|
||||
"### 阶段2:方案设计\n"
|
||||
"- 给出可落地的实现方案,并说明关键取舍。\n\n"
|
||||
"### 阶段3:对齐复核\n"
|
||||
"- 核对是否满足用户目标与约束。\n"
|
||||
"- 仅在必要时提出澄清问题。\n\n"
|
||||
"### 阶段4:最终计划\n"
|
||||
"- 输出结构:背景、改动点、关键文件、验证方式。\n"
|
||||
"- 推荐一个主方案,不要只堆备选项。\n\n"
|
||||
"### 阶段5:执行切换\n"
|
||||
"- 明确提示用户先切换到 agent 模式并确认后再执行。\n"
|
||||
"- 在 plan 模式下不要执行实现。\n\n"
|
||||
"## 计划模式工具(生命周期)\n"
|
||||
"- 内置工具 `enter_plan_mode_v2` 与 `exit_plan_mode_v2` 对应 cc-mini 风格的进入/退出计划模式,用于绑定或释放本会话的 plan 状态。\n"
|
||||
"- 优先在 plan 模式下就地更新计划文件;仅在用户明确要求新开计划文档时,才对 `enter_plan_mode_v2` 使用 `force_new_plan: true`。\n"
|
||||
"- 计划文档写完后,可继续保持 plan 模式并给用户摘要;也可调用 `exit_plan_mode_v2`。仅在用户已明确同意按该计划执行时使用 "
|
||||
"`confirm: true`;若只是结束规划、尚未批准执行,使用 `confirm: false`。\n"
|
||||
"- 若用户输入信息量低的续写(如「继续」「好的」),不要重复大段套话,应小幅修订计划文件或提出一个具体问题。"
|
||||
)
|
||||
|
||||
|
||||
def inject_plan_context(*, base_system: str, state: PlanAgentStateV2, lang: str, max_chars: int = 3000) -> str:
|
||||
plan_text = str(state.plan_content or "").strip()
|
||||
if not plan_text:
|
||||
return base_system
|
||||
if len(plan_text) > max_chars:
|
||||
plan_text = plan_text[:max_chars] + "\n...<plan_truncated>"
|
||||
is_en = str(lang or "").startswith("en")
|
||||
header = "Approved plan context:\n" if is_en else "已确认计划上下文:\n"
|
||||
return f"{header}{plan_text}\n\n{str(base_system or '').strip()}".strip()
|
||||
|
||||
|
||||
__all__ = ["build_plan_mode_prefix", "inject_plan_context"]
|
||||
|
||||
60
runtime/plan_agent_v2/state_store.py
Normal file
60
runtime/plan_agent_v2/state_store.py
Normal file
|
|
@ -0,0 +1,60 @@
|
|||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import time
|
||||
from typing import Any
|
||||
|
||||
from .models import PLAN_MODE_NORMAL, PlanAgentStateV2
|
||||
|
||||
|
||||
def _state_key(session_id: str) -> str:
|
||||
return f"AIA_PLAN_AGENT_V2_STATE:{str(session_id or '').strip()}"
|
||||
|
||||
|
||||
class PlanAgentStateStoreV2:
|
||||
def __init__(self, store: Any):
|
||||
self._store = store
|
||||
|
||||
def load(self, *, session_id: str) -> PlanAgentStateV2:
|
||||
sid = str(session_id or "").strip()
|
||||
if not sid:
|
||||
return PlanAgentStateV2()
|
||||
raw = str(self._store.get_setting(_state_key(sid)) or "").strip()
|
||||
if not raw:
|
||||
return PlanAgentStateV2()
|
||||
try:
|
||||
obj = json.loads(raw)
|
||||
except Exception:
|
||||
return PlanAgentStateV2()
|
||||
return PlanAgentStateV2.from_dict(obj if isinstance(obj, dict) else None)
|
||||
|
||||
def save(self, *, session_id: str, state: PlanAgentStateV2) -> PlanAgentStateV2:
|
||||
sid = str(session_id or "").strip()
|
||||
if not sid:
|
||||
return state
|
||||
now_ms = int(time.time() * 1000)
|
||||
next_state = PlanAgentStateV2(
|
||||
mode=state.mode,
|
||||
owner_specialist=state.owner_specialist,
|
||||
plan_id=state.plan_id,
|
||||
plan_path=state.plan_path,
|
||||
plan_content=state.plan_content,
|
||||
plan_confirmed=bool(state.plan_confirmed),
|
||||
entered_at_ms=int(state.entered_at_ms or 0),
|
||||
updated_at_ms=now_ms,
|
||||
last_user_text_norm=str(state.last_user_text_norm or "").strip().lower(),
|
||||
plan_loop_count=int(state.plan_loop_count or 0),
|
||||
)
|
||||
self._store.set_setting(_state_key(sid), json.dumps(next_state.to_dict(), ensure_ascii=False))
|
||||
return next_state
|
||||
|
||||
def reset(self, *, session_id: str) -> PlanAgentStateV2:
|
||||
sid = str(session_id or "").strip()
|
||||
if not sid:
|
||||
return PlanAgentStateV2()
|
||||
self._store.delete_setting(_state_key(sid))
|
||||
return PlanAgentStateV2(mode=PLAN_MODE_NORMAL)
|
||||
|
||||
|
||||
__all__ = ["PlanAgentStateStoreV2"]
|
||||
|
||||
35
runtime/plan_agent_v2/switch.py
Normal file
35
runtime/plan_agent_v2/switch.py
Normal file
|
|
@ -0,0 +1,35 @@
|
|||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
from typing import Any
|
||||
|
||||
|
||||
def _is_truthy(raw: str | None) -> bool:
|
||||
return str(raw or "").strip().lower() in {"1", "true", "yes", "on"}
|
||||
|
||||
|
||||
def v2_feature_enabled(*, store: Any | None = None) -> bool:
|
||||
# Default off for shadow path safety.
|
||||
raw = ""
|
||||
try:
|
||||
if store is not None:
|
||||
raw = str(store.get_setting("AIA_EXPERT_PLAN_AGENT_V2_ENABLED") or "").strip()
|
||||
except Exception:
|
||||
raw = ""
|
||||
if not raw:
|
||||
raw = str(os.getenv("AIA_EXPERT_PLAN_AGENT_V2_ENABLED") or "").strip()
|
||||
if not raw:
|
||||
return False
|
||||
return _is_truthy(raw)
|
||||
|
||||
|
||||
def should_route_to_v2(*, store: Any | None, interaction_mode: str, force_flag: bool = False) -> bool:
|
||||
if str(interaction_mode or "").strip().lower() != "expert":
|
||||
return False
|
||||
if force_flag:
|
||||
return True
|
||||
return v2_feature_enabled(store=store)
|
||||
|
||||
|
||||
__all__ = ["should_route_to_v2", "v2_feature_enabled"]
|
||||
|
||||
52
runtime/plan_agent_v2/tool_policy.py
Normal file
52
runtime/plan_agent_v2/tool_policy.py
Normal file
|
|
@ -0,0 +1,52 @@
|
|||
from __future__ import annotations
|
||||
|
||||
from typing import Iterable
|
||||
|
||||
from .models import PLAN_MODE_PLAN
|
||||
from oclaw.runtime.tools.base import ToolRegistry, ToolSpec
|
||||
|
||||
|
||||
_DEFAULT_PLAN_ALLOWLIST = frozenset(
|
||||
{
|
||||
"read_file",
|
||||
"search_files",
|
||||
"glob",
|
||||
"list_directory",
|
||||
"list_workspace_tree",
|
||||
"search_files_context",
|
||||
"system_time",
|
||||
# Plan-mode control tools are non-read-only by design, but must stay callable.
|
||||
"enter_plan_mode_v2",
|
||||
"exit_plan_mode_v2",
|
||||
}
|
||||
)
|
||||
|
||||
|
||||
def plan_mode_allowed_tool_names(extra_allowed: Iterable[str] | None = None) -> set[str]:
|
||||
out = set(_DEFAULT_PLAN_ALLOWLIST)
|
||||
for x in (extra_allowed or []):
|
||||
n = str(x or "").strip()
|
||||
if n:
|
||||
out.add(n)
|
||||
return out
|
||||
|
||||
|
||||
def filter_tools_for_mode(
|
||||
*,
|
||||
registry: ToolRegistry,
|
||||
mode: str,
|
||||
extra_allowed: Iterable[str] | None = None,
|
||||
) -> list[ToolSpec]:
|
||||
tools = list(registry.list())
|
||||
if str(mode or "").strip().lower() != PLAN_MODE_PLAN:
|
||||
return tools
|
||||
allow = plan_mode_allowed_tool_names(extra_allowed=extra_allowed)
|
||||
out: list[ToolSpec] = []
|
||||
for t in tools:
|
||||
if t.name in allow or bool(t.is_read_only()):
|
||||
out.append(t)
|
||||
return out
|
||||
|
||||
|
||||
__all__ = ["filter_tools_for_mode", "plan_mode_allowed_tool_names"]
|
||||
|
||||
161
runtime/plan_agent_v2/tool_specs.py
Normal file
161
runtime/plan_agent_v2/tool_specs.py
Normal file
|
|
@ -0,0 +1,161 @@
|
|||
from __future__ import annotations
|
||||
|
||||
from typing import Any
|
||||
|
||||
from .manager import PlanModeManagerV2
|
||||
from .models import PLAN_MODE_PLAN
|
||||
from .trace import emit_plan_agent_v2_trace
|
||||
from oclaw.runtime.tools.base import ToolSpec
|
||||
|
||||
DEFAULT_SESSION_KEY = "AIA_PLAN_AGENT_V2_DEFAULT_SESSION_ID"
|
||||
|
||||
|
||||
def _emit_tool_trace(
|
||||
*,
|
||||
store: Any,
|
||||
session_id: str,
|
||||
args: dict[str, Any],
|
||||
event_type: str,
|
||||
payload: dict[str, Any] | None = None,
|
||||
) -> None:
|
||||
trace_id = str(args.get("trace_id") or "").strip()
|
||||
if not trace_id:
|
||||
return
|
||||
parent_raw = str(args.get("parent_span_id") or "").strip()
|
||||
emit_plan_agent_v2_trace(
|
||||
store=store,
|
||||
session_id=session_id,
|
||||
trace_id=trace_id,
|
||||
parent_span_id=parent_raw or None,
|
||||
event_type=event_type,
|
||||
payload=payload,
|
||||
)
|
||||
|
||||
|
||||
def _resolve_session_id(*, store: Any, args: dict[str, Any]) -> str:
|
||||
sid = str(args.get("session_id") or "").strip()
|
||||
if sid:
|
||||
return sid
|
||||
try:
|
||||
return str(store.get_setting(DEFAULT_SESSION_KEY) or "").strip()
|
||||
except Exception:
|
||||
return ""
|
||||
|
||||
|
||||
def enter_plan_mode_v2_tool(*, store: Any) -> ToolSpec:
|
||||
mgr = PlanModeManagerV2(store=store)
|
||||
|
||||
def _handler(args: dict[str, Any]) -> dict[str, Any]:
|
||||
session_id = _resolve_session_id(store=store, args=args)
|
||||
if not session_id:
|
||||
return {"ok": False, "error_code": "session_id_required", "error": "session_id_required"}
|
||||
specialist = str(args.get("owner_specialist") or "generalist").strip().lower() or "generalist"
|
||||
force_new = bool(args.get("force_new_plan"))
|
||||
st = mgr.enter(session_id=session_id, owner_specialist=specialist, force_new_plan=force_new)
|
||||
_emit_tool_trace(
|
||||
store=store,
|
||||
session_id=session_id,
|
||||
args=args,
|
||||
event_type="plan_mode_tool_enter",
|
||||
payload={
|
||||
"tool": "enter_plan_mode_v2",
|
||||
"owner_specialist": specialist,
|
||||
"force_new_plan": force_new,
|
||||
"plan_id": st.plan_id,
|
||||
},
|
||||
)
|
||||
return {"ok": True, "state": st.to_dict()}
|
||||
|
||||
return ToolSpec(
|
||||
name="enter_plan_mode_v2",
|
||||
description=(
|
||||
"Enter plan mode for this session (shadow v2), cc-mini-style: binds a dedicated plan file path. "
|
||||
"Use when the user wants structured planning. Set force_new_plan=true only when starting a brand-new "
|
||||
"plan document; otherwise reuse the existing plan when possible. Optional trace_id / parent_span_id "
|
||||
"attach observability to the current trace."
|
||||
),
|
||||
parameters={
|
||||
"type": "object",
|
||||
"properties": {
|
||||
"session_id": {"type": "string"},
|
||||
"owner_specialist": {"type": "string"},
|
||||
"force_new_plan": {"type": "boolean", "default": False},
|
||||
"trace_id": {"type": "string"},
|
||||
"parent_span_id": {"type": "string"},
|
||||
},
|
||||
"required": [],
|
||||
"additionalProperties": False,
|
||||
},
|
||||
handler=_handler,
|
||||
tags=frozenset({"plan_mode", "shadow_v2", "read"}),
|
||||
read_only=True,
|
||||
risk_level="low",
|
||||
)
|
||||
|
||||
|
||||
def exit_plan_mode_v2_tool(*, store: Any) -> ToolSpec:
|
||||
mgr = PlanModeManagerV2(store=store)
|
||||
|
||||
def _handler(args: dict[str, Any]) -> dict[str, Any]:
|
||||
session_id = _resolve_session_id(store=store, args=args)
|
||||
if not session_id:
|
||||
return {"ok": False, "error_code": "session_id_required", "error": "session_id_required"}
|
||||
confirm = bool(args.get("confirm"))
|
||||
st = mgr.confirm(session_id=session_id) if confirm else mgr.exit_without_confirm(session_id=session_id)
|
||||
_emit_tool_trace(
|
||||
store=store,
|
||||
session_id=session_id,
|
||||
args=args,
|
||||
event_type="plan_mode_tool_exit",
|
||||
payload={
|
||||
"tool": "exit_plan_mode_v2",
|
||||
"confirmed": bool(confirm),
|
||||
"plan_id": st.plan_id,
|
||||
"plan_confirmed": bool(st.plan_confirmed),
|
||||
},
|
||||
)
|
||||
return {"ok": True, "confirmed": bool(confirm), "state": st.to_dict()}
|
||||
|
||||
return ToolSpec(
|
||||
name="exit_plan_mode_v2",
|
||||
description=(
|
||||
"Exit plan mode for this session (shadow v2). confirm=true marks the plan as approved for execution "
|
||||
"(same intent as the user confirming in agent mode). confirm=false leaves plan mode without approving "
|
||||
"execution—use when ending planning without a run approval yet. Optional trace_id / parent_span_id for "
|
||||
"trace correlation."
|
||||
),
|
||||
parameters={
|
||||
"type": "object",
|
||||
"properties": {
|
||||
"session_id": {"type": "string"},
|
||||
"confirm": {"type": "boolean", "default": False},
|
||||
"trace_id": {"type": "string"},
|
||||
"parent_span_id": {"type": "string"},
|
||||
},
|
||||
"required": [],
|
||||
"additionalProperties": False,
|
||||
},
|
||||
handler=_handler,
|
||||
tags=frozenset({"plan_mode", "shadow_v2", "write"}),
|
||||
read_only=False,
|
||||
risk_level="high",
|
||||
)
|
||||
|
||||
|
||||
def materialize_plan_mode_v2_tools(*, store: Any) -> list[ToolSpec]:
|
||||
return [enter_plan_mode_v2_tool(store=store), exit_plan_mode_v2_tool(store=store)]
|
||||
|
||||
|
||||
def is_plan_mode_v2_active(*, store: Any, session_id: str) -> bool:
|
||||
mgr = PlanModeManagerV2(store=store)
|
||||
return mgr.load_state(session_id=session_id).mode == PLAN_MODE_PLAN
|
||||
|
||||
|
||||
__all__ = [
|
||||
"enter_plan_mode_v2_tool",
|
||||
"exit_plan_mode_v2_tool",
|
||||
"materialize_plan_mode_v2_tools",
|
||||
"is_plan_mode_v2_active",
|
||||
"DEFAULT_SESSION_KEY",
|
||||
]
|
||||
|
||||
37
runtime/plan_agent_v2/trace.py
Normal file
37
runtime/plan_agent_v2/trace.py
Normal file
|
|
@ -0,0 +1,37 @@
|
|||
from __future__ import annotations
|
||||
|
||||
import time
|
||||
from typing import Any
|
||||
|
||||
|
||||
def emit_plan_agent_v2_trace(
|
||||
*,
|
||||
store: Any,
|
||||
session_id: str,
|
||||
trace_id: str | None,
|
||||
parent_span_id: str | None,
|
||||
event_type: str,
|
||||
payload: dict[str, Any] | None = None,
|
||||
) -> None:
|
||||
if not str(trace_id or "").strip():
|
||||
return
|
||||
merged = dict(payload or {})
|
||||
merged.setdefault("pipeline", "plan_agent_v2")
|
||||
merged.setdefault("ts_ms", int(time.time() * 1000))
|
||||
try:
|
||||
from oclaw.runtime.orchestration.trace import new_span_id
|
||||
|
||||
store.add_trace_event(
|
||||
session_id=str(session_id or ""),
|
||||
trace_id=str(trace_id),
|
||||
span_id=new_span_id(),
|
||||
parent_span_id=parent_span_id,
|
||||
event_type=str(event_type or "plan_agent_v2"),
|
||||
payload=merged,
|
||||
)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
__all__ = ["emit_plan_agent_v2_trace"]
|
||||
|
||||
2
runtime/plan_agent_v2_adapter.py
Normal file
2
runtime/plan_agent_v2_adapter.py
Normal file
|
|
@ -0,0 +1,2 @@
|
|||
from oclaw.runtime.plan_agent_v2.adapter import * # noqa: F403
|
||||
|
||||
2
runtime/plan_agent_v2_compat.py
Normal file
2
runtime/plan_agent_v2_compat.py
Normal file
|
|
@ -0,0 +1,2 @@
|
|||
from oclaw.runtime.plan_agent_v2.compat import * # noqa: F403
|
||||
|
||||
2
runtime/plan_agent_v2_gateway_adapter.py
Normal file
2
runtime/plan_agent_v2_gateway_adapter.py
Normal file
|
|
@ -0,0 +1,2 @@
|
|||
from oclaw.runtime.plan_agent_v2.gateway_adapter import * # noqa: F403
|
||||
|
||||
86
runtime/plan_agent_v2_gateway_cutover.py
Normal file
86
runtime/plan_agent_v2_gateway_cutover.py
Normal file
|
|
@ -0,0 +1,86 @@
|
|||
from __future__ import annotations
|
||||
|
||||
import time
|
||||
import uuid
|
||||
from dataclasses import dataclass
|
||||
from typing import Any
|
||||
|
||||
from oclaw.runtime.gateway import OclawGatewayResult
|
||||
from oclaw.runtime.plan_agent_v2 import (
|
||||
build_shadow_gateway_result,
|
||||
evaluate_gateway_expert_turn_shadow,
|
||||
)
|
||||
from oclaw.runtime.types import StandardMessage
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class GatewayCutoverDraftOutput:
|
||||
handled: bool
|
||||
result: OclawGatewayResult | None
|
||||
system_prompt_override: str = ""
|
||||
decision_action: str = ""
|
||||
|
||||
|
||||
def maybe_handle_expert_turn_v2_draft(
|
||||
*,
|
||||
store: Any,
|
||||
msg: StandardMessage,
|
||||
lang: str,
|
||||
interaction_mode: str,
|
||||
requested_specialist: str,
|
||||
base_system_prompt: str,
|
||||
force_flag: bool = False,
|
||||
) -> GatewayCutoverDraftOutput:
|
||||
"""Draft-only helper for future gateway cutover.
|
||||
|
||||
Important:
|
||||
- This module is intentionally NOT wired into `runtime/gateway.py`.
|
||||
- It documents and validates the minimal cutover behavior in isolation.
|
||||
"""
|
||||
t0 = time.perf_counter()
|
||||
trace_id = str(uuid.uuid4())
|
||||
run_id = str(uuid.uuid4())
|
||||
|
||||
shadow = evaluate_gateway_expert_turn_shadow(
|
||||
store=store,
|
||||
msg=msg,
|
||||
lang=lang,
|
||||
interaction_mode=interaction_mode,
|
||||
requested_specialist=requested_specialist,
|
||||
base_system_prompt=base_system_prompt,
|
||||
force_flag=force_flag,
|
||||
trace_id=trace_id,
|
||||
parent_span_id=None,
|
||||
)
|
||||
if not shadow.used_v2 or shadow.decision is None:
|
||||
return GatewayCutoverDraftOutput(handled=False, result=None)
|
||||
|
||||
action = str(shadow.decision.action or "")
|
||||
elapsed_ms = int((time.perf_counter() - t0) * 1000)
|
||||
if action in {"enter_plan", "stay_plan"}:
|
||||
row = build_shadow_gateway_result(
|
||||
decision=shadow.decision,
|
||||
run_id=run_id,
|
||||
trace_id=trace_id,
|
||||
elapsed_ms=elapsed_ms,
|
||||
requested_specialist=requested_specialist,
|
||||
)
|
||||
result = OclawGatewayResult(**row)
|
||||
return GatewayCutoverDraftOutput(
|
||||
handled=True,
|
||||
result=result,
|
||||
decision_action=action,
|
||||
system_prompt_override="",
|
||||
)
|
||||
|
||||
# run_agent: draft suggests continuing legacy execution with injected prompt.
|
||||
return GatewayCutoverDraftOutput(
|
||||
handled=False,
|
||||
result=None,
|
||||
decision_action=action,
|
||||
system_prompt_override=str(shadow.decision.system_prompt_override or ""),
|
||||
)
|
||||
|
||||
|
||||
__all__ = ["GatewayCutoverDraftOutput", "maybe_handle_expert_turn_v2_draft"]
|
||||
|
||||
2
runtime/plan_agent_v2_manager.py
Normal file
2
runtime/plan_agent_v2_manager.py
Normal file
|
|
@ -0,0 +1,2 @@
|
|||
from oclaw.runtime.plan_agent_v2.manager import * # noqa: F403
|
||||
|
||||
2
runtime/plan_agent_v2_models.py
Normal file
2
runtime/plan_agent_v2_models.py
Normal file
|
|
@ -0,0 +1,2 @@
|
|||
from oclaw.runtime.plan_agent_v2.models import * # noqa: F403
|
||||
|
||||
2
runtime/plan_agent_v2_prompt_injector.py
Normal file
2
runtime/plan_agent_v2_prompt_injector.py
Normal file
|
|
@ -0,0 +1,2 @@
|
|||
from oclaw.runtime.plan_agent_v2.prompt_injector import * # noqa: F403
|
||||
|
||||
2
runtime/plan_agent_v2_state_store.py
Normal file
2
runtime/plan_agent_v2_state_store.py
Normal file
|
|
@ -0,0 +1,2 @@
|
|||
from oclaw.runtime.plan_agent_v2.state_store import * # noqa: F403
|
||||
|
||||
2
runtime/plan_agent_v2_switch.py
Normal file
2
runtime/plan_agent_v2_switch.py
Normal file
|
|
@ -0,0 +1,2 @@
|
|||
from oclaw.runtime.plan_agent_v2.switch import * # noqa: F403
|
||||
|
||||
2
runtime/plan_agent_v2_tool_policy.py
Normal file
2
runtime/plan_agent_v2_tool_policy.py
Normal file
|
|
@ -0,0 +1,2 @@
|
|||
from oclaw.runtime.plan_agent_v2.tool_policy import * # noqa: F403
|
||||
|
||||
2
runtime/plan_agent_v2_tool_specs.py
Normal file
2
runtime/plan_agent_v2_tool_specs.py
Normal file
|
|
@ -0,0 +1,2 @@
|
|||
from oclaw.runtime.plan_agent_v2.tool_specs import * # noqa: F403
|
||||
|
||||
2
runtime/plan_agent_v2_trace.py
Normal file
2
runtime/plan_agent_v2_trace.py
Normal file
|
|
@ -0,0 +1,2 @@
|
|||
from oclaw.runtime.plan_agent_v2.trace import * # noqa: F403
|
||||
|
||||
Loading…
Add table
Add a link
Reference in a new issue