# How to Implement a Custom Transport for OpenFlux: A Complete Guide

> Learn to implement a custom transport for OpenFlux by embedding BaseTransport and implementing key interface methods. Follow this guide to extend OpenFlux functionality.

- Repository: [p1neappleXpress/OpenFlux](https://github.com/p1neappleXpress/OpenFlux)
- Tags: how-to-guide
- Published: 2026-09-14

---

**To implement a custom transport for OpenFlux, create a struct that embeds `*transport.BaseTransport` from [`transport/transport.go`](https://github.com/p1neappleXpress/OpenFlux/blob/main/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`](https://github.com/p1neappleXpress/OpenFlux/blob/main/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`](https://github.com/p1neappleXpress/OpenFlux/blob/main/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.

```go
type MyProtoTransport struct {
	*transport.BaseTransport
	addr   string
	conn   net.Conn
	quitCh chan struct{}
}

```

Provide a constructor that initializes the base transport with a configuration:

```go
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.

```go
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:

```go
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.

```go
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:

```go
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:

```go
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.

```go
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`](https://github.com/p1neappleXpress/OpenFlux/blob/main/transport/encrypted.go). You can wrap any custom transport to add AES-256-GCM encryption without modifying the underlying implementation.

```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.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`](https://github.com/p1neappleXpress/OpenFlux/blob/main/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`](https://github.com/p1neappleXpress/OpenFlux/blob/main/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`](https://github.com/p1neappleXpress/OpenFlux/blob/main/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`](https://github.com/p1neappleXpress/OpenFlux/blob/main/transport/encrypted.go) to add end-to-end encryption without code duplication.
- Reference [`transport/yandex/yandex.go`](https://github.com/p1neappleXpress/OpenFlux/blob/main/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`](https://github.com/p1neappleXpress/OpenFlux/blob/main/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`](https://github.com/p1neappleXpress/OpenFlux/blob/main/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`](https://github.com/p1neappleXpress/OpenFlux/blob/main/transport/yandex/yandex.go)** – WebSocket transport with exponential backoff reconnects and keep-alive.
- **[`transport/oneme/max_transport.go`](https://github.com/p1neappleXpress/OpenFlux/blob/main/transport/oneme/max_transport.go)** – Example of wrapping an external client library to satisfy the interface.
- **[`transport/compressor.go`](https://github.com/p1neappleXpress/OpenFlux/blob/main/transport/compressor.go)** – Utility for optional packet compression that can be layered atop any transport.