Skip to content

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

func NewHttpClient(baseURL string, options ...HttpClientOption) (*HttpClient, error)
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

type ArrowSerializable interface {
    ArrowSchema() *arrow.Schema
}

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

type KV struct {
    Key   string
    Value string
}

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

type HookToken interface{}

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)