Documentation
¶
Overview ¶
Package wt provides a high-level framework for building WebTransport applications in Go.
WebTransport is a modern protocol built on QUIC and HTTP/3 that provides multiplexed bidirectional streams, unidirectional streams, and unreliable datagrams — all without TCP's head-of-line blocking.
This package sits on top of quic-go/webtransport-go and provides:
- Path-based routing with parameter extraction
- Middleware stack (auth, logging, rate limiting, compression, metrics)
- Session management with rooms and pub/sub
- Message framing (length-prefixed) over streams
- Type-safe stream handlers using Go generics
- WebSocket fallback for browsers without WebTransport support
- Self-signed certificate generation for development
Quick Start ¶
server := wt.New(
wt.WithAddr(":4433"),
wt.WithSelfSignedTLS(),
)
server.Handle("/echo", func(c *wt.Context) {
for {
stream, err := c.AcceptStream()
if err != nil {
return
}
go func() {
defer stream.Close()
io.Copy(stream, stream)
}()
}
})
server.ListenAndServe()
Architecture ¶
The framework follows familiar Go patterns inspired by Gin, Echo, and Chi:
- Handlers receive a Context wrapping the WebTransport session
- Middleware uses the func(c *Context, next HandlerFunc) signature
- Route groups share prefixes and middleware
- Sessions are tracked in a SessionStore for enumeration and broadcast
Streams vs Datagrams ¶
WebTransport offers two data transport modes:
Streams are reliable and ordered, similar to TCP connections. Use them for messages that must arrive: chat messages, game events, file transfers.
Datagrams are unreliable and unordered, similar to UDP packets. Use them for data where the latest value matters more than every value: player positions, sensor readings, cursor locations.
WebSocket Fallback ¶
The [fallback] sub-package provides transparent WebSocket fallback for browsers that don't support WebTransport (notably Safari as of early 2026). Stream multiplexing is simulated over the single WebSocket connection.
Package wt provides a high-level framework for building WebTransport applications in Go. It sits on top of quic-go/webtransport-go and provides routing, middleware, session management, and codec support.
Example (Middleware) ¶
package main
import (
"fmt"
"github.com/rarebek/wt"
"github.com/rarebek/wt/middleware"
)
func main() {
server := wt.New(wt.WithAddr(":4433"), wt.WithSelfSignedTLS())
// Stack middleware
server.Use(middleware.DefaultLogger())
server.Use(middleware.Recover(nil))
server.Use(middleware.RateLimit(100))
server.Handle("/app", func(c *wt.Context) {
_ = c
})
fmt.Println("Middleware configured")
}
Output: Middleware configured
Index ¶
- Constants
- func AltSvcHeader(port int) string
- func AltSvcMiddleware(port int) func(http.Handler) http.Handler
- func CallTyped[T any](c *RPCClient, method string, params any) (T, error)
- func CertFingerprint(certDER []byte) string
- func Datagrams(c *Context) iter.Seq[[]byte]
- func DebugMux(s *Server) *http.ServeMux
- func DecodeBatch(data []byte) [][]byte
- func DefaultErrorPage(w http.ResponseWriter, _ *http.Request, code int, msg string)
- func HTMLErrorPage(w http.ResponseWriter, _ *http.Request, code int, msg string)
- func Hash(data []byte) string
- func IsConnectionError(err error) bool
- func IsMessageError(err error) bool
- func IsSessionClosed(err error) bool
- func IsStreamClosed(err error) bool
- func IsUpgradeError(err error) bool
- func JoinPath(segments ...string) string
- func KeepAlive(c *Context, interval time.Duration) func()
- func ListenAndServeWithGracefulShutdown(s *Server, drainTimeout time.Duration) error
- func Messages(s *Stream) iter.Seq[[]byte]
- func Must[T any](val T, err error) T
- func PProfMux() *http.ServeMux
- func Pipe(a, b *Stream) error
- func PipeRaw(rw io.ReadWriteCloser, s *Stream) error
- func RequiredFields(v any, fields ...string) error
- func Retry(ctx context.Context, cfg RetryConfig, fn func() error) error
- func ServerInfo() map[string]string
- func SetAltSvcHeader(w http.ResponseWriter, port int)
- func Streams(c *Context) iter.Seq[*Stream]
- func ValidateDatagramSize(data []byte) error
- func ValidateMessage(msg any) error
- type BackpressureWriter
- type BatcherOption
- type BufferedReader
- type CertRotator
- type CompressedStream
- type CompressionStats
- type ConnInfo
- type ConnectionError
- type Context
- func (c *Context) AcceptStream() (*Stream, error)
- func (c *Context) AcceptUniStream() (*ReceiveStream, error)
- func (c *Context) Close() error
- func (c *Context) CloseWithError(code uint32, msg string) error
- func (c *Context) Context() context.Context
- func (c *Context) Get(key string) (any, bool)
- func (c *Context) GetString(key string) string
- func (c *Context) ID() string
- func (c *Context) Info() ConnInfo
- func (c *Context) InfoJSON() string
- func (c *Context) LocalAddr() net.Addr
- func (c *Context) MustGet(key string) any
- func (c *Context) OpenStream() (*Stream, error)
- func (c *Context) OpenStreamSync() (*Stream, error)
- func (c *Context) OpenUniStream() (*SendStream, error)
- func (c *Context) Param(name string) string
- func (c *Context) Params() map[string]string
- func (c *Context) ReceiveDatagram() ([]byte, error)
- func (c *Context) ReceiveDatagramContext(ctx context.Context) ([]byte, error)
- func (c *Context) RemoteAddr() net.Addr
- func (c *Context) Request() *http.Request
- func (c *Context) SendDatagram(data []byte) error
- func (c *Context) SendDatagramSafe(data []byte) error
- func (c *Context) Server() *Server
- func (c *Context) Session() *webtransport.Session
- func (c *Context) Set(key string, value any)
- type ContextStream
- type DatagramBatcher
- type ErrorPageHandler
- type ErrorResponse
- type Event
- type EventBus
- type EventHandler
- type EventType
- type FlowControlMonitor
- type FlowStats
- type Group
- type HandlerFunc
- type HealthCheck
- type HealthResponse
- type InterceptedStream
- type InterceptorOption
- type KVSync
- func (kv *KVSync) Delete(key string)
- func (kv *KVSync) Get(key string, v any) error
- func (kv *KVSync) GetRaw(key string) (json.RawMessage, bool)
- func (kv *KVSync) Keys() []string
- func (kv *KVSync) Len() int
- func (kv *KVSync) OnChange(fn func(key string, value json.RawMessage))
- func (kv *KVSync) Set(key string, value any) error
- func (kv *KVSync) Snapshot() map[string]json.RawMessage
- type MessageError
- type MiddlewareFunc
- type MigrationEvent
- type MigrationWatcher
- type Option
- func WithAddr(addr string) Option
- func WithAutoCert(domain string, cacheDir string) Option
- func WithAutoCertMulti(domains []string, cacheDir string) Option
- func WithCertRotator(cr *CertRotator) Option
- func WithCheckOrigin(fn func(r *http.Request) bool) Option
- func WithIdleTimeout(d time.Duration) Option
- func WithQUICConfig(cfg QUICConfig) Option
- func WithSelfSignedTLS() Option
- func WithTLS(certFile, keyFile string) Option
- type PreflightResult
- type PresenceInfo
- type PresenceTracker
- func (pt *PresenceTracker) Count(room string) int
- func (pt *PresenceTracker) GetPresence(room string) []PresenceInfo
- func (pt *PresenceTracker) GetPresenceJSON(room string) []byte
- func (pt *PresenceTracker) Join(room string, c *Context)
- func (pt *PresenceTracker) Leave(room string, c *Context)
- func (pt *PresenceTracker) OnChange(fn func(room string, info PresenceInfo, event string))
- func (pt *PresenceTracker) SetMetadata(room, sessionID string, metadata map[string]any)
- func (pt *PresenceTracker) UpdateStatus(room, sessionID, status string)
- type Priority
- type PubSub
- func (ps *PubSub) Publish(topic string, data []byte)
- func (ps *PubSub) PublishExcept(topic string, data []byte, excludeID string)
- func (ps *PubSub) Subscribe(topic string, c *Context)
- func (ps *PubSub) SubscriberCount(topic string) int
- func (ps *PubSub) Topics() []string
- func (ps *PubSub) TopicsForSession(sessionID string) []string
- func (ps *PubSub) Unsubscribe(topic string, c *Context)
- func (ps *PubSub) UnsubscribeAll(c *Context)
- type QUICConfig
- type RPCClient
- type RPCError
- type RPCHandler
- type RPCRequest
- type RPCResponse
- type RPCServer
- type ReceiveStream
- type ReliableDatagram
- type ReliableOption
- type ResumeStore
- type ResumeToken
- type RetryConfig
- type RingBuffer
- type Room
- func (r *Room) Broadcast(data []byte)
- func (r *Room) BroadcastExcept(data []byte, excludeID string)
- func (r *Room) BroadcastStream(data []byte)
- func (r *Room) BroadcastStreamExcept(data []byte, excludeID string)
- func (r *Room) Count() int
- func (r *Room) ForEach(fn func(*Context))
- func (r *Room) Join(ctx *Context)
- func (r *Room) Leave(ctx *Context)
- func (r *Room) Members() []*Context
- func (r *Room) Name() string
- func (r *Room) OnJoin(fn func(*Context))
- func (r *Room) OnLeave(fn func(*Context))
- func (r *Room) SafeBroadcast(data []byte, logger *slog.Logger)
- func (r *Room) SafeBroadcastExcept(data []byte, excludeID string, logger *slog.Logger)
- type RoomManager
- type RoomMessage
- type RoomWithHistory
- func (r *RoomWithHistory) BroadcastAndRecord(senderID string, data []byte)
- func (r *RoomWithHistory) BroadcastExceptAndRecord(senderID string, data []byte)
- func (r *RoomWithHistory) ClearHistory()
- func (r *RoomWithHistory) History() []RoomMessage
- func (r *RoomWithHistory) HistorySize() int
- func (r *RoomWithHistory) ReplayHistory(c *Context)
- func (r *RoomWithHistory) ReplayHistorySince(c *Context, since time.Time)
- type RotatorOption
- type Route
- type Router
- type SendStream
- type Server
- func (s *Server) Addr() string
- func (s *Server) Broadcast(data []byte)
- func (s *Server) BroadcastExcept(data []byte, excludeID string)
- func (s *Server) CertHash() string
- func (s *Server) Close() error
- func (s *Server) Group(prefix string, mw ...MiddlewareFunc) *Group
- func (s *Server) Handle(pattern string, handler HandlerFunc, mw ...MiddlewareFunc)
- func (s *Server) ListenAndServe() error
- func (s *Server) Multicast(data []byte, filter func(*Context) bool)
- func (s *Server) MulticastStream(data []byte, filter func(*Context) bool)
- func (s *Server) OnConnect(fn func(*Context))
- func (s *Server) OnDisconnect(fn func(*Context))
- func (s *Server) OnShutdown(fn ShutdownHook)
- func (s *Server) Preflight() []string
- func (s *Server) PreflightCheck() PreflightResult
- func (s *Server) SessionCount() int
- func (s *Server) Sessions() *SessionStore
- func (s *Server) Shutdown(ctx context.Context) error
- func (s *Server) Use(mw ...MiddlewareFunc)
- type SessionCloseError
- type SessionStore
- func (ss *SessionStore) Add(ctx *Context)
- func (ss *SessionStore) Broadcast(data []byte)
- func (ss *SessionStore) CloseAll()
- func (ss *SessionStore) Count() int
- func (ss *SessionStore) Each(fn func(*Context))
- func (ss *SessionStore) FindByValue(key string, value any) []*Context
- func (ss *SessionStore) Get(id string) (*Context, bool)
- func (ss *SessionStore) IDs() []string
- func (ss *SessionStore) Remove(id string)
- type ShutdownHook
- type Stream
- func (s *Stream) CancelRead(code uint32)
- func (s *Stream) CancelWrite(code uint32)
- func (s *Stream) Close() error
- func (s *Stream) Raw() *webtransport.Stream
- func (s *Stream) Read(b []byte) (int, error)
- func (s *Stream) ReadMessage() ([]byte, error)
- func (s *Stream) SessionContext() *Context
- func (s *Stream) SetDeadline(t time.Time) error
- func (s *Stream) SetReadDeadline(t time.Time) error
- func (s *Stream) SetWriteDeadline(t time.Time) error
- func (s *Stream) WithContext(ctx context.Context) *ContextStream
- func (s *Stream) WithDeadline(deadline time.Time) *ContextStream
- func (s *Stream) WithTimeout(d time.Duration) *ContextStream
- func (s *Stream) Write(b []byte) (int, error)
- func (s *Stream) WriteMessage(data []byte) error
- type StreamCloseError
- type StreamConfig
- type StreamHandler
- type StreamInterceptor
- type StreamMux
- type StreamOptions
- type Tags
- func (t *Tags) AllTags() []string
- func (t *Tags) Count(tag string) int
- func (t *Tags) HasTag(sessionID, tag string) bool
- func (t *Tags) SessionsWithTag(tag string) []string
- func (t *Tags) Tag(sessionID, tag string)
- func (t *Tags) TagsForSession(sessionID string) []string
- func (t *Tags) Untag(sessionID, tag string)
- func (t *Tags) UntagAll(sessionID string)
- type Ticker
- type TypedDatagram
- type TypedPubSub
- func (tp *TypedPubSub[T]) Publish(topic string, msg T) error
- func (tp *TypedPubSub[T]) PublishExcept(topic string, msg T, excludeID string) error
- func (tp *TypedPubSub[T]) Subscribe(topic string, c *Context)
- func (tp *TypedPubSub[T]) Unsubscribe(topic string, c *Context)
- func (tp *TypedPubSub[T]) UnsubscribeAll(c *Context)
- type TypedRoom
- type TypedStream
- type UpgradeError
- type Validator
Examples ¶
Constants ¶
const ( CodeOK uint32 = 0 // Success / normal closure CodeForbidden uint32 = 403 // Authenticated but not allowed CodeNotFound uint32 = 404 // Route not found CodeTimeout uint32 = 408 // Session or stream timeout CodeTooManyRequests uint32 = 429 // Rate limit exceeded CodeInternalError uint32 = 500 // Server-side error // Framework error codes (0x1000+) CodeProtocolError uint32 = 0x1000 // Invalid protocol message CodeMessageTooLarge uint32 = 0x1001 // Message exceeds MaxMessageSize CodeInvalidCodec uint32 = 0x1002 // Unknown or misconfigured codec CodeStreamLimit uint32 = 0x1003 // Too many concurrent streams CodeSessionExpired uint32 = 0x1004 // Session TTL expired CodeBadRequest uint32 = 0x1005 // Malformed request data CodeShuttingDown uint32 = 0x1006 // Server is draining connections )
Standard error codes for WebTransport sessions and streams. HTTP-like codes (0-999) for familiar semantics. Framework codes (0x1000+) for wt-specific errors. QUIC transport codes (0x100+) are defined by RFC 9000 Section 20.
const DefaultKeepAliveInterval = 15 * time.Second
DefaultKeepAliveInterval is 15 seconds, chosen to be well under typical UDP NAT timeout of 20-30 seconds.
const MaxDatagramSize = 1200
MaxDatagramSize is the recommended maximum datagram payload size. QUIC datagrams are limited by the path MTU minus QUIC overhead. Typical safe size is ~1200 bytes (minimum QUIC MTU of 1280 minus headers). Larger datagrams may be fragmented or dropped.
const MaxMessageSize = 16 * 1024 * 1024
MaxMessageSize is the maximum size of a length-prefixed message (16 MB).
const Version = "0.1.0-dev"
Version is the current framework version.
Variables ¶
This section is empty.
Functions ¶
func AltSvcHeader ¶
AltSvcHeader returns the Alt-Svc HTTP header value that tells browsers to upgrade from HTTP/2 to HTTP/3 for WebTransport.
Browsers use this header to discover that a server supports HTTP/3. Include it in your HTTP/1.1 or HTTP/2 responses.
Usage:
// On your HTTP/1.1 or HTTP/2 server:
http.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) {
wt.SetAltSvcHeader(w, 4433) // WebTransport on port 4433
// ... serve regular HTTP
})
func AltSvcMiddleware ¶
AltSvcMiddleware returns an HTTP middleware that adds the Alt-Svc header to every response, advertising HTTP/3 availability.
func CertFingerprint ¶
CertFingerprint returns the hex-encoded SHA-256 fingerprint of a DER-encoded certificate.
func Datagrams ¶
Datagrams returns an iterator over incoming datagrams.
for data := range wt.Datagrams(c) {
process(data)
}
func DebugMux ¶
DebugMux returns an http.ServeMux with pprof, health check, and metrics. One endpoint for all debugging/monitoring needs.
Usage:
server := wt.New(...)
go http.ListenAndServe(":6060", wt.DebugMux(server))
func DecodeBatch ¶
DecodeBatch decodes a batch-encoded datagram into individual messages. Pre-allocates the result slice based on estimated message count.
func DefaultErrorPage ¶
DefaultErrorPage returns a JSON error response.
func HTMLErrorPage ¶
HTMLErrorPage returns an HTML error page.
func IsConnectionError ¶
IsConnectionError checks if an error is a connection-level failure.
func IsMessageError ¶
IsMessageError checks if an error is a message read/write failure.
func IsSessionClosed ¶
IsSessionClosed checks if an error is a session closure.
func IsStreamClosed ¶
IsStreamClosed checks if an error is a stream closure.
func IsUpgradeError ¶
IsUpgradeError checks if an error is an upgrade failure.
func KeepAlive ¶
KeepAlive sends periodic datagrams to keep the QUIC connection alive. This prevents NAT mappings from expiring (typically 20-30 seconds for UDP).
Usage:
server.Handle("/app", func(c *wt.Context) {
stop := wt.KeepAlive(c, 15*time.Second)
defer stop()
// ... handle session
})
Returns a stop function that cancels the keep-alive.
func ListenAndServeWithGracefulShutdown ¶
ListenAndServeWithGracefulShutdown starts the server and handles SIGTERM/SIGINT for graceful shutdown. Active sessions are drained within the given timeout.
Usage:
server := wt.New(...)
server.Handle("/app", handler)
wt.ListenAndServeWithGracefulShutdown(server, 30*time.Second)
func Messages ¶
Messages returns an iterator over length-prefixed messages from a stream.
for msg := range wt.Messages(stream) {
process(msg)
}
func PProfMux ¶
PProfMux returns an http.ServeMux with pprof endpoints registered. Serve this on a separate port for profiling in production.
Usage:
go http.ListenAndServe(":6060", wt.PProfMux())
// Then: go tool pprof http://localhost:6060/debug/pprof/profile?seconds=30
func Pipe ¶
Pipe bidirectionally copies data between two streams. Useful for proxying one WebTransport stream to another. Returns when either stream closes or encounters an error.
func PipeRaw ¶
func PipeRaw(rw io.ReadWriteCloser, s *Stream) error
PipeRaw bidirectionally copies between an io.ReadWriteCloser and a Stream.
func RequiredFields ¶
RequiredFields checks that the given struct fields are non-zero. Useful for simple message validation without a full validation library.
Usage:
type ChatMsg struct {
User string
Text string
}
func (m ChatMsg) Validate() error {
return wt.RequiredFields(m, "User", "Text")
}
func Retry ¶
func Retry(ctx context.Context, cfg RetryConfig, fn func() error) error
Retry executes fn up to MaxAttempts times with exponential backoff. Returns the last error if all attempts fail.
func ServerInfo ¶
ServerInfo returns information about the framework and runtime.
func SetAltSvcHeader ¶
func SetAltSvcHeader(w http.ResponseWriter, port int)
SetAltSvcHeader sets the Alt-Svc header on an HTTP response.
func Streams ¶
Streams returns an iterator over incoming streams. Use with Go 1.23+ range-over-func:
for stream := range wt.Streams(c) {
go handleStream(stream)
}
func ValidateDatagramSize ¶
ValidateDatagramSize checks if a datagram payload is within safe limits. Returns an error if the payload is too large.
func ValidateMessage ¶
ValidateMessage checks if a decoded message implements Validator and validates it. Returns nil if the message doesn't implement Validator.
Types ¶
type BackpressureWriter ¶
type BackpressureWriter struct {
// contains filtered or unexported fields
}
BackpressureWriter wraps stream writes with buffering and backpressure tracking. Use this when you need to detect slow consumers without blocking the sender goroutine.
func NewBackpressureWriter ¶
func NewBackpressureWriter(s *Stream, bufferSize int) *BackpressureWriter
NewBackpressureWriter creates a writer with the given buffer size. Messages are dropped (not queued indefinitely) when the buffer is full, preventing slow consumers from causing memory leaks.
func (*BackpressureWriter) BufferUsage ¶
func (bw *BackpressureWriter) BufferUsage() float64
BufferUsage returns the current buffer utilization as a fraction (0.0 to 1.0).
func (*BackpressureWriter) Close ¶
func (bw *BackpressureWriter) Close()
Close stops the writer and drains remaining messages.
func (*BackpressureWriter) IsFull ¶
func (bw *BackpressureWriter) IsFull() bool
IsFull returns true if the send buffer is at capacity.
func (*BackpressureWriter) Send ¶
func (bw *BackpressureWriter) Send(msg []byte) bool
Send attempts to queue a message for sending. Returns true if queued, false if the buffer is full (message dropped).
func (*BackpressureWriter) Stats ¶
func (bw *BackpressureWriter) Stats() (sent, dropped uint64)
Stats returns the number of messages sent and dropped.
type BatcherOption ¶
type BatcherOption func(*DatagramBatcher)
BatcherOption configures the DatagramBatcher.
func WithBatchEncoder ¶
func WithBatchEncoder(fn func(batch [][]byte) []byte) BatcherOption
WithBatchEncoder sets a custom function to encode a batch into a single datagram. Default: concatenates with 2-byte length prefix per item.
func WithBatchInterval ¶
func WithBatchInterval(d time.Duration) BatcherOption
WithBatchInterval sets the maximum time to wait before flushing (default: 16ms ≈ 60Hz).
func WithBatchSize ¶
func WithBatchSize(n int) BatcherOption
WithBatchSize sets the maximum number of datagrams per batch (default: 10).
type BufferedReader ¶
type BufferedReader struct {
*Stream
// contains filtered or unexported fields
}
BufferedReader wraps a Stream with a larger read buffer.
func NewBufferedReader ¶
func NewBufferedReader(s *Stream, bufSize int) *BufferedReader
NewBufferedReader creates a stream reader with a custom buffer size.
func (*BufferedReader) ReadBuffered ¶
func (br *BufferedReader) ReadBuffered() ([]byte, error)
ReadBuffered reads into the internal buffer and returns available bytes. More efficient than multiple small Read calls.
type CertRotator ¶
type CertRotator struct {
// contains filtered or unexported fields
}
CertRotator watches TLS certificate files and reloads them without server restart. QUIC connections established before rotation continue using the old cert. New connections use the new cert.
Usage:
rotator := wt.NewCertRotator("cert.pem", "key.pem",
wt.WithRotationInterval(1*time.Hour),
)
server := wt.New(
wt.WithAddr(":443"),
wt.WithCertRotator(rotator),
)
func NewCertRotator ¶
func NewCertRotator(certFile, keyFile string, opts ...RotatorOption) (*CertRotator, error)
NewCertRotator creates a certificate rotator.
func (*CertRotator) GetCertificate ¶
func (cr *CertRotator) GetCertificate(*tls.ClientHelloInfo) (*tls.Certificate, error)
GetCertificate returns the current certificate. Implements tls.Config.GetCertificate.
func (*CertRotator) TLSConfig ¶
func (cr *CertRotator) TLSConfig() *tls.Config
TLSConfig returns a tls.Config that uses the rotator for certificates.
type CompressedStream ¶
type CompressedStream struct {
*Stream
// contains filtered or unexported fields
}
CompressedStream wraps a Stream with per-message gzip compression. Messages are compressed before writing and decompressed after reading. Only useful for messages > ~100 bytes; smaller messages may increase in size.
func NewCompressedStream ¶
func NewCompressedStream(s *Stream, threshold int) *CompressedStream
NewCompressedStream wraps a stream with gzip compression. Messages smaller than threshold bytes are sent uncompressed. A 1-byte header indicates compression: 0x00 = raw, 0x01 = gzip.
func (*CompressedStream) ReadMessage ¶
func (cs *CompressedStream) ReadMessage() ([]byte, error)
ReadMessage reads and decompresses a message.
func (*CompressedStream) WriteMessage ¶
func (cs *CompressedStream) WriteMessage(data []byte) error
WriteMessage compresses and writes a message. Small messages are sent raw with a 0x00 prefix. Large messages are gzip-compressed with a 0x01 prefix.
type CompressionStats ¶
type CompressionStats struct {
RawBytes int64
CompressedBytes int64
MessagesRaw int64
MessagesGzip int64
}
CompressionStats tracks compression effectiveness.
func (CompressionStats) Ratio ¶
func (cs CompressionStats) Ratio() float64
Ratio returns the compression ratio (0.0 = perfect, 1.0 = no compression).
type ConnInfo ¶
type ConnInfo struct {
SessionID string `json:"session_id"`
RemoteAddr string `json:"remote_addr"`
LocalAddr string `json:"local_addr"`
Path string `json:"path"`
Params map[string]string `json:"params,omitempty"`
ConnectedAt time.Time `json:"connected_at"`
Transport string `json:"transport"` // "webtransport" or "websocket"
UserAgent string `json:"user_agent,omitempty"`
Origin string `json:"origin,omitempty"`
}
ConnInfo provides detailed connection information for a session.
type ConnectionError ¶
type ConnectionError struct {
Op string // "dial", "accept", "handshake"
Addr string
Wrapped error
}
ConnectionError represents a connection-level failure.
func (*ConnectionError) Error ¶
func (e *ConnectionError) Error() string
func (*ConnectionError) Unwrap ¶
func (e *ConnectionError) Unwrap() error
type Context ¶
type Context struct {
// contains filtered or unexported fields
}
Context wraps a WebTransport session with routing info, metadata, and helpers.
func (*Context) AcceptStream ¶
AcceptStream accepts the next incoming bidirectional stream.
func (*Context) AcceptUniStream ¶
func (c *Context) AcceptUniStream() (*ReceiveStream, error)
AcceptUniStream accepts the next incoming unidirectional stream (receive only).
func (*Context) CloseWithError ¶
CloseWithError closes the session with an error code and message.
func (*Context) OpenStream ¶
OpenStream opens a new bidirectional stream to the client.
func (*Context) OpenStreamSync ¶
OpenStreamSync opens a new bidirectional stream, blocking until flow control allows.
func (*Context) OpenUniStream ¶
func (c *Context) OpenUniStream() (*SendStream, error)
OpenUniStream opens a unidirectional stream to the client (send only).
func (*Context) Param ¶
Param returns a path parameter value by name. For pattern "/chat/{room}" and path "/chat/general", Param("room") returns "general".
func (*Context) ReceiveDatagram ¶
ReceiveDatagram receives the next datagram from the client.
func (*Context) ReceiveDatagramContext ¶
ReceiveDatagramContext receives a datagram with explicit context for cancellation/timeout.
func (*Context) RemoteAddr ¶
RemoteAddr returns the client's address.
func (*Context) Request ¶
Request returns the original HTTP request that initiated the WebTransport session.
func (*Context) SendDatagram ¶
SendDatagram sends an unreliable datagram to the client.
func (*Context) SendDatagramSafe ¶
SendDatagramSafe sends a datagram with size validation. Returns an error if the payload exceeds MaxDatagramSize.
func (*Context) Session ¶
func (c *Context) Session() *webtransport.Session
Session returns the underlying webtransport.Session.
type ContextStream ¶
type ContextStream struct {
*Stream
// contains filtered or unexported fields
}
ContextStream wraps a Stream with context-aware read/write operations. When the context is cancelled, all pending reads and writes are unblocked.
func (*ContextStream) Close ¶
func (cs *ContextStream) Close() error
Close cancels the context and closes the stream.
func (*ContextStream) Context ¶
func (cs *ContextStream) Context() context.Context
Context returns the stream's context.
func (*ContextStream) ReadMessageContext ¶
func (cs *ContextStream) ReadMessageContext() ([]byte, error)
ReadMessageContext reads a length-prefixed message with context support. Returns context.Canceled or context.DeadlineExceeded if the context is done.
func (*ContextStream) WriteMessageContext ¶
func (cs *ContextStream) WriteMessageContext(data []byte) error
WriteMessageContext writes a length-prefixed message with context support.
type DatagramBatcher ¶
type DatagramBatcher struct {
// contains filtered or unexported fields
}
DatagramBatcher collects datagrams and sends them in batches. This reduces the number of individual QUIC packets, improving throughput at the cost of slightly increased latency.
Useful for high-frequency updates (game state, sensor readings) where batching many small messages into fewer packets improves efficiency.
func NewDatagramBatcher ¶
func NewDatagramBatcher(c *Context, opts ...BatcherOption) *DatagramBatcher
NewDatagramBatcher creates a batcher that flushes either when maxSize is reached or interval elapses, whichever comes first.
func (*DatagramBatcher) Add ¶
func (b *DatagramBatcher) Add(data []byte)
Add adds a datagram to the current batch. If the batch is full, it's flushed immediately.
func (*DatagramBatcher) Close ¶
func (b *DatagramBatcher) Close()
Close stops the batcher and flushes remaining data.
func (*DatagramBatcher) Flush ¶
func (b *DatagramBatcher) Flush()
Flush sends any buffered datagrams immediately.
type ErrorPageHandler ¶
ErrorPageHandler is called when a non-WebTransport HTTP request hits a WebTransport route. Customize this to return helpful error messages.
type ErrorResponse ¶
type ErrorResponse struct {
Error string `json:"error"`
Code int `json:"code"`
Message string `json:"message,omitempty"`
}
ErrorResponse is the JSON body returned when WebTransport upgrade fails.
type EventBus ¶
type EventBus struct {
// contains filtered or unexported fields
}
EventBus provides a publish/subscribe system for session lifecycle events.
func (*EventBus) Emit ¶
Emit publishes an event to all registered handlers. Handlers are called synchronously in registration order.
func (*EventBus) EmitAsync ¶
EmitAsync publishes an event to all registered handlers asynchronously.
func (*EventBus) On ¶
func (eb *EventBus) On(eventType EventType, handler EventHandler)
On registers a handler for the given event type.
type FlowControlMonitor ¶
type FlowControlMonitor struct {
StreamsOpened atomic.Int64
StreamsClosed atomic.Int64
DatagramsSent atomic.Int64
DatagramsRecvd atomic.Int64
BytesSent atomic.Int64
BytesReceived atomic.Int64
WriteBlocks atomic.Int64 // times a write was blocked by flow control
}
FlowControlMonitor tracks stream and datagram flow control metrics. Useful for monitoring backpressure and identifying slow consumers.
func NewFlowControlMonitor ¶
func NewFlowControlMonitor() *FlowControlMonitor
NewFlowControlMonitor creates a new monitor.
func (*FlowControlMonitor) Stats ¶
func (fc *FlowControlMonitor) Stats() FlowStats
Stats returns current metrics.
type FlowStats ¶
type FlowStats struct {
StreamsOpened int64 `json:"streams_opened"`
StreamsClosed int64 `json:"streams_closed"`
StreamsActive int64 `json:"streams_active"`
DatagramsSent int64 `json:"datagrams_sent"`
DatagramsRecvd int64 `json:"datagrams_received"`
BytesSent int64 `json:"bytes_sent"`
BytesReceived int64 `json:"bytes_received"`
WriteBlocks int64 `json:"write_blocks"`
}
FlowStats returns a snapshot of flow control metrics.
type Group ¶
type Group struct {
// contains filtered or unexported fields
}
Group is a collection of routes that share a path prefix and middleware.
func (*Group) Handle ¶
func (g *Group) Handle(pattern string, handler HandlerFunc, mw ...MiddlewareFunc)
Handle registers a handler in the group.
type HandlerFunc ¶
type HandlerFunc func(*Context)
HandlerFunc is the function signature for WebTransport session handlers.
func HandleBoth ¶
func HandleBoth(streamFn StreamHandler, datagramFn func([]byte, *Context) []byte) HandlerFunc
HandleBoth is a convenience for handlers that process both streams and datagrams. Datagrams are handled in a background goroutine, streams in the main loop.
Usage:
server.Handle("/game", wt.HandleBoth(
func(s *wt.Stream, c *wt.Context) {
// handle game event stream
},
func(data []byte, c *wt.Context) []byte {
// handle position datagram
return nil // no reply
},
))
func HandleDatagram ¶
func HandleDatagram(fn func(data []byte, c *Context) []byte) HandlerFunc
HandleDatagram is a convenience for handlers that only process datagrams. It loops receiving datagrams and calls the handler for each.
Usage:
server.Handle("/ping", wt.HandleDatagram(func(data []byte, c *wt.Context) []byte {
return data // echo
}))
Example ¶
package main
import (
"fmt"
"github.com/rarebek/wt"
)
func main() {
server := wt.New(wt.WithAddr(":4433"), wt.WithSelfSignedTLS())
// HandleDatagram auto-receives datagrams and echoes reply
server.Handle("/ping", wt.HandleDatagram(func(data []byte, c *wt.Context) []byte {
return append([]byte("pong:"), data...)
}))
fmt.Println("Datagram handler registered")
}
Output: Datagram handler registered
func HandleStream ¶
func HandleStream(fn StreamHandler) HandlerFunc
HandleStream is a convenience for handlers that process one stream at a time. It auto-accepts streams and calls the handler for each in a new goroutine.
Usage:
server.Handle("/echo", wt.HandleStream(func(s *wt.Stream, c *wt.Context) {
defer s.Close()
msg, _ := s.ReadMessage()
s.WriteMessage(msg)
}))
Example ¶
package main
import (
"fmt"
"github.com/rarebek/wt"
)
func main() {
server := wt.New(wt.WithAddr(":4433"), wt.WithSelfSignedTLS())
// HandleStream auto-accepts streams and calls handler for each
server.Handle("/echo", wt.HandleStream(func(s *wt.Stream, c *wt.Context) {
defer s.Close()
msg, _ := s.ReadMessage()
s.WriteMessage(msg)
}))
fmt.Println("Stream handler registered")
}
Output: Stream handler registered
type HealthCheck ¶
type HealthCheck struct {
// contains filtered or unexported fields
}
HealthCheck provides an HTTP health check endpoint that reports server status. Serve this alongside your WebTransport server on an HTTP port for load balancers and monitoring systems.
func NewHealthCheck ¶
func NewHealthCheck(s *Server) *HealthCheck
NewHealthCheck creates a health check handler for the given server.
func (*HealthCheck) Handler ¶
func (h *HealthCheck) Handler() http.Handler
Handler returns an http.Handler for health checks.
func (*HealthCheck) ServeHTTP ¶
func (h *HealthCheck) ServeHTTP(w http.ResponseWriter, r *http.Request)
ServeHTTP implements http.Handler for health checks.
type HealthResponse ¶
type HealthResponse struct {
Status string `json:"status"`
ActiveSessions int `json:"active_sessions"`
Uptime string `json:"uptime"`
Transport string `json:"transport"`
}
HealthResponse is the JSON response from the health check endpoint.
type InterceptedStream ¶
type InterceptedStream struct {
*Stream
// contains filtered or unexported fields
}
InterceptedStream wraps a Stream with read/write interceptors.
func Intercept ¶
func Intercept(s *Stream, opts ...InterceptorOption) *InterceptedStream
Intercept wraps a stream with the given interceptors.
Usage:
stream := wt.Intercept(rawStream,
wt.OnRead(func(data []byte) ([]byte, error) {
log.Printf("received %d bytes", len(data))
return data, nil
}),
wt.OnWrite(func(data []byte) ([]byte, error) {
log.Printf("sending %d bytes", len(data))
return data, nil
}),
)
func (*InterceptedStream) ReadMessage ¶
func (is *InterceptedStream) ReadMessage() ([]byte, error)
ReadMessage reads a message through the read interceptor.
func (*InterceptedStream) WriteMessage ¶
func (is *InterceptedStream) WriteMessage(data []byte) error
WriteMessage writes a message through the write interceptor.
type InterceptorOption ¶
type InterceptorOption func(*StreamInterceptor)
InterceptorOption configures a StreamInterceptor.
type KVSync ¶
type KVSync struct {
// contains filtered or unexported fields
}
KVSync provides a synchronized key-value store that can be shared between server and client via a stream. Useful for game state sync, config sync, or shared document state.
func (*KVSync) GetRaw ¶
func (kv *KVSync) GetRaw(key string) (json.RawMessage, bool)
GetRaw retrieves the raw JSON for a key.
func (*KVSync) OnChange ¶
func (kv *KVSync) OnChange(fn func(key string, value json.RawMessage))
OnChange sets a callback that fires when a key is updated.
type MessageError ¶
MessageError occurs during message read/write operations.
func (*MessageError) Error ¶
func (e *MessageError) Error() string
func (*MessageError) Unwrap ¶
func (e *MessageError) Unwrap() error
type MiddlewareFunc ¶
type MiddlewareFunc func(c *Context, next HandlerFunc)
MiddlewareFunc is the function signature for middleware. Call next(c) to pass control to the next middleware or final handler.
type MigrationEvent ¶
type MigrationEvent struct {
SessionID string
OldAddr net.Addr
NewAddr net.Addr
MigratedAt time.Time
}
MigrationEvent represents a connection migration (IP address change).
type MigrationWatcher ¶
type MigrationWatcher struct {
// contains filtered or unexported fields
}
MigrationWatcher monitors sessions for address changes. QUIC handles migration transparently, but this watcher lets you react to migrations (logging, analytics, security checks).
func NewMigrationWatcher ¶
func NewMigrationWatcher(store *SessionStore, onMigrate func(MigrationEvent)) *MigrationWatcher
NewMigrationWatcher creates a watcher that polls session addresses.
func (*MigrationWatcher) Stop ¶
func (mw *MigrationWatcher) Stop()
Stop stops the migration watcher.
type Option ¶
type Option func(*Server)
Option configures the Server.
func WithAutoCert ¶
WithAutoCert configures automatic TLS certificate management via Let's Encrypt. Certificates are automatically obtained and renewed.
Requirements:
- The server must be publicly accessible on port 443
- DNS must point to this server
- A cache directory stores certificates (e.g., "/var/cache/certs")
Note: ACME validation uses TLS-ALPN-01 challenge, which requires port 443 TCP. The WebTransport server itself runs on UDP, so you need both:
- TCP port 443 for ACME challenges (handled by autocert)
- UDP port 443 for QUIC/WebTransport (handled by the framework)
Usage:
server := wt.New(
wt.WithAddr(":443"),
wt.WithAutoCert("example.com", "/var/cache/certs"),
)
func WithAutoCertMulti ¶
WithAutoCertMulti is like WithAutoCert but supports multiple domains.
func WithCertRotator ¶
func WithCertRotator(cr *CertRotator) Option
WithCertRotator configures the server to use a CertRotator for TLS.
func WithCheckOrigin ¶
WithCheckOrigin sets a function to validate the request origin.
func WithIdleTimeout ¶
WithIdleTimeout sets the session idle timeout (default 30s).
func WithQUICConfig ¶
func WithQUICConfig(cfg QUICConfig) Option
WithQUICConfig sets QUIC-level transport options.
func WithSelfSignedTLS ¶
func WithSelfSignedTLS() Option
WithSelfSignedTLS generates a self-signed certificate for development. Returns the certificate hash that browsers need for serverCertificateHashes.
type PreflightResult ¶
PreflightResult holds the result of a preflight check.
type PresenceInfo ¶
type PresenceInfo struct {
UserID string `json:"user_id"`
SessionID string `json:"session_id"`
Status string `json:"status"` // "online", "idle", "away", "typing"
Metadata map[string]any `json:"metadata,omitempty"`
UpdatedAt time.Time `json:"updated_at"`
}
PresenceInfo represents a user's current state in a room.
type PresenceTracker ¶
type PresenceTracker struct {
// contains filtered or unexported fields
}
PresenceTracker tracks user presence across rooms. Integrates with RoomManager to automatically track join/leave and supports custom status updates.
func NewPresenceTracker ¶
func NewPresenceTracker() *PresenceTracker
NewPresenceTracker creates a new presence tracker.
func (*PresenceTracker) Count ¶
func (pt *PresenceTracker) Count(room string) int
Count returns the number of present sessions in a room.
func (*PresenceTracker) GetPresence ¶
func (pt *PresenceTracker) GetPresence(room string) []PresenceInfo
GetPresence returns all presence info for a room.
func (*PresenceTracker) GetPresenceJSON ¶
func (pt *PresenceTracker) GetPresenceJSON(room string) []byte
GetPresenceJSON returns presence info as a JSON byte slice.
func (*PresenceTracker) Join ¶
func (pt *PresenceTracker) Join(room string, c *Context)
Join records a session joining a room.
func (*PresenceTracker) Leave ¶
func (pt *PresenceTracker) Leave(room string, c *Context)
Leave records a session leaving a room.
func (*PresenceTracker) OnChange ¶
func (pt *PresenceTracker) OnChange(fn func(room string, info PresenceInfo, event string))
OnChange sets a callback for presence changes.
func (*PresenceTracker) SetMetadata ¶
func (pt *PresenceTracker) SetMetadata(room, sessionID string, metadata map[string]any)
SetMetadata sets custom metadata for a session's presence.
func (*PresenceTracker) UpdateStatus ¶
func (pt *PresenceTracker) UpdateStatus(room, sessionID, status string)
UpdateStatus updates a session's presence status in a room.
type Priority ¶
type Priority int
Priority represents a stream urgency level. Higher priority streams are serviced first when resources are constrained. Based on HTTP/3 priority signaling (RFC 9218).
const ( // PriorityBackground is for non-urgent data (analytics, telemetry). PriorityBackground Priority = 0 // PriorityLow is for bulk transfers (file uploads, log streaming). PriorityLow Priority = 1 // PriorityNormal is the default priority. PriorityNormal Priority = 3 // PriorityHigh is for interactive data (chat messages, game events). PriorityHigh Priority = 5 // PriorityCritical is for control messages (auth, heartbeat, disconnect). PriorityCritical Priority = 7 )
type PubSub ¶
type PubSub struct {
// contains filtered or unexported fields
}
PubSub provides topic-based publish/subscribe messaging. Sessions subscribe to topics and receive datagrams published to those topics. More granular than rooms — a session can subscribe to any combination of topics.
func (*PubSub) PublishExcept ¶
PublishExcept sends to all subscribers except the given session.
func (*PubSub) SubscriberCount ¶
SubscriberCount returns the number of subscribers for a topic.
func (*PubSub) TopicsForSession ¶
TopicsForSession returns all topics a session is subscribed to.
func (*PubSub) Unsubscribe ¶
Unsubscribe removes a session from a topic.
func (*PubSub) UnsubscribeAll ¶
UnsubscribeAll removes a session from all topics.
type QUICConfig ¶
type QUICConfig struct {
// InitialStreamReceiveWindow is the initial flow control window for each stream (default: 512KB).
// Increase for high-throughput streams. quic-go auto-tunes up to MaxStreamReceiveWindow.
InitialStreamReceiveWindow uint64
// MaxStreamReceiveWindow is the maximum stream flow control window (default: 6MB).
MaxStreamReceiveWindow uint64
// InitialConnectionReceiveWindow is the initial connection-level flow control window (default: 768KB).
InitialConnectionReceiveWindow uint64
// MaxConnectionReceiveWindow is the max connection-level flow control window (default: 15MB).
MaxConnectionReceiveWindow uint64
// MaxIncomingStreams is the maximum number of concurrent incoming streams per session.
// Default: 100 in quic-go.
MaxIncomingStreams int64
// MaxIncomingUniStreams is the max number of concurrent incoming unidirectional streams.
MaxIncomingUniStreams int64
}
QUICConfig exposes QUIC-level tuning options. These are passed to the underlying quic-go transport.
func DefaultQUICConfig ¶
func DefaultQUICConfig() QUICConfig
DefaultQUICConfig returns sensible defaults for most applications.
func GameServerQUICConfig ¶
func GameServerQUICConfig() QUICConfig
GameServerQUICConfig returns config optimized for game servers: smaller windows (less buffering = lower latency), more streams.
func HighThroughputQUICConfig ¶
func HighThroughputQUICConfig() QUICConfig
HighThroughputQUICConfig returns config optimized for large transfers: bigger windows, fewer streams.
type RPCClient ¶
type RPCClient struct {
// contains filtered or unexported fields
}
RPCClient sends JSON-RPC requests over a WebTransport stream.
func NewRPCClient ¶
NewRPCClient wraps a stream for RPC calls.
type RPCHandler ¶
type RPCHandler func(params json.RawMessage) (any, error)
RPCHandler handles a single RPC method.
type RPCRequest ¶
type RPCRequest struct {
ID uint64 `json:"id"`
Method string `json:"method"`
Params json.RawMessage `json:"params,omitempty"`
}
RPCRequest represents a JSON-RPC-like request over a stream.
type RPCResponse ¶
type RPCResponse struct {
ID uint64 `json:"id"`
Result json.RawMessage `json:"result,omitempty"`
Error *RPCError `json:"error,omitempty"`
}
RPCResponse represents a JSON-RPC-like response.
type RPCServer ¶
type RPCServer struct {
// contains filtered or unexported fields
}
RPCServer handles JSON-RPC requests over a WebTransport stream.
func (*RPCServer) Register ¶
func (rpc *RPCServer) Register(method string, handler RPCHandler)
Register adds a handler for the given method name.
type ReceiveStream ¶
type ReceiveStream struct {
// contains filtered or unexported fields
}
ReceiveStream wraps a unidirectional receive stream.
func (*ReceiveStream) CancelRead ¶
func (s *ReceiveStream) CancelRead(code uint32)
CancelRead cancels reading.
func (*ReceiveStream) Read ¶
func (s *ReceiveStream) Read(b []byte) (int, error)
Read reads bytes from the receive stream.
func (*ReceiveStream) ReadMessage ¶
func (s *ReceiveStream) ReadMessage() ([]byte, error)
ReadMessage reads a length-prefixed message.
func (*ReceiveStream) SetReadDeadline ¶
func (s *ReceiveStream) SetReadDeadline(t time.Time) error
SetReadDeadline sets the read deadline.
type ReliableDatagram ¶
type ReliableDatagram struct {
// contains filtered or unexported fields
}
ReliableDatagram adds optional reliability on top of unreliable datagrams. It uses sequence numbers and acknowledgments to detect and retransmit lost messages.
This is useful when you want datagram-like semantics (no head-of-line blocking, independent from streams) but need delivery guarantees.
Trade-off: adds 6 bytes of overhead per datagram (2 byte seq + 4 byte timestamp) and introduces retransmission latency for lost messages.
func NewReliableDatagram ¶
func NewReliableDatagram(c *Context, onReceive func(data []byte), opts ...ReliableOption) *ReliableDatagram
NewReliableDatagram creates a reliable datagram layer over a session's datagrams. The onReceive callback is called for each reliably-delivered message.
func (*ReliableDatagram) PendingCount ¶
func (rd *ReliableDatagram) PendingCount() int
PendingCount returns the number of unacknowledged messages.
func (*ReliableDatagram) Send ¶
func (rd *ReliableDatagram) Send(data []byte) error
Send sends a datagram with reliability guarantees.
type ReliableOption ¶
type ReliableOption func(*ReliableDatagram)
ReliableOption configures the ReliableDatagram.
func WithMaxRetries ¶
func WithMaxRetries(n int) ReliableOption
WithMaxRetries sets the maximum retransmission attempts (default: 5).
func WithRetryTimeout ¶
func WithRetryTimeout(d time.Duration) ReliableOption
WithRetryTimeout sets the retransmission timeout (default: 100ms).
type ResumeStore ¶
type ResumeStore struct {
// contains filtered or unexported fields
}
ResumeStore persists session state for reconnection. When a client disconnects and reconnects with a resume token, the framework can restore their context store values (user info, etc.) without requiring re-authentication.
func NewResumeStore ¶
func NewResumeStore(ttl time.Duration) *ResumeStore
NewResumeStore creates a store with the given TTL for saved states. States are automatically expired after TTL.
func (*ResumeStore) Count ¶
func (rs *ResumeStore) Count() int
Count returns the number of stored resume states.
func (*ResumeStore) Restore ¶
func (rs *ResumeStore) Restore(c *Context, token ResumeToken) bool
Restore applies saved state to a new session context. Returns true if the token was valid and state was restored.
func (*ResumeStore) Save ¶
func (rs *ResumeStore) Save(c *Context) ResumeToken
Save stores the session's context values and returns a resume token. The client should store this token and present it on reconnection.
type ResumeToken ¶
type ResumeToken string
ResumeToken is an opaque token that clients can use to restore session state after a reconnect.
type RetryConfig ¶
type RetryConfig struct {
MaxAttempts int
InitDelay time.Duration
MaxDelay time.Duration
Jitter bool // Add random jitter to delays
}
RetryConfig configures retry behavior for stream operations.
func DefaultRetryConfig ¶
func DefaultRetryConfig() RetryConfig
DefaultRetryConfig returns sensible retry defaults.
type RingBuffer ¶
type RingBuffer[T any] struct { // contains filtered or unexported fields }
RingBuffer is a fixed-size, lock-free ring buffer for messages. When full, the oldest message is overwritten. Useful for storing recent messages (chat history, event log).
func NewRingBuffer ¶
func NewRingBuffer[T any](capacity int) *RingBuffer[T]
NewRingBuffer creates a ring buffer with the given capacity.
func (*RingBuffer[T]) Items ¶
func (rb *RingBuffer[T]) Items() []T
Items returns all items in order (oldest first).
func (*RingBuffer[T]) Last ¶
func (rb *RingBuffer[T]) Last() (T, bool)
Last returns the most recently added item.
func (*RingBuffer[T]) Len ¶
func (rb *RingBuffer[T]) Len() int
Len returns the number of items in the buffer.
func (*RingBuffer[T]) Push ¶
func (rb *RingBuffer[T]) Push(item T)
Push adds an item to the buffer. Overwrites oldest if full.
type Room ¶
type Room struct {
// contains filtered or unexported fields
}
Room represents a named group of sessions for pub/sub messaging.
func (*Room) BroadcastExcept ¶
BroadcastExcept sends a datagram to all members except the specified session.
func (*Room) BroadcastStream ¶
BroadcastStream sends a message to all room members via reliable streams. Unlike Broadcast (datagrams, unreliable), this guarantees delivery. Each member gets a new stream with the message.
func (*Room) BroadcastStreamExcept ¶
BroadcastStreamExcept sends a reliable message to all except the given session.
func (*Room) ForEach ¶
ForEach iterates over all members without allocating a slice. The callback runs under a read lock — do not block for long.
func (*Room) SafeBroadcast ¶
SafeBroadcast sends a datagram to all room members, recovering from panics. Logs and skips any member that causes a panic (e.g., closed connection).
type RoomManager ¶
type RoomManager struct {
// contains filtered or unexported fields
}
RoomManager manages named rooms.
func NewRoomManager ¶
func NewRoomManager() *RoomManager
NewRoomManager creates a new RoomManager.
Example ¶
package main
import (
"fmt"
"log"
"github.com/rarebek/wt"
)
func main() {
rooms := wt.NewRoomManager()
lobby := rooms.GetOrCreate("lobby")
lobby.OnJoin(func(c *wt.Context) {
log.Printf("user %s joined lobby", c.ID())
})
fmt.Println("Room manager created")
}
Output: Room manager created
func (*RoomManager) Get ¶
func (rm *RoomManager) Get(name string) (*Room, bool)
Get returns a room by name if it exists.
func (*RoomManager) GetOrCreate ¶
func (rm *RoomManager) GetOrCreate(name string) *Room
GetOrCreate returns a room by name, creating it if it doesn't exist.
type RoomMessage ¶
type RoomMessage struct {
SenderID string `json:"sender_id"`
Data []byte `json:"data"`
Timestamp time.Time `json:"timestamp"`
}
RoomMessage represents a message stored in room history.
type RoomWithHistory ¶
type RoomWithHistory struct {
*Room
// contains filtered or unexported fields
}
RoomWithHistory wraps a Room with message history support. New members can receive recent messages when they join.
func NewRoomWithHistory ¶
func NewRoomWithHistory(room *Room, historySize int) *RoomWithHistory
NewRoomWithHistory creates a room wrapper with history of the given capacity.
func (*RoomWithHistory) BroadcastAndRecord ¶
func (r *RoomWithHistory) BroadcastAndRecord(senderID string, data []byte)
BroadcastAndRecord sends data to all room members AND records it in history.
func (*RoomWithHistory) BroadcastExceptAndRecord ¶
func (r *RoomWithHistory) BroadcastExceptAndRecord(senderID string, data []byte)
BroadcastExceptAndRecord sends to all except sender AND records in history.
func (*RoomWithHistory) ClearHistory ¶
func (r *RoomWithHistory) ClearHistory()
ClearHistory removes all stored messages.
func (*RoomWithHistory) History ¶
func (r *RoomWithHistory) History() []RoomMessage
History returns recent messages (oldest first).
func (*RoomWithHistory) HistorySize ¶
func (r *RoomWithHistory) HistorySize() int
HistorySize returns the number of messages in history.
func (*RoomWithHistory) ReplayHistory ¶
func (r *RoomWithHistory) ReplayHistory(c *Context)
ReplayHistory sends all history messages to a specific session via datagrams. Useful for sending catch-up data when a new member joins.
func (*RoomWithHistory) ReplayHistorySince ¶
func (r *RoomWithHistory) ReplayHistorySince(c *Context, since time.Time)
ReplayHistorySince sends messages newer than the given timestamp.
type RotatorOption ¶
type RotatorOption func(*CertRotator)
RotatorOption configures the CertRotator.
func WithRotationInterval ¶
func WithRotationInterval(d time.Duration) RotatorOption
WithRotationInterval sets how often to check for new certificates (default: 1 hour).
func WithRotationLogger ¶
func WithRotationLogger(logger *slog.Logger) RotatorOption
WithRotationLogger sets the logger for rotation events.
type Route ¶
type Route struct {
Pattern string
Handler HandlerFunc
Middleware []MiddlewareFunc
// contains filtered or unexported fields
}
Route represents a registered path pattern with its handler and middleware.
type Router ¶
type Router struct {
// contains filtered or unexported fields
}
Router handles path-based routing for WebTransport sessions.
func (*Router) Add ¶
func (r *Router) Add(pattern string, handler HandlerFunc, mw ...MiddlewareFunc)
Add registers a handler for the given path pattern. Patterns support parameters like "/chat/{room}" and "/game/{id}/input".
func (*Router) ExtractParams ¶
ExtractParams extracts path parameters from a URL path given a pattern.
type SendStream ¶
type SendStream struct {
// contains filtered or unexported fields
}
SendStream wraps a unidirectional send stream.
func (*SendStream) CancelWrite ¶
func (s *SendStream) CancelWrite(code uint32)
CancelWrite cancels writing.
func (*SendStream) SetWriteDeadline ¶
func (s *SendStream) SetWriteDeadline(t time.Time) error
SetWriteDeadline sets the write deadline.
func (*SendStream) Write ¶
func (s *SendStream) Write(b []byte) (int, error)
Write writes bytes to the send stream.
func (*SendStream) WriteMessage ¶
func (s *SendStream) WriteMessage(data []byte) error
WriteMessage writes a length-prefixed message.
type Server ¶
type Server struct {
// contains filtered or unexported fields
}
Server is the main WebTransport framework server.
func New ¶
New creates a new WebTransport server with the given options.
Example ¶
package main
import (
"fmt"
"io"
"github.com/rarebek/wt"
)
func main() {
server := wt.New(
wt.WithAddr(":4433"),
wt.WithSelfSignedTLS(),
)
server.Handle("/echo", func(c *wt.Context) {
for {
stream, err := c.AcceptStream()
if err != nil {
return
}
go func() {
defer stream.Close()
io.Copy(stream, stream)
}()
}
})
fmt.Println("Server configured on", server.Addr())
}
Output: Server configured on :4433
func (*Server) BroadcastExcept ¶
BroadcastExcept sends a datagram to all active sessions except the specified one.
func (*Server) CertHash ¶
CertHash returns the SHA-256 hash of the self-signed certificate, needed for the browser's serverCertificateHashes option. Returns empty string if not using self-signed TLS.
func (*Server) Group ¶
func (s *Server) Group(prefix string, mw ...MiddlewareFunc) *Group
Group creates a route group with shared middleware.
Example ¶
package main
import (
"fmt"
"github.com/rarebek/wt"
)
func main() {
server := wt.New(wt.WithAddr(":4433"), wt.WithSelfSignedTLS())
// Create a group with auth middleware
api := server.Group("/api", func(c *wt.Context, next wt.HandlerFunc) {
// Check auth header
if c.Request().Header.Get("Authorization") == "" {
c.CloseWithError(401, "unauthorized")
return
}
next(c)
})
api.Handle("/data", func(c *wt.Context) {
// Only authenticated users reach here
_ = c
})
fmt.Println("Group registered")
}
Output: Group registered
func (*Server) Handle ¶
func (s *Server) Handle(pattern string, handler HandlerFunc, mw ...MiddlewareFunc)
Handle registers a session handler for the given path pattern. The handler receives a Context with the full session — accept streams, receive datagrams, access path params, etc.
Example ¶
package main
import (
"fmt"
"github.com/rarebek/wt"
)
func main() {
server := wt.New(wt.WithAddr(":4433"), wt.WithSelfSignedTLS())
// Simple echo handler
server.Handle("/echo", func(c *wt.Context) {
for {
stream, err := c.AcceptStream()
if err != nil {
return
}
go func() {
defer stream.Close()
msg, _ := stream.ReadMessage()
stream.WriteMessage(msg)
}()
}
})
fmt.Println("Handler registered")
}
Output: Handler registered
func (*Server) ListenAndServe ¶
ListenAndServe starts the WebTransport server.
func (*Server) Multicast ¶
Multicast sends a datagram to sessions matching a filter function. More flexible than Broadcast — only sends to sessions that match.
Usage:
// Send to all sessions with user role "admin"
server.Multicast(data, func(c *Context) bool {
role, _ := c.Get("role")
return role == "admin"
})
func (*Server) MulticastStream ¶
MulticastStream sends a reliable message via streams to matching sessions.
func (*Server) OnDisconnect ¶
OnDisconnect registers a callback for closed sessions.
func (*Server) OnShutdown ¶
func (s *Server) OnShutdown(fn ShutdownHook)
OnShutdown registers a function to be called when the server shuts down. Hooks run in registration order before connections are drained.
func (*Server) Preflight ¶
PreflightCheck verifies the server configuration before starting. Returns a list of issues found. Empty list = ready to start.
Usage:
server := wt.New(...)
if issues := server.Preflight(); len(issues) > 0 {
for _, issue := range issues {
log.Printf("WARN: %s", issue)
}
}
func (*Server) PreflightCheck ¶
func (s *Server) PreflightCheck() PreflightResult
PreflightCheck runs the preflight check and returns a structured result.
func (*Server) SessionCount ¶
SessionCount returns the number of active sessions.
func (*Server) Sessions ¶
func (s *Server) Sessions() *SessionStore
Sessions returns the session store for looking up active sessions.
func (*Server) Shutdown ¶
Shutdown gracefully shuts down the server. It stops accepting new connections, then waits for active sessions to finish or for the context to be cancelled, whichever comes first.
func (*Server) Use ¶
func (s *Server) Use(mw ...MiddlewareFunc)
Use adds global middleware that runs on every session.
type SessionCloseError ¶
SessionCloseError represents a session closure with a code and message.
func (*SessionCloseError) Error ¶
func (e *SessionCloseError) Error() string
type SessionStore ¶
type SessionStore struct {
// contains filtered or unexported fields
}
SessionStore tracks active sessions and provides lookup/broadcast capabilities.
func NewSessionStore ¶
func NewSessionStore() *SessionStore
NewSessionStore creates a new SessionStore.
func (*SessionStore) Broadcast ¶
func (ss *SessionStore) Broadcast(data []byte)
Broadcast sends a datagram to all active sessions.
func (*SessionStore) CloseAll ¶
func (ss *SessionStore) CloseAll()
CloseAll closes all active sessions.
func (*SessionStore) Count ¶
func (ss *SessionStore) Count() int
Count returns the number of active sessions.
func (*SessionStore) Each ¶
func (ss *SessionStore) Each(fn func(*Context))
Each iterates over all active sessions. The callback should not block for long.
func (*SessionStore) FindByValue ¶
func (ss *SessionStore) FindByValue(key string, value any) []*Context
FindByValue returns all sessions where the given key matches the given value. Useful for finding all sessions for a specific user, role, etc.
func (*SessionStore) Get ¶
func (ss *SessionStore) Get(id string) (*Context, bool)
Get returns a session by ID.
func (*SessionStore) IDs ¶
func (ss *SessionStore) IDs() []string
IDs returns all active session IDs.
func (*SessionStore) Remove ¶
func (ss *SessionStore) Remove(id string)
Remove unregisters a session.
type ShutdownHook ¶
type ShutdownHook func()
ShutdownHook is a function called during server shutdown.
type Stream ¶
type Stream struct {
// contains filtered or unexported fields
}
Stream wraps a bidirectional WebTransport stream with framing and helpers.
func OpenTypedStream ¶
OpenTypedStream opens a stream with the given type header. The remote end's StreamMux will route it to the matching handler.
func RetryStream ¶
RetryStream attempts to open a stream with retry.
func (*Stream) CancelRead ¶
CancelRead cancels the read side of the stream.
func (*Stream) CancelWrite ¶
CancelWrite cancels the write side of the stream.
func (*Stream) Raw ¶
func (s *Stream) Raw() *webtransport.Stream
Raw returns the underlying webtransport.Stream.
func (*Stream) ReadMessage ¶
ReadMessage reads a length-prefixed message from the stream.
func (*Stream) SessionContext ¶
Context returns the session context this stream belongs to.
func (*Stream) SetDeadline ¶
SetDeadline sets read and write deadlines.
func (*Stream) SetReadDeadline ¶
SetReadDeadline sets the read deadline.
func (*Stream) SetWriteDeadline ¶
SetWriteDeadline sets the write deadline.
func (*Stream) WithContext ¶
func (s *Stream) WithContext(ctx context.Context) *ContextStream
WithContext creates a ContextStream that respects the given context. When ctx is cancelled, the stream is closed automatically.
func (*Stream) WithDeadline ¶
func (s *Stream) WithDeadline(deadline time.Time) *ContextStream
WithDeadline creates a stream that automatically closes at the given deadline.
func (*Stream) WithTimeout ¶
func (s *Stream) WithTimeout(d time.Duration) *ContextStream
WithTimeout creates a stream that automatically closes after the given duration.
func (*Stream) WriteMessage ¶
WriteMessage writes a length-prefixed message to the stream. Format: [4 bytes big-endian length][payload]
type StreamCloseError ¶
StreamCloseError represents a stream closure with an error code.
func (*StreamCloseError) Error ¶
func (e *StreamCloseError) Error() string
type StreamConfig ¶
type StreamConfig struct {
// Priority hint for this stream (not enforced by QUIC, but useful for
// application-level scheduling).
Priority Priority
// TypeID is the StreamMux type identifier (0 = no mux).
TypeID uint16
}
StreamConfig holds configuration for a new stream.
func DefaultStreamConfig ¶
func DefaultStreamConfig() StreamConfig
DefaultStreamConfig returns default stream configuration.
type StreamHandler ¶
StreamHandler handles a single stream.
type StreamInterceptor ¶
type StreamInterceptor struct {
// contains filtered or unexported fields
}
StreamInterceptor intercepts stream message reads and writes. Useful for logging, metrics, validation, or transformation of messages.
type StreamMux ¶
type StreamMux struct {
// contains filtered or unexported fields
}
StreamMux multiplexes different stream types within a single session. Each stream's first 2 bytes identify its type, and the mux routes it to the appropriate handler.
This solves the problem of "I have one WebTransport session but I need different handlers for chat streams vs game state streams vs file uploads."
Usage:
mux := wt.NewStreamMux()
mux.Handle(1, handleChat) // type=1 → chat handler
mux.Handle(2, handleGameInput) // type=2 → game input handler
mux.Handle(3, handleFileUpload) // type=3 → file upload handler
server.Handle("/app", func(c *wt.Context) {
mux.Serve(c) // auto-routes streams by type
})
func NewStreamMux ¶
func NewStreamMux() *StreamMux
NewStreamMux creates a new stream multiplexer.
Example ¶
package main
import (
"fmt"
"github.com/rarebek/wt"
)
func main() {
mux := wt.NewStreamMux()
const (
TypeChat uint16 = 1
TypeGame uint16 = 2
)
mux.Handle(TypeChat, func(s *wt.Stream, c *wt.Context) {
// Handle chat stream
})
mux.Handle(TypeGame, func(s *wt.Stream, c *wt.Context) {
// Handle game stream
})
fmt.Println("StreamMux configured with 2 handlers")
}
Output: StreamMux configured with 2 handlers
func (*StreamMux) Fallback ¶
func (m *StreamMux) Fallback(handler StreamHandler)
Fallback sets a handler for unrecognized stream types.
func (*StreamMux) Handle ¶
func (m *StreamMux) Handle(typeID uint16, handler StreamHandler)
Handle registers a handler for the given stream type ID.
type StreamOptions ¶
type StreamOptions struct {
// ReadBufferSize sets the size of the read buffer (default: 4096).
ReadBufferSize int
// WriteBufferSize sets the size of the write buffer (default: 4096).
WriteBufferSize int
// MaxMessageSize overrides the default maximum message size for this stream.
MaxMessageSize int
}
StreamOptions configures stream behavior.
func DefaultStreamOptions ¶
func DefaultStreamOptions() StreamOptions
DefaultStreamOptions returns default stream configuration.
type Tags ¶
type Tags struct {
// contains filtered or unexported fields
}
Tags provides a thread-safe tagging system for sessions. Tags can be used to categorize sessions for targeted operations like "send to all sessions tagged 'premium'" or "disconnect all 'guest' sessions".
func (*Tags) SessionsWithTag ¶
SessionsWithTag returns all session IDs with the given tag.
func (*Tags) TagsForSession ¶
TagsForSession returns all tags for a session.
type Ticker ¶
type Ticker struct {
// contains filtered or unexported fields
}
Ticker sends periodic datagrams to a session at a fixed interval. Useful for game state updates, heartbeats, or periodic data pushes.
Usage:
server.Handle("/game", func(c *wt.Context) {
ticker := wt.NewTicker(c, 50*time.Millisecond, func() []byte {
return getGameState()
})
defer ticker.Stop()
// ... handle other streams
})
type TypedDatagram ¶
type TypedDatagram[T any] struct { // contains filtered or unexported fields }
TypedDatagram provides type-safe datagram read/write on a Context.
func NewTypedDatagram ¶
func NewTypedDatagram[T any](ctx *Context, c codec.Codec) *TypedDatagram[T]
NewTypedDatagram wraps a Context with typed datagram encoding/decoding.
func (*TypedDatagram[T]) Receive ¶
func (td *TypedDatagram[T]) Receive() (T, error)
Receive receives and decodes a datagram.
func (*TypedDatagram[T]) Send ¶
func (td *TypedDatagram[T]) Send(v T) error
Send encodes and sends a datagram.
type TypedPubSub ¶
type TypedPubSub[T any] struct { // contains filtered or unexported fields }
TypedPubSub provides type-safe publish/subscribe with automatic encoding.
func NewTypedPubSub ¶
func NewTypedPubSub[T any](ps *PubSub, c codec.Codec) *TypedPubSub[T]
NewTypedPubSub wraps a PubSub with typed messages.
func (*TypedPubSub[T]) Publish ¶
func (tp *TypedPubSub[T]) Publish(topic string, msg T) error
Publish encodes and publishes a message to all subscribers of a topic.
func (*TypedPubSub[T]) PublishExcept ¶
func (tp *TypedPubSub[T]) PublishExcept(topic string, msg T, excludeID string) error
PublishExcept encodes and publishes to all except one session.
func (*TypedPubSub[T]) Subscribe ¶
func (tp *TypedPubSub[T]) Subscribe(topic string, c *Context)
Subscribe delegates to the underlying PubSub.
func (*TypedPubSub[T]) Unsubscribe ¶
func (tp *TypedPubSub[T]) Unsubscribe(topic string, c *Context)
Unsubscribe delegates to the underlying PubSub.
func (*TypedPubSub[T]) UnsubscribeAll ¶
func (tp *TypedPubSub[T]) UnsubscribeAll(c *Context)
UnsubscribeAll delegates to the underlying PubSub.
type TypedRoom ¶
type TypedRoom[T any] struct { // contains filtered or unexported fields }
TypedRoom provides type-safe broadcast over a Room using a codec.
func NewTypedRoom ¶
NewTypedRoom wraps a Room with typed broadcast support.
func (*TypedRoom[T]) Broadcast ¶
Broadcast encodes and broadcasts a value to all room members via datagrams.
func (*TypedRoom[T]) BroadcastExcept ¶
BroadcastExcept broadcasts to all members except the specified session.
type TypedStream ¶
TypedStream provides type-safe read/write over a Stream using a codec.
func NewTypedStream ¶
NewTypedStream wraps a Stream with typed encoding/decoding.
Example ¶
package main
import (
"fmt"
"github.com/rarebek/wt/codec"
)
func main() {
type Input struct {
Action string `json:"action"`
}
type Output struct {
Result string `json:"result"`
}
// In a real handler:
// stream, _ := c.AcceptStream()
// typed := wt.NewTypedStream[Input, Output](stream, codec.JSON{})
// input, _ := typed.Read()
// typed.Write(Output{Result: "ok"})
_ = codec.JSON{}
fmt.Println("TypedStream example")
}
Output: TypedStream example
func (*TypedStream[R, W]) Close ¶
func (ts *TypedStream[R, W]) Close() error
Close closes the underlying stream.
func (*TypedStream[R, W]) Read ¶
func (ts *TypedStream[R, W]) Read() (R, error)
Read reads and decodes a message from the stream.
func (*TypedStream[R, W]) Stream ¶
func (ts *TypedStream[R, W]) Stream() *Stream
Stream returns the underlying Stream.
func (*TypedStream[R, W]) Write ¶
func (ts *TypedStream[R, W]) Write(v W) error
Write encodes and writes a message to the stream.
type UpgradeError ¶
UpgradeError occurs when WebTransport upgrade fails.
func (*UpgradeError) Error ¶
func (e *UpgradeError) Error() string
Source Files
¶
- altsvc.go
- autocert.go
- backpressure.go
- batch.go
- certrotation.go
- compressed_stream.go
- conninfo.go
- context.go
- convenience.go
- ctxstream.go
- datagram.go
- doc.go
- errorpage.go
- errors.go
- events.go
- flowcontrol.go
- health.go
- interceptor.go
- iterator.go
- keepalive.go
- kv.go
- middleware.go
- migration.go
- migration_hooks.go
- multicast.go
- mux.go
- pipe.go
- pprof.go
- preflight.go
- presence.go
- priority.go
- pubsub.go
- quicconfig.go
- reliable_datagram.go
- resume.go
- retry.go
- ringbuf.go
- room_broadcast.go
- room_history.go
- router.go
- rpc.go
- safe_room.go
- session.go
- shutdown_hook.go
- signal.go
- stream.go
- streamconfig.go
- tags.go
- ticker.go
- typed.go
- typed_pubsub.go
- typed_room.go
- util.go
- validate.go
- version.go
- wt.go
Directories
¶
| Path | Synopsis |
|---|---|
|
Package client provides a WebTransport client with reconnection support.
|
Package client provides a WebTransport client with reconnection support. |
|
Package codec provides message encoding/decoding for the wt framework.
|
Package codec provides message encoding/decoding for the wt framework. |
|
examples
|
|
|
aistream
command
Example: AI token streaming server.
|
Example: AI token streaming server. |
|
chat
command
Example: WebTransport chat server with rooms.
|
Example: WebTransport chat server with rooms. |
|
collab
command
Example: Real-time collaboration server.
|
Example: Real-time collaboration server. |
|
dualprotocol
command
Example: Dual-protocol server (HTTP/2 + WebTransport).
|
Example: Dual-protocol server (HTTP/2 + WebTransport). |
|
echo
command
Example: minimal WebTransport echo server.
|
Example: minimal WebTransport echo server. |
|
fallback_demo
command
Example: Combined WebTransport + WebSocket + SSE fallback.
|
Example: Combined WebTransport + WebSocket + SSE fallback. |
|
gameserver
command
Example: WebTransport game server.
|
Example: WebTransport game server. |
|
iot
command
Example: IoT telemetry gateway.
|
Example: IoT telemetry gateway. |
|
minimal
command
The simplest possible wt server — 15 lines.
|
The simplest possible wt server — 15 lines. |
|
notification
command
Example: Server-push notification system.
|
Example: Server-push notification system. |
|
proxy
command
Example: WebTransport proxy.
|
Example: WebTransport proxy. |
|
SSE provides a Server-Sent Events fallback for environments where both WebTransport (UDP) and WebSocket (TCP upgrade) are blocked.
|
SSE provides a Server-Sent Events fallback for environments where both WebTransport (UDP) and WebSocket (TCP upgrade) are blocked. |
|
Package middleware provides built-in middleware for the wt WebTransport framework.
|
Package middleware provides built-in middleware for the wt WebTransport framework. |
|
Command wtbench is a load testing tool for WebTransport servers.
|
Command wtbench is a load testing tool for WebTransport servers. |