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
37 changes: 37 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,42 @@ function UserProfile({ userId }) {

The query keys follow the pattern `[ServiceName, methodName, request?]` and are fully type-safe with `as const` assertions.

### Streaming client style

For RIDL methods declared with `=> stream (...)`, the default generated TypeScript client uses the existing callback API:

```typescript
client.subscribeMessages(req, {
onMessage(message) {},
onError(error, reconnect) {},
})
```

Pass `-streamClient=asyncIterable` to generate stream methods that return `AsyncIterable` instead:

```typescript
const stream = client.subscribeMessages(req, { signal })

for await (const message of stream) {
console.log(message)
}
```

This keeps the generated WebRPC client independent of any UI/cache library while making it easy to use with TanStack Query's `streamedQuery` helper:

```typescript
import { experimental_streamedQuery as streamedQuery } from '@tanstack/react-query'

useQuery({
queryKey: client.queryKey.subscribeMessages(req),
queryFn: streamedQuery({
streamFn: ({ signal }) => client.subscribeMessages(req, { signal }),
initialValue: [],
reducer: (messages, chunk) => [...messages, chunk.message],
}),
})
```

### Enum style: `enum` vs `union`

TypeScript best practices are moving away from `enum` declarations — they emit runtime
Expand Down Expand Up @@ -142,6 +178,7 @@ Change any of the following values by passing `-option="Value"` CLI flag to `web
| `-webrpcHeader` | send Webrpc header in all HTTP requests | `true` | v0.15.0 |
| `-schemaHash=false` | don't emit schema hash + version consts | `true` | v0.28.0 |
| `-enumStyle` | enum codegen style: `enum` or `union` | `enum` | v0.29.0 |
| `-streamClient` | streaming client style: `callback` or `asyncIterable` | `callback` | next |

**Note:** Generated code requires ES2022+ runtime environment.

Expand Down
11 changes: 9 additions & 2 deletions client.go.tmpl
Original file line number Diff line number Diff line change
Expand Up @@ -56,8 +56,14 @@ export class {{$service.Name}} implements {{$service.Name}}Client {
{{- $methodRespName = printf "%sResponse" $method.Name -}}
{{- end -}}
{{- end}}
{{firstLetterToLower .Name}} = ({{template "methodInputs" dict "Method" . "Opts" $opts "TypeMap" $typeMap}}): {{if $method.StreamOutput}}WebrpcStreamController{{else}}{{if $method.Succinct}}Promise<{{(index $method.Outputs 0).Type}}>{{else}}Promise<{{$method.Name}}{{if $opts.compat}}Return{{else}}Response{{end}}>{{end}}{{end}} => {
{{firstLetterToLower .Name}} = ({{template "methodInputs" dict "Method" . "Opts" $opts "TypeMap" $typeMap}}): {{if $method.StreamOutput}}{{if $opts.streamClientAsyncIterable}}WebrpcStream<{{$methodRespName}}>{{else}}WebrpcStreamController{{end}}{{else}}{{if $method.Succinct}}Promise<{{(index $method.Outputs 0).Type}}>{{else}}Promise<{{$method.Name}}{{if $opts.compat}}Return{{else}}Response{{end}}>{{end}}{{end}} => {
{{- if $method.StreamOutput }}
{{- if $opts.streamClientAsyncIterable }}
return webrpcAsyncIterable<{{$methodRespName}}>(() => this.fetch(this.url('{{.Name}}'),
{{if .Inputs | len }}createHttpRequest(JsonEncode(req), options?.headers, options?.signal){{- else}}createHttpRequest('{}', options?.headers, options?.signal){{end }}
), '{{$methodRespName}}', options?.signal)
}
{{- else }}
const abortController = new AbortController()
const abortSignal = abortController.signal

Expand All @@ -81,6 +87,7 @@ export class {{$service.Name}} implements {{$service.Name}}Client {
closed: resp
}
}
{{- end }}
{{- else }}
return this.fetch(
this.url('{{.Name}}'),
Expand All @@ -102,6 +109,6 @@ export class {{$service.Name}} implements {{$service.Name}}Client {
{{- end -}}
{{end -}}
{{if $opts.streaming}}
{{template "sse"}}
{{template "sse" dict "Opts" $opts}}
{{end}}
{{- end -}}
9 changes: 9 additions & 0 deletions clientHelpers.go.tmpl
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,14 @@ const buildResponse = (res: Response): Promise<any> => {
export type Fetch = (input: RequestInfo, init?: RequestInit) => Promise<Response>

{{if $opts.streaming}}
{{- if $opts.streamClientAsyncIterable }}
export interface WebrpcOptions {
headers?: HeadersInit;
signal?: AbortSignal;
}

export type WebrpcStream<T> = AsyncIterable<T>
{{- else }}
export interface WebrpcStreamOptions<T> extends WebrpcOptions {
onMessage: (message: T) => void;
onError: (error: WebrpcError, reconnect: () => void) => void;
Expand All @@ -47,5 +55,6 @@ export interface WebrpcStreamController {
abort: (reason?: any) => void;
closed: Promise<void>;
}
{{- end }}
{{end}}
{{end}}
2 changes: 1 addition & 1 deletion clientInterface.go.tmpl
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,7 @@ export interface {{$service.Name}}Client {
* @deprecated {{ $deprecated.Value }}
*/
{{- end }}
{{firstLetterToLower $method.Name}}({{template "methodInputs" dict "Method" $method "TypeMap" $typeMap "Opts" $opts}}): {{if $method.StreamOutput}}WebrpcStreamController{{else}}{{if $method.Succinct}}Promise<{{(index $method.Outputs 0).Type}}>{{else}}Promise<{{$method.Name}}{{if $opts.compat}}Return{{else}}Response{{end}}>{{end}}{{end}}
{{firstLetterToLower $method.Name}}({{template "methodInputs" dict "Method" $method "TypeMap" $typeMap "Opts" $opts}}): {{if $method.StreamOutput}}{{if $opts.streamClientAsyncIterable}}WebrpcStream<{{if $method.Succinct}}{{(index $method.Outputs 0).Type}}{{else}}{{$method.Name}}{{if $opts.compat}}Return{{else}}Response{{end}}{{end}}>{{else}}WebrpcStreamController{{end}}{{else}}{{if $method.Succinct}}Promise<{{(index $method.Outputs 0).Type}}>{{else}}Promise<{{$method.Name}}{{if $opts.compat}}Return{{else}}Response{{end}}>{{end}}{{end}}
{{- if lt (add $i 1) (len $service.Methods)}}{{"\n"}}{{end}}
{{- end}}
}
Expand Down
168 changes: 168 additions & 0 deletions clientSSE.go.tmpl
Original file line number Diff line number Diff line change
@@ -1,4 +1,171 @@
{{ define "sse" }}
{{- $opts := .Opts -}}
{{- if $opts.streamClientAsyncIterable }}
const WebrpcStreamReadTimeoutMs = (10 + 1) * 1000

const webrpcAsyncIterable = <T>(
fetchResponse: () => Promise<Response>,
responseType: string,
signal?: AbortSignal
): WebrpcStream<T> => ({
async *[Symbol.asyncIterator]() {
let res: Response
try {
res = await fetchResponse()
} catch (error) {
if (signal?.aborted) {
return
}
throw WebrpcRequestFailedError.new({ cause: `fetch(): ${error instanceof Error ? error.message : String(error)}` })
}

yield* readWebrpcStream<T>(res, responseType, signal)
}
})

const readWebrpcStream = async function* <T>(
res: Response,
responseType: string,
signal?: AbortSignal
): AsyncGenerator<T> {
if (!res.ok) {
await buildResponse(res)
return
}

if (!res.body) {
throw WebrpcBadResponseError.new({
status: res.status,
cause: "Invalid response, missing body",
})
}

const reader = res.body.getReader()
const decoder = new TextDecoder()
let buffer = ""

try {
while (true) {
const { value, done } = await readWebrpcStreamChunk(reader, signal)

if (done) {
buffer += decoder.decode()
if (buffer.trim().length > 0) {
yield decodeWebrpcStreamLine<T>(buffer, responseType, res.status)
}
return
}

buffer += decoder.decode(value, { stream: true })

const lines = buffer.split("\n")
buffer = lines.pop() || ""

for (const line of lines) {
if (line.trim().length === 0) {
continue
}
yield decodeWebrpcStreamLine<T>(line, responseType, res.status)
}
}
} catch (error) {
if (error instanceof WebrpcError) {
try {
await reader.cancel(error)
} catch (_) {
// Ignore cleanup errors and report the protocol error below.
}
throw error
}

if (
signal?.aborted ||
(error instanceof DOMException && error.name === "AbortError")
) {
return
}

try {
await reader.cancel(error)
} catch (_) {
// Ignore cleanup errors and report the original stream error below.
}

throw WebrpcStreamLostError.new({
cause: `reader.read(): ${error instanceof Error ? error.message : String(error)}`,
})
} finally {
try {
reader.releaseLock()
} catch (_) {
// Ignore releaseLock errors; the stream is already finishing.
}
}
}

const readWebrpcStreamChunk = (
reader: ReadableStreamDefaultReader<Uint8Array>,
signal?: AbortSignal
): Promise<ReadableStreamReadResult<Uint8Array>> => {
return new Promise((resolve, reject) => {
const timeoutId = setTimeout(() => {
cleanup()
reject(new Error("Timeout, no data or heartbeat received"))
}, WebrpcStreamReadTimeoutMs)

const abort = () => {
cleanup()
reject(new DOMException("AbortError", "AbortError"))
}

const cleanup = () => {
clearTimeout(timeoutId)
signal?.removeEventListener("abort", abort)
}

if (signal?.aborted) {
abort()
return
}

signal?.addEventListener("abort", abort, { once: true })

reader.read().then(
(result) => {
cleanup()
resolve(result)
},
(error) => {
cleanup()
reject(error)
}
)
})
}

const decodeWebrpcStreamLine = <T>(line: string, responseType: string, status: number): T => {
let data: any
try {
data = JsonDecode<any>(line, responseType)
} catch (error) {
if (error instanceof WebrpcError) {
throw error
}
throw WebrpcBadResponseError.new({
status,
cause: `JsonDecode(): ${error instanceof Error ? error.message : String(error)}`,
})
}

if (data && Object.prototype.hasOwnProperty.call(data, "webrpcError")) {
const error = data.webrpcError
const code: number = typeof error.code === "number" ? error.code : 0
throw (webrpcErrorByCode[code] || WebrpcError).new(error)
}

return data as T
}
{{- else }}
const sseResponse = async (
res: Response,
options: WebrpcStreamOptions<any>,
Expand Down Expand Up @@ -120,4 +287,5 @@ const sseResponse = async (
return;
}
};
{{- end }}
{{ end }}
9 changes: 9 additions & 0 deletions main.go.tmpl
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
{{- set $opts "webrpcHeader" (ternary (eq (default .Opts.webrpcHeader "true") "false") false true) -}}
{{- set $opts "schemaHash" (ternary (eq (default .Opts.schemaHash "true") "false") false true) -}}
{{- set $opts "enumStyle" (default .Opts.enumStyle "enum") -}}
{{- set $opts "streamClient" (default .Opts.streamClient "callback") -}}

{{- /* Print help on -help. */ -}}
{{- if exists .Opts "help" -}}
Expand All @@ -29,6 +30,14 @@
{{- exit 1 -}}
{{- end -}}

{{- if not (in $opts.streamClient "callback" "asyncIterable") -}}
{{- stderrPrintf "-streamClient=%q is not supported, must be \"callback\" or \"asyncIterable\"\n" $opts.streamClient -}}
{{- exit 1 -}}
{{- end -}}

{{- set $opts "streamClientCallback" (eq $opts.streamClient "callback") -}}
{{- set $opts "streamClientAsyncIterable" (eq $opts.streamClient "asyncIterable") -}}

{{- if ne .WebrpcVersion "v1" -}}
{{- stderrPrintf "%s generator error: unsupported Webrpc version %s\n" .WebrpcTarget .WebrpcVersion -}}
{{- exit 1 -}}
Expand Down
12 changes: 12 additions & 0 deletions methodInputs.go.tmpl
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,11 @@

{{- if gt (len $method.Inputs) 0}}req: {{(index $method.Inputs 0).Type}}, {{end}}
{{- if $method.StreamOutput -}}
{{- if $opts.streamClientAsyncIterable -}}
options?: WebrpcOptions
{{- else -}}
options: WebrpcStreamOptions<{{(index $method.Outputs 0).Type}}>
{{- end -}}
{{- else -}}
headers?: object, signal?: AbortSignal
{{- end -}}
Expand All @@ -17,7 +21,11 @@

{{- if gt (len $method.Inputs) 0}}req: {{$method.Name}}Args, {{end}}
{{- if $method.StreamOutput -}}
{{- if $opts.streamClientAsyncIterable -}}
options?: WebrpcOptions
{{- else -}}
options: WebrpcStreamOptions<{{$method.Name}}Return>
{{- end -}}
{{- else -}}
headers?: object, signal?: AbortSignal
{{- end -}}
Expand All @@ -26,7 +34,11 @@

{{- if gt (len $method.Inputs) 0}}req: {{$method.Name}}Request, {{end}}
{{- if $method.StreamOutput -}}
{{- if $opts.streamClientAsyncIterable -}}
options?: WebrpcOptions
{{- else -}}
options: WebrpcStreamOptions<{{$method.Name}}Response>
{{- end -}}
{{- else -}}
headers?: object, signal?: AbortSignal
{{- end -}}
Expand Down
3 changes: 2 additions & 1 deletion tests-unit/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,8 @@
"test:watch": "vitest",
"test:update": "vitest run -u",
"generate": "webrpc-gen -schema=service.ridl -target=../ -out=client.gen.ts",
"generate:enum": "webrpc-gen -schema=enum.ridl -target=../ -client -out=enum-default.gen.ts && webrpc-gen -schema=enum.ridl -target=../ -client -enumStyle=union -out=enum-union.gen.ts"
"generate:enum": "webrpc-gen -schema=enum.ridl -target=../ -client -out=enum-default.gen.ts && webrpc-gen -schema=enum.ridl -target=../ -client -enumStyle=union -out=enum-union.gen.ts",
"generate:stream": "webrpc-gen -schema=stream.ridl -target=../ -client -streamClient=asyncIterable -out=stream-async.gen.ts"
},
"devDependencies": {
"@tsconfig/strictest": "^2.0.7",
Expand Down
Loading
Loading