From 86cd52bf2114391631b74ec4c86462987ce3ba90 Mon Sep 17 00:00:00 2001 From: Haitao Pan Date: Sat, 6 Jun 2026 11:26:22 +0800 Subject: [PATCH] Fix OpenClaw smoke recovery on main --- internal/acp/orchestrator.go | 57 +++++++++++- internal/acp/routing_test.go | 61 +++++++++++++ .../validate-openclaw-session.sh | 86 ++++++++++++++++--- 3 files changed, 191 insertions(+), 13 deletions(-) diff --git a/internal/acp/orchestrator.go b/internal/acp/orchestrator.go index 50d6608..3f043e7 100644 --- a/internal/acp/orchestrator.go +++ b/internal/acp/orchestrator.go @@ -435,8 +435,8 @@ func (o *SessionOrchestrator) startOpenClawGatewayTask( "progress": running["progress"], })) } - return running, nil - } + return running, nil +} func openClawGatewayCompletedResultUpdate(sessionID string, threadID string, turnID string, result map[string]any) map[string]any { success := true @@ -579,6 +579,17 @@ func (o *SessionOrchestrator) openClawArtifactPrepare( notify, ) if !prepareResult.OK { + if openClawPrepareUnsupported(prepareResult.Error) { + prepared := openClawLegacyPreparedArtifactScope(params, sessionKey, runID) + log.Printf( + "level=warn component=openclaw_gateway event=session_prepare_legacy_fallback provider=%q sessionId=%q runId=%q artifactScope=%q", + gatewayProvider, + sessionKey, + runID, + prepared.ArtifactScope, + ) + return prepared, nil + } return nil, gatewayRPCError(prepareResult.Error, "openclaw artifact prepare failed") } prepared := openClawPreparedArtifactScopeFromPayload(shared.AsMap(prepareResult.Payload)) @@ -588,6 +599,48 @@ func (o *SessionOrchestrator) openClawArtifactPrepare( return prepared, nil } +func openClawPrepareUnsupported(errorPayload map[string]any) bool { + code := strings.ToUpper(strings.TrimSpace(shared.StringArg(errorPayload, "code", ""))) + message := strings.ToLower(strings.TrimSpace(shared.StringArg(errorPayload, "message", ""))) + if !strings.Contains(message, "xworkmate.session.prepare") { + return false + } + return code == "INVALID_REQUEST" || + code == "METHOD_NOT_FOUND" || + code == "UNKNOWN_METHOD" || + strings.Contains(message, "unknown method") || + strings.Contains(message, "method not found") +} + +func openClawLegacyPreparedArtifactScope(params map[string]any, sessionKey string, runID string) *openClawPreparedArtifactScope { + sessionKey = strings.TrimSpace(sessionKey) + runID = strings.TrimSpace(runID) + artifactScope := "tasks/" + sessionKey + "/" + runID + workspaceRoot := openClawLegacyArtifactWorkspaceRoot(params) + return &openClawPreparedArtifactScope{ + RemoteWorkingDirectory: workspaceRoot, + RemoteWorkspaceRefKind: "remotePath", + ArtifactScope: artifactScope, + ArtifactDirectory: filepath.Join(workspaceRoot, filepath.FromSlash(artifactScope)), + RelativeArtifactDirectory: artifactScope, + ScopeKind: "task", + } +} + +func openClawLegacyArtifactWorkspaceRoot(params map[string]any) string { + for _, key := range []string{"remoteWorkingDirectoryHint", "remoteWorkingDirectory"} { + value := strings.TrimSpace(shared.StringArg(params, key, "")) + if value == "" { + continue + } + cleaned := filepath.Clean(value) + if strings.HasPrefix(cleaned, "/home/ubuntu/.openclaw/workspace") { + return strings.TrimRight(cleaned, string(os.PathSeparator)) + } + } + return "/home/ubuntu/.openclaw/workspace" +} + func openClawSessionPrepareParams(params map[string]any, openClawSessionKey string, runID string, artifactContract openClawArtifactContract) map[string]any { appThreadKey := openClawAppThreadKey(params) result := map[string]any{ diff --git a/internal/acp/routing_test.go b/internal/acp/routing_test.go index 8546c43..60677b5 100644 --- a/internal/acp/routing_test.go +++ b/internal/acp/routing_test.go @@ -742,6 +742,54 @@ func TestGatewayRequestSkillsStatusAutoConnectsOpenClaw(t *testing.T) { } } +func TestExecuteSessionTaskGatewayFallsBackWhenPrepareUnsupported(t *testing.T) { + gateway := newAcpFakeOpenClawGateway(t) + gateway.unsupportedSessionPrepare.Store(true) + defer gateway.Close() + + t.Setenv("GATEWAY_RPC_URL", gateway.URL()) + t.Setenv("BRIDGE_AUTH_TOKEN", "bridge-token") + + server := NewServer() + response, rpcErr := server.executeSessionTask(task{ + req: shared.RPCRequest{ + Method: "session.start", + Params: map[string]any{ + "sessionId": "session-openclaw-legacy-prepare", + "threadId": "thread-openclaw-legacy-prepare", + "taskPrompt": "say pong", + "workingDirectory": t.TempDir(), + "routing": map[string]any{ + "routingMode": "explicit", + "explicitExecutionTarget": "gateway", + "preferredGatewayProviderId": "openclaw", + }, + }, + }, + }) + if rpcErr != nil { + t.Fatalf("expected legacy prepare fallback response, got rpc error: %#v", rpcErr) + } + if got := response["output"]; got != "gateway pong" { + t.Fatalf("expected gateway pong output after legacy prepare fallback, got %#v", response) + } + if got := gateway.Methods(); !sameMethods(got, []string{"connect", "xworkmate.session.prepare", "chat.send", "xworkmate.tasks.get"}) { + t.Fatalf("expected legacy prepare attempt to continue through chat/send and native task lookup, got %#v", got) + } + chatParams := gateway.LastChatSendParams() + receipt := strings.TrimSpace(shared.StringArg(chatParams, "systemProvenanceReceipt", "")) + sessionKey := shared.StringArg(chatParams, "sessionKey", "") + runID := shared.StringArg(chatParams, "idempotencyKey", "") + for _, expected := range []string{ + "artifactDirectory: /home/ubuntu/.openclaw/workspace/tasks/" + sessionKey + "/" + runID, + "artifactScope: tasks/" + sessionKey + "/" + runID, + } { + if !strings.Contains(receipt, expected) { + t.Fatalf("expected fallback provenance receipt to include %q, got %q", expected, receipt) + } + } +} + func TestExecuteSessionTaskGatewayNoDisplayableOutputFails(t *testing.T) { gateway := newAcpFakeOpenClawGateway(t) defer gateway.Close() @@ -2616,6 +2664,7 @@ type acpFakeOpenClawGateway struct { artifactWorkspaceRoot string alternateRunID string alternateSessionKey string + unsupportedSessionPrepare atomic.Bool } func newAcpFakeOpenClawGateway(t *testing.T) *acpFakeOpenClawGateway { @@ -2754,6 +2803,18 @@ func newAcpFakeOpenClawGateway(t *testing.T) *acpFakeOpenClawGateway { fake.artifactPrepareCount.Add(1) params := shared.AsMap(frame["params"]) fake.lastArtifactPrepareParams.Store(params) + if fake.unsupportedSessionPrepare.Load() { + _ = conn.WriteJSON(map[string]any{ + "type": "res", + "id": id, + "ok": false, + "error": map[string]any{ + "code": "INVALID_REQUEST", + "message": "unknown method: xworkmate.session.prepare", + }, + }) + continue + } runID := strings.TrimSpace(shared.StringArg(params, "runId", "fake-run")) sessionKey := strings.TrimSpace(shared.StringArg(params, "openclawSessionKey", "main")) artifactScope := "tasks/" + sessionKey + "/" + runID diff --git a/scripts/github-actions/validate-openclaw-session.sh b/scripts/github-actions/validate-openclaw-session.sh index 6c9f100..3b250c4 100755 --- a/scripts/github-actions/validate-openclaw-session.sh +++ b/scripts/github-actions/validate-openclaw-session.sh @@ -94,7 +94,8 @@ def terminal_result(payload): if not isinstance(payload, dict): return {} nested = payload.get("result") - if isinstance(nested, dict) and str(payload.get("status", "")).lower() in { + status = str(payload.get("status", "")).lower() + if isinstance(nested, dict) and nested and status in { "completed", "failed", "cancelled", @@ -136,6 +137,60 @@ def require_nonempty(payload, key): raise SystemExit(f"OpenClaw smoke result missing {key}: {json.dumps(payload, ensure_ascii=False, sort_keys=True)[:1000]}") +def first_nonempty(payload, *keys): + if not isinstance(payload, dict): + return "" + for key in keys: + value = payload.get(key) + if isinstance(value, str) and value.strip(): + return value.strip() + return "" + + +def openclaw_session_key_for_app_thread(app_thread_key): + app_thread_key = str(app_thread_key or "").strip() + if not app_thread_key: + app_thread_key = "main" + return "agent:main:" + app_thread_key + + +def task_handle_from_payload(payload): + if not isinstance(payload, dict): + return {} + candidates = [] + for key in ("result", "payload", "params"): + if isinstance(payload.get(key), dict): + candidates.append(payload[key]) + candidates.append(payload) + for candidate in candidates: + if not isinstance(candidate, dict): + continue + session_id = first_nonempty(candidate, "sessionId") + thread_id = first_nonempty(candidate, "threadId") + turn_id = first_nonempty(candidate, "turnId") + run_id = first_nonempty(candidate, "runId") + app_thread_key = first_nonempty(candidate, "appThreadKey") + openclaw_session_key = first_nonempty(candidate, "openclawSessionKey") + if (app_thread_key and openclaw_session_key and run_id) or (session_id and thread_id and (turn_id or run_id)): + return { + "sessionId": session_id, + "threadId": thread_id, + "turnId": turn_id, + "runId": run_id or turn_id, + "appThreadKey": app_thread_key or session_id or thread_id, + "openclawSessionKey": openclaw_session_key or openclaw_session_key_for_app_thread(app_thread_key or session_id or thread_id), + } + return {} + + +def find_task_handle(payloads, final): + for payload in reversed(payloads): + handle = task_handle_from_payload(payload) + if handle: + return handle + return task_handle_from_payload(final) + + def is_valid_no_displayable_contract(payload): if not isinstance(payload, dict): return False @@ -143,8 +198,12 @@ def is_valid_no_displayable_contract(payload): return False if payload.get("resolvedGatewayProviderId") != "openclaw": return False - for key in ("sessionId", "threadId", "runId", "openclawSessionKey", "artifactScope"): + for key in ("sessionId", "threadId", "runId", "artifactScope"): require_nonempty(payload, key) + if not first_nonempty(payload, "openclawSessionKey", "sessionKey"): + artifact_scope = first_nonempty(payload, "artifactScope") + if not artifact_scope.startswith("tasks/"): + raise SystemExit(f"OpenClaw smoke result missing session scope: {json.dumps(payload, ensure_ascii=False, sort_keys=True)[:1000]}") return True @@ -158,17 +217,20 @@ if not payloads or payloads[-1].get("done") is not True: raise SystemExit("missing SSE done marker") result = terminal_result(final.get("result") or final.get("payload") or {}) -if result.get("status") == "running": - run_id = result.get("runId") - app_thread_key = result.get("appThreadKey") - openclaw_session_key = result.get("openclawSessionKey") +handle = result +if result.get("status") != "running": + handle = find_task_handle(payloads, final) +if handle.get("status") == "running" or (not result and handle): + run_id = first_nonempty(handle, "runId", "turnId") + app_thread_key = first_nonempty(handle, "appThreadKey", "sessionId", "threadId") + openclaw_session_key = first_nonempty(handle, "openclawSessionKey") or openclaw_session_key_for_app_thread(app_thread_key) if not app_thread_key: - raise SystemExit(f"OpenClaw smoke running handle missing appThreadKey: {json.dumps(result, ensure_ascii=False, sort_keys=True)[:1000]}") + raise SystemExit(f"OpenClaw smoke running handle missing appThreadKey: {json.dumps(handle, ensure_ascii=False, sort_keys=True)[:1000]}") if not openclaw_session_key: - raise SystemExit(f"OpenClaw smoke running handle missing openclawSessionKey: {json.dumps(result, ensure_ascii=False, sort_keys=True)[:1000]}") + raise SystemExit(f"OpenClaw smoke running handle missing openclawSessionKey: {json.dumps(handle, ensure_ascii=False, sort_keys=True)[:1000]}") if not run_id: - raise SystemExit(f"OpenClaw smoke running handle missing runId: {json.dumps(result, ensure_ascii=False, sort_keys=True)[:1000]}") + raise SystemExit(f"OpenClaw smoke running handle missing runId: {json.dumps(handle, ensure_ascii=False, sort_keys=True)[:1000]}") deadline = time.time() + poll_timeout while time.time() < deadline: @@ -197,7 +259,8 @@ if result.get("status") == "running": poll_result = resp_data.get("result") or {} status = poll_result.get("status") if status in ("completed", "failed", "cancelled"): - result = terminal_result(poll_result) + terminal = terminal_result(poll_result) + result = terminal if terminal else poll_result final["result"] = poll_result break except Exception as exc: @@ -232,7 +295,8 @@ if "pong" not in output_text.lower(): print("OpenClaw smoke OK: session contract completed without displayable output") sys.exit(0) result_preview = json.dumps(result, ensure_ascii=False, sort_keys=True)[:1000] - raise SystemExit(f"OpenClaw smoke did not return pong: {output_text[:500]}\nresult preview: {result_preview}") + payload_preview = json.dumps(payloads[:6], ensure_ascii=False, sort_keys=True)[:1500] + raise SystemExit(f"OpenClaw smoke did not return pong: {output_text[:500]}\nresult preview: {result_preview}\nSSE preview: {payload_preview}") print("OpenClaw smoke OK: pong received from session contract") PY