Skip to content

Commit fe0fb62

Browse files
committed
feat(rpc): expose native shared state for custom channels
Tracked in dvcol/devkit-extension#7. Scoped lint, type checks, native RPC tests and browser artifact checks pass.
1 parent 9aa752a commit fe0fb62

13 files changed

Lines changed: 182 additions & 20 deletions

File tree

‎alias.ts‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@ export const alias = {
1616
'devframe/rpc/transports/ws-client': r('devframe/src/rpc/transports/ws-client.ts'),
1717
'devframe/rpc/client': r('devframe/src/rpc/client.ts'),
1818
'devframe/rpc/dump': r('devframe/src/rpc/dump/index.ts'),
19+
'devframe/rpc/shared-state': r('devframe/src/rpc/shared-state.ts'),
1920
'devframe/rpc/server': r('devframe/src/rpc/server.ts'),
2021
'devframe/rpc': r('devframe/src/rpc'),
2122
'devframe/types': r('devframe/src/types/index.ts'),

‎docs/content/6.errors/DF0013.md‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -17,4 +17,4 @@ Pass `initialValue` on the first call: `ctx.rpc.sharedState.get(key, { initialVa
1717

1818
## Source
1919

20-
- [`packages/devframe/src/node/rpc-shared-state.ts`](https://github.com/devframes/devframe/blob/main/packages/devframe/src/node/rpc-shared-state.ts): `RpcSharedStateHost.get()` throws `DF0013` when neither an existing entry nor an `initialValue` is provided for a key.
20+
- [`packages/devframe/src/rpc/shared-state-server.ts`](https://github.com/devframes/devframe/blob/main/packages/devframe/src/rpc/shared-state-server.ts): `RpcSharedStateHost.get()` throws `DF0013` when neither an existing entry nor an `initialValue` is provided for a key.

‎docs/content/8.references/5.browser-api.md‎

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -72,3 +72,16 @@ The `error.code` values of `InPageChannelError`: [Errors and fallbacks](/guide/i
7272
| `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 |
7373
| `invalid-args` | Incoming arguments failed their Standard-Schema validation | The message lists the schema issues |
7474
| `state-uninitialized` | The page script read a shared state before providing its `initialValue` | Initialize on first access |
75+
76+
## Shared state over custom RPC channels
77+
78+
`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()`.
79+
80+
| Factory | Required connection members | Result |
81+
|---------|-----------------------------|--------|
82+
| `createRpcSharedStateClientHost(rpc)` | `call`, `callEvent`, function registration through `client`, trust events, `isTrusted`, and `connectionMeta.backend` | Native `get`, `keys`, `onKeyAdded` and `delete`. |
83+
| `createRpcSharedStateServerHost(rpc)` | Function registration and filtered broadcasts | Native state publication, including `get(key, { sharedState })`. |
84+
85+
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.
86+
87+
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.

‎packages/devframe/package.json‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -38,6 +38,7 @@
3838
"./rpc/client": "./dist/rpc/client.mjs",
3939
"./rpc/dump": "./dist/rpc/dump.mjs",
4040
"./rpc/server": "./dist/rpc/server.mjs",
41+
"./rpc/shared-state": "./dist/rpc/shared-state.mjs",
4142
"./rpc/transports/sse-client": "./dist/rpc/transports/sse-client.mjs",
4243
"./rpc/transports/sse-server": "./dist/rpc/transports/sse-server.mjs",
4344
"./rpc/transports/ws-bun": "./dist/rpc/transports/ws-bun.mjs",

‎packages/devframe/src/client/rpc-shared-state.ts‎

Lines changed: 11 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,10 +1,17 @@
1-
import type { RpcSharedStateGetOptions, RpcSharedStateHost } from 'devframe/types'
1+
import type { RpcFunctionsCollector } from 'devframe/rpc'
2+
import type { ConnectionMeta, DevframeRpcClientFunctions, RpcSharedStateGetOptions, RpcSharedStateHost } from 'devframe/types'
23
import type { SharedState, SharedStatePatch } from 'devframe/utils/shared-state'
34
import type { DevframeRpcClient } from './rpc'
45
import { createSharedState } from 'devframe/utils/shared-state'
56
import { DEVFRAME_EVENTS } from '../events'
67

7-
export function createRpcSharedStateClientHost(rpc: DevframeRpcClient): RpcSharedStateHost {
8+
/** Native shared-state synchronization over an authenticated RPC connection. */
9+
export function createRpcSharedStateClientHost<Context>(
10+
rpc: Pick<DevframeRpcClient, 'call' | 'callEvent' | 'events' | 'isTrusted'> & {
11+
client: Pick<RpcFunctionsCollector<DevframeRpcClientFunctions, Context>, 'register'>
12+
connectionMeta: Pick<ConnectionMeta, 'backend'>
13+
},
14+
): RpcSharedStateHost {
815
const sharedState = new Map<string, SharedState<any>>()
916
const stateDisposers = new Map<string, () => void>()
1017
const initialValues = new Map<string, any>()
@@ -121,7 +128,7 @@ export function createRpcSharedStateClientHost(rpc: DevframeRpcClient): RpcShare
121128
}
122129
}
123130

124-
return new Promise<SharedState<T>>((resolve) => {
131+
return new Promise<SharedState<T>>((resolve, reject) => {
125132
if (!rpc.isTrusted) {
126133
resolve(state)
127134
let initialized = false
@@ -133,7 +140,7 @@ export function createRpcSharedStateClientHost(rpc: DevframeRpcClient): RpcShare
133140
})
134141
}
135142
else {
136-
initSharedState().then(resolve)
143+
initSharedState().then(resolve, reject)
137144
}
138145
})
139146
},

‎packages/devframe/src/node/diagnostics.ts‎

Lines changed: 0 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -20,9 +20,6 @@ export const diagnostics = defineDiagnostics({
2020
DF0012: {
2121
why: (p: { filepath: string }) => `Failed to parse storage file: ${p.filepath}, falling back to defaults.`,
2222
},
23-
DF0013: {
24-
why: (p: { key: string }) => `Shared state of "${p.key}" is not found, please provide an initial value for the first time`,
25-
},
2623
DF0014: {
2724
why: (p: { name: string }) => `RPC function "${p.name}" has an invalid \`agent\` field: \`description\` must be a non-empty string.`,
2825
fix: 'Provide a short description (~1–3 sentences) explaining what the tool does and when agents should invoke it.',

‎packages/devframe/src/node/host-functions.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -3,9 +3,9 @@ import type { DevframeNodeContext, DevframeNodeRpcSession, DevframeNodeRpcSessio
33
import type { AsyncLocalStorage } from 'node:async_hooks'
44
import { RpcFunctionsCollectorBase } from 'devframe/rpc'
55
import { createDebug } from 'obug'
6+
import { createRpcSharedStateServerHost } from '../rpc/shared-state-server'
67
import { removeClientAgentSession } from './client-agent'
78
import { diagnostics } from './diagnostics'
8-
import { createRpcSharedStateServerHost } from './rpc-shared-state'
99
import { createRpcStreamingServerHost } from './rpc-streaming'
1010

1111
const debugBroadcast = createDebug('devframe:rpc:broadcast')

‎packages/devframe/src/rpc/diagnostics.ts‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,9 @@ import { defineDiagnostics } from 'devframe/utils/nostics'
33
export const diagnostics = defineDiagnostics({
44
docsBase: 'https://devfra.me/errors',
55
codes: {
6+
DF0013: {
7+
why: (p: { key: string }) => `Shared state of "${p.key}" is not found, please provide an initial value for the first time`,
8+
},
69
DF0019: {
710
why: (p: { name: string }) =>
811
`RPC function "${p.name}" has \`agent\` set but \`jsonSerializable\` is \`false\`; MCP requires JSON-serializable data.`,

packages/devframe/src/node/rpc-shared-state.ts renamed to packages/devframe/src/rpc/shared-state-server.ts

Lines changed: 15 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
1-
import type { RpcFunctionsHost, RpcSharedStateGetOptions, RpcSharedStateHost } from 'devframe/types'
1+
import type { RpcFunctionsCollector } from 'devframe/rpc'
2+
import type { DevframeNodeRpcSessionMeta, DevframeRpcServerFunctions, RpcFunctionsHost, RpcSharedStateGetOptions, RpcSharedStateHost } from 'devframe/types'
23
import type { SharedState, SharedStatePatch } from 'devframe/utils/shared-state'
34
import { createSharedState } from 'devframe/utils/shared-state'
45
import { createDebug } from 'obug'
@@ -8,8 +9,13 @@ import { diagnostics } from './diagnostics'
89
const debug = createDebug('devframe:rpc:state:changed')
910
const debugSubscribe = createDebug('devframe:rpc:state:subscribe')
1011

11-
export function createRpcSharedStateServerHost(
12-
rpc: RpcFunctionsHost,
12+
/**
13+
* Publish native shared state over registered RPC functions and broadcasts.
14+
* Each channel supplies session metadata with a `subscribedStates` set and
15+
* retains birpc's default RPC `this` binding for subscription handlers.
16+
*/
17+
export function createRpcSharedStateServerHost<Context>(
18+
rpc: Pick<RpcFunctionsCollector<DevframeRpcServerFunctions, Context>, 'register'> & Pick<RpcFunctionsHost, 'broadcast'>,
1319
): RpcSharedStateHost {
1420
const sharedState = new Map<string, SharedState<any>>()
1521
const stateDisposers = new Map<string, () => void>()
@@ -93,12 +99,13 @@ export function createRpcSharedStateServerHost(
9399
rpc.register({
94100
name: 'devframe:rpc:server-state:subscribe',
95101
type: 'event',
96-
handler(key: string) {
97-
const session = rpc.getCurrentRpcSession()
98-
if (!session)
102+
handler(this: { $meta?: DevframeNodeRpcSessionMeta } | undefined, key: string) {
103+
/** birpc binds this handler to the calling connection across transports. */
104+
const meta = this?.$meta
105+
if (!meta)
99106
return
100-
debugSubscribe('subscribe', { key, session: session.meta.id })
101-
session.meta.subscribedStates.add(key)
107+
debugSubscribe('subscribe', { key, session: meta.id })
108+
meta.subscribedStates.add(key)
102109
},
103110
})
104111

Lines changed: 123 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,123 @@
1+
import type { ChannelOptions } from 'birpc'
2+
import type { RpcClientEvents } from 'devframe/client'
3+
import type { DevframeRpcClientFunctions, DevframeRpcServerFunctions, RpcFunctionsHost } from 'devframe/types'
4+
import type { MessagePort } from 'node:worker_threads'
5+
import { MessageChannel } from 'node:worker_threads'
6+
import { RpcFunctionsCollectorBase } from 'devframe/rpc'
7+
import { createRpcClient } from 'devframe/rpc/client'
8+
import { createRpcServer } from 'devframe/rpc/server'
9+
import { createRpcSharedStateClientHost, createRpcSharedStateServerHost } from 'devframe/rpc/shared-state'
10+
import { createEventEmitter } from 'devframe/utils/events'
11+
import { structuredCloneDeserialize, structuredCloneSerialize } from 'devframe/utils/structured-clone'
12+
import { expect, it } from 'vitest'
13+
14+
/** JSON records also travel through transports that cannot clone native values. */
15+
function channelFor(port: MessagePort): ChannelOptions {
16+
return {
17+
post: message => port.postMessage(message),
18+
on: (handler) => {
19+
port.on('message', handler)
20+
},
21+
off: (handler) => {
22+
port.off('message', handler)
23+
},
24+
serialize: value => JSON.stringify(structuredCloneSerialize(value)),
25+
deserialize: value => structuredCloneDeserialize(JSON.parse(value)),
26+
}
27+
}
28+
29+
it('shares native state across independent channels with per-connection subscriptions', async () => {
30+
expect.assertions(13)
31+
const collector = new RpcFunctionsCollectorBase<DevframeRpcServerFunctions, undefined>(undefined)
32+
const group = createRpcServer<DevframeRpcClientFunctions, DevframeRpcServerFunctions>(collector.functions)
33+
const broadcast: RpcFunctionsHost['broadcast'] = async (options) => {
34+
await Promise.all(group.clients
35+
.filter(client => options.filter?.(client) !== false)
36+
.map(client => client.$callRaw({ ...options, optional: true, event: true })))
37+
}
38+
const sharedState = createRpcSharedStateServerHost({ register: collector.register.bind(collector), broadcast })
39+
const counter = await sharedState.get('counter', { initialValue: { count: 1 } })
40+
const channels = [new MessageChannel(), new MessageChannel()]
41+
const peers = channels.map((channel, index) => {
42+
const meta = { id: index, subscribedStates: new Set<string>() }
43+
const serverChannel = { ...channelFor(channel.port1), meta }
44+
group.updateChannels(current => current.push(serverChannel))
45+
const client = new RpcFunctionsCollectorBase<DevframeRpcClientFunctions, undefined>(undefined)
46+
const rpc = createRpcClient<DevframeRpcServerFunctions, DevframeRpcClientFunctions>(client.functions, {
47+
channel: channelFor(channel.port2),
48+
})
49+
const state = createRpcSharedStateClientHost({
50+
call: rpc.$call,
51+
callEvent: rpc.$callEvent,
52+
client,
53+
isTrusted: true,
54+
events: createEventEmitter<RpcClientEvents>(),
55+
connectionMeta: { backend: 'none' },
56+
})
57+
return { meta, rpc, state, serverChannel }
58+
})
59+
const [first, second] = peers
60+
try {
61+
const firstCounter = await first.state.get<{ count: number }>('counter')
62+
expect(firstCounter.value()).toEqual({ count: 1 })
63+
expect(first.meta.subscribedStates.has('counter')).toBe(true)
64+
expect(second.meta.subscribedStates.has('counter')).toBe(false)
65+
66+
const secondCounter = await second.state.get<{ count: number }>('counter')
67+
expect(second.meta.subscribedStates.has('counter')).toBe(true)
68+
firstCounter.mutate((draft) => {
69+
draft.count = 2
70+
})
71+
await expect.poll(() => counter.value().count).toBe(2)
72+
await expect.poll(() => secondCounter.value().count).toBe(2)
73+
74+
group.clients.find(client => client.$meta === first.meta)?.$close()
75+
group.updateChannels(current => current.splice(current.indexOf(first.serverChannel), 1))
76+
first.rpc.$close()
77+
await expect(first.rpc.$call('devframe:rpc:server-state:get', 'counter')).rejects.toThrow()
78+
secondCounter.mutate((draft) => {
79+
draft.count = 3
80+
})
81+
await expect.poll(() => counter.value().count).toBe(3)
82+
expect(firstCounter.value().count).toBe(2)
83+
expect(sharedState.keys()).toEqual(['counter'])
84+
expect(first.state.delete('counter')).toBe(true)
85+
expect(sharedState.delete('counter')).toBe(true)
86+
expect(sharedState.keys()).toEqual([])
87+
}
88+
finally {
89+
for (const peer of peers) peer.rpc.$close()
90+
for (const client of group.clients) client.$close()
91+
group.updateChannels(current => current.splice(0))
92+
for (const channel of channels) {
93+
channel.port1.close()
94+
channel.port2.close()
95+
}
96+
}
97+
})
98+
99+
it('rejects an initial snapshot when its RPC connection closes', async () => {
100+
expect.assertions(1)
101+
const client = new RpcFunctionsCollectorBase<DevframeRpcClientFunctions, undefined>(undefined)
102+
const channel = new MessageChannel()
103+
const rpc = createRpcClient<DevframeRpcServerFunctions, DevframeRpcClientFunctions>(client.functions, {
104+
channel: channelFor(channel.port1),
105+
})
106+
const state = createRpcSharedStateClientHost({
107+
call: rpc.$call,
108+
callEvent: rpc.$callEvent,
109+
client,
110+
isTrusted: true,
111+
events: createEventEmitter<RpcClientEvents>(),
112+
connectionMeta: { backend: 'none' },
113+
})
114+
try {
115+
const pending = state.get('counter')
116+
rpc.$close()
117+
await expect(pending).rejects.toThrow('closed')
118+
}
119+
finally {
120+
channel.port1.close()
121+
channel.port2.close()
122+
}
123+
}, 1000)

0 commit comments

Comments
 (0)