How to Implement a Custom Transport for OpenFlux: A Complete Developer Guide
To implement a custom transport for OpenFlux, create a struct that embeds *transport.BaseTransport from transport/transport.go, implement the Transport interface methods (Start, Stop, Send, Receive, IsConnected), and utilize helper methods like RecordSend and CallReceive for statistics and packet routing.
OpenFlux defines a pluggable transport layer in the p1neappleXpress/OpenFlux repository that abstracts packet transmission between clients and exit nodes. All custom transports must satisfy the Transport interface declared in transport/transport.go and typically embed the shared BaseTransport helper to inherit common connection management, statistics tracking, and lifecycle controls. This architecture allows developers to add support for custom protocols—whether TCP variants, WebSockets, or proprietary binary formats—while reusing the framework's robust state management and reconnection logic.
Understanding the Transport Interface Architecture
The foundation of every OpenFlux transport resides in transport/transport.go, which defines the contract that all transport implementations must fulfill.
The Transport Interface
The Transport interface declares six core methods that manage the complete lifecycle of a connection:
Start() error– Opens underlying connections and launches background goroutinesStop() error– Gracefully terminates connections and cleanup resourcesSend(data []byte) error– Transmits packets to the remote peerReceive(callback func([]byte))– Registers the callback for incoming packetsIsConnected() bool– Returns the current connection stateStats() TransportStats– Provides transmission statistics
The BaseTransport Helper
Rather than implementing these methods from scratch, custom transports embed *BaseTransport, a concrete struct defined in transport/transport.go that provides:
- State management: Thread-safe
runningandconnectedflags with mutex protection - Statistics tracking: Automatic byte counters via
RecordSendandRecordReceive - Callback routing: The
CallReceivemethod to dispatch packets to the registered handler - Reconnection bookkeeping:
RecordReconnectfor tracking retry attempts - Configuration access:
GetConfig()exposingTransportConfigvalues likeMaxQueueSizeandReconnectDelay
Step-by-Step Implementation Guide
Building a custom transport involves defining your protocol-specific state while delegating common operations to the base implementation.
1. Define the Transport Struct
Create a new Go package (conventionally under transport/<yourproto>/) and define a struct that embeds the base transport alongside protocol-specific fields:
type MyProtoTransport struct {
*transport.BaseTransport
address string
conn net.Conn
quitCh chan struct{}
}
2. Implement the Constructor
Use transport.NewBaseTransport(cfg) to initialize the embedded helper with your chosen configuration:
func NewMyProtoTransport(addr string, cfg transport.TransportConfig) *MyProtoTransport {
return &MyProtoTransport{
BaseTransport: transport.NewBaseTransport(cfg),
address: addr,
quitCh: make(chan struct{}),
}
}
3. Handle Connection Lifecycle
The Start and Stop methods manage the transition between states.
Start must:
- Call
t.BaseTransport.Start()to mark the transport as running - Establish the underlying connection (TCP, WebSocket, etc.)
- Set
t.SetConnected(true)via the base transport - Launch background readers using
utils.SafeGo
func (t *MyProtoTransport) Start() error {
if err := t.BaseTransport.Start(); err != nil {
return err
}
c, err := net.Dial("tcp", t.address)
if err != nil {
t.RecordReconnect()
return err
}
t.conn = c
t.SetConnected(true)
utils.SafeGo("myproto.reader", t.readLoop)
return nil
}
Stop must:
- Signal background goroutines to exit
- Close underlying connections
- Call
t.BaseTransport.Stop()to update shared state
func (t *MyProtoTransport) Stop() error {
close(t.quitCh)
if t.conn != nil {
_ = t.conn.Close()
}
return t.BaseTransport.Stop()
}
4. Implement Data Transmission
The Send method writes data to the underlying channel while respecting back-pressure and recording statistics:
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
}
For production use, implement a bounded write channel sized according to t.GetConfig().MaxQueueSize to prevent unbounded memory growth during network congestion.
5. Implement Packet Reception
The Receive method registration delegates to the base transport, while your read loop invokes CallReceive to route packets:
func (t *MyProtoTransport) Receive(cb func([]byte)) {
t.BaseTransport.Receive(cb)
}
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])
}
}
}
6. Expose Connection State
Delegate the IsConnected check to the base transport:
func (t *MyProtoTransport) IsConnected() bool {
return t.BaseTransport.IsConnected()
}
Complete Minimal Example
Below is a functional TCP-based transport demonstrating the full integration pattern from transport/transport.go:
package myproto
import (
"fmt"
"net"
"universal-bypass-tool/transport"
"universal-bypass-tool/utils"
)
type MyProtoTransport struct {
*transport.BaseTransport
addr string
conn net.Conn
quitCh chan struct{}
}
func NewMyProtoTransport(addr string, cfg transport.TransportConfig) *MyProtoTransport {
return &MyProtoTransport{
BaseTransport: transport.NewBaseTransport(cfg),
addr: addr,
quitCh: make(chan struct{}),
}
}
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)
utils.SafeGo("myproto.reader", t.readLoop)
return nil
}
func (t *MyProtoTransport) Stop() error {
close(t.quitCh)
if t.conn != nil {
_ = t.conn.Close()
}
return t.BaseTransport.Stop()
}
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
}
func (t *MyProtoTransport) Receive(cb func([]byte)) {
t.BaseTransport.Receive(cb)
}
func (t *MyProtoTransport) IsConnected() bool {
return t.BaseTransport.IsConnected()
}
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])
}
}
}
Advanced Patterns
Adding Encryption with EncryptedTransport
For end-to-end encryption, wrap your custom transport with the EncryptedTransport layer defined in transport/encrypted.go:
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.Fatal(err)
}
enc.Start()
This preserves your transport's interface while adding AES-256-GCM encryption without modifying your implementation.
Implementing Reconnect Logic
The transport/yandex/yandex.go file demonstrates sophisticated reconnection handling. Respect the configuration parameters from TransportConfig:
MaxReconnectAttempts: Limit retry loopsReconnectDelay: Initial delay between attemptsReconnectMultiplier: Exponential backoff factor
Use t.RecordReconnect() to increment the attempt counter and update statistics when connection failures occur.
Handling Back-Pressure
As shown in the Yandex transport implementation, size your write queues using cfg.MaxQueueSize to prevent memory exhaustion. The base transport tracks these statistics, allowing upstream components to monitor queue depth via Stats().
Key Reference Files
transport/transport.go– DefinesTransportinterface,BaseTransportimplementation, andTransportConfig(source)transport/encrypted.go– Demonstrates transport wrapping and AES-256-GCM encryption (source)transport/yandex/yandex.go– Complex WebSocket transport with reconnect and keep-alive logic (source)transport/oneme/max_transport.go– Example of delegating to external client libraries (source)transport/compressor.go– Optional packet compression utilities (source)
Summary
- Embed
*BaseTransportfromtransport/transport.goto inherit connection state, statistics, and lifecycle management - Implement six core methods:
Start,Stop,Send,Receive,IsConnected, andStats(thoughStatsis often handled by the base) - Use helper methods:
RecordSend,RecordReceive,CallReceive,SetConnected, andRecordReconnectfor consistent state tracking - Respect configuration: Honor
MaxQueueSize,ReconnectDelay, andMaxReconnectAttemptsfromTransportConfigfor consistent behavior across transports - Wrap for encryption: Use
EncryptedTransportintransport/encrypted.goto add security without rewriting transport logic
Frequently Asked Questions
What methods must I implement for a custom OpenFlux transport?
You must satisfy the Transport interface defined in transport/transport.go, which requires Start() error, Stop() error, Send(data []byte) error, Receive(callback func([]byte)), IsConnected() bool, and Stats() TransportStats. When embedding BaseTransport, you inherit default implementations for Receive and Stats, but you must override Start, Stop, Send, and IsConnected to handle your protocol-specific logic.
How do I handle reconnection logic in OpenFlux transports?
Invoke t.RecordReconnect() from the base transport whenever a connection attempt fails, then respect the ReconnectDelay and MaxReconnectAttempts fields from TransportConfig. The transport/yandex/yandex.go file provides a reference implementation using exponential backoff based on ReconnectMultiplier. Always update the connected state via t.SetConnected(false) when connections drop and t.SetConnected(true) when re-established.
Can I add encryption to my custom transport without modifying its code?
Yes. Create your transport normally, then wrap it using transport.NewEncryptedTransport() as demonstrated in transport/encrypted.go. This accepts any Transport implementation and returns an encrypted wrapper that handles AES-256-GCM encryption/decryption transparently, preserving the same interface while adding cryptographic protection.
Where should I place my custom transport code in the OpenFlux repository?
Create a new subdirectory under transport/ (e.g., transport/myproto/) containing your implementation files. This convention keeps transport implementations organized and allows the build system to locate dependencies consistently. Place configuration structs in the same file as your transport or in a separate config.go file within that directory.
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 →