How to Implement a Custom Transport for OpenFlux: A Complete Guide
To implement a custom transport for OpenFlux, create a struct that embeds *transport.BaseTransport from transport/transport.go and implements the Transport interface methods: Start, Stop, Send, and Receive.
OpenFlux exposes a pluggable transport layer that abstracts how packets travel between clients and exit nodes. When you implement a custom transport for OpenFlux, you can tunnel traffic over proprietary protocols, WebSockets, or specialized TCP implementations while reusing the framework’s connection lifecycle, statistics, and reconnect logic. The architecture centers on a minimal interface declaration and a reusable BaseTransport helper that manages thread-safe state and bookkeeping.
Understanding the Transport Interface
All transports must satisfy the Transport interface declared in transport/transport.go. This contract ensures that OpenFlux can manage connections uniformly regardless of the underlying protocol.
Core Interface Requirements
The Transport interface in transport/transport.go declares six methods:
Start()– Initializes the underlying connection and launches background goroutines.Stop()– Gracefully terminates connections and cleans up resources.Send(data []byte)– Transmits packets to the remote peer.Receive(callback)– Registers a callback for incoming packets.IsConnected()– Reports the current connection state.Stats()– Returns transmission metrics.
The BaseTransport Helper
Rather than implementing state management from scratch, embed *transport.BaseTransport (defined in the same file). This helper provides:
- State tracking –
runningandconnectedflags with mutex protection viaMu. - Statistics –
RecordSend(bytes),RecordReceive(bytes), andRecordReconnect()counters. - Callbacks –
CallReceive(payload)to route data through the registered callback. - Configuration – Access to
TransportConfigfields includingMaxQueueSize,ReconnectDelay, andMaxReconnectAttempts.
Embedding BaseTransport allows your implementation to focus on protocol-specific logic while inheriting thread-safe state transitions and metric collection.
Step-by-Step Implementation Guide
Follow this pattern to build a production-ready custom transport.
1. Create Your Transport Package
Create a new subdirectory under transport/ (e.g., transport/myproto). This keeps your implementation isolated and follows the existing project structure used by transport/yandex/ and transport/oneme/.
2. Define the Struct and Constructor
Define a struct that embeds *transport.BaseTransport and stores protocol-specific state such as net.Conn or WebSocket clients.
type MyProtoTransport struct {
*transport.BaseTransport
addr string
conn net.Conn
quitCh chan struct{}
}
Provide a constructor that initializes the base transport with a configuration:
func NewMyProtoTransport(addr string, cfg transport.TransportConfig) *MyProtoTransport {
return &MyProtoTransport{
BaseTransport: transport.NewBaseTransport(cfg),
addr: addr,
quitCh: make(chan struct{}),
}
}
3. Implement Start and Stop Lifecycle
The Start method must call t.BaseTransport.Start() first to set the running flag, then establish the connection and set t.SetConnected(true). Launch background readers using goroutines that respect the quitCh signal.
func (t *MyProtoTransport) Start() error {
if err := t.BaseTransport.Start(); err != nil {
return err
}
c, err := net.Dial("tcp", t.addr)
if err != nil {
t.RecordReconnect()
return err
}
t.conn = c
t.SetConnected(true)
// Launch background reader
utils.SafeGo("myproto.reader", t.readLoop)
return nil
}
For Stop, close the quit channel, close the connection, and delegate to the base implementation:
func (t *MyProtoTransport) Stop() error {
close(t.quitCh)
if t.conn != nil {
_ = t.conn.Close()
}
return t.BaseTransport.Stop()
}
4. Implement Send with Back-pressure
The Send method writes data to the underlying channel and records statistics via t.RecordSend(len(data)). Respect back-pressure by using bounded channels sized according to t.GetConfig().MaxQueueSize to prevent unbounded memory growth.
func (t *MyProtoTransport) Send(data []byte) error {
if !t.IsConnected() {
return fmt.Errorf("transport not connected")
}
_, err := t.conn.Write(data)
if err == nil {
t.RecordSend(len(data))
}
return err
}
5. Implement Receive and Callbacks
Override Receive to register the callback with the base transport:
func (t *MyProtoTransport) Receive(cb func([]byte)) {
t.BaseTransport.Receive(cb)
}
In your read loop, invoke t.CallReceive(payload) whenever a full packet arrives. This routes the payload through the base transport’s registered callback mechanism:
func (t *MyProtoTransport) readLoop() {
buf := make([]byte, 4096)
for {
select {
case <-t.quitCh:
return
default:
}
n, err := t.conn.Read(buf)
if err != nil {
t.SetConnected(false)
return
}
if n > 0 {
t.RecordReceive(n)
t.CallReceive(buf[:n])
}
}
}
Complete Working Example
Below is a minimal TCP-based transport that demonstrates the full skeleton. This example uses the shared base transport for bookkeeping and implements graceful shutdown.
package myproto
import (
"fmt"
"net"
"universal-bypass-tool/transport"
"universal-bypass-tool/utils"
)
// MyProtoTransport is a simple TCP-based transport.
type MyProtoTransport struct {
*transport.BaseTransport
addr string
conn net.Conn
quitCh chan struct{}
}
// NewMyProtoTransport creates a transport that will connect to addr.
func NewMyProtoTransport(addr string, cfg transport.TransportConfig) *MyProtoTransport {
return &MyProtoTransport{
BaseTransport: transport.NewBaseTransport(cfg),
addr: addr,
quitCh: make(chan struct{}),
}
}
// Start opens the TCP connection and launches the read loop.
func (t *MyProtoTransport) Start() error {
if err := t.BaseTransport.Start(); err != nil {
return err
}
c, err := net.Dial("tcp", t.addr)
if err != nil {
t.RecordReconnect()
return err
}
t.conn = c
t.SetConnected(true)
// Launch a background reader.
utils.SafeGo("myproto.reader", t.readLoop)
return nil
}
// Stop closes the connection and stops the background goroutine.
func (t *MyProtoTransport) Stop() error {
close(t.quitCh)
if t.conn != nil {
_ = t.conn.Close()
}
return t.BaseTransport.Stop()
}
// Send writes data to the socket, respecting back-pressure.
func (t *MyProtoTransport) Send(data []byte) error {
if !t.IsConnected() {
return fmt.Errorf("transport not connected")
}
// Simple write; a real implementation would use a bounded channel.
_, err := t.conn.Write(data)
if err == nil {
t.RecordSend(len(data))
}
return err
}
// readLoop reads from the socket and forwards packets to the base transport.
func (t *MyProtoTransport) readLoop() {
buf := make([]byte, 4096)
for {
select {
case <-t.quitCh:
return
default:
}
n, err := t.conn.Read(buf)
if err != nil {
t.SetConnected(false)
return
}
if n > 0 {
// Record the reception and forward the payload.
t.RecordReceive(n)
t.CallReceive(buf[:n])
}
}
}
Advanced Patterns and Encryption
Layering with EncryptedTransport
OpenFlux supports transparent encryption via EncryptedTransport defined in transport/encrypted.go. You can wrap any custom transport to add AES-256-GCM encryption without modifying the underlying implementation.
cfg := transport.DefaultConfig()
tcp := myproto.NewMyProtoTransport("example.com:443", cfg)
enc, err := transport.NewEncryptedTransport(tcp, "my-shared-secret", "session-123", false)
if err != nil {
log.Fatalf("failed to create encrypted transport: %v", err)
}
if err := enc.Start(); err != nil {
log.Fatalf("transport start error: %v", err)
}
Production Patterns from Yandex Transport
For complex implementations involving WebSockets, keep-alive loops, and exponential backoff, study transport/yandex/yandex.go. This file demonstrates:
- Reconnect logic – Using
RecordReconnect()and calculating retry delays viacfg.ReconnectDelay * cfg.ReconnectMultiplier^attempt. - Keep-alive – Periodic ping frames to maintain NAT mappings.
- External libraries –
transport/oneme/max_transport.goshows how to delegate work to third-party client libraries while still satisfying theTransportinterface.
Summary
- Implement the
Transportinterface fromtransport/transport.goby embedding*transport.BaseTransportto inherit state management and statistics. - Use
SetConnected,RecordSend,RecordReceive, andRecordReconnectto keep the base transport state synchronized. - Respect
TransportConfigfields such asMaxQueueSizeandMaxReconnectAttemptsto ensure consistent behavior across all OpenFlux transports. - Wrap custom transports with
NewEncryptedTransportfromtransport/encrypted.goto add end-to-end encryption without code duplication. - Reference
transport/yandex/yandex.gofor production-grade examples of reconnection and back-pressure handling.
Frequently Asked Questions
What methods must a custom transport implement?
You must implement Start, Stop, Send, Receive, and IsConnected. The Stats method is typically inherited from BaseTransport. These are defined in transport/transport.go and ensure OpenFlux can manage your transport uniformly.
How does BaseTransport help with statistics?
BaseTransport provides RecordSend and RecordReceive methods that update internal counters automatically. When you wrap your transport with EncryptedTransport or expose metrics via Stats(), these counters report accurate byte and packet counts without manual bookkeeping.
Can I add encryption to my custom transport?
Yes. Import transport/encrypted.go and wrap your transport with NewEncryptedTransport(yourTransport, secret, sessionID, isClient). This returns a Transport interface that encrypts all Send data and decrypts Receive payloads transparently.
Where can I find production examples of custom transports?
The OpenFlux repository contains several reference implementations:
transport/yandex/yandex.go– WebSocket transport with exponential backoff reconnects and keep-alive.transport/oneme/max_transport.go– Example of wrapping an external client library to satisfy the interface.transport/compressor.go– Utility for optional packet compression that can be layered atop any transport.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →