API Reference¶
Full Go documentation is available on pkg.go.dev. For the protocol specification, see the wire protocol page.
Registration Functions¶
These generic functions register RPC methods on a Server:
| Function | Description |
|---|---|
Unary[P, R](s, name, handler) |
Register a unary method returning a result |
UnaryVoid[P](s, name, handler) |
Register a unary method with no result |
Producer[P](s, name, outputSchema, handler) |
Register a producer stream |
ProducerWithHeader[P](s, name, outputSchema, headerSchema, handler) |
Register a producer with a header |
Exchange[P](s, name, outputSchema, inputSchema, handler) |
Register an exchange stream |
ExchangeWithHeader[P](s, name, outputSchema, inputSchema, headerSchema, handler) |
Register an exchange with a header |
DynamicStreamWithHeader[P](s, name, headerSchema, handler) |
Register a stream where producer/exchange mode is determined at runtime |
Handler signatures:
// Unary
func(ctx context.Context, callCtx *CallContext, params P) (R, error)
// UnaryVoid
func(ctx context.Context, callCtx *CallContext, params P) error
// Producer / Exchange
func(ctx context.Context, callCtx *CallContext, params P) (*StreamResult, error)
Server¶
func NewServer() *Server
func NewProtocol(name string) *Server // an additional protocol, for Server.AddProtocol
| Method | Description |
|---|---|
SetServerID(id string) |
Set the server identifier included in response metadata |
SetServiceName(name string) |
Set a logical service name used by observability hooks |
ServiceName() string |
Returns the logical service name |
SetDispatchHook(hook DispatchHook) |
Register a hook called around each RPC dispatch |
AddProtocol(protocol *Server) error |
Host an additional application protocol (built with NewProtocol) on every transport; call before serving; refused once the server has served |
SetIncludeTracebacks(enabled bool) |
Tracebacks are on by default on every transport; false turns them off server-wide |
RunStdio() |
Run the server loop on stdin/stdout |
Serve(r io.Reader, w io.Writer) |
Run the server on any reader/writer pair |
ServeWithContext(ctx context.Context, r io.Reader, w io.Writer) |
Run the server with a context for cancellation |
HttpServer¶
func NewHttpServer(server *Server) *HttpServer
func NewHttpServerWithKey(server *Server, signingKey []byte) (*HttpServer, error)
func RegisterStateType(v interface{})
| Method | Description |
|---|---|
SetTokenTTL(d time.Duration) |
Set state token maximum age |
ServeHTTP(w http.ResponseWriter, r *http.Request) |
Implements http.Handler |
HttpClient¶
| Method | Description |
|---|---|
CallUnary(ctx, method, params, outputSchema) |
Invoke a unary RPC and return an owned ClientBatch |
OpenProducer(ctx, method, params, schema) |
Initialize a producer stream |
OpenExchange(ctx, method, params, schema) |
Initialize a typed exchange stream |
ListProtocols(ctx) |
vgirpc.ListProtocols on this client's connection |
DescribeProtocol(ctx, name) |
vgirpc.DescribeProtocol on this client's connection |
Describe(ctx) |
Describe the first application protocol the server hosts |
Close() |
Close client-owned idle HTTP connections; local and idempotent |
HttpClientStream.Next receives the next producer batch,
HttpClientStream.Exchange sends exactly one declared-schema batch, and
HttpClientStream.Cancel explicitly asks the server to cancel. Returned
ClientBatch values and stream headers are caller-owned and must be released.
HttpClientStream.Close only releases local state.
WithClientProtocol("my.Service.v1") is required: the routing key rides
both as vgi_rpc.protocol and as the path's protocol segment, so a client that
names no protocol is refused at construction rather than on arrival.
Client options otherwise configure the underlying net/http client, URL prefix,
request headers, protocol version, request and response size limits, and the
client-directed log handler.
WithClientTCPProxy("socks5h://127.0.0.1:1055") selects an explicit SOCKS5h
transport for Tailscale userspace networking. It uses proxy-side target-name
resolution, supports domain/IPv4/IPv6 targets, accepts only NO AUTH, and
never falls back to direct TCP.
Reflection client¶
func ListProtocols(ctx context.Context, target ReflectionTarget) ([]HostedProtocol, error)
func DescribeProtocol(ctx context.Context, target ReflectionTarget, name string) (*ClientServiceDescription, error)
ReflectionTarget is any *HttpClient or *TcpClient (TCP, Unix, raw Iroh,
HTTP-over-Iroh), bound to any protocol; its connection is reused and never
closed. HostedProtocol{Name, Version, Hash, Deprecated, DeprecationMessage,
Features} comes back in the server's order. A server that does not host
vgi_rpc.Reflection.v1 (reflection is opt-in via RegisterReflection) returns
*ReflectionNotSupportedError, which embeds the server's *RpcError; an
unknown protocol name is a plain *RpcError with kind
protocol_not_supported. See Introspection.
Stream Interfaces¶
ProducerState¶
type ProducerState interface {
Produce(ctx context.Context, out *OutputCollector, callCtx *CallContext) error
}
ExchangeState¶
type ExchangeState interface {
Exchange(ctx context.Context, input arrow.RecordBatch, out *OutputCollector, callCtx *CallContext) error
}
StreamResult¶
Returned by producer/exchange init handlers:
type StreamResult struct {
OutputSchema *arrow.Schema
State interface{} // ProducerState or ExchangeState
InputSchema *arrow.Schema // exchange only; nil for producers
Header ArrowSerializable // optional header sent before data
}
OutputCollector¶
| Method | Description |
|---|---|
Emit(batch arrow.RecordBatch) error |
Emit a pre-built RecordBatch |
EmitArrays(arrays []arrow.Array, numRows int64) error |
Build and emit a batch from arrays |
EmitMap(data map[string][]interface{}) error |
Build and emit a batch from column maps |
Finish() error |
Signal end-of-stream (producer only) |
Finished() bool |
Whether Finish() has been called |
ClientLog(level LogLevel, message string, extras ...KV) |
Emit a log batch to the client |
ArrowSerializable¶
CallContext¶
type CallContext struct {
Ctx context.Context
RequestID string
ServerID string
Method string
LogLevel LogLevel
}
| Method | Description |
|---|---|
ClientLog(level LogLevel, msg string, extras ...KV) |
Record a log message for the client |
RespondWithExternalRef(ref ExternalRef) error |
Answer this unary call with a pre-published ref instead of the returned value (see External Storage) |
External Storage¶
type ExternalStorage interface {
Upload(data []byte, schema *arrow.Schema, contentEncoding string) (string, error)
}
server.SetExternalLocation(vgirpc.DefaultExternalLocationConfig(storage))
Pre-published results¶
type ExternalRef struct { /* unexported: url, sha256 */ }
func NewExternalRef(url, sha256Hex string) (ExternalRef, error)
func PublishExternal(batch arrow.RecordBatch, storage ExternalStorage,
compression *Compression, includeSHA256 bool) (ExternalRef, error)
func PublishExternalResult[R any](value R, storage ExternalStorage,
compression *Compression, includeSHA256 bool) (ExternalRef, error)
| Function / method | Description |
|---|---|
NewExternalRef(url, sha256Hex) |
Validated ref; url non-empty, sha256Hex 64 lowercase hex or "" (no digest: clients skip the content check) |
ExternalRef.URL() / SHA256() |
The ref's URL and digest ("" when none) |
ExternalRef.PointerBatch(schema) |
The zero-row pointer batch + metadata the dispatcher writes |
PublishExternal(batch, ...) |
Serialize a 1-row result batch exactly as the per-call externalizer does, hash the raw bytes, compress, upload once, return the ref |
PublishExternalResult(value, ...) |
PublishExternal for a Go value, building the result batch as Unary would for result type R |
CallContext.RespondWithExternalRef(ref) |
Answer the current unary call with the ref's pointer batch: no build, serialization or upload; never inlined; not counted toward the externalized-response cap |
RpcError¶
type RpcError struct {
Type string
Message string
Traceback string
RequestID string
Kind string // vgi_rpc.error_kind, "" when absent
Code string // vgi_rpc.error_code, "" when the server sent none
Details []map[string]any // vgi_rpc.error_details as received
}
| Method | Description |
|---|---|
Error() string |
Returns error string |
Is(target error) bool |
Supports errors.Is |
ErrorCode() Code |
The code; CodeUnknown when absent or unrecognised |
IsRetryable() bool |
UNAVAILABLE, or RESOURCE_EXHAUSTED with RetryInfo |
TypedDetails() []ErrorDetail |
Catalog details, unknown types skipped |
RetryInfo(), ErrorInfo(), BadRequest(), PreconditionFailure(), QuotaFailure(), ResourceInfo(), Help(), LocalizedMessage() |
One catalog detail and whether it was present |
Servers raise *StatusError{Code, Kind, Message, Details}; see the
error handling guide.
Sentinel: ErrRpc — use with errors.Is(err, vgirpc.ErrRpc)
Request¶
type Request struct {
Method string
Version string
RequestID string
LogLevel string
Batch arrow.RecordBatch
Metadata map[string]string
}
Logging¶
LogLevel¶
type LogLevel string
const (
LogException LogLevel = "EXCEPTION"
LogError LogLevel = "ERROR"
LogWarn LogLevel = "WARN"
LogInfo LogLevel = "INFO"
LogDebug LogLevel = "DEBUG"
LogTrace LogLevel = "TRACE"
)
KV¶
Dispatch Hook¶
type DispatchHook interface {
OnDispatchStart(ctx context.Context, info DispatchInfo) (context.Context, HookToken)
OnDispatchEnd(ctx context.Context, token HookToken, info DispatchInfo, stats *CallStatistics, err error)
}
DispatchInfo¶
type DispatchInfo struct {
Method string // RPC method name
MethodType string // "unary" or "stream"
ServerID string
RequestID string
TransportMetadata map[string]string // IPC custom metadata or HTTP headers
}
CallStatistics¶
type CallStatistics struct {
InputBatches int64
OutputBatches int64
InputRows int64
OutputRows int64
InputBytes int64
OutputBytes int64
}
| Method | Description |
|---|---|
RecordInput(numRows, bufferBytes int64) |
Record one input batch |
RecordOutput(numRows, bufferBytes int64) |
Record one output batch |
HookToken¶
Opaque value returned by OnDispatchStart and passed to OnDispatchEnd.
Method Types¶
type MethodType int
const (
MethodUnary MethodType = iota
MethodProducer
MethodExchange
MethodDynamic
)
Batch Kinds¶
type BatchKind int
const (
BatchData BatchKind = iota
BatchLog
BatchError
BatchExternalPointer
BatchShmPointer
BatchStateToken
)
Metadata Keys¶
| Constant | Value |
|---|---|
MetaMethod |
vgi_rpc.method |
MetaRequestVersion |
vgi_rpc.request_version |
MetaRequestID |
vgi_rpc.request_id |
MetaLogLevel |
vgi_rpc.log_level |
MetaLogMessage |
vgi_rpc.log_message |
MetaLogExtra |
vgi_rpc.log_extra |
MetaServerID |
vgi_rpc.server_id |
MetaErrorCode |
vgi_rpc.error_code |
MetaErrorKind |
vgi_rpc.error_kind |
MetaErrorDetails |
vgi_rpc.error_details |
MetaStreamState |
vgi_rpc.stream_state#b64 |
MetaShmOffset |
vgi_rpc.shm_offset |
MetaShmLength |
vgi_rpc.shm_length |
MetaLocation |
vgi_rpc.location |
MetaTraceparent |
traceparent |
MetaTracestate |
tracestate |
MetaProtocolName |
vgi_rpc.protocol_name |
MetaDescribeVersion |
vgi_rpc.describe_version |
ProtocolVersion |
"1" |
DescribeVersion |
"4" |
Wire Functions¶
| Function | Description |
|---|---|
ReadRequest(r io.Reader) (*Request, error) |
Read one IPC stream and parse the request |
WriteRequest(w, method, params, protocolVersion) |
Frame a request for forwarding (intermediaries) |
WriteUnaryResponse(w, schema, logs, result, serverID, requestID) |
Write a unary response |
WriteErrorResponse(w, schema, err, serverID, requestID) |
Write an error response (nil schema = empty) |
WriteVoidResponse(w, logs, serverID, requestID) |
Write a void response |
FindStateToken(data) []byte |
Extract the stream-state cursor token (intermediaries) |
FindCallStateToken(data) []byte |
Extract the stream's call token, if the server split its state |
FindStreamTokens(data) (state, callState []byte) |
Extract both continuation tokens in one walk — what forwarding a continuation needs |
FindProtocolVersion(data) string |
Recover the stamped application protocol_version (intermediaries) |
ReadUnaryResult(data) (schema, result, ok) |
Unwrap a unary response envelope without a typed decode |
WriteUnaryResult(w, envelopeSchema, resultBytes) |
Rewrap raw result bytes as a unary response |
DecodeContentEncoding(data, contentEncoding, maxOutputSize) |
Decode an HTTP body per its Content-Encoding (zstd/gzip) |