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 – running and connected flags with mutex protection via Mu.
  • Statistics – RecordSend(bytes), RecordReceive(bytes), and RecordReconnect() counters.
  • Callbacks – CallReceive(payload) to route data through the registered callback.
  • Configuration – Access to TransportConfig fields including MaxQueueSize, ReconnectDelay, and MaxReconnectAttempts.

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 via cfg.ReconnectDelay * cfg.ReconnectMultiplier^attempt.
  • Keep-alive – Periodic ping frames to maintain NAT mappings.
  • External libraries – transport/oneme/max_transport.go shows how to delegate work to third-party client libraries while still satisfying the Transport interface.

Summary

  • Implement the Transport interface from transport/transport.go by embedding *transport.BaseTransport to inherit state management and statistics.
  • Use SetConnected, RecordSend, RecordReceive, and RecordReconnect to keep the base transport state synchronized.
  • Respect TransportConfig fields such as MaxQueueSize and MaxReconnectAttempts to ensure consistent behavior across all OpenFlux transports.
  • Wrap custom transports with NewEncryptedTransport from transport/encrypted.go to add end-to-end encryption without code duplication.
  • Reference transport/yandex/yandex.go for 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:

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 →