Compare commits
4 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 6fdc9ea918 | |||
| d869589663 | |||
| 072362a8e6 | |||
| b628a5f964 |
18
changelog.md
18
changelog.md
@@ -1,5 +1,23 @@
|
|||||||
# Changelog
|
# Changelog
|
||||||
|
|
||||||
|
## 2026-02-18 - 3.2.0 - feat(remoteingress (edge/hub/protocol))
|
||||||
|
add dynamic port configuration: handshake, FRAME_CONFIG frames, and hot-reloadable listeners
|
||||||
|
|
||||||
|
- Introduce a JSON handshake from hub -> edge with initial listen ports and stun interval so edges can configure listeners at connect time.
|
||||||
|
- Add FRAME_CONFIG (0x06) to the protocol and implement runtime config updates pushed from hub to connected edges.
|
||||||
|
- Edge now applies initial ports and supports hot-reloading: spawn/abort listeners when ports change, and emit PortsAssigned / PortsUpdated events.
|
||||||
|
- Hub now stores allowed edge metadata (listen_ports, stun_interval_secs), sends handshake responses on auth, and forwards config updates to connected edges.
|
||||||
|
- TypeScript bridge/client updated to emit new port events and periodically log status; updateAllowedEdges API accepts listenPorts and stunIntervalSecs.
|
||||||
|
- Stun interval handling moved to use handshake-provided/stored value instead of config.listen_ports being static.
|
||||||
|
|
||||||
|
## 2026-02-18 - 3.1.1 - fix(readme)
|
||||||
|
update README: add issue reporting/security section, document connection tokens and token utilities, clarify architecture/API and improve examples/formatting
|
||||||
|
|
||||||
|
- Added an 'Issue Reporting and Security' section linking to community.foss.global for bug/security reports and contributor onboarding.
|
||||||
|
- Documented connection tokens: encodeConnectionToken/decodeConnectionToken utilities, token format (base64url), and examples for hub and edge provisioning.
|
||||||
|
- Clarified Hub/Edge usage and examples: condensed event handlers, added token-based start() example, and provided explicit config alternative.
|
||||||
|
- Improved README formatting: added emojis, rephrased architecture descriptions, fixed wording and license path capitalization, and expanded example scenarios and interfaces.
|
||||||
|
|
||||||
## 2026-02-17 - 3.1.0 - feat(edge)
|
## 2026-02-17 - 3.1.0 - feat(edge)
|
||||||
support connection tokens when starting an edge and add token encode/decode utilities
|
support connection tokens when starting an edge and add token encode/decode utilities
|
||||||
|
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
{
|
{
|
||||||
"name": "@serve.zone/remoteingress",
|
"name": "@serve.zone/remoteingress",
|
||||||
"version": "3.1.0",
|
"version": "3.2.0",
|
||||||
"private": false,
|
"private": false,
|
||||||
"description": "Edge ingress tunnel for DcRouter - accepts incoming TCP connections at network edge and tunnels them to DcRouter SmartProxy preserving client IP via PROXY protocol v1.",
|
"description": "Edge ingress tunnel for DcRouter - accepts incoming TCP connections at network edge and tunnels them to DcRouter SmartProxy preserving client IP via PROXY protocol v1.",
|
||||||
"main": "dist_ts/index.js",
|
"main": "dist_ts/index.js",
|
||||||
|
|||||||
199
readme.md
199
readme.md
@@ -2,13 +2,17 @@
|
|||||||
|
|
||||||
Edge ingress tunnel for DcRouter — accepts incoming TCP connections at the network edge and tunnels them to a DcRouter SmartProxy instance, preserving the original client IP via PROXY protocol v1.
|
Edge ingress tunnel for DcRouter — accepts incoming TCP connections at the network edge and tunnels them to a DcRouter SmartProxy instance, preserving the original client IP via PROXY protocol v1.
|
||||||
|
|
||||||
|
## Issue Reporting and Security
|
||||||
|
|
||||||
|
For reporting bugs, issues, or security vulnerabilities, please visit [community.foss.global/](https://community.foss.global/). This is the central community hub for all issue reporting. Developers who sign and comply with our contribution agreement and go through identification can also get a [code.foss.global/](https://code.foss.global/) account to submit Pull Requests directly.
|
||||||
|
|
||||||
## Install
|
## Install
|
||||||
|
|
||||||
```sh
|
```sh
|
||||||
pnpm install @serve.zone/remoteingress
|
pnpm install @serve.zone/remoteingress
|
||||||
```
|
```
|
||||||
|
|
||||||
## Architecture
|
## 🏗️ Architecture
|
||||||
|
|
||||||
`@serve.zone/remoteingress` uses a **Hub/Edge** topology with a high-performance Rust core and a TypeScript API surface:
|
`@serve.zone/remoteingress` uses a **Hub/Edge** topology with a high-performance Rust core and a TypeScript API surface:
|
||||||
|
|
||||||
@@ -17,8 +21,8 @@ pnpm install @serve.zone/remoteingress
|
|||||||
│ Network Edge │ ◄══════════════════════════► │ Private Cluster │
|
│ Network Edge │ ◄══════════════════════════► │ Private Cluster │
|
||||||
│ │ (multiplexed frames + │ │
|
│ │ (multiplexed frames + │ │
|
||||||
│ RemoteIngressEdge │ shared-secret auth) │ RemoteIngressHub │
|
│ RemoteIngressEdge │ shared-secret auth) │ RemoteIngressHub │
|
||||||
│ Listens on :80,:443│ │ Forwards to │
|
│ Accepts client TCP │ │ Forwards to │
|
||||||
│ Accepts client TCP │ │ SmartProxy on │
|
│ connections │ │ SmartProxy on │
|
||||||
│ │ │ local ports │
|
│ │ │ local ports │
|
||||||
└─────────────────────┘ └─────────────────────┘
|
└─────────────────────┘ └─────────────────────┘
|
||||||
▲ │
|
▲ │
|
||||||
@@ -28,26 +32,27 @@ pnpm install @serve.zone/remoteingress
|
|||||||
|
|
||||||
| Component | Role |
|
| Component | Role |
|
||||||
|-----------|------|
|
|-----------|------|
|
||||||
| **RemoteIngressEdge** | Deployed at the network edge (e.g. a VPS or cloud instance). Listens on public ports, accepts raw TCP connections, and multiplexes them over a single TLS tunnel to the hub. |
|
| **RemoteIngressEdge** | Deployed at the network edge (e.g. a VPS or cloud instance). Accepts raw TCP connections and multiplexes them over a single TLS tunnel to the hub. |
|
||||||
| **RemoteIngressHub** | Deployed alongside DcRouter/SmartProxy in a private cluster. Accepts edge connections, demuxes streams, and forwards each to SmartProxy with a PROXY protocol v1 header so the real client IP is preserved. |
|
| **RemoteIngressHub** | Deployed alongside DcRouter/SmartProxy in a private cluster. Accepts edge connections, demuxes streams, and forwards each to SmartProxy with a PROXY protocol v1 header so the real client IP is preserved. |
|
||||||
| **Rust Binary** (`remoteingress-bin`) | The performance-critical networking core. Managed via `@push.rocks/smartrust` RustBridge IPC — you never interact with it directly. Cross-compiled for `linux/amd64` and `linux/arm64`. |
|
| **Rust Binary** (`remoteingress-bin`) | The performance-critical networking core. Managed via `@push.rocks/smartrust` RustBridge IPC — you never interact with it directly. Cross-compiled for `linux/amd64` and `linux/arm64`. |
|
||||||
|
|
||||||
### Key Features
|
### ✨ Key Features
|
||||||
|
|
||||||
- **TLS-encrypted tunnel** between edge and hub (auto-generated self-signed cert or bring your own)
|
- 🔒 **TLS-encrypted tunnel** between edge and hub (auto-generated self-signed cert or bring your own)
|
||||||
- **Multiplexed streams** — thousands of client connections flow over a single tunnel
|
- 🔀 **Multiplexed streams** — thousands of client connections flow over a single tunnel
|
||||||
- **PROXY protocol v1** — SmartProxy sees the real client IP, not the tunnel IP
|
- 🌐 **PROXY protocol v1** — SmartProxy sees the real client IP, not the tunnel IP
|
||||||
- **Shared-secret authentication** — edges must present valid credentials to connect
|
- 🔑 **Shared-secret authentication** — edges must present valid credentials to connect
|
||||||
- **STUN-based public IP discovery** — the edge automatically discovers its public IP via Cloudflare STUN
|
- 🎫 **Connection tokens** — encode all connection details into a single opaque string
|
||||||
- **Auto-reconnect** with exponential backoff if the tunnel drops
|
- 📡 **STUN-based public IP discovery** — the edge automatically discovers its public IP via Cloudflare STUN
|
||||||
- **Event-driven** — both Hub and Edge extend `EventEmitter` for real-time monitoring
|
- 🔄 **Auto-reconnect** with exponential backoff if the tunnel drops
|
||||||
- **Rust core** — all frame encoding, TLS, and TCP proxying happen in native code for maximum throughput
|
- 📣 **Event-driven** — both Hub and Edge extend `EventEmitter` for real-time monitoring
|
||||||
|
- ⚡ **Rust core** — all frame encoding, TLS, and TCP proxying happen in native code for maximum throughput
|
||||||
|
|
||||||
## Usage
|
## 🚀 Usage
|
||||||
|
|
||||||
Both classes are imported from the package and communicate with the Rust binary under the hood. All you need to do is configure and start them.
|
Both classes are imported from the package and communicate with the Rust binary under the hood. All you need to do is configure and start them.
|
||||||
|
|
||||||
### Setting up the Hub (private cluster side)
|
### Setting Up the Hub (Private Cluster Side)
|
||||||
|
|
||||||
```typescript
|
```typescript
|
||||||
import { RemoteIngressHub } from '@serve.zone/remoteingress';
|
import { RemoteIngressHub } from '@serve.zone/remoteingress';
|
||||||
@@ -95,32 +100,45 @@ console.log(status);
|
|||||||
await hub.stop();
|
await hub.stop();
|
||||||
```
|
```
|
||||||
|
|
||||||
### Setting up the Edge (network edge side)
|
### Setting Up the Edge (Network Edge Side)
|
||||||
|
|
||||||
|
The edge can be configured in two ways: with an **opaque connection token** (recommended) or with explicit config fields.
|
||||||
|
|
||||||
|
#### Option A: Connection Token (Recommended)
|
||||||
|
|
||||||
|
A single token encodes all connection details — ideal for provisioning edges at scale:
|
||||||
|
|
||||||
```typescript
|
```typescript
|
||||||
import { RemoteIngressEdge } from '@serve.zone/remoteingress';
|
import { RemoteIngressEdge } from '@serve.zone/remoteingress';
|
||||||
|
|
||||||
const edge = new RemoteIngressEdge();
|
const edge = new RemoteIngressEdge();
|
||||||
|
|
||||||
// Listen for events
|
edge.on('tunnelConnected', () => console.log('Tunnel established'));
|
||||||
edge.on('tunnelConnected', () => {
|
edge.on('tunnelDisconnected', () => console.log('Tunnel lost — will auto-reconnect'));
|
||||||
console.log('Tunnel to hub established');
|
edge.on('publicIpDiscovered', ({ ip }) => console.log(`Public IP: ${ip}`));
|
||||||
});
|
|
||||||
edge.on('tunnelDisconnected', () => {
|
// Single token contains hubHost, hubPort, edgeId, and secret
|
||||||
console.log('Tunnel to hub lost — will auto-reconnect');
|
await edge.start({
|
||||||
});
|
token: 'eyJoIjoiaHViLmV4YW1wbGUuY29tIiwicCI6ODQ0MywiZSI6ImVkZ2UtbnljLTAxIiwicyI6InN1cGVyc2VjcmV0dG9rZW4xIn0',
|
||||||
edge.on('publicIpDiscovered', ({ ip }) => {
|
});
|
||||||
console.log(`Public IP: ${ip}`);
|
```
|
||||||
});
|
|
||||||
|
#### Option B: Explicit Config
|
||||||
|
|
||||||
|
```typescript
|
||||||
|
import { RemoteIngressEdge } from '@serve.zone/remoteingress';
|
||||||
|
|
||||||
|
const edge = new RemoteIngressEdge();
|
||||||
|
|
||||||
|
edge.on('tunnelConnected', () => console.log('Tunnel established'));
|
||||||
|
edge.on('tunnelDisconnected', () => console.log('Tunnel lost — will auto-reconnect'));
|
||||||
|
edge.on('publicIpDiscovered', ({ ip }) => console.log(`Public IP: ${ip}`));
|
||||||
|
|
||||||
// Start the edge — it connects to the hub and starts listening for clients
|
|
||||||
await edge.start({
|
await edge.start({
|
||||||
hubHost: 'hub.example.com', // hostname or IP of the hub
|
hubHost: 'hub.example.com', // hostname or IP of the hub
|
||||||
hubPort: 8443, // must match hub's tunnelPort (default: 8443)
|
hubPort: 8443, // must match hub's tunnelPort (default: 8443)
|
||||||
edgeId: 'edge-nyc-01', // unique edge identifier
|
edgeId: 'edge-nyc-01', // unique edge identifier
|
||||||
secret: 'supersecrettoken1', // must match the hub's allowed edge secret
|
secret: 'supersecrettoken1', // must match the hub's allowed edge secret
|
||||||
listenPorts: [80, 443], // public ports to accept TCP connections on
|
|
||||||
stunIntervalSecs: 300, // STUN refresh interval in seconds (optional)
|
|
||||||
});
|
});
|
||||||
|
|
||||||
// Check status at any time
|
// Check status at any time
|
||||||
@@ -138,9 +156,39 @@ console.log(edgeStatus);
|
|||||||
await edge.stop();
|
await edge.stop();
|
||||||
```
|
```
|
||||||
|
|
||||||
### API Reference
|
### 🎫 Connection Tokens
|
||||||
|
|
||||||
#### `RemoteIngressHub`
|
Connection tokens let you distribute a single opaque string instead of four separate config values. The hub operator generates the token; the edge operator just pastes it in.
|
||||||
|
|
||||||
|
```typescript
|
||||||
|
import { encodeConnectionToken, decodeConnectionToken } from '@serve.zone/remoteingress';
|
||||||
|
|
||||||
|
// Hub side: generate a token for a new edge
|
||||||
|
const token = encodeConnectionToken({
|
||||||
|
hubHost: 'hub.example.com',
|
||||||
|
hubPort: 8443,
|
||||||
|
edgeId: 'edge-nyc-01',
|
||||||
|
secret: 'supersecrettoken1',
|
||||||
|
});
|
||||||
|
console.log(token);
|
||||||
|
// => 'eyJoIjoiaHViLmV4YW1wbGUuY29tIiwi...'
|
||||||
|
|
||||||
|
// Edge side: inspect a token (optional — start() does this automatically)
|
||||||
|
const data = decodeConnectionToken(token);
|
||||||
|
console.log(data);
|
||||||
|
// {
|
||||||
|
// hubHost: 'hub.example.com',
|
||||||
|
// hubPort: 8443,
|
||||||
|
// edgeId: 'edge-nyc-01',
|
||||||
|
// secret: 'supersecrettoken1'
|
||||||
|
// }
|
||||||
|
```
|
||||||
|
|
||||||
|
Tokens are base64url-encoded (URL-safe, no padding) — safe to pass as environment variables, CLI arguments, or store in config files.
|
||||||
|
|
||||||
|
## 📖 API Reference
|
||||||
|
|
||||||
|
### `RemoteIngressHub`
|
||||||
|
|
||||||
| Method / Property | Description |
|
| Method / Property | Description |
|
||||||
|-------------------|-------------|
|
|-------------------|-------------|
|
||||||
@@ -152,18 +200,48 @@ await edge.stop();
|
|||||||
|
|
||||||
**Events:** `edgeConnected`, `edgeDisconnected`, `streamOpened`, `streamClosed`
|
**Events:** `edgeConnected`, `edgeDisconnected`, `streamOpened`, `streamClosed`
|
||||||
|
|
||||||
#### `RemoteIngressEdge`
|
### `RemoteIngressEdge`
|
||||||
|
|
||||||
| Method / Property | Description |
|
| Method / Property | Description |
|
||||||
|-------------------|-------------|
|
|-------------------|-------------|
|
||||||
| `start(config)` | Spawns the Rust binary, connects to the hub, and starts listening on the specified ports. |
|
| `start(config)` | Spawns the Rust binary and connects to the hub. Accepts `{ token: string }` or `IEdgeConfig`. |
|
||||||
| `stop()` | Gracefully shuts down the edge and kills the Rust process. |
|
| `stop()` | Gracefully shuts down the edge and kills the Rust process. |
|
||||||
| `getStatus()` | Returns current edge status including connection state, public IP, and active streams. |
|
| `getStatus()` | Returns current edge status including connection state, public IP, and active streams. |
|
||||||
| `running` | `boolean` — whether the Rust binary is alive. |
|
| `running` | `boolean` — whether the Rust binary is alive. |
|
||||||
|
|
||||||
**Events:** `tunnelConnected`, `tunnelDisconnected`, `publicIpDiscovered`
|
**Events:** `tunnelConnected`, `tunnelDisconnected`, `publicIpDiscovered`
|
||||||
|
|
||||||
### Wire Protocol
|
### Token Utilities
|
||||||
|
|
||||||
|
| Function | Description |
|
||||||
|
|----------|-------------|
|
||||||
|
| `encodeConnectionToken(data)` | Encodes `IConnectionTokenData` into a base64url token string. |
|
||||||
|
| `decodeConnectionToken(token)` | Decodes a token back into `IConnectionTokenData`. Throws on malformed or incomplete tokens. |
|
||||||
|
|
||||||
|
### Interfaces
|
||||||
|
|
||||||
|
```typescript
|
||||||
|
interface IHubConfig {
|
||||||
|
tunnelPort?: number; // default: 8443
|
||||||
|
targetHost?: string; // default: '127.0.0.1'
|
||||||
|
}
|
||||||
|
|
||||||
|
interface IEdgeConfig {
|
||||||
|
hubHost: string;
|
||||||
|
hubPort?: number; // default: 8443
|
||||||
|
edgeId: string;
|
||||||
|
secret: string;
|
||||||
|
}
|
||||||
|
|
||||||
|
interface IConnectionTokenData {
|
||||||
|
hubHost: string;
|
||||||
|
hubPort: number;
|
||||||
|
edgeId: string;
|
||||||
|
secret: string;
|
||||||
|
}
|
||||||
|
```
|
||||||
|
|
||||||
|
## 🔌 Wire Protocol
|
||||||
|
|
||||||
The tunnel uses a custom binary frame protocol over TLS:
|
The tunnel uses a custom binary frame protocol over TLS:
|
||||||
|
|
||||||
@@ -173,37 +251,64 @@ The tunnel uses a custom binary frame protocol over TLS:
|
|||||||
|
|
||||||
| Frame Type | Value | Direction | Purpose |
|
| Frame Type | Value | Direction | Purpose |
|
||||||
|------------|-------|-----------|---------|
|
|------------|-------|-----------|---------|
|
||||||
| `OPEN` | `0x01` | Edge -> Hub | Open a new stream; payload is PROXY v1 header |
|
| `OPEN` | `0x01` | Edge → Hub | Open a new stream; payload is PROXY v1 header |
|
||||||
| `DATA` | `0x02` | Edge -> Hub | Client data flowing upstream |
|
| `DATA` | `0x02` | Edge → Hub | Client data flowing upstream |
|
||||||
| `CLOSE` | `0x03` | Edge -> Hub | Client closed the connection |
|
| `CLOSE` | `0x03` | Edge → Hub | Client closed the connection |
|
||||||
| `DATA_BACK` | `0x04` | Hub -> Edge | Response data flowing downstream |
|
| `DATA_BACK` | `0x04` | Hub → Edge | Response data flowing downstream |
|
||||||
| `CLOSE_BACK` | `0x05` | Hub -> Edge | Upstream (SmartProxy) closed the connection |
|
| `CLOSE_BACK` | `0x05` | Hub → Edge | Upstream (SmartProxy) closed the connection |
|
||||||
|
|
||||||
Max payload size per frame: **16 MB**.
|
Max payload size per frame: **16 MB**.
|
||||||
|
|
||||||
### Example Scenarios
|
## 💡 Example Scenarios
|
||||||
|
|
||||||
1. **Expose a private Kubernetes cluster to the internet** — Deploy an Edge on a public VPS, configure your DNS to point to the VPS IP. The Edge tunnels all traffic to the Hub running inside the cluster, which hands it off to SmartProxy/DcRouter. Your cluster stays fully private — no public-facing ports needed.
|
### 1. Expose a Private Kubernetes Cluster to the Internet
|
||||||
|
|
||||||
2. **Multi-region edge ingress** — Run multiple Edges in different geographic regions (NYC, Frankfurt, Tokyo) all connecting to a single Hub. Use GeoDNS to route users to their nearest Edge. The Hub sees the real client IPs via PROXY protocol regardless of which edge they connected through.
|
Deploy an Edge on a public VPS, point your DNS to the VPS IP. The Edge tunnels all traffic to the Hub running inside the cluster, which hands it off to SmartProxy/DcRouter. Your cluster stays fully private — no public-facing ports needed.
|
||||||
|
|
||||||
3. **Secure API exposure** — Your backend runs on a private network with no direct internet access. An Edge on a minimal cloud instance acts as the only public entry point. TLS tunnel + shared-secret auth ensure only your authorized Edge can forward traffic.
|
### 2. Multi-Region Edge Ingress
|
||||||
|
|
||||||
|
Run multiple Edges in different geographic regions (NYC, Frankfurt, Tokyo) all connecting to a single Hub. Use GeoDNS to route users to their nearest Edge. The Hub sees the real client IPs via PROXY protocol regardless of which edge they connected through.
|
||||||
|
|
||||||
|
### 3. Secure API Exposure
|
||||||
|
|
||||||
|
Your backend runs on a private network with no direct internet access. An Edge on a minimal cloud instance acts as the only public entry point. TLS tunnel + shared-secret auth ensure only your authorized Edge can forward traffic.
|
||||||
|
|
||||||
|
### 4. Token-Based Edge Provisioning
|
||||||
|
|
||||||
|
Generate connection tokens on the hub side and distribute them to edge operators. Each edge only needs a single token string to connect — no manual configuration of host, port, ID, and secret.
|
||||||
|
|
||||||
|
```typescript
|
||||||
|
// Hub operator generates token
|
||||||
|
const token = encodeConnectionToken({
|
||||||
|
hubHost: 'hub.prod.example.com',
|
||||||
|
hubPort: 8443,
|
||||||
|
edgeId: 'edge-tokyo-01',
|
||||||
|
secret: 'generated-secret-abc123',
|
||||||
|
});
|
||||||
|
// Send `token` to the edge operator via secure channel
|
||||||
|
|
||||||
|
// Edge operator starts with just the token
|
||||||
|
const edge = new RemoteIngressEdge();
|
||||||
|
await edge.start({ token });
|
||||||
|
```
|
||||||
|
|
||||||
## License and Legal Information
|
## License and Legal Information
|
||||||
|
|
||||||
This repository contains open-source code that is licensed under the MIT License. A copy of the MIT License can be found in the [license](license) file within this repository.
|
This repository contains open-source code licensed under the MIT License. A copy of the license can be found in the [LICENSE](./LICENSE) file.
|
||||||
|
|
||||||
**Please note:** The MIT License does not grant permission to use the trade names, trademarks, service marks, or product names of the project, except as required for reasonable and customary use in describing the origin of the work and reproducing the content of the NOTICE file.
|
**Please note:** The MIT License does not grant permission to use the trade names, trademarks, service marks, or product names of the project, except as required for reasonable and customary use in describing the origin of the work and reproducing the content of the NOTICE file.
|
||||||
|
|
||||||
### Trademarks
|
### Trademarks
|
||||||
|
|
||||||
This project is owned and maintained by Task Venture Capital GmbH. The names and logos associated with Task Venture Capital GmbH and any related products or services are trademarks of Task Venture Capital GmbH and are not included within the scope of the MIT license granted herein. Use of these trademarks must comply with Task Venture Capital GmbH's Trademark Guidelines, and any usage must be approved in writing by Task Venture Capital GmbH.
|
This project is owned and maintained by Task Venture Capital GmbH. The names and logos associated with Task Venture Capital GmbH and any related products or services are trademarks of Task Venture Capital GmbH or third parties, and are not included within the scope of the MIT license granted herein.
|
||||||
|
|
||||||
|
Use of these trademarks must comply with Task Venture Capital GmbH's Trademark Guidelines or the guidelines of the respective third-party owners, and any usage must be approved in writing. Third-party trademarks used herein are the property of their respective owners and used only in a descriptive manner, e.g. for an implementation of an API or similar.
|
||||||
|
|
||||||
### Company Information
|
### Company Information
|
||||||
|
|
||||||
Task Venture Capital GmbH
|
Task Venture Capital GmbH
|
||||||
Registered at District court Bremen HRB 35230 HB, Germany
|
Registered at District Court Bremen HRB 35230 HB, Germany
|
||||||
|
|
||||||
For any legal inquiries or if you require further information, please contact us via email at hello@task.vc.
|
For any legal inquiries or further information, please contact us via email at hello@task.vc.
|
||||||
|
|
||||||
By using this repository, you acknowledge that you have read this section, agree to comply with its terms, and understand that the licensing of the code does not imply endorsement by Task Venture Capital GmbH of any derivative works.
|
By using this repository, you acknowledge that you have read this section, agree to comply with its terms, and understand that the licensing of the code does not imply endorsement by Task Venture Capital GmbH of any derivative works.
|
||||||
|
|||||||
@@ -301,6 +301,18 @@ async fn handle_request(
|
|||||||
serde_json::json!({ "ip": ip }),
|
serde_json::json!({ "ip": ip }),
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
EdgeEvent::PortsAssigned { listen_ports } => {
|
||||||
|
send_event(
|
||||||
|
"portsAssigned",
|
||||||
|
serde_json::json!({ "listenPorts": listen_ports }),
|
||||||
|
);
|
||||||
|
}
|
||||||
|
EdgeEvent::PortsUpdated { listen_ports } => {
|
||||||
|
send_event(
|
||||||
|
"portsUpdated",
|
||||||
|
serde_json::json!({ "listenPorts": listen_ports }),
|
||||||
|
);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
|
|||||||
@@ -1,15 +1,16 @@
|
|||||||
use std::collections::HashMap;
|
use std::collections::HashMap;
|
||||||
use std::sync::atomic::{AtomicU32, Ordering};
|
use std::sync::atomic::{AtomicU32, Ordering};
|
||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
use tokio::io::{AsyncReadExt, AsyncWriteExt};
|
use tokio::io::{AsyncBufReadExt, AsyncReadExt, AsyncWriteExt, BufReader};
|
||||||
use tokio::net::{TcpListener, TcpStream};
|
use tokio::net::{TcpListener, TcpStream};
|
||||||
use tokio::sync::{mpsc, Mutex, RwLock};
|
use tokio::sync::{mpsc, Mutex, RwLock};
|
||||||
|
use tokio::task::JoinHandle;
|
||||||
use tokio_rustls::TlsConnector;
|
use tokio_rustls::TlsConnector;
|
||||||
use serde::{Deserialize, Serialize};
|
use serde::{Deserialize, Serialize};
|
||||||
|
|
||||||
use remoteingress_protocol::*;
|
use remoteingress_protocol::*;
|
||||||
|
|
||||||
/// Edge configuration.
|
/// Edge configuration (hub-host + credentials only; ports come from hub).
|
||||||
#[derive(Debug, Clone, Deserialize, Serialize)]
|
#[derive(Debug, Clone, Deserialize, Serialize)]
|
||||||
#[serde(rename_all = "camelCase")]
|
#[serde(rename_all = "camelCase")]
|
||||||
pub struct EdgeConfig {
|
pub struct EdgeConfig {
|
||||||
@@ -17,8 +18,26 @@ pub struct EdgeConfig {
|
|||||||
pub hub_port: u16,
|
pub hub_port: u16,
|
||||||
pub edge_id: String,
|
pub edge_id: String,
|
||||||
pub secret: String,
|
pub secret: String,
|
||||||
pub listen_ports: Vec<u16>,
|
}
|
||||||
pub stun_interval_secs: Option<u64>,
|
|
||||||
|
/// Handshake config received from hub after authentication.
|
||||||
|
#[derive(Debug, Clone, Deserialize)]
|
||||||
|
#[serde(rename_all = "camelCase")]
|
||||||
|
struct HandshakeConfig {
|
||||||
|
listen_ports: Vec<u16>,
|
||||||
|
#[serde(default = "default_stun_interval")]
|
||||||
|
stun_interval_secs: u64,
|
||||||
|
}
|
||||||
|
|
||||||
|
fn default_stun_interval() -> u64 {
|
||||||
|
300
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Runtime config update received from hub via FRAME_CONFIG.
|
||||||
|
#[derive(Debug, Clone, Deserialize)]
|
||||||
|
#[serde(rename_all = "camelCase")]
|
||||||
|
struct ConfigUpdate {
|
||||||
|
listen_ports: Vec<u16>,
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Events emitted by the edge.
|
/// Events emitted by the edge.
|
||||||
@@ -30,6 +49,10 @@ pub enum EdgeEvent {
|
|||||||
TunnelDisconnected,
|
TunnelDisconnected,
|
||||||
#[serde(rename_all = "camelCase")]
|
#[serde(rename_all = "camelCase")]
|
||||||
PublicIpDiscovered { ip: String },
|
PublicIpDiscovered { ip: String },
|
||||||
|
#[serde(rename_all = "camelCase")]
|
||||||
|
PortsAssigned { listen_ports: Vec<u16> },
|
||||||
|
#[serde(rename_all = "camelCase")]
|
||||||
|
PortsUpdated { listen_ports: Vec<u16> },
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Edge status response.
|
/// Edge status response.
|
||||||
@@ -54,6 +77,7 @@ pub struct TunnelEdge {
|
|||||||
public_ip: Arc<RwLock<Option<String>>>,
|
public_ip: Arc<RwLock<Option<String>>>,
|
||||||
active_streams: Arc<AtomicU32>,
|
active_streams: Arc<AtomicU32>,
|
||||||
next_stream_id: Arc<AtomicU32>,
|
next_stream_id: Arc<AtomicU32>,
|
||||||
|
listen_ports: Arc<RwLock<Vec<u16>>>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl TunnelEdge {
|
impl TunnelEdge {
|
||||||
@@ -69,6 +93,7 @@ impl TunnelEdge {
|
|||||||
public_ip: Arc::new(RwLock::new(None)),
|
public_ip: Arc::new(RwLock::new(None)),
|
||||||
active_streams: Arc::new(AtomicU32::new(0)),
|
active_streams: Arc::new(AtomicU32::new(0)),
|
||||||
next_stream_id: Arc::new(AtomicU32::new(1)),
|
next_stream_id: Arc::new(AtomicU32::new(1)),
|
||||||
|
listen_ports: Arc::new(RwLock::new(Vec::new())),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -84,7 +109,7 @@ impl TunnelEdge {
|
|||||||
connected: *self.connected.read().await,
|
connected: *self.connected.read().await,
|
||||||
public_ip: self.public_ip.read().await.clone(),
|
public_ip: self.public_ip.read().await.clone(),
|
||||||
active_streams: self.active_streams.load(Ordering::Relaxed) as usize,
|
active_streams: self.active_streams.load(Ordering::Relaxed) as usize,
|
||||||
listen_ports: self.config.read().await.listen_ports.clone(),
|
listen_ports: self.listen_ports.read().await.clone(),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -100,6 +125,7 @@ impl TunnelEdge {
|
|||||||
let active_streams = self.active_streams.clone();
|
let active_streams = self.active_streams.clone();
|
||||||
let next_stream_id = self.next_stream_id.clone();
|
let next_stream_id = self.next_stream_id.clone();
|
||||||
let event_tx = self.event_tx.clone();
|
let event_tx = self.event_tx.clone();
|
||||||
|
let listen_ports = self.listen_ports.clone();
|
||||||
|
|
||||||
tokio::spawn(async move {
|
tokio::spawn(async move {
|
||||||
edge_main_loop(
|
edge_main_loop(
|
||||||
@@ -109,6 +135,7 @@ impl TunnelEdge {
|
|||||||
active_streams,
|
active_streams,
|
||||||
next_stream_id,
|
next_stream_id,
|
||||||
event_tx,
|
event_tx,
|
||||||
|
listen_ports,
|
||||||
shutdown_rx,
|
shutdown_rx,
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
@@ -124,6 +151,7 @@ impl TunnelEdge {
|
|||||||
}
|
}
|
||||||
*self.running.write().await = false;
|
*self.running.write().await = false;
|
||||||
*self.connected.write().await = false;
|
*self.connected.write().await = false;
|
||||||
|
self.listen_ports.write().await.clear();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -134,6 +162,7 @@ async fn edge_main_loop(
|
|||||||
active_streams: Arc<AtomicU32>,
|
active_streams: Arc<AtomicU32>,
|
||||||
next_stream_id: Arc<AtomicU32>,
|
next_stream_id: Arc<AtomicU32>,
|
||||||
event_tx: mpsc::UnboundedSender<EdgeEvent>,
|
event_tx: mpsc::UnboundedSender<EdgeEvent>,
|
||||||
|
listen_ports: Arc<RwLock<Vec<u16>>>,
|
||||||
mut shutdown_rx: mpsc::Receiver<()>,
|
mut shutdown_rx: mpsc::Receiver<()>,
|
||||||
) {
|
) {
|
||||||
let mut backoff_ms: u64 = 1000;
|
let mut backoff_ms: u64 = 1000;
|
||||||
@@ -148,6 +177,7 @@ async fn edge_main_loop(
|
|||||||
&active_streams,
|
&active_streams,
|
||||||
&next_stream_id,
|
&next_stream_id,
|
||||||
&event_tx,
|
&event_tx,
|
||||||
|
&listen_ports,
|
||||||
&mut shutdown_rx,
|
&mut shutdown_rx,
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
@@ -155,6 +185,7 @@ async fn edge_main_loop(
|
|||||||
*connected.write().await = false;
|
*connected.write().await = false;
|
||||||
let _ = event_tx.send(EdgeEvent::TunnelDisconnected);
|
let _ = event_tx.send(EdgeEvent::TunnelDisconnected);
|
||||||
active_streams.store(0, Ordering::Relaxed);
|
active_streams.store(0, Ordering::Relaxed);
|
||||||
|
listen_ports.write().await.clear();
|
||||||
|
|
||||||
match result {
|
match result {
|
||||||
EdgeLoopResult::Shutdown => break,
|
EdgeLoopResult::Shutdown => break,
|
||||||
@@ -182,6 +213,7 @@ async fn connect_to_hub_and_run(
|
|||||||
active_streams: &Arc<AtomicU32>,
|
active_streams: &Arc<AtomicU32>,
|
||||||
next_stream_id: &Arc<AtomicU32>,
|
next_stream_id: &Arc<AtomicU32>,
|
||||||
event_tx: &mpsc::UnboundedSender<EdgeEvent>,
|
event_tx: &mpsc::UnboundedSender<EdgeEvent>,
|
||||||
|
listen_ports: &Arc<RwLock<Vec<u16>>>,
|
||||||
shutdown_rx: &mut mpsc::Receiver<()>,
|
shutdown_rx: &mut mpsc::Receiver<()>,
|
||||||
) -> EdgeLoopResult {
|
) -> EdgeLoopResult {
|
||||||
// Build TLS connector that skips cert verification (auth is via secret)
|
// Build TLS connector that skips cert verification (auth is via secret)
|
||||||
@@ -220,12 +252,47 @@ async fn connect_to_hub_and_run(
|
|||||||
return EdgeLoopResult::Reconnect;
|
return EdgeLoopResult::Reconnect;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Read handshake response line from hub (JSON with initial config)
|
||||||
|
let mut buf_reader = BufReader::new(read_half);
|
||||||
|
let mut handshake_line = String::new();
|
||||||
|
match buf_reader.read_line(&mut handshake_line).await {
|
||||||
|
Ok(0) => {
|
||||||
|
log::error!("Hub rejected connection (EOF before handshake)");
|
||||||
|
return EdgeLoopResult::Reconnect;
|
||||||
|
}
|
||||||
|
Ok(_) => {}
|
||||||
|
Err(e) => {
|
||||||
|
log::error!("Failed to read handshake response: {}", e);
|
||||||
|
return EdgeLoopResult::Reconnect;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
let handshake: HandshakeConfig = match serde_json::from_str(handshake_line.trim()) {
|
||||||
|
Ok(h) => h,
|
||||||
|
Err(e) => {
|
||||||
|
log::error!("Invalid handshake response: {}", e);
|
||||||
|
return EdgeLoopResult::Reconnect;
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
log::info!(
|
||||||
|
"Handshake from hub: ports {:?}, stun_interval {}s",
|
||||||
|
handshake.listen_ports,
|
||||||
|
handshake.stun_interval_secs
|
||||||
|
);
|
||||||
|
|
||||||
*connected.write().await = true;
|
*connected.write().await = true;
|
||||||
let _ = event_tx.send(EdgeEvent::TunnelConnected);
|
let _ = event_tx.send(EdgeEvent::TunnelConnected);
|
||||||
log::info!("Connected to hub at {}", addr);
|
log::info!("Connected to hub at {}", addr);
|
||||||
|
|
||||||
|
// Store initial ports and emit event
|
||||||
|
*listen_ports.write().await = handshake.listen_ports.clone();
|
||||||
|
let _ = event_tx.send(EdgeEvent::PortsAssigned {
|
||||||
|
listen_ports: handshake.listen_ports.clone(),
|
||||||
|
});
|
||||||
|
|
||||||
// Start STUN discovery
|
// Start STUN discovery
|
||||||
let stun_interval = config.stun_interval_secs.unwrap_or(300);
|
let stun_interval = handshake.stun_interval_secs;
|
||||||
let public_ip_clone = public_ip.clone();
|
let public_ip_clone = public_ip.clone();
|
||||||
let event_tx_clone = event_tx.clone();
|
let event_tx_clone = event_tx.clone();
|
||||||
let stun_handle = tokio::spawn(async move {
|
let stun_handle = tokio::spawn(async move {
|
||||||
@@ -249,14 +316,112 @@ async fn connect_to_hub_and_run(
|
|||||||
// Shared tunnel writer
|
// Shared tunnel writer
|
||||||
let tunnel_writer = Arc::new(Mutex::new(write_half));
|
let tunnel_writer = Arc::new(Mutex::new(write_half));
|
||||||
|
|
||||||
// Start TCP listeners for each port
|
// Start TCP listeners for initial ports (hot-reloadable)
|
||||||
let mut listener_handles = Vec::new();
|
let mut port_listeners: HashMap<u16, JoinHandle<()>> = HashMap::new();
|
||||||
for &port in &config.listen_ports {
|
apply_port_config(
|
||||||
|
&handshake.listen_ports,
|
||||||
|
&mut port_listeners,
|
||||||
|
&tunnel_writer,
|
||||||
|
&client_writers,
|
||||||
|
active_streams,
|
||||||
|
next_stream_id,
|
||||||
|
&config.edge_id,
|
||||||
|
);
|
||||||
|
|
||||||
|
// Read frames from hub
|
||||||
|
let mut frame_reader = FrameReader::new(buf_reader);
|
||||||
|
let result = loop {
|
||||||
|
tokio::select! {
|
||||||
|
frame_result = frame_reader.next_frame() => {
|
||||||
|
match frame_result {
|
||||||
|
Ok(Some(frame)) => {
|
||||||
|
match frame.frame_type {
|
||||||
|
FRAME_DATA_BACK => {
|
||||||
|
let writers = client_writers.lock().await;
|
||||||
|
if let Some(tx) = writers.get(&frame.stream_id) {
|
||||||
|
let _ = tx.send(frame.payload).await;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
FRAME_CLOSE_BACK => {
|
||||||
|
let mut writers = client_writers.lock().await;
|
||||||
|
writers.remove(&frame.stream_id);
|
||||||
|
}
|
||||||
|
FRAME_CONFIG => {
|
||||||
|
if let Ok(update) = serde_json::from_slice::<ConfigUpdate>(&frame.payload) {
|
||||||
|
log::info!("Config update from hub: ports {:?}", update.listen_ports);
|
||||||
|
*listen_ports.write().await = update.listen_ports.clone();
|
||||||
|
let _ = event_tx.send(EdgeEvent::PortsUpdated {
|
||||||
|
listen_ports: update.listen_ports.clone(),
|
||||||
|
});
|
||||||
|
apply_port_config(
|
||||||
|
&update.listen_ports,
|
||||||
|
&mut port_listeners,
|
||||||
|
&tunnel_writer,
|
||||||
|
&client_writers,
|
||||||
|
active_streams,
|
||||||
|
next_stream_id,
|
||||||
|
&config.edge_id,
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
_ => {
|
||||||
|
log::warn!("Unexpected frame type {} from hub", frame.frame_type);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
Ok(None) => {
|
||||||
|
log::info!("Hub disconnected (EOF)");
|
||||||
|
break EdgeLoopResult::Reconnect;
|
||||||
|
}
|
||||||
|
Err(e) => {
|
||||||
|
log::error!("Hub frame error: {}", e);
|
||||||
|
break EdgeLoopResult::Reconnect;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
_ = shutdown_rx.recv() => {
|
||||||
|
break EdgeLoopResult::Shutdown;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
// Cleanup
|
||||||
|
stun_handle.abort();
|
||||||
|
for (_, h) in port_listeners.drain() {
|
||||||
|
h.abort();
|
||||||
|
}
|
||||||
|
|
||||||
|
result
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Apply a new port configuration: spawn listeners for added ports, abort removed ports.
|
||||||
|
fn apply_port_config(
|
||||||
|
new_ports: &[u16],
|
||||||
|
port_listeners: &mut HashMap<u16, JoinHandle<()>>,
|
||||||
|
tunnel_writer: &Arc<Mutex<tokio::io::WriteHalf<tokio_rustls::client::TlsStream<TcpStream>>>>,
|
||||||
|
client_writers: &Arc<Mutex<HashMap<u32, mpsc::Sender<Vec<u8>>>>>,
|
||||||
|
active_streams: &Arc<AtomicU32>,
|
||||||
|
next_stream_id: &Arc<AtomicU32>,
|
||||||
|
edge_id: &str,
|
||||||
|
) {
|
||||||
|
let new_set: std::collections::HashSet<u16> = new_ports.iter().copied().collect();
|
||||||
|
let old_set: std::collections::HashSet<u16> = port_listeners.keys().copied().collect();
|
||||||
|
|
||||||
|
// Remove ports no longer needed
|
||||||
|
for &port in old_set.difference(&new_set) {
|
||||||
|
if let Some(handle) = port_listeners.remove(&port) {
|
||||||
|
log::info!("Stopping listener on port {}", port);
|
||||||
|
handle.abort();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Add new ports
|
||||||
|
for &port in new_set.difference(&old_set) {
|
||||||
let tunnel_writer = tunnel_writer.clone();
|
let tunnel_writer = tunnel_writer.clone();
|
||||||
let client_writers = client_writers.clone();
|
let client_writers = client_writers.clone();
|
||||||
let active_streams = active_streams.clone();
|
let active_streams = active_streams.clone();
|
||||||
let next_stream_id = next_stream_id.clone();
|
let next_stream_id = next_stream_id.clone();
|
||||||
let edge_id = config.edge_id.clone();
|
let edge_id = edge_id.to_string();
|
||||||
|
|
||||||
let handle = tokio::spawn(async move {
|
let handle = tokio::spawn(async move {
|
||||||
let listener = match TcpListener::bind(("0.0.0.0", port)).await {
|
let listener = match TcpListener::bind(("0.0.0.0", port)).await {
|
||||||
@@ -299,55 +464,8 @@ async fn connect_to_hub_and_run(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
listener_handles.push(handle);
|
port_listeners.insert(port, handle);
|
||||||
}
|
}
|
||||||
|
|
||||||
// Read frames from hub
|
|
||||||
let mut frame_reader = FrameReader::new(read_half);
|
|
||||||
let result = loop {
|
|
||||||
tokio::select! {
|
|
||||||
frame_result = frame_reader.next_frame() => {
|
|
||||||
match frame_result {
|
|
||||||
Ok(Some(frame)) => {
|
|
||||||
match frame.frame_type {
|
|
||||||
FRAME_DATA_BACK => {
|
|
||||||
let writers = client_writers.lock().await;
|
|
||||||
if let Some(tx) = writers.get(&frame.stream_id) {
|
|
||||||
let _ = tx.send(frame.payload).await;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
FRAME_CLOSE_BACK => {
|
|
||||||
let mut writers = client_writers.lock().await;
|
|
||||||
writers.remove(&frame.stream_id);
|
|
||||||
}
|
|
||||||
_ => {
|
|
||||||
log::warn!("Unexpected frame type {} from hub", frame.frame_type);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
Ok(None) => {
|
|
||||||
log::info!("Hub disconnected (EOF)");
|
|
||||||
break EdgeLoopResult::Reconnect;
|
|
||||||
}
|
|
||||||
Err(e) => {
|
|
||||||
log::error!("Hub frame error: {}", e);
|
|
||||||
break EdgeLoopResult::Reconnect;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
_ = shutdown_rx.recv() => {
|
|
||||||
break EdgeLoopResult::Shutdown;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
};
|
|
||||||
|
|
||||||
// Cleanup
|
|
||||||
stun_handle.abort();
|
|
||||||
for h in listener_handles {
|
|
||||||
h.abort();
|
|
||||||
}
|
|
||||||
|
|
||||||
result
|
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn handle_client_connection(
|
async fn handle_client_connection(
|
||||||
|
|||||||
@@ -37,6 +37,24 @@ impl Default for HubConfig {
|
|||||||
pub struct AllowedEdge {
|
pub struct AllowedEdge {
|
||||||
pub id: String,
|
pub id: String,
|
||||||
pub secret: String,
|
pub secret: String,
|
||||||
|
#[serde(default)]
|
||||||
|
pub listen_ports: Vec<u16>,
|
||||||
|
pub stun_interval_secs: Option<u64>,
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Handshake response sent to edge after authentication.
|
||||||
|
#[derive(Debug, Clone, Serialize)]
|
||||||
|
#[serde(rename_all = "camelCase")]
|
||||||
|
struct HandshakeResponse {
|
||||||
|
listen_ports: Vec<u16>,
|
||||||
|
stun_interval_secs: u64,
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Configuration update pushed to a connected edge at runtime.
|
||||||
|
#[derive(Debug, Clone, Serialize)]
|
||||||
|
#[serde(rename_all = "camelCase")]
|
||||||
|
pub struct EdgeConfigUpdate {
|
||||||
|
pub listen_ports: Vec<u16>,
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Runtime status of a connected edge.
|
/// Runtime status of a connected edge.
|
||||||
@@ -75,7 +93,7 @@ pub struct HubStatus {
|
|||||||
/// The tunnel hub that accepts edge connections and demuxes streams to SmartProxy.
|
/// The tunnel hub that accepts edge connections and demuxes streams to SmartProxy.
|
||||||
pub struct TunnelHub {
|
pub struct TunnelHub {
|
||||||
config: RwLock<HubConfig>,
|
config: RwLock<HubConfig>,
|
||||||
allowed_edges: Arc<RwLock<HashMap<String, String>>>, // id -> secret
|
allowed_edges: Arc<RwLock<HashMap<String, AllowedEdge>>>,
|
||||||
connected_edges: Arc<Mutex<HashMap<String, ConnectedEdgeInfo>>>,
|
connected_edges: Arc<Mutex<HashMap<String, ConnectedEdgeInfo>>>,
|
||||||
event_tx: mpsc::UnboundedSender<HubEvent>,
|
event_tx: mpsc::UnboundedSender<HubEvent>,
|
||||||
event_rx: Mutex<Option<mpsc::UnboundedReceiver<HubEvent>>>,
|
event_rx: Mutex<Option<mpsc::UnboundedReceiver<HubEvent>>>,
|
||||||
@@ -86,6 +104,7 @@ pub struct TunnelHub {
|
|||||||
struct ConnectedEdgeInfo {
|
struct ConnectedEdgeInfo {
|
||||||
connected_at: u64,
|
connected_at: u64,
|
||||||
active_streams: Arc<Mutex<HashMap<u32, mpsc::Sender<Vec<u8>>>>>,
|
active_streams: Arc<Mutex<HashMap<u32, mpsc::Sender<Vec<u8>>>>>,
|
||||||
|
config_tx: mpsc::Sender<EdgeConfigUpdate>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl TunnelHub {
|
impl TunnelHub {
|
||||||
@@ -108,12 +127,35 @@ impl TunnelHub {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/// Update the list of allowed edges.
|
/// Update the list of allowed edges.
|
||||||
|
/// For any currently-connected edge whose ports changed, push a config update.
|
||||||
pub async fn update_allowed_edges(&self, edges: Vec<AllowedEdge>) {
|
pub async fn update_allowed_edges(&self, edges: Vec<AllowedEdge>) {
|
||||||
let mut map = self.allowed_edges.write().await;
|
let mut map = self.allowed_edges.write().await;
|
||||||
map.clear();
|
|
||||||
for edge in edges {
|
// Build new map
|
||||||
map.insert(edge.id, edge.secret);
|
let mut new_map = HashMap::new();
|
||||||
|
for edge in &edges {
|
||||||
|
new_map.insert(edge.id.clone(), edge.clone());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Push config updates to connected edges whose ports changed
|
||||||
|
let connected = self.connected_edges.lock().await;
|
||||||
|
for edge in &edges {
|
||||||
|
if let Some(info) = connected.get(&edge.id) {
|
||||||
|
// Check if ports changed compared to old config
|
||||||
|
let ports_changed = match map.get(&edge.id) {
|
||||||
|
Some(old) => old.listen_ports != edge.listen_ports,
|
||||||
|
None => true, // newly allowed edge that's already connected
|
||||||
|
};
|
||||||
|
if ports_changed {
|
||||||
|
let update = EdgeConfigUpdate {
|
||||||
|
listen_ports: edge.listen_ports.clone(),
|
||||||
|
};
|
||||||
|
let _ = info.config_tx.try_send(update);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
*map = new_map;
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Get the current hub status.
|
/// Get the current hub status.
|
||||||
@@ -208,13 +250,13 @@ impl TunnelHub {
|
|||||||
async fn handle_edge_connection(
|
async fn handle_edge_connection(
|
||||||
stream: TcpStream,
|
stream: TcpStream,
|
||||||
acceptor: TlsAcceptor,
|
acceptor: TlsAcceptor,
|
||||||
allowed: Arc<RwLock<HashMap<String, String>>>,
|
allowed: Arc<RwLock<HashMap<String, AllowedEdge>>>,
|
||||||
connected: Arc<Mutex<HashMap<String, ConnectedEdgeInfo>>>,
|
connected: Arc<Mutex<HashMap<String, ConnectedEdgeInfo>>>,
|
||||||
event_tx: mpsc::UnboundedSender<HubEvent>,
|
event_tx: mpsc::UnboundedSender<HubEvent>,
|
||||||
target_host: String,
|
target_host: String,
|
||||||
) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
|
) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
|
||||||
let tls_stream = acceptor.accept(stream).await?;
|
let tls_stream = acceptor.accept(stream).await?;
|
||||||
let (read_half, write_half) = tokio::io::split(tls_stream);
|
let (read_half, mut write_half) = tokio::io::split(tls_stream);
|
||||||
let mut buf_reader = BufReader::new(read_half);
|
let mut buf_reader = BufReader::new(read_half);
|
||||||
|
|
||||||
// Read auth line: "EDGE <edgeId> <secret>\n"
|
// Read auth line: "EDGE <edgeId> <secret>\n"
|
||||||
@@ -230,26 +272,36 @@ async fn handle_edge_connection(
|
|||||||
let edge_id = parts[1].to_string();
|
let edge_id = parts[1].to_string();
|
||||||
let secret = parts[2];
|
let secret = parts[2];
|
||||||
|
|
||||||
// Verify credentials
|
// Verify credentials and extract edge config
|
||||||
{
|
let (listen_ports, stun_interval_secs) = {
|
||||||
let edges = allowed.read().await;
|
let edges = allowed.read().await;
|
||||||
match edges.get(&edge_id) {
|
match edges.get(&edge_id) {
|
||||||
Some(expected) => {
|
Some(edge) => {
|
||||||
if !constant_time_eq(secret.as_bytes(), expected.as_bytes()) {
|
if !constant_time_eq(secret.as_bytes(), edge.secret.as_bytes()) {
|
||||||
return Err(format!("invalid secret for edge {}", edge_id).into());
|
return Err(format!("invalid secret for edge {}", edge_id).into());
|
||||||
}
|
}
|
||||||
|
(edge.listen_ports.clone(), edge.stun_interval_secs.unwrap_or(300))
|
||||||
}
|
}
|
||||||
None => {
|
None => {
|
||||||
return Err(format!("unknown edge {}", edge_id).into());
|
return Err(format!("unknown edge {}", edge_id).into());
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
};
|
||||||
|
|
||||||
log::info!("Edge {} authenticated", edge_id);
|
log::info!("Edge {} authenticated", edge_id);
|
||||||
let _ = event_tx.send(HubEvent::EdgeConnected {
|
let _ = event_tx.send(HubEvent::EdgeConnected {
|
||||||
edge_id: edge_id.clone(),
|
edge_id: edge_id.clone(),
|
||||||
});
|
});
|
||||||
|
|
||||||
|
// Send handshake response with initial config before frame protocol begins
|
||||||
|
let handshake = HandshakeResponse {
|
||||||
|
listen_ports: listen_ports.clone(),
|
||||||
|
stun_interval_secs,
|
||||||
|
};
|
||||||
|
let mut handshake_json = serde_json::to_string(&handshake)?;
|
||||||
|
handshake_json.push('\n');
|
||||||
|
write_half.write_all(handshake_json.as_bytes()).await?;
|
||||||
|
|
||||||
// Track this edge
|
// Track this edge
|
||||||
let streams: Arc<Mutex<HashMap<u32, mpsc::Sender<Vec<u8>>>>> =
|
let streams: Arc<Mutex<HashMap<u32, mpsc::Sender<Vec<u8>>>>> =
|
||||||
Arc::new(Mutex::new(HashMap::new()));
|
Arc::new(Mutex::new(HashMap::new()));
|
||||||
@@ -258,6 +310,9 @@ async fn handle_edge_connection(
|
|||||||
.unwrap_or_default()
|
.unwrap_or_default()
|
||||||
.as_secs();
|
.as_secs();
|
||||||
|
|
||||||
|
// Create config update channel
|
||||||
|
let (config_tx, mut config_rx) = mpsc::channel::<EdgeConfigUpdate>(16);
|
||||||
|
|
||||||
{
|
{
|
||||||
let mut edges = connected.lock().await;
|
let mut edges = connected.lock().await;
|
||||||
edges.insert(
|
edges.insert(
|
||||||
@@ -265,6 +320,7 @@ async fn handle_edge_connection(
|
|||||||
ConnectedEdgeInfo {
|
ConnectedEdgeInfo {
|
||||||
connected_at: now,
|
connected_at: now,
|
||||||
active_streams: streams.clone(),
|
active_streams: streams.clone(),
|
||||||
|
config_tx,
|
||||||
},
|
},
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
@@ -272,6 +328,23 @@ async fn handle_edge_connection(
|
|||||||
// Shared writer for sending frames back to edge
|
// Shared writer for sending frames back to edge
|
||||||
let write_half = Arc::new(Mutex::new(write_half));
|
let write_half = Arc::new(Mutex::new(write_half));
|
||||||
|
|
||||||
|
// Spawn task to forward config updates as FRAME_CONFIG frames
|
||||||
|
let config_writer = write_half.clone();
|
||||||
|
let config_edge_id = edge_id.clone();
|
||||||
|
let config_handle = tokio::spawn(async move {
|
||||||
|
while let Some(update) = config_rx.recv().await {
|
||||||
|
if let Ok(payload) = serde_json::to_vec(&update) {
|
||||||
|
let frame = encode_frame(0, FRAME_CONFIG, &payload);
|
||||||
|
let mut w = config_writer.lock().await;
|
||||||
|
if w.write_all(&frame).await.is_err() {
|
||||||
|
log::error!("Failed to send config update to edge {}", config_edge_id);
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
log::info!("Sent config update to edge {}: ports {:?}", config_edge_id, update.listen_ports);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
// Frame reading loop
|
// Frame reading loop
|
||||||
let mut frame_reader = FrameReader::new(buf_reader);
|
let mut frame_reader = FrameReader::new(buf_reader);
|
||||||
|
|
||||||
@@ -398,6 +471,7 @@ async fn handle_edge_connection(
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Cleanup
|
// Cleanup
|
||||||
|
config_handle.abort();
|
||||||
{
|
{
|
||||||
let mut edges = connected.lock().await;
|
let mut edges = connected.lock().await;
|
||||||
edges.remove(&edge_id);
|
edges.remove(&edge_id);
|
||||||
|
|||||||
@@ -6,6 +6,7 @@ pub const FRAME_DATA: u8 = 0x02;
|
|||||||
pub const FRAME_CLOSE: u8 = 0x03;
|
pub const FRAME_CLOSE: u8 = 0x03;
|
||||||
pub const FRAME_DATA_BACK: u8 = 0x04;
|
pub const FRAME_DATA_BACK: u8 = 0x04;
|
||||||
pub const FRAME_CLOSE_BACK: u8 = 0x05;
|
pub const FRAME_CLOSE_BACK: u8 = 0x05;
|
||||||
|
pub const FRAME_CONFIG: u8 = 0x06; // Hub -> Edge: configuration update
|
||||||
|
|
||||||
// Frame header size: 4 (stream_id) + 1 (type) + 4 (length) = 9 bytes
|
// Frame header size: 4 (stream_id) + 1 (type) + 4 (length) = 9 bytes
|
||||||
pub const FRAME_HEADER_SIZE: usize = 9;
|
pub const FRAME_HEADER_SIZE: usize = 9;
|
||||||
|
|||||||
@@ -3,6 +3,6 @@
|
|||||||
*/
|
*/
|
||||||
export const commitinfo = {
|
export const commitinfo = {
|
||||||
name: '@serve.zone/remoteingress',
|
name: '@serve.zone/remoteingress',
|
||||||
version: '3.1.0',
|
version: '3.2.0',
|
||||||
description: 'Edge ingress tunnel for DcRouter - accepts incoming TCP connections at network edge and tunnels them to DcRouter SmartProxy preserving client IP via PROXY protocol v1.'
|
description: 'Edge ingress tunnel for DcRouter - accepts incoming TCP connections at network edge and tunnels them to DcRouter SmartProxy preserving client IP via PROXY protocol v1.'
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -43,6 +43,7 @@ export interface IEdgeConfig {
|
|||||||
export class RemoteIngressEdge extends EventEmitter {
|
export class RemoteIngressEdge extends EventEmitter {
|
||||||
private bridge: InstanceType<typeof plugins.smartrust.RustBridge<TEdgeCommands>>;
|
private bridge: InstanceType<typeof plugins.smartrust.RustBridge<TEdgeCommands>>;
|
||||||
private started = false;
|
private started = false;
|
||||||
|
private statusInterval: ReturnType<typeof setInterval> | undefined;
|
||||||
|
|
||||||
constructor() {
|
constructor() {
|
||||||
super();
|
super();
|
||||||
@@ -79,6 +80,14 @@ export class RemoteIngressEdge extends EventEmitter {
|
|||||||
this.bridge.on('management:publicIpDiscovered', (data: { ip: string }) => {
|
this.bridge.on('management:publicIpDiscovered', (data: { ip: string }) => {
|
||||||
this.emit('publicIpDiscovered', data);
|
this.emit('publicIpDiscovered', data);
|
||||||
});
|
});
|
||||||
|
this.bridge.on('management:portsAssigned', (data: { listenPorts: number[] }) => {
|
||||||
|
console.log(`[RemoteIngressEdge] Ports assigned by hub: ${data.listenPorts.join(', ')}`);
|
||||||
|
this.emit('portsAssigned', data);
|
||||||
|
});
|
||||||
|
this.bridge.on('management:portsUpdated', (data: { listenPorts: number[] }) => {
|
||||||
|
console.log(`[RemoteIngressEdge] Ports updated by hub: ${data.listenPorts.join(', ')}`);
|
||||||
|
this.emit('portsUpdated', data);
|
||||||
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -113,12 +122,30 @@ export class RemoteIngressEdge extends EventEmitter {
|
|||||||
});
|
});
|
||||||
|
|
||||||
this.started = true;
|
this.started = true;
|
||||||
|
|
||||||
|
// Start periodic status logging
|
||||||
|
this.statusInterval = setInterval(async () => {
|
||||||
|
try {
|
||||||
|
const status = await this.getStatus();
|
||||||
|
console.log(
|
||||||
|
`[RemoteIngressEdge] Status: connected=${status.connected}, ` +
|
||||||
|
`streams=${status.activeStreams}, ports=[${status.listenPorts.join(',')}], ` +
|
||||||
|
`publicIp=${status.publicIp ?? 'unknown'}`
|
||||||
|
);
|
||||||
|
} catch {
|
||||||
|
// Bridge may be shutting down
|
||||||
|
}
|
||||||
|
}, 60_000);
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Stop the edge and kill the Rust process.
|
* Stop the edge and kill the Rust process.
|
||||||
*/
|
*/
|
||||||
public async stop(): Promise<void> {
|
public async stop(): Promise<void> {
|
||||||
|
if (this.statusInterval) {
|
||||||
|
clearInterval(this.statusInterval);
|
||||||
|
this.statusInterval = undefined;
|
||||||
|
}
|
||||||
if (this.started) {
|
if (this.started) {
|
||||||
try {
|
try {
|
||||||
await this.bridge.sendCommand('stopEdge', {} as Record<string, never>);
|
await this.bridge.sendCommand('stopEdge', {} as Record<string, never>);
|
||||||
|
|||||||
@@ -20,7 +20,7 @@ type THubCommands = {
|
|||||||
};
|
};
|
||||||
updateAllowedEdges: {
|
updateAllowedEdges: {
|
||||||
params: {
|
params: {
|
||||||
edges: Array<{ id: string; secret: string }>;
|
edges: Array<{ id: string; secret: string; listenPorts?: number[]; stunIntervalSecs?: number }>;
|
||||||
};
|
};
|
||||||
result: { updated: boolean };
|
result: { updated: boolean };
|
||||||
};
|
};
|
||||||
@@ -122,7 +122,7 @@ export class RemoteIngressHub extends EventEmitter {
|
|||||||
/**
|
/**
|
||||||
* Update the list of allowed edges that can connect to this hub.
|
* Update the list of allowed edges that can connect to this hub.
|
||||||
*/
|
*/
|
||||||
public async updateAllowedEdges(edges: Array<{ id: string; secret: string }>): Promise<void> {
|
public async updateAllowedEdges(edges: Array<{ id: string; secret: string; listenPorts?: number[]; stunIntervalSecs?: number }>): Promise<void> {
|
||||||
await this.bridge.sendCommand('updateAllowedEdges', { edges });
|
await this.bridge.sendCommand('updateAllowedEdges', { edges });
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user