-
Notifications
You must be signed in to change notification settings - Fork 22
/
rpc.ts
130 lines (123 loc) · 3.6 KB
/
rpc.ts
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
import { createEventBuffer } from "./async/event-buffer.ts";
import { defer } from "./async/observer.ts";
export type Method<
TMetadata = any,
THeader = any,
TTrailer = any,
TServiceName extends string = string,
TMethodName extends string = string,
TRequestStream extends boolean = boolean,
TResponseStream extends boolean = boolean,
TReq = any,
TRes = any,
> = [
MethodDescriptor<
TReq,
TRes,
TMethodName,
TServiceName,
TRequestStream,
TResponseStream
>,
MethodImpl<TReq, TRes, TMetadata, THeader, TTrailer>,
];
export type RpcClientImpl<TMetadata = any, THeader = any, TTrailer = any> = <
TReq,
TRes,
>(
methodDescriptor: MethodDescriptor<TReq, TRes>,
) => MethodImpl<TReq, TRes, TMetadata, THeader, TTrailer>;
export interface MethodDescriptor<
TReq,
TRes,
TMethodName extends string = string,
TServiceName extends string = string,
TRequestStream extends boolean = boolean,
TResponseStream extends boolean = boolean,
> {
methodName: TMethodName;
service: { serviceName: TServiceName };
requestStream: TRequestStream;
responseStream: TResponseStream;
requestType: {
serializeBinary: (value: TReq) => Uint8Array;
deserializeBinary: (value: Uint8Array) => TReq;
serializeJson: (value: TReq) => string;
};
responseType: {
serializeBinary: (value: TRes) => Uint8Array;
deserializeBinary: (value: Uint8Array) => TRes;
serializeJson: (value: TRes) => string;
};
}
type ThenArg<T> = T extends Promise<infer U> ? U : T;
export type RpcReturnType<TRes, TResArgs extends any[]> = (
Promise<TResArgs extends [] ? ThenArg<TRes> : [ThenArg<TRes>, ...TResArgs]>
);
export interface MethodImpl<
TReq,
TRes,
TMetadata = any,
THeader = any,
TTrailer = any,
> {
(
req: AsyncGenerator<TReq>,
metadata?: TMetadata,
): [AsyncGenerator<TRes>, Promise<THeader>, Promise<TTrailer>];
}
export interface MethodImplHandlerReq<TReq, TMetadata> {
metadata?: TMetadata;
messages: AsyncGenerator<TReq>;
drainEnd: Promise<void>;
}
export interface MethodImplHandlerRes<TRes, THeader, TTrailer> {
header(value: THeader): void;
send(value: TRes): void;
end(value: TTrailer): void;
}
export interface MethodImplHandler<TReq, TRes, TMetadata, THeader, TTrailer> {
(
req: MethodImplHandlerReq<TReq, TMetadata>,
res: MethodImplHandlerRes<TRes, THeader, TTrailer>,
): void;
}
export function getMethodImpl<
TReq,
TRes,
TMetadata,
THeader,
TTrailer,
>(
handler: MethodImplHandler<TReq, TRes, TMetadata, THeader, TTrailer>,
): MethodImpl<TReq, TRes, TMetadata, THeader, TTrailer> {
return (messages: AsyncGenerator<TReq>, metadata?: TMetadata) => {
const headerPromise = defer<THeader>();
const trailerPromise = defer<TTrailer>();
const drainEnd = defer<void>();
const eventBuffer = createEventBuffer<TRes>({
onDrainEnd: drainEnd.resolve,
});
const header = headerPromise.resolve;
const send = eventBuffer.push;
const end = (value: TTrailer) => {
eventBuffer.finish();
trailerPromise.resolve(value);
};
handler({ metadata, messages, drainEnd }, { header, send, end });
return [eventBuffer.drain(), headerPromise, trailerPromise];
};
}
export function createServerImplBuilder<TMetadata, THeader, TTrailer>() {
const buffer = createEventBuffer<Method<TMetadata, THeader, TTrailer>>();
return {
register<TReq, TRes>(
methodDescriptor: MethodDescriptor<TReq, TRes>,
handler: MethodImplHandler<TReq, TRes, TMetadata, THeader, TTrailer>,
) {
buffer.push([methodDescriptor, getMethodImpl(handler)]);
},
finish: buffer.finish,
drain: buffer.drain,
};
}