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 goroutines
  • Stop() error – Gracefully terminates connections and cleanup resources
  • Send(data []byte) error – Transmits packets to the remote peer
  • Receive(callback func([]byte)) – Registers the callback for incoming packets
  • IsConnected() bool – Returns the current connection state
  • Stats() 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 running and connected flags with mutex protection
  • Statistics tracking: Automatic byte counters via RecordSend and RecordReceive
  • Callback routing: The CallReceive method to dispatch packets to the registered handler
  • Reconnection bookkeeping: RecordReconnect for tracking retry attempts
  • Configuration access: GetConfig() exposing TransportConfig values like MaxQueueSize and ReconnectDelay

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:

  1. Call t.BaseTransport.Start() to mark the transport as running
  2. Establish the underlying connection (TCP, WebSocket, etc.)
  3. Set t.SetConnected(true) via the base transport
  4. 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:

  1. Signal background goroutines to exit
  2. Close underlying connections
  3. 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 loops
  • ReconnectDelay: Initial delay between attempts
  • ReconnectMultiplier: 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

Summary

  • Embed *BaseTransport from transport/transport.go to inherit connection state, statistics, and lifecycle management
  • Implement six core methods: Start, Stop, Send, Receive, IsConnected, and Stats (though Stats is often handled by the base)
  • Use helper methods: RecordSend, RecordReceive, CallReceive, SetConnected, and RecordReconnect for consistent state tracking
  • Respect configuration: Honor MaxQueueSize, ReconnectDelay, and MaxReconnectAttempts from TransportConfig for consistent behavior across transports
  • Wrap for encryption: Use EncryptedTransport in transport/encrypted.go to 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:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →