diff --git a/packages/gateway-protocol/README.md b/packages/gateway-protocol/README.md index cd7f302e0141..11cf1ee7d2c9 100644 --- a/packages/gateway-protocol/README.md +++ b/packages/gateway-protocol/README.md @@ -81,6 +81,14 @@ console.log(frame.id, frame.method); the validator for the selected method's `params`; the root entry point exports those validators as `validate*Params` functions. +External lifecycle controllers can validate suspension responses with +`validateGatewaySuspendPrepareResult` and `validateGatewaySuspendStatusResult`. +These use the canonical result schemas without changing the payload. Preserve +optional `writeCustody`: absence means unknown custody, not an empty list. Phase +names are open strings; consumers must not discard an unfamiliar owner phase. +Validation does not authorize a stop or replace lease, process-identity, and +readiness checks owned by the controller. + ## Guard an event without TypeBox Use the lightweight guards when code only needs safe frame discrimination. They diff --git a/packages/gateway-protocol/src/gateway-suspend.test.ts b/packages/gateway-protocol/src/gateway-suspend.test.ts index 4ef183b5a99e..a89eb39dee40 100644 --- a/packages/gateway-protocol/src/gateway-suspend.test.ts +++ b/packages/gateway-protocol/src/gateway-suspend.test.ts @@ -2,8 +2,8 @@ import { Value } from "typebox/value"; import { describe, expect, it } from "vitest"; import { GatewaySuspendBlockerSchema, - GatewaySuspendPrepareResultSchema, - GatewaySuspendStatusResultSchema, + validateGatewaySuspendPrepareResult, + validateGatewaySuspendStatusResult, validateGatewaySuspendPrepareParams, validateGatewaySuspendHandoffParams, } from "./index.js"; @@ -81,12 +81,10 @@ describe("gateway suspension protocol", () => { blockers: [{ kind: "terminal-session", count: 1, message: "1 open terminal session" }], }; - expect(Value.Check(GatewaySuspendPrepareResultSchema, draining)).toBe(true); - expect(Value.Check(GatewaySuspendPrepareResultSchema, { ...draining, unexpected: true })).toBe( - false, - ); + expect(validateGatewaySuspendPrepareResult(draining)).toBe(true); + expect(validateGatewaySuspendPrepareResult({ ...draining, unexpected: true })).toBe(false); expect( - Value.Check(GatewaySuspendPrepareResultSchema, { + validateGatewaySuspendPrepareResult({ status: "busy", reason: "active-work", retryAfterMs: 250, @@ -95,7 +93,7 @@ describe("gateway suspension protocol", () => { }), ).toBe(true); expect( - Value.Check(GatewaySuspendPrepareResultSchema, { + validateGatewaySuspendPrepareResult({ status: "ready", suspensionId: "suspension-1", expiresAtMs: 2_000, @@ -114,13 +112,90 @@ describe("gateway suspension protocol", () => { blockers: [{ kind: "terminal-persistence", count: 1, message: "1 pending terminal write" }], }; - expect(Value.Check(GatewaySuspendStatusResultSchema, draining)).toBe(true); - expect(Value.Check(GatewaySuspendStatusResultSchema, { ...draining, suspensionId: "id" })).toBe( - false, - ); - expect(Value.Check(GatewaySuspendStatusResultSchema, { status: "running" })).toBe(true); - expect( - Value.Check(GatewaySuspendStatusResultSchema, { status: "ready", expiresAtMs: 2_000 }), - ).toBe(true); + expect(validateGatewaySuspendStatusResult(draining)).toBe(true); + expect(validateGatewaySuspendStatusResult({ ...draining, suspensionId: "id" })).toBe(false); + expect(validateGatewaySuspendStatusResult({ status: "running" })).toBe(true); + expect(validateGatewaySuspendStatusResult({ status: "ready", expiresAtMs: 2_000 })).toBe(true); + }); + + it.each([ + { + name: "prepare busy", + validate: validateGatewaySuspendPrepareResult, + result: { + status: "busy", + reason: "active-work", + retryAfterMs: 250, + activeCount: 1, + blockers: [{ kind: "session-mutation", count: 1, message: "1 pending write" }], + }, + }, + { + name: "prepare draining", + validate: validateGatewaySuspendPrepareResult, + result: { + status: "draining", + suspensionId: "held-lease", + expiresAtMs: 2_000, + retryAfterMs: 250, + activeCount: 1, + blockers: [{ kind: "session-mutation", count: 1, message: "1 pending write" }], + }, + }, + { + name: "prepare ready", + validate: validateGatewaySuspendPrepareResult, + result: { + status: "ready", + suspensionId: "held-lease", + expiresAtMs: 2_000, + activeCount: 0, + blockers: [], + }, + }, + { + name: "status draining", + validate: validateGatewaySuspendStatusResult, + result: { + status: "draining", + expiresAtMs: 2_000, + retryAfterMs: 250, + activeCount: 1, + blockers: [{ kind: "session-mutation", count: 1, message: "1 pending write" }], + }, + }, + { + name: "status ready", + validate: validateGatewaySuspendStatusResult, + result: { status: "ready", expiresAtMs: 2_000 }, + }, + ])("validates $name custody without erasing unknown or held evidence", ({ validate, result }) => { + expect(validate(result)).toBe(true); + expect(result).not.toHaveProperty("writeCustody"); + for (const writeCustody of [ + [], + [{ phase: "backup", count: 0 }], + [{ phase: "migration", count: 1 }], + [{ phase: "future-owner-phase", count: 2 }], + ]) { + const response = { ...result, writeCustody }; + const original = structuredClone(response); + expect(validate(response)).toBe(true); + expect(response).toEqual(original); + } + for (const writeCustody of [ + null, + {}, + [null], + [{ count: 1 }], + [{ phase: "", count: 1 }], + [{ phase: "backup", count: -1 }], + [{ phase: "backup", count: 0.5 }], + [{ phase: "backup", count: "1" }], + [{ phase: "backup", count: 1, extra: true }], + ]) { + expect(validate({ ...result, writeCustody })).toBe(false); + expect(validate.errors).not.toBeNull(); + } }); }); diff --git a/packages/gateway-protocol/src/validator-registry.ts b/packages/gateway-protocol/src/validator-registry.ts index f3219930237c..b9cd1c34e4a1 100644 --- a/packages/gateway-protocol/src/validator-registry.ts +++ b/packages/gateway-protocol/src/validator-registry.ts @@ -91,7 +91,9 @@ export const validateWorkerLiveEventParams = compile( checkWorkerProtocolJson, ); export const validateGatewaySuspendPrepareParams = compile(S.GatewaySuspendPrepareParamsSchema); +export const validateGatewaySuspendPrepareResult = compile(S.GatewaySuspendPrepareResultSchema); export const validateGatewaySuspendStatusParams = compile(S.GatewaySuspendStatusParamsSchema); +export const validateGatewaySuspendStatusResult = compile(S.GatewaySuspendStatusResultSchema); export const validateGatewaySuspendResumeParams = compile(S.GatewaySuspendResumeParamsSchema); export const validateGatewaySuspendHandoffParams = compile(S.GatewaySuspendHandoffParamsSchema); export const validateRequestFrame = compile(S.RequestFrameSchema); diff --git a/src/infra/device-auth-store.ts b/src/infra/device-auth-store.ts index f23f4837d180..c632d0e48aca 100644 --- a/src/infra/device-auth-store.ts +++ b/src/infra/device-auth-store.ts @@ -125,18 +125,16 @@ async function executeDeviceAuth( return result; } -/** Open the shared actor during request preparation without reading or caching token facts. */ +/** Prepare the command runtime before connection work, without reading or caching token facts. */ export async function prepareDeviceAuthStore( params: DeviceAuthOperation & { readOnly?: boolean }, ): Promise { - const { context, assertActive } = captureDeviceAuthOperation(params); - assertActive(); - const prepare = async () => {}; - const options = { assertCurrent: assertActive }; - await (params.readOnly - ? runOpenClawStateWorkerOperation(context, prepare, { ...options, existingOnly: true }) - : runOpenClawStateWorkerOperation(context, prepare, options)); - assertActive(); + await executeDeviceAuth( + captureDeviceAuthOperation(params), + "deviceAuth.prepare", + undefined, + params.readOnly === true, + ); } async function readDeviceAuth( diff --git a/src/infra/device-auth-store.worker.test.ts b/src/infra/device-auth-store.worker.test.ts index 33f8d014d148..46722da7cf4e 100644 --- a/src/infra/device-auth-store.worker.test.ts +++ b/src/infra/device-auth-store.worker.test.ts @@ -80,6 +80,38 @@ it("keeps cold, warm, read-only, ordered token-data operations and cleanup off t }); }); +it.each([false, true])( + "prepares the device command runtime without reading token facts (readOnly: %s)", + async (readOnly) => { + await withOpenClawTestState({ label: "device-token-command-preparation" }, async (state) => { + const lookup = { deviceId: "synthetic-device", role: "operator", env: state.env }; + await tokens.storeDeviceAuthToken({ ...lookup, token: "synthetic-stored" }); + await closeOpenClawStateDatabaseAsync(); + const commands: string[] = []; + const runOperation = workerStore.runOpenClawStateWorkerOperation; + vi.spyOn(workerStore, "runOpenClawStateWorkerOperation").mockImplementation( + (context, operation, options) => + runOperation( + context, + (scope) => + operation({ + execute(command, commandOptions) { + commands.push(command.type); + return scope.execute(command, commandOptions); + }, + }), + options, + ), + ); + await tokens.prepareDeviceAuthStore({ env: state.env, readOnly }); + expect(commands).toEqual(["deviceAuth.prepare"]); + expect(await tokens.loadDeviceAuthTokenReadOnly(lookup)).toMatchObject({ + token: "synthetic-stored", + }); + }); + }, +); + it("rejects canceled loads, retired sources and expired schema scopes without publishing observations", async () => { await withOpenClawTestState({ label: "device-token-admission" }, async (state) => { const lookup = { deviceId: "synthetic-device", role: "operator", env: state.env }; diff --git a/src/state/openclaw-state-worker-contract.ts b/src/state/openclaw-state-worker-contract.ts index 7d00d483d99b..fe9acd66e340 100644 --- a/src/state/openclaw-state-worker-contract.ts +++ b/src/state/openclaw-state-worker-contract.ts @@ -129,6 +129,7 @@ export type OpenClawStateWorkerOperations = McpOAuthReadOperations & }; output: OpenClawStateLeaseAcquisition; }; + "deviceAuth.prepare": { input: undefined; output: void }; "deviceAuth.list": { input: { deviceId: string }; output: DeviceAuthEntry[] }; "deviceAuth.read": { input: Parameters[1] & { diff --git a/src/state/openclaw-state-worker-runtime.ts b/src/state/openclaw-state-worker-runtime.ts index 8cddbd328653..804a5aac41c4 100644 --- a/src/state/openclaw-state-worker-runtime.ts +++ b/src/state/openclaw-state-worker-runtime.ts @@ -153,6 +153,10 @@ export function executeSharedStateCommand( open: () => OpenClawStateDatabase, hasNativeDatabase: boolean, ): Operations[keyof Operations]["output"] { + // Dispatch preparation has loaded this module; do not open or observe token state. + if (command.type === "deviceAuth.prepare") { + return undefined; + } if (command.type === "mcpOAuth.read") { return readMcpOAuthStoreInDatabase(open().db, command.input); }