gorpc

package module
v1.0.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 3, 2026 License: MIT Imports: 22 Imported by: 0

README

GoRPC

Go Reference CI

Small Go-to-Go RPC for internal services. Share ordinary Go types, register functions, and use one long-lived connection for calls, notifications, and streams. Either end can initiate work.

No schema files, generated stubs, or separate DTO models. GoRPC uses length-prefixed MessagePack frames, with optional gzip compression.

Version note: v1.0.0 is the first stable release, with negotiated stream flow control and hardening changes validated in v1.0.0-rc.4. Read the documentation at your dependency's tag.

Why GoRPC?

GoRPC fits systems where you control both ends and both are written in Go. An agent can connect to a controller, publish status, accept commands, and stream results over the same connection. The dialing side does not need a separate listener for callbacks.

Choose it when you want ordinary Go contracts and independent calls in either direction, without a schema compiler or a broader microservices framework. Prefer another tool when you need browser clients, cross-language APIs, or built-in service discovery and load balancing.

See Choosing GoRPC for use cases, design tradeoffs, and a comparison with other RPC libraries.

Start here

This checkout requires Go 1.26.6 or newer. To try it, run these commands from the repository root.

Start the inventory server:

go run ./examples/inventory/server

In another terminal:

go run ./examples/inventory/client

The client makes a normal call, receives a server notification, makes an async call, and handles a deliberate not_found error. Stop the server with Ctrl-C.

The inventory walkthrough has expected output, shared types, address options, and optional authentication. For your own module, install the version you intend to use explicitly, for example:

go get github.com/dan-sherwin/gorpc@v1.0.0

The API in brief

Both applications import the same request and response types:

type GetItemRequest struct {
    ID string
}

type GetItemResponse struct {
    ID   string
    Name string
}

Register a handler on a server (or on a client that accepts callbacks):

gorpc.MustRegister(server, "inventory.get", func(ctx *gorpc.Context, req GetItemRequest) (GetItemResponse, error) {
    return GetItemResponse{ID: req.ID, Name: "Widget Pack"}, nil
})

Call it over a connected client, accepted connection, or managed peer:

ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()

var item GetItemResponse
if err := client.CallContext(ctx, "inventory.get", GetItemRequest{ID: "widget-001"}, &item); err != nil {
    return err
}

These are excerpts; the runnable examples include imports, connection setup, error handling, and shutdown. The function string is a wire dispatch name, not a required Go function name. See calls and notifications.

Runnable examples

Example What it demonstrates
Inventory Shared types, sync and async calls, server push, remote errors, shutdown
Server streaming A slow reader, unrelated calls on the same connection, early cancellation
Client streaming Bounded chunks, a final response, byte counts and SHA-256 verification
Bidirectional streaming Concurrent sending and receiving, half-close, a final summary
Peers and reconnect Calls in both directions, authentication, connection replacement without replay

The last four each run with one command and choose their own loopback ports. See the example index for commands. CI runs the programs and checks their documented output, in addition to testing the package.

What GoRPC handles

  • Unary calls, async callbacks, one-way notifications, and server broadcasts.
  • Server, client, and bidirectional streaming with typed helpers.
  • Automatic client reconnect, ping/pong monitoring, and optional peer arbitration.
  • Independent per-stream item and byte receive windows, negotiated between peers.
  • Deadline propagation, cancellation, structured remote errors, and bounded shutdown.
  • Optional shared-secret authentication, compression, interceptors, and admission limits.
  • TCP, Unix sockets, and existing net.Listener implementations.

“Client” and “server” describe who dials and who accepts. Once connected, both can register handlers and initiate requests or streams.

Know the boundaries

Authentication is not encryption. Shared-secret HMAC authenticates the dialer to the accepting peer. It does not authenticate the accepting server, encrypt traffic, or authorize individual functions. The dialing helpers do not configure TLS. Use a trusted network or an externally secured tunnel and apply application authorization where needed.

Reconnect is not retry. New work can use a replacement connection. In-flight calls and streams fail with ErrUnavailable and are not replayed. The remote side may already have acted; safe retries need application-level idempotency or resume checkpoints.

Flow control is not a processing acknowledgment. It bounds queued encoded payloads. A successful Send does not mean the receiver committed the item. Notifications likewise report local write success, not remote completion.

Receive windows default to 16 items and 64 MiB per stream. Frame sizes are limited too. Connection counts, inbound handler concurrency, and decoded application objects still need application-level limits.

GoRPC is not a cross-language gRPC replacement. Service discovery, pub/sub, load balancing, and generated code are outside its scope.

Documentation

The hosted API reference follows published versions. Select the tag that matches your dependency.

Development

go build ./...
go vet ./...
go test -race -tags=integration ./... -count=1
go test -run '^$' -fuzz=FuzzReadFrame -fuzztime=2000000x -parallel=4 -timeout=5m
golangci-lint run --build-tags integration
govulncheck ./...

CI also checks module tidiness. The integration suite builds released peers and the runnable examples, so it needs the Go toolchain and may need module-proxy access. See testing for narrower checks.

Versioning and license

Semantic Versioning; see the changelog for release history. Public API compatibility is preserved within v1. Breaking API changes require a new major version. MIT licensed; see LICENSE.

Documentation

Overview

Package gorpc provides a small Go-to-Go RPC transport for internal services.

It is intentionally not a protobuf, gRPC, Connect, or IDL replacement. Both sides share normal Go request and response types, and the wire protocol uses length-prefixed MessagePack frames over a single full-duplex connection. Once connected, either side can send unary requests, receive responses, send one-way notifications, and open server-streaming, client-streaming, or bidirectional-streaming calls. Streaming peers negotiate independent item and byte receive windows. Send waits for receiver credit without blocking delivery on other streams.

Optional features include HMAC shared-secret authentication, gzip payload compression, inbound interceptors, explicit singleflight calls, server broadcast notifications, and backpressure limits for pending calls, active streams, and concurrent writes. PeerManager can coordinate all listeners and dial paths in a process so two full-duplex peers retain one physical connection regardless of which side initiated it.

The dialing Client reconnects aggressively after network loss. Calls and streams already in flight fail with ErrUnavailable instead of being replayed, because the remote peer may already have processed the request or some stream items. New calls and new streams can use the re-established connection.

For a first call, see the executable package example. The repository's examples directory includes complete unary, streaming, and managed-peer programs, with run instructions and expected output. Shared-secret auth authenticates the dialer only; it does not encrypt traffic or authorize individual functions. See the repository's service checklist before deployment.

Example
package main

import (
	"context"
	"fmt"
	"net"
	"time"

	"github.com/dan-sherwin/gorpc"
)

func main() {
	type request struct {
		Name string
	}
	type response struct {
		Greeting string
	}

	server := gorpc.NewServer(gorpc.ServerOptions{})
	gorpc.MustRegister(server, "greet", func(_ *gorpc.Context, req request) (response, error) {
		return response{Greeting: "Hello, " + req.Name}, nil
	})
	listener, err := net.Listen("tcp", "127.0.0.1:0")
	if err != nil {
		panic(err)
	}
	served := make(chan error, 1)
	go func() { served <- server.ServeListener(listener) }()
	defer func() {
		ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
		defer cancel()
		if err := server.Shutdown(ctx); err != nil {
			panic(err)
		}
		if err := <-served; err != nil {
			panic(err)
		}
	}()

	ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
	defer cancel()
	client := gorpc.NewTCPClient(listener.Addr().String(), "example")
	defer func() { _ = client.Close() }()
	if err := client.Connect(ctx); err != nil {
		panic(err)
	}
	var reply response
	if err := client.CallContext(ctx, "greet", request{Name: "Go"}, &reply); err != nil {
		panic(err)
	}
	fmt.Println(reply.Greeting)
}
Output:
Hello, Go

Index

Examples

Constants

View Source
const (
	// BackpressureSideClient means the dialing side rejected new local work.
	BackpressureSideClient = "client"
	// BackpressureSideServer means the accepting side rejected new local work.
	BackpressureSideServer = "server"

	// BackpressureReasonPendingCalls means the pending request limit was reached.
	BackpressureReasonPendingCalls = "pending_calls"
	// BackpressureReasonActiveStreams means the active stream limit was reached.
	BackpressureReasonActiveStreams = "active_streams"
	// BackpressureReasonConcurrentWrites means the concurrent write limit was reached.
	BackpressureReasonConcurrentWrites = "concurrent_writes"
)
View Source
const (
	DefaultDialTimeout       = 5 * time.Second
	DefaultWriteTimeout      = 10 * time.Second
	DefaultReconnectMinDelay = 100 * time.Millisecond
	DefaultReconnectMaxDelay = 5 * time.Second
	DefaultReconnectJitter   = 0.2
	DefaultPingInterval      = 10 * time.Second
	DefaultPingTimeout       = 3 * time.Second
)

Reconnect defaults used by Client when options are unset.

View Source
const (
	ErrorCodeCanceled         = "canceled"
	ErrorCodeDeadlineExceeded = "deadline_exceeded"
	ErrorCodeInternal         = "internal"
	ErrorCodeInvalidRequest   = "invalid_request"
	ErrorCodeNotFound         = "not_found"
	ErrorCodeUnauthorized     = "unauthorized"
	ErrorCodeUnavailable      = "unavailable"
	ErrorCodeBackpressure     = "backpressure"
	ErrorCodePeerConnected    = "peer_connected"
)

Remote error codes used by the built-in server and helpers.

View Source
const CodecMessagePack = "msgpack"

CodecMessagePack is the v1 MessagePack codec name used during handshake.

View Source
const CompressionGzip = "gzip"

CompressionGzip is the built-in gzip compressor name used during handshake.

View Source
const DefaultHandshakeTimeout = 5 * time.Second

DefaultHandshakeTimeout is the default timeout for the initial protocol handshake.

View Source
const DefaultMaxFrameSize int64 = 64 * 1024 * 1024

DefaultMaxFrameSize limits both encoded frames and uncompressed payloads.

View Source
const ProtocolVersion uint16 = 1

ProtocolVersion is the current GoRPC wire protocol version.

Variables

View Source
var (
	ErrClosed               = errors.New("gorpc: closed")
	ErrAuthentication       = errors.New("gorpc: authentication failed")
	ErrDuplicateFunction    = errors.New("gorpc: duplicate function")
	ErrInvalidFunction      = errors.New("gorpc: invalid function")
	ErrInvalidHandler       = errors.New("gorpc: invalid handler")
	ErrInvalidResponse      = errors.New("gorpc: invalid response")
	ErrUnavailable          = errors.New("gorpc: unavailable")
	ErrBackpressure         = errors.New("gorpc: backpressure")
	ErrPeerConnected        = errors.New("gorpc: peer already connected")
	ErrPeerConfiguration    = errors.New("gorpc: peer configuration conflict")
	ErrPeerIdentityRequired = errors.New("gorpc: peer identity is required")
	ErrPeerSelfConnection   = errors.New("gorpc: peer cannot connect to itself")
)

Common GoRPC errors.

View Source
var (
	ErrFrameTooLarge = errors.New("gorpc: frame too large")
	ErrProtocol      = errors.New("gorpc: protocol error")
)

Frame read/write errors.

Functions

func Call

func Call[Req, Resp any](ctx context.Context, client *Client, function string, req Req) (Resp, error)

Call performs a typed unary request/response call.

func MustRegister

func MustRegister[Req, Resp any](target any, function string, fn HandlerFunc[Req, Resp])

MustRegister is Register that panics on error.

func MustRegisterBidiStream added in v0.5.0

func MustRegisterBidiStream[Recv, Send any](target any, function string, fn BidiStreamHandlerFunc[Recv, Send])

MustRegisterBidiStream is RegisterBidiStream that panics on error.

func MustRegisterClientStream added in v0.5.0

func MustRegisterClientStream[Item, Resp any](target any, function string, fn ClientStreamHandlerFunc[Item, Resp])

MustRegisterClientStream is RegisterClientStream that panics on error.

func MustRegisterNotify added in v0.4.0

func MustRegisterNotify[Req any](target any, function string, fn NotifyHandlerFunc[Req])

MustRegisterNotify is RegisterNotify that panics on error.

func MustRegisterServerStream added in v0.5.0

func MustRegisterServerStream[Req, Item any](target any, function string, fn ServerStreamHandlerFunc[Req, Item])

MustRegisterServerStream is RegisterServerStream that panics on error.

func Notify added in v0.4.0

func Notify[Req any](ctx context.Context, client *Client, function string, req Req) error

Notify sends a typed one-way notification.

func Register

func Register[Req, Resp any](target any, function string, fn HandlerFunc[Req, Resp]) error

Register binds a typed unary handler to a function name. The target can be a *Server, for functions the accepted side handles, or a *Client, for functions the dialing side handles after the connection is established.

func RegisterBidiStream added in v0.5.0

func RegisterBidiStream[Recv, Send any](target any, function string, fn BidiStreamHandlerFunc[Recv, Send]) error

RegisterBidiStream binds a typed bidirectional-streaming handler to a function name. Both sides can send and receive stream items.

func RegisterClientStream added in v0.5.0

func RegisterClientStream[Item, Resp any](target any, function string, fn ClientStreamHandlerFunc[Item, Resp]) error

RegisterClientStream binds a typed client-streaming handler to a function name. The caller sends zero or more request items and receives one response.

func RegisterNotify added in v0.4.0

func RegisterNotify[Req any](target any, function string, fn NotifyHandlerFunc[Req]) error

RegisterNotify binds a typed one-way notification handler to a function name. The sender gets write success/failure only; handler errors are local to the receiver.

func RegisterServerStream added in v0.5.0

func RegisterServerStream[Req, Item any](target any, function string, fn ServerStreamHandlerFunc[Req, Item]) error

RegisterServerStream binds a typed server-streaming handler to a function name. The caller sends one request and receives zero or more response items.

Types

type Auth added in v0.2.0

type Auth struct {
	// contains filtered or unexported fields
}

Auth configures optional connection authentication.

func SharedSecret added in v0.2.0

func SharedSecret(secret string) Auth

SharedSecret enables HMAC-SHA256 challenge/response authentication.

type BackpressureInfo added in v1.0.0

type BackpressureInfo struct {
	Side      string
	Reason    string
	Limit     int
	RequestID uint64
	Function  string
	FrameType FrameType
}

BackpressureInfo describes a rejected operation.

type BackpressureOptions added in v1.0.0

type BackpressureOptions struct {
	MaxPendingCalls     int
	MaxActiveStreams    int
	MaxConcurrentWrites int
	OnBackpressure      func(BackpressureInfo)
}

BackpressureOptions controls optional limits that reject new work before memory grows without bound. Zero values keep existing unlimited behavior.

type BidiStreamHandle added in v0.5.0

type BidiStreamHandle[Send, Recv any] struct {
	// contains filtered or unexported fields
}

BidiStreamHandle is a typed bidirectional stream wrapper. Send and Recv can be used concurrently by different goroutines.

func BidiStream added in v0.5.0

func BidiStream[Send, Recv any](ctx context.Context, target any, function string) (*BidiStreamHandle[Send, Recv], error)

BidiStream opens a bidirectional stream. Send and Recv can be used concurrently by different goroutines. Recv returns io.EOF when the remote send side closes cleanly.

The target can be either *Client or an accepted *Conn, so either connected side can open a stream to the other side.

func BidiStreamWithOptions added in v1.0.0

func BidiStreamWithOptions[Send, Recv any](ctx context.Context, target any, function string, opts StreamOptions) (*BidiStreamHandle[Send, Recv], error)

BidiStreamWithOptions opens a bidirectional stream with stream options.

func (*BidiStreamHandle[Send, Recv]) Cancel added in v0.5.0

func (s *BidiStreamHandle[Send, Recv]) Cancel() error

Cancel cancels the stream and sends a best-effort cancel frame.

func (*BidiStreamHandle[Send, Recv]) CloseSend added in v0.5.0

func (s *BidiStreamHandle[Send, Recv]) CloseSend() error

CloseSend closes the local sending side of the stream.

func (*BidiStreamHandle[Send, Recv]) Recv added in v0.5.0

func (s *BidiStreamHandle[Send, Recv]) Recv() (Recv, error)

Recv receives one typed stream item. It returns io.EOF when the remote side cleanly closes its send side.

func (*BidiStreamHandle[Send, Recv]) Send added in v0.5.0

func (s *BidiStreamHandle[Send, Recv]) Send(item Send) error

Send sends one typed stream item.

func (*BidiStreamHandle[Send, Recv]) Stream added in v0.5.0

func (s *BidiStreamHandle[Send, Recv]) Stream() *Stream

Stream returns the raw stream.

type BidiStreamHandlerFunc added in v0.5.0

type BidiStreamHandlerFunc[Recv, Send any] func(*Context, *BidiStreamHandle[Send, Recv]) error

BidiStreamHandlerFunc receives and sends stream items independently. Closing one sending direction leaves the other direction open.

type BroadcastResult added in v1.0.0

type BroadcastResult struct {
	Total  int
	Sent   int
	Failed int
	Errors map[*Conn]error
}

BroadcastResult reports the outcome of Server.NotifyAll.

func (BroadcastResult) OK added in v1.0.0

func (r BroadcastResult) OK() bool

OK reports whether every connection accepted the broadcast frame.

type Client

type Client struct {
	// contains filtered or unexported fields
}

Client is the dialing side of a long-lived full-duplex GoRPC connection. It can send requests, register functions for the accepted side to call, and reconnects automatically after connection loss until Close is called.

func Dial

func Dial(ctx context.Context, network, address string, opts ClientOptions) (*Client, error)

Dial connects to a GoRPC server, completes the protocol handshake, and starts background reconnect monitoring.

func NewClient added in v0.3.0

func NewClient(network, address string, opts ClientOptions) *Client

NewClient creates a client without connecting it. Use this when the dialing side needs to register functions before the accepted side can call them.

func NewTCPClient added in v0.3.0

func NewTCPClient(address, clientName string, opts ...ClientOptions) *Client

NewTCPClient creates a TCP client without connecting it.

func NewUnixClient added in v0.3.0

func NewUnixClient(path, clientName string, opts ...ClientOptions) *Client

NewUnixClient creates a Unix socket client without connecting it.

func NewUnixPacketClient added in v0.3.0

func NewUnixPacketClient(path, clientName string, opts ...ClientOptions) *Client

NewUnixPacketClient creates a Unix packet socket client without connecting it.

func TCPDial added in v0.2.0

func TCPDial(address, clientName string, opts ...ClientOptions) (*Client, error)

TCPDial connects to address using TCP and reconnects automatically until Close is called.

func UnixDial added in v0.2.0

func UnixDial(path, clientName string, opts ...ClientOptions) (*Client, error)

UnixDial connects to path using a Unix socket and reconnects automatically until Close is called.

func UnixPacketDial added in v0.2.0

func UnixPacketDial(path, clientName string, opts ...ClientOptions) (*Client, error)

UnixPacketDial connects to path using a Unix packet socket and reconnects automatically until Close is called.

func (*Client) AsyncCall added in v0.2.0

func (c *Client) AsyncCall(function string, req any, handler any, correlationID string) error

AsyncCall sends a unary request and invokes handler when the response arrives.

func (*Client) AsyncCallContext added in v0.2.0

func (c *Client) AsyncCallContext(ctx context.Context, function string, req any, handler any, correlationID string) error

AsyncCallContext sends a unary request and invokes handler when the response arrives. The context only controls waiting for a connection and writing the request frame.

func (*Client) AsyncCallWithTimeout added in v0.2.0

func (c *Client) AsyncCallWithTimeout(function string, req any, handler any, correlationID string, timeout time.Duration) error

AsyncCallWithTimeout sends a unary request using a timeout while waiting for a connection and writing the request frame. The response handler runs later.

func (*Client) Call added in v0.2.0

func (c *Client) Call(function string, req any, resp any) error

Call performs a unary request/response call using context.Background.

func (*Client) CallContext added in v0.2.0

func (c *Client) CallContext(ctx context.Context, function string, req any, resp any) error

CallContext performs a unary request/response call. If the client is reconnecting, CallContext waits for the next connection until ctx is canceled.

func (*Client) CallSingleflight added in v1.0.0

func (c *Client) CallSingleflight(function string, key string, req any, resp any) error

CallSingleflight performs a unary request/response call and collapses concurrent calls with the same function and key into one remote request.

func (*Client) CallSingleflightContext added in v1.0.0

func (c *Client) CallSingleflightContext(ctx context.Context, function string, key string, req any, resp any) error

CallSingleflightContext performs a unary call and shares one in-flight remote request with concurrent callers using the same function and key. If key is empty, GoRPC builds a key from the encoded request payload.

func (*Client) CallSingleflightWithTimeout added in v1.0.0

func (c *Client) CallSingleflightWithTimeout(function string, key string, req any, resp any, timeout time.Duration) error

CallSingleflightWithTimeout performs a singleflight unary call with a timeout.

func (*Client) CallWithTimeout added in v0.2.0

func (c *Client) CallWithTimeout(function string, req any, resp any, timeout time.Duration) error

CallWithTimeout performs a unary request/response call with a timeout.

func (*Client) Close

func (c *Client) Close() error

Close closes the client and stops reconnect attempts.

func (*Client) Connect added in v0.3.0

func (c *Client) Connect(ctx context.Context) error

Connect establishes the first connection and starts background reconnect monitoring. It is called automatically by Dial and the TCPDial helpers.

func (*Client) Notify added in v0.4.0

func (c *Client) Notify(function string, req any) error

Notify sends a one-way typed notification using context.Background.

func (*Client) NotifyContext added in v0.4.0

func (c *Client) NotifyContext(ctx context.Context, function string, req any) error

NotifyContext sends a one-way typed notification. Success means the frame was written locally; GoRPC does not wait for remote handler completion or remote errors.

func (*Client) NotifyWithTimeout added in v0.4.0

func (c *Client) NotifyWithTimeout(function string, req any, timeout time.Duration) error

NotifyWithTimeout sends a one-way typed notification with a timeout while waiting for a connection and writing the notification frame.

func (*Client) SupportsStreamFlowControl added in v1.0.0

func (c *Client) SupportsStreamFlowControl() bool

SupportsStreamFlowControl reports whether the active connection negotiated stream credit. It returns false while disconnected or connected to an older peer.

func (*Client) WaitReady added in v0.2.0

func (c *Client) WaitReady(ctx context.Context) error

WaitReady blocks until the client has an active connection or ctx is canceled.

type ClientContext added in v0.2.0

type ClientContext interface {
	CorrelationID() string
	RequestID() uint64
	Function() string
	Error() error
}

ClientContext is passed to asynchronous response handlers for requests made by either a Client or an accepted Conn.

type ClientFunc added in v0.2.0

type ClientFunc[Req, Resp any] func(context.Context, Req) (Resp, error)

ClientFunc is the typed function shape returned by Function.

func Function added in v0.2.0

func Function[Req, Resp any](client *Client, function string) ClientFunc[Req, Resp]

Function returns a typed client function bound to a remote function name.

type ClientOptions

type ClientOptions struct {
	ClientName   string
	Codec        Codec
	Compression  Compressor
	MaxFrameSize int64
	// HandshakeTimeout bounds protocol negotiation after dialing.
	HandshakeTimeout time.Duration
	Auth             Auth
	// DialTimeout bounds network dialing, independently of HandshakeTimeout.
	DialTimeout  time.Duration
	WriteTimeout time.Duration
	Logger       *slog.Logger
	Dialer       *net.Dialer

	ReconnectMinDelay time.Duration
	ReconnectMaxDelay time.Duration
	ReconnectJitter   float64
	PingInterval      time.Duration
	PingTimeout       time.Duration

	Backpressure      BackpressureOptions
	StreamOptions     StreamOptions
	UnaryInterceptor  UnaryInterceptor
	NotifyInterceptor NotifyInterceptor
	StreamInterceptor StreamInterceptor
}

ClientOptions configures Dial and the network-specific dial helpers.

type ClientStreamHandle added in v0.5.0

type ClientStreamHandle[Item, Resp any] struct {
	// contains filtered or unexported fields
}

ClientStreamHandle is returned by ClientStream. It lets the caller send many request items and then receive one final response.

func ClientStream added in v0.5.0

func ClientStream[Item, Resp any](ctx context.Context, target any, function string) (*ClientStreamHandle[Item, Resp], error)

ClientStream opens a client-streaming call. The caller sends zero or more typed items and then calls CloseAndRecv for the final typed response.

The target can be either *Client or an accepted *Conn, so either connected side can open a stream to the other side.

func ClientStreamWithOptions added in v1.0.0

func ClientStreamWithOptions[Item, Resp any](ctx context.Context, target any, function string, opts StreamOptions) (*ClientStreamHandle[Item, Resp], error)

ClientStreamWithOptions opens a client-streaming call with stream options.

func (*ClientStreamHandle[Item, Resp]) Cancel added in v0.5.0

func (s *ClientStreamHandle[Item, Resp]) Cancel() error

Cancel cancels the stream and sends a best-effort cancel frame.

func (*ClientStreamHandle[Item, Resp]) CloseAndRecv added in v0.5.0

func (s *ClientStreamHandle[Item, Resp]) CloseAndRecv() (Resp, error)

CloseAndRecv closes the local sending side and waits for the final typed response.

func (*ClientStreamHandle[Item, Resp]) Send added in v0.5.0

func (s *ClientStreamHandle[Item, Resp]) Send(item Item) error

Send sends one typed request item.

func (*ClientStreamHandle[Item, Resp]) Stream added in v0.5.0

func (s *ClientStreamHandle[Item, Resp]) Stream() *Stream

Stream returns the raw stream.

type ClientStreamHandlerFunc added in v0.5.0

type ClientStreamHandlerFunc[Item, Resp any] func(*Context, *StreamReader[Item]) (Resp, error)

ClientStreamHandlerFunc receives zero or more request items and returns one final response.

type Codec

type Codec interface {
	Name() string
	Marshal(v any) ([]byte, error)
	Unmarshal(data []byte, v any) error
}

Codec marshals frame envelopes and function payloads. Implementations must be safe for concurrent use and ignore unknown fields when decoding envelopes and handshakes, so optional protocol extensions can coexist with older peers.

type Compressor added in v1.0.0

type Compressor interface {
	Name() string
	Compress(data []byte) ([]byte, error)
	Decompress(data []byte) ([]byte, error)
}

Compressor compresses and decompresses frame payloads after the GoRPC handshake negotiates a matching compressor on both peers. Implementations must be safe for concurrent use.

func GzipCompression added in v1.0.0

func GzipCompression() Compressor

GzipCompression returns the built-in gzip compressor.

type Conn added in v0.3.0

type Conn struct {
	// contains filtered or unexported fields
}

Conn is one accepted GoRPC connection. A Conn can receive requests through server-registered functions and can also initiate requests back to the client over the same full-duplex connection.

func (*Conn) AsyncCall added in v0.3.0

func (c *Conn) AsyncCall(function string, req any, handler any, correlationID string) error

AsyncCall sends a unary request to the connected client and invokes handler when the response arrives.

func (*Conn) AsyncCallContext added in v0.3.0

func (c *Conn) AsyncCallContext(ctx context.Context, function string, req any, handler any, correlationID string) error

AsyncCallContext sends a unary request to the connected client and invokes handler when the response arrives.

func (*Conn) AsyncCallWithTimeout added in v0.3.0

func (c *Conn) AsyncCallWithTimeout(function string, req any, handler any, correlationID string, timeout time.Duration) error

AsyncCallWithTimeout sends a unary request to the connected client using a timeout while writing the request frame. The response handler runs later.

func (*Conn) Call added in v0.3.0

func (c *Conn) Call(function string, req any, resp any) error

Call performs a unary request/response call to the connected client.

func (*Conn) CallContext added in v0.3.0

func (c *Conn) CallContext(ctx context.Context, function string, req any, resp any) error

CallContext performs a unary request/response call to the connected client.

func (*Conn) CallSingleflight added in v1.0.0

func (c *Conn) CallSingleflight(function string, key string, req any, resp any) error

CallSingleflight performs a unary request/response call to the connected client and collapses concurrent calls with the same function and key into one remote request.

func (*Conn) CallSingleflightContext added in v1.0.0

func (c *Conn) CallSingleflightContext(ctx context.Context, function string, key string, req any, resp any) error

CallSingleflightContext performs a unary call to the connected client and shares one in-flight remote request with concurrent callers using the same function and key. If key is empty, GoRPC builds a key from the encoded request payload.

func (*Conn) CallSingleflightWithTimeout added in v1.0.0

func (c *Conn) CallSingleflightWithTimeout(function string, key string, req any, resp any, timeout time.Duration) error

CallSingleflightWithTimeout performs a singleflight call with a timeout.

func (*Conn) CallWithTimeout added in v0.3.0

func (c *Conn) CallWithTimeout(function string, req any, resp any, timeout time.Duration) error

CallWithTimeout performs a unary request/response call to the connected client with a timeout.

func (*Conn) ClientName added in v0.3.0

func (c *Conn) ClientName() string

ClientName returns the self-reported client name from the connection handshake. It is useful for logs and metrics, but is not authenticated.

func (*Conn) Close added in v0.3.0

func (c *Conn) Close() error

Close closes the accepted connection, cancels active inbound handlers, and fails pending outbound calls.

func (*Conn) ConnectionGeneration added in v1.0.0

func (c *Conn) ConnectionGeneration() uint64

ConnectionGeneration returns the opaque, process-local generation assigned to this physical accepted connection. The value is nonzero and is never reused by another live connection in the process.

func (*Conn) Done added in v1.0.0

func (c *Conn) Done() <-chan struct{}

Done returns a channel that is closed when this physical accepted connection closes. It remains open for the lifetime of the active connection. Done returns nil when called on a nil Conn.

func (*Conn) LocalAddr added in v0.3.0

func (c *Conn) LocalAddr() net.Addr

LocalAddr returns the local address for the connection.

func (*Conn) Notify added in v0.4.0

func (c *Conn) Notify(function string, req any) error

Notify sends a one-way typed notification to the connected client.

func (*Conn) NotifyContext added in v0.4.0

func (c *Conn) NotifyContext(ctx context.Context, function string, req any) error

NotifyContext sends a one-way typed notification to the connected client. Success means the frame was written locally; GoRPC does not wait for remote handler completion or remote errors.

func (*Conn) NotifyWithTimeout added in v0.4.0

func (c *Conn) NotifyWithTimeout(function string, req any, timeout time.Duration) error

NotifyWithTimeout sends a one-way typed notification to the connected client with a timeout while writing the notification frame.

func (*Conn) RemoteAddr added in v0.3.0

func (c *Conn) RemoteAddr() net.Addr

RemoteAddr returns the peer address for the connection.

func (*Conn) SupportsStreamFlowControl added in v1.0.0

func (c *Conn) SupportsStreamFlowControl() bool

SupportsStreamFlowControl reports whether this connection negotiated stream credit.

type Context added in v0.2.0

type Context struct {
	context.Context
	// contains filtered or unexported fields
}

Context is the message-scoped context passed to request and notification handlers.

func (*Context) Call added in v0.3.0

func (c *Context) Call(function string, req any, resp any) error

Call performs a unary request/response call back over the same accepted connection that delivered this request.

func (*Context) CallContext added in v0.3.0

func (c *Context) CallContext(ctx context.Context, function string, req any, resp any) error

CallContext performs a unary request/response call back over the same accepted connection.

func (*Context) CallWithTimeout added in v0.3.0

func (c *Context) CallWithTimeout(function string, req any, resp any, timeout time.Duration) error

CallWithTimeout performs a unary request/response call back over the same accepted connection with a timeout.

func (*Context) ClientName added in v0.2.0

func (c *Context) ClientName() string

ClientName returns the self-reported client name from the connection handshake. It is useful for logs and metrics, but is not authenticated.

func (*Context) Conn added in v0.3.0

func (c *Context) Conn() *Conn

Conn returns the accepted connection that delivered this request when the handler is running on a Server. Client-side handlers return nil here because they can already call back through their Client.

func (*Context) ConnectionGeneration added in v1.0.0

func (c *Context) ConnectionGeneration() uint64

ConnectionGeneration returns the opaque, process-local generation of the physical connection that delivered this message. It is always nonzero for a dispatched handler and changes when a dialing Client reconnects, even though the logical Client or Peer remains the same.

func (*Context) Function added in v0.2.0

func (c *Context) Function() string

Function returns the remote function name for the request.

func (*Context) IsNotify added in v0.4.0

func (c *Context) IsNotify() bool

IsNotify reports whether the inbound message is a one-way notification.

func (*Context) IsStream added in v0.5.0

func (c *Context) IsStream() bool

IsStream reports whether the inbound message opened a stream.

func (*Context) LocalAddr added in v0.2.0

func (c *Context) LocalAddr() net.Addr

LocalAddr returns the local address for the connection.

func (*Context) Notify added in v0.4.0

func (c *Context) Notify(function string, req any) error

Notify sends a one-way typed notification back over the same accepted connection that delivered this request.

func (*Context) NotifyContext added in v0.4.0

func (c *Context) NotifyContext(ctx context.Context, function string, req any) error

NotifyContext sends a one-way typed notification back over the same accepted connection.

func (*Context) NotifyWithTimeout added in v0.4.0

func (c *Context) NotifyWithTimeout(function string, req any, timeout time.Duration) error

NotifyWithTimeout sends a one-way typed notification back over the same accepted connection with a timeout while writing the notification frame.

func (*Context) RemoteAddr added in v0.2.0

func (c *Context) RemoteAddr() net.Addr

RemoteAddr returns the peer address for the connection.

func (*Context) RequestID added in v0.2.0

func (c *Context) RequestID() uint64

RequestID returns the request or notification ID from the GoRPC frame.

func (*Context) StreamKind added in v0.5.0

func (c *Context) StreamKind() StreamKind

StreamKind returns the stream shape for streaming handlers. For non-stream handlers it returns zero.

type Frame

type Frame struct {
	Version          uint16     `msgpack:"version"`
	Type             FrameType  `msgpack:"type"`
	RequestID        uint64     `msgpack:"request_id,omitempty"`
	Function         string     `msgpack:"function,omitempty"`
	StreamKind       StreamKind `msgpack:"stream_kind,omitempty"`
	Compression      string     `msgpack:"compression,omitempty"`
	DeadlineUnixNano int64      `msgpack:"deadline_unix_nano,omitempty"`
	Payload          []byte     `msgpack:"payload,omitempty"`
	WindowItems      uint64     `msgpack:"window_items,omitempty"`
	WindowBytes      uint64     `msgpack:"window_bytes,omitempty"`
}

Frame is the v1 wire envelope. It is MessagePack-encoded and written with a 4-byte big-endian length prefix.

type FrameType

type FrameType uint8

FrameType identifies the kind of message carried by a frame.

const (
	FrameHello FrameType = iota + 1
	FrameHelloAck
	FrameRequest
	FrameResponse
	FrameError
	FrameCancel
	FramePing
	FramePong
	FrameStreamItem
	FrameStreamEnd
	FrameAuth
	FrameAuthAck
	FrameNotify
	FrameStreamStart
	FrameStreamWindow
)

Frame types used by the v1 protocol.

func (FrameType) String

func (t FrameType) String() string

type HandlerFunc

type HandlerFunc[Req, Resp any] func(*Context, Req) (Resp, error)

HandlerFunc is the typed function shape used by registered unary functions.

type LimitedDecompressor added in v1.0.0

type LimitedDecompressor interface {
	Compressor
	DecompressLimit(data []byte, maxSize int64) ([]byte, error)
}

LimitedDecompressor bounds decompressed output while it is being produced. Custom compressors should implement it to avoid allocating an oversized payload before MaxFrameSize can be checked.

type MessagePackCodec

type MessagePackCodec struct{}

MessagePackCodec is the default v1 codec. It accepts one complete MessagePack value with at most 64 nested arrays or maps.

func (MessagePackCodec) Marshal

func (MessagePackCodec) Marshal(v any) ([]byte, error)

Marshal encodes v as MessagePack.

func (MessagePackCodec) Name

func (MessagePackCodec) Name() string

Name returns the handshake name for MessagePackCodec.

func (MessagePackCodec) Unmarshal

func (MessagePackCodec) Unmarshal(data []byte, v any) error

Unmarshal decodes MessagePack data into v.

type NotifyFunc added in v0.4.0

type NotifyFunc[Req any] func(context.Context, Req) error

NotifyFunc is the typed function shape returned by Notification.

func Notification added in v0.4.0

func Notification[Req any](client *Client, function string) NotifyFunc[Req]

Notification returns a typed client notification function bound to a remote function name.

type NotifyHandler added in v1.0.0

type NotifyHandler func(*Context, NotifyRequest) error

NotifyHandler is the raw handler shape used by notification interceptors.

type NotifyHandlerFunc added in v0.4.0

type NotifyHandlerFunc[Req any] func(*Context, Req) error

NotifyHandlerFunc is the typed function shape used by registered one-way notification handlers.

type NotifyInterceptor added in v1.0.0

type NotifyInterceptor func(*Context, NotifyRequest, NotifyHandler) error

NotifyInterceptor wraps an inbound notification handler.

type NotifyRequest added in v1.0.0

type NotifyRequest struct {
	Payload []byte
}

NotifyRequest is passed to notification interceptors.

type Peer added in v1.0.0

type Peer struct {
	// contains filtered or unexported fields
}

Peer is one logical, full-duplex relationship. Calls are routed over the single active physical connection regardless of which side initiated it.

func (*Peer) AsyncCallContext added in v1.0.0

func (p *Peer) AsyncCallContext(ctx context.Context, function string, req any, handler any, correlationID string) error

AsyncCallContext starts an asynchronous unary call over the active connection.

func (*Peer) Call added in v1.0.0

func (p *Peer) Call(function string, req any, resp any) error

Call invokes a unary function over the active physical connection.

func (*Peer) CallContext added in v1.0.0

func (p *Peer) CallContext(ctx context.Context, function string, req any, resp any) error

CallContext invokes a unary function using ctx for readiness and call cancellation.

func (*Peer) CallWithTimeout added in v1.0.0

func (p *Peer) CallWithTimeout(function string, req any, resp any, timeout time.Duration) error

CallWithTimeout invokes a unary function and bounds the complete wait and call.

func (*Peer) EndpointForGeneration added in v1.0.0

func (p *Peer) EndpointForGeneration(generation uint64) (*PeerEndpoint, bool)

EndpointForGeneration returns a callback endpoint only when generation identifies the peer's exact current physical connection. The returned endpoint remains bound to that connection; after reconnect, its calls fail with ErrUnavailable instead of waiting for or switching to the replacement.

func (*Peer) Name added in v1.0.0

func (p *Peer) Name() string

Name returns the remote peer identity.

func (*Peer) Notify added in v1.0.0

func (p *Peer) Notify(function string, req any) error

Notify sends a one-way notification over the active physical connection.

func (*Peer) NotifyContext added in v1.0.0

func (p *Peer) NotifyContext(ctx context.Context, function string, req any) error

NotifyContext sends a one-way notification using ctx for readiness and cancellation.

func (*Peer) Status added in v1.0.0

func (p *Peer) Status() PeerStatus

Status returns the current logical connection state.

func (*Peer) WaitReady added in v1.0.0

func (p *Peer) WaitReady(ctx context.Context) error

WaitReady waits until the peer has one active physical connection.

type PeerClient added in v1.0.0

type PeerClient struct {
	// contains filtered or unexported fields
}

PeerClient is a caller-owned lease on a managed peer relationship. Closing one lease never interrupts another user of the same peer.

func (*PeerClient) AsyncCallContext added in v1.0.0

func (c *PeerClient) AsyncCallContext(ctx context.Context, function string, req any, handler any, correlationID string) error

AsyncCallContext starts an asynchronous unary call over this lease.

func (*PeerClient) Call added in v1.0.0

func (c *PeerClient) Call(function string, req any, resp any) error

Call invokes a unary function over this lease's active connection.

func (*PeerClient) CallContext added in v1.0.0

func (c *PeerClient) CallContext(ctx context.Context, function string, req any, resp any) error

CallContext invokes a unary function using ctx for readiness and cancellation.

func (*PeerClient) Close added in v1.0.0

func (c *PeerClient) Close() error

Close releases this caller's interest in keeping an outbound connection.

func (*PeerClient) Notify added in v1.0.0

func (c *PeerClient) Notify(function string, req any) error

Notify sends a one-way notification over this lease's active connection.

func (*PeerClient) NotifyContext added in v1.0.0

func (c *PeerClient) NotifyContext(ctx context.Context, function string, req any) error

NotifyContext sends a one-way notification using ctx for readiness and cancellation.

func (*PeerClient) Peer added in v1.0.0

func (c *PeerClient) Peer() *Peer

Peer returns the shared logical peer behind this lease.

func (*PeerClient) Status added in v1.0.0

func (c *PeerClient) Status() PeerStatus

Status returns a point-in-time snapshot of this lease's logical peer.

func (*PeerClient) WaitReady added in v1.0.0

func (c *PeerClient) WaitReady(ctx context.Context) error

WaitReady waits until this lease has an active physical connection.

type PeerDialOptions added in v1.0.0

type PeerDialOptions struct {
	PeerName         string
	Network          string
	Address          string
	ClientOptions    ClientOptions
	RegisterHandlers func(*Client) error
}

PeerDialOptions configure a managed outbound connection. PeerName is the authenticated application identity expected at the remote endpoint.

type PeerDirection added in v1.0.0

type PeerDirection string

PeerDirection identifies which side created the active physical connection.

const (
	// PeerDirectionInbound means the remote peer established the physical connection.
	PeerDirectionInbound PeerDirection = "inbound"
	// PeerDirectionOutbound means the local peer established the physical connection.
	PeerDirectionOutbound PeerDirection = "outbound"
)

type PeerEndpoint added in v1.0.0

type PeerEndpoint struct {
	// contains filtered or unexported fields
}

PeerEndpoint is a callback endpoint bound to one exact physical connection. It never waits for or switches to a replacement connection after reconnect. Obtain one with Peer.EndpointForGeneration.

func (*PeerEndpoint) Call added in v1.0.0

func (e *PeerEndpoint) Call(function string, req any, resp any) error

Call invokes a unary function on the endpoint's bound physical connection.

func (*PeerEndpoint) CallContext added in v1.0.0

func (e *PeerEndpoint) CallContext(ctx context.Context, function string, req any, resp any) error

CallContext invokes a unary function on the endpoint's bound physical connection. It returns ErrUnavailable when that connection is no longer current and never waits for a replacement connection.

func (*PeerEndpoint) CallWithTimeout added in v1.0.0

func (e *PeerEndpoint) CallWithTimeout(function string, req any, resp any, timeout time.Duration) error

CallWithTimeout invokes a unary function on the endpoint's bound physical connection with a timeout.

func (*PeerEndpoint) ConnectionGeneration added in v1.0.0

func (e *PeerEndpoint) ConnectionGeneration() uint64

ConnectionGeneration returns the physical connection generation captured by this endpoint. It returns zero for a nil endpoint.

func (*PeerEndpoint) Notify added in v1.0.0

func (e *PeerEndpoint) Notify(function string, req any) error

Notify sends a one-way notification on the endpoint's bound physical connection.

func (*PeerEndpoint) NotifyContext added in v1.0.0

func (e *PeerEndpoint) NotifyContext(ctx context.Context, function string, req any) error

NotifyContext sends a one-way notification on the endpoint's bound physical connection. It returns ErrUnavailable when that connection is no longer current and never waits for a replacement connection.

func (*PeerEndpoint) NotifyWithTimeout added in v1.0.0

func (e *PeerEndpoint) NotifyWithTimeout(function string, req any, timeout time.Duration) error

NotifyWithTimeout sends a one-way notification on the endpoint's bound physical connection with a timeout.

type PeerManager added in v1.0.0

type PeerManager struct {
	// contains filtered or unexported fields
}

PeerManager owns all GoRPC connections for one local application identity. A manager must be shared by every GoRPC server and dial path in the process.

func NewPeerManager added in v1.0.0

func NewPeerManager(localName string) *PeerManager

NewPeerManager creates a connection manager for localName.

func (*PeerManager) Close added in v1.0.0

func (m *PeerManager) Close() error

Close stops all pending dials and closes all managed connections.

func (*PeerManager) Dial added in v1.0.0

func (m *PeerManager) Dial(ctx context.Context, opts PeerDialOptions) (*PeerClient, error)

Dial acquires one logical peer and waits for either an existing inbound connection or one shared outbound dial to become ready.

func (*PeerManager) LocalName added in v1.0.0

func (m *PeerManager) LocalName() string

LocalName returns the identity placed in managed outbound handshakes.

func (*PeerManager) Peer added in v1.0.0

func (m *PeerManager) Peer(peerName string) (*Peer, bool)

Peer returns a known logical peer, including accepted-only peers.

func (*PeerManager) Peers added in v1.0.0

func (m *PeerManager) Peers() []*Peer

Peers returns a snapshot of all known logical peers.

type PeerStatus added in v1.0.0

type PeerStatus struct {
	PeerName      string
	Active        bool
	Direction     PeerDirection
	Network       string
	LocalAddress  string
	RemoteAddress string
	ConnectedAt   time.Time
	// ConnectionGeneration is the opaque, nonzero process-local identity of
	// the current physical connection. It changes after an automatic reconnect.
	ConnectionGeneration uint64
	// StreamFlowControl is true when the active connection negotiated stream credit.
	StreamFlowControl bool
	Dialing           bool
	LastError         string
}

PeerStatus is a point-in-time snapshot of one logical peer relationship.

type RemoteError

type RemoteError struct {
	Code    string         `msgpack:"code" json:"code"`
	Message string         `msgpack:"message" json:"message"`
	Details map[string]any `msgpack:"details,omitempty" json:"details,omitempty"`
}

RemoteError is sent in FrameError payloads and returned by callers when the server handled the request but rejected or failed it.

Example
package main

import (
	"errors"
	"fmt"

	"github.com/dan-sherwin/gorpc"
)

func main() {
	err := fmt.Errorf("upload: %w", gorpc.NewRemoteError(
		gorpc.ErrorCodeBackpressure, "too many streams", nil,
	))
	if errors.Is(err, gorpc.ErrBackpressure) {
		fmt.Println("backpressure: caller decides whether and when to retry")
	}
	var remote *gorpc.RemoteError
	if errors.As(err, &remote) {
		fmt.Println(remote.Code, remote.Message)
	}
}
Output:
backpressure: caller decides whether and when to retry
backpressure too many streams

func NewRemoteError

func NewRemoteError(code, message string, details map[string]any) *RemoteError

NewRemoteError creates a structured error suitable for returning from a handler.

func (*RemoteError) Error

func (e *RemoteError) Error() string

func (*RemoteError) Is added in v1.0.0

func (e *RemoteError) Is(target error) bool

Is matches remote transport and context errors to their local sentinels.

type Server

type Server struct {
	// contains filtered or unexported fields
}

Server accepts GoRPC connections, dispatches registered functions, and exposes accepted connections that can initiate requests back to the dialing side.

func NewServer

func NewServer(opts ServerOptions) *Server

NewServer creates a Server with default codec and limits where options are unset.

func (*Server) Connections added in v0.3.0

func (s *Server) Connections() []*Conn

Connections returns a snapshot of currently accepted connections.

func (*Server) NotifyAll added in v1.0.0

func (s *Server) NotifyAll(function string, req any) BroadcastResult

NotifyAll sends a one-way notification to every currently accepted connection. It snapshots the connection list before sending.

func (*Server) NotifyAllContext added in v1.0.0

func (s *Server) NotifyAllContext(ctx context.Context, function string, req any) BroadcastResult

NotifyAllContext sends a notification to every currently accepted connection. The context controls each write attempt; it does not wait for remote handler completion because notifications do not have responses.

func (*Server) NotifyAllWithTimeout added in v1.0.0

func (s *Server) NotifyAllWithTimeout(function string, req any, timeout time.Duration) BroadcastResult

NotifyAllWithTimeout sends a notification to every currently accepted connection with a timeout for the broadcast write loop.

func (*Server) ServeListener added in v0.2.0

func (s *Server) ServeListener(ln net.Listener) error

ServeListener accepts GoRPC connections from ln until Shutdown is called or the listener returns an unrecoverable error. It always closes ln before returning. A Server may serve multiple listeners concurrently.

func (*Server) ServeTCP added in v0.2.0

func (s *Server) ServeTCP(address string) error

ServeTCP listens on address with the "tcp" network and serves GoRPC connections.

func (*Server) ServeUnix added in v0.2.0

func (s *Server) ServeUnix(path string) error

ServeUnix listens on path with the "unix" network and serves GoRPC connections.

func (*Server) ServeUnixPacket added in v0.2.0

func (s *Server) ServeUnixPacket(path string) error

ServeUnixPacket listens on path with the "unixpacket" network and serves GoRPC connections.

func (*Server) Shutdown

func (s *Server) Shutdown(ctx context.Context) error

Shutdown closes all listeners and connections, and waits for handlers to exit.

type ServerOptions

type ServerOptions struct {
	Codec            Codec
	Compression      Compressor
	MaxFrameSize     int64
	HandshakeTimeout time.Duration
	Auth             Auth
	WriteTimeout     time.Duration
	Logger           *slog.Logger
	OnConnect        func(*Conn)
	OnDisconnect     func(*Conn)
	PeerManager      *PeerManager

	Backpressure      BackpressureOptions
	StreamOptions     StreamOptions
	UnaryInterceptor  UnaryInterceptor
	NotifyInterceptor NotifyInterceptor
	StreamInterceptor StreamInterceptor
}

ServerOptions configures a GoRPC server.

type ServerStreamHandlerFunc added in v0.5.0

type ServerStreamHandlerFunc[Req, Item any] func(*Context, Req, *StreamWriter[Item]) error

ServerStreamHandlerFunc handles one request and sends zero or more response items before returning.

type Stream added in v0.5.0

type Stream struct {
	// contains filtered or unexported fields
}

Stream is the raw bidirectional item stream used by the typed streaming helpers. Most callers should prefer ServerStream, ClientStream, BidiStream, and the typed handler registration functions.

func (*Stream) Cancel added in v0.5.0

func (s *Stream) Cancel() error

Cancel cancels the whole stream locally before sending a best-effort cancel frame to the remote side.

func (*Stream) CloseSend added in v0.5.0

func (s *Stream) CloseSend() error

CloseSend closes the local sending side. Receiving can continue. A blocked Send is interrupted; an item already being written precedes the end frame.

func (*Stream) Context added in v0.5.0

func (s *Stream) Context() context.Context

Context returns the stream context. It is canceled when the stream ends, the connection closes, or either peer cancels the stream.

func (*Stream) Function added in v0.5.0

func (s *Stream) Function() string

Function returns the remote function name for the stream.

func (*Stream) Recv added in v0.5.0

func (s *Stream) Recv(item any) error

Recv reads one stream item into item. It returns io.EOF after the remote side closes its send side and all buffered items have been consumed.

func (*Stream) RequestID added in v0.5.0

func (s *Stream) RequestID() uint64

RequestID returns the stream request ID.

func (*Stream) Send added in v0.5.0

func (s *Stream) Send(item any) error

Send encodes and writes one item. With negotiated flow control, it waits for receiver credit or cancellation. Calls to Send are serialized.

type StreamHandler added in v1.0.0

type StreamHandler func(*Context, StreamRequest, *Stream) ([]byte, error)

StreamHandler is the raw handler shape used by stream interceptors. Client streaming handlers return the final response payload; server and bidi stream handlers return nil payloads.

type StreamInterceptor added in v1.0.0

type StreamInterceptor func(*Context, StreamRequest, *Stream, StreamHandler) ([]byte, error)

StreamInterceptor wraps an inbound stream handler.

type StreamKind added in v0.5.0

type StreamKind uint8

StreamKind identifies the shape of a streaming function.

const (
	// StreamKindServer means the caller sends one request and the handler sends zero or more items.
	StreamKindServer StreamKind = iota + 1
	// StreamKindClient means the caller sends zero or more items and the handler sends one response.
	StreamKindClient
	// StreamKindBidi means both sides can send and receive stream items.
	StreamKindBidi
)

func (StreamKind) String added in v0.5.0

func (k StreamKind) String() string

type StreamOptions added in v1.0.0

type StreamOptions struct {
	// RecvBuffer limits queued items. The default is 16.
	RecvBuffer int
	// RecvBytes limits queued, uncompressed item payloads. The default is 64 MiB.
	RecvBytes int64
}

StreamOptions configures newly opened streams. Zero values keep GoRPC defaults.

type StreamReader added in v0.5.0

type StreamReader[T any] struct {
	// contains filtered or unexported fields
}

StreamReader is a typed receive-only stream wrapper.

func ServerStream added in v0.5.0

func ServerStream[Req, Item any](ctx context.Context, target any, function string, req Req) (*StreamReader[Item], error)

ServerStream opens a server-streaming call. The caller sends one request and receives zero or more typed items until Recv returns io.EOF or an error.

The target can be either *Client or an accepted *Conn, so either connected side can open a stream to the other side.

func ServerStreamWithOptions added in v1.0.0

func ServerStreamWithOptions[Req, Item any](ctx context.Context, target any, function string, req Req, opts StreamOptions) (*StreamReader[Item], error)

ServerStreamWithOptions opens a server-streaming call with stream options.

func (*StreamReader[T]) Cancel added in v0.5.0

func (r *StreamReader[T]) Cancel() error

Cancel cancels the stream and sends a best-effort cancel frame.

func (*StreamReader[T]) Recv added in v0.5.0

func (r *StreamReader[T]) Recv() (T, error)

Recv receives one typed stream item. It returns io.EOF when the remote side cleanly closes its send side.

func (*StreamReader[T]) Stream added in v0.5.0

func (r *StreamReader[T]) Stream() *Stream

Stream returns the raw stream.

type StreamRequest added in v1.0.0

type StreamRequest struct {
	Kind    StreamKind
	Payload []byte
}

StreamRequest is passed to stream interceptors.

type StreamWriter added in v0.5.0

type StreamWriter[T any] struct {
	// contains filtered or unexported fields
}

StreamWriter is a typed send-only stream wrapper.

func (*StreamWriter[T]) Close added in v0.5.0

func (w *StreamWriter[T]) Close() error

Close closes the local sending side of the stream.

func (*StreamWriter[T]) Send added in v0.5.0

func (w *StreamWriter[T]) Send(item T) error

Send sends one typed stream item.

func (*StreamWriter[T]) Stream added in v0.5.0

func (w *StreamWriter[T]) Stream() *Stream

Stream returns the raw stream.

type UnaryHandler added in v1.0.0

type UnaryHandler func(*Context, UnaryRequest) ([]byte, error)

UnaryHandler is the raw handler shape used by unary interceptors.

type UnaryInterceptor added in v1.0.0

type UnaryInterceptor func(*Context, UnaryRequest, UnaryHandler) ([]byte, error)

UnaryInterceptor wraps an inbound unary handler.

type UnaryRequest added in v1.0.0

type UnaryRequest struct {
	Payload []byte
}

UnaryRequest is passed to unary interceptors.

Directories

Path Synopsis
examples
bidistream command
Package main sends work and receives results concurrently on one stream.
Package main sends work and receives results concurrently on one stream.
clientstream command
Package main uploads chunks and verifies the receiver's final checksum.
Package main uploads chunks and verifies the receiver's final checksum.
inventory/api
Package api contains the inventory contract shared by the example programs.
Package api contains the inventory contract shared by the example programs.
inventory/client command
Package main runs the inventory example client.
Package main runs the inventory example client.
inventory/server command
Package main runs the inventory example server.
Package main runs the inventory example server.
peers command
Package main demonstrates full-duplex managed peers and connection replacement.
Package main demonstrates full-duplex managed peers and connection replacement.
serverstream command
Package main demonstrates a slow stream reader and early cancellation.
Package main demonstrates a slow stream reader and early cancellation.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL