똥코드제거
yoonhyunwoo committed
Nov 17, 2025 at 18:58 UTC
93ba0a9c190fc3fbd9b3433ca186c580d6619044
2 files changed
+30
-124
cmd/portal-tunnel/config.go
+9
-4
@@ -8,6 +8,8 @@ import (
8
"gopkg.in/yaml.v3"
9
)
10
11
+var defaultProtocols = []string{"http/1.1", "h2"}
12
+
13
// RelayConfig describes a named relay endpoint and its bootstrap URLs.
14
type RelayConfig struct {
15
Name string `yaml:"name"`
@@ -171,10 +173,13 @@ func (cfg *TunnelConfig) validate() error {
173
}
174
175
func (cfg *TunnelConfig) applyDefaults() {
174
- const defaultProtocol = "http/1.1"
176
for i := range cfg.Services {
176
- if len(cfg.Services[i].Protocols) == 0 {
177
- cfg.Services[i].Protocols = []string{defaultProtocol}
178
- }
177
+ applyServiceDefaults(&cfg.Services[i])
178
+ }
179
+}
180
+
181
+func applyServiceDefaults(svc *ServiceConfig) {
182
+ if len(svc.Protocols) == 0 {
183
+ svc.Protocols = append([]string(nil), defaultProtocols...)
184
}
185
}
cmd/portal-tunnel/main.go
+21
-120
@@ -11,7 +11,6 @@ import (
11
"strings"
12
"sync"
13
"syscall"
14
- "time"
14
15
"github.com/rs/zerolog/log"
16
"gosuda.org/portal/sdk"
@@ -25,12 +24,6 @@ var (
24
flagName string
25
)
26
28
-type serviceContext struct {
29
- Name string
30
- LocalAddr string
31
- RelayServers []string
32
-}
33
-
27
func main() {
28
if len(os.Args) < 2 {
29
printTunnelUsage()
@@ -64,7 +57,7 @@ func printTunnelUsage() {
57
fmt.Println()
58
fmt.Println("Usage:")
59
fmt.Println(" portal-tunnel expose --config <file>")
67
- fmt.Println(" portal-tunnel expose [--relay URL,URL...]] [--host HOST] [--port PORT] [--name NAME]")
60
+ fmt.Println(" portal-tunnel expose [--relay URL1,URL2] [--host HOST] [--port PORT] [--name NAME]")
61
}
62
63
func runExpose() error {
@@ -79,11 +72,6 @@ func runExposeWithConfig() error {
72
if err != nil {
73
return fmt.Errorf("load config: %w", err)
74
}
82
- services, err := selectServices(cfg)
83
- if err != nil {
84
- return err
85
- }
86
-
75
relayDir := NewRelayDirectory(cfg.Relays)
76
ctx, cancel := context.WithCancel(context.Background())
77
defer cancel()
@@ -98,11 +86,11 @@ func runExposeWithConfig() error {
86
cancel()
87
}()
88
101
- errCh := make(chan error, len(services))
89
+ errCh := make(chan error, len(cfg.Services))
90
var wg sync.WaitGroup
91
104
- for _, svc := range services {
105
- service := svc
92
+ for i := range cfg.Services {
93
+ service := &cfg.Services[i]
94
wg.Add(1)
95
go func() {
96
defer wg.Done()
@@ -138,22 +126,13 @@ func runExposeWithFlags() error {
126
return fmt.Errorf("--relay must include at least one non-empty URL when --config is not provided")
127
}
128
141
- host := strings.TrimSpace(flagHost)
142
- if host == "" {
143
- host = "localhost"
144
- }
145
- port := strings.TrimSpace(flagPort)
146
- if port == "" {
147
- return fmt.Errorf("--port is required when --config is not provided")
148
- }
149
-
150
- target := net.JoinHostPort(host, port)
129
+ target := net.JoinHostPort(flagHost, flagPort)
130
service := &ServiceConfig{
131
Name: strings.TrimSpace(flagName),
132
Target: target,
154
- Protocols: []string{"http/1.1", "h2"},
133
RelayPreference: []string{"flags"},
134
}
135
+ applyServiceDefaults(service)
136
137
relayDir := NewRelayDirectory([]RelayConfig{
138
{
@@ -182,49 +161,39 @@ func runExposeWithFlags() error {
161
return nil
162
}
163
185
-func proxyConnection(ctx context.Context, svcCtx *serviceContext, relayConn net.Conn, connNum int) error {
164
+func proxyConnection(ctx context.Context, localAddr string, relayConn net.Conn) error {
165
defer relayConn.Close()
166
188
- // Connect to local service
189
- localConn, err := net.Dial("tcp", svcCtx.LocalAddr)
167
+ localConn, err := net.Dial("tcp", localAddr)
168
if err != nil {
191
- return fmt.Errorf("failed to connect to local service %s: %w", svcCtx.LocalAddr, err)
169
+ return fmt.Errorf("failed to connect to local service %s: %w", localAddr, err)
170
}
171
defer localConn.Close()
172
195
- // Bidirectional copy
173
errCh := make(chan error, 2)
197
- cancelCopy := make(chan struct{})
174
+ stopCh := make(chan struct{})
175
go func() {
176
select {
177
case <-ctx.Done():
178
relayConn.Close()
179
localConn.Close()
203
- case <-cancelCopy:
180
+ case <-stopCh:
181
}
182
}()
183
207
- // Relay -> Local
184
go func() {
185
_, err := io.Copy(localConn, relayConn)
186
errCh <- err
187
}()
188
213
- // Local -> Relay
189
go func() {
190
_, err := io.Copy(relayConn, localConn)
191
errCh <- err
192
}()
193
219
- // Wait for one direction to finish
194
err = <-errCh
221
-
222
- // Close both connections to stop the other goroutine
195
+ close(stopCh)
196
relayConn.Close()
224
- localConn.Close()
225
- close(cancelCopy)
226
-
227
- // Wait for other goroutine
197
<-errCh
198
199
return err
@@ -244,18 +213,7 @@ func runServiceTunnel(ctx context.Context, relayDir *RelayDirectory, service *Se
213
serviceName = fmt.Sprintf("tunnel-%s", leaseID[:8])
214
log.Info().Str("service", serviceName).Msg("No service name provided; generated automatically")
215
}
247
- svcCtx := &serviceContext{
248
- Name: serviceName,
249
- LocalAddr: localAddr,
250
- RelayServers: bootstrapServers,
251
- }
252
-
253
- log.Info().Str("service", serviceName).Msgf("Waiting for local service at %s (interval=%v)...", localAddr, time.Second)
254
- if err := waitForLocalService(localAddr, 0, time.Second); err != nil {
255
- return fmt.Errorf("service %s: %w", serviceName, err)
256
- }
257
- log.Info().Str("service", serviceName).Msgf("✓ Local service is reachable at %s", localAddr)
258
-
216
+ log.Info().Str("service", serviceName).Msgf("Local service is reachable at %s", localAddr)
217
log.Info().Str("service", serviceName).Msgf("Starting Portal Tunnel (%s)...", origin)
218
log.Info().Str("service", serviceName).Msgf(" Local: %s", localAddr)
219
log.Info().Str("service", serviceName).Msgf(" Relays: %s", strings.Join(bootstrapServers, ", "))
@@ -281,12 +239,11 @@ func runServiceTunnel(ctx context.Context, relayDir *RelayDirectory, service *Se
239
}()
240
241
log.Info().Str("service", serviceName).Msg("")
284
- log.Info().Str("service", serviceName).Msg("=== Service is now publicly accessible ===")
242
log.Info().Str("service", serviceName).Msg("Access via:")
243
log.Info().Str("service", serviceName).Msgf("- Name: /peer/%s", serviceName)
244
log.Info().Str("service", serviceName).Msgf("- Lease ID: /peer/%s", leaseID)
288
- relayHost := extractHost(bootstrapServers[0])
289
- log.Info().Str("service", serviceName).Msgf("- Example: http://%s/peer/%s", relayHost, serviceName)
245
+ log.Info().Str("service", serviceName).Msgf("- Example: http://%s/peer/%s", bootstrapServers[0], serviceName)
246
+
247
log.Info().Str("service", serviceName).Msg("")
248
249
connCount := 0
@@ -311,29 +268,17 @@ func runServiceTunnel(ctx context.Context, relayDir *RelayDirectory, service *Se
268
}
269
270
connCount++
314
- currentConnCount := connCount
315
- log.Info().Str("service", serviceName).Msgf("→ [#%d] New connection from %s", currentConnCount, relayConn.RemoteAddr())
271
+ log.Info().Str("service", serviceName).Msgf("→ [#%d] New connection from %s", connCount, relayConn.RemoteAddr())
272
273
connWG.Add(1)
318
- go func(relayConn net.Conn, connNum int) {
274
+ go func(relayConn net.Conn) {
275
defer connWG.Done()
320
- if err := proxyConnection(ctx, svcCtx, relayConn, connNum); err != nil {
321
- log.Error().Str("service", serviceName).Err(err).Int("conn", connNum).Msg("Proxy error")
276
+ if err := proxyConnection(ctx, localAddr, relayConn); err != nil {
277
+ log.Error().Str("service", serviceName).Err(err).Msg("Proxy error")
278
}
323
- log.Info().Str("service", serviceName).Msgf("← [#%d] Connection closed", connNum)
324
- }(relayConn, currentConnCount)
325
- }
326
-}
327
-
328
-func selectServices(cfg *TunnelConfig) ([]*ServiceConfig, error) {
329
- if len(cfg.Services) == 0 {
330
- return nil, fmt.Errorf("config has no services")
331
- }
332
- services := make([]*ServiceConfig, len(cfg.Services))
333
- for i := range cfg.Services {
334
- services[i] = &cfg.Services[i]
279
+ log.Info().Str("service", serviceName).Msg("Connection closed")
280
+ }(relayConn)
281
}
336
- return services, nil
282
}
283
284
func parseCommaSeparatedURLs(raw string) []string {
@@ -354,47 +299,3 @@ func parseCommaSeparatedURLs(raw string) []string {
299
300
return out
301
}
357
-
358
-func extractHost(wsURL string) string {
359
- // Simple extraction: ws://host:port/path -> host:port
360
- // Remove ws:// or wss://
361
- host := wsURL
362
- if len(host) > 5 && host[:5] == "ws://" {
363
- host = host[5:]
364
- } else if len(host) > 6 && host[:6] == "wss://" {
365
- host = host[6:]
366
- }
367
-
368
- // Remove path
369
- if idx := len(host); idx > 0 {
370
- for i, c := range host {
371
- if c == '/' {
372
- idx = i
373
- break
374
- }
375
- }
376
- host = host[:idx]
377
- }
378
-
379
- return host
380
-}
381
-
382
-// waitForLocalService tries to connect repeatedly until success or timeout.
383
-// If timeout == 0, it waits indefinitely.
384
-func waitForLocalService(localAddr string, timeout, interval time.Duration) error {
385
- deadline := time.Time{}
386
- if timeout > 0 {
387
- deadline = time.Now().Add(timeout)
388
- }
389
- for {
390
- conn, err := net.DialTimeout("tcp", localAddr, 2*time.Second)
391
- if err == nil {
392
- conn.Close()
393
- return nil
394
- }
395
- if !deadline.IsZero() && time.Now().After(deadline) {
396
- return fmt.Errorf("timeout waiting for local service at %s: %w", localAddr, err)
397
- }
398
- time.Sleep(interval)
399
- }
400
-}