oclaw/openclaw_runtime/worker.py
oliver ba3836f00f 初始化:独立 oclaw 仓库首提交
- 在 oclaw/ 下重新初始化 Git 仓库
- 补齐子仓库 .gitignore,避免提交本地运行态数据(_local、node_modules、logs 等)
- 提交当前工程代码与配置

Made-with: Cursor
2026-04-24 22:31:22 +08:00

236 lines
10 KiB
Python

from __future__ import annotations
import json
import threading
import time
import uuid
from typing import Any
from oclaw.agents.factory import build_gateway_executor
from oclaw.openclaw_runtime.agent_core_run import AgentCoreRunInput, run_agent_core
from oclaw.openclaw_runtime.memory_stage import build_memory_context
from oclaw.openclaw_runtime.relay_pointer import build_acp_relay_result, validate_relay_share_envelope
from oclaw.openclaw_runtime.types import StandardMessage
_LOCK = threading.Lock()
_THREAD: threading.Thread | None = None
def ensure_worker_started(*, store: Any, poll_interval_s: float = 1.0) -> str:
global _THREAD
with _LOCK:
if _THREAD and _THREAD.is_alive():
return _THREAD.name
wid = f"openclaw-worker-{uuid.uuid4().hex[:8]}"
t = threading.Thread(
target=_worker_loop,
name=wid,
kwargs={"store": store, "worker_id": wid, "poll_interval_s": max(0.3, float(poll_interval_s or 1.0))},
daemon=True,
)
t.start()
_THREAD = t
return wid
def _worker_loop(*, store: Any, worker_id: str, poll_interval_s: float) -> None:
while True:
task = None
try:
task = store.openclaw_task_claim(worker_id=worker_id, lease_seconds=90)
except Exception:
task = None
if not task:
time.sleep(poll_interval_s)
continue
payload: dict[str, Any] = {}
try:
payload = json.loads(task.payload or "{}")
if not isinstance(payload, dict):
payload = {}
except Exception:
payload = {}
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:
store.add_trace_event(
session_id=session_id,
trace_id=trace_id,
span_id=str(uuid.uuid4()),
parent_span_id=None,
event_type="task_claimed",
payload={"task_id": task.id, "worker_id": worker_id, "attempt_count": int(task.attempt_count or 0)},
)
except Exception:
pass
try:
lang = str(payload.get("lang") or "zh")
session_id = str(payload.get("session_id") or task.session_id or "")
user_text = str(payload.get("text") or "")
attachments = payload.get("attachments") or []
metadata = payload.get("metadata") if isinstance(payload.get("metadata"), dict) else {}
relay_share_envelope = payload.get("relay_share_envelope") if isinstance(payload.get("relay_share_envelope"), dict) else {}
if relay_share_envelope and "relay_share_envelope" not in metadata:
metadata["relay_share_envelope"] = relay_share_envelope
acp_parent_run_id = str(payload.get("acp_parent_run_id") or "").strip()
acp_child_run_id = str(payload.get("acp_child_run_id") or "").strip()
if acp_parent_run_id or acp_child_run_id:
ok_env, env_err, env_norm = validate_relay_share_envelope(relay_share_envelope)
if not ok_env:
fail_result = {
"ok": False,
"error_code": str(env_err or "relay_envelope_invalid"),
"retryable": False,
"acp_parent_run_id": acp_parent_run_id,
"acp_child_run_id": acp_child_run_id,
}
store.openclaw_task_fail(task_id=task.id, error=str(env_err or "relay_envelope_invalid"), result=fail_result)
if trace_id:
store.add_trace_event(
session_id=session_id,
trace_id=trace_id,
span_id=str(uuid.uuid4()),
parent_span_id=None,
event_type="task_failed",
payload={
"task_id": task.id,
"ok": False,
"error": str(env_err or "relay_envelope_invalid"),
"error_code": str(env_err or "relay_envelope_invalid"),
"retryable": False,
"acp_parent_run_id": acp_parent_run_id,
"acp_child_run_id": acp_child_run_id,
},
)
continue
relay_share_envelope = env_norm
metadata["relay_share_envelope"] = relay_share_envelope
tenant_id = str(payload.get("tenant_id") or "")
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 "")
memory_ctx = build_memory_context(
store=store,
session_id=session_id,
tenant_id=tenant_id,
user_id=user_id,
query_text=user_text,
)
executor = build_gateway_executor(
store,
lang=lang,
specialist=selected_specialist or "generalist",
profile_id=model_profile_id,
viewer_user_id=user_id or None,
viewer_username=viewer_username or None,
viewer_tenant_id=tenant_id or None,
policy_session_id=session_id or None,
path_policy_tenant_id=tenant_id or None,
path_policy_user_id=user_id or None,
)
system_prompt = ""
if hasattr(executor, "_compose_system_prompt"):
try:
system_prompt = str(executor._compose_system_prompt() or "")
except Exception:
system_prompt = ""
if not system_prompt:
system_prompt = str(getattr(executor, "system_prompt", "") or "")
max_messages = int(store.get_setting("AIA_TURN_MAX_CONTEXT_MESSAGES") or 80)
max_tool_rounds = int(store.get_setting("AIA_TURN_MAX_TOOL_ROUNDS") or 8)
max_tool_workers = int(store.get_setting("AIA_TURN_MAX_TOOL_WORKERS") or 8)
msg = StandardMessage(
session_id=session_id,
tenant_id=tenant_id,
user_id=user_id,
role=str(payload.get("role") or "member"),
channel=str(payload.get("channel") or "admin_chat"), # type: ignore[arg-type]
text=user_text,
attachments=list(attachments or []),
metadata=dict(metadata or {}),
)
run_out = run_agent_core(
store=store,
data=AgentCoreRunInput(
msg=msg,
lang=lang,
system_prompt=system_prompt,
model=executor.model,
tools=executor.tools,
trace_id=trace_id or None,
parent_span_id=None,
run_id=run_id,
max_messages=max(10, min(max_messages, 400)),
max_tool_rounds=max(1, min(max_tool_rounds, 30)),
max_tool_workers=max(1, min(max_tool_workers, 32)),
max_attempts=2,
memory_context=memory_ctx,
openclaw_task_id=str(task.id),
openclaw_worker_id=worker_id,
),
)
base_result = {
"run_id": str(run_out.run_id or ""),
"reply_text": run_out.outcome.final_text,
"tool_trace_count": len(run_out.outcome.tool_traces),
"relay_pointer_count": int(payload.get("relay_pointer_count") or 0),
"relay_envelope_present": bool(isinstance(relay_share_envelope, dict) and bool(relay_share_envelope)),
}
if acp_parent_run_id or acp_child_run_id:
base_result.update(
build_acp_relay_result(
parent_run_id=acp_parent_run_id,
child_run_id=acp_child_run_id,
relay_envelope=relay_share_envelope,
)
)
store.openclaw_task_finish(task_id=task.id, result=base_result)
if trace_id:
trace_payload = {
"task_id": task.id,
"ok": True,
"tool_trace_count": len(run_out.outcome.tool_traces),
"relay_pointer_count": int(payload.get("relay_pointer_count") or 0),
"relay_envelope_present": bool(isinstance(relay_share_envelope, dict) and bool(relay_share_envelope)),
}
if acp_parent_run_id or acp_child_run_id:
trace_payload.update(
{
"acp_parent_run_id": acp_parent_run_id,
"acp_child_run_id": acp_child_run_id,
}
)
store.add_trace_event(
session_id=session_id,
trace_id=trace_id,
span_id=str(uuid.uuid4()),
parent_span_id=None,
event_type="task_finished",
payload=trace_payload,
)
except Exception as exc:
store.openclaw_task_fail(task_id=task.id, error=str(exc), result={"ok": False})
try:
if trace_id:
store.add_trace_event(
session_id=session_id,
trace_id=trace_id,
span_id=str(uuid.uuid4()),
parent_span_id=None,
event_type="task_failed",
payload={"task_id": task.id, "ok": False, "error": str(exc)[:500]},
)
except Exception:
pass
__all__ = ["ensure_worker_started"]