Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 14 additions & 1 deletion docs/content/1.guide/11.client.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`)

Expand Down
1 change: 1 addition & 0 deletions docs/content/8.references/3.events.md
Original file line number Diff line number Diff line change
Expand Up @@ -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. |
Expand Down
187 changes: 187 additions & 0 deletions packages/devframe/src/client/rpc-cache.test.ts
Original file line number Diff line number Diff line change
@@ -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<void>(resolve => server.listen(0, '127.0.0.1', resolve))
onTestFinished(async () => {
await devframe.close()
await new Promise<void>(resolve => server.close(() => resolve()))
})
const ctx = await devframe.context
if (legacy)
ctx.rpc.definitions.delete('devframe:rpc:cacheable-functions')
const counts: Record<string, number> = {}
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<unknown> }
}

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<string>()
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' })
})
})
47 changes: 43 additions & 4 deletions packages/devframe/src/client/rpc.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand Down Expand Up @@ -99,6 +99,7 @@ export interface DevframeRpcClientOptions extends SetupDevframeConnectionOptions
/** Channel overrides for the SSE transport, the `wsOptions` counterpart. */
sseOptions?: Partial<SseRpcChannelOptions>
rpcOptions?: Partial<BirpcOptions<DevframeRpcServerFunctions, DevframeRpcClientFunctions, boolean>>
/** Cache declared static/opted-in query methods with `true`, or select an explicit function list. */
cacheOptions?: boolean | Partial<RpcCacheOptions>
/**
* Mirror `agent`-flagged client RPC functions (functions registered on
Expand Down Expand Up @@ -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<void> | 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<any> {
const candidates = [
Expand Down Expand Up @@ -385,6 +406,7 @@ export async function getDevframeRpcClient(
})
}

let mode: DevframeRpcClientMode
const liveModeOptions = {
authToken,
connectionMeta,
Expand All @@ -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)
Expand All @@ -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,
})
Expand Down Expand Up @@ -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?.()
Expand Down
1 change: 1 addition & 0 deletions packages/devframe/src/events.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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',
Expand Down
13 changes: 13 additions & 0 deletions packages/devframe/src/node/host-functions.ts
Original file line number Diff line number Diff line change
@@ -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'
Expand Down Expand Up @@ -40,6 +42,17 @@ export class RpcFunctionsHostImpl extends RpcFunctionsCollectorBase<DevframeRpcS

this.sharedState = createRpcSharedStateServerHost(this)
this.streaming = createRpcStreamingServerHost(this)

this.register(defineRpcFunction({
name: 'devframe:rpc:cacheable-functions',
type: 'query',
handler: () => [...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
Expand Down
4 changes: 4 additions & 0 deletions packages/devframe/src/types/rpc-augments.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,8 @@
* To be extended
*/
export interface DevframeRpcClientFunctions {
/** Clear cached RPC results and refresh cache eligibility. */
'devframe:rpc:cache:invalidate': () => Promise<void>
/** Invoke a tool registered in this browser document. @internal */
'devframe:agent:invoke-client-tool': (id: string, args: Record<string, unknown>) => Promise<unknown>
/**
Expand Down Expand Up @@ -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<string[]>
/** 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<void>
/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down
2 changes: 2 additions & 0 deletions tests/__snapshots__/tsnapi/devframe/index.snapshot.d.ts
Original file line number Diff line number Diff line change
Expand Up @@ -227,6 +227,7 @@ export interface DevframeNodeRpcSessionMeta {
uploadingStreams?: Set<string>;
}
export interface DevframeRpcClientFunctions {
'devframe:rpc:cache:invalidate': () => Promise<void>;
'devframe:agent:invoke-client-tool': (_: string, _: Record<string, unknown>) => Promise<unknown>;
'devframe:auth:revoked': () => Promise<void>;
'devframe:streaming:chunk': (_: string, _: string, _: number, _: any) => Promise<void>;
Expand Down Expand Up @@ -256,6 +257,7 @@ export interface DevframeRpcOptions {
snapshot?: DevframeSnapshotRpcEntry[];
}
export interface DevframeRpcServerFunctions {
'devframe:rpc:cacheable-functions': () => Promise<string[]>;
'devframe:agent:sync-client-tools': (_: string, _: BrowserAgentToolManifest[]) => Promise<void>;
'anonymous:devframe:auth': (_: {
authToken: string;
Expand Down
Loading