From fde4658bf6bbf74ddfb54a1b41ca23f403645dfa Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Wed, 29 Jul 2026 10:09:51 -0400 Subject: [PATCH] fix(ui): recover sessions and dashboards after prolonged outages (#115654) * fix(ui): recover long-running sessions and dashboards * fix(ui): expose browser-safe gateway timers * refactor(ui): keep browser timer surface minimal * refactor(gateway): extract pending request state * fix(ci): refresh browser runtime and format baseline * fix(ci): exclude generated plugin outputs from targeted lint * fix(ui): bound long-running connection and pane lifecycles * refactor(ui): split long-running lifecycle regression coverage * fix(ui): bound long-running media and lint ownership * fix(security): preserve explicit workspace disable denylist * fix(ci): preserve plugin manifest fallback lint --- .../modules/copilot-runtime.js | 2 +- packages/gateway-client/src/browser.ts | 2 +- .../src/client.watchdog.test.ts | 91 +++++++++ .../gateway-client/src/event-listeners.ts | 24 +++ .../gateway-client/src/pending-request.ts | 12 ++ .../src/protocol-client.handshake.test.ts | 95 +++++++++ .../gateway-client/src/protocol-client.ts | 50 ++--- packages/gateway-client/src/timeouts.ts | 16 ++ scripts/check-changed.mjs | 46 ++++- src/infra/dotenv.ts | 1 - ...ged-lanes-generated-extension-lint.test.ts | 50 +++++ test/scripts/changed-lanes.test.ts | 42 ++++ ui/src/api/gateway.node.test.ts | 111 +++++++++++ ui/src/api/gateway.ts | 45 ++++- .../app-sidebar-session-catalog-live.ts | 8 + ui/src/components/app-sidebar.test.ts | 1 + ui/src/components/exec-approval.test.ts | 25 +++ ui/src/components/exec-approval.ts | 9 +- .../session-data-controller-catalog.ts | 9 + ui/src/components/session-data-controller.ts | 6 +- ui/src/e2e/chat-flow.media-files.e2e.test.ts | 29 ++- ui/src/lib/board/gateway-provider.test.ts | 181 +++++++++++++++++- ui/src/lib/board/gateway-provider.ts | 29 ++- ...chat-history-subscription-disposal.test.ts | 114 +++++++++++ ui/src/pages/chat/chat-history.ts | 42 ++++ ui/src/pages/chat/chat-state-controller.ts | 3 + .../chat-message-media-lifecycle.test.ts | 45 +++++ .../chat/components/chat-message-media.ts | 49 ++--- .../chat-message-pairing-qr-lifecycle.test.ts | 104 ++++++++++ .../app-sidebar-cases/catalog-ownership.ts | 150 ++++++++------- .../app-sidebar-cases/catalog-reconnect.ts | 62 ++++++ ui/src/test-helpers/control-ui-e2e.ts | 10 + 32 files changed, 1297 insertions(+), 166 deletions(-) create mode 100644 packages/gateway-client/src/event-listeners.ts create mode 100644 packages/gateway-client/src/pending-request.ts create mode 100644 packages/gateway-client/src/protocol-client.handshake.test.ts create mode 100644 test/scripts/changed-lanes-generated-extension-lint.test.ts create mode 100644 ui/src/pages/chat/chat-history-subscription-disposal.test.ts create mode 100644 ui/src/pages/chat/components/chat-message-pairing-qr-lifecycle.test.ts create mode 100644 ui/src/test-helpers/app-sidebar-cases/catalog-reconnect.ts diff --git a/extensions/browser/chrome-extension/modules/copilot-runtime.js b/extensions/browser/chrome-extension/modules/copilot-runtime.js index b120b37acd84..3954c1c1b6c1 100644 --- a/extensions/browser/chrome-extension/modules/copilot-runtime.js +++ b/extensions/browser/chrome-extension/modules/copilot-runtime.js @@ -1 +1 @@ -function normalizeDeviceMetadataForAuth(value){if(typeof value!="string")return"";let trimmed=value.trim();return trimmed?trimmed.replace(/[A-Z]/g,char=>String.fromCharCode(char.charCodeAt(0)+32)):""}function buildDeviceAuthPayloadV3(params){let scopes=params.scopes.join(","),token=params.token??"",platform=normalizeDeviceMetadataForAuth(params.platform),deviceFamily=normalizeDeviceMetadataForAuth(params.deviceFamily);return["v3",params.deviceId,params.clientId,params.clientMode,params.role,scopes,String(params.signedAtMs),token,params.nonce,platform,deviceFamily].join("|")}function normalized(value){return typeof value=="string"&&value.trim()||void 0}function selectGatewayConnectAuth(params){let authToken=normalized(params.token),bootstrapToken=normalized(params.bootstrapToken),explicitDeviceToken=normalized(params.deviceToken),authPassword=normalized(params.password),storedToken=normalized(params.storedToken),stored={storedToken,storedScopes:params.storedScopes};if(params.preferBootstrapToken&&bootstrapToken)return{authBootstrapToken:bootstrapToken,authPassword,...stored};let useRetryToken=params.pendingDeviceTokenRetry===!0&&!explicitDeviceToken&&!!(authToken&&storedToken&¶ms.trustedDeviceTokenRetry),resolvedDeviceToken=explicitDeviceToken??(useRetryToken||!(authToken||authPassword)&&(!bootstrapToken||storedToken)?storedToken:void 0),usingStoredDeviceToken=!!(resolvedDeviceToken&&!explicitDeviceToken&&storedToken)&&resolvedDeviceToken===storedToken,selectedToken=authToken??resolvedDeviceToken,authBootstrapToken=!authToken&&!resolvedDeviceToken&&!authPassword?bootstrapToken:void 0;return{authToken:selectedToken,authBootstrapToken,authDeviceToken:useRetryToken?storedToken:void 0,authPassword,authApprovalRuntimeToken:normalized(params.approvalRuntimeToken),authAgentRuntimeIdentityToken:normalized(params.agentRuntimeIdentityToken),signatureToken:selectedToken??authBootstrapToken,resolvedDeviceToken,usingStoredDeviceToken,...stored}}function buildGatewayConnectAuth(selected){let auth={token:selected.authToken,bootstrapToken:selected.authBootstrapToken,deviceToken:selected.authDeviceToken??selected.resolvedDeviceToken,password:selected.authPassword,approvalRuntimeToken:selected.authApprovalRuntimeToken,agentRuntimeIdentityToken:selected.authAgentRuntimeIdentityToken};return Object.values(auth).some(Boolean)?auth:void 0}function resolveGatewayConnectScopes(params){return params.requestedScopes??(params.usingStoredDeviceToken&¶ms.storedScopes?.length?params.storedScopes:[...params.defaultScopes])}var GatewayBrowserDeviceAuthLifecycle=class{constructor(deps){this.deps=deps}async buildPlan(params){let identity=await this.deps.loadIdentity(),stored=identity?await this.deps.tokenStore.load({clientId:params.client.id,deviceId:identity.deviceId,role:params.role}):null,storedValue=stored?.token,selectedAuth=selectGatewayConnectAuth({token:params.token,bootstrapToken:params.bootstrapToken,password:params.password,storedToken:storedValue,storedScopes:stored?.scopes,pendingDeviceTokenRetry:params.pendingDeviceTokenRetry,trustedDeviceTokenRetry:params.trustedDeviceTokenRetry,preferBootstrapToken:params.preferBootstrapToken}),{usingStoredDeviceToken}=selectedAuth,scopes=resolveGatewayConnectScopes({requestedScopes:selectedAuth.authBootstrapToken&¶ms.bootstrapScopes?[...params.bootstrapScopes]:void 0,usingStoredDeviceToken,storedScopes:selectedAuth.storedScopes,defaultScopes:params.defaultScopes});if(!identity)return{clientId:params.client.id,role:params.role,identity,selectedAuth,scopes,auth:buildGatewayConnectAuth(selectedAuth)};let signedAtMs=this.deps.nowMs?.()??Date.now(),nonce=params.nonce??"",{authBootstrapToken:primary,signatureToken:signed}=selectedAuth,token=null;primary?token=primary:signed&&(token=signed);let payload=buildDeviceAuthPayloadV3({deviceId:identity.deviceId,clientId:params.client.id,clientMode:params.client.mode,role:params.role,scopes,signedAtMs,token,nonce,platform:params.client.platform,deviceFamily:params.client.deviceFamily});return{clientId:params.client.id,role:params.role,identity,selectedAuth,scopes,auth:buildGatewayConnectAuth(selectedAuth),device:{id:identity.deviceId,publicKey:identity.publicKey,signature:await identity.sign(payload),signedAt:signedAtMs,nonce}}}async acceptHello(hello,plan){let token=hello.auth?.deviceToken?.trim();!token||!plan.identity||await this.deps.tokenStore.store({clientId:plan.clientId,deviceId:plan.identity.deviceId,role:hello.auth?.role??plan.role,token,scopes:hello.auth?.scopes??[]})}async clearStoredToken(plan){plan.identity&&await this.deps.tokenStore.clear({clientId:plan.clientId,deviceId:plan.identity.deviceId,role:plan.role})}};function isRecord(value){return!!value&&typeof value=="object"&&!Array.isArray(value)}function isNonEmptyString(value){return typeof value=="string"&&value.length>0}function isNonNegativeInteger(value){return typeof value=="number"&&Number.isInteger(value)&&value>=0}function isGatewayErrorShape(value){return!isRecord(value)||!isNonEmptyString(value.code)||!isNonEmptyString(value.message)||value.retryable!==void 0&&typeof value.retryable!="boolean"?!1:value.retryAfterMs===void 0||isNonNegativeInteger(value.retryAfterMs)}function isGatewayEventFrame(value){return!isRecord(value)||value.type!=="event"||!isNonEmptyString(value.event)?!1:value.seq===void 0||isNonNegativeInteger(value.seq)}function isGatewayResponseFrame(value){return!isRecord(value)||value.type!=="res"||!isNonEmptyString(value.id)||typeof value.ok!="boolean"?!1:value.error===void 0||isGatewayErrorShape(value.error)}function computeBackoff(policy,attempt){let base=Math.min(policy.maxMs,policy.initialMs*policy.factor**Math.max(attempt-1,0)),jitter=base*policy.jitter*Math.random();return Math.min(policy.maxMs,Math.round(base+jitter))}async function sleepWithAbort(ms,abortSignal,options={}){if(!Number.isFinite(ms)||ms<=0)return;let delayMs=Math.min(Math.max(Math.floor(ms),1),2147e6);await new Promise((resolve,reject)=>{let settled=!1,timer=null,cleanup=()=>abortSignal?.removeEventListener("abort",onAbort),onAbort=()=>{settled||(settled=!0,timer&&clearTimeout(timer),timer=null,cleanup(),reject(new Error("aborted",{cause:abortSignal?.reason??new Error("aborted")})))};if(abortSignal?.addEventListener("abort",onAbort,{once:!0}),abortSignal?.aborted){onAbort();return}timer=setTimeout(()=>{settled=!0,cleanup(),timer=null,resolve()},delayMs),options.ref===!1&&timer.unref?.(),abortSignal?.aborted&&onAbort()})}var RetrySupervisor=class{constructor(policy,maxAttempts=Number.POSITIVE_INFINITY){this.policy=policy;this.maxAttempts=maxAttempts;this.attempts=0;this.initialMs=policy.initialMs}reset(initialMs=this.policy.initialMs){this.cancel(),this.attempts=0,this.initialMs=initialMs,this.nextDelayOverrideMs=void 0}cancel(reason=new Error("retry cancelled")){this.pendingAbort?.abort(reason),this.pendingAbort=void 0}next(abortSignal){let override=this.nextDelayOverrideMs;if(this.nextDelayOverrideMs=void 0,override===void 0&&++this.attempts>Math.ceil(this.maxAttempts))return;let attempt=Math.max(this.attempts,1),delayMs=override??computeBackoff({...this.policy,initialMs:this.initialMs},attempt);this.cancel();let pendingAbort=new AbortController;return this.pendingAbort=pendingAbort,{attempt,delayMs,signal:abortSignal?AbortSignal.any([pendingAbort.signal,abortSignal]):pendingAbort.signal}}},DEFAULT_RETRY_CONFIG={attempts:3,minDelayMs:300,maxDelayMs:3e4,jitter:0},defaultSleep=ms=>new Promise(resolve=>{setTimeout(resolve,ms)});function asFiniteNumber(value){return typeof value=="number"&&Number.isFinite(value)?value:void 0}function clampNumber(value,fallback,min,max){let next=asFiniteNumber(value);return next===void 0?fallback:Math.min(Math.max(next,min??Number.NEGATIVE_INFINITY),max??Number.POSITIVE_INFINITY)}function resolveAttemptCount(value,fallback){return Math.max(1,Math.round(asFiniteNumber(value)??fallback))}function resolveRetryDelayMs(value){let finite=value===Number.POSITIVE_INFINITY?2147e6:asFiniteNumber(value)??0;return Math.min(Math.max(Math.round(finite),0),2147e6)}function resolveJitterConfig(value,fallback){if(value==="full")return"full";let fraction=asFiniteNumber(value);return fraction===void 0?fallback:Math.min(Math.max(fraction,0),1)}function resolveRetryConfig(defaults=DEFAULT_RETRY_CONFIG,overrides){let attempts=resolveAttemptCount(overrides?.attempts,defaults.attempts),minDelayMs=resolveRetryDelayMs(clampNumber(overrides?.minDelayMs,defaults.minDelayMs,0)),maxDelayMs=Math.max(minDelayMs,resolveRetryDelayMs(clampNumber(overrides?.maxDelayMs,defaults.maxDelayMs,0)));return{attempts,minDelayMs,maxDelayMs,jitter:resolveJitterConfig(overrides?.jitter,defaults.jitter)}}function applyJitter(delayMs,jitter,mode,random){if(jitter==="full")return mode==="symmetric"?Math.max(0,Math.round(delayMs*(.5+random()*.5))):Math.max(0,Math.ceil(delayMs*(1+random())));if(jitter<=0)return mode==="positive"?Math.ceil(delayMs):delayMs;let fraction=random(),offset=mode==="positive"?fraction*jitter:(fraction*2-1)*jitter,raw=delayMs*(1+offset);return Math.max(0,mode==="positive"?Math.ceil(raw):Math.round(raw))}function toRetryError(value,fallbackMessage="Non-Error thrown"){if(value instanceof Error)return value;if(typeof value=="string")return new Error(value);let error=new Error(fallbackMessage,{cause:value});return(typeof value=="object"&&value!==null||typeof value=="function")&&Object.assign(error,value),error}function createRetryRunner(runtime={}){let runtimeSleep=runtime.sleep??defaultSleep,runtimeRandom=runtime.random??Math.random,createFailure=runtime.createFailure??(errors=>toRetryError(errors.at(-1)??new Error("Retry failed")));return async function(fn,attemptsOrOptions=3,initialDelayMs=300){let attemptErrors=[];if(typeof attemptsOrOptions=="number"){let attempts=resolveAttemptCount(attemptsOrOptions,DEFAULT_RETRY_CONFIG.attempts);for(let index=0;index0?resolved.maxDelayMs:Number.POSITIVE_INFINITY,retryAfterMaxDelayMs=options.retryAfterMaxDelayMs===void 0?maxDelayMs:Math.max(minDelayMs,resolveRetryDelayMs(clampNumber(options.retryAfterMaxDelayMs,maxDelayMs,0))),random=options.random??runtimeRandom,sleep=options.sleep??runtimeSleep,shouldRetry=options.shouldRetry??(()=>!0);for(let attempt=1;attempt<=maxAttempts;attempt+=1)try{return await fn()}catch(err2){if(attemptErrors.push(err2),attempt>=maxAttempts||!shouldRetry(err2,attempt))break;let context={attempt,maxAttempts,err:err2,label:options.label},retryAfterMs=options.retryAfterMs?.(err2),hasRetryAfter=typeof retryAfterMs=="number"&&Number.isFinite(retryAfterMs),configuredDelay=typeof options.delayMs=="function"?options.delayMs(context):options.delayMs,resolvedConfiguredDelay=configuredDelay===void 0?void 0:resolveRetryDelayMs(configuredDelay),baseDelay=hasRetryAfter?Math.max(retryAfterMs,minDelayMs):resolvedConfiguredDelay===void 0?minDelayMs*2**(attempt-1):Math.max(resolvedConfiguredDelay,minDelayMs),delayCap=hasRetryAfter?retryAfterMaxDelayMs:maxDelayMs,delay=Math.min(baseDelay,delayCap),canHonorRetryAfter=hasRetryAfter&&(retryAfterMs??0)<=delayCap,wantsPositiveDraw=resolved.jitter==="full"&&!hasRetryAfter||canHonorRetryAfter;delay=applyJitter(delay,resolved.jitter,wantsPositiveDraw?"positive":"symmetric",random),delay=Math.min(Math.max(delay,minDelayMs),delayCap),await options.onRetry?.({...context,delayMs:delay}),delay>0&&await sleep(delay)}throw createFailure(attemptErrors)}}var retryAsync=createRetryRunner();var GatewayProtocolRequestError=class extends Error{constructor(error){super(error.message??"request failed"),this.name="GatewayProtocolRequestError",this.code=error.code??"UNAVAILABLE",this.gatewayCode=this.code,this.details=error.details,this.retryable=error.retryable===!0,this.retryAfterMs=error.retryAfterMs}},GatewayProtocolClient=class{constructor(opts){this.opts=opts;this.socket=null;this.pending=new Map;this.listeners=new Set;this.stopped=!0;this.generation=0;this.lastSeq=null;this.connectNonce=null;this.connectSent=!1;this.connectRequestSent=!1;this.handshakeTimer=null;this.reconnectSignal=null;this.socketOpened=!1;this.helloReceived=!1;this.connectTiming=null;this.reconnectSupervisor=new RetrySupervisor({initialMs:opts.reconnect.initialMs,maxMs:opts.reconnect.maxMs,factor:opts.reconnect.multiplier,jitter:0})}get connected(){return this.socket?.isOpen()??!1}get hasPendingRequests(){return this.pending.size>0}get connecting(){return this.connectSent&&!this.helloReceived}get hasUnboundedPendingRequests(){return[...this.pending.values()].some(pending=>pending.unbounded)}start(){this.socket||this.reconnectSignal||(this.stopped=!1,this.reconnectSupervisor.cancel(),this.connect())}stop(){this.stopped=!0,this.clearHandshakeTimer(),this.reconnectSignal=null,this.reconnectSupervisor.reset();let socket=this.socket;socket&&this.opts.notifyStoppedClose&&(this.stoppedSocket={socket,context:this.closeContext()}),this.socket=null,this.connectFailure=void 0,this.connectTiming=null,this.flushRequests(new Error("gateway client stopped")),socket?.close()}request(method,params,options){let socket=this.socket;if(!socket?.isOpen())return Promise.reject(new Error("gateway not connected"));if(typeof method!="string"||method.length===0)return Promise.reject(new Error("invalid request frame: method must be a non-empty string"));let id=this.opts.createRequestId(),timeoutMs=options?.timeoutMs===null?void 0:options?.timeoutMs??this.opts.requestTimeoutMs;return new Promise((resolve,reject)=>{let timeout,pending={resolve:value=>resolve(value),reject,expectFinal:options?.expectFinal===!0,acceptedNotified:!1,onAccepted:options?.onAccepted,unbounded:timeoutMs===void 0,method,startedAtMs:this.nowMs()},onAbort=()=>{this.pending.delete(id),timeout&&clearTimeout(timeout),this.finishRequestTiming(id,pending,!1,"CLIENT_ABORTED"),reject(this.opts.createRequestAbortError?.(method)??new Error(`gateway request aborted for ${method}`))},cleanup=()=>{timeout&&clearTimeout(timeout),options?.signal?.removeEventListener("abort",onAbort)};if(options?.signal?.aborted){reject(this.opts.createRequestAbortError?.(method)??new Error(`gateway request aborted for ${method}`));return}pending.cleanup=cleanup,timeoutMs!==void 0&&timeoutMs>=0&&(timeout=setTimeout(()=>{this.pending.delete(id),options?.signal?.removeEventListener("abort",onAbort),this.finishRequestTiming(id,pending,!1,"CLIENT_TIMEOUT"),reject(this.opts.createRequestTimeoutError?.(method,timeoutMs)??new Error(`gateway request timed out after ${timeoutMs}ms: ${method}`))},timeoutMs),timeout.unref?.()),options?.signal?.addEventListener("abort",onAbort,{once:!0}),this.pending.set(id,pending);try{socket.send(JSON.stringify({type:"req",id,method,params})),this.invoke("sent",()=>options?.onSent?.())}catch(error){this.pending.delete(id),cleanup(),this.finishRequestTiming(id,pending,!1,"CLIENT_SEND_ERROR"),reject(error instanceof Error?error:new Error(String(error)))}})}addEventListener(listener){return this.listeners.add(listener),()=>this.listeners.delete(listener)}closeSocket(code,reason){this.socket?.close(code,reason)}resetReconnectBackoff(initialMs){this.reconnectSignal=null,this.reconnectSupervisor.reset(initialMs)}recordTiming(phase,generation,plan,detail){let now=this.nowMs(),state=this.connectTiming;!state||state.generation!==generation||(state.hasChallenge||=phase==="challenge",state.usedFallback||=phase==="fallback",this.invoke("connect timing",()=>this.opts.onTiming?.({phase,generation,durationMs:Math.max(0,now-state.startedAtMs),phaseDurationMs:Math.max(0,now-state.lastAtMs),hasChallenge:state.hasChallenge,usedFallback:state.usedFallback,plan,detail})),state.lastAtMs=now,(phase==="hello"||phase==="failed")&&(this.connectTiming=null))}connect(){if(this.stopped)return;let generation=this.generation+1;this.connectNonce=null,this.connectSent=!1,this.connectRequestSent=!1,this.socketOpened=!1,this.helloReceived=!1,this.connectFailure=void 0;let socket;try{socket=this.opts.createSocket({open:()=>this.handleOpen(socket,generation),message:data=>this.handleMessage(socket,generation,data),close:(code,reason)=>this.handleClose(socket,generation,code,reason),error:error=>this.handleSocketError(socket,generation,error)})}catch(error){let normalized2=error instanceof Error?error:new Error(String(error));if(this.opts.onSocketFactoryError?.(normalized2),this.opts.onConnectError?.(normalized2),this.opts.rethrowSocketFactoryError?.(normalized2))throw normalized2;this.opts.shouldRetrySocketFactoryError?.(normalized2)&&!this.stopped&&!this.socket&&!this.reconnectSignal&&this.scheduleReconnect();return}this.generation=generation,this.socket=socket;let now=this.nowMs();this.connectTiming={generation,startedAtMs:now,lastAtMs:now,hasChallenge:!1,usedFallback:!1}}handleOpen(socket,generation){if(this.isActive(socket,generation)){if(this.socketOpened=!0,this.recordTiming("socket-open",generation),this.connectNonce){this.sendConnect(socket,generation);return}this.armHandshakeTimer(socket,generation)}}armHandshakeTimer(socket,generation){this.clearHandshakeTimer();let armedAt=Date.now();this.handshakeTimer=setTimeout(()=>{if(this.handshakeTimer=null,!this.isActive(socket,generation)||this.connectSent||!socket.isOpen())return;if(this.opts.handshake.mode==="fallback"){this.recordTiming("fallback",generation),this.sendConnect(socket,generation);return}let elapsedMs=Date.now()-armedAt,error=new Error(this.opts.handshake.timeoutMessage?.(elapsedMs)??`gateway connect challenge timeout after ${elapsedMs}ms`);this.opts.onConnectError?.(error),socket.close(1008,"connect challenge timeout")},this.opts.handshake.timeoutMs),this.handshakeTimer.unref?.()}sendConnect(socket,generation){if(!this.isActive(socket,generation)||!socket.isOpen()||this.connectSent)return;this.connectSent=!0,this.clearHandshakeTimer();let planOrPromise;try{planOrPromise=this.opts.buildConnectPlan({nonce:this.connectNonce,generation})}catch(error){this.handleConnectPlanError(socket,generation,error);return}if(planOrPromise instanceof Promise){planOrPromise.then(plan=>this.sendConnectPlan(socket,generation,plan)).catch(error=>this.handleConnectPlanError(socket,generation,error));return}this.sendConnectPlan(socket,generation,planOrPromise)}handleConnectPlanError(socket,generation,error){if(!this.isActive(socket,generation))return;let normalized2=error instanceof Error?error:new Error(String(error)),outcome=this.opts.onConnectPlanError?.(normalized2)??{closeCode:1008,closeReason:"connect failed"};this.opts.onConnectError?.(outcome.error??normalized2),outcome.stop&&(this.stopped=!0),socket.close(outcome.closeCode,outcome.closeReason)}sendConnectPlan(socket,generation,plan){if(!this.isActive(socket,generation)||!socket.isOpen())return;let context={generation,nonce:this.connectNonce,plan};this.recordTiming("connect-plan-ready",generation,plan),this.recordTiming("request-sent",generation,plan),this.connectRequestSent=!0,this.request("connect",this.opts.buildConnectParams(plan)).then(hello=>{this.isActive(socket,generation)&&(this.helloReceived=!0,this.connectFailure=void 0,this.reconnectSupervisor.reset(),this.recordTiming("hello",generation,plan),this.opts.onConnectHello?.(hello,context),this.invoke("hello",()=>this.opts.onHello?.(hello)))}).catch(error=>{if(!this.isActive(socket,generation))return;let requestError=error instanceof GatewayProtocolRequestError?error:new GatewayProtocolRequestError({message:String(error)}),outcome=this.opts.onConnectFailure?.(requestError,context)??{closeCode:1008,closeReason:"connect failed"};this.connectFailure={error:requestError,reconnectDelayMs:outcome.reconnectDelayMs},outcome.stop&&(this.stopped=!0),socket.close(outcome.closeCode,outcome.closeReason)})}handleMessage(socket,generation,raw){if(!this.isActive(socket,generation))return;let parsed;try{parsed=JSON.parse(raw)}catch(error){this.opts.onParseError?.(error);return}if(isGatewayEventFrame(parsed)){if(this.opts.onActivity?.(),parsed.event==="connect.challenge"){let payload=parsed.payload,nonce=typeof payload?.nonce=="string"?payload.nonce.trim():"";if(!nonce){if(this.opts.handshake.mode==="require-challenge"){let error=new Error("gateway connect challenge missing nonce");this.opts.onConnectError?.(error),socket.close(1008,"connect challenge missing nonce")}return}this.connectNonce=nonce,this.recordTiming("challenge",generation),this.sendConnect(socket,generation);return}let seq=typeof parsed.seq=="number"?parsed.seq:null;if(seq!==null){if(this.lastSeq!==null&&seq>this.lastSeq+1){let expected=this.lastSeq+1;if(this.invoke("gap",()=>this.opts.onGap?.({expected,received:seq})),!this.isActive(socket,generation))return}this.lastSeq=seq}this.invoke("event",()=>this.opts.onEvent?.(parsed));for(let listener of this.listeners)this.invoke("event listener",()=>listener(parsed));return}isGatewayResponseFrame(parsed)&&(this.opts.onActivity?.(),this.handleResponse(parsed))}handleResponse(frame){let pending=this.pending.get(frame.id);if(!pending)return;let status=frame.payload?.status;if(pending.expectFinal&&status==="accepted"){pending.acceptedNotified||(pending.acceptedNotified=!0,this.invoke("accepted",()=>pending.onAccepted?.(frame.payload)));return}if(this.pending.delete(frame.id),pending.cleanup?.(),frame.ok){this.finishRequestTiming(frame.id,pending,!0),pending.resolve(frame.payload);return}this.finishRequestTiming(frame.id,pending,!1,frame.error?.code),pending.reject(this.opts.createRequestError?.(frame.error??{})??new GatewayProtocolRequestError(frame.error??{}))}handleClose(socket,generation,code,reason){if(this.socket!==socket){if(this.stoppedSocket?.socket===socket){let context2={...this.stoppedSocket.context,code,reason};this.stoppedSocket=void 0,this.invoke("close",()=>this.opts.onClose?.(context2,{retry:!1,notify:!0}))}return}this.socket=null,this.clearHandshakeTimer();let context={...this.closeContext(),code,reason,generation};this.connectFailure=void 0;let decision=this.opts.resolveClose(context);this.flushRequests(decision.pendingError??context.connectFailure?.error??new Error(`gateway closed (${code}): ${reason}`)),this.invoke("close",()=>this.opts.onClose?.(context,decision)),decision.retry&&!this.stopped&&this.scheduleReconnect(decision.reconnectDelayMs??context.connectFailure?.reconnectDelayMs)}handleSocketError(socket,generation,error){!this.isActive(socket,generation)||this.connectSent||this.opts.onConnectError?.(error)}flushRequests(error){for(let[id,pending]of this.pending)this.finishRequestTiming(id,pending,!1,"CLIENT_CLOSED"),pending.cleanup?.(),pending.reject(error);this.pending.clear()}finishRequestTiming(id,pending,ok,errorCode){let endedAtMs=this.nowMs();this.invoke("request timing",()=>this.opts.onRequestTiming?.({id,method:pending.method,ok,durationMs:Math.max(0,endedAtMs-pending.startedAtMs),startedAtMs:pending.startedAtMs,endedAtMs,errorCode}))}scheduleReconnect(overrideMs){overrideMs!==void 0&&(this.reconnectSupervisor.nextDelayOverrideMs=overrideMs);let retry=this.reconnectSupervisor.next();retry&&(this.reconnectSignal=retry.signal,sleepWithAbort(retry.delayMs,retry.signal).then(()=>{this.reconnectSignal===retry.signal&&(this.reconnectSignal=null,this.connect())},()=>{this.reconnectSignal===retry.signal&&(this.reconnectSignal=null)}))}closeContext(){return{generation:this.generation,socketOpened:this.socketOpened,helloReceived:this.helloReceived,connectRequestSent:this.connectRequestSent,connectFailure:this.connectFailure}}isActive(socket,generation){return!this.stopped&&this.socket===socket&&this.generation===generation}nowMs(){return this.opts.nowMs?.()??Date.now()}clearHandshakeTimer(){this.handshakeTimer&&(clearTimeout(this.handshakeTimer),this.handshakeTimer=null)}invoke(label,callback){try{callback()}catch(error){this.opts.onCallbackError?.(label,error)}}};var GATEWAY_CLIENT_IDS={WEBCHAT_UI:"webchat-ui",CONTROL_UI:"openclaw-control-ui",BROWSER_COPILOT:"openclaw-browser-copilot",TUI:"openclaw-tui",WEBCHAT:"webchat",CLI:"cli",GATEWAY_CLIENT:"gateway-client",MACOS_APP:"openclaw-macos",LINUX_APP:"openclaw-linux",IOS_APP:"openclaw-ios",WATCHOS_APP:"openclaw-watchos",ANDROID_APP:"openclaw-android",NODE_HOST:"node-host",WORKER:"openclaw-worker",TEST:"test",FINGERPRINT:"fingerprint",PROBE:"openclaw-probe"};var GATEWAY_CLIENT_MODES={WEBCHAT:"webchat",CLI:"cli",UI:"ui",BACKEND:"backend",NODE:"node",WORKER:"worker",PROBE:"probe",TEST:"test"},GATEWAY_CLIENT_CAPS={AGENT_KIND:"agent-kind",APPROVALS:"approvals",EXEC_APPROVALS:"exec-approvals",INLINE_WIDGETS:"inline-widgets",RUN_TOOL_BINDINGS:"run-tool-bindings",SESSION_SCOPED_EVENTS:"session-scoped-events",PLUGIN_APPROVALS:"plugin-approvals",TASK_SUGGESTIONS:"task-suggestions",TERMINAL_OFFSET_SEQ:"terminal-offset-seq",TOOL_EVENTS:"tool-events",UI_COMMANDS:"ui-commands"},GATEWAY_CLIENT_ID_SET=new Set(Object.values(GATEWAY_CLIENT_IDS)),GATEWAY_CLIENT_MODE_SET=new Set(Object.values(GATEWAY_CLIENT_MODES));var PROTOCOL_VERSION=4,MIN_CLIENT_PROTOCOL_VERSION=4;/*! noble-ed25519 - MIT License (c) 2019 Paul Miller (paulmillr.com) */var ed25519_CURVE=Object.freeze({p:0x7fffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffedn,n:0x1000000000000000000000000000000014def9dea2f79cd65812631a5cf5d3edn,h:8n,a:0x7fffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffecn,d:0x52036cee2b6ffe738cc740797779e89800700a4d4141d8ab75eb4dca135978a3n,Gx:0x216936d3cd6e53fec0a4e231fdd6dc5c692cc7609525a7b2c9562d608f25d51an,Gy:0x6666666666666666666666666666666666666666666666666666666666666658n}),{p:P,n:N,Gx,Gy,a:_a,d:_d,h}=ed25519_CURVE,L=32,captureTrace=(...args)=>{"captureStackTrace"in Error&&typeof Error.captureStackTrace=="function"&&Error.captureStackTrace(...args)},err=(message="")=>{let e=new Error(message);throw captureTrace(e,err),e},isBig=n=>typeof n=="bigint",isStr=s=>typeof s=="string",isBytes=a=>a instanceof Uint8Array||ArrayBuffer.isView(a)&&a.constructor.name==="Uint8Array"&&"BYTES_PER_ELEMENT"in a&&a.BYTES_PER_ELEMENT===1,abytes=(value,length,title="")=>{let bytes=isBytes(value),len=value?.length,needsLen=length!==void 0;if(!bytes||needsLen&&len!==length){let prefix=title&&`"${title}" `,ofLen=needsLen?` of length ${length}`:"",got=bytes?`length=${len}`:`type=${typeof value}`,msg=prefix+"expected Uint8Array"+ofLen+", got "+got;throw bytes?new RangeError(msg):new TypeError(msg)}return value},u8n=len=>new Uint8Array(len),u8fr=buf=>Uint8Array.from(buf),padh=(n,pad)=>n.toString(16).padStart(pad,"0"),bytesToHex=b=>Array.from(abytes(b)).map(e=>padh(e,2)).join(""),C={_0:48,_9:57,A:65,F:70,a:97,f:102},_ch=ch=>{if(ch>=C._0&&ch<=C._9)return ch-C._0;if(ch>=C.A&&ch<=C.F)return ch-(C.A-10);if(ch>=C.a&&ch<=C.f)return ch-(C.a-10)},hexToBytes=hex=>{let e="hex invalid";if(!isStr(hex))return err(e);let hl=hex.length,al=hl/2;if(hl%2)return err(e);let array=u8n(al);for(let ai=0,hi=0;aiglobalThis?.crypto,subtle=()=>cr()?.subtle??err("crypto.subtle must be defined, consider polyfill"),concatBytes=(...arrs)=>{let len=0;for(let a of arrs)len+=abytes(a).length;let r=u8n(len),pad=0;return arrs.forEach(a=>{r.set(a,pad),pad+=a.length}),r},randomBytes=(len=L)=>cr().getRandomValues(u8n(len)),big=BigInt,assertRange=(n,min,max,msg="bad number: out of range")=>{if(!isBig(n))throw new TypeError(msg);if(min<=n&&n{let r=a%b;return r>=0n?r:b+r},P_MASK=(1n<<255n)-1n,modP=num=>{num<0n&&err("negative coordinate");let r=(num>>255n)*19n+(num&P_MASK);return r=(r>>255n)*19n+(r&P_MASK),r%P},modN=a=>M(a,N),invert=(num,md)=>{(num===0n||md<=0n)&&err("no inverse n="+num+" mod="+md);let a=M(num,md),b=md,x=0n,y=1n,u=1n,v=0n;for(;a!==0n;){let q=b/a,r=b%a,m=x-u*q,n=y-v*q;b=a,a=r,x=u,y=v,u=m,v=n}return b===1n?M(x,md):err("no inverse")},callHash=name=>{let fn=hashes[name];return typeof fn!="function"&&err("hashes."+name+" not set"),fn},checkDigest=value=>abytes(value,64,"digest");var apoint=p=>p instanceof Point?p:err("Point expected"),B256=2n**256n,Point=class _Point{static BASE;static ZERO;X;Y;Z;T;constructor(X,Y,Z,T){let max=B256;this.X=assertRange(X,0n,max),this.Y=assertRange(Y,0n,max),this.Z=assertRange(Z,1n,max),this.T=assertRange(T,0n,max),Object.freeze(this)}static CURVE(){return ed25519_CURVE}static fromAffine(p){return new _Point(p.x,p.y,1n,modP(p.x*p.y))}static fromBytes(hex,zip215=!1){let d=_d,normed=u8fr(abytes(hex,L)),lastByte=hex[31];normed[31]=lastByte&-129;let y=bytesToNumberLE(normed);assertRange(y,0n,zip215?B256:P);let y2=modP(y*y),u=M(y2-1n),v=modP(d*y2+1n),{isValid,value:x}=uvRatio(u,v);isValid||err("bad point: y not sqrt");let isXOdd=(x&1n)===1n,isLastByteOdd=(lastByte&128)!==0;return!zip215&&x===0n&&isLastByteOdd&&err("bad point: x==0, isLastByteOdd"),isLastByteOdd!==isXOdd&&(x=M(-x)),new _Point(x,y,1n,modP(x*y))}static fromHex(hex,zip215){return _Point.fromBytes(hexToBytes(hex),zip215)}get x(){return this.toAffine().x}get y(){return this.toAffine().y}assertValidity(){let a=_a,d=_d,p=this;if(p.is0())return err("bad point: ZERO");let{X,Y,Z,T}=p,X2=modP(X*X),Y2=modP(Y*Y),Z2=modP(Z*Z),Z4=modP(Z2*Z2),aX2=modP(X2*a),left=modP(Z2*(aX2+Y2)),right=M(Z4+modP(d*modP(X2*Y2)));if(left!==right)return err("bad point: equation left != right (1)");let XY=modP(X*Y),ZT=modP(Z*T);return XY!==ZT?err("bad point: equation left != right (2)"):this}equals(other){let{X:X1,Y:Y1,Z:Z1}=this,{X:X2,Y:Y2,Z:Z2}=apoint(other),X1Z2=modP(X1*Z2),X2Z1=modP(X2*Z1),Y1Z2=modP(Y1*Z2),Y2Z1=modP(Y2*Z1);return X1Z2===X2Z1&&Y1Z2===Y2Z1}is0(){return this.equals(I)}negate(){return new _Point(M(-this.X),this.Y,this.Z,M(-this.T))}double(){let{X:X1,Y:Y1,Z:Z1}=this,a=_a,A=modP(X1*X1),B=modP(Y1*Y1),C2=modP(2n*Z1*Z1),D=modP(a*A),x1y1=M(X1+Y1),E=M(modP(x1y1*x1y1)-A-B),G2=M(D+B),F=M(G2-C2),H=M(D-B),X3=modP(E*F),Y3=modP(G2*H),T3=modP(E*H),Z3=modP(F*G2);return new _Point(X3,Y3,Z3,T3)}add(other){let{X:X1,Y:Y1,Z:Z1,T:T1}=this,{X:X2,Y:Y2,Z:Z2,T:T2}=apoint(other),a=_a,d=_d,A=modP(X1*X2),B=modP(Y1*Y2),C2=modP(modP(T1*d)*T2),D=modP(Z1*Z2),E=M(modP(M(X1+Y1)*M(X2+Y2))-A-B),F=M(D-C2),G2=M(D+C2),H=M(B-modP(a*A)),X3=modP(E*F),Y3=modP(G2*H),T3=modP(E*H),Z3=modP(F*G2);return new _Point(X3,Y3,Z3,T3)}subtract(other){return this.add(apoint(other).negate())}multiply(n,safe=!0){if(!safe&&n===0n||(assertRange(n,1n,N),!safe&&this.is0()))return I;if(n===1n)return this;if(this.equals(G))return wNAF(n).p;let p=I,f=G;for(let d=this;n>0n;d=d.double(),n>>=1n)n&1n?p=p.add(d):safe&&(f=f.add(d));return p}multiplyUnsafe(scalar){return this.multiply(scalar,!1)}toAffine(){let{X,Y,Z}=this;if(this.equals(I))return{x:0n,y:1n};let iz=invert(Z,P);modP(Z*iz)!==1n&&err("invalid inverse");let x=modP(X*iz),y=modP(Y*iz);return{x,y}}toBytes(){let{x,y}=this.toAffine(),b=numTo32bLE(y);return b[31]|=x&1n?128:0,b}toHex(){return bytesToHex(this.toBytes())}clearCofactor(){return this.multiply(big(h),!1)}isSmallOrder(){return this.clearCofactor().is0()}isTorsionFree(){let p=this.multiply(N/2n,!1).double();return N%2n&&(p=p.add(this)),p.is0()}},G=new Point(Gx,Gy,1n,M(Gx*Gy)),I=new Point(0n,1n,1n,0n);Point.BASE=G;Point.ZERO=I;var numTo32bLE=num=>hexToBytes(padh(assertRange(num,0n,B256),64)).reverse(),bytesToNumberLE=b=>big("0x"+bytesToHex(u8fr(abytes(b)).reverse())),pow2=(x,power)=>{let r=x;for(;power-- >0n;)r=modP(r*r);return r},pow_2_252_3=x=>{let x2=modP(x*x),b2=modP(x2*x),b4=modP(pow2(b2,2n)*b2),b5=modP(pow2(b4,1n)*x),b10=modP(pow2(b5,5n)*b5),b20=modP(pow2(b10,10n)*b10),b40=modP(pow2(b20,20n)*b20),b80=modP(pow2(b40,40n)*b40),b160=modP(pow2(b80,80n)*b80),b240=modP(pow2(b160,80n)*b80),b250=modP(pow2(b240,10n)*b10);return{pow_p_5_8:modP(pow2(b250,2n)*x),b2}},RM1=0x2b8324804fc1df0b2b4d00993dfbd7a72f431806ad2fe478c4ee1b274a0ea0b0n,uvRatio=(u,v)=>{let v3=modP(v*modP(v*v)),v7=modP(modP(v3*v3)*v),pow=pow_2_252_3(modP(u*v7)).pow_p_5_8,x=modP(u*modP(v3*pow)),vx2=modP(v*modP(x*x)),root1=x,root2=modP(x*RM1),useRoot1=vx2===u,useRoot2=vx2===M(-u),noRoot=vx2===M(-u*RM1);return useRoot1&&(x=root1),(useRoot2||noRoot)&&(x=root2),(M(x)&1n)===1n&&(x=M(-x)),{isValid:useRoot1||useRoot2,value:x}},modL_LE=hash=>modN(bytesToNumberLE(hash)),sha512a=(...m)=>Promise.resolve(callHash("sha512Async")(concatBytes(...m))).then(checkDigest),sha512s=(...m)=>checkDigest(callHash("sha512")(concatBytes(...m))),hash2extK=hashed=>{let copy=u8fr(hashed),head=copy.slice(0,32);head[0]&=248,head[31]&=127,head[31]|=64;let prefix=copy.slice(32,64),scalar=modL_LE(head),point=G.multiply(scalar),pointBytes=point.toBytes();return{head,prefix,scalar,point,pointBytes}},getExtendedPublicKeyAsync=secretKey=>sha512a(abytes(secretKey,L)).then(hash2extK),getExtendedPublicKey=secretKey=>hash2extK(sha512s(abytes(secretKey,L))),getPublicKeyAsync=secretKey=>getExtendedPublicKeyAsync(secretKey).then(p=>p.pointBytes);var hashFinishA=res=>sha512a(res.hashable).then(res.finish);var _sign=(e,rBytes,msg)=>{let{pointBytes:P2,scalar:s}=e,r=modL_LE(rBytes),R=G.multiply(r).toBytes();return{hashable:concatBytes(R,P2,msg),finish:hashed=>{let S=modN(r+modL_LE(hashed)*s);return abytes(concatBytes(R,numTo32bLE(S)),64)}}},signAsync=async(message,secretKey)=>{let m=abytes(message),e=await getExtendedPublicKeyAsync(secretKey),rBytes=await sha512a(e.prefix,m);return hashFinishA(_sign(e,rBytes,m))};var hashes={sha512Async:async message=>{let s=subtle(),m=concatBytes(message);return u8n(await s.digest("SHA-512",m.buffer))},sha512:void 0},randomSecretKey=seed=>(seed=seed===void 0?randomBytes(L):seed,abytes(seed,L));var utils=Object.freeze({getExtendedPublicKeyAsync,getExtendedPublicKey,randomSecretKey}),W=8,scalarBits=256,pwindows=Math.ceil(scalarBits/W)+1,pwindowSize=2**(W-1),precompute=()=>{let points=[],p=G,b=p;for(let w=0;w{let n=p.negate();return cnd?n:p},wNAF=n=>{let comp=Gpows||(Gpows=precompute()),p=I,f=G,pow_2_w=2**W,maxNum=pow_2_w,mask=big(pow_2_w-1),shiftBy=big(W);for(let w=0;w>=shiftBy,wbits>pwindowSize&&(wbits-=maxNum,n+=1n);let off=w*pwindowSize,offF=off,offP=off+Math.abs(wbits)-1,isEven=w%2!==0,isNeg=wbits<0;wbits===0?f=f.add(ctneg(isEven,comp[offF])):p=p.add(ctneg(isNeg,comp[offP]))}return n!==0n&&err("invalid wnaf"),{p,f}};export{GATEWAY_CLIENT_CAPS,GATEWAY_CLIENT_IDS,GATEWAY_CLIENT_MODES,GatewayBrowserDeviceAuthLifecycle,GatewayProtocolClient,GatewayProtocolRequestError,MIN_CLIENT_PROTOCOL_VERSION,PROTOCOL_VERSION,utils as ed25519Utils,getPublicKeyAsync,signAsync}; +function normalizeDeviceMetadataForAuth(value){if(typeof value!="string")return"";let trimmed=value.trim();return trimmed?trimmed.replace(/[A-Z]/g,char=>String.fromCharCode(char.charCodeAt(0)+32)):""}function buildDeviceAuthPayloadV3(params){let scopes=params.scopes.join(","),token=params.token??"",platform=normalizeDeviceMetadataForAuth(params.platform),deviceFamily=normalizeDeviceMetadataForAuth(params.deviceFamily);return["v3",params.deviceId,params.clientId,params.clientMode,params.role,scopes,String(params.signedAtMs),token,params.nonce,platform,deviceFamily].join("|")}function normalized(value){return typeof value=="string"&&value.trim()||void 0}function selectGatewayConnectAuth(params){let authToken=normalized(params.token),bootstrapToken=normalized(params.bootstrapToken),explicitDeviceToken=normalized(params.deviceToken),authPassword=normalized(params.password),storedToken=normalized(params.storedToken),stored={storedToken,storedScopes:params.storedScopes};if(params.preferBootstrapToken&&bootstrapToken)return{authBootstrapToken:bootstrapToken,authPassword,...stored};let useRetryToken=params.pendingDeviceTokenRetry===!0&&!explicitDeviceToken&&!!(authToken&&storedToken&¶ms.trustedDeviceTokenRetry),resolvedDeviceToken=explicitDeviceToken??(useRetryToken||!(authToken||authPassword)&&(!bootstrapToken||storedToken)?storedToken:void 0),usingStoredDeviceToken=!!(resolvedDeviceToken&&!explicitDeviceToken&&storedToken)&&resolvedDeviceToken===storedToken,selectedToken=authToken??resolvedDeviceToken,authBootstrapToken=!authToken&&!resolvedDeviceToken&&!authPassword?bootstrapToken:void 0;return{authToken:selectedToken,authBootstrapToken,authDeviceToken:useRetryToken?storedToken:void 0,authPassword,authApprovalRuntimeToken:normalized(params.approvalRuntimeToken),authAgentRuntimeIdentityToken:normalized(params.agentRuntimeIdentityToken),signatureToken:selectedToken??authBootstrapToken,resolvedDeviceToken,usingStoredDeviceToken,...stored}}function buildGatewayConnectAuth(selected){let auth={token:selected.authToken,bootstrapToken:selected.authBootstrapToken,deviceToken:selected.authDeviceToken??selected.resolvedDeviceToken,password:selected.authPassword,approvalRuntimeToken:selected.authApprovalRuntimeToken,agentRuntimeIdentityToken:selected.authAgentRuntimeIdentityToken};return Object.values(auth).some(Boolean)?auth:void 0}function resolveGatewayConnectScopes(params){return params.requestedScopes??(params.usingStoredDeviceToken&¶ms.storedScopes?.length?params.storedScopes:[...params.defaultScopes])}var GatewayBrowserDeviceAuthLifecycle=class{constructor(deps){this.deps=deps}async buildPlan(params){let identity=await this.deps.loadIdentity(),stored=identity?await this.deps.tokenStore.load({clientId:params.client.id,deviceId:identity.deviceId,role:params.role}):null,storedValue=stored?.token,selectedAuth=selectGatewayConnectAuth({token:params.token,bootstrapToken:params.bootstrapToken,password:params.password,storedToken:storedValue,storedScopes:stored?.scopes,pendingDeviceTokenRetry:params.pendingDeviceTokenRetry,trustedDeviceTokenRetry:params.trustedDeviceTokenRetry,preferBootstrapToken:params.preferBootstrapToken}),{usingStoredDeviceToken}=selectedAuth,scopes=resolveGatewayConnectScopes({requestedScopes:selectedAuth.authBootstrapToken&¶ms.bootstrapScopes?[...params.bootstrapScopes]:void 0,usingStoredDeviceToken,storedScopes:selectedAuth.storedScopes,defaultScopes:params.defaultScopes});if(!identity)return{clientId:params.client.id,role:params.role,identity,selectedAuth,scopes,auth:buildGatewayConnectAuth(selectedAuth)};let signedAtMs=this.deps.nowMs?.()??Date.now(),nonce=params.nonce??"",{authBootstrapToken:primary,signatureToken:signed}=selectedAuth,token=null;primary?token=primary:signed&&(token=signed);let payload=buildDeviceAuthPayloadV3({deviceId:identity.deviceId,clientId:params.client.id,clientMode:params.client.mode,role:params.role,scopes,signedAtMs,token,nonce,platform:params.client.platform,deviceFamily:params.client.deviceFamily});return{clientId:params.client.id,role:params.role,identity,selectedAuth,scopes,auth:buildGatewayConnectAuth(selectedAuth),device:{id:identity.deviceId,publicKey:identity.publicKey,signature:await identity.sign(payload),signedAt:signedAtMs,nonce}}}async acceptHello(hello,plan){let token=hello.auth?.deviceToken?.trim();!token||!plan.identity||await this.deps.tokenStore.store({clientId:plan.clientId,deviceId:plan.identity.deviceId,role:hello.auth?.role??plan.role,token,scopes:hello.auth?.scopes??[]})}async clearStoredToken(plan){plan.identity&&await this.deps.tokenStore.clear({clientId:plan.clientId,deviceId:plan.identity.deviceId,role:plan.role})}};function isRecord(value){return!!value&&typeof value=="object"&&!Array.isArray(value)}function isNonEmptyString(value){return typeof value=="string"&&value.length>0}function isNonNegativeInteger(value){return typeof value=="number"&&Number.isInteger(value)&&value>=0}function isGatewayErrorShape(value){return!isRecord(value)||!isNonEmptyString(value.code)||!isNonEmptyString(value.message)||value.retryable!==void 0&&typeof value.retryable!="boolean"?!1:value.retryAfterMs===void 0||isNonNegativeInteger(value.retryAfterMs)}function isGatewayEventFrame(value){return!isRecord(value)||value.type!=="event"||!isNonEmptyString(value.event)?!1:value.seq===void 0||isNonNegativeInteger(value.seq)}function isGatewayResponseFrame(value){return!isRecord(value)||value.type!=="res"||!isNonEmptyString(value.id)||typeof value.ok!="boolean"?!1:value.error===void 0||isGatewayErrorShape(value.error)}function computeBackoff(policy,attempt){let base=Math.min(policy.maxMs,policy.initialMs*policy.factor**Math.max(attempt-1,0)),jitter=base*policy.jitter*Math.random();return Math.min(policy.maxMs,Math.round(base+jitter))}async function sleepWithAbort(ms,abortSignal,options={}){if(!Number.isFinite(ms)||ms<=0)return;let delayMs=Math.min(Math.max(Math.floor(ms),1),2147e6);await new Promise((resolve,reject)=>{let settled=!1,timer=null,cleanup=()=>abortSignal?.removeEventListener("abort",onAbort),onAbort=()=>{settled||(settled=!0,timer&&clearTimeout(timer),timer=null,cleanup(),reject(new Error("aborted",{cause:abortSignal?.reason??new Error("aborted")})))};if(abortSignal?.addEventListener("abort",onAbort,{once:!0}),abortSignal?.aborted){onAbort();return}timer=setTimeout(()=>{settled=!0,cleanup(),timer=null,resolve()},delayMs),options.ref===!1&&timer.unref?.(),abortSignal?.aborted&&onAbort()})}var RetrySupervisor=class{constructor(policy,maxAttempts=Number.POSITIVE_INFINITY){this.policy=policy;this.maxAttempts=maxAttempts;this.attempts=0;this.initialMs=policy.initialMs}reset(initialMs=this.policy.initialMs){this.cancel(),this.attempts=0,this.initialMs=initialMs,this.nextDelayOverrideMs=void 0}cancel(reason=new Error("retry cancelled")){this.pendingAbort?.abort(reason),this.pendingAbort=void 0}next(abortSignal){let override=this.nextDelayOverrideMs;if(this.nextDelayOverrideMs=void 0,override===void 0&&++this.attempts>Math.ceil(this.maxAttempts))return;let attempt=Math.max(this.attempts,1),delayMs=override??computeBackoff({...this.policy,initialMs:this.initialMs},attempt);this.cancel();let pendingAbort=new AbortController;return this.pendingAbort=pendingAbort,{attempt,delayMs,signal:abortSignal?AbortSignal.any([pendingAbort.signal,abortSignal]):pendingAbort.signal}}},DEFAULT_RETRY_CONFIG={attempts:3,minDelayMs:300,maxDelayMs:3e4,jitter:0},defaultSleep=ms=>new Promise(resolve=>{setTimeout(resolve,ms)});function asFiniteNumber(value){return typeof value=="number"&&Number.isFinite(value)?value:void 0}function clampNumber(value,fallback,min,max){let next=asFiniteNumber(value);return next===void 0?fallback:Math.min(Math.max(next,min??Number.NEGATIVE_INFINITY),max??Number.POSITIVE_INFINITY)}function resolveAttemptCount(value,fallback){return Math.max(1,Math.round(asFiniteNumber(value)??fallback))}function resolveRetryDelayMs(value){let finite=value===Number.POSITIVE_INFINITY?2147e6:asFiniteNumber(value)??0;return Math.min(Math.max(Math.round(finite),0),2147e6)}function resolveJitterConfig(value,fallback){if(value==="full")return"full";let fraction=asFiniteNumber(value);return fraction===void 0?fallback:Math.min(Math.max(fraction,0),1)}function resolveRetryConfig(defaults=DEFAULT_RETRY_CONFIG,overrides){let attempts=resolveAttemptCount(overrides?.attempts,defaults.attempts),minDelayMs=resolveRetryDelayMs(clampNumber(overrides?.minDelayMs,defaults.minDelayMs,0)),maxDelayMs=Math.max(minDelayMs,resolveRetryDelayMs(clampNumber(overrides?.maxDelayMs,defaults.maxDelayMs,0)));return{attempts,minDelayMs,maxDelayMs,jitter:resolveJitterConfig(overrides?.jitter,defaults.jitter)}}function applyJitter(delayMs,jitter,mode,random){if(jitter==="full")return mode==="symmetric"?Math.max(0,Math.round(delayMs*(.5+random()*.5))):Math.max(0,Math.ceil(delayMs*(1+random())));if(jitter<=0)return mode==="positive"?Math.ceil(delayMs):delayMs;let fraction=random(),offset=mode==="positive"?fraction*jitter:(fraction*2-1)*jitter,raw=delayMs*(1+offset);return Math.max(0,mode==="positive"?Math.ceil(raw):Math.round(raw))}function toRetryError(value,fallbackMessage="Non-Error thrown"){if(value instanceof Error)return value;if(typeof value=="string")return new Error(value);let error=new Error(fallbackMessage,{cause:value});return(typeof value=="object"&&value!==null||typeof value=="function")&&Object.assign(error,value),error}function createRetryRunner(runtime={}){let runtimeSleep=runtime.sleep??defaultSleep,runtimeRandom=runtime.random??Math.random,createFailure=runtime.createFailure??(errors=>toRetryError(errors.at(-1)??new Error("Retry failed")));return async function(fn,attemptsOrOptions=3,initialDelayMs=300){let attemptErrors=[];if(typeof attemptsOrOptions=="number"){let attempts=resolveAttemptCount(attemptsOrOptions,DEFAULT_RETRY_CONFIG.attempts);for(let index=0;index0?resolved.maxDelayMs:Number.POSITIVE_INFINITY,retryAfterMaxDelayMs=options.retryAfterMaxDelayMs===void 0?maxDelayMs:Math.max(minDelayMs,resolveRetryDelayMs(clampNumber(options.retryAfterMaxDelayMs,maxDelayMs,0))),random=options.random??runtimeRandom,sleep=options.sleep??runtimeSleep,shouldRetry=options.shouldRetry??(()=>!0);for(let attempt=1;attempt<=maxAttempts;attempt+=1)try{return await fn()}catch(err2){if(attemptErrors.push(err2),attempt>=maxAttempts||!shouldRetry(err2,attempt))break;let context={attempt,maxAttempts,err:err2,label:options.label},retryAfterMs=options.retryAfterMs?.(err2),hasRetryAfter=typeof retryAfterMs=="number"&&Number.isFinite(retryAfterMs),configuredDelay=typeof options.delayMs=="function"?options.delayMs(context):options.delayMs,resolvedConfiguredDelay=configuredDelay===void 0?void 0:resolveRetryDelayMs(configuredDelay),baseDelay=hasRetryAfter?Math.max(retryAfterMs,minDelayMs):resolvedConfiguredDelay===void 0?minDelayMs*2**(attempt-1):Math.max(resolvedConfiguredDelay,minDelayMs),delayCap=hasRetryAfter?retryAfterMaxDelayMs:maxDelayMs,delay=Math.min(baseDelay,delayCap),canHonorRetryAfter=hasRetryAfter&&(retryAfterMs??0)<=delayCap,wantsPositiveDraw=resolved.jitter==="full"&&!hasRetryAfter||canHonorRetryAfter;delay=applyJitter(delay,resolved.jitter,wantsPositiveDraw?"positive":"symmetric",random),delay=Math.min(Math.max(delay,minDelayMs),delayCap),await options.onRetry?.({...context,delayMs:delay}),delay>0&&await sleep(delay)}throw createFailure(attemptErrors)}}var retryAsync=createRetryRunner();var GatewayEventListeners=class{constructor(){this.listeners=new Map}add(listener){let subscription=this.listeners.get(listener)??{};return this.listeners.set(listener,subscription),()=>{this.listeners.get(listener)===subscription&&this.listeners.delete(listener)}}snapshot(){return[...this.listeners]}isCurrent(listener,subscription){return this.listeners.get(listener)===subscription}};var DEFAULT_PREAUTH_HANDSHAKE_TIMEOUT_MS=15e3;function startGatewayConnectTimeout(onTimeout){let timer=setTimeout(onTimeout,DEFAULT_PREAUTH_HANDSHAKE_TIMEOUT_MS);return timer.unref?.(),timer}function clearGatewayConnectTimeout(timer){return timer!==null&&clearTimeout(timer),null}var GatewayProtocolRequestError=class extends Error{constructor(error){super(error.message??"request failed"),this.name="GatewayProtocolRequestError",this.code=error.code??"UNAVAILABLE",this.gatewayCode=this.code,this.details=error.details,this.retryable=error.retryable===!0,this.retryAfterMs=error.retryAfterMs}},GatewayProtocolClient=class{constructor(opts){this.opts=opts;this.socket=null;this.pending=new Map;this.listeners=new GatewayEventListeners;this.stopped=!0;this.generation=0;this.lastSeq=null;this.connectNonce=null;this.connectSent=!1;this.connectRequestSent=!1;this.handshakeTimer=null;this.reconnectSignal=null;this.socketOpened=!1;this.helloReceived=!1;this.connectTiming=null;this.reconnectSupervisor=new RetrySupervisor({initialMs:opts.reconnect.initialMs,maxMs:opts.reconnect.maxMs,factor:opts.reconnect.multiplier,jitter:0})}get connected(){return this.socket?.isOpen()??!1}get hasPendingRequests(){return this.pending.size>0}get connecting(){return this.connectSent&&!this.helloReceived}get hasUnboundedPendingRequests(){return[...this.pending.values()].some(pending=>pending.unbounded)}start(){this.socket||this.reconnectSignal||(this.stopped=!1,this.reconnectSupervisor.cancel(),this.connect())}stop(){this.stopped=!0,this.clearHandshakeTimer(),this.reconnectSignal=null,this.reconnectSupervisor.reset();let socket=this.socket;socket&&this.opts.notifyStoppedClose&&(this.stoppedSocket={socket,context:this.closeContext()}),this.socket=null,this.connectFailure=void 0,this.connectTiming=null,this.flushRequests(new Error("gateway client stopped")),socket?.close()}request(method,params,options){let socket=this.socket;if(!socket?.isOpen())return Promise.reject(new Error("gateway not connected"));if(typeof method!="string"||method.length===0)return Promise.reject(new Error("invalid request frame: method must be a non-empty string"));let id=this.opts.createRequestId(),timeoutMs=options?.timeoutMs===null?void 0:options?.timeoutMs??this.opts.requestTimeoutMs;return new Promise((resolve,reject)=>{let timeout,pending={resolve:value=>resolve(value),reject,expectFinal:options?.expectFinal===!0,acceptedNotified:!1,onAccepted:options?.onAccepted,unbounded:timeoutMs===void 0,method,startedAtMs:this.nowMs()},onAbort=()=>{this.pending.delete(id),timeout&&clearTimeout(timeout),this.finishRequestTiming(id,pending,!1,"CLIENT_ABORTED"),reject(this.opts.createRequestAbortError?.(method)??new Error(`gateway request aborted for ${method}`))},cleanup=()=>{timeout&&clearTimeout(timeout),options?.signal?.removeEventListener("abort",onAbort)};if(options?.signal?.aborted){reject(this.opts.createRequestAbortError?.(method)??new Error(`gateway request aborted for ${method}`));return}pending.cleanup=cleanup,timeoutMs!==void 0&&timeoutMs>=0&&(timeout=setTimeout(()=>{this.pending.delete(id),options?.signal?.removeEventListener("abort",onAbort),this.finishRequestTiming(id,pending,!1,"CLIENT_TIMEOUT"),reject(this.opts.createRequestTimeoutError?.(method,timeoutMs)??new Error(`gateway request timed out after ${timeoutMs}ms: ${method}`))},timeoutMs),timeout.unref?.()),options?.signal?.addEventListener("abort",onAbort,{once:!0}),this.pending.set(id,pending);try{socket.send(JSON.stringify({type:"req",id,method,params})),this.invoke("sent",()=>options?.onSent?.())}catch(error){this.pending.delete(id),cleanup(),this.finishRequestTiming(id,pending,!1,"CLIENT_SEND_ERROR"),reject(error instanceof Error?error:new Error(String(error)))}})}addEventListener(listener){return this.listeners.add(listener)}closeSocket(code,reason){this.socket?.close(code,reason)}resetReconnectBackoff(initialMs){this.reconnectSignal=null,this.reconnectSupervisor.reset(initialMs)}recordTiming(phase,generation,plan,detail){let now=this.nowMs(),state=this.connectTiming;!state||state.generation!==generation||(state.hasChallenge||=phase==="challenge",state.usedFallback||=phase==="fallback",this.invoke("connect timing",()=>this.opts.onTiming?.({phase,generation,durationMs:Math.max(0,now-state.startedAtMs),phaseDurationMs:Math.max(0,now-state.lastAtMs),hasChallenge:state.hasChallenge,usedFallback:state.usedFallback,plan,detail})),state.lastAtMs=now,(phase==="hello"||phase==="failed")&&(this.connectTiming=null))}connect(){if(this.stopped)return;let generation=this.generation+1;this.connectNonce=null,this.connectSent=!1,this.connectRequestSent=!1,this.socketOpened=!1,this.helloReceived=!1,this.connectFailure=void 0;let socket;try{socket=this.opts.createSocket({open:()=>this.handleOpen(socket,generation),message:data=>this.handleMessage(socket,generation,data),close:(code,reason)=>this.handleClose(socket,generation,code,reason),error:error=>this.handleSocketError(socket,generation,error)})}catch(error){let normalized2=error instanceof Error?error:new Error(String(error));if(this.opts.onSocketFactoryError?.(normalized2),this.opts.onConnectError?.(normalized2),this.opts.rethrowSocketFactoryError?.(normalized2))throw normalized2;this.opts.shouldRetrySocketFactoryError?.(normalized2)&&!this.stopped&&!this.socket&&!this.reconnectSignal&&this.scheduleReconnect();return}this.generation=generation,this.socket=socket;let now=this.nowMs();this.connectTiming={generation,startedAtMs:now,lastAtMs:now,hasChallenge:!1,usedFallback:!1}}handleOpen(socket,generation){if(this.isActive(socket,generation)){if(this.socketOpened=!0,this.recordTiming("socket-open",generation),this.connectNonce){this.sendConnect(socket,generation);return}this.armHandshakeTimer(socket,generation)}}armHandshakeTimer(socket,generation){this.clearHandshakeTimer();let armedAt=Date.now();this.handshakeTimer=setTimeout(()=>{if(this.handshakeTimer=null,!this.isActive(socket,generation)||this.connectSent||!socket.isOpen())return;if(this.opts.handshake.mode==="fallback"){this.recordTiming("fallback",generation),this.sendConnect(socket,generation);return}let elapsedMs=Date.now()-armedAt,error=new Error(this.opts.handshake.timeoutMessage?.(elapsedMs)??`gateway connect challenge timeout after ${elapsedMs}ms`);this.opts.onConnectError?.(error),socket.close(1008,"connect challenge timeout")},this.opts.handshake.timeoutMs),this.handshakeTimer.unref?.()}sendConnect(socket,generation){if(!this.isActive(socket,generation)||!socket.isOpen()||this.connectSent)return;this.connectSent=!0,this.clearHandshakeTimer(),this.handshakeTimer=startGatewayConnectTimeout(()=>{this.isActive(socket,generation)&&!this.helloReceived&&socket.close(4e3,"connect timeout")});let planOrPromise;try{planOrPromise=this.opts.buildConnectPlan({nonce:this.connectNonce,generation})}catch(error){this.handleConnectPlanError(socket,generation,error);return}if(planOrPromise instanceof Promise){planOrPromise.then(plan=>this.sendConnectPlan(socket,generation,plan)).catch(error=>this.handleConnectPlanError(socket,generation,error));return}this.sendConnectPlan(socket,generation,planOrPromise)}handleConnectPlanError(socket,generation,error){if(!this.isActive(socket,generation))return;let normalized2=error instanceof Error?error:new Error(String(error)),outcome=this.opts.onConnectPlanError?.(normalized2)??{closeCode:1008,closeReason:"connect failed"};this.opts.onConnectError?.(outcome.error??normalized2),outcome.stop&&(this.stopped=!0),socket.close(outcome.closeCode,outcome.closeReason)}sendConnectPlan(socket,generation,plan){if(!this.isActive(socket,generation)||!socket.isOpen())return;let context={generation,nonce:this.connectNonce,plan};this.recordTiming("connect-plan-ready",generation,plan),this.recordTiming("request-sent",generation,plan),this.connectRequestSent=!0,this.request("connect",this.opts.buildConnectParams(plan)).then(hello=>{this.isActive(socket,generation)&&(this.helloReceived=!0,this.clearHandshakeTimer(),this.connectFailure=void 0,this.reconnectSupervisor.reset(),this.recordTiming("hello",generation,plan),this.opts.onConnectHello?.(hello,context),this.invoke("hello",()=>this.opts.onHello?.(hello)))}).catch(error=>{if(!this.isActive(socket,generation))return;let requestError=error instanceof GatewayProtocolRequestError?error:new GatewayProtocolRequestError({message:String(error)}),outcome=this.opts.onConnectFailure?.(requestError,context)??{closeCode:1008,closeReason:"connect failed"};this.connectFailure={error:requestError,reconnectDelayMs:outcome.reconnectDelayMs},outcome.stop&&(this.stopped=!0),socket.close(outcome.closeCode,outcome.closeReason)})}handleMessage(socket,generation,raw){if(!this.isActive(socket,generation))return;let parsed;try{parsed=JSON.parse(raw)}catch(error){this.opts.onParseError?.(error);return}if(isGatewayEventFrame(parsed)){if(this.opts.onActivity?.(),parsed.event==="connect.challenge"){let payload=parsed.payload,nonce=typeof payload?.nonce=="string"?payload.nonce.trim():"";if(!nonce){if(this.opts.handshake.mode==="require-challenge"){let error=new Error("gateway connect challenge missing nonce");this.opts.onConnectError?.(error),socket.close(1008,"connect challenge missing nonce")}return}this.connectNonce=nonce,this.recordTiming("challenge",generation),this.sendConnect(socket,generation);return}let seq=typeof parsed.seq=="number"?parsed.seq:null;if(seq!==null){if(this.lastSeq!==null&&seq>this.lastSeq+1){let expected=this.lastSeq+1;if(this.invoke("gap",()=>this.opts.onGap?.({expected,received:seq})),!this.isActive(socket,generation))return}this.lastSeq=seq}let listeners=this.listeners.snapshot();this.invoke("event",()=>this.opts.onEvent?.(parsed));for(let[listener,subscription]of listeners){if(!this.isActive(socket,generation))return;this.listeners.isCurrent(listener,subscription)&&this.invoke("event listener",()=>listener(parsed))}return}isGatewayResponseFrame(parsed)&&(this.opts.onActivity?.(),this.handleResponse(parsed))}handleResponse(frame){let pending=this.pending.get(frame.id);if(!pending)return;let status=frame.payload?.status;if(pending.expectFinal&&status==="accepted"){pending.acceptedNotified||(pending.acceptedNotified=!0,this.invoke("accepted",()=>pending.onAccepted?.(frame.payload)));return}if(this.pending.delete(frame.id),pending.cleanup?.(),frame.ok){this.finishRequestTiming(frame.id,pending,!0),pending.resolve(frame.payload);return}this.finishRequestTiming(frame.id,pending,!1,frame.error?.code),pending.reject(this.opts.createRequestError?.(frame.error??{})??new GatewayProtocolRequestError(frame.error??{}))}handleClose(socket,generation,code,reason){if(this.socket!==socket){if(this.stoppedSocket?.socket===socket){let context2={...this.stoppedSocket.context,code,reason};this.stoppedSocket=void 0,this.invoke("close",()=>this.opts.onClose?.(context2,{retry:!1,notify:!0}))}return}this.socket=null,this.clearHandshakeTimer();let context={...this.closeContext(),code,reason,generation};this.connectFailure=void 0;let decision=this.opts.resolveClose(context);this.flushRequests(decision.pendingError??context.connectFailure?.error??new Error(`gateway closed (${code}): ${reason}`)),this.invoke("close",()=>this.opts.onClose?.(context,decision)),decision.retry&&!this.stopped&&this.scheduleReconnect(decision.reconnectDelayMs??context.connectFailure?.reconnectDelayMs)}handleSocketError(socket,generation,error){!this.isActive(socket,generation)||this.connectSent||this.opts.onConnectError?.(error)}flushRequests(error){for(let[id,pending]of this.pending)this.finishRequestTiming(id,pending,!1,"CLIENT_CLOSED"),pending.cleanup?.(),pending.reject(error);this.pending.clear()}finishRequestTiming(id,pending,ok,errorCode){let endedAtMs=this.nowMs();this.invoke("request timing",()=>this.opts.onRequestTiming?.({id,method:pending.method,ok,durationMs:Math.max(0,endedAtMs-pending.startedAtMs),startedAtMs:pending.startedAtMs,endedAtMs,errorCode}))}scheduleReconnect(overrideMs){overrideMs!==void 0&&(this.reconnectSupervisor.nextDelayOverrideMs=overrideMs);let retry=this.reconnectSupervisor.next();retry&&(this.reconnectSignal=retry.signal,sleepWithAbort(retry.delayMs,retry.signal).then(()=>{this.reconnectSignal===retry.signal&&(this.reconnectSignal=null,this.connect())},()=>{this.reconnectSignal===retry.signal&&(this.reconnectSignal=null)}))}closeContext(){return{generation:this.generation,socketOpened:this.socketOpened,helloReceived:this.helloReceived,connectRequestSent:this.connectRequestSent,connectFailure:this.connectFailure}}isActive(socket,generation){return!this.stopped&&this.socket===socket&&this.generation===generation}nowMs(){return this.opts.nowMs?.()??Date.now()}clearHandshakeTimer(){this.handshakeTimer=clearGatewayConnectTimeout(this.handshakeTimer)}invoke(label,callback){try{callback()}catch(error){this.opts.onCallbackError?.(label,error)}}};var GATEWAY_CLIENT_IDS={WEBCHAT_UI:"webchat-ui",CONTROL_UI:"openclaw-control-ui",BROWSER_COPILOT:"openclaw-browser-copilot",TUI:"openclaw-tui",WEBCHAT:"webchat",CLI:"cli",GATEWAY_CLIENT:"gateway-client",MACOS_APP:"openclaw-macos",LINUX_APP:"openclaw-linux",IOS_APP:"openclaw-ios",WATCHOS_APP:"openclaw-watchos",ANDROID_APP:"openclaw-android",NODE_HOST:"node-host",WORKER:"openclaw-worker",TEST:"test",FINGERPRINT:"fingerprint",PROBE:"openclaw-probe"};var GATEWAY_CLIENT_MODES={WEBCHAT:"webchat",CLI:"cli",UI:"ui",BACKEND:"backend",NODE:"node",WORKER:"worker",PROBE:"probe",TEST:"test"},GATEWAY_CLIENT_CAPS={AGENT_KIND:"agent-kind",APPROVALS:"approvals",EXEC_APPROVALS:"exec-approvals",INLINE_WIDGETS:"inline-widgets",RUN_TOOL_BINDINGS:"run-tool-bindings",SESSION_SCOPED_EVENTS:"session-scoped-events",PLUGIN_APPROVALS:"plugin-approvals",TASK_SUGGESTIONS:"task-suggestions",TERMINAL_OFFSET_SEQ:"terminal-offset-seq",TOOL_EVENTS:"tool-events",UI_COMMANDS:"ui-commands"},GATEWAY_CLIENT_ID_SET=new Set(Object.values(GATEWAY_CLIENT_IDS)),GATEWAY_CLIENT_MODE_SET=new Set(Object.values(GATEWAY_CLIENT_MODES));var PROTOCOL_VERSION=4,MIN_CLIENT_PROTOCOL_VERSION=4;/*! noble-ed25519 - MIT License (c) 2019 Paul Miller (paulmillr.com) */var ed25519_CURVE=Object.freeze({p:0x7fffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffedn,n:0x1000000000000000000000000000000014def9dea2f79cd65812631a5cf5d3edn,h:8n,a:0x7fffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffecn,d:0x52036cee2b6ffe738cc740797779e89800700a4d4141d8ab75eb4dca135978a3n,Gx:0x216936d3cd6e53fec0a4e231fdd6dc5c692cc7609525a7b2c9562d608f25d51an,Gy:0x6666666666666666666666666666666666666666666666666666666666666658n}),{p:P,n:N,Gx,Gy,a:_a,d:_d,h}=ed25519_CURVE,L=32,captureTrace=(...args)=>{"captureStackTrace"in Error&&typeof Error.captureStackTrace=="function"&&Error.captureStackTrace(...args)},err=(message="")=>{let e=new Error(message);throw captureTrace(e,err),e},isBig=n=>typeof n=="bigint",isStr=s=>typeof s=="string",isBytes=a=>a instanceof Uint8Array||ArrayBuffer.isView(a)&&a.constructor.name==="Uint8Array"&&"BYTES_PER_ELEMENT"in a&&a.BYTES_PER_ELEMENT===1,abytes=(value,length,title="")=>{let bytes=isBytes(value),len=value?.length,needsLen=length!==void 0;if(!bytes||needsLen&&len!==length){let prefix=title&&`"${title}" `,ofLen=needsLen?` of length ${length}`:"",got=bytes?`length=${len}`:`type=${typeof value}`,msg=prefix+"expected Uint8Array"+ofLen+", got "+got;throw bytes?new RangeError(msg):new TypeError(msg)}return value},u8n=len=>new Uint8Array(len),u8fr=buf=>Uint8Array.from(buf),padh=(n,pad)=>n.toString(16).padStart(pad,"0"),bytesToHex=b=>Array.from(abytes(b)).map(e=>padh(e,2)).join(""),C={_0:48,_9:57,A:65,F:70,a:97,f:102},_ch=ch=>{if(ch>=C._0&&ch<=C._9)return ch-C._0;if(ch>=C.A&&ch<=C.F)return ch-(C.A-10);if(ch>=C.a&&ch<=C.f)return ch-(C.a-10)},hexToBytes=hex=>{let e="hex invalid";if(!isStr(hex))return err(e);let hl=hex.length,al=hl/2;if(hl%2)return err(e);let array=u8n(al);for(let ai=0,hi=0;aiglobalThis?.crypto,subtle=()=>cr()?.subtle??err("crypto.subtle must be defined, consider polyfill"),concatBytes=(...arrs)=>{let len=0;for(let a of arrs)len+=abytes(a).length;let r=u8n(len),pad=0;return arrs.forEach(a=>{r.set(a,pad),pad+=a.length}),r},randomBytes=(len=L)=>cr().getRandomValues(u8n(len)),big=BigInt,assertRange=(n,min,max,msg="bad number: out of range")=>{if(!isBig(n))throw new TypeError(msg);if(min<=n&&n{let r=a%b;return r>=0n?r:b+r},P_MASK=(1n<<255n)-1n,modP=num=>{num<0n&&err("negative coordinate");let r=(num>>255n)*19n+(num&P_MASK);return r=(r>>255n)*19n+(r&P_MASK),r%P},modN=a=>M(a,N),invert=(num,md)=>{(num===0n||md<=0n)&&err("no inverse n="+num+" mod="+md);let a=M(num,md),b=md,x=0n,y=1n,u=1n,v=0n;for(;a!==0n;){let q=b/a,r=b%a,m=x-u*q,n=y-v*q;b=a,a=r,x=u,y=v,u=m,v=n}return b===1n?M(x,md):err("no inverse")},callHash=name=>{let fn=hashes[name];return typeof fn!="function"&&err("hashes."+name+" not set"),fn},checkDigest=value=>abytes(value,64,"digest");var apoint=p=>p instanceof Point?p:err("Point expected"),B256=2n**256n,Point=class _Point{static BASE;static ZERO;X;Y;Z;T;constructor(X,Y,Z,T){let max=B256;this.X=assertRange(X,0n,max),this.Y=assertRange(Y,0n,max),this.Z=assertRange(Z,1n,max),this.T=assertRange(T,0n,max),Object.freeze(this)}static CURVE(){return ed25519_CURVE}static fromAffine(p){return new _Point(p.x,p.y,1n,modP(p.x*p.y))}static fromBytes(hex,zip215=!1){let d=_d,normed=u8fr(abytes(hex,L)),lastByte=hex[31];normed[31]=lastByte&-129;let y=bytesToNumberLE(normed);assertRange(y,0n,zip215?B256:P);let y2=modP(y*y),u=M(y2-1n),v=modP(d*y2+1n),{isValid,value:x}=uvRatio(u,v);isValid||err("bad point: y not sqrt");let isXOdd=(x&1n)===1n,isLastByteOdd=(lastByte&128)!==0;return!zip215&&x===0n&&isLastByteOdd&&err("bad point: x==0, isLastByteOdd"),isLastByteOdd!==isXOdd&&(x=M(-x)),new _Point(x,y,1n,modP(x*y))}static fromHex(hex,zip215){return _Point.fromBytes(hexToBytes(hex),zip215)}get x(){return this.toAffine().x}get y(){return this.toAffine().y}assertValidity(){let a=_a,d=_d,p=this;if(p.is0())return err("bad point: ZERO");let{X,Y,Z,T}=p,X2=modP(X*X),Y2=modP(Y*Y),Z2=modP(Z*Z),Z4=modP(Z2*Z2),aX2=modP(X2*a),left=modP(Z2*(aX2+Y2)),right=M(Z4+modP(d*modP(X2*Y2)));if(left!==right)return err("bad point: equation left != right (1)");let XY=modP(X*Y),ZT=modP(Z*T);return XY!==ZT?err("bad point: equation left != right (2)"):this}equals(other){let{X:X1,Y:Y1,Z:Z1}=this,{X:X2,Y:Y2,Z:Z2}=apoint(other),X1Z2=modP(X1*Z2),X2Z1=modP(X2*Z1),Y1Z2=modP(Y1*Z2),Y2Z1=modP(Y2*Z1);return X1Z2===X2Z1&&Y1Z2===Y2Z1}is0(){return this.equals(I)}negate(){return new _Point(M(-this.X),this.Y,this.Z,M(-this.T))}double(){let{X:X1,Y:Y1,Z:Z1}=this,a=_a,A=modP(X1*X1),B=modP(Y1*Y1),C2=modP(2n*Z1*Z1),D=modP(a*A),x1y1=M(X1+Y1),E=M(modP(x1y1*x1y1)-A-B),G2=M(D+B),F=M(G2-C2),H=M(D-B),X3=modP(E*F),Y3=modP(G2*H),T3=modP(E*H),Z3=modP(F*G2);return new _Point(X3,Y3,Z3,T3)}add(other){let{X:X1,Y:Y1,Z:Z1,T:T1}=this,{X:X2,Y:Y2,Z:Z2,T:T2}=apoint(other),a=_a,d=_d,A=modP(X1*X2),B=modP(Y1*Y2),C2=modP(modP(T1*d)*T2),D=modP(Z1*Z2),E=M(modP(M(X1+Y1)*M(X2+Y2))-A-B),F=M(D-C2),G2=M(D+C2),H=M(B-modP(a*A)),X3=modP(E*F),Y3=modP(G2*H),T3=modP(E*H),Z3=modP(F*G2);return new _Point(X3,Y3,Z3,T3)}subtract(other){return this.add(apoint(other).negate())}multiply(n,safe=!0){if(!safe&&n===0n||(assertRange(n,1n,N),!safe&&this.is0()))return I;if(n===1n)return this;if(this.equals(G))return wNAF(n).p;let p=I,f=G;for(let d=this;n>0n;d=d.double(),n>>=1n)n&1n?p=p.add(d):safe&&(f=f.add(d));return p}multiplyUnsafe(scalar){return this.multiply(scalar,!1)}toAffine(){let{X,Y,Z}=this;if(this.equals(I))return{x:0n,y:1n};let iz=invert(Z,P);modP(Z*iz)!==1n&&err("invalid inverse");let x=modP(X*iz),y=modP(Y*iz);return{x,y}}toBytes(){let{x,y}=this.toAffine(),b=numTo32bLE(y);return b[31]|=x&1n?128:0,b}toHex(){return bytesToHex(this.toBytes())}clearCofactor(){return this.multiply(big(h),!1)}isSmallOrder(){return this.clearCofactor().is0()}isTorsionFree(){let p=this.multiply(N/2n,!1).double();return N%2n&&(p=p.add(this)),p.is0()}},G=new Point(Gx,Gy,1n,M(Gx*Gy)),I=new Point(0n,1n,1n,0n);Point.BASE=G;Point.ZERO=I;var numTo32bLE=num=>hexToBytes(padh(assertRange(num,0n,B256),64)).reverse(),bytesToNumberLE=b=>big("0x"+bytesToHex(u8fr(abytes(b)).reverse())),pow2=(x,power)=>{let r=x;for(;power-- >0n;)r=modP(r*r);return r},pow_2_252_3=x=>{let x2=modP(x*x),b2=modP(x2*x),b4=modP(pow2(b2,2n)*b2),b5=modP(pow2(b4,1n)*x),b10=modP(pow2(b5,5n)*b5),b20=modP(pow2(b10,10n)*b10),b40=modP(pow2(b20,20n)*b20),b80=modP(pow2(b40,40n)*b40),b160=modP(pow2(b80,80n)*b80),b240=modP(pow2(b160,80n)*b80),b250=modP(pow2(b240,10n)*b10);return{pow_p_5_8:modP(pow2(b250,2n)*x),b2}},RM1=0x2b8324804fc1df0b2b4d00993dfbd7a72f431806ad2fe478c4ee1b274a0ea0b0n,uvRatio=(u,v)=>{let v3=modP(v*modP(v*v)),v7=modP(modP(v3*v3)*v),pow=pow_2_252_3(modP(u*v7)).pow_p_5_8,x=modP(u*modP(v3*pow)),vx2=modP(v*modP(x*x)),root1=x,root2=modP(x*RM1),useRoot1=vx2===u,useRoot2=vx2===M(-u),noRoot=vx2===M(-u*RM1);return useRoot1&&(x=root1),(useRoot2||noRoot)&&(x=root2),(M(x)&1n)===1n&&(x=M(-x)),{isValid:useRoot1||useRoot2,value:x}},modL_LE=hash=>modN(bytesToNumberLE(hash)),sha512a=(...m)=>Promise.resolve(callHash("sha512Async")(concatBytes(...m))).then(checkDigest),sha512s=(...m)=>checkDigest(callHash("sha512")(concatBytes(...m))),hash2extK=hashed=>{let copy=u8fr(hashed),head=copy.slice(0,32);head[0]&=248,head[31]&=127,head[31]|=64;let prefix=copy.slice(32,64),scalar=modL_LE(head),point=G.multiply(scalar),pointBytes=point.toBytes();return{head,prefix,scalar,point,pointBytes}},getExtendedPublicKeyAsync=secretKey=>sha512a(abytes(secretKey,L)).then(hash2extK),getExtendedPublicKey=secretKey=>hash2extK(sha512s(abytes(secretKey,L))),getPublicKeyAsync=secretKey=>getExtendedPublicKeyAsync(secretKey).then(p=>p.pointBytes);var hashFinishA=res=>sha512a(res.hashable).then(res.finish);var _sign=(e,rBytes,msg)=>{let{pointBytes:P2,scalar:s}=e,r=modL_LE(rBytes),R=G.multiply(r).toBytes();return{hashable:concatBytes(R,P2,msg),finish:hashed=>{let S=modN(r+modL_LE(hashed)*s);return abytes(concatBytes(R,numTo32bLE(S)),64)}}},signAsync=async(message,secretKey)=>{let m=abytes(message),e=await getExtendedPublicKeyAsync(secretKey),rBytes=await sha512a(e.prefix,m);return hashFinishA(_sign(e,rBytes,m))};var hashes={sha512Async:async message=>{let s=subtle(),m=concatBytes(message);return u8n(await s.digest("SHA-512",m.buffer))},sha512:void 0},randomSecretKey=seed=>(seed=seed===void 0?randomBytes(L):seed,abytes(seed,L));var utils=Object.freeze({getExtendedPublicKeyAsync,getExtendedPublicKey,randomSecretKey}),W=8,scalarBits=256,pwindows=Math.ceil(scalarBits/W)+1,pwindowSize=2**(W-1),precompute=()=>{let points=[],p=G,b=p;for(let w=0;w{let n=p.negate();return cnd?n:p},wNAF=n=>{let comp=Gpows||(Gpows=precompute()),p=I,f=G,pow_2_w=2**W,maxNum=pow_2_w,mask=big(pow_2_w-1),shiftBy=big(W);for(let w=0;w>=shiftBy,wbits>pwindowSize&&(wbits-=maxNum,n+=1n);let off=w*pwindowSize,offF=off,offP=off+Math.abs(wbits)-1,isEven=w%2!==0,isNeg=wbits<0;wbits===0?f=f.add(ctneg(isEven,comp[offF])):p=p.add(ctneg(isNeg,comp[offP]))}return n!==0n&&err("invalid wnaf"),{p,f}};export{GATEWAY_CLIENT_CAPS,GATEWAY_CLIENT_IDS,GATEWAY_CLIENT_MODES,GatewayBrowserDeviceAuthLifecycle,GatewayProtocolClient,GatewayProtocolRequestError,MIN_CLIENT_PROTOCOL_VERSION,PROTOCOL_VERSION,utils as ed25519Utils,getPublicKeyAsync,signAsync}; diff --git a/packages/gateway-client/src/browser.ts b/packages/gateway-client/src/browser.ts index bba095c467b6..791146902982 100644 --- a/packages/gateway-client/src/browser.ts +++ b/packages/gateway-client/src/browser.ts @@ -7,7 +7,7 @@ export * from "./protocol-client.js"; export * from "./reconnect-policy.js"; export * from "./session-projection.js"; export * from "./session-subscriptions.js"; -export { DEFAULT_PREAUTH_HANDSHAKE_TIMEOUT_MS } from "./timeouts.js"; +export { DEFAULT_PREAUTH_HANDSHAKE_TIMEOUT_MS, resolveSafeTimeoutDelayMs } from "./timeouts.js"; export * from "@openclaw/gateway-protocol/client-info"; export * from "@openclaw/gateway-protocol/connect-error-details"; export * from "@openclaw/gateway-protocol/gateway-error-details"; diff --git a/packages/gateway-client/src/client.watchdog.test.ts b/packages/gateway-client/src/client.watchdog.test.ts index 404c9df130e4..04645953a592 100644 --- a/packages/gateway-client/src/client.watchdog.test.ts +++ b/packages/gateway-client/src/client.watchdog.test.ts @@ -223,6 +223,97 @@ describe("GatewayClient", () => { } }); + test.each([ + { retirement: "event owner", firstListenerCalls: 0 }, + { retirement: "first direct listener", firstListenerCalls: 1 }, + ])( + "does not deliver a retired frame after the $retirement closes its socket", + ({ retirement, firstListenerCalls }) => { + const onEvent = vi.fn(() => { + if (retirement === "event owner") { + client.stop(); + } + }); + const firstListener = vi.fn(() => { + if (retirement === "first direct listener") { + client.stop(); + } + }); + const staleListener = vi.fn(); + const { client, connections } = createSyntheticGatewayProtocol({ onEvent }); + client.addEventListener(firstListener); + client.addEventListener(staleListener); + client.start(); + const connection = connections[0]; + if (!connection) { + throw new Error("synthetic protocol connection missing"); + } + + connection.handlers.message( + JSON.stringify({ + type: "event", + event: "board.command", + payload: { command: "retired" }, + seq: 1, + }), + ); + + expect(onEvent).toHaveBeenCalledOnce(); + expect(firstListener).toHaveBeenCalledTimes(firstListenerCalls); + expect(staleListener).not.toHaveBeenCalled(); + expect(connection.close).toHaveBeenCalledOnce(); + }, + ); + + test.each([ + { replacement: "a new callback", reuseCallback: false }, + { replacement: "the same callback", reuseCallback: true }, + ])("does not revive a removed subscription replaced with $replacement", ({ reuseCallback }) => { + const removedListener = vi.fn(); + const addedListener = reuseCallback ? removedListener : vi.fn(); + let removeListener = () => {}; + let isFirstEvent = true; + const onEvent = vi.fn(() => { + if (isFirstEvent) { + isFirstEvent = false; + removeListener(); + client.addEventListener(addedListener); + } + }); + const { client, connections } = createSyntheticGatewayProtocol({ onEvent }); + removeListener = client.addEventListener(removedListener); + client.start(); + const connection = connections[0]; + if (!connection) { + throw new Error("synthetic protocol connection missing"); + } + + connection.handlers.message( + JSON.stringify({ type: "event", event: "board.changed", payload: {}, seq: 1 }), + ); + + expect(removedListener).not.toHaveBeenCalled(); + expect(addedListener).not.toHaveBeenCalled(); + + connection.handlers.message( + JSON.stringify({ type: "event", event: "board.changed", payload: {}, seq: 2 }), + ); + + expect(addedListener).toHaveBeenCalledOnce(); + if (!reuseCallback) { + expect(removedListener).not.toHaveBeenCalled(); + } + + // Calling the retired subscription's disposer cannot remove its replacement. + removeListener(); + connection.handlers.message( + JSON.stringify({ type: "event", event: "board.changed", payload: {}, seq: 3 }), + ); + + expect(addedListener).toHaveBeenCalledTimes(2); + client.stop(); + }); + test.each([ { recovery: "stops the socket", restart: false }, { recovery: "replaces the socket", restart: true }, diff --git a/packages/gateway-client/src/event-listeners.ts b/packages/gateway-client/src/event-listeners.ts new file mode 100644 index 000000000000..669d6a54a7ce --- /dev/null +++ b/packages/gateway-client/src/event-listeners.ts @@ -0,0 +1,24 @@ +type GatewayEventListener = (event: TEvent) => void; + +/** Subscription identity prevents old frames and disposers from reviving callbacks. */ +export class GatewayEventListeners { + private readonly listeners = new Map, object>(); + + add(listener: GatewayEventListener): () => void { + const subscription = this.listeners.get(listener) ?? {}; + this.listeners.set(listener, subscription); + return () => { + if (this.listeners.get(listener) === subscription) { + this.listeners.delete(listener); + } + }; + } + + snapshot(): Array<[GatewayEventListener, object]> { + return [...this.listeners]; + } + + isCurrent(listener: GatewayEventListener, subscription: object): boolean { + return this.listeners.get(listener) === subscription; + } +} diff --git a/packages/gateway-client/src/pending-request.ts b/packages/gateway-client/src/pending-request.ts new file mode 100644 index 000000000000..5b5f11ed86b0 --- /dev/null +++ b/packages/gateway-client/src/pending-request.ts @@ -0,0 +1,12 @@ +/** Owned settlement, cleanup, and timing state for one Gateway wire request. */ +export type GatewayPendingRequest = { + resolve: (value: unknown) => void; + reject: (error: Error) => void; + expectFinal: boolean; + acceptedNotified: boolean; + onAccepted?: (payload: unknown) => void; + cleanup?: () => void; + unbounded: boolean; + method: string; + startedAtMs: number; +}; diff --git a/packages/gateway-client/src/protocol-client.handshake.test.ts b/packages/gateway-client/src/protocol-client.handshake.test.ts new file mode 100644 index 000000000000..5413a3b055d5 --- /dev/null +++ b/packages/gateway-client/src/protocol-client.handshake.test.ts @@ -0,0 +1,95 @@ +import { afterEach, describe, expect, it, vi } from "vitest"; +import { GatewayProtocolClient, type GatewayProtocolSocketHandlers } from "./protocol-client.js"; +import { DEFAULT_PREAUTH_HANDSHAKE_TIMEOUT_MS } from "./timeouts.js"; + +type HandshakeConnection = { + handlers: GatewayProtocolSocketHandlers; + send: ReturnType void>>; + close: ReturnType void>>; +}; + +function createHandshakeClient( + buildConnectPlan: () => Record | Promise> = () => ({}), +) { + const connections: HandshakeConnection[] = []; + let nextRequestId = 0; + const client = new GatewayProtocolClient>({ + createSocket: (handlers) => { + let open = true; + const send = vi.fn<(data: string) => void>(); + const close = vi.fn<(code?: number, reason?: string) => void>((code, reason) => { + open = false; + handlers.close(code ?? 1000, reason ?? ""); + }); + connections.push({ handlers, send, close }); + return { isOpen: () => open, send, close }; + }, + createRequestId: () => `request-${++nextRequestId}`, + buildConnectPlan, + buildConnectParams: (plan) => plan, + resolveClose: () => ({ retry: true, notify: true }), + handshake: { mode: "require-challenge", timeoutMs: 100 }, + reconnect: { initialMs: 10, multiplier: 2, maxMs: 100 }, + }); + return { client, connections }; +} + +function receiveConnectChallenge(connection: HandshakeConnection): void { + connection.handlers.open(); + connection.handlers.message( + JSON.stringify({ + type: "event", + event: "connect.challenge", + payload: { nonce: "synthetic-nonce" }, + }), + ); +} + +describe("GatewayProtocolClient connect handshake", () => { + afterEach(() => vi.useRealTimers()); + + it("reconnects when an open Gateway never responds to connect", async () => { + vi.useFakeTimers(); + const { client, connections } = createHandshakeClient(); + client.start(); + const connection = connections[0]; + expect(connection).toBeDefined(); + if (!connection) { + return; + } + receiveConnectChallenge(connection); + expect(connection.send).toHaveBeenCalledOnce(); + + await vi.advanceTimersByTimeAsync(DEFAULT_PREAUTH_HANDSHAKE_TIMEOUT_MS); + + expect(connection.close).toHaveBeenCalledWith(4000, "connect timeout"); + await vi.advanceTimersByTimeAsync(10); + expect(connections).toHaveLength(2); + client.stop(); + }); + + it("retires device preparation that outlives the connect handshake", async () => { + vi.useFakeTimers(); + let resolvePlan: (plan: Record) => void = () => undefined; + const plan = new Promise>((resolve) => { + resolvePlan = resolve; + }); + const { client, connections } = createHandshakeClient(() => plan); + client.start(); + const connection = connections[0]; + expect(connection).toBeDefined(); + if (!connection) { + return; + } + receiveConnectChallenge(connection); + expect(connection.send).not.toHaveBeenCalled(); + + await vi.advanceTimersByTimeAsync(DEFAULT_PREAUTH_HANDSHAKE_TIMEOUT_MS); + + expect(connection.close).toHaveBeenCalledWith(4000, "connect timeout"); + resolvePlan({}); + await vi.advanceTimersByTimeAsync(0); + expect(connection.send).not.toHaveBeenCalled(); + client.stop(); + }); +}); diff --git a/packages/gateway-client/src/protocol-client.ts b/packages/gateway-client/src/protocol-client.ts index e22af43ecc71..ed28aea02470 100644 --- a/packages/gateway-client/src/protocol-client.ts +++ b/packages/gateway-client/src/protocol-client.ts @@ -4,6 +4,9 @@ import { isGatewayResponseFrame, } from "@openclaw/gateway-protocol/frame-guards"; import { RetrySupervisor, sleepWithAbort } from "@openclaw/retry"; +import { GatewayEventListeners } from "./event-listeners.js"; +import type { GatewayPendingRequest } from "./pending-request.js"; +import { clearGatewayConnectTimeout, startGatewayConnectTimeout } from "./timeouts.js"; export type GatewayProtocolSocket = { isOpen: () => boolean; @@ -146,17 +149,6 @@ type ConnectTimingState = { usedFallback: boolean; }; type CloseSnapshot = Omit; -type PendingRequest = { - resolve: (value: unknown) => void; - reject: (error: Error) => void; - expectFinal: boolean; - acceptedNotified: boolean; - onAccepted?: (payload: unknown) => void; - cleanup?: () => void; - unbounded: boolean; - method: string; - startedAtMs: number; -}; /** * Browser-safe gateway wire client. Environment adapters own transport and auth @@ -164,8 +156,8 @@ type PendingRequest = { */ export class GatewayProtocolClient { private socket: GatewayProtocolSocket | null = null; - private readonly pending = new Map(); - private listeners = new Set<(event: EventFrame) => void>(); + private readonly pending = new Map(); + private readonly listeners = new GatewayEventListeners(); private stopped = true; private generation = 0; private lastSeq: number | null = null; @@ -250,7 +242,7 @@ export class GatewayProtocolClient { options?.timeoutMs === null ? undefined : (options?.timeoutMs ?? this.opts.requestTimeoutMs); return new Promise((resolve, reject) => { let timeout: ReturnType | undefined; - const pending: PendingRequest = { + const pending: GatewayPendingRequest = { resolve: (value) => resolve(value as T), reject, expectFinal: options?.expectFinal === true, @@ -312,8 +304,7 @@ export class GatewayProtocolClient { } addEventListener(listener: (event: EventFrame) => void): () => void { - this.listeners.add(listener); - return () => this.listeners.delete(listener); + return this.listeners.add(listener); } closeSocket(code?: number, reason?: string): void { @@ -448,6 +439,13 @@ export class GatewayProtocolClient { } this.connectSent = true; this.clearHandshakeTimer(); + // The challenge timer ends before asynchronous device preparation. Keep + // the same socket supervised until hello so a silent peer cannot strand it. + this.handshakeTimer = startGatewayConnectTimeout(() => { + if (this.isActive(socket, generation) && !this.helloReceived) { + socket.close(4000, "connect timeout"); + } + }); let planOrPromise: TPlan | Promise; try { planOrPromise = this.opts.buildConnectPlan({ @@ -501,6 +499,7 @@ export class GatewayProtocolClient { return; } this.helloReceived = true; + this.clearHandshakeTimer(); this.connectFailure = undefined; this.reconnectSupervisor.reset(); this.recordTiming("hello", generation, plan); @@ -572,9 +571,17 @@ export class GatewayProtocolClient { } this.lastSeq = seq; } + // An owner may replace the socket while handling this frame. Snapshot + // first so replacement listeners cannot inherit a retired event. + const listeners = this.listeners.snapshot(); this.invoke("event", () => this.opts.onEvent?.(parsed)); - for (const listener of this.listeners) { - this.invoke("event listener", () => listener(parsed)); + for (const [listener, subscription] of listeners) { + if (!this.isActive(socket, generation)) { + return; + } + if (this.listeners.isCurrent(listener, subscription)) { + this.invoke("event listener", () => listener(parsed)); + } } return; } @@ -665,7 +672,7 @@ export class GatewayProtocolClient { private finishRequestTiming( id: string, - pending: PendingRequest, + pending: GatewayPendingRequest, ok: boolean, errorCode?: string, ): void { @@ -730,10 +737,7 @@ export class GatewayProtocolClient { } private clearHandshakeTimer(): void { - if (this.handshakeTimer) { - clearTimeout(this.handshakeTimer); - this.handshakeTimer = null; - } + this.handshakeTimer = clearGatewayConnectTimeout(this.handshakeTimer); } private invoke(label: string, callback: () => void): void { diff --git a/packages/gateway-client/src/timeouts.ts b/packages/gateway-client/src/timeouts.ts index bf2db25d522b..3ab214b6dd5a 100644 --- a/packages/gateway-client/src/timeouts.ts +++ b/packages/gateway-client/src/timeouts.ts @@ -28,6 +28,22 @@ function isTestRuntimeEnv(env: NodeJS.ProcessEnv): boolean { export const MAX_SAFE_TIMEOUT_DELAY_MS = 2_147_483_647; /** Default server-side window for gateway preauth handshakes. */ export const DEFAULT_PREAUTH_HANDSHAKE_TIMEOUT_MS = 15_000; + +/** Starts the browser-safe deadline that covers Gateway connect preparation and hello. */ +export function startGatewayConnectTimeout(onTimeout: () => void): ReturnType { + const timer = setTimeout(onTimeout, DEFAULT_PREAUTH_HANDSHAKE_TIMEOUT_MS); + timer.unref?.(); + return timer; +} + +/** Clears either pending Gateway handshake phase without retaining its timer. */ +export function clearGatewayConnectTimeout(timer: ReturnType | null): null { + if (timer !== null) { + clearTimeout(timer); + } + return null; +} + /** Default deadline for a single non-streaming Gateway request. */ export const DEFAULT_GATEWAY_REQUEST_TIMEOUT_MS = 30_000; /** Minimum client watchdog delay for connect challenge setup. */ diff --git a/scripts/check-changed.mjs b/scripts/check-changed.mjs index 8f26b063409c..305ef9fabecd 100644 --- a/scripts/check-changed.mjs +++ b/scripts/check-changed.mjs @@ -28,6 +28,7 @@ import { resolveLocalHeavyCheckEnv, } from "./lib/local-heavy-check-runtime.mjs"; import { runManagedCommand } from "./lib/managed-child-process.mjs"; +import { listGeneratedExtensionAssetSources } from "./lib/static-extension-assets.mjs"; import { createSparseTsgoSkipEnv } from "./lib/tsgo-sparse-guard.mjs"; const NPM_LOCK_POLICY_PATH_RE = @@ -74,6 +75,7 @@ const MACOS_APP_CI_PATH_RE = /^(?:apps\/(?:macos|macos-mlx-tts|shared|swabble)\/|Swabble\/|scripts\/(?:codesign-mac-app|create-dmg|notarize-mac-artifact|package-mac-app|package-mac-dist)\.sh$|scripts\/lib\/(?:plistbuddy|swift-toolchain)\.sh$|test\/scripts\/(?:codesign-mac-app|create-dmg|notarize-mac-artifact|package-mac-app|package-mac-dist)\.test\.ts$)/u; let corepackPnpmShimDir; let corepackPnpmShimCleanupRegistered = false; +let cachedGeneratedExtensionAssetPaths; let npmLockPackageDirsForChangedPaths; async function ensureChangedCheckRuntimeDependencies(paths) { @@ -421,6 +423,11 @@ async function runChangedCheckViaCrabbox(argv = [], env = process.env) { export function createChangedCheckPlan(result, options = {}) { const commands = []; const baseEnv = createChangedCheckChildEnv(options.env ?? process.env); + const generatedExtensionAssetPaths = result.paths.some((changedPath) => + LINTABLE_EXTENSION_PATH_RE.test(changedPath), + ) + ? (cachedGeneratedExtensionAssetPaths ??= new Set(listGeneratedExtensionAssetSources())) + : new Set(); const add = (name, args, env) => { if (!commands.some((command) => command.name === name && sameArgs(command.args, args))) { commands.push({ name, args, ...(env ? { env } : {}) }); @@ -437,9 +444,18 @@ export function createChangedCheckPlan(result, options = {}) { }; const addTypecheck = (name, args) => add(name, args, createSparseTsgoSkipEnv(baseEnv)); const addLint = (name, args) => add(name, args, baseEnv); - const addTargetedLint = (createCommand, lintablePathRe, fallbackName, fallbackArgs) => { - const targets = result.paths.filter((changedPath) => lintablePathRe.test(changedPath)); - const otherPaths = result.paths.filter((changedPath) => !lintablePathRe.test(changedPath)); + const addTargetedLint = ( + createCommand, + lintablePathRe, + fallbackName, + fallbackArgs, + ignoredPaths, + ) => { + const candidatePaths = ignoredPaths + ? result.paths.filter((changedPath) => !ignoredPaths.has(changedPath)) + : result.paths; + const targets = candidatePaths.filter((changedPath) => lintablePathRe.test(changedPath)); + const otherPaths = candidatePaths.filter((changedPath) => !lintablePathRe.test(changedPath)); const targetedCommands = []; for (let offset = 0; offset < targets.length; offset += TARGETED_LINT_PATH_LIMIT) { @@ -685,12 +701,24 @@ export function createChangedCheckPlan(result, options = {}) { addLint("lint core", ["lint:core"]); } if (lanes.extensions || lanes.extensionTests) { - addTargetedLint( - createTargetedExtensionLintCommand, - LINTABLE_EXTENSION_PATH_RE, - "lint extensions", - ["lint:extensions"], - ); + // Generated plugin outputs have their own asset-integrity gate and are + // intentionally ignored by oxlint; manifests still need full-lane fallback. + if ( + !result.paths.some((changedPath) => generatedExtensionAssetPaths.has(changedPath)) || + result.paths.some( + (changedPath) => + getChangedPathFacts(changedPath).surface === "extension" && + !generatedExtensionAssetPaths.has(changedPath), + ) + ) { + addTargetedLint( + createTargetedExtensionLintCommand, + LINTABLE_EXTENSION_PATH_RE, + "lint extensions", + ["lint:extensions"], + generatedExtensionAssetPaths, + ); + } } if (lanes.tooling || lanes.liveDockerTooling) { if ( diff --git a/src/infra/dotenv.ts b/src/infra/dotenv.ts index db39a132126b..cfed6dd4f51a 100644 --- a/src/infra/dotenv.ts +++ b/src/infra/dotenv.ts @@ -229,7 +229,6 @@ const BLOCKED_WORKSPACE_DOTENV_PREFIXES = [ // Workspace .env is untrusted; reserve the full OpenClaw runtime namespace // for shell/global config so new OPENCLAW_* controls are fail-closed by default. "OPENCLAW_", - "OPENCLAW_CLAWHUB_", "OPENCLAW_DISABLE_", "OPENCLAW_SKIP_", "OPENCLAW_UPDATE_", diff --git a/test/scripts/changed-lanes-generated-extension-lint.test.ts b/test/scripts/changed-lanes-generated-extension-lint.test.ts new file mode 100644 index 000000000000..7930c0d200b3 --- /dev/null +++ b/test/scripts/changed-lanes-generated-extension-lint.test.ts @@ -0,0 +1,50 @@ +import { describe, expect, it } from "vitest"; +import { detectChangedLanes } from "../../scripts/changed-lanes.mjs"; +import { createChangedCheckPlan } from "../../scripts/check-changed.mjs"; + +describe("generated extension asset lint planning", () => { + it("still lints extension tests alongside a generated browser asset", () => { + const generatedAsset = "extensions/browser/chrome-extension/modules/copilot-runtime.js"; + const extensionTest = "extensions/browser/chrome-extension/modules/copilot-gateway.test.ts"; + const result = detectChangedLanes([generatedAsset, extensionTest]); + const plan = createChangedCheckPlan(result, { env: { PATH: "/usr/bin" } }); + + expect(result.lanes.extensionTests).toBe(true); + expect(plan.commands).toContainEqual( + expect.objectContaining({ + name: "lint extension changed file", + args: [ + "scripts/run-oxlint.mjs", + "--tsconfig", + "config/tsconfig/oxlint.extensions.json", + extensionTest, + ], + }), + ); + expect( + plan.commands + .filter((command) => command.args[0] === "scripts/run-oxlint.mjs") + .flatMap((command) => command.args), + ).not.toContain(generatedAsset); + }); + + it("keeps fallback extension lint for a manifest beside a generated browser asset", () => { + const generatedAsset = "extensions/browser/chrome-extension/modules/copilot-runtime.js"; + const manifest = "extensions/browser/openclaw.plugin.json"; + const result = detectChangedLanes([generatedAsset, manifest]); + const plan = createChangedCheckPlan(result, { env: { PATH: "/usr/bin" } }); + + expect(result.lanes.extensions).toBe(true); + expect(plan.commands).toContainEqual( + expect.objectContaining({ + name: "lint extensions", + args: ["lint:extensions"], + }), + ); + expect( + plan.commands + .filter((command) => command.args[0] === "scripts/run-oxlint.mjs") + .flatMap((command) => command.args), + ).not.toContain(generatedAsset); + }); +}); diff --git a/test/scripts/changed-lanes.test.ts b/test/scripts/changed-lanes.test.ts index 99ddfc1eece6..e52d11f9b5e1 100644 --- a/test/scripts/changed-lanes.test.ts +++ b/test/scripts/changed-lanes.test.ts @@ -804,6 +804,48 @@ describe("scripts/changed-lanes", () => { } }); + it("keeps manifest-declared generated browser assets out of targeted extension lint", () => { + const generatedAsset = "extensions/browser/chrome-extension/modules/copilot-runtime.js"; + const result = detectChangedLanes([ + generatedAsset, + "packages/gateway-client/src/protocol-client.ts", + ]); + const plan = createChangedCheckPlan(result, { env: { PATH: "/usr/bin" } }); + + expect(result.lanes.extensions).toBe(true); + expect(plan.commands.map((command) => command.args[0])).toContain("tsgo:extensions"); + expect(plan.commands.map((command) => command.args[0])).not.toContain("lint:extensions"); + expect( + plan.commands + .filter((command) => command.args[0] === "scripts/run-oxlint.mjs") + .flatMap((command) => command.args), + ).not.toContain(generatedAsset); + }); + + it("still lints extension source alongside its generated browser asset", () => { + const generatedAsset = "extensions/browser/chrome-extension/modules/copilot-runtime.js"; + const source = "extensions/browser/scripts/copilot-runtime-entry.ts"; + const result = detectChangedLanes([generatedAsset, source]); + const plan = createChangedCheckPlan(result, { env: { PATH: "/usr/bin" } }); + + expect(plan.commands).toContainEqual( + expect.objectContaining({ + name: "lint extension changed file", + args: [ + "scripts/run-oxlint.mjs", + "--tsconfig", + "config/tsconfig/oxlint.extensions.json", + source, + ], + }), + ); + expect( + plan.commands + .filter((command) => command.args[0] === "scripts/run-oxlint.mjs") + .flatMap((command) => command.args), + ).not.toContain(generatedAsset); + }); + it.each([ { owner: "core", diff --git a/ui/src/api/gateway.node.test.ts b/ui/src/api/gateway.node.test.ts index c0814c3e9a04..ee6c82cc5a09 100644 --- a/ui/src/api/gateway.node.test.ts +++ b/ui/src/api/gateway.node.test.ts @@ -633,6 +633,117 @@ describe("GatewayBrowserClient", () => { expect(ws.lastClose).toEqual({ code: 4000, reason: "terminal liveness timeout" }); }); + it("reconnects a silently stalled socket using its advertised Gateway heartbeat", async () => { + useNodeFakeTimers(); + const client = new GatewayBrowserClient({ url: DEFAULT_GATEWAY_URL }); + try { + const { ws, connectFrame } = await startConnect(client); + ws.emitMessage({ + type: "res", + id: connectFrame.id, + ok: true, + payload: { + type: "hello-ok", + protocol: 4, + auth: { role: "operator", scopes: [] }, + policy: { tickIntervalMs: 1_000 }, + }, + }); + + await vi.advanceTimersByTimeAsync(3_000); + + expect(ws.lastClose).toEqual({ code: 4000, reason: "tick timeout" }); + } finally { + client.stop(); + } + }); + + it.each([Number.MAX_SAFE_INTEGER, 2 ** 32 + 1])( + "clamps the advertised heartbeat %d before scheduling its browser timer", + async (advertisedTickIntervalMs) => { + useNodeFakeTimers(); + const setIntervalSpy = vi.spyOn(globalThis, "setInterval"); + const client = new GatewayBrowserClient({ url: DEFAULT_GATEWAY_URL }); + + try { + const { ws, connectFrame } = await startConnect(client); + ws.emitMessage({ + type: "res", + id: connectFrame.id, + ok: true, + payload: { + type: "hello-ok", + protocol: 4, + auth: { role: "operator", scopes: [] }, + policy: { tickIntervalMs: advertisedTickIntervalMs }, + }, + }); + await vi.advanceTimersByTimeAsync(0); + + expect(setIntervalSpy).toHaveBeenLastCalledWith(expect.any(Function), 2_147_483_647); + await vi.advanceTimersByTimeAsync(5_000); + expect(ws.lastClose).toBeNull(); + } finally { + client.stop(); + } + }, + ); + + it("keeps a healthy heartbeat and explicitly unbounded request alive", async () => { + useNodeFakeTimers(); + const client = new GatewayBrowserClient({ url: DEFAULT_GATEWAY_URL }); + try { + const { ws, connectFrame } = await startConnect(client); + ws.emitMessage({ + type: "res", + id: connectFrame.id, + ok: true, + payload: { + type: "hello-ok", + protocol: 4, + auth: { role: "operator", scopes: [] }, + policy: { tickIntervalMs: 1_000 }, + }, + }); + const request = client.request("wizard.next", {}, { timeoutMs: null }); + const requestFrame = JSON.parse(ws.sent.at(-1) ?? "{}") as { id?: string }; + + for (let seq = 1; seq <= 4; seq += 1) { + await vi.advanceTimersByTimeAsync(1_000); + ws.emitMessage({ type: "event", event: "tick", seq, payload: {} }); + } + + expect(ws.lastClose).toBeNull(); + ws.emitMessage({ type: "res", id: requestFrame.id, ok: true, payload: { done: true } }); + await expect(request).resolves.toEqual({ done: true }); + } finally { + client.stop(); + } + }); + + it("disposes the Gateway heartbeat when its browser client stops", async () => { + useNodeFakeTimers(); + const client = new GatewayBrowserClient({ url: DEFAULT_GATEWAY_URL }); + const { ws, connectFrame } = await startConnect(client); + ws.emitMessage({ + type: "res", + id: connectFrame.id, + ok: true, + payload: { + type: "hello-ok", + protocol: 4, + auth: { role: "operator", scopes: [] }, + policy: { tickIntervalMs: 1_000 }, + }, + }); + + client.stop(); + await vi.advanceTimersByTimeAsync(5_000); + + expect(ws.lastClose).toEqual({ code: undefined, reason: undefined }); + expect(vi.getTimerCount()).toBe(0); + }); + it("reports failed request timing without including request params", async () => { const onRequestTiming = vi.fn(); const client = new GatewayBrowserClient({ diff --git a/ui/src/api/gateway.ts b/ui/src/api/gateway.ts index a406da29e59f..018c534d6e7c 100644 --- a/ui/src/api/gateway.ts +++ b/ui/src/api/gateway.ts @@ -27,6 +27,7 @@ import { shouldRetryGatewayWithDeviceToken, isRetryableGatewayStartupUnavailableError, resolveGatewayStartupRetryAfterMs, + resolveSafeTimeoutDelayMs, MIN_CLIENT_PROTOCOL_VERSION, PROTOCOL_VERSION, } from "@openclaw/gateway-client/browser"; @@ -206,6 +207,8 @@ const STARTUP_RETRY_CLOSE_CODE = 4013; const BROWSER_WEBSOCKET_CLOSE_CODE = 1006; const BROWSER_WEBSOCKET_CONSTRUCTOR_ERROR_CODE = "BROWSER_WEBSOCKET_CONSTRUCTOR_ERROR"; const BROWSER_WEBSOCKET_SECURITY_ERROR_CODE = "BROWSER_WEBSOCKET_SECURITY_ERROR"; +const DEFAULT_GATEWAY_TICK_INTERVAL_MS = 30_000; +const MIN_GATEWAY_TICK_WATCH_INTERVAL_MS = 1_000; function getErrorMessage(err: unknown): string { return err instanceof Error && err.message ? err.message : String(err); } @@ -298,6 +301,8 @@ async function buildGatewayConnectDevice(params: { export class GatewayBrowserClient { private readonly client: GatewayProtocolClient; inboundActivitySeq = 0; + private lastInboundActivityAtMs: number | null = null; + private tickWatchTimer: ReturnType | null = null; private pendingDeviceTokenRetry = false; private deviceTokenRetryBudgetUsed = false; private readonly recoveryScopeTracker = new GatewayRecoveryScopeTracker(); @@ -326,6 +331,7 @@ export class GatewayBrowserClient { }, resolveClose: (context) => this.resolveClose(context), onClose: (context, decision) => { + this.stopTickWatch(); const error = context.connectFailure?.error; this.client.recordTiming("failed", context.generation, undefined, { errorCode: error instanceof GatewayRequestError ? error.code : "SOCKET_CLOSED", @@ -342,7 +348,10 @@ export class GatewayBrowserClient { onSocketFactoryError: (error) => this.handleSocketFactoryError(error), onEvent: (event) => this.opts.onEvent?.(event), onGap: (info) => this.opts.onGap?.(info), - onActivity: () => (this.inboundActivitySeq += 1), + onActivity: () => { + this.inboundActivitySeq += 1; + this.lastInboundActivityAtMs = Date.now(); + }, onTiming: ({ plan, detail, ...timing }) => { this.opts.onConnectTiming?.({ ...timing, @@ -374,6 +383,7 @@ export class GatewayBrowserClient { } stop() { + this.stopTickWatch(); this.client.stop(); this.pendingDeviceTokenRetry = false; this.deviceTokenRetryBudgetUsed = false; @@ -491,6 +501,7 @@ export class GatewayBrowserClient { } private handleConnectHello(hello: GatewayHelloOk, plan: ConnectPlan) { + this.startTickWatch(hello); this.pendingDeviceTokenRetry = false; this.deviceTokenRetryBudgetUsed = false; this.opts.bootstrapToken = undefined; @@ -506,6 +517,38 @@ export class GatewayBrowserClient { void this.updateRecoveryScopeForHello(hello, plan); } + private startTickWatch(hello: GatewayHelloOk): void { + this.stopTickWatch(); + const advertisedTickIntervalMs = hello.policy?.tickIntervalMs; + // Gateway policy is remote input; use the shared timer clamp so an + // oversized interval cannot wrap into a resource-exhausting hot loop. + const tickIntervalMs = resolveSafeTimeoutDelayMs( + typeof advertisedTickIntervalMs === "number" && + Number.isFinite(advertisedTickIntervalMs) && + advertisedTickIntervalMs > 0 + ? advertisedTickIntervalMs + : DEFAULT_GATEWAY_TICK_INTERVAL_MS, + { minMs: MIN_GATEWAY_TICK_WATCH_INTERVAL_MS }, + ); + this.lastInboundActivityAtMs = Date.now(); + this.tickWatchTimer = setInterval(() => { + const lastActivityAtMs = this.lastInboundActivityAtMs; + // Preserve long-running requests while real Gateway heartbeats arrive; + // only a silent socket should enter the shared reconnect lifecycle. + if (lastActivityAtMs !== null && Date.now() - lastActivityAtMs > tickIntervalMs * 2) { + this.forceReconnect("tick timeout"); + } + }, tickIntervalMs); + } + + private stopTickWatch(): void { + if (this.tickWatchTimer !== null) { + clearInterval(this.tickWatchTimer); + this.tickWatchTimer = null; + } + this.lastInboundActivityAtMs = null; + } + private async updateRecoveryScopeForHello(hello: GatewayHelloOk, plan: ConnectPlan) { if ( await this.recoveryScopeTracker.resolve({ diff --git a/ui/src/components/app-sidebar-session-catalog-live.ts b/ui/src/components/app-sidebar-session-catalog-live.ts index 479f172cef45..c88ffcdf65ef 100644 --- a/ui/src/components/app-sidebar-session-catalog-live.ts +++ b/ui/src/components/app-sidebar-session-catalog-live.ts @@ -190,6 +190,14 @@ export class SessionCatalogLiveState { this.progressive = true; } + retireConnection(reset = false): void { + if (reset) { + this.resetConnection(); + return; + } + this.clear(); + } + async requestList( client: GatewayBrowserClient, agentId: string, diff --git a/ui/src/components/app-sidebar.test.ts b/ui/src/components/app-sidebar.test.ts index 54739535e45d..3741588b7a46 100644 --- a/ui/src/components/app-sidebar.test.ts +++ b/ui/src/components/app-sidebar.test.ts @@ -8,6 +8,7 @@ import "../test-helpers/app-sidebar-cases/catalog-compat.ts"; import "../test-helpers/app-sidebar-cases/catalog-live-events.ts"; import "../test-helpers/app-sidebar-cases/catalog-project-activity.ts"; import "../test-helpers/app-sidebar-cases/catalog-live.ts"; +import "../test-helpers/app-sidebar-cases/catalog-reconnect.ts"; import "../test-helpers/app-sidebar-cases/catalog-live-errors.ts"; import "../test-helpers/app-sidebar-cases/catalog-live-state.ts"; import "../test-helpers/app-sidebar-cases/catalog-ownership.ts"; diff --git a/ui/src/components/exec-approval.test.ts b/ui/src/components/exec-approval.test.ts index 62e0d8ddfd61..c11f05caf4ab 100644 --- a/ui/src/components/exec-approval.test.ts +++ b/ui/src/components/exec-approval.test.ts @@ -271,6 +271,31 @@ describe("openclaw-exec-approval", () => { expect(onDecision).not.toHaveBeenCalled(); }); + it.each([ + { reason: "a decision is in flight", busy: true, allowedDecisions: undefined }, + { reason: "denial is unavailable", busy: false, allowedDecisions: ["allow-once"] as const }, + ])("keeps the pending approval visible when $reason", async ({ busy, allowedDecisions }) => { + const { onDecision } = await renderApproval( + createExecRequest({ + request: { + command: "echo hello", + ...(allowedDecisions ? { allowedDecisions: [...allowedDecisions] } : {}), + }, + }), + { busy }, + ); + const { modal } = await getRenderedModalDialog(container); + const cancel = new CustomEvent("modal-cancel", { + bubbles: true, + composed: true, + cancelable: true, + }); + + expect(modal.dispatchEvent(cancel)).toBe(false); + expect(cancel.defaultPrevented).toBe(true); + expect(onDecision).not.toHaveBeenCalled(); + }); + it("suppresses the automatic modal for the inline request but opens it on demand", async () => { const { approval } = await renderApproval(createExecRequest(), { inlineApprovalId: "approval-1", diff --git a/ui/src/components/exec-approval.ts b/ui/src/components/exec-approval.ts index f31d2b62750e..3653d6c66685 100644 --- a/ui/src/components/exec-approval.ts +++ b/ui/src/components/exec-approval.ts @@ -151,10 +151,13 @@ class ExecApproval extends OpenClawLightDomContentsElement { return nothing; } const decisions = resolveApprovalDecisions(active); - const handleCancel = () => { - if (!props.busy && decisions.includes("deny")) { - void props.onDecision(active.id, "deny"); + const handleCancel = (event: Event) => { + if (props.busy || !decisions.includes("deny")) { + // Dismissal must never hide an approval that cannot yet be resolved. + event.preventDefault(); + return; } + void props.onDecision(active.id, "deny"); }; return html` { ), ) .toBe(64); + const initialBlobUrls = await page + .locator("img.chat-message-image") + .evaluateAll((images) => images.map((image) => image.getAttribute("src"))); + const retainedRecentBlobUrl = expectDefined( + initialBlobUrls[0], + "recent managed image Blob URL", + ); await replaceHistory( historyFor([0], "Recently viewed managed image"), @@ -611,26 +618,30 @@ suite.define(() => { await replaceHistory(historyFor([64], "Overflow managed image"), "Overflow managed image 65"); await expect.poll(async () => (await readBlobProof()).created.length).toBe(65); const overflowProof = await readBlobProof(); - const retainedRecentBlobUrl = expectDefined( - overflowProof.created[0], - "recent managed image Blob URL", - ); + // Concurrent image fetches can resolve in any order. Find the real LRU + // rather than assuming that creation order matches transcript order. const evictedBlobUrl = expectDefined( - overflowProof.created[1], + overflowProof.created.find((blobUrl) => blobUrl !== retainedRecentBlobUrl), "evicted managed image Blob URL", ); + const evictedImageIndex = initialBlobUrls.indexOf(evictedBlobUrl); + expect(evictedImageIndex).toBeGreaterThanOrEqual(0); expect(overflowProof.revoked).toContain(evictedBlobUrl); expect(overflowProof.revoked).not.toContain(retainedRecentBlobUrl); const evictedPath = new URL( - expectDefined(imageUrls[1], "evicted managed image URL"), + expectDefined(imageUrls[evictedImageIndex], "evicted managed image URL"), suite.server.baseUrl, ).pathname; const fetchesBeforeRevisit = fetchedMedia.filter( (request) => request.pathname === evictedPath, ).length; - await replaceHistory(historyFor([1], "Refetched managed image"), "Refetched managed image 2"); - const revisitedImage = page.getByAltText("Refetched managed image 2"); + const revisitedImageAlt = `Refetched managed image ${evictedImageIndex + 1}`; + await replaceHistory( + historyFor([evictedImageIndex], "Refetched managed image"), + revisitedImageAlt, + ); + const revisitedImage = page.getByAltText(revisitedImageAlt); await expect .poll(() => revisitedImage.evaluate((image) => @@ -655,7 +666,7 @@ suite.define(() => { const proofSummary = { cacheCapacity: 64, createdBlobUrls: finalProof.created.length, - evictedBlobIndex: 1, + evictedBlobIndex: evictedImageIndex, evictedImageFetches, refetchedImageNaturalWidth: await revisitedImage.evaluate( (image) => (image as HTMLImageElement).naturalWidth, diff --git a/ui/src/lib/board/gateway-provider.test.ts b/ui/src/lib/board/gateway-provider.test.ts index 75d744b104a1..dd10fbd0f816 100644 --- a/ui/src/lib/board/gateway-provider.test.ts +++ b/ui/src/lib/board/gateway-provider.test.ts @@ -1,4 +1,5 @@ // @vitest-environment node +import type { EventFrame } from "@openclaw/gateway-protocol"; import { afterEach, describe, expect, it, vi } from "vitest"; import { GatewayBoardProvider, type BoardProvider } from "./provider.ts"; @@ -154,6 +155,168 @@ describe("gateway board provider lifecycle", () => { expect(provider.snapshot$.value).toEqual(snapshot); }); + it("pauses board retries during an outage and refreshes once after reconnect", async () => { + vi.useFakeTimers(); + let connected = true; + const snapshot = { + sessionKey: "agent:main:paused-reconnect", + revision: 1, + tabs: [], + widgets: [], + }; + const request = vi + .fn() + .mockRejectedValueOnce(new Error("temporarily unavailable")) + .mockImplementation(async () => { + if (!connected) { + throw new Error("gateway not connected"); + } + return snapshot; + }); + const client = { + request: request as never, + addEventListener: () => () => {}, + }; + const provider = new GatewayBoardProvider(snapshot.sessionKey, client); + await vi.advanceTimersByTimeAsync(0); + expect(request).toHaveBeenCalledOnce(); + + connected = false; + provider.attachClient(client, false); + await vi.advanceTimersByTimeAsync(31_000); + + expect(request).toHaveBeenCalledOnce(); + expect(vi.getTimerCount()).toBe(0); + + connected = true; + provider.attachClient(client, true); + await vi.advanceTimersByTimeAsync(0); + + expect(request).toHaveBeenCalledTimes(2); + expect(provider.snapshot$.value).toEqual(snapshot); + provider.dispose(); + }); + + it("pauses a failed user-requested widget refresh until reconnect", async () => { + vi.useFakeTimers(); + const snapshot = { + sessionKey: "agent:main:offline-widget-refresh", + revision: 1, + tabs: [], + widgets: [], + }; + const request = vi + .fn() + .mockRejectedValueOnce(new Error("gateway not connected")) + .mockResolvedValue(snapshot); + const client = { + request: request as never, + addEventListener: () => () => {}, + }; + const provider = new GatewayBoardProvider(snapshot.sessionKey, client, false); + + const refresh = provider.refreshWidgetFrame("status"); + await vi.advanceTimersByTimeAsync(0); + await refresh; + await vi.advanceTimersByTimeAsync(31_000); + + expect(request).toHaveBeenCalledOnce(); + expect(vi.getTimerCount()).toBe(0); + + provider.attachClient(client, true); + await vi.advanceTimersByTimeAsync(0); + + expect(request).toHaveBeenCalledTimes(2); + expect(provider.snapshot$.value).toEqual(snapshot); + provider.dispose(); + }); + + it("preserves a manual offline refresh that joins an existing retry loop", async () => { + vi.useFakeTimers(); + const snapshot = { + sessionKey: "agent:main:coalesced-offline-widget-refresh", + revision: 1, + tabs: [], + widgets: [], + }; + const request = vi + .fn() + .mockRejectedValueOnce(new Error("temporarily unavailable")) + .mockRejectedValueOnce(new Error("gateway not connected")) + .mockResolvedValue(snapshot); + const client = { + request: request as never, + addEventListener: () => () => {}, + }; + const provider = new GatewayBoardProvider(snapshot.sessionKey, client); + await vi.advanceTimersByTimeAsync(0); + expect(request).toHaveBeenCalledOnce(); + + provider.attachClient(client, false); + const refresh = provider.refreshWidgetFrame("status"); + await vi.advanceTimersByTimeAsync(0); + await refresh; + await vi.advanceTimersByTimeAsync(31_000); + + expect(request).toHaveBeenCalledTimes(2); + expect(vi.getTimerCount()).toBe(0); + + provider.attachClient(client, true); + await vi.advanceTimersByTimeAsync(0); + + expect(request).toHaveBeenCalledTimes(3); + expect(provider.snapshot$.value).toEqual(snapshot); + provider.dispose(); + }); + + it("does not reread a queued manual refresh after its gateway disconnects", async () => { + vi.useFakeTimers(); + const initial = { + sessionKey: "agent:main:offline-queued-refresh", + revision: 1, + tabs: [], + widgets: [], + }; + const changed = { ...initial, revision: 2 }; + let listener: ((event: { event: string; payload: unknown }) => void) | undefined; + const resolvers: Array<(value: typeof initial) => void> = []; + const request = vi.fn( + () => + new Promise((resolve) => { + resolvers.push(resolve); + }), + ); + const client = { + request: request as never, + addEventListener: (next: (event: EventFrame) => void) => { + listener = next as typeof listener; + return () => {}; + }, + }; + const provider = new GatewayBoardProvider(initial.sessionKey, client, true); + const refresh = provider.refreshWidgetFrame("status"); + + listener?.({ + event: "board.changed", + payload: { sessionKey: initial.sessionKey, revision: changed.revision }, + }); + provider.attachClient(client, false); + resolvers[0]?.(initial); + await refresh; + await vi.advanceTimersByTimeAsync(31_000); + + expect(request).toHaveBeenCalledOnce(); + expect(vi.getTimerCount()).toBe(0); + + provider.attachClient(client, true); + resolvers[1]?.(changed); + await vi.advanceTimersByTimeAsync(0); + + expect(request).toHaveBeenCalledTimes(2); + expect(provider.snapshot$.value).toEqual(changed); + provider.dispose(); + }); + it("retries a transient board.changed refresh failure", async () => { vi.useFakeTimers(); let listener: ((event: { event: string; payload: unknown }) => void) | undefined; @@ -240,17 +403,15 @@ describe("gateway board provider lifecycle", () => { resolvers.push(resolve); }), ); - const provider = new GatewayBoardProvider( - initial.sessionKey, - { - request: request as never, - addEventListener: (next) => { - listener = next as typeof listener; - return () => {}; - }, + const client = { + request: request as never, + addEventListener: (next: (event: EventFrame) => void) => { + listener = next as typeof listener; + return () => {}; }, - false, - ); + }; + const provider = new GatewayBoardProvider(initial.sessionKey, client, false); + provider.attachClient(client, true); const refresh = provider.refreshWidgetFrame("status"); listener?.({ diff --git a/ui/src/lib/board/gateway-provider.ts b/ui/src/lib/board/gateway-provider.ts index 0685d3cb559b..837f03a31c94 100644 --- a/ui/src/lib/board/gateway-provider.ts +++ b/ui/src/lib/board/gateway-provider.ts @@ -37,6 +37,7 @@ export class GatewayBoardProvider implements BoardProvider { private unsubscribe: (() => void) | undefined; private refreshLoop: Promise | undefined; private refreshRequested = false; + private userRefreshRequested = false; private readonly changedWidgets = new Set(); private stateGeneration = 0; private connected = false; @@ -71,6 +72,9 @@ export class GatewayBoardProvider implements BoardProvider { } const connectionActivated = connected && !this.connected; this.connected = connected; + if (!connected) { + this.wakeRetryDelay?.(); + } if (client === this.client) { if (connectionActivated) { void this.activate(); @@ -112,6 +116,7 @@ export class GatewayBoardProvider implements BoardProvider { this.clientGeneration += 1; this.stateGeneration += 1; this.refreshRequested = false; + this.userRefreshRequested = false; this.changedWidgets.clear(); this.appViews.clear(); this.wakeRetryDelay?.(); @@ -189,7 +194,7 @@ export class GatewayBoardProvider implements BoardProvider { } refreshWidgetFrame(name: string): Promise { - return this.requestRefresh(name); + return this.requestRefresh(name, true); } async widgetAppView(name: string, revision: number): Promise { @@ -258,18 +263,19 @@ export class GatewayBoardProvider implements BoardProvider { ); } - private requestRefresh(changedWidget?: string): Promise { + private requestRefresh(changedWidget?: string, userRequested = false): Promise { if (this.disposed) { return Promise.resolve(); } this.refreshRequested = true; + this.userRefreshRequested ||= userRequested; if (changedWidget) { this.changedWidgets.add(changedWidget); } this.wakeRetryDelay?.(); this.refreshLoop ??= this.runRefreshLoop().finally(() => { this.refreshLoop = undefined; - if (this.refreshRequested) { + if (this.refreshRequested && (this.connected || this.userRefreshRequested)) { void this.requestRefresh(); } }); @@ -283,6 +289,14 @@ export class GatewayBoardProvider implements BoardProvider { this.refreshRequested = false; return; } + // Preserve pending board changes without polling a disconnected client; + // the next attached live connection owns the replacement refresh. + if (!this.connected && !this.userRefreshRequested) { + return; + } + // A manual refresh permits one offline request, never a follow-up after + // that request completes and discovers a disconnected Gateway. + this.userRefreshRequested = false; const changedWidgets = new Set(this.changedWidgets); this.changedWidgets.clear(); const client = this.client; @@ -298,6 +312,9 @@ export class GatewayBoardProvider implements BoardProvider { this.refreshRequested = true; continue; } + // A completed request satisfies any manual refresh that joined it; + // only a newer board change should require a follow-up request. + this.userRefreshRequested = false; if (stateGeneration !== this.stateGeneration) { this.refreshRequested = true; for (const name of changedWidgets) { @@ -325,6 +342,12 @@ export class GatewayBoardProvider implements BoardProvider { for (const name of changedWidgets) { this.changedWidgets.add(name); } + if (!this.connected) { + if (this.userRefreshRequested) { + continue; + } + return; + } const delayMs = retry.delayMs; // Carry backoff across failed loop iterations; successful refreshes reset it above. retry.delayMs = Math.min(delayMs * 2, 30_000); diff --git a/ui/src/pages/chat/chat-history-subscription-disposal.test.ts b/ui/src/pages/chat/chat-history-subscription-disposal.test.ts new file mode 100644 index 000000000000..006cbd0bad7d --- /dev/null +++ b/ui/src/pages/chat/chat-history-subscription-disposal.test.ts @@ -0,0 +1,114 @@ +// @vitest-environment node +import { afterEach, describe, expect, it, vi } from "vitest"; +import type { GatewayBrowserClient } from "../../api/gateway.ts"; +import type { SessionCapability } from "../../lib/sessions/index.ts"; +import { + disposeSelectedSessionMessageSubscription, + syncSelectedSessionMessageSubscription, + type ChatState, +} from "./chat-history.ts"; + +const subscription = { key: "agent:main:main", agentId: null }; + +function createSubscriptionState( + unsubscribeMessages: ReturnType>, + subscribeMessages: ReturnType> = vi.fn< + SessionCapability["subscribeMessages"] + >(), +): ChatState { + return { + client: {} as GatewayBrowserClient, + connected: true, + connectionEpoch: 1, + sessionKey: subscription.key, + chatLoading: false, + chatMessages: [], + chatThinkingLevel: null, + chatVerboseLevel: null, + chatSending: false, + chatMessage: "", + chatAttachments: [], + chatQueue: [], + chatRunId: null, + chatStream: null, + chatStreamStartedAt: null, + lastError: null, + hello: null, + sessions: { subscribeMessages, unsubscribeMessages }, + }; +} + +describe("disposed chat message subscriptions", () => { + afterEach(() => vi.useRealTimers()); + + it("releases an active message subscription when its pane is disposed", () => { + const unsubscribeMessages = vi + .fn() + .mockResolvedValue(undefined); + const state = createSubscriptionState(unsubscribeMessages); + state.chatSessionMessageSubscriptionRequestedKey = subscription.key; + state.chatSessionMessageSubscription = subscription; + + disposeSelectedSessionMessageSubscription(state); + + expect(unsubscribeMessages).toHaveBeenCalledExactlyOnceWith(subscription); + expect(state.chatSessionMessageSubscriptionRequestedKey).toBeNull(); + expect(state.chatSessionMessageSubscription).toBeNull(); + }); + + it("releases a subscription that resolves after its pane is disposed", async () => { + let resolveSubscription: (value: typeof subscription) => void = () => undefined; + const pendingSubscription = new Promise((resolve) => { + resolveSubscription = resolve; + }); + const unsubscribeMessages = vi + .fn() + .mockResolvedValue(undefined); + const state = createSubscriptionState( + unsubscribeMessages, + vi.fn().mockReturnValue(pendingSubscription), + ); + + const sync = syncSelectedSessionMessageSubscription(state as never); + await Promise.resolve(); + disposeSelectedSessionMessageSubscription(state); + resolveSubscription(subscription); + await sync; + + expect(unsubscribeMessages).toHaveBeenCalledExactlyOnceWith(subscription); + expect(state.chatSessionMessageSubscription).toBeNull(); + }); + + it("retries a temporary release failure without another pane synchronization", async () => { + vi.useFakeTimers(); + const unsubscribeMessages = vi + .fn() + .mockRejectedValueOnce(new Error("temporary observer release failure")) + .mockResolvedValueOnce(undefined); + const state = createSubscriptionState(unsubscribeMessages); + state.chatSessionMessageSubscription = subscription; + + disposeSelectedSessionMessageSubscription(state); + await vi.advanceTimersByTimeAsync(250); + + expect(unsubscribeMessages).toHaveBeenCalledTimes(2); + expect(unsubscribeMessages).toHaveBeenLastCalledWith(subscription); + expect(state.chatSessionMessageSubscription).toBeNull(); + }); + + it("bounds permanently failing releases without leaking retry timers", async () => { + vi.useFakeTimers(); + const unsubscribeMessages = vi + .fn() + .mockRejectedValue(new Error("observer unavailable")); + const state = createSubscriptionState(unsubscribeMessages); + state.chatSessionMessageSubscription = subscription; + + disposeSelectedSessionMessageSubscription(state); + await vi.runAllTimersAsync(); + + expect(unsubscribeMessages).toHaveBeenCalledTimes(3); + expect(vi.getTimerCount()).toBe(0); + expect(state.chatSessionMessageSubscription).toBeNull(); + }); +}); diff --git a/ui/src/pages/chat/chat-history.ts b/ui/src/pages/chat/chat-history.ts index 3a263be187f9..a5d77a975174 100644 --- a/ui/src/pages/chat/chat-history.ts +++ b/ui/src/pages/chat/chat-history.ts @@ -92,6 +92,8 @@ const SYNTHETIC_TRANSCRIPT_REPAIR_RESULT = "[openclaw] missing tool result in session history; inserted synthetic error result for transcript repair."; const CHAT_HISTORY_REQUEST_LIMIT = 100; const STARTUP_CHAT_HISTORY_RETRY_TIMEOUT_MS = 60_000; +const SESSION_MESSAGE_RELEASE_RETRY_MS = 250; +const MAX_SESSION_MESSAGE_RELEASE_ATTEMPTS = 3; type ChatHistoryPaneRequests = { historyVersion: number; @@ -298,6 +300,8 @@ export type ChatState = { canvasPluginSurfaceUrl?: string | null; settings?: { chatPersistCommentary?: boolean; gatewayUrl?: string | null }; sessions?: Partial; + chatSessionMessageSubscriptionRequestedKey?: string | null; + chatSessionMessageSubscription?: SessionMessageSubscription | null; chatBranches?: SessionBranch[]; chatBranchesSessionKey?: string | null; chatBranchesConnectionEpoch?: number | null; @@ -714,6 +718,44 @@ async function retryPendingSessionMessageSubscriptionReleases( ); } +export function disposeSelectedSessionMessageSubscription(state: ChatState): void { + const requests = getChatHistoryPaneRequests(state); + requests.subscriptionGeneration += 1; + const subscriptions = new Set(requests.pendingSubscriptionReleases); + requests.pendingSubscriptionReleases.clear(); + if (state.chatSessionMessageSubscription) { + subscriptions.add(state.chatSessionMessageSubscription); + } + state.chatSessionMessageSubscriptionRequestedKey = null; + state.chatSessionMessageSubscription = null; + const sessions = state.sessions; + if (!sessions?.unsubscribeMessages) { + return; + } + const unsubscribeMessages = sessions.unsubscribeMessages.bind(sessions); + for (const subscription of subscriptions) { + // A detached pane cannot drain another queue. Retry on its longer-lived + // session owner, but stop after terminal failures so timers cannot leak. + void (async () => { + let retryDelayMs = SESSION_MESSAGE_RELEASE_RETRY_MS; + for (let attempt = 0; attempt < MAX_SESSION_MESSAGE_RELEASE_ATTEMPTS; attempt += 1) { + try { + await unsubscribeMessages(subscription); + return; + } catch { + if (attempt + 1 === MAX_SESSION_MESSAGE_RELEASE_ATTEMPTS) { + return; + } + await new Promise((resolve) => { + globalThis.setTimeout(resolve, retryDelayMs); + }); + retryDelayMs = Math.min(retryDelayMs * 2, 30_000); + } + } + })(); + } +} + export async function syncSelectedSessionMessageSubscription( state: ChatSessionMessageSubscriptionState, opts?: { force?: boolean }, diff --git a/ui/src/pages/chat/chat-state-controller.ts b/ui/src/pages/chat/chat-state-controller.ts index a4c0d61c5cfb..630a82489043 100644 --- a/ui/src/pages/chat/chat-state-controller.ts +++ b/ui/src/pages/chat/chat-state-controller.ts @@ -1,5 +1,6 @@ import type { ReactiveController, ReactiveControllerHost } from "lit"; import type { ChatAttachment } from "../../lib/chat/chat-types.ts"; +import { disposeSelectedSessionMessageSubscription } from "./chat-history.ts"; import { subscribeChatOutboxProjection } from "./chat-queue.ts"; import type { ChatPageHost } from "./chat-state-host.ts"; import { invalidateImageLightbox } from "./chat-state-page.ts"; @@ -75,6 +76,7 @@ export class ChatStateController implements Reactiv attach(state: TState) { if (this.stateValue && this.stateValue !== state) { + disposeSelectedSessionMessageSubscription(this.stateValue); releaseChatMediaResourceSubscriber(this.stateValue.requestUpdate); this.attachmentReads.abortReads(); this.composerPersistence.stop(); @@ -360,6 +362,7 @@ export class ChatStateController implements Reactiv } const state = this.stateValue; if (state) { + disposeSelectedSessionMessageSubscription(state); releaseChatMediaResourceSubscriber(state.requestUpdate); cancelChatStreamRenderFrame(state); cancelChatScroll(state); diff --git a/ui/src/pages/chat/components/chat-message-media-lifecycle.test.ts b/ui/src/pages/chat/components/chat-message-media-lifecycle.test.ts index 0bf0a098b0a5..199d9e41aa85 100644 --- a/ui/src/pages/chat/components/chat-message-media-lifecycle.test.ts +++ b/ui/src/pages/chat/components/chat-message-media-lifecycle.test.ts @@ -10,6 +10,7 @@ import { observeChatMediaResource, readManagedImageBlobUrl, releaseChatMediaResourceSubscriber, + schedulePairingQrExpiryRefresh, type ImageRenderOptions, type RenderableImageBlock, } from "./chat-message-media.ts"; @@ -77,6 +78,50 @@ function observeSubscriber(subscriber: () => void): () => void { } describe("chat media resource lifecycle", () => { + it("refreshes every split pane when a shared pairing QR expires", async () => { + const message = { + content: [ + { + type: "openclaw_pairing_qr", + image_url: "data:image/png;base64,cXJwbmc=", + expiresAtMs: Date.now() + 1_000, + }, + ], + }; + const refreshFirst = observeSubscriber(vi.fn()); + const refreshSecond = observeSubscriber(vi.fn()); + + schedulePairingQrExpiryRefresh("shared-pairing-qr", message, refreshFirst); + schedulePairingQrExpiryRefresh("shared-pairing-qr", message, refreshSecond); + await vi.advanceTimersByTimeAsync(1_000); + + expect(refreshFirst).toHaveBeenCalledOnce(); + expect(refreshSecond).toHaveBeenCalledOnce(); + }); + + it("releases a pairing QR expiry timer when its chat pane disconnects", async () => { + const refresh = observeSubscriber(vi.fn()); + schedulePairingQrExpiryRefresh( + "disconnected-pairing-qr", + { + content: [ + { + type: "openclaw_pairing_qr", + image_url: "data:image/png;base64,cXJwbmc=", + expiresAtMs: Date.now() + 1_000, + }, + ], + }, + refresh, + ); + + releaseChatMediaResourceSubscriber(refresh); + await vi.advanceTimersByTimeAsync(1_000); + + expect(refresh).not.toHaveBeenCalled(); + expect(vi.getTimerCount()).toBe(0); + }); + it("wakes a managed image after one transient failure without an external render", async () => { const source = managedImageSource(); const blobUrl = installManagedImageUrls(); diff --git a/ui/src/pages/chat/components/chat-message-media.ts b/ui/src/pages/chat/components/chat-message-media.ts index 39929c9ae17b..a4877083a8b8 100644 --- a/ui/src/pages/chat/components/chat-message-media.ts +++ b/ui/src/pages/chat/components/chat-message-media.ts @@ -11,12 +11,6 @@ export type PairingQrExpiryNotice = { title: string; reason: string; }; -type PairingQrExpiryRefreshTimer = { - expiresAtMs: number; - onRequestUpdate: () => void; - timer: ReturnType; -}; -const pairingQrExpiryRefreshTimers = new Map(); export type ImageBlock = { url: string; @@ -48,7 +42,7 @@ export type RenderableImageBlock = ImageBlock & { export type AttachmentItem = Extract; -type ChatMediaResourceKind = "assistant-attachment" | "managed-image"; +type ChatMediaResourceKind = "assistant-attachment" | "managed-image" | "pairing-qr"; export type ChatMediaResource = { kind: ChatMediaResourceKind; @@ -530,41 +524,30 @@ function resolveNearestFuturePairingQrExpiresAtMs( return nearestExpiresAtMs; } -function clearPairingQrExpiryRefreshTimer(messageKey: string) { - const existing = pairingQrExpiryRefreshTimers.get(messageKey); - if (!existing) { - return; - } - clearTimeout(existing.timer); - pairingQrExpiryRefreshTimers.delete(messageKey); -} - export function schedulePairingQrExpiryRefresh( messageKey: string, message: unknown, onRequestUpdate: (() => void) | undefined, ) { - const nowMs = Date.now(); - const expiresAtMs = resolveNearestFuturePairingQrExpiresAtMs(message, nowMs); - const existing = pairingQrExpiryRefreshTimers.get(messageKey); - if (!expiresAtMs || !onRequestUpdate) { - if (existing) { - clearPairingQrExpiryRefreshTimer(messageKey); + if (!onRequestUpdate) { + return; + } + const refreshAt = resolveNearestFuturePairingQrExpiresAtMs(message); + if (refreshAt === undefined) { + const subscriber = chatMediaSubscribers.get(onRequestUpdate); + const resourceKey = chatMediaResourceKey("pairing-qr", messageKey); + const resource = subscriber?.resources.get(resourceKey); + if (subscriber && resource) { + subscriber.resources.delete(resourceKey); + detachChatMediaResourceSubscriber(resource, onRequestUpdate); + pruneChatMediaSubscriber(onRequestUpdate, subscriber); } return; } - if (existing?.expiresAtMs === expiresAtMs && existing.onRequestUpdate === onRequestUpdate) { - return; - } - clearPairingQrExpiryRefreshTimer(messageKey); - const timer = setTimeout( - () => { - pairingQrExpiryRefreshTimers.delete(messageKey); - onRequestUpdate(); - }, - Math.max(0, expiresAtMs - nowMs), + const resource = observeChatMediaResource("pairing-qr", messageKey, onRequestUpdate); + scheduleChatMediaResourceRefresh(resource, refreshAt, () => + notifyChatMediaResourceSubscribers(resource), ); - pairingQrExpiryRefreshTimers.set(messageKey, { expiresAtMs, onRequestUpdate, timer }); } export function extractTranscriptAttachments(message: unknown): AttachmentItem[] { diff --git a/ui/src/pages/chat/components/chat-message-pairing-qr-lifecycle.test.ts b/ui/src/pages/chat/components/chat-message-pairing-qr-lifecycle.test.ts new file mode 100644 index 000000000000..c8351d1f5173 --- /dev/null +++ b/ui/src/pages/chat/components/chat-message-pairing-qr-lifecycle.test.ts @@ -0,0 +1,104 @@ +/* @vitest-environment jsdom */ + +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; +import { + isChatMediaResourceCurrent, + observeChatMediaResource, + releaseChatMediaResourceSubscriber, + schedulePairingQrExpiryRefresh, +} from "./chat-message-media.ts"; + +const subscribers = new Set<() => void>(); + +beforeEach(() => { + vi.useFakeTimers(); + vi.setSystemTime(new Date("2026-07-28T00:00:00.000Z")); +}); + +afterEach(() => { + for (const subscriber of subscribers) { + releaseChatMediaResourceSubscriber(subscriber); + } + subscribers.clear(); + vi.clearAllTimers(); + vi.useRealTimers(); + vi.restoreAllMocks(); +}); + +function observeSubscriber(subscriber: () => void): () => void { + subscribers.add(subscriber); + return subscriber; +} + +function pairingQrMessage(expiresAtMs: number): { content: Record[] } { + return { + content: [ + { + type: "openclaw_pairing_qr", + image_url: "data:image/png;base64,cXJwbmc=", + expiresAtMs, + }, + ], + }; +} + +describe("pairing QR expiry resource lifecycle", () => { + it.each([ + { name: "ordinary messages", message: { content: [{ type: "text", text: "hello" }] } }, + { + name: "expired pairing QR messages", + message: () => pairingQrMessage(Date.now() - 1), + }, + ])("does not subscribe or schedule expiry for $name", ({ message }) => { + const messageKey = crypto.randomUUID(); + const refresh = observeSubscriber(vi.fn()); + const resolvedMessage = typeof message === "function" ? message() : message; + + schedulePairingQrExpiryRefresh(messageKey, resolvedMessage, refresh); + + expect(vi.getTimerCount()).toBe(0); + expect(observeChatMediaResource("pairing-qr", messageKey).subscribers.size).toBe(0); + + // Reattach the probe so normal subscriber cleanup also removes its resource. + schedulePairingQrExpiryRefresh(messageKey, pairingQrMessage(Date.now() + 1_000), refresh); + }); + + it("removes the last subscription when a previously active QR loses its expiry", async () => { + const messageKey = crypto.randomUUID(); + const refresh = observeSubscriber(vi.fn()); + + schedulePairingQrExpiryRefresh(messageKey, pairingQrMessage(Date.now() + 1_000), refresh); + const resource = observeChatMediaResource("pairing-qr", messageKey); + + schedulePairingQrExpiryRefresh(messageKey, { content: [] }, refresh); + await vi.advanceTimersByTimeAsync(1_000); + + expect(isChatMediaResourceCurrent(resource)).toBe(false); + expect(resource.subscribers.size).toBe(0); + expect(refresh).not.toHaveBeenCalled(); + expect(vi.getTimerCount()).toBe(0); + }); + + it("keeps the other split pane and unrelated media subscriptions alive", async () => { + const messageKey = crypto.randomUUID(); + const first = observeSubscriber(vi.fn()); + const second = observeSubscriber(vi.fn()); + const message = pairingQrMessage(Date.now() + 1_000); + const independentResource = observeChatMediaResource( + "assistant-attachment", + crypto.randomUUID(), + first, + ); + + schedulePairingQrExpiryRefresh(messageKey, message, first); + schedulePairingQrExpiryRefresh(messageKey, message, second); + schedulePairingQrExpiryRefresh(messageKey, { content: [] }, first); + await vi.advanceTimersByTimeAsync(1_000); + + expect(first).not.toHaveBeenCalled(); + expect(second).toHaveBeenCalledOnce(); + expect(independentResource.subscribers.has(first)).toBe(true); + expect(isChatMediaResourceCurrent(independentResource)).toBe(true); + expect(vi.getTimerCount()).toBe(0); + }); +}); diff --git a/ui/src/test-helpers/app-sidebar-cases/catalog-ownership.ts b/ui/src/test-helpers/app-sidebar-cases/catalog-ownership.ts index ce3b67ee576e..1a89d654f934 100644 --- a/ui/src/test-helpers/app-sidebar-cases/catalog-ownership.ts +++ b/ui/src/test-helpers/app-sidebar-cases/catalog-ownership.ts @@ -5,74 +5,86 @@ import { catalogPage, createGatewayHarness, createSessions, mountSidebar } from import "../../components/app-sidebar.ts"; describe("AppSidebar session catalog ownership", () => { - it("retires catalog rows and creation after reconnect loses its catalog owner", async () => { - vi.useFakeTimers(); - let provider: HTMLElement | undefined; - try { - const firstPage = catalogPage([{ threadId: "thread-1", name: "Retired session" }], "page-2"); - const catalog = firstPage.catalogs[0]; - if (!catalog) { - throw new Error("expected a session catalog"); + it.each([ + { owner: "the selected agent", assistantAgentId: null }, + { owner: "the advertised catalog capability", assistantAgentId: "main" }, + ])( + "retires catalog rows and creation after reconnect loses $owner", + async ({ assistantAgentId }) => { + vi.useFakeTimers(); + let provider: HTMLElement | undefined; + try { + const firstPage = catalogPage( + [{ threadId: "thread-1", name: "Retired session" }], + "page-2", + ); + const catalog = firstPage.catalogs[0]; + if (!catalog) { + throw new Error("expected a session catalog"); + } + catalog.capabilities.createSession = { model: "anthropic/claude-opus-4-8" }; + const expandedPage = catalogPage([{ threadId: "thread-2", name: "Retired page" }]); + const expandedCatalog = expandedPage.catalogs[0]; + if (!expandedCatalog) { + throw new Error("expected an expanded session catalog"); + } + expandedCatalog.capabilities.createSession = catalog.capabilities.createSession; + const request = vi + .fn() + .mockResolvedValueOnce(firstPage) + .mockResolvedValueOnce(expandedPage); + const gateway = createGatewayHarness({ request } as unknown as GatewayBrowserClient); + const catalogHello = { + type: "hello-ok", + protocol: 1, + auth: { role: "operator", scopes: ["operator.admin"] }, + features: { methods: ["sessions.catalog.list"] }, + } satisfies NonNullable; + gateway.publish({ hello: catalogHello }); + const mounted = await mountSidebar( + gateway.gateway, + createSessions("main", ["agent:main:main"]), + ); + const { sidebar } = mounted; + provider = mounted.provider; + sidebar.connected = true; + await sidebar.updateComplete; + await vi.advanceTimersByTimeAsync(0); + + await sidebar.sessionData.loadMoreSessionCatalog("codex"); + await sidebar.updateComplete; + expect(sidebar.textContent).toContain("Retired session"); + expect(sidebar.textContent).toContain("Retired page"); + expect(sidebar.querySelector(".sidebar-session-catalog-new")).not.toBeNull(); + expect(sidebar.sessionData.sessionCatalogPageDepths.size).toBe(1); + expect(sidebar.sessionData.sessionCatalogRevisions.size).toBe(1); + + gateway.publish({ phase: "reconnecting", hello: null }); + await sidebar.updateComplete; + expect(sidebar.textContent).toContain("Retired page"); + expect(sidebar.sessionData.sessionCatalogPageDepths.size).toBe(1); + + gateway.publish({ + phase: "connected", + assistantAgentId, + hello: { ...catalogHello, features: { ...catalogHello.features, methods: [] } }, + }); + await sidebar.updateComplete; + await vi.advanceTimersByTimeAsync(0); + await sidebar.updateComplete; + + expect(sidebar.sessionData.sessionCatalogAgentId).toBeNull(); + expect(sidebar.sessionData.sessionCatalogs).toEqual([]); + expect(sidebar.sessionData.sessionCatalogPageDepths.size).toBe(0); + expect(sidebar.sessionData.sessionCatalogRevisions.size).toBe(0); + expect(sidebar.textContent).not.toContain("Retired session"); + expect(sidebar.textContent).not.toContain("Retired page"); + expect(sidebar.querySelector(".sidebar-session-catalog-new")).toBeNull(); + expect(request).toHaveBeenCalledTimes(2); + } finally { + provider?.remove(); + vi.useRealTimers(); } - catalog.capabilities.createSession = { model: "anthropic/claude-opus-4-8" }; - const expandedPage = catalogPage([{ threadId: "thread-2", name: "Retired page" }]); - const expandedCatalog = expandedPage.catalogs[0]; - if (!expandedCatalog) { - throw new Error("expected an expanded session catalog"); - } - expandedCatalog.capabilities.createSession = catalog.capabilities.createSession; - const request = vi.fn().mockResolvedValueOnce(firstPage).mockResolvedValueOnce(expandedPage); - const gateway = createGatewayHarness({ request } as unknown as GatewayBrowserClient); - const catalogHello = { - type: "hello-ok", - protocol: 1, - auth: { role: "operator", scopes: ["operator.admin"] }, - features: { methods: ["sessions.catalog.list"] }, - } satisfies NonNullable; - gateway.publish({ hello: catalogHello }); - const mounted = await mountSidebar( - gateway.gateway, - createSessions("main", ["agent:main:main"]), - ); - const { sidebar } = mounted; - provider = mounted.provider; - sidebar.connected = true; - await sidebar.updateComplete; - await vi.advanceTimersByTimeAsync(0); - - await sidebar.sessionData.loadMoreSessionCatalog("codex"); - await sidebar.updateComplete; - expect(sidebar.textContent).toContain("Retired session"); - expect(sidebar.textContent).toContain("Retired page"); - expect(sidebar.querySelector(".sidebar-session-catalog-new")).not.toBeNull(); - expect(sidebar.sessionData.sessionCatalogPageDepths.size).toBe(1); - expect(sidebar.sessionData.sessionCatalogRevisions.size).toBe(1); - - gateway.publish({ phase: "reconnecting", hello: null }); - await sidebar.updateComplete; - expect(sidebar.textContent).toContain("Retired page"); - expect(sidebar.sessionData.sessionCatalogPageDepths.size).toBe(1); - - gateway.publish({ - phase: "connected", - assistantAgentId: null, - hello: { ...catalogHello, features: { ...catalogHello.features, methods: [] } }, - }); - await sidebar.updateComplete; - await vi.advanceTimersByTimeAsync(0); - await sidebar.updateComplete; - - expect(sidebar.sessionData.sessionCatalogAgentId).toBeNull(); - expect(sidebar.sessionData.sessionCatalogs).toEqual([]); - expect(sidebar.sessionData.sessionCatalogPageDepths.size).toBe(0); - expect(sidebar.sessionData.sessionCatalogRevisions.size).toBe(0); - expect(sidebar.textContent).not.toContain("Retired session"); - expect(sidebar.textContent).not.toContain("Retired page"); - expect(sidebar.querySelector(".sidebar-session-catalog-new")).toBeNull(); - expect(request).toHaveBeenCalledTimes(2); - } finally { - provider?.remove(); - vi.useRealTimers(); - } - }); + }, + ); }); diff --git a/ui/src/test-helpers/app-sidebar-cases/catalog-reconnect.ts b/ui/src/test-helpers/app-sidebar-cases/catalog-reconnect.ts new file mode 100644 index 000000000000..1e47730bb1bf --- /dev/null +++ b/ui/src/test-helpers/app-sidebar-cases/catalog-reconnect.ts @@ -0,0 +1,62 @@ +import { describe, expect, it, vi } from "vitest"; +import type { SessionsCatalogListResult } from "../../../../packages/gateway-protocol/src/index.ts"; +import { GatewayRequestError, type GatewayBrowserClient } from "../../api/gateway.ts"; +import type { ApplicationGatewaySnapshot } from "../../app/context.ts"; +import { + catalogPage, + createGatewayHarness, + createSessions, + deferred, + mountSidebar, +} from "../app-sidebar.ts"; + +describe("AppSidebar catalog reconnect", () => { + it("keeps progressive catalogs after a stale fallback and same-client reconnect", async () => { + vi.useFakeTimers(); + try { + const legacyFallback = deferred(); + const request = vi + .fn() + .mockRejectedValueOnce( + new GatewayRequestError({ + code: "INVALID_REQUEST", + message: + "invalid sessions.catalog.list params: at root: unexpected property 'progressId'", + }), + ) + .mockReturnValueOnce(legacyFallback.promise) + .mockResolvedValue(catalogPage([])); + const gateway = createGatewayHarness({ request } as unknown as GatewayBrowserClient); + const hello = { + features: { methods: ["sessions.catalog.list"] }, + } as ApplicationGatewaySnapshot["hello"]; + gateway.publish({ hello }); + const { sidebar } = await mountSidebar( + gateway.gateway, + createSessions("main", ["agent:main:main"]), + ); + sidebar.connected = true; + await sidebar.updateComplete; + await vi.advanceTimersByTimeAsync(0); + expect(request).toHaveBeenCalledTimes(2); + + gateway.publish({ phase: "reconnecting", hello: null }); + await sidebar.updateComplete; + gateway.publish({ phase: "connected", hello }); + await sidebar.updateComplete; + await vi.advanceTimersByTimeAsync(50); + + legacyFallback.resolve(catalogPage([])); + await vi.advanceTimersByTimeAsync(0); + await vi.advanceTimersByTimeAsync(30_000); + + expect(request).toHaveBeenLastCalledWith("sessions.catalog.list", { + agentId: "main", + limitPerHost: 40, + progressId: expect.any(String), + }); + } finally { + vi.useRealTimers(); + } + }); +}); diff --git a/ui/src/test-helpers/control-ui-e2e.ts b/ui/src/test-helpers/control-ui-e2e.ts index 5de9259164ad..192b36f6aa98 100644 --- a/ui/src/test-helpers/control-ui-e2e.ts +++ b/ui/src/test-helpers/control-ui-e2e.ts @@ -1409,6 +1409,7 @@ function installControlUiMockGateway( readonly protocol = ""; readyState = MockWebSocket.CONNECTING; readonly url: string; + private tickTimer: number | null = null; constructor(url: string | URL) { super(); @@ -1452,6 +1453,10 @@ function installControlUiMockGateway( return; } this.readyState = MockWebSocket.CLOSED; + if (this.tickTimer !== null) { + window.clearInterval(this.tickTimer); + this.tickTimer = null; + } sessionMessageSubscriptions.clear(); stopRepeatingSessionEvents(); this.dispatchEvent(new CloseEvent("close", { code, reason })); @@ -1481,6 +1486,11 @@ function installControlUiMockGateway( ? { id, ok: false, error: mockError, type: "res" } : { id, ok: true, payload, type: "res" }, ); + if (!mockError && method === "connect" && this.readyState === MockWebSocket.OPEN) { + this.tickTimer = window.setInterval(() => { + this.deliver({ event: "tick", payload: {}, seq: ++seq, type: "event" }); + }, 30_000); + } if (!mockError) { updateSessionMessageSubscription(method, frame.params); }