From b7ef5b99d4172b11f5410ab43353f9ab2fc1b0fd Mon Sep 17 00:00:00 2001 From: Sakana <15715093608@163.com> Date: Mon, 5 Oct 2026 14:28:10 +0800 Subject: [PATCH 1/3] fix(rpc): honor automatic client caching Discover cacheable methods from server declarations and invalidate client caches when declarations or connection state change. Preserve explicit cache options and reject stale in-flight cache writes. --- docs/content/1.guide/11.client.md | 15 +- docs/content/8.references/3.events.md | 1 + .../devframe/src/client/rpc-cache.test.ts | 187 ++++++++++++++++++ packages/devframe/src/client/rpc.ts | 47 ++++- packages/devframe/src/events.ts | 1 + packages/devframe/src/node/host-functions.ts | 13 ++ packages/devframe/src/types/rpc-augments.ts | 4 + .../tsnapi/devframe/constants.snapshot.d.ts | 1 + .../tsnapi/devframe/index.snapshot.d.ts | 2 + 9 files changed, 266 insertions(+), 5 deletions(-) create mode 100644 packages/devframe/src/client/rpc-cache.test.ts diff --git a/docs/content/1.guide/11.client.md b/docs/content/1.guide/11.client.md index c26e6a53e..621c430ce 100644 --- a/docs/content/1.guide/11.client.md +++ b/docs/content/1.guide/11.client.md @@ -195,7 +195,20 @@ Set `cacheOptions: true` (or an object): const rpc = await connectDevframe({ cacheOptions: true }) ``` -`query` / `static` responses are memoized per argument hash; the `rpc:cache:invalidate` broadcast clears entries after a mutation. +With `true`, the client reads the server's RPC declarations and memoizes `static` responses and `query` responses marked `cacheable: true`, per argument hash. Actions and events continue to reach the server. An options object selects an explicit list, for example `{ functions: ['my:query'] }`, and can provide a `keySerializer`. + +After changing data that a cached query reads, clear connected clients' caches: + +```ts +import { DEVFRAME_EVENTS } from 'devframe/constants' + +await ctx.rpc.broadcast({ + method: DEVFRAME_EVENTS.broadcast.cacheInvalidate, + args: [], +}) +``` + +The client also clears its cache when RPC declarations or connection trust change, and when it disconnects or closes. A server without cache discovery support serves automatic-cache calls without caching. ## Discovery (`__connection.json`) diff --git a/docs/content/8.references/3.events.md b/docs/content/8.references/3.events.md index 96fb817f8..41464657c 100644 --- a/docs/content/8.references/3.events.md +++ b/docs/content/8.references/3.events.md @@ -105,6 +105,7 @@ Pushed to subscribed RPC clients, wired by the core node side. | Name | Carries | |---|---| +| `devframe:rpc:cache:invalidate` | Clear cached RPC results and refresh automatic cache eligibility. Sent when RPC declarations change, or by hosts after data changes. | | `devframe:auth:revoked` | This connection's bearer token was revoked; the RPC client drops to untrusted. | | `devframe:rpc:client-state:updated` | Full shared-state snapshot for a key. | | `devframe:rpc:client-state:patch` | Incremental shared-state patch for a key. | diff --git a/packages/devframe/src/client/rpc-cache.test.ts b/packages/devframe/src/client/rpc-cache.test.ts new file mode 100644 index 000000000..796800e0a --- /dev/null +++ b/packages/devframe/src/client/rpc-cache.test.ts @@ -0,0 +1,187 @@ +import type { AddressInfo } from 'node:net' +import type { DevframeRpcClientOptions } from './rpc' +import { createServer } from 'node:http' +import { defineDevframe, defineRpcFunction } from 'devframe' +import { DEVFRAME_EVENTS } from 'devframe/constants' +import { afterEach, beforeEach, describe, expect, it, onTestFinished, vi } from 'vitest' +import { initDevframe } from '../adapters/initiate' +import { connectDevframe } from './index' + +const connectionGlobals = ['__DEVFRAME_CONNECTION_META__', '__DEVFRAME_CONNECTION_AUTH_TOKEN__', '__DEVFRAME_CONNECTION__'] + +beforeEach(() => { + vi.stubGlobal('navigator', { userAgent: 'vitest' }) + for (const key of connectionGlobals) delete (globalThis as any)[key] +}) + +afterEach(() => { + vi.unstubAllGlobals() + for (const key of connectionGlobals) delete (globalThis as any)[key] +}) + +async function setup(transport: 'websocket' | 'sse', cacheOptions: DevframeRpcClientOptions['cacheOptions'] = true, legacy = false) { + const devframe = initDevframe(defineDevframe({ + id: 'cache-test', + name: 'Cache test', + version: '0.0.0', + packageName: 'cache-test', + homepage: 'https://example.test', + description: 'RPC cache regression tests.', + setup: () => {}, + }), { base: '/__cache/', auth: false }) + const server = createServer((req, res) => devframe.nodeMiddleware(req, res)) + devframe.attach(server) + await new Promise(resolve => server.listen(0, '127.0.0.1', resolve)) + onTestFinished(async () => { + await devframe.close() + await new Promise(resolve => server.close(() => resolve())) + }) + const ctx = await devframe.context + if (legacy) + ctx.rpc.definitions.delete('devframe:rpc:cacheable-functions') + const counts: Record = {} + for (const [name, type, cacheable] of [ + ['static', 'static', false], + ['query', 'query', true], + ['uncached', 'query', false], + ['default', 'query', undefined], + ['action', 'action', true], + ['event', 'event', true], + ] as const) { + ctx.rpc.register(defineRpcFunction({ + name: `test:${name}`, + type, + cacheable, + handler: (value: unknown) => { + counts[name] = (counts[name] ?? 0) + 1 + return value + }, + })) + } + const origin = `http://127.0.0.1:${(server.address() as AddressInfo).port}` + vi.stubGlobal('location', new URL(`${origin}/__cache/`)) + const client = await connectDevframe({ + baseURL: `${origin}/__cache/`, + transport, + cacheOptions, + otpParam: false, + simpleAuth: false, + webmcp: false, + callTimeout: 3000, + }) + onTestFinished(() => client.close?.()) + return { ctx, client, counts, call: client.call as (name: string, value?: unknown) => Promise } +} + +describe.each(['websocket', 'sse'] as const)('automatic RPC cache over %s', (transport) => { + it('caches static and opted-in queries from the first call, by arguments', async () => { + const { call, counts } = await setup(transport) + for (const name of ['static', 'query', 'uncached', 'default', 'action', 'event']) { + for (const value of [1, 1, 2]) + await expect(call(`test:${name}`, value)).resolves.toBe(value) + } + expect(counts).toEqual({ static: 2, query: 2, uncached: 3, default: 3, action: 3, event: 3 }) + }) + + it('caches falsy results and keeps event calls reaching the server', async () => { + const { client, call, counts } = await setup(transport) + for (const value of [false, 0, '', null, undefined]) { + await expect(call('test:query', value)).resolves.toBe(value) + await expect(call('test:query', value)).resolves.toBe(value) + } + expect(counts.query).toBe(5) + await client.callEvent('test:query' as any, 0) + await client.callEvent('test:query' as any, 0) + await vi.waitFor(() => expect(counts.query).toBe(7)) + }) + + it('preserves explicit function lists and custom serializers', async () => { + const keySerializer = vi.fn(() => 'same-key') + const { call, counts } = await setup(transport, { functions: ['test:uncached'], keySerializer }) + await expect(call('test:uncached', 1)).resolves.toBe(1) + await expect(call('test:uncached', 2)).resolves.toBe(1) + await call('test:query', 1) + await call('test:query', 1) + expect(counts).toEqual({ uncached: 1, query: 2 }) + expect(keySerializer).toHaveBeenCalled() + }) + + it('keeps caching disabled with false', async () => { + const { call, counts } = await setup(transport, false) + await call('test:query', 1) + await call('test:query', 1) + expect(counts.query).toBe(2) + }) + + it('does not cache rejected results', async () => { + const { ctx, call } = await setup(transport) + const handler = vi.fn() + .mockRejectedValueOnce(new Error('retry me')) + .mockResolvedValue('ok') + ctx.rpc.register(defineRpcFunction({ name: 'test:retry', type: 'query', cacheable: true, handler })) + await expect(call('test:retry')).rejects.toThrow('retry me') + await expect(call('test:retry')).resolves.toBe('ok') + await expect(call('test:retry')).resolves.toBe('ok') + expect(handler).toHaveBeenCalledTimes(2) + }) + + it('falls back to uncached calls with an older server', async () => { + const { call, counts } = await setup(transport, true, true) + await call('test:query', 1) + await call('test:query', 1) + expect(counts.query).toBe(2) + }) + + it('refreshes eligibility when a function is registered or updated', async () => { + const { ctx, client, call } = await setup(transport) + await call('test:query', 1) + const clear = vi.spyOn(client.cacheManager, 'clear') + const handler = vi.fn(() => 'new') + ctx.rpc.register(defineRpcFunction({ name: 'test:late', type: 'query', cacheable: true, handler })) + await vi.waitFor(() => expect(clear).toHaveBeenCalled()) + await expect(call('test:late')).resolves.toBe('new') + await expect(call('test:late')).resolves.toBe('new') + expect(handler).toHaveBeenCalledTimes(1) + clear.mockClear() + ctx.rpc.update(defineRpcFunction({ name: 'test:late', type: 'action', handler })) + await vi.waitFor(() => expect(clear).toHaveBeenCalled()) + await call('test:late') + await call('test:late') + expect(handler).toHaveBeenCalledTimes(3) + }) + + it('clears results on invalidation and rejects late cache writes', async () => { + const { ctx, client, call, counts } = await setup(transport) + const pending = Promise.withResolvers() + const handler = vi.fn(() => pending.promise) + ctx.rpc.register(defineRpcFunction({ name: 'test:slow', type: 'query', cacheable: true, handler })) + await call('test:query', 1) + const slow = call('test:slow') + await vi.waitFor(() => expect(handler).toHaveBeenCalledTimes(1)) + const clear = vi.spyOn(client.cacheManager, 'clear') + await ctx.rpc.broadcast({ method: DEVFRAME_EVENTS.broadcast.cacheInvalidate, args: [] }) + await vi.waitFor(() => expect(clear).toHaveBeenCalled()) + pending.resolve('old') + await slow + await call('test:slow') + expect(handler).toHaveBeenCalledTimes(2) + await call('test:query', 1) + expect(counts.query).toBe(2) + }) + + it('clears cached results when trust is revoked and when closed', async () => { + const { ctx, client, call } = await setup(transport) + await call('test:query', 1) + expect(client.cacheManager.has('test:query', [1])).toBe(true) + await client.services.state() + await ctx.rpc.broadcast({ method: DEVFRAME_EVENTS.broadcast.authRevoked, args: [] }) + await vi.waitFor(() => expect(client.isTrusted).toBe(false)) + expect(client.cacheManager.has('test:query', [1])).toBe(false) + await expect(call('test:query', 1)).rejects.toThrow(/Not authorized/) + await client.requestTrust() + await call('test:query', 1) + client.close?.() + expect(client.cacheManager.has('test:query', [1])).toBe(false) + await expect(call('test:query', 1)).rejects.toMatchObject({ kind: 'connection' }) + }) +}) diff --git a/packages/devframe/src/client/rpc.ts b/packages/devframe/src/client/rpc.ts index 4f09a287d..b76a2e345 100644 --- a/packages/devframe/src/client/rpc.ts +++ b/packages/devframe/src/client/rpc.ts @@ -7,7 +7,7 @@ import type { DevframeConnection, DevframeConnectionStatus, SetupDevframeConnect import type { DevframeServicesClient } from './rpc-services' import type { RpcStreamingClientHost } from './rpc-streaming' import type { DevframeScopedClientContext } from './scope' -import { DEVFRAME_OTP_URL_PARAM } from 'devframe/constants' +import { DEVFRAME_EVENTS, DEVFRAME_OTP_URL_PARAM } from 'devframe/constants' import { RpcCacheManager, RpcFunctionsCollectorBase } from 'devframe/rpc' import { createEventEmitter } from 'devframe/utils/events' import { withBase } from 'devframe/utils/url' @@ -99,6 +99,7 @@ export interface DevframeRpcClientOptions extends SetupDevframeConnectionOptions /** Channel overrides for the SSE transport, the `wsOptions` counterpart. */ sseOptions?: Partial rpcOptions?: Partial> + /** Cache declared static/opted-in query methods with `true`, or select an explicit function list. */ cacheOptions?: boolean | Partial /** * Mirror `agent`-flagged client RPC functions (functions registered on @@ -358,6 +359,26 @@ export async function getDevframeRpcClient( const disposeWebMcp = options.webmcp === false ? undefined : registerWebMcpTools(clientRpc) let disposeBrowserAgentBridge: (() => void) | undefined let closed = false + let cacheGeneration = 0 + let cacheFunctions: Promise | undefined + + function invalidateCache(): void { + cacheGeneration++ + cacheFunctions = undefined + cacheManager.clear() + if (cacheOptions === true) + cacheManager.updateOptions({ functions: [] }) + } + + if (cacheOptions) { + clientRpc.register({ + name: DEVFRAME_EVENTS.broadcast.cacheInvalidate, + type: 'event', + handler: invalidateCache, + }) + events.on(DEVFRAME_EVENTS.client.connectionStatus, invalidateCache) + events.on(DEVFRAME_EVENTS.client.isTrustedUpdated, invalidateCache) + } async function fetchJsonFromBases(path: string): Promise { const candidates = [ @@ -385,6 +406,7 @@ export async function getDevframeRpcClient( }) } + let mode: DevframeRpcClientMode const liveModeOptions = { authToken, connectionMeta, @@ -396,12 +418,28 @@ export async function getDevframeRpcClient( ...rpcOptions, async onRequest(req, next, resolve) { await rpcOptions.onRequest?.call(this, req, next, resolve) - if (cacheOptions && cacheManager?.validate(req.m)) { + // Auth and protocol calls must be able to run before cache discovery. + if (cacheOptions === true && req.i && mode.isTrusted && !req.m.startsWith('anonymous:') && req.m !== 'devframe:rpc:cacheable-functions') { + while (!cacheFunctions) { + if (closed || mode.status !== 'connected') + break + const generation = cacheGeneration + cacheFunctions = mode.callOptional('devframe:rpc:cacheable-functions').then((functions) => { + if (generation === cacheGeneration) + cacheManager.updateOptions({ functions: functions ?? [] }) + }) + await cacheFunctions + } + await cacheFunctions + } + const generation = cacheGeneration + if (cacheOptions && req.i && !closed && mode.isTrusted && cacheManager.validate(req.m)) { if (cacheManager.has(req.m, req.a)) { return resolve(cacheManager.cached(req.m, req.a)) } const res = await next(req) - cacheManager.apply(req, res) + if (generation === cacheGeneration && !closed) + cacheManager.apply(req, res) } else { await next(req) @@ -411,7 +449,7 @@ export async function getDevframeRpcClient( } const transport = resolveClientTransport(options.transport ?? 'auto', connectionMeta) - const mode = transport === 'static' + mode = transport === 'static' ? await createStaticRpcClientMode({ fetchJsonFromBases, }) @@ -450,6 +488,7 @@ export async function getDevframeRpcClient( /** Release authentication and transport resources even if another disposer fails. */ function closeRpcClient(): void { closed = true + invalidateCache() try { disposeBrowserAgentBridge?.() disposeWebMcp?.() diff --git a/packages/devframe/src/events.ts b/packages/devframe/src/events.ts index ebc958d67..61ad90022 100644 --- a/packages/devframe/src/events.ts +++ b/packages/devframe/src/events.ts @@ -51,6 +51,7 @@ export const DEVFRAME_EVENTS = { */ broadcast: { authRevoked: 'devframe:auth:revoked', + cacheInvalidate: 'devframe:rpc:cache:invalidate', clientStateUpdated: 'devframe:rpc:client-state:updated', clientStatePatch: 'devframe:rpc:client-state:patch', streamingChunk: 'devframe:streaming:chunk', diff --git a/packages/devframe/src/node/host-functions.ts b/packages/devframe/src/node/host-functions.ts index ea2b2c9c7..242c59e4f 100644 --- a/packages/devframe/src/node/host-functions.ts +++ b/packages/devframe/src/node/host-functions.ts @@ -1,6 +1,8 @@ import type { BirpcGroup } from 'birpc' import type { DevframeNodeContext, DevframeNodeRpcSession, DevframeNodeRpcSessionMeta, DevframeRpcClientFunctions, DevframeRpcServerFunctions, RpcBroadcastOptions, RpcFunctionsHost as RpcFunctionsHostType, RpcSharedStateHost, RpcStreamingHost } from 'devframe/types' import type { AsyncLocalStorage } from 'node:async_hooks' +import { defineRpcFunction } from 'devframe' +import { DEVFRAME_EVENTS } from 'devframe/constants' import { RpcFunctionsCollectorBase } from 'devframe/rpc' import { createDebug } from 'obug' import { removeClientAgentSession } from './client-agent' @@ -40,6 +42,17 @@ export class RpcFunctionsHostImpl extends RpcFunctionsCollectorBase [...this.definitions.values()] + .filter(fn => fn.type === 'static' || (fn.type === 'query' && fn.cacheable === true)) + .map(fn => fn.name), + })) + this.onChanged(() => { + void this.broadcast({ method: DEVFRAME_EVENTS.broadcast.cacheInvalidate, args: [] }) + }) } sharedState: RpcSharedStateHost diff --git a/packages/devframe/src/types/rpc-augments.ts b/packages/devframe/src/types/rpc-augments.ts index 05a5d2cf2..5f0c58094 100644 --- a/packages/devframe/src/types/rpc-augments.ts +++ b/packages/devframe/src/types/rpc-augments.ts @@ -2,6 +2,8 @@ * To be extended */ export interface DevframeRpcClientFunctions { + /** Clear cached RPC results and refresh cache eligibility. */ + 'devframe:rpc:cache:invalidate': () => Promise /** Invoke a tool registered in this browser document. @internal */ 'devframe:agent:invoke-client-tool': (id: string, args: Record) => Promise /** @@ -53,6 +55,8 @@ export interface DevframeRpcClientFunctions { * To be extended */ export interface DevframeRpcServerFunctions { + /** Methods eligible for automatic client caching. @internal */ + 'devframe:rpc:cacheable-functions': () => Promise /** Replace this connection's browser-agent tool manifest, tagged with the calling tab's stable client id. @internal */ 'devframe:agent:sync-client-tools': (clientId: string, tools: import('../client/browser-agent').BrowserAgentToolManifest[]) => Promise /** diff --git a/tests/__snapshots__/tsnapi/devframe/constants.snapshot.d.ts b/tests/__snapshots__/tsnapi/devframe/constants.snapshot.d.ts index d3f0cf78c..970f9e0b0 100644 --- a/tests/__snapshots__/tsnapi/devframe/constants.snapshot.d.ts +++ b/tests/__snapshots__/tsnapi/devframe/constants.snapshot.d.ts @@ -27,6 +27,7 @@ export declare const DEVFRAME_EVENTS: { }; readonly broadcast: { readonly authRevoked: "devframe:auth:revoked"; + readonly cacheInvalidate: "devframe:rpc:cache:invalidate"; readonly clientStateUpdated: "devframe:rpc:client-state:updated"; readonly clientStatePatch: "devframe:rpc:client-state:patch"; readonly streamingChunk: "devframe:streaming:chunk"; diff --git a/tests/__snapshots__/tsnapi/devframe/index.snapshot.d.ts b/tests/__snapshots__/tsnapi/devframe/index.snapshot.d.ts index 9d0ebbbfd..44ca6d835 100644 --- a/tests/__snapshots__/tsnapi/devframe/index.snapshot.d.ts +++ b/tests/__snapshots__/tsnapi/devframe/index.snapshot.d.ts @@ -227,6 +227,7 @@ export interface DevframeNodeRpcSessionMeta { uploadingStreams?: Set; } export interface DevframeRpcClientFunctions { + 'devframe:rpc:cache:invalidate': () => Promise; 'devframe:agent:invoke-client-tool': (_: string, _: Record) => Promise; 'devframe:auth:revoked': () => Promise; 'devframe:streaming:chunk': (_: string, _: string, _: number, _: any) => Promise; @@ -256,6 +257,7 @@ export interface DevframeRpcOptions { snapshot?: DevframeSnapshotRpcEntry[]; } export interface DevframeRpcServerFunctions { + 'devframe:rpc:cacheable-functions': () => Promise; 'devframe:agent:sync-client-tools': (_: string, _: BrowserAgentToolManifest[]) => Promise; 'anonymous:devframe:auth': (_: { authToken: string; From f6ae3e665def22aa9c60aa8ba442d526b369b198 Mon Sep 17 00:00:00 2001 From: Sakana <15715093608@163.com> Date: Wed, 7 Oct 2026 13:14:37 +0800 Subject: [PATCH 2/3] fix(rpc): retry cache discovery after failures Clear failed discovery promises so later RPC calls can retry. Preserve newer discovery after invalidation and keep timed-out actions from running. Cover rejection, timeout recovery, and invalidation races over WebSocket and SSE. Align caching documentation with repository terminology. --- docs/content/1.guide/11.client.md | 6 +- docs/content/8.references/3.events.md | 2 +- .../devframe/src/client/rpc-cache.test.ts | 73 +++++++++++++++++++ packages/devframe/src/client/rpc.ts | 5 ++ 4 files changed, 82 insertions(+), 4 deletions(-) diff --git a/docs/content/1.guide/11.client.md b/docs/content/1.guide/11.client.md index 621c430ce..6cf9a6618 100644 --- a/docs/content/1.guide/11.client.md +++ b/docs/content/1.guide/11.client.md @@ -195,9 +195,9 @@ Set `cacheOptions: true` (or an object): const rpc = await connectDevframe({ cacheOptions: true }) ``` -With `true`, the client reads the server's RPC declarations and memoizes `static` responses and `query` responses marked `cacheable: true`, per argument hash. Actions and events continue to reach the server. An options object selects an explicit list, for example `{ functions: ['my:query'] }`, and can provide a `keySerializer`. +With `true`, the RPC client reads declarations from the node side. It memoizes `static` responses and `query` responses marked `cacheable: true`, per argument hash. Actions and events continue to reach the node side. An options object selects an explicit list, for example `{ functions: ['my:query'] }`, and can provide a `keySerializer`. -After changing data that a cached query reads, clear connected clients' caches: +After changing data that a cached query reads, clear connected RPC clients' caches: ```ts import { DEVFRAME_EVENTS } from 'devframe/constants' @@ -208,7 +208,7 @@ await ctx.rpc.broadcast({ }) ``` -The client also clears its cache when RPC declarations or connection trust change, and when it disconnects or closes. A server without cache discovery support serves automatic-cache calls without caching. +The RPC client also clears its cache when RPC declarations or connection trust change, and when it disconnects or closes. Older devframes serve automatic-cache calls without caching when their node side lacks cache discovery support. If discovery fails, the waiting calls reject and a later RPC call retries discovery. ## Discovery (`__connection.json`) diff --git a/docs/content/8.references/3.events.md b/docs/content/8.references/3.events.md index 41464657c..fa05a6229 100644 --- a/docs/content/8.references/3.events.md +++ b/docs/content/8.references/3.events.md @@ -105,7 +105,7 @@ Pushed to subscribed RPC clients, wired by the core node side. | Name | Carries | |---|---| -| `devframe:rpc:cache:invalidate` | Clear cached RPC results and refresh automatic cache eligibility. Sent when RPC declarations change, or by hosts after data changes. | +| `devframe:rpc:cache:invalidate` | Clear cached RPC results and refresh automatic cache eligibility. Sent when RPC declarations change, or by the devframe's node side after data changes. | | `devframe:auth:revoked` | This connection's bearer token was revoked; the RPC client drops to untrusted. | | `devframe:rpc:client-state:updated` | Full shared-state snapshot for a key. | | `devframe:rpc:client-state:patch` | Incremental shared-state patch for a key. | diff --git a/packages/devframe/src/client/rpc-cache.test.ts b/packages/devframe/src/client/rpc-cache.test.ts index 796800e0a..6329323eb 100644 --- a/packages/devframe/src/client/rpc-cache.test.ts +++ b/packages/devframe/src/client/rpc-cache.test.ts @@ -132,6 +132,79 @@ describe.each(['websocket', 'sse'] as const)('automatic RPC cache over %s', (tra expect(counts.query).toBe(2) }) + it('retries discovery after a rejected request without blocking actions', async () => { + const { ctx, client, call, counts } = await setup(transport) + await call('test:query', 1) + const clear = vi.spyOn(client.cacheManager, 'clear') + const discovery = vi.fn() + .mockRejectedValueOnce(new Error('temporary discovery failure')) + .mockResolvedValue(['test:query']) + ctx.rpc.update(defineRpcFunction({ name: 'devframe:rpc:cacheable-functions', type: 'query', handler: discovery })) + await vi.waitFor(() => expect(clear).toHaveBeenCalled()) + + await expect(call('test:action', 1)).rejects.toThrow('temporary discovery failure') + expect(counts.action).toBeUndefined() + await expect(call('test:action', 2)).resolves.toBe(2) + await expect(call('test:action', 2)).resolves.toBe(2) + await call('test:query', 1) + await call('test:query', 1) + expect(counts).toEqual({ action: 2, query: 2 }) + expect(discovery).toHaveBeenCalledTimes(2) + }) + + it('retries discovery after a timeout and ignores its late response', async () => { + const { ctx, client, call, counts } = await setup(transport) + await call('test:query', 1) + const clear = vi.spyOn(client.cacheManager, 'clear') + const pending = Promise.withResolvers() + const discovery = vi.fn() + .mockReturnValueOnce(pending.promise) + .mockResolvedValue(['test:query']) + const timedOut = Promise.withResolvers() + client.events.on(DEVFRAME_EVENTS.client.error, (_error, method) => { + if (method === 'devframe:rpc:cacheable-functions') + timedOut.resolve() + }) + ctx.rpc.update(defineRpcFunction({ name: 'devframe:rpc:cacheable-functions', type: 'query', handler: discovery })) + await vi.waitFor(() => expect(clear).toHaveBeenCalled()) + + await expect(call('test:action', 1)).rejects.toMatchObject({ kind: 'timeout' }) + await timedOut.promise + // The error event fires before the discovery promise rejects. + await new Promise(resolve => setTimeout(resolve, 0)) + pending.resolve(['test:action']) + expect(counts.action).toBeUndefined() + await expect(call('test:action', 2)).resolves.toBe(2) + await expect(call('test:action', 2)).resolves.toBe(2) + await call('test:query', 1) + await call('test:query', 1) + expect(counts).toEqual({ action: 2, query: 2 }) + expect(discovery).toHaveBeenCalledTimes(2) + }) + + it('keeps refreshed discovery when an invalidated request rejects', async () => { + const { ctx, client, call } = await setup(transport) + await call('test:query', 1) + const clear = vi.spyOn(client.cacheManager, 'clear') + const pending = Promise.withResolvers() + const discovery = vi.fn() + .mockReturnValueOnce(pending.promise) + .mockResolvedValue(['test:query']) + ctx.rpc.update(defineRpcFunction({ name: 'devframe:rpc:cacheable-functions', type: 'query', handler: discovery })) + await vi.waitFor(() => expect(clear).toHaveBeenCalled()) + + const rejected = expect(call('test:action', 1)).rejects.toThrow('outdated discovery') + await vi.waitFor(() => expect(discovery).toHaveBeenCalledTimes(1)) + clear.mockClear() + await ctx.rpc.broadcast({ method: DEVFRAME_EVENTS.broadcast.cacheInvalidate, args: [] }) + await vi.waitFor(() => expect(clear).toHaveBeenCalled()) + await call('test:query', 1) + pending.reject(new Error('outdated discovery')) + await rejected + await call('test:query', 1) + expect(discovery).toHaveBeenCalledTimes(2) + }) + it('refreshes eligibility when a function is registered or updated', async () => { const { ctx, client, call } = await setup(transport) await call('test:query', 1) diff --git a/packages/devframe/src/client/rpc.ts b/packages/devframe/src/client/rpc.ts index b76a2e345..e27716705 100644 --- a/packages/devframe/src/client/rpc.ts +++ b/packages/devframe/src/client/rpc.ts @@ -427,6 +427,11 @@ export async function getDevframeRpcClient( cacheFunctions = mode.callOptional('devframe:rpc:cacheable-functions').then((functions) => { if (generation === cacheGeneration) cacheManager.updateOptions({ functions: functions ?? [] }) + }).catch((error) => { + // Retry on a later call without resetting discovery after invalidation. + if (generation === cacheGeneration) + cacheFunctions = undefined + throw error }) await cacheFunctions } From 0fa73c09a88286816f22bd895a7f3a7724665563 Mon Sep 17 00:00:00 2001 From: Sakana <15715093608@163.com> Date: Wed, 7 Oct 2026 13:36:11 +0800 Subject: [PATCH 3/3] fix(rpc): stop expired calls before cache discovery dispatch Keep each call within its original timeout across cache discovery refreshes. Stop expired calls before dispatch while allowing later retries to use the shared discovery result. Cover call and callOptional over WebSocket and SSE. All 28 cache tests and the full suite pass, along with lint, knip, typecheck, and build. --- .../devframe/src/client/rpc-cache.test.ts | 34 +++++++++++++++++-- packages/devframe/src/client/rpc.ts | 12 ++++++- 2 files changed, 43 insertions(+), 3 deletions(-) diff --git a/packages/devframe/src/client/rpc-cache.test.ts b/packages/devframe/src/client/rpc-cache.test.ts index 6329323eb..bf78a3bbb 100644 --- a/packages/devframe/src/client/rpc-cache.test.ts +++ b/packages/devframe/src/client/rpc-cache.test.ts @@ -19,7 +19,7 @@ afterEach(() => { for (const key of connectionGlobals) delete (globalThis as any)[key] }) -async function setup(transport: 'websocket' | 'sse', cacheOptions: DevframeRpcClientOptions['cacheOptions'] = true, legacy = false) { +async function setup(transport: 'websocket' | 'sse', cacheOptions: DevframeRpcClientOptions['cacheOptions'] = true, legacy = false, callTimeout = 3000) { const devframe = initDevframe(defineDevframe({ id: 'cache-test', name: 'Cache test', @@ -67,7 +67,7 @@ async function setup(transport: 'websocket' | 'sse', cacheOptions: DevframeRpcCl otpParam: false, simpleAuth: false, webmcp: false, - callTimeout: 3000, + callTimeout, }) onTestFinished(() => client.close?.()) return { ctx, client, counts, call: client.call as (name: string, value?: unknown) => Promise } @@ -182,6 +182,36 @@ describe.each(['websocket', 'sse'] as const)('automatic RPC cache over %s', (tra expect(discovery).toHaveBeenCalledTimes(2) }) + it.each(['call', 'callOptional'] as const)('does not send an expired %s after rediscovery, but allows a later retry', async (method) => { + const { ctx, client, call, counts } = await setup(transport, true, false, 800) + await call('test:query', 1) + const clear = vi.spyOn(client.cacheManager, 'clear') + const first = Promise.withResolvers() + const second = Promise.withResolvers() + const discovery = vi.fn() + .mockReturnValueOnce(first.promise) + .mockReturnValueOnce(second.promise) + ctx.rpc.update(defineRpcFunction({ name: 'devframe:rpc:cacheable-functions', type: 'query', handler: discovery })) + await vi.waitFor(() => expect(clear).toHaveBeenCalled()) + + const expired = expect(client[method]('test:action' as any, 1)).rejects.toMatchObject({ kind: 'timeout' }) + await vi.waitFor(() => expect(discovery).toHaveBeenCalledTimes(1)) + clear.mockClear() + await ctx.rpc.broadcast({ method: DEVFRAME_EVENTS.broadcast.cacheInvalidate, args: [] }) + await vi.waitFor(() => expect(clear).toHaveBeenCalled()) + await new Promise(resolve => setTimeout(resolve, 400)) + first.resolve(['test:query']) + await vi.waitFor(() => expect(discovery).toHaveBeenCalledTimes(2)) + + await expired + expect(counts.action).toBeUndefined() + const retry = client[method]('test:action' as any, 1) + second.resolve(['test:query']) + await expect(retry).resolves.toBe(1) + expect(counts.action).toBe(1) + expect(discovery).toHaveBeenCalledTimes(2) + }) + it('keeps refreshed discovery when an invalidated request rejects', async () => { const { ctx, client, call } = await setup(transport) await call('test:query', 1) diff --git a/packages/devframe/src/client/rpc.ts b/packages/devframe/src/client/rpc.ts index e27716705..65b2bb758 100644 --- a/packages/devframe/src/client/rpc.ts +++ b/packages/devframe/src/client/rpc.ts @@ -11,7 +11,7 @@ import { DEVFRAME_EVENTS, DEVFRAME_OTP_URL_PARAM } from 'devframe/constants' import { RpcCacheManager, RpcFunctionsCollectorBase } from 'devframe/rpc' import { createEventEmitter } from 'devframe/utils/events' import { withBase } from 'devframe/utils/url' -import { setupDevframeConnection } from './connection' +import { DevframeConnectionError, setupDevframeConnection } from './connection' import { storeAuthToken } from './connection-storage' import { authenticateWithUrlOtp } from './otp' import { createDevframeServicesClient } from './rpc-services' @@ -406,6 +406,8 @@ export async function getDevframeRpcClient( }) } + const timeout = options.callTimeout ?? 0 + const requestBudget = timeout > 0 ? timeout : Infinity let mode: DevframeRpcClientMode const liveModeOptions = { authToken, @@ -417,10 +419,17 @@ export async function getDevframeRpcClient( rpcOptions: { ...rpcOptions, async onRequest(req, next, resolve) { + const deadline = performance.now() + requestBudget await rpcOptions.onRequest?.call(this, req, next, resolve) // Auth and protocol calls must be able to run before cache discovery. if (cacheOptions === true && req.i && mode.isTrusted && !req.m.startsWith('anonymous:') && req.m !== 'devframe:rpc:cacheable-functions') { + // Rediscovery shares this call's deadline, even when another call keeps waiting. + function checkDeadline(): void { + if (performance.now() >= deadline) + throw new DevframeConnectionError('timeout', `[devframe] RPC call "${req.m}" timed out after ${timeout}ms`) + } while (!cacheFunctions) { + checkDeadline() if (closed || mode.status !== 'connected') break const generation = cacheGeneration @@ -436,6 +445,7 @@ export async function getDevframeRpcClient( await cacheFunctions } await cacheFunctions + checkDeadline() } const generation = cacheGeneration if (cacheOptions && req.i && !closed && mode.isTrusted && cacheManager.validate(req.m)) {