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
1 change: 1 addition & 0 deletions alias.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ export const alias = {
'devframe/rpc/transports/ws-client': r('devframe/src/rpc/transports/ws-client.ts'),
'devframe/rpc/client': r('devframe/src/rpc/client.ts'),
'devframe/rpc/dump': r('devframe/src/rpc/dump/index.ts'),
'devframe/rpc/shared-state': r('devframe/src/rpc/shared-state.ts'),
'devframe/rpc/server': r('devframe/src/rpc/server.ts'),
'devframe/rpc': r('devframe/src/rpc'),
'devframe/types': r('devframe/src/types/index.ts'),
Expand Down
13 changes: 13 additions & 0 deletions docs/content/8.references/5.browser-api.md
Original file line number Diff line number Diff line change
Expand Up @@ -72,3 +72,16 @@ The `error.code` values of `InPageChannelError`: [Errors and fallbacks](/guide/i
| `not-cloneable` | The port refused to clone a payload (`DataCloneError`) | Strip functions/DOM nodes/reactivity proxies, or declare `jsonSerializable: true` for the precise error above |
| `invalid-args` | Incoming arguments failed their Standard-Schema validation | The message lists the schema issues |
| `state-uninitialized` | The page script read a shared state before providing its `initialValue` | Initialize on first access |

## Shared state over custom RPC channels

`devframe/rpc/shared-state` exports the same `RpcSharedStateHost` implementations used by the RPC client and node context. Custom channel bindings can compose them with `RpcFunctionsCollectorBase`, `createRpcClient()` and `createRpcServer()`.

| Factory | Required connection members | Result |
|---------|-----------------------------|--------|
| `createRpcSharedStateClientHost(rpc)` | `call`, `callEvent`, function registration through `client`, trust events, `isTrusted`, and `connectionMeta.backend` | Native `get`, `keys`, `onKeyAdded` and `delete`. |
| `createRpcSharedStateServerHost(rpc)` | Function registration and filtered broadcasts | Native state publication, including `get(key, { sharedState })`. |

The channel binding owns peer authentication, serialization and disconnection. Each serving channel supplies birpc metadata with a `subscribedStates: Set<string>` and retains the default RPC `this` binding. Subscription handlers use that calling connection's metadata; filtered broadcasts reach its subscribed keys. Close the birpc connection when its transport disconnects and delete mirrored keys when disposing their owner.

A snapshot requested without an initial value rejects if its RPC call fails. Supplying an initial value retains the RPC client's existing immediate-state behavior.
2 changes: 1 addition & 1 deletion knip.jsonc
Original file line number Diff line number Diff line change
Expand Up @@ -137,7 +137,7 @@
"src/node/index.ts",
"src/node/{auth,hub-internals}/index.ts",
"src/recipes/interactive-auth.ts",
"src/rpc/{index,client,server}.ts",
"src/rpc/{index,client,server,shared-state}.ts",
"src/rpc/dump/index.ts",
"src/rpc/transports/{sse-client,sse-server,ws-bun,ws-deno,ws-client,ws-server}.ts",
"src/types/index.ts",
Expand Down
1 change: 1 addition & 0 deletions packages/devframe/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@
"./rpc/client": "./dist/rpc/client.mjs",
"./rpc/dump": "./dist/rpc/dump.mjs",
"./rpc/server": "./dist/rpc/server.mjs",
"./rpc/shared-state": "./dist/rpc/shared-state.mjs",
"./rpc/transports/sse-client": "./dist/rpc/transports/sse-client.mjs",
"./rpc/transports/sse-server": "./dist/rpc/transports/sse-server.mjs",
"./rpc/transports/ws-bun": "./dist/rpc/transports/ws-bun.mjs",
Expand Down
15 changes: 11 additions & 4 deletions packages/devframe/src/client/rpc-shared-state.ts
Original file line number Diff line number Diff line change
@@ -1,10 +1,17 @@
import type { RpcSharedStateGetOptions, RpcSharedStateHost } from 'devframe/types'
import type { RpcFunctionsCollector } from 'devframe/rpc'
import type { ConnectionMeta, DevframeRpcClientFunctions, RpcSharedStateGetOptions, RpcSharedStateHost } from 'devframe/types'
import type { SharedState, SharedStatePatch } from 'devframe/utils/shared-state'
import type { DevframeRpcClient } from './rpc'
import { createSharedState } from 'devframe/utils/shared-state'
import { DEVFRAME_EVENTS } from '../events'

export function createRpcSharedStateClientHost(rpc: DevframeRpcClient): RpcSharedStateHost {
/** Native shared-state synchronization over an authenticated RPC connection. */
export function createRpcSharedStateClientHost<Context>(
rpc: Pick<DevframeRpcClient, 'call' | 'callEvent' | 'events' | 'isTrusted'> & {
client: Pick<RpcFunctionsCollector<DevframeRpcClientFunctions, Context>, 'register'>
connectionMeta: Pick<ConnectionMeta, 'backend'>
},
): RpcSharedStateHost {
const sharedState = new Map<string, SharedState<any>>()
const stateDisposers = new Map<string, () => void>()
const initialValues = new Map<string, any>()
Expand Down Expand Up @@ -121,7 +128,7 @@ export function createRpcSharedStateClientHost(rpc: DevframeRpcClient): RpcShare
}
}

return new Promise<SharedState<T>>((resolve) => {
return new Promise<SharedState<T>>((resolve, reject) => {
if (!rpc.isTrusted) {
resolve(state)
let initialized = false
Expand All @@ -133,7 +140,7 @@ export function createRpcSharedStateClientHost(rpc: DevframeRpcClient): RpcShare
})
}
else {
initSharedState().then(resolve)
initSharedState().then(resolve, reject)
}
})
},
Expand Down
23 changes: 15 additions & 8 deletions packages/devframe/src/node/rpc-shared-state.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import type { RpcFunctionsHost, RpcSharedStateGetOptions, RpcSharedStateHost } from 'devframe/types'
import type { RpcFunctionsCollector } from 'devframe/rpc'
import type { DevframeNodeRpcSessionMeta, DevframeRpcServerFunctions, RpcFunctionsHost, RpcSharedStateGetOptions, RpcSharedStateHost } from 'devframe/types'
import type { SharedState, SharedStatePatch } from 'devframe/utils/shared-state'
import { createSharedState } from 'devframe/utils/shared-state'
import { createDebug } from 'obug'
Expand All @@ -8,8 +9,13 @@ import { diagnostics } from './diagnostics'
const debug = createDebug('devframe:rpc:state:changed')
const debugSubscribe = createDebug('devframe:rpc:state:subscribe')

export function createRpcSharedStateServerHost(
rpc: RpcFunctionsHost,
/**
* Publish native shared state over registered RPC functions and broadcasts.
* Each channel supplies session metadata with a `subscribedStates` set and
* retains birpc's default RPC `this` binding for subscription handlers.
*/
export function createRpcSharedStateServerHost<Context>(
rpc: Pick<RpcFunctionsCollector<DevframeRpcServerFunctions, Context>, 'register'> & Pick<RpcFunctionsHost, 'broadcast'>,
): RpcSharedStateHost {
const sharedState = new Map<string, SharedState<any>>()
const stateDisposers = new Map<string, () => void>()
Expand Down Expand Up @@ -93,12 +99,13 @@ export function createRpcSharedStateServerHost(
rpc.register({
name: 'devframe:rpc:server-state:subscribe',
type: 'event',
handler(key: string) {
const session = rpc.getCurrentRpcSession()
if (!session)
handler(this: { $meta?: DevframeNodeRpcSessionMeta } | undefined, key: string) {
/** birpc binds this handler to the calling connection across transports. */
const meta = this?.$meta
if (!meta)
return
debugSubscribe('subscribe', { key, session: session.meta.id })
session.meta.subscribedStates.add(key)
debugSubscribe('subscribe', { key, session: meta.id })
meta.subscribedStates.add(key)
},
})

Expand Down
144 changes: 144 additions & 0 deletions packages/devframe/src/rpc/shared-state.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,144 @@
import type { ChannelOptions } from 'birpc'
import type { RpcClientEvents } from 'devframe/client'
import type { DevframeRpcClientFunctions, DevframeRpcServerFunctions, RpcFunctionsHost } from 'devframe/types'
import type { MessagePort } from 'node:worker_threads'
import { MessageChannel } from 'node:worker_threads'
import { createHostContext } from 'devframe/node'
import { RpcFunctionsCollectorBase } from 'devframe/rpc'
import { createRpcClient } from 'devframe/rpc/client'
import { createRpcServer } from 'devframe/rpc/server'
import { createRpcSharedStateClientHost, createRpcSharedStateServerHost } from 'devframe/rpc/shared-state'
import { createEventEmitter } from 'devframe/utils/events'
import { structuredCloneDeserialize, structuredCloneSerialize } from 'devframe/utils/structured-clone'
import { expect, it } from 'vitest'
import { createContextRpcServer } from '../node/rpc-core'

/** JSON records also travel through transports that cannot clone native values. */
function channelFor(port: MessagePort): ChannelOptions {
return {
post: message => port.postMessage(message),
on: (handler) => {
port.on('message', handler)
},
off: (handler) => {
port.off('message', handler)
},
serialize: value => JSON.stringify(structuredCloneSerialize(value)),
deserialize: value => structuredCloneDeserialize(JSON.parse(value)),
}
}

async function createStateServer(mode: string) {
if (mode === 'node context') {
const context = await createHostContext({
cwd: process.cwd(),
mode: 'dev',
host: {
mountStatic() {},
resolveOrigin: () => 'http://localhost',
getStorageDir: () => process.cwd(),
},
})
const { rpcGroup } = createContextRpcServer({ context, auth: false })
return { group: rpcGroup, sharedState: context.rpc.sharedState }
}

const collector = new RpcFunctionsCollectorBase<DevframeRpcServerFunctions, undefined>(undefined)
const group = createRpcServer<DevframeRpcClientFunctions, DevframeRpcServerFunctions>(collector.functions)
const broadcast: RpcFunctionsHost['broadcast'] = async (options) => {
await Promise.all(group.clients
.filter(client => options.filter?.(client) !== false)
.map(client => client.$callRaw({ ...options, optional: true, event: true })))
}
const sharedState = createRpcSharedStateServerHost({ register: collector.register.bind(collector), broadcast })
return { group, sharedState }
}

it.each(['custom channels', 'node context'])('shares state across %s with per-connection subscriptions', async (mode) => {
expect.assertions(13)
const { group, sharedState } = await createStateServer(mode)
const counter = await sharedState.get('counter', { initialValue: { count: 1 } })
const channels = [new MessageChannel(), new MessageChannel()]
const peers = channels.map((channel, index) => {
const meta = { id: index, subscribedStates: new Set<string>() }
const serverChannel = { ...channelFor(channel.port1), meta }
group.updateChannels(current => current.push(serverChannel))
const client = new RpcFunctionsCollectorBase<DevframeRpcClientFunctions, undefined>(undefined)
const rpc = createRpcClient<DevframeRpcServerFunctions, DevframeRpcClientFunctions>(client.functions, {
channel: channelFor(channel.port2),
})
const state = createRpcSharedStateClientHost({
call: rpc.$call,
callEvent: rpc.$callEvent,
client,
isTrusted: true,
events: createEventEmitter<RpcClientEvents>(),
connectionMeta: { backend: 'none' },
})
return { meta, rpc, state, serverChannel }
})
const [first, second] = peers
try {
const firstCounter = await first.state.get<{ count: number }>('counter')
expect(firstCounter.value()).toEqual({ count: 1 })
expect(first.meta.subscribedStates.has('counter')).toBe(true)
expect(second.meta.subscribedStates.has('counter')).toBe(false)

const secondCounter = await second.state.get<{ count: number }>('counter')
expect(second.meta.subscribedStates.has('counter')).toBe(true)
firstCounter.mutate((draft) => {
draft.count = 2
})
await expect.poll(() => counter.value().count).toBe(2)
await expect.poll(() => secondCounter.value().count).toBe(2)

group.clients.find(client => client.$meta === first.meta)?.$close()
group.updateChannels(current => current.splice(current.indexOf(first.serverChannel), 1))
first.rpc.$close()
await expect(first.rpc.$call('devframe:rpc:server-state:get', 'counter')).rejects.toThrow()
secondCounter.mutate((draft) => {
draft.count = 3
})
await expect.poll(() => counter.value().count).toBe(3)
expect(firstCounter.value().count).toBe(2)
expect(sharedState.keys()).toEqual(['counter'])
expect(first.state.delete('counter')).toBe(true)
expect(sharedState.delete('counter')).toBe(true)
expect(sharedState.keys()).toEqual([])
}
finally {
for (const peer of peers) peer.rpc.$close()
for (const client of group.clients) client.$close()
group.updateChannels(current => current.splice(0))
for (const channel of channels) {
channel.port1.close()
channel.port2.close()
}
}
})

it('rejects an initial snapshot when its RPC connection closes', async () => {
expect.assertions(1)
const client = new RpcFunctionsCollectorBase<DevframeRpcClientFunctions, undefined>(undefined)
const channel = new MessageChannel()
const rpc = createRpcClient<DevframeRpcServerFunctions, DevframeRpcClientFunctions>(client.functions, {
channel: channelFor(channel.port1),
})
const state = createRpcSharedStateClientHost({
call: rpc.$call,
callEvent: rpc.$callEvent,
client,
isTrusted: true,
events: createEventEmitter<RpcClientEvents>(),
connectionMeta: { backend: 'none' },
})
try {
const pending = state.get('counter')
rpc.$close()
await expect(pending).rejects.toThrow('closed')
}
finally {
channel.port1.close()
channel.port2.close()
}
}, 1000)
2 changes: 2 additions & 0 deletions packages/devframe/src/rpc/shared-state.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
export { createRpcSharedStateClientHost } from '../client/rpc-shared-state'
export { createRpcSharedStateServerHost } from '../node/rpc-shared-state'
11 changes: 8 additions & 3 deletions packages/devframe/tsdown.config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,10 @@ const nodeDeps = {
const clientEntries = {
'client/index': 'src/client/index.ts',
'in-page-channel/index': 'src/in-page-channel/index.ts',
'rpc/index': 'src/rpc/index.ts',
'rpc/client': 'src/rpc/client.ts',
'rpc/server': 'src/rpc/server.ts',
'rpc/shared-state': 'src/rpc/shared-state.ts',
'utils/agent-tool-name': 'src/utils/agent-tool-name.ts',
'utils/colors': 'src/utils/colors.ts',
'utils/crypto-token': 'src/utils/crypto-token.ts',
Expand All @@ -75,10 +79,7 @@ const serverEntries = {
'index': 'src/index.ts',
'constants': 'src/constants.ts',
'types/index': 'src/types/index.ts',
'rpc/index': 'src/rpc/index.ts',
'rpc/client': 'src/rpc/client.ts',
'rpc/dump': 'src/rpc/dump/index.ts',
'rpc/server': 'src/rpc/server.ts',
'rpc/transports/sse-client': 'src/rpc/transports/sse-client.ts',
'rpc/transports/sse-server': 'src/rpc/transports/sse-server.ts',
'rpc/transports/ws-bun': 'src/rpc/transports/ws-bun.ts',
Expand Down Expand Up @@ -140,6 +141,10 @@ export default defineConfig([
await checkClientDist({
entries: [
resolve(distDir, 'client/index.mjs'),
resolve(distDir, 'rpc/index.mjs'),
resolve(distDir, 'rpc/client.mjs'),
resolve(distDir, 'rpc/server.mjs'),
resolve(distDir, 'rpc/shared-state.mjs'),
resolve(distDir, 'in-page-channel/index.mjs'),
resolve(distDir, 'utils/agent-tool-name.mjs'),
resolve(distDir, 'utils/colors.mjs'),
Expand Down
Loading
Loading