Skip to content

Commit fdc268c

Browse files
committed
fix(rpc): preserve diagnostics and test native subscriptions
1 parent 5677ecd commit fdc268c

9 files changed

Lines changed: 85 additions & 195 deletions

File tree

‎knip.jsonc‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -137,7 +137,7 @@
137137
"src/node/index.ts",
138138
"src/node/{auth,hub-internals}/index.ts",
139139
"src/recipes/interactive-auth.ts",
140-
"src/rpc/{index,client,server}.ts",
140+
"src/rpc/{index,client,server,shared-state}.ts",
141141
"src/rpc/dump/index.ts",
142142
"src/rpc/transports/{sse-client,sse-server,ws-bun,ws-deno,ws-client,ws-server}.ts",
143143
"src/types/index.ts",

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

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
import { defineDiagnostics } from 'devframe/utils/nostics'
2+
import { sharedStateNotFound } from '../rpc/diagnostics'
23

34
/**
45
* DF00xx codes are allocated across packages (e.g. @devframes/json-render
@@ -20,6 +21,7 @@ export const diagnostics = defineDiagnostics({
2021
DF0012: {
2122
why: (p: { filepath: string }) => `Failed to parse storage file: ${p.filepath}, falling back to defaults.`,
2223
},
24+
DF0013: sharedStateNotFound,
2325
DF0014: {
2426
why: (p: { name: string }) => `RPC function "${p.name}" has an invalid \`agent\` field: \`description\` must be a non-empty string.`,
2527
fix: 'Provide a short description (~1–3 sentences) explaining what the tool does and when agents should invoke it.',

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

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,11 +1,13 @@
11
import { defineDiagnostics } from 'devframe/utils/nostics'
22

3+
export const sharedStateNotFound = {
4+
why: (parameters: { key: string }) => `Shared state of "${parameters.key}" is not found, please provide an initial value for the first time`,
5+
} as const
6+
37
export const diagnostics = defineDiagnostics({
48
docsBase: 'https://devfra.me/errors',
59
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-
},
10+
DF0013: sharedStateNotFound,
911
DF0019: {
1012
why: (p: { name: string }) =>
1113
`RPC function "${p.name}" has \`agent\` set but \`jsonSerializable\` is \`false\`; MCP requires JSON-serializable data.`,

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

Lines changed: 23 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -3,13 +3,15 @@ import type { RpcClientEvents } from 'devframe/client'
33
import type { DevframeRpcClientFunctions, DevframeRpcServerFunctions, RpcFunctionsHost } from 'devframe/types'
44
import type { MessagePort } from 'node:worker_threads'
55
import { MessageChannel } from 'node:worker_threads'
6+
import { createHostContext } from 'devframe/node'
67
import { RpcFunctionsCollectorBase } from 'devframe/rpc'
78
import { createRpcClient } from 'devframe/rpc/client'
89
import { createRpcServer } from 'devframe/rpc/server'
910
import { createRpcSharedStateClientHost, createRpcSharedStateServerHost } from 'devframe/rpc/shared-state'
1011
import { createEventEmitter } from 'devframe/utils/events'
1112
import { structuredCloneDeserialize, structuredCloneSerialize } from 'devframe/utils/structured-clone'
1213
import { expect, it } from 'vitest'
14+
import { createContextRpcServer } from '../node/rpc-core'
1315

1416
/** JSON records also travel through transports that cannot clone native values. */
1517
function channelFor(port: MessagePort): ChannelOptions {
@@ -26,8 +28,21 @@ function channelFor(port: MessagePort): ChannelOptions {
2628
}
2729
}
2830

29-
it('shares native state across independent channels with per-connection subscriptions', async () => {
30-
expect.assertions(13)
31+
async function createStateServer(mode: string) {
32+
if (mode === 'node context') {
33+
const context = await createHostContext({
34+
cwd: process.cwd(),
35+
mode: 'dev',
36+
host: {
37+
mountStatic() {},
38+
resolveOrigin: () => 'http://localhost',
39+
getStorageDir: () => process.cwd(),
40+
},
41+
})
42+
const { rpcGroup } = createContextRpcServer({ context, auth: false })
43+
return { group: rpcGroup, sharedState: context.rpc.sharedState }
44+
}
45+
3146
const collector = new RpcFunctionsCollectorBase<DevframeRpcServerFunctions, undefined>(undefined)
3247
const group = createRpcServer<DevframeRpcClientFunctions, DevframeRpcServerFunctions>(collector.functions)
3348
const broadcast: RpcFunctionsHost['broadcast'] = async (options) => {
@@ -36,6 +51,12 @@ it('shares native state across independent channels with per-connection subscrip
3651
.map(client => client.$callRaw({ ...options, optional: true, event: true })))
3752
}
3853
const sharedState = createRpcSharedStateServerHost({ register: collector.register.bind(collector), broadcast })
54+
return { group, sharedState }
55+
}
56+
57+
it.each(['custom channels', 'node context'])('shares state across %s with per-connection subscriptions', async (mode) => {
58+
expect.assertions(13)
59+
const { group, sharedState } = await createStateServer(mode)
3960
const counter = await sharedState.get('counter', { initialValue: { count: 1 } })
4061
const channels = [new MessageChannel(), new MessageChannel()]
4162
const peers = channels.map((channel, index) => {

‎tests/__snapshots__/tsnapi/devframe/client.snapshot.d.ts‎

Lines changed: 30 additions & 157 deletions
Original file line numberDiff line numberDiff line change
@@ -2,137 +2,9 @@
22
* Generated by tsnapi — public API snapshot of `devframe/client`
33
*/
44
// #region Interfaces
5-
export interface DevframeConnection {
6-
connectionMeta: ConnectionMeta;
7-
metaBaseUrl: string;
8-
authToken?: string;
9-
}
10-
export interface DevframeRpcClient {
11-
events: EventEmitter<RpcClientEvents>;
12-
readonly isTrusted: boolean | null;
13-
readonly status: DevframeConnectionStatus;
14-
readonly connectionError: Error | null;
15-
readonly transport: 'websocket' | 'sse' | 'static';
16-
readonly connection: DevframeConnection;
17-
readonly connectionMeta: ConnectionMeta;
18-
ensureTrusted: (_?: number) => Promise<boolean>;
19-
requestTrust: () => Promise<boolean>;
20-
requestTrustWithToken: (_: string) => Promise<boolean>;
21-
requestTrustWithCode: (_: string) => Promise<boolean>;
22-
requestAuthCode: (_?: {
23-
reissue?: boolean;
24-
}) => Promise<void>;
25-
call: DevframeRpcClientCall;
26-
callEvent: DevframeRpcClientCallEvent;
27-
callOptional: DevframeRpcClientCallOptional;
28-
client: DevframeClientRpcHost;
29-
sharedState: RpcSharedStateHost;
30-
services: DevframeServicesClient;
31-
streaming: RpcStreamingClientHost;
32-
cacheManager: RpcCacheManager;
33-
scope: {
34-
<NS extends string>(_: NS): DevframeScopedClientContext<NS, SettingsForNamespace<NS>>;
35-
(_?: null | ''): DevframeRpcClient;
36-
};
37-
close?: () => void;
38-
}
39-
export interface DevframeRpcClientMode {
40-
readonly transport?: 'websocket' | 'sse' | 'static';
41-
readonly isTrusted: boolean;
42-
readonly status: DevframeConnectionStatus;
43-
readonly connectionError: Error | null;
44-
ensureTrusted: DevframeRpcClient['ensureTrusted'];
45-
requestTrust: DevframeRpcClient['requestTrust'];
46-
requestTrustWithToken: DevframeRpcClient['requestTrustWithToken'];
47-
requestTrustWithCode: (_: string) => Promise<string | null>;
48-
requestAuthCode: DevframeRpcClient['requestAuthCode'];
49-
call: DevframeRpcClient['call'];
50-
callEvent: DevframeRpcClient['callEvent'];
51-
callOptional: DevframeRpcClient['callOptional'];
52-
close?: () => void;
53-
}
54-
export interface DevframeRpcClientOptions extends SetupDevframeConnectionOptions {
55-
authToken?: string;
56-
otpParam?: string | false;
57-
simpleAuth?: boolean;
58-
transport?: 'auto' | 'websocket' | 'sse';
59-
wsOptions?: Partial<WsRpcChannelOptions>;
60-
sseOptions?: Partial<SseRpcChannelOptions>;
61-
rpcOptions?: Partial<BirpcOptions<DevframeRpcServerFunctions, DevframeRpcClientFunctions, boolean>>;
62-
cacheOptions?: boolean | Partial<RpcCacheOptions>;
63-
webmcp?: boolean;
64-
callTimeout?: number;
65-
}
66-
export interface DevframeRpcContext {
67-
readonly rpc: DevframeRpcClient;
68-
}
69-
export interface DevframeScopedClientContext<NS extends string = string, Settings extends Record<string, any> = Record<string, any>> {
70-
readonly namespace: NS;
71-
readonly base: DevframeRpcClient;
72-
rpc: DevframeScopedClientRpc<NS>;
73-
settings: DevframeSettings<Settings>;
74-
scope: DevframeRpcClient['scope'];
75-
}
76-
export interface DevframeScopedClientRpc<NS extends string = string> {
77-
readonly namespace: NS;
78-
register: (_: RpcFunctionDefinition<string, any, any, any, any, any, DevframeRpcContext>) => void;
79-
call: {
80-
<T extends keyof ScopedServerFunctions<NS> & string>(_: T, ..._: Parameters<ScopedRpcFn<DevframeRpcServerFunctions, NS, T>>): Promise<Awaited<ReturnType<ScopedRpcFn<DevframeRpcServerFunctions, NS, T>>>>;
81-
<T extends keyof DevframeRpcServerFunctions & string>(_: T, ..._: Parameters<Extract<DevframeRpcServerFunctions[T], AnyRpcFn>>): Promise<Awaited<ReturnType<Extract<DevframeRpcServerFunctions[T], AnyRpcFn>>>>;
82-
(_: string, ..._: any[]): Promise<any>;
83-
};
84-
callEvent: {
85-
<T extends keyof ScopedServerFunctions<NS> & string>(_: T, ..._: Parameters<ScopedRpcFn<DevframeRpcServerFunctions, NS, T>>): void;
86-
<T extends keyof DevframeRpcServerFunctions & string>(_: T, ..._: Parameters<Extract<DevframeRpcServerFunctions[T], AnyRpcFn>>): void;
87-
(_: string, ..._: any[]): void;
88-
};
89-
callOptional: {
90-
<T extends keyof ScopedServerFunctions<NS> & string>(_: T, ..._: Parameters<ScopedRpcFn<DevframeRpcServerFunctions, NS, T>>): Promise<Awaited<ReturnType<ScopedRpcFn<DevframeRpcServerFunctions, NS, T>>> | undefined>;
91-
<T extends keyof DevframeRpcServerFunctions & string>(_: T, ..._: Parameters<Extract<DevframeRpcServerFunctions[T], AnyRpcFn>>): Promise<Awaited<ReturnType<Extract<DevframeRpcServerFunctions[T], AnyRpcFn>>> | undefined>;
92-
(_: string, ..._: any[]): Promise<any>;
93-
};
94-
sharedState: {
95-
<T extends keyof ScopedSharedStates<NS> & string>(_: T, _?: RpcSharedStateGetOptions<ScopedSharedStates<NS>[T]>): Promise<SharedState<ScopedSharedStates<NS>[T]>>;
96-
<T extends Record<string, any> = Record<string, any>>(_: string, _?: RpcSharedStateGetOptions<T>): Promise<SharedState<T>>;
97-
};
98-
streaming: DevframeScopedClientStreamingHost;
99-
}
100-
export interface DevframeScopedClientStreamingHost {
101-
subscribe: <T = unknown>(_: string, _: string, _?: StreamingSubscribeOptions) => StreamReader<T>;
102-
upload: <T = unknown>(_: string, _: string) => StreamSink<T>;
103-
}
104-
export interface DevframeServiceClientHandle<NS extends string = string> extends DevframeServiceMeta {
105-
readonly scope: NS;
106-
readonly rpc: DevframeScopedClientRpc<NS>;
107-
}
108-
export interface DevframeServicesClient {
109-
has: (_: string) => boolean;
110-
get: <PKG extends string>(_: PKG) => DevframeServiceClientHandle<DevframeServiceScopeOf<PKG>> | undefined;
111-
keys: () => string[];
112-
state: () => Promise<SharedState<DevframeServicesState>>;
113-
}
1145
export interface RegisterWebMcpToolsOptions {
1156
modelContext?: WebMcpModelContext;
1167
}
117-
export interface RpcClientEvents {
118-
'rpc:is-trusted:updated': (_: boolean) => void;
119-
'connection:status': (_: DevframeConnectionStatus, _: DevframeConnectionStatus) => void;
120-
'connection:error': (_: Error) => void;
121-
'rpc:error': (_: Error, _: string) => void;
122-
}
123-
export interface RpcStreamingClientHost {
124-
subscribe: <T = unknown>(_: string, _: string, _?: StreamingSubscribeOptions) => StreamReader<T>;
125-
upload: <T = unknown>(_: string, _: string) => StreamSink<T>;
126-
}
127-
export interface SetupDevframeConnectionOptions {
128-
connection?: DevframeConnection;
129-
connectionMeta?: ConnectionMeta;
130-
baseURL?: string | string[];
131-
authToken?: string;
132-
}
133-
export interface StreamingSubscribeOptions {
134-
highWaterMark?: number;
135-
}
1368
export interface WebMcpModelContext {
1379
registerTool: (_: WebMcpToolDescriptor, _?: {
13810
signal?: AbortSignal;
@@ -178,50 +50,51 @@ export interface WsUrlLocation {
17850
}
17951
// #endregion
18052

181-
// #region Types
182-
export type DevframeClientRpcHost = RpcFunctionsCollector<DevframeRpcClientFunctions, DevframeRpcContext>;
183-
export type DevframeConnectionErrorKind = 'connection' | 'auth' | 'timeout';
184-
export type DevframeConnectionStatus = 'connecting' | 'connected' | 'unauthorized' | 'disconnected' | 'error';
185-
export type DevframeRpcClientCall = BirpcReturn<DevframeRpcServerFunctions, DevframeRpcClientFunctions>['$call'];
186-
export type DevframeRpcClientCallEvent = BirpcReturn<DevframeRpcServerFunctions, DevframeRpcClientFunctions>['$callEvent'];
187-
export type DevframeRpcClientCallOptional = BirpcReturn<DevframeRpcServerFunctions, DevframeRpcClientFunctions>['$callOptional'];
188-
// #endregion
189-
190-
// #region Classes
191-
export declare class DevframeConnectionError extends Error {
192-
name: string;
193-
readonly kind: DevframeConnectionErrorKind;
194-
constructor(_: DevframeConnectionErrorKind, _: string, _?: {
195-
cause?: unknown;
196-
});
197-
}
198-
// #endregion
199-
20053
// #region Functions
20154
export declare function authenticateWithUrlOtp(_: Pick<DevframeRpcClient, 'isTrusted' | 'requestTrustWithCode'>, _?: {
20255
param?: string;
20356
}): Promise<boolean>;
20457
export declare function consumeOtpFromUrl(_?: string): string | undefined;
20558
export declare function createClientSettings<T extends Record<string, any> = Record<string, any>>(_: DevframeRpcClient, _: string): DevframeSettings<T>;
206-
export declare function createRpcStreamingClientHost(_: DevframeRpcClient): RpcStreamingClientHost;
207-
export declare function createScopedClientContext<NS extends string = string>(_: DevframeRpcClient, _: NS): DevframeScopedClientContext<NS>;
208-
export declare function getDevframeConnection(): DevframeConnection | undefined;
209-
export declare function getDevframeRpcClient(_?: DevframeRpcClientOptions): Promise<DevframeRpcClient>;
210-
export declare function isCallableStatus(_: DevframeConnectionStatus): boolean;
21159
export declare function readOtpFromUrl(_?: string): string | undefined;
212-
export declare function registerDevframeViewerOrigin(_: DevframeConnection, _?: any): Promise<boolean>;
21360
export declare function registerWebMcpTools<LocalFunctions, SetupContext>(_: RpcFunctionsCollector<LocalFunctions, SetupContext>, _?: RegisterWebMcpToolsOptions): () => void;
214-
export declare function resolveClientTransport(_: 'auto' | 'websocket' | 'sse', _: ConnectionMeta): 'websocket' | 'sse' | 'static';
21561
export declare function resolveSseUrl(_: ConnectionMeta['sse'], _: string, _: WsUrlLocation): string;
21662
export declare function resolveWebMcpModelContext(): WebMcpModelContext | undefined;
21763
export declare function resolveWsUrl(_: ConnectionMeta['websocket'], _: string, _: WsUrlLocation): string;
218-
export declare function setupDevframeConnection(_?: SetupDevframeConnectionOptions): Promise<DevframeConnection>;
21964
// #endregion
22065

22166
// #region Variables
22267
export declare const connectDevframe: typeof getDevframeRpcClient;
22368
// #endregion
22469

225-
// #region Referenced (internal)
226-
type AnyRpcFn = (..._: any[]) => any;
70+
// #region Other
71+
export { createRpcStreamingClientHost }
72+
export { createScopedClientContext }
73+
export { DevframeClientRpcHost }
74+
export { DevframeConnection }
75+
export { DevframeConnectionError }
76+
export { DevframeConnectionErrorKind }
77+
export { DevframeConnectionStatus }
78+
export { DevframeRpcClient }
79+
export { DevframeRpcClientCall }
80+
export { DevframeRpcClientCallEvent }
81+
export { DevframeRpcClientCallOptional }
82+
export { DevframeRpcClientMode }
83+
export { DevframeRpcClientOptions }
84+
export { DevframeRpcContext }
85+
export { DevframeScopedClientContext }
86+
export { DevframeScopedClientRpc }
87+
export { DevframeScopedClientStreamingHost }
88+
export { DevframeServiceClientHandle }
89+
export { DevframeServicesClient }
90+
export { getDevframeConnection }
91+
export { getDevframeRpcClient }
92+
export { isCallableStatus }
93+
export { registerDevframeViewerOrigin }
94+
export { resolveClientTransport }
95+
export { RpcClientEvents }
96+
export { RpcStreamingClientHost }
97+
export { setupDevframeConnection }
98+
export { SetupDevframeConnectionOptions }
99+
export { StreamingSubscribeOptions }
227100
// #endregion

‎tests/__snapshots__/tsnapi/devframe/internal.snapshot.d.ts‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -80,7 +80,7 @@ export declare const diagnostics: import("nostics").Diagnostics<{
8080
}) => string;
8181
};
8282
readonly DF0013: {
83-
readonly why: (p: {
83+
readonly why: (parameters: {
8484
key: string;
8585
}) => string;
8686
};
@@ -373,7 +373,7 @@ export declare const diagnostics: import("nostics").Diagnostics<{
373373
};
374374
}, readonly [(d: import("nostics").Diagnostic, { method }?: {
375375
method?: "log" | "warn" | "error";
376-
}) => void]>;
376+
}) => void], never>;
377377
// #endregion
378378

379379
// #region Other

‎tests/__snapshots__/tsnapi/devframe/rpc.snapshot.js‎

Lines changed: 2 additions & 30 deletions
Original file line numberDiff line numberDiff line change
@@ -1,41 +1,13 @@
11
/**
22
* Generated by tsnapi — public API snapshot of `devframe/rpc`
33
*/
4-
// #region Classes
5-
export class RpcCacheManager {
6-
cacheMap
7-
options
8-
keySerializer
9-
constructor(_) {}
10-
updateOptions(_) {}
11-
cached(_, _) {}
12-
has(_, _) {}
13-
apply(_, _) {}
14-
validate(_) {}
15-
clear(_) {}
16-
}
17-
export class RpcFunctionsCollectorBase {
18-
context
19-
definitions
20-
functions
21-
_onChanged
22-
constructor(_) {}
23-
register(_, _) {}
24-
update(_, _) {}
25-
onChanged(_) {}
26-
async getHandler(_) {}
27-
getSchema(_) {}
28-
has(_) {}
29-
get(_) {}
30-
list() {}
31-
}
32-
// #endregion
33-
344
// #region Other
355
export { createDefineWrapperWithContext }
366
export { defineRpcFunction }
377
export { getRpcHandler }
388
export { getRpcResolvedSetupResult }
9+
export { RpcCacheManager }
10+
export { RpcFunctionsCollectorBase }
3911
export { strictJsonStringify }
4012
export { STRUCTURED_CLONE_PREFIX }
4113
export { validateDefinition }
Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,10 @@
1+
/**
2+
* Generated by tsnapi — public API snapshot of `devframe/rpc/shared-state`
3+
*/
4+
// #region Functions
5+
export declare function createRpcSharedStateClientHost<Context>(_: Pick<DevframeRpcClient, 'call' | 'callEvent' | 'events' | 'isTrusted'> & {
6+
client: Pick<RpcFunctionsCollector<DevframeRpcClientFunctions, Context>, 'register'>;
7+
connectionMeta: Pick<ConnectionMeta, 'backend'>;
8+
}): RpcSharedStateHost;
9+
export declare function createRpcSharedStateServerHost<Context>(_: Pick<RpcFunctionsCollector<DevframeRpcServerFunctions, Context>, 'register'> & Pick<RpcFunctionsHost, 'broadcast'>): RpcSharedStateHost;
10+
// #endregion
Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,10 @@
1+
/**
2+
* Generated by tsnapi — public API snapshot of `devframe/rpc/shared-state`
3+
*/
4+
// #region Functions
5+
export function createRpcSharedStateServerHost(_) {}
6+
// #endregion
7+
8+
// #region Other
9+
export { createRpcSharedStateClientHost }
10+
// #endregion

0 commit comments

Comments
 (0)