wt

package module
v0.1.0 Latest Latest
Warning

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

Go to latest
Published: Mar 28, 2026 License: MIT Imports: 37 Imported by: 0

README

wt

A high-level WebTransport framework for Go. Build real-time applications with multiplexed streams, unreliable datagrams, and zero boilerplate.

Built on top of quic-go/webtransport-go.

Why

WebTransport is the successor to WebSocket — faster, more flexible, built on QUIC/HTTP3. But the existing Go libraries give you raw protocol access. wt gives you a framework:

  • Routing — path-based handlers with parameters (/chat/{room})
  • Middleware — auth, logging, rate limiting, compression, metrics
  • Rooms — pub/sub groups with broadcast and presence
  • Typed streams — generics-based message encoding/decoding
  • WebSocket fallback — transparent fallback for browsers without WebTransport support
  • Self-signed certs — one-line dev setup, outputs the hash for browser serverCertificateHashes

Install

go get github.com/rarebek/wt

Requires Go 1.22+.

Quick Start

package main

import (
    "io"
    "log"

    "github.com/rarebek/wt"
    "github.com/rarebek/wt/middleware"
)

func main() {
    server := wt.New(
        wt.WithAddr(":4433"),
        wt.WithSelfSignedTLS(),
    )

    log.Printf("Cert hash: %s", server.CertHash())

    server.Use(middleware.DefaultLogger())
    server.Use(middleware.Recover(nil))

    // Echo server: streams echo back, datagrams echo back
    server.Handle("/echo", func(c *wt.Context) {
        // Echo datagrams
        go func() {
            for {
                data, err := c.ReceiveDatagram()
                if err != nil {
                    return
                }
                c.SendDatagram(data)
            }
        }()

        // Echo streams
        for {
            stream, err := c.AcceptStream()
            if err != nil {
                return
            }
            go func() {
                defer stream.Close()
                io.Copy(stream, stream)
            }()
        }
    })

    log.Fatal(server.ListenAndServe())
}

Connect from the browser:

const transport = new WebTransport("https://localhost:4433/echo", {
    serverCertificateHashes: [{
        algorithm: "sha-256",
        value: hexToBytes("PASTE_CERT_HASH_HERE")
    }]
});
await transport.ready;

Features

Path Routing with Parameters
server.Handle("/game/{id}", func(c *wt.Context) {
    gameID := c.Param("id")
    // ...
})
Middleware
// Global
server.Use(middleware.DefaultLogger())
server.Use(middleware.Recover(nil))
server.Use(middleware.RateLimit(100))

// Per-group
admin := server.Group("/admin", middleware.BearerAuth(validateToken))
admin.Handle("/control", handleAdmin)

// Per-route
server.Handle("/public", handler, middleware.CORS(middleware.CORSConfig{
    AllowedOrigins: []string{"https://example.com"},
}))

Built-in middleware:

  • Logger — structured logging via slog
  • Recover — panic recovery
  • RateLimit — per-IP connection limiting
  • TokenBucket — per-session message rate limiting
  • BearerAuth / QueryAuth / RequireKey — authentication
  • MaxSessions — global session limit
  • CORS — origin validation
  • Timeout — session timeout
  • Compress — gzip/deflate compression
  • Metrics — session counting and duration tracking
Rooms (Pub/Sub)
rooms := wt.NewRoomManager()

server.Handle("/chat/{room}", func(c *wt.Context) {
    room := rooms.GetOrCreate(c.Param("room"))
    room.Join(c)
    defer room.Leave(c)

    room.Broadcast([]byte("someone joined"))
    room.BroadcastExcept([]byte("hello"), c.ID())
})
Typed Streams (Generics)
stream, _ := c.AcceptStream()
typed := wt.NewTypedStream[InputMsg, OutputMsg](stream, codec.JSON{})

input, err := typed.Read()   // auto-decoded InputMsg
typed.Write(OutputMsg{...})  // auto-encoded OutputMsg
Datagrams (Unreliable, Fast)
// Send — fire and forget, no head-of-line blocking
c.SendDatagram(positionData)

// Receive
data, err := c.ReceiveDatagram()
Message Framing
stream, _ := c.AcceptStream()

// Length-prefixed messages (4-byte header + payload)
stream.WriteMessage([]byte("hello"))
msg, _ := stream.ReadMessage()
WebSocket Fallback
import "github.com/rarebek/wt/fallback"

// Serve WebSocket fallback on a separate HTTP port
http.Handle("/ws", fallback.Handler(func(conn *fallback.WSConn) {
    // Same stream/datagram API as WebTransport
    stream, _ := conn.AcceptStream()
    data, _ := conn.ReceiveDatagram()
}))
Self-Signed Certificates
server := wt.New(wt.WithSelfSignedTLS())
fmt.Println("Browser cert hash:", server.CertHash())

Architecture

wt (core)
├── Server          — config, lifecycle, route registration
├── Router          — path matching with {param} extraction
├── Context         — session wrapper with params, store, helpers
├── Stream          — bidirectional stream with message framing
├── SendStream      — unidirectional send
├── ReceiveStream   — unidirectional receive
├── TypedStream     — generics-based typed read/write
├── SessionStore    — active session tracking
├── RoomManager     — named rooms with pub/sub
├── Room            — member management, broadcast
│
├── codec/          — message encoding (JSON, MsgPack)
├── middleware/      — auth, logging, rate limit, compress, metrics
├── client/         — Go client with auto-reconnect
└── fallback/       — WebSocket fallback with stream multiplexing

Performance

Benchmarks on Intel i9-14900K:

Operation ns/op allocs/op
Router match 130 2
Middleware chain (3 layers) 5.9 0
Context Get 10.5 0
Session store lookup 11 0
E2E echo (64B over QUIC) 53,580 26

Browser Support

Browser WebTransport Fallback
Chrome 97+ Yes -
Edge 97+ Yes -
Firefox 114+ Yes -
Safari Coming (Interop 2026) WebSocket

License

MIT

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

Examples

Constants

View Source
const (
	CodeOK              uint32 = 0   // Success / normal closure
	CodeUnauthorized    uint32 = 401 // Missing or invalid authentication
	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
	CodeServiceUnavail  uint32 = 503 // Server at capacity / circuit open

	// 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.

View Source
const DefaultKeepAliveInterval = 15 * time.Second

DefaultKeepAliveInterval is 15 seconds, chosen to be well under typical UDP NAT timeout of 20-30 seconds.

View Source
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.

View Source
const MaxMessageSize = 16 * 1024 * 1024

MaxMessageSize is the maximum size of a length-prefixed message (16 MB).

View Source
const Version = "0.1.0-dev"

Version is the current framework version.

Variables

This section is empty.

Functions

func AltSvcHeader

func AltSvcHeader(port int) string

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

func AltSvcMiddleware(port int) func(http.Handler) http.Handler

AltSvcMiddleware returns an HTTP middleware that adds the Alt-Svc header to every response, advertising HTTP/3 availability.

func CallTyped

func CallTyped[T any](c *RPCClient, method string, params any) (T, error)

CallTyped sends an RPC request and unmarshals the result into the given type.

func CertFingerprint

func CertFingerprint(certDER []byte) string

CertFingerprint returns the hex-encoded SHA-256 fingerprint of a DER-encoded certificate.

func Datagrams

func Datagrams(c *Context) iter.Seq[[]byte]

Datagrams returns an iterator over incoming datagrams.

for data := range wt.Datagrams(c) {
    process(data)
}

func DebugMux

func DebugMux(s *Server) *http.ServeMux

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

func DecodeBatch(data []byte) [][]byte

DecodeBatch decodes a batch-encoded datagram into individual messages. Pre-allocates the result slice based on estimated message count.

func DefaultErrorPage

func DefaultErrorPage(w http.ResponseWriter, _ *http.Request, code int, msg string)

DefaultErrorPage returns a JSON error response.

func HTMLErrorPage

func HTMLErrorPage(w http.ResponseWriter, _ *http.Request, code int, msg string)

HTMLErrorPage returns an HTML error page.

func Hash

func Hash(data []byte) string

Hash returns the SHA-256 hash of the given data as a hex string.

func IsConnectionError

func IsConnectionError(err error) bool

IsConnectionError checks if an error is a connection-level failure.

func IsMessageError

func IsMessageError(err error) bool

IsMessageError checks if an error is a message read/write failure.

func IsSessionClosed

func IsSessionClosed(err error) bool

IsSessionClosed checks if an error is a session closure.

func IsStreamClosed

func IsStreamClosed(err error) bool

IsStreamClosed checks if an error is a stream closure.

func IsUpgradeError

func IsUpgradeError(err error) bool

IsUpgradeError checks if an error is an upgrade failure.

func JoinPath

func JoinPath(segments ...string) string

JoinPath joins path segments with / separators, cleaning doubles.

func KeepAlive

func KeepAlive(c *Context, interval time.Duration) func()

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

func ListenAndServeWithGracefulShutdown(s *Server, drainTimeout time.Duration) error

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

func Messages(s *Stream) iter.Seq[[]byte]

Messages returns an iterator over length-prefixed messages from a stream.

for msg := range wt.Messages(stream) {
    process(msg)
}

func Must

func Must[T any](val T, err error) T

Must panics if err is non-nil. Useful for initialization.

func PProfMux

func PProfMux() *http.ServeMux

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

func Pipe(a, b *Stream) error

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

func RequiredFields(v any, fields ...string) error

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

func ServerInfo() map[string]string

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

func Streams(c *Context) iter.Seq[*Stream]

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

func ValidateDatagramSize(data []byte) error

ValidateDatagramSize checks if a datagram payload is within safe limits. Returns an error if the payload is too large.

func ValidateMessage

func ValidateMessage(msg any) error

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) Stop

func (cr *CertRotator) Stop()

Stop stops the certificate watcher.

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

func (c *Context) AcceptStream() (*Stream, error)

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) Close

func (c *Context) Close() error

Close closes the session with a success code.

func (*Context) CloseWithError

func (c *Context) CloseWithError(code uint32, msg string) error

CloseWithError closes the session with an error code and message.

func (*Context) Context

func (c *Context) Context() context.Context

Context returns the session's context (for cancellation/deadline).

func (*Context) Get

func (c *Context) Get(key string) (any, bool)

Get retrieves a value from the context store.

func (*Context) GetString

func (c *Context) GetString(key string) string

GetString retrieves a string value from the context store.

func (*Context) ID

func (c *Context) ID() string

ID returns a unique identifier for this session.

func (*Context) Info

func (c *Context) Info() ConnInfo

Info returns connection information for the session.

func (*Context) InfoJSON

func (c *Context) InfoJSON() string

InfoJSON returns connection info as a JSON string.

func (*Context) LocalAddr

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

LocalAddr returns the server's local address.

func (*Context) MustGet

func (c *Context) MustGet(key string) any

MustGet retrieves a value or panics if not found.

func (*Context) OpenStream

func (c *Context) OpenStream() (*Stream, error)

OpenStream opens a new bidirectional stream to the client.

func (*Context) OpenStreamSync

func (c *Context) OpenStreamSync() (*Stream, error)

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

func (c *Context) Param(name string) string

Param returns a path parameter value by name. For pattern "/chat/{room}" and path "/chat/general", Param("room") returns "general".

func (*Context) Params

func (c *Context) Params() map[string]string

Params returns all path parameters.

func (*Context) ReceiveDatagram

func (c *Context) ReceiveDatagram() ([]byte, error)

ReceiveDatagram receives the next datagram from the client.

func (*Context) ReceiveDatagramContext

func (c *Context) ReceiveDatagramContext(ctx context.Context) ([]byte, error)

ReceiveDatagramContext receives a datagram with explicit context for cancellation/timeout.

func (*Context) RemoteAddr

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

RemoteAddr returns the client's address.

func (*Context) Request

func (c *Context) Request() *http.Request

Request returns the original HTTP request that initiated the WebTransport session.

func (*Context) SendDatagram

func (c *Context) SendDatagram(data []byte) error

SendDatagram sends an unreliable datagram to the client.

func (*Context) SendDatagramSafe

func (c *Context) SendDatagramSafe(data []byte) error

SendDatagramSafe sends a datagram with size validation. Returns an error if the payload exceeds MaxDatagramSize.

func (*Context) Server

func (c *Context) Server() *Server

Server returns the Server instance.

func (*Context) Session

func (c *Context) Session() *webtransport.Session

Session returns the underlying webtransport.Session.

func (*Context) Set

func (c *Context) Set(key string, value any)

Set stores a key-value pair in the context (thread-safe).

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

type ErrorPageHandler func(w http.ResponseWriter, r *http.Request, code int, msg string)

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 Event

type Event struct {
	Type    EventType
	Session *Context
	Room    string // Only set for room events
}

Event represents a session lifecycle event.

type EventBus

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

EventBus provides a publish/subscribe system for session lifecycle events.

func NewEventBus

func NewEventBus() *EventBus

NewEventBus creates a new event bus.

func (*EventBus) Emit

func (eb *EventBus) Emit(event Event)

Emit publishes an event to all registered handlers. Handlers are called synchronously in registration order.

func (*EventBus) EmitAsync

func (eb *EventBus) EmitAsync(event Event)

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 EventHandler

type EventHandler func(Event)

EventHandler is a callback for session events.

type EventType

type EventType int

EventType represents a session lifecycle event.

const (
	EventConnect    EventType = iota // Session opened
	EventDisconnect                  // Session closed
	EventJoinRoom                    // Session joined a room
	EventLeaveRoom                   // Session left a room
)

func (EventType) String

func (e EventType) String() string

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.

func (*Group) Use

func (g *Group) Use(mw ...MiddlewareFunc)

Use adds middleware to 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.

func OnRead

func OnRead(fn func(data []byte) ([]byte, error)) InterceptorOption

OnRead sets a function that processes messages after reading from the stream. The function can transform or validate the message.

func OnWrite

func OnWrite(fn func(data []byte) ([]byte, error)) InterceptorOption

OnWrite sets a function that processes messages before writing to the stream.

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 NewKVSync

func NewKVSync() *KVSync

NewKVSync creates a new synchronized key-value store.

func (*KVSync) Delete

func (kv *KVSync) Delete(key string)

Delete removes a key.

func (*KVSync) Get

func (kv *KVSync) Get(key string, v any) error

Get retrieves a value by key and unmarshals it into v.

func (*KVSync) GetRaw

func (kv *KVSync) GetRaw(key string) (json.RawMessage, bool)

GetRaw retrieves the raw JSON for a key.

func (*KVSync) Keys

func (kv *KVSync) Keys() []string

Keys returns all keys.

func (*KVSync) Len

func (kv *KVSync) Len() int

Len returns the number of keys.

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.

func (*KVSync) Set

func (kv *KVSync) Set(key string, value any) error

Set sets a key-value pair and notifies the onChange callback.

func (*KVSync) Snapshot

func (kv *KVSync) Snapshot() map[string]json.RawMessage

Snapshot returns a copy of all key-value pairs as raw JSON.

type MessageError

type MessageError struct {
	Op      string // "read", "write"
	Size    int
	Wrapped error
}

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 WithAddr

func WithAddr(addr string) Option

WithAddr sets the listen address (default ":4433").

func WithAutoCert

func WithAutoCert(domain string, cacheDir string) Option

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

func WithAutoCertMulti(domains []string, cacheDir string) Option

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

func WithCheckOrigin(fn func(r *http.Request) bool) Option

WithCheckOrigin sets a function to validate the request origin.

func WithIdleTimeout

func WithIdleTimeout(d time.Duration) Option

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.

func WithTLS

func WithTLS(certFile, keyFile string) Option

WithTLS sets TLS certificate and key files.

type PreflightResult

type PreflightResult struct {
	Ready  bool
	Issues []string
}

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 NewPubSub

func NewPubSub() *PubSub

NewPubSub creates a new pub/sub hub.

func (*PubSub) Publish

func (ps *PubSub) Publish(topic string, data []byte)

Publish sends a datagram to all subscribers of a topic.

func (*PubSub) PublishExcept

func (ps *PubSub) PublishExcept(topic string, data []byte, excludeID string)

PublishExcept sends to all subscribers except the given session.

func (*PubSub) Subscribe

func (ps *PubSub) Subscribe(topic string, c *Context)

Subscribe adds a session to a topic.

func (*PubSub) SubscriberCount

func (ps *PubSub) SubscriberCount(topic string) int

SubscriberCount returns the number of subscribers for a topic.

func (*PubSub) Topics

func (ps *PubSub) Topics() []string

Topics returns all active topics.

func (*PubSub) TopicsForSession

func (ps *PubSub) TopicsForSession(sessionID string) []string

TopicsForSession returns all topics a session is subscribed to.

func (*PubSub) Unsubscribe

func (ps *PubSub) Unsubscribe(topic string, c *Context)

Unsubscribe removes a session from a topic.

func (*PubSub) UnsubscribeAll

func (ps *PubSub) UnsubscribeAll(c *Context)

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

func NewRPCClient(s *Stream) *RPCClient

NewRPCClient wraps a stream for RPC calls.

func (*RPCClient) Call

func (c *RPCClient) Call(method string, params any) (json.RawMessage, error)

Call sends an RPC request and waits for the response.

func (*RPCClient) Close

func (c *RPCClient) Close() error

Close closes the underlying stream.

type RPCError

type RPCError struct {
	Code    int    `json:"code"`
	Message string `json:"message"`
}

RPCError represents an error in the RPC response.

func (*RPCError) Error

func (e *RPCError) Error() string

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 NewRPCServer

func NewRPCServer() *RPCServer

NewRPCServer creates a new RPC server.

func (*RPCServer) Register

func (rpc *RPCServer) Register(method string, handler RPCHandler)

Register adds a handler for the given method name.

func (*RPCServer) Serve

func (rpc *RPCServer) Serve(s *Stream)

Serve handles RPC requests on the given stream. Each request-response pair uses the stream's message framing. Call this from a StreamHandler or HandleStream.

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]) Cap

func (rb *RingBuffer[T]) Cap() int

Cap returns the buffer capacity.

func (*RingBuffer[T]) Clear

func (rb *RingBuffer[T]) Clear()

Clear empties the buffer.

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) Broadcast

func (r *Room) Broadcast(data []byte)

Broadcast sends a datagram to all members in the room.

func (*Room) BroadcastExcept

func (r *Room) BroadcastExcept(data []byte, excludeID string)

BroadcastExcept sends a datagram to all members except the specified session.

func (*Room) BroadcastStream

func (r *Room) BroadcastStream(data []byte)

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

func (r *Room) BroadcastStreamExcept(data []byte, excludeID string)

BroadcastStreamExcept sends a reliable message to all except the given session.

func (*Room) Count

func (r *Room) Count() int

Count returns the number of members in the room.

func (*Room) ForEach

func (r *Room) ForEach(fn func(*Context))

ForEach iterates over all members without allocating a slice. The callback runs under a read lock — do not block for long.

func (*Room) Join

func (r *Room) Join(ctx *Context)

Join adds a session to the room.

func (*Room) Leave

func (r *Room) Leave(ctx *Context)

Leave removes a session from the room.

func (*Room) Members

func (r *Room) Members() []*Context

Members returns all sessions in the room.

func (*Room) Name

func (r *Room) Name() string

Name returns the room name.

func (*Room) OnJoin

func (r *Room) OnJoin(fn func(*Context))

OnJoin sets a callback for when sessions join the room.

func (*Room) OnLeave

func (r *Room) OnLeave(fn func(*Context))

OnLeave sets a callback for when sessions leave the room.

func (*Room) SafeBroadcast

func (r *Room) SafeBroadcast(data []byte, logger *slog.Logger)

SafeBroadcast sends a datagram to all room members, recovering from panics. Logs and skips any member that causes a panic (e.g., closed connection).

func (*Room) SafeBroadcastExcept

func (r *Room) SafeBroadcastExcept(data []byte, excludeID string, logger *slog.Logger)

SafeBroadcastExcept is SafeBroadcast but excludes the given session.

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.

func (*RoomManager) Remove

func (rm *RoomManager) Remove(name string)

Remove deletes a room.

func (*RoomManager) Rooms

func (rm *RoomManager) Rooms() []string

Rooms returns all room names.

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 NewRouter

func NewRouter() *Router

NewRouter creates a new Router.

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

func (r *Router) ExtractParams(pattern, path string) map[string]string

ExtractParams extracts path parameters from a URL path given a pattern.

func (*Router) Match

func (r *Router) Match(path string) (*Route, map[string]string)

Match finds the route matching the given path and extracts parameters.

func (*Router) Routes

func (r *Router) Routes() []*Route

Routes returns all registered routes.

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) Close

func (s *SendStream) Close() error

Close closes the send stream.

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

func New(opts ...Option) *Server

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) Addr

func (s *Server) Addr() string

Addr returns the configured listen address.

func (*Server) Broadcast

func (s *Server) Broadcast(data []byte)

Broadcast sends a datagram to all active sessions.

func (*Server) BroadcastExcept

func (s *Server) BroadcastExcept(data []byte, excludeID string)

BroadcastExcept sends a datagram to all active sessions except the specified one.

func (*Server) CertHash

func (s *Server) CertHash() string

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) Close

func (s *Server) Close() error

Close gracefully shuts down the server.

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

func (s *Server) ListenAndServe() error

ListenAndServe starts the WebTransport server.

func (*Server) Multicast

func (s *Server) Multicast(data []byte, filter func(*Context) bool)

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

func (s *Server) MulticastStream(data []byte, filter func(*Context) bool)

MulticastStream sends a reliable message via streams to matching sessions.

func (*Server) OnConnect

func (s *Server) OnConnect(fn func(*Context))

OnConnect registers a callback for new sessions (after middleware).

func (*Server) OnDisconnect

func (s *Server) OnDisconnect(fn func(*Context))

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

func (s *Server) Preflight() []string

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

func (s *Server) SessionCount() int

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

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

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

type SessionCloseError struct {
	Code    uint32
	Message string
}

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) Add

func (ss *SessionStore) Add(ctx *Context)

Add registers a session.

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

func OpenTypedStream(c *Context, typeID uint16) (*Stream, error)

OpenTypedStream opens a stream with the given type header. The remote end's StreamMux will route it to the matching handler.

func RetryStream

func RetryStream(ctx context.Context, c *Context, cfg RetryConfig) (*Stream, error)

RetryStream attempts to open a stream with retry.

func (*Stream) CancelRead

func (s *Stream) CancelRead(code uint32)

CancelRead cancels the read side of the stream.

func (*Stream) CancelWrite

func (s *Stream) CancelWrite(code uint32)

CancelWrite cancels the write side of the stream.

func (*Stream) Close

func (s *Stream) Close() error

Close closes the stream.

func (*Stream) Raw

func (s *Stream) Raw() *webtransport.Stream

Raw returns the underlying webtransport.Stream.

func (*Stream) Read

func (s *Stream) Read(b []byte) (int, error)

Read reads raw bytes from the stream.

func (*Stream) ReadMessage

func (s *Stream) ReadMessage() ([]byte, error)

ReadMessage reads a length-prefixed message from the stream.

func (*Stream) SessionContext

func (s *Stream) SessionContext() *Context

Context returns the session context this stream belongs to.

func (*Stream) SetDeadline

func (s *Stream) SetDeadline(t time.Time) error

SetDeadline sets read and write deadlines.

func (*Stream) SetReadDeadline

func (s *Stream) SetReadDeadline(t time.Time) error

SetReadDeadline sets the read deadline.

func (*Stream) SetWriteDeadline

func (s *Stream) SetWriteDeadline(t time.Time) error

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) Write

func (s *Stream) Write(b []byte) (int, error)

Write writes raw bytes to the stream.

func (*Stream) WriteMessage

func (s *Stream) WriteMessage(data []byte) error

WriteMessage writes a length-prefixed message to the stream. Format: [4 bytes big-endian length][payload]

type StreamCloseError

type StreamCloseError struct {
	Code   uint32
	Remote bool // true if the remote side closed
}

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

type StreamHandler func(s *Stream, c *Context)

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.

func (*StreamMux) Serve

func (m *StreamMux) Serve(c *Context)

Serve accepts streams from the context and routes them by type. Blocks until the session is closed.

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 NewTags

func NewTags() *Tags

NewTags creates a new tag registry.

func (*Tags) AllTags

func (t *Tags) AllTags() []string

AllTags returns all known tags.

func (*Tags) Count

func (t *Tags) Count(tag string) int

Count returns the number of sessions with the given tag.

func (*Tags) HasTag

func (t *Tags) HasTag(sessionID, tag string) bool

HasTag checks if a session has a specific tag.

func (*Tags) SessionsWithTag

func (t *Tags) SessionsWithTag(tag string) []string

SessionsWithTag returns all session IDs with the given tag.

func (*Tags) Tag

func (t *Tags) Tag(sessionID, tag string)

Tag adds a tag to a session.

func (*Tags) TagsForSession

func (t *Tags) TagsForSession(sessionID string) []string

TagsForSession returns all tags for a session.

func (*Tags) Untag

func (t *Tags) Untag(sessionID, tag string)

Untag removes a tag from a session.

func (*Tags) UntagAll

func (t *Tags) UntagAll(sessionID string)

UntagAll removes 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
})

func NewTicker

func NewTicker(c *Context, interval time.Duration, getData func() []byte) *Ticker

NewTicker starts sending datagrams at the given interval. The getData function is called each tick to produce the datagram payload. Return nil to skip a tick (no datagram sent).

func (*Ticker) Stop

func (t *Ticker) Stop()

Stop stops the ticker.

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

func NewTypedRoom[T any](room *Room, c codec.Codec) *TypedRoom[T]

NewTypedRoom wraps a Room with typed broadcast support.

func (*TypedRoom[T]) Broadcast

func (tr *TypedRoom[T]) Broadcast(v T) error

Broadcast encodes and broadcasts a value to all room members via datagrams.

func (*TypedRoom[T]) BroadcastExcept

func (tr *TypedRoom[T]) BroadcastExcept(v T, excludeID string) error

BroadcastExcept broadcasts to all members except the specified session.

func (*TypedRoom[T]) Room

func (tr *TypedRoom[T]) Room() *Room

Room returns the underlying Room.

type TypedStream

type TypedStream[R any, W any] struct {
	// contains filtered or unexported fields
}

TypedStream provides type-safe read/write over a Stream using a codec.

func NewTypedStream

func NewTypedStream[R any, W any](s *Stream, c codec.Codec) *TypedStream[R, W]

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

type UpgradeError struct {
	StatusCode int
	Message    string
}

UpgradeError occurs when WebTransport upgrade fails.

func (*UpgradeError) Error

func (e *UpgradeError) Error() string

type Validator

type Validator interface {
	Validate() error
}

Validator can validate itself. Implement on message types for automatic validation.

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.

Jump to

Keyboard shortcuts

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