portal-tunnel: add config-driven multi-service tunneling

yoonhyunwoo committed Nov 13, 2025 at 16:44 UTC cbf43e367389f5fe69916e8cb46da656d9849781
5 files changed +378 -88
cmd/portal-tunnel/config.go new
+180
@@ -0,0 +1,180 @@
1 +package main
2 +
3 +import (
4 + "fmt"
5 + "os"
6 + "strings"
7 +
8 + "gopkg.in/yaml.v3"
9 +)
10 +
11 +// RelayConfig describes a named relay endpoint and its bootstrap URLs.
12 +type RelayConfig struct {
13 + Name string `yaml:"name"`
14 + URLs []string `yaml:"urls"`
15 +}
16 +
17 +// ServiceConfig describes a local service exposed through the tunnel.
18 +type ServiceConfig struct {
19 + Name string `yaml:"name"`
20 + RelayPreference []string `yaml:"relayPreference"`
21 + Target string `yaml:"target"`
22 + Protocols []string `yaml:"protocols"`
23 +}
24 +
25 +// TunnelConfig represents the YAML configuration schema for portal-tunnel.
26 +type TunnelConfig struct {
27 + Relays []RelayConfig `yaml:"relays"`
28 + Services []ServiceConfig `yaml:"services"`
29 +}
30 +
31 +// RelayDirectory provides lookup helpers for relay definitions.
32 +type RelayDirectory struct {
33 + entries map[string]RelayConfig
34 +}
35 +
36 +// LoadConfig reads the YAML file at path, parses it into TunnelConfig, and validates it.
37 +func LoadConfig(path string) (*TunnelConfig, error) {
38 + data, err := os.ReadFile(path)
39 + if err != nil {
40 + return nil, fmt.Errorf("read config: %w", err)
41 + }
42 +
43 + var cfg TunnelConfig
44 + if err := yaml.Unmarshal(data, &cfg); err != nil {
45 + return nil, fmt.Errorf("parse config: %w", err)
46 + }
47 + cfg.applyDefaults()
48 +
49 + if err := cfg.validate(); err != nil {
50 + return nil, err
51 + }
52 +
53 + return &cfg, nil
54 +}
55 +
56 +// NewRelayDirectory builds a lookup structure for relay definitions.
57 +func NewRelayDirectory(relays []RelayConfig) *RelayDirectory {
58 + idx := make(map[string]RelayConfig, len(relays))
59 + for _, relay := range relays {
60 + idx[relay.Name] = relay
61 + }
62 + return &RelayDirectory{entries: idx}
63 +}
64 +
65 +// BootstrapServers aggregates URLs for the given relay preference list.
66 +// Preference order is preserved and duplicate URLs are removed.
67 +func (rd *RelayDirectory) BootstrapServers(preferences []string) ([]string, error) {
68 + if len(preferences) == 0 {
69 + return nil, fmt.Errorf("relayPreference must contain at least one relay name")
70 + }
71 +
72 + seen := map[string]struct{}{}
73 + var servers []string
74 + for _, relayName := range preferences {
75 + relayName = strings.TrimSpace(relayName)
76 + if relayName == "" {
77 + continue
78 + }
79 + relay, ok := rd.entries[relayName]
80 + if !ok {
81 + continue
82 + }
83 + for _, url := range relay.URLs {
84 + url = strings.TrimSpace(url)
85 + if url == "" {
86 + continue
87 + }
88 + if _, exists := seen[url]; exists {
89 + continue
90 + }
91 + seen[url] = struct{}{}
92 + servers = append(servers, url)
93 + }
94 + }
95 +
96 + if len(servers) == 0 {
97 + return nil, fmt.Errorf("no bootstrap servers resolved for relayPreference %v", preferences)
98 + }
99 +
100 + return servers, nil
101 +}
102 +
103 +func (cfg *TunnelConfig) validate() error {
104 + var errs []string
105 +
106 + relayIdx := map[string]RelayConfig{}
107 + if len(cfg.Relays) == 0 {
108 + errs = append(errs, "at least one relay must be defined")
109 + }
110 + for i, relay := range cfg.Relays {
111 + prefix := fmt.Sprintf("relays[%d]", i)
112 + name := strings.TrimSpace(relay.Name)
113 + if name == "" {
114 + errs = append(errs, fmt.Sprintf("%s: name is required", prefix))
115 + } else {
116 + if _, exists := relayIdx[name]; exists {
117 + errs = append(errs, fmt.Sprintf("%s: duplicate relay name %q", prefix, name))
118 + } else {
119 + relayIdx[name] = relay
120 + }
121 + }
122 +
123 + if len(relay.URLs) == 0 {
124 + errs = append(errs, fmt.Sprintf("%s: at least one url is required", prefix))
125 + }
126 + for j, url := range relay.URLs {
127 + if strings.TrimSpace(url) == "" {
128 + errs = append(errs, fmt.Sprintf("%s.urls[%d]: url cannot be empty", prefix, j))
129 + }
130 + }
131 + }
132 +
133 + if len(cfg.Services) == 0 {
134 + errs = append(errs, "at least one service must be defined")
135 + }
136 + for i, service := range cfg.Services {
137 + prefix := fmt.Sprintf("services[%d]", i)
138 + name := strings.TrimSpace(service.Name)
139 + if name == "" {
140 + errs = append(errs, fmt.Sprintf("%s: name is required", prefix))
141 + }
142 + target := strings.TrimSpace(service.Target)
143 + if target == "" {
144 + errs = append(errs, fmt.Sprintf("%s: target is required", prefix))
145 + }
146 + for j, proto := range service.Protocols {
147 + if strings.TrimSpace(proto) == "" {
148 + errs = append(errs, fmt.Sprintf("%s.protocols[%d]: protocol cannot be empty", prefix, j))
149 + }
150 + }
151 + if len(service.RelayPreference) == 0 {
152 + errs = append(errs, fmt.Sprintf("%s: relayPreference must list at least one relay name", prefix))
153 + }
154 + for j, relayName := range service.RelayPreference {
155 + relayName = strings.TrimSpace(relayName)
156 + if relayName == "" {
157 + errs = append(errs, fmt.Sprintf("%s.relayPreference[%d]: relay name cannot be empty", prefix, j))
158 + continue
159 + }
160 + if _, exists := relayIdx[relayName]; !exists {
161 + errs = append(errs, fmt.Sprintf("%s.relayPreference[%d]: relay %q is not defined", prefix, j, relayName))
162 + }
163 + }
164 + }
165 +
166 + if len(errs) > 0 {
167 + return fmt.Errorf("invalid config:\n - %s", strings.Join(errs, "\n - "))
168 + }
169 +
170 + return nil
171 +}
172 +
173 +func (cfg *TunnelConfig) applyDefaults() {
174 + const defaultProtocol = "http/1.1"
175 + for i := range cfg.Services {
176 + if len(cfg.Services[i].Protocols) == 0 {
177 + cfg.Services[i].Protocols = []string{defaultProtocol}
178 + }
179 + }
180 +}
cmd/portal-tunnel/config.yaml.example new
+25
@@ -0,0 +1,25 @@
1 +relays:
2 + - name: thumbgo
3 + urls:
4 + - wss://portal.thumbgo.kr/relay
5 + - name: gosuda
6 + urls:
7 + - wss://portal.gosuda.org/relay
8 +
9 +services:
10 + - name: test-thumbgo
11 + relayPreference:
12 + - thumbgo
13 + - gosuda
14 + target: localhost:1013
15 + protocols:
16 + - http/1.1
17 + - h2
18 + - name: test-thumbgo2
19 + relayPreference:
20 + - thumbgo
21 + - gosuda
22 + target: localhost:1014
23 + protocols:
24 + - http/1.1
25 + - h2
cmd/portal-tunnel/main.go
+171 -88
@@ -8,6 +8,8 @@ import (
8 "net"
9 "os"
10 "os/signal"
11 + "strings"
12 + "sync"
13 "syscall"
14 "time"
15
@@ -16,12 +18,16 @@ import (
18 )
19
20 var (
19 - flagRelayURL string
20 - flagHost string
21 - flagPort string
22 - flagName string
21 + flagConfigPath string
22 + flagService string
23 )
24
25 +type serviceContext struct {
26 + Name string
27 + LocalAddr string
28 + RelayServers []string
29 +}
30 +
31 func main() {
32 if len(os.Args) < 2 {
33 printTunnelUsage()
@@ -31,10 +37,8 @@ func main() {
37 switch os.Args[1] {
38 case "expose":
39 fs := flag.NewFlagSet("expose", flag.ExitOnError)
34 - fs.StringVar(&flagRelayURL, "relay", "ws://localhost:4017/relay", "Portal relay server URL")
35 - fs.StringVar(&flagHost, "host", "localhost", "Local host to proxy to")
36 - fs.StringVar(&flagPort, "port", "4018", "Local port to proxy to")
37 - fs.StringVar(&flagName, "name", "", "Service name (will be generated if not provided)")
40 + fs.StringVar(&flagConfigPath, "config", "", "Path to portal-tunnel config file")
41 + fs.StringVar(&flagService, "service", "", "Specific service name to expose (defaults to first entry)")
42 _ = fs.Parse(os.Args[2:])
43
44 if err := runExpose(); err != nil {
@@ -53,124 +57,92 @@ func printTunnelUsage() {
57 fmt.Println("portal-tunnel — Expose local services through Portal relay")
58 fmt.Println()
59 fmt.Println("Usage:")
56 - fmt.Println(" portal-tunnel expose [port PORT] [--relay URL] [--name NAME] [--host HOST]")
60 + fmt.Println(" portal-tunnel expose --config <file> [--service <name>]")
61 }
62
63 func runExpose() error {
60 - localAddr := net.JoinHostPort(flagHost, flagPort)
61 -
62 - // Always wait until the local service is available
63 - log.Info().Msgf("Waiting for local service at %s (interval=%v)...", localAddr, time.Second)
64 - if err := waitForLocalService(localAddr, 0, time.Second); err != nil {
65 - return err
66 - }
67 - log.Info().Msgf("✓ Local service is reachable at %s", localAddr)
68 -
69 - // Create credential
70 - cred := sdk.NewCredential()
71 - leaseID := cred.ID()
72 -
73 - // Use provided name or generate from lease ID
74 - if flagName == "" {
75 - flagName = fmt.Sprintf("tunnel-%s", leaseID[:8])
64 + if flagConfigPath == "" {
65 + return fmt.Errorf("--config is required")
66 }
67
78 - log.Info().Msgf("Starting Portal Tunnel...")
79 - log.Info().Msgf(" Local: %s", localAddr)
80 - log.Info().Msgf(" Relay: %s", flagRelayURL)
81 - log.Info().Msgf(" Name: %s", flagName)
82 - log.Info().Msgf(" Lease ID: %s", leaseID)
83 -
84 - // Create SDK client
85 - client, err := sdk.NewClient(func(c *sdk.RDClientConfig) {
86 - c.BootstrapServers = []string{flagRelayURL}
87 - })
68 + cfg, err := LoadConfig(flagConfigPath)
69 if err != nil {
89 - return fmt.Errorf("failed to connect to relay: %w", err)
70 + return fmt.Errorf("load config: %w", err)
71 }
91 - defer client.Close()
92 -
93 - // Register listener
94 - listener, err := client.Listen(cred, flagName, []string{"http/1.1", "h2"})
72 + services, err := selectServices(cfg, flagService)
73 if err != nil {
96 - return fmt.Errorf("failed to register service: %w", err)
74 + return err
75 }
98 - defer listener.Close()
76
100 - log.Info().Msg("")
101 - log.Info().Msg("=== Service is now publicly accessible ===")
102 - log.Info().Msg("Access via:")
103 - log.Info().Msgf("- Name: /peer/%s", flagName)
104 - log.Info().Msgf("- Lease ID: /peer/%s", leaseID)
105 - relayHost := extractHost(flagRelayURL)
106 - log.Info().Msgf("- Example: http://%s/peer/%s", relayHost, flagName)
107 - log.Info().Msg("")
108 - log.Info().Msg("Press Ctrl+C to stop...")
109 - log.Info().Msg("")
110 -
111 - // Handle connections
77 + relayDir := NewRelayDirectory(cfg.Relays)
78 ctx, cancel := context.WithCancel(context.Background())
79 defer cancel()
80
81 // Graceful shutdown
82 sigCh := make(chan os.Signal, 1)
83 signal.Notify(sigCh, syscall.SIGINT, syscall.SIGTERM)
118 -
84 go func() {
85 <-sigCh
86 log.Info().Msg("")
122 - log.Info().Msg("Shutting down tunnel...")
87 + log.Info().Msg("Shutting down tunnels...")
88 cancel()
89 }()
90
126 - // Accept connections and proxy them
127 - connCount := 0
128 - for {
129 - select {
130 - case <-ctx.Done():
131 - log.Info().Msg("Tunnel stopped")
132 - return nil
133 - default:
134 - }
91 + errCh := make(chan error, len(services))
92 + var wg sync.WaitGroup
93
136 - relayConn, err := listener.Accept()
137 - if err != nil {
138 - // Check if context was cancelled
139 - select {
140 - case <-ctx.Done():
141 - return nil
142 - default:
143 - log.Error().Err(err).Msg("Failed to accept connection")
144 - continue
94 + for _, svc := range services {
95 + service := svc
96 + wg.Add(1)
97 + go func() {
98 + defer wg.Done()
99 + if err := runServiceTunnel(ctx, relayDir, service); err != nil {
100 + errCh <- err
101 }
146 - }
102 + }()
103 + }
104
148 - connCount++
149 - currentConnCount := connCount
150 - log.Info().Msgf("→ [#%d] New connection from %s", currentConnCount, relayConn.RemoteAddr())
105 + doneCh := make(chan struct{})
106 + go func() {
107 + wg.Wait()
108 + close(doneCh)
109 + }()
110
152 - // Handle connection in goroutine
153 - go func(relayConn net.Conn, connNum int) {
154 - if err := proxyConnection(relayConn, localAddr, connNum); err != nil {
155 - log.Error().Err(err).Int("conn", connNum).Msg("Proxy error")
156 - }
157 - log.Info().Msgf("← [#%d] Connection closed", connNum)
158 - }(relayConn, currentConnCount)
111 + select {
112 + case err := <-errCh:
113 + cancel()
114 + <-doneCh
115 + return err
116 + case <-ctx.Done():
117 + <-doneCh
118 + log.Info().Msg("Tunnel stopped")
119 + return nil
120 + case <-doneCh:
121 + return nil
122 }
123 }
124
162 -func proxyConnection(relayConn net.Conn, localAddr string, connNum int) error {
125 +func proxyConnection(ctx context.Context, svcCtx *serviceContext, relayConn net.Conn, connNum int) error {
126 defer relayConn.Close()
127
128 // Connect to local service
166 - localConn, err := net.Dial("tcp", localAddr)
129 + localConn, err := net.Dial("tcp", svcCtx.LocalAddr)
130 if err != nil {
168 - return fmt.Errorf("failed to connect to local service: %w", err)
131 + return fmt.Errorf("failed to connect to local service %s: %w", svcCtx.LocalAddr, err)
132 }
133 defer localConn.Close()
134
135 // Bidirectional copy
136 errCh := make(chan error, 2)
137 + cancelCopy := make(chan struct{})
138 + go func() {
139 + select {
140 + case <-ctx.Done():
141 + relayConn.Close()
142 + localConn.Close()
143 + case <-cancelCopy:
144 + }
145 + }()
146
147 // Relay -> Local
148 go func() {
@@ -190,6 +162,7 @@ func proxyConnection(relayConn net.Conn, localAddr string, connNum int) error {
162 // Close both connections to stop the other goroutine
163 relayConn.Close()
164 localConn.Close()
165 + close(cancelCopy)
166
167 // Wait for other goroutine
168 <-errCh
@@ -197,6 +170,116 @@ func proxyConnection(relayConn net.Conn, localAddr string, connNum int) error {
170 return err
171 }
172
173 +func runServiceTunnel(ctx context.Context, relayDir *RelayDirectory, service *ServiceConfig) error {
174 + localAddr := service.Target
175 + bootstrapServers, err := relayDir.BootstrapServers(service.RelayPreference)
176 + if err != nil {
177 + return fmt.Errorf("service %s: resolve relay servers: %w", service.Name, err)
178 + }
179 +
180 + svcCtx := &serviceContext{
181 + Name: service.Name,
182 + LocalAddr: localAddr,
183 + RelayServers: bootstrapServers,
184 + }
185 +
186 + log.Info().Str("service", service.Name).Msgf("Waiting for local service at %s (interval=%v)...", localAddr, time.Second)
187 + if err := waitForLocalService(localAddr, 0, time.Second); err != nil {
188 + return fmt.Errorf("service %s: %w", service.Name, err)
189 + }
190 + log.Info().Str("service", service.Name).Msgf("✓ Local service is reachable at %s", localAddr)
191 +
192 + cred := sdk.NewCredential()
193 + leaseID := cred.ID()
194 +
195 + log.Info().Str("service", service.Name).Msgf("Starting Portal Tunnel (config=%s)...", flagConfigPath)
196 + log.Info().Str("service", service.Name).Msgf(" Local: %s", localAddr)
197 + log.Info().Str("service", service.Name).Msgf(" Relays: %s", strings.Join(bootstrapServers, ", "))
198 + log.Info().Str("service", service.Name).Msgf(" Lease ID: %s", leaseID)
199 +
200 + client, err := sdk.NewClient(func(c *sdk.RDClientConfig) {
201 + c.BootstrapServers = bootstrapServers
202 + })
203 + if err != nil {
204 + return fmt.Errorf("service %s: failed to connect to relay: %w", service.Name, err)
205 + }
206 + defer client.Close()
207 +
208 + listener, err := client.Listen(cred, service.Name, service.Protocols)
209 + if err != nil {
210 + return fmt.Errorf("service %s: failed to register service: %w", service.Name, err)
211 + }
212 + defer listener.Close()
213 +
214 + go func() {
215 + <-ctx.Done()
216 + _ = listener.Close()
217 + }()
218 +
219 + log.Info().Str("service", service.Name).Msg("")
220 + log.Info().Str("service", service.Name).Msg("=== Service is now publicly accessible ===")
221 + log.Info().Str("service", service.Name).Msg("Access via:")
222 + log.Info().Str("service", service.Name).Msgf("- Name: /peer/%s", service.Name)
223 + log.Info().Str("service", service.Name).Msgf("- Lease ID: /peer/%s", leaseID)
224 + relayHost := extractHost(bootstrapServers[0])
225 + log.Info().Str("service", service.Name).Msgf("- Example: http://%s/peer/%s", relayHost, service.Name)
226 + log.Info().Str("service", service.Name).Msg("")
227 +
228 + connCount := 0
229 + var connWG sync.WaitGroup
230 + defer connWG.Wait()
231 + for {
232 + select {
233 + case <-ctx.Done():
234 + return nil
235 + default:
236 + }
237 +
238 + relayConn, err := listener.Accept()
239 + if err != nil {
240 + select {
241 + case <-ctx.Done():
242 + return nil
243 + default:
244 + log.Error().Str("service", service.Name).Err(err).Msg("Failed to accept connection")
245 + continue
246 + }
247 + }
248 +
249 + connCount++
250 + currentConnCount := connCount
251 + log.Info().Str("service", service.Name).Msgf("→ [#%d] New connection from %s", currentConnCount, relayConn.RemoteAddr())
252 +
253 + connWG.Add(1)
254 + go func(relayConn net.Conn, connNum int) {
255 + defer connWG.Done()
256 + if err := proxyConnection(ctx, svcCtx, relayConn, connNum); err != nil {
257 + log.Error().Str("service", service.Name).Err(err).Int("conn", connNum).Msg("Proxy error")
258 + }
259 + log.Info().Str("service", service.Name).Msgf("← [#%d] Connection closed", connNum)
260 + }(relayConn, currentConnCount)
261 + }
262 +}
263 +
264 +func selectServices(cfg *TunnelConfig, name string) ([]*ServiceConfig, error) {
265 + if len(cfg.Services) == 0 {
266 + return nil, fmt.Errorf("config has no services")
267 + }
268 + if name == "" {
269 + services := make([]*ServiceConfig, len(cfg.Services))
270 + for i := range cfg.Services {
271 + services[i] = &cfg.Services[i]
272 + }
273 + return services, nil
274 + }
275 + for i := range cfg.Services {
276 + if cfg.Services[i].Name == name {
277 + return []*ServiceConfig{&cfg.Services[i]}, nil
278 + }
279 + }
280 + return nil, fmt.Errorf("service %q not found in config", name)
281 +}
282 +
283 func extractHost(wsURL string) string {
284 // Simple extraction: ws://host:port/path -> host:port
285 // Remove ws:// or wss://
go.mod
+1
@@ -19,4 +19,5 @@ require (
19 github.com/mattn/go-isatty v0.0.20 // indirect
20 golang.org/x/sys v0.37.0 // indirect
21 golang.org/x/text v0.30.0 // indirect
22 + gopkg.in/yaml.v3 v3.0.1 // indirect
23 )
go.sum
+1
@@ -48,4 +48,5 @@ golang.org/x/text v0.30.0/go.mod h1:yDdHFIX9t+tORqspjENWgzaCVXgk0yYnYuSZ8UzzBVM=
48 google.golang.org/protobuf v1.36.10 h1:AYd7cD/uASjIL6Q9LiTjz8JLcrh/88q5UObnmY3aOOE=
49 google.golang.org/protobuf v1.36.10/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco=
50 gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
51 +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
52 gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=