feat: restore portal tunnel

rabbitprincess committed Dec 5, 2025 at 20:05 UTC 08ad51d8af599a48c4c59d785aaf7619238075c0
4 files changed +401 -1
cmd/portal-tunnel/config.go new
+90
@@ -0,0 +1,90 @@
1 +package main
2 +
3 +import (
4 + "fmt"
5 + "os"
6 + "strings"
7 +
8 + "gopkg.in/yaml.v3"
9 + "gosuda.org/portal/sdk"
10 +)
11 +
12 +var defaultProtocols = []string{"http/1.1", "h2"}
13 +
14 +// ServiceConfig describes a local service exposed through the tunnel.
15 +type ServiceConfig struct {
16 + Name string `yaml:"name"`
17 + Target string `yaml:"target"`
18 + Protocols []string `yaml:"protocols"`
19 + Metadata sdk.Metadata `yaml:"metadata,omitempty"`
20 +}
21 +
22 +// TunnelConfig represents the YAML configuration schema for portal-tunnel.
23 +type TunnelConfig struct {
24 + Relays []string `yaml:"relays"`
25 + Service ServiceConfig `yaml:"service"`
26 +}
27 +
28 +// LoadConfig reads the YAML file at path, parses it into TunnelConfig, and validates it for single-service use.
29 +func LoadConfig(path string) (*TunnelConfig, error) {
30 + data, err := os.ReadFile(path)
31 + if err != nil {
32 + return nil, fmt.Errorf("read config: %w", err)
33 + }
34 +
35 + var cfg TunnelConfig
36 + if err := yaml.Unmarshal(data, &cfg); err != nil {
37 + return nil, fmt.Errorf("parse config: %w", err)
38 + }
39 + cfg.applyDefaults()
40 +
41 + if err := cfg.validate(); err != nil {
42 + return nil, err
43 + }
44 +
45 + return &cfg, nil
46 +}
47 +
48 +func (cfg *TunnelConfig) validate() error {
49 + var errs []string
50 +
51 + if len(cfg.Relays) == 0 {
52 + errs = append(errs, "at least one relay must be defined")
53 + }
54 + for i, url := range cfg.Relays {
55 + if strings.TrimSpace(url) == "" {
56 + errs = append(errs, fmt.Sprintf("relays[%d]: url cannot be empty", i))
57 + }
58 + }
59 +
60 + service := cfg.Service
61 + name := strings.TrimSpace(service.Name)
62 + if name == "" {
63 + errs = append(errs, "service: name is required")
64 + }
65 + target := strings.TrimSpace(service.Target)
66 + if target == "" {
67 + errs = append(errs, "service: target is required")
68 + }
69 + for i, proto := range service.Protocols {
70 + if strings.TrimSpace(proto) == "" {
71 + errs = append(errs, fmt.Sprintf("service.protocols[%d]: protocol cannot be empty", i))
72 + }
73 + }
74 +
75 + if len(errs) > 0 {
76 + return fmt.Errorf("invalid config:\n - %s", strings.Join(errs, "\n - "))
77 + }
78 +
79 + return nil
80 +}
81 +
82 +func (cfg *TunnelConfig) applyDefaults() {
83 + applyServiceDefaults(&cfg.Service)
84 +}
85 +
86 +func applyServiceDefaults(svc *ServiceConfig) {
87 + if len(svc.Protocols) == 0 {
88 + svc.Protocols = append([]string(nil), defaultProtocols...)
89 + }
90 +}
cmd/portal-tunnel/main.go new
+298
@@ -0,0 +1,298 @@
1 +package main
2 +
3 +import (
4 + "context"
5 + "fmt"
6 + "io"
7 + "net"
8 + "os"
9 + "os/signal"
10 + "strings"
11 + "sync"
12 + "syscall"
13 +
14 + "github.com/rs/zerolog/log"
15 + "github.com/spf13/cobra"
16 + "gosuda.org/portal/sdk"
17 + "gosuda.org/portal/utils"
18 +)
19 +
20 +var (
21 + flagConfigPath string
22 + flagRelayURLs string
23 + flagHost string
24 + flagPort string
25 + flagName string
26 + flagDesc string
27 + flagTags string
28 + flagThumbnail string
29 + flagOwner string
30 + flagHide bool
31 +)
32 +
33 +var rootCmd = &cobra.Command{
34 + Use: "portal-tunnel",
35 + Short: "Expose local services through Portal relay",
36 + RunE: func(cmd *cobra.Command, args []string) error {
37 + if flagConfigPath == "" {
38 + return runExposeWithFlags()
39 + }
40 + return runExposeWithConfig()
41 + },
42 +}
43 +
44 +func init() {
45 + rootCmd.Flags().StringVar(&flagConfigPath, "config", "", "Path to portal-tunnel config file")
46 + rootCmd.Flags().StringVar(&flagRelayURLs, "relay", "ws://localhost:4017/relay", "Portal relay server URLs when config is not provided (comma-separated)")
47 + rootCmd.Flags().StringVar(&flagHost, "host", "localhost", "Local host to proxy to when config is not provided")
48 + rootCmd.Flags().StringVar(&flagPort, "port", "4018", "Local port to proxy to when config is not provided")
49 + rootCmd.Flags().StringVar(&flagName, "name", "", "Service name when config is not provided (auto-generated if empty)")
50 + rootCmd.Flags().StringVar(&flagDesc, "description", "", "Service description metadata")
51 + rootCmd.Flags().StringVar(&flagTags, "tags", "", "Service tags metadata (comma-separated)")
52 + rootCmd.Flags().StringVar(&flagThumbnail, "thumbnail", "", "Service thumbnail URL metadata")
53 + rootCmd.Flags().StringVar(&flagOwner, "owner", "", "Service owner metadata")
54 + rootCmd.Flags().BoolVar(&flagHide, "hide", false, "Hide service from discovery (metadata)")
55 +}
56 +
57 +func main() {
58 + if err := rootCmd.Execute(); err != nil {
59 + os.Exit(1)
60 + }
61 +}
62 +
63 +func runExposeWithConfig() error {
64 + cfg, err := LoadConfig(flagConfigPath)
65 + if err != nil {
66 + return fmt.Errorf("load config: %w", err)
67 + }
68 +
69 + relayURLs := normalizeRelayURLs(cfg.Relays)
70 + if len(relayURLs) == 0 {
71 + return fmt.Errorf("config: relays must include at least one URL")
72 + }
73 +
74 + ctx, cancel := context.WithCancel(context.Background())
75 + defer cancel()
76 +
77 + // Graceful shutdown
78 + sigCh := make(chan os.Signal, 1)
79 + signal.Notify(sigCh, syscall.SIGINT, syscall.SIGTERM)
80 + go func() {
81 + <-sigCh
82 + log.Info().Msg("")
83 + log.Info().Msg("Shutting down tunnel...")
84 + cancel()
85 + }()
86 +
87 + if err := runServiceTunnel(ctx, relayURLs, &cfg.Service, fmt.Sprintf("config=%s", flagConfigPath)); err != nil {
88 + return err
89 + }
90 +
91 + log.Info().Msg("Tunnel stopped")
92 + return nil
93 +}
94 +
95 +func runExposeWithFlags() error {
96 + relayURLs := utils.ParseURLs(flagRelayURLs)
97 + if len(relayURLs) == 0 {
98 + return fmt.Errorf("--relay must include at least one non-empty URL when --config is not provided")
99 + }
100 +
101 + var metadata sdk.Metadata
102 + if strings.TrimSpace(flagDesc) != "" {
103 + metadata.Description = flagDesc
104 + }
105 + if strings.TrimSpace(flagTags) != "" {
106 + tags := strings.Split(flagTags, ",")
107 + for i := range tags {
108 + tags[i] = strings.TrimSpace(tags[i])
109 + }
110 + filtered := tags[:0]
111 + for _, t := range tags {
112 + if t != "" {
113 + filtered = append(filtered, t)
114 + }
115 + }
116 + metadata.Tags = filtered
117 + }
118 + if strings.TrimSpace(flagThumbnail) != "" {
119 + metadata.Thumbnail = flagThumbnail
120 + }
121 + if strings.TrimSpace(flagOwner) != "" {
122 + metadata.Owner = flagOwner
123 + }
124 + if flagHide {
125 + metadata.Hide = flagHide
126 + }
127 +
128 + target := net.JoinHostPort(flagHost, flagPort)
129 + service := &ServiceConfig{
130 + Name: strings.TrimSpace(flagName),
131 + Target: target,
132 + Metadata: metadata,
133 + }
134 + applyServiceDefaults(service)
135 +
136 + ctx, cancel := context.WithCancel(context.Background())
137 + defer cancel()
138 +
139 + sigCh := make(chan os.Signal, 1)
140 + signal.Notify(sigCh, syscall.SIGINT, syscall.SIGTERM)
141 + go func() {
142 + <-sigCh
143 + log.Info().Msg("")
144 + log.Info().Msg("Shutting down tunnel...")
145 + cancel()
146 + }()
147 +
148 + if err := runServiceTunnel(ctx, relayURLs, service, "flags"); err != nil {
149 + return err
150 + }
151 +
152 + log.Info().Msg("Tunnel stopped")
153 + return nil
154 +}
155 +
156 +func proxyConnection(ctx context.Context, localAddr string, relayConn net.Conn) error {
157 + defer relayConn.Close()
158 +
159 + localConn, err := net.Dial("tcp", localAddr)
160 + if err != nil {
161 + return fmt.Errorf("failed to connect to local service %s: %w", localAddr, err)
162 + }
163 + defer localConn.Close()
164 +
165 + errCh := make(chan error, 2)
166 + stopCh := make(chan struct{})
167 + go func() {
168 + select {
169 + case <-ctx.Done():
170 + relayConn.Close()
171 + localConn.Close()
172 + case <-stopCh:
173 + }
174 + }()
175 +
176 + go func() {
177 + _, err := io.Copy(localConn, relayConn)
178 + errCh <- err
179 + }()
180 +
181 + go func() {
182 + _, err := io.Copy(relayConn, localConn)
183 + errCh <- err
184 + }()
185 +
186 + err = <-errCh
187 + close(stopCh)
188 + relayConn.Close()
189 + <-errCh
190 +
191 + return err
192 +}
193 +
194 +func runServiceTunnel(ctx context.Context, relayURLs []string, service *ServiceConfig, origin string) error {
195 + localAddr := service.Target
196 + serviceName := strings.TrimSpace(service.Name)
197 + if len(relayURLs) == 0 {
198 + return fmt.Errorf("no relay URLs provided")
199 + }
200 + bootstrapServers := relayURLs
201 +
202 + cred := sdk.NewCredential()
203 + leaseID := cred.ID()
204 + if serviceName == "" {
205 + serviceName = fmt.Sprintf("tunnel-%s", leaseID[:8])
206 + log.Info().Str("service", serviceName).Msg("No service name provided; generated automatically")
207 + }
208 + log.Info().Str("service", serviceName).Msgf("Local service is reachable at %s", localAddr)
209 + log.Info().Str("service", serviceName).Msgf("Starting Portal Tunnel (%s)...", origin)
210 + log.Info().Str("service", serviceName).Msgf(" Local: %s", localAddr)
211 + log.Info().Str("service", serviceName).Msgf(" Relays: %s", strings.Join(bootstrapServers, ", "))
212 + log.Info().Str("service", serviceName).Msgf(" Lease ID: %s", leaseID)
213 +
214 + client, err := sdk.NewClient(func(c *sdk.ClientConfig) {
215 + c.BootstrapServers = bootstrapServers
216 + })
217 + if err != nil {
218 + return fmt.Errorf("service %s: failed to connect to relay: %w", serviceName, err)
219 + }
220 + defer client.Close()
221 +
222 + listener, err := client.Listen(cred, serviceName, service.Protocols,
223 + sdk.WithDescription(service.Metadata.Description),
224 + sdk.WithTags(service.Metadata.Tags),
225 + sdk.WithOwner(service.Metadata.Owner),
226 + sdk.WithThumbnail(service.Metadata.Thumbnail),
227 + sdk.WithHide(service.Metadata.Hide),
228 + )
229 + if err != nil {
230 + return fmt.Errorf("service %s: failed to register service: %w", serviceName, err)
231 + }
232 + defer listener.Close()
233 +
234 + go func() {
235 + <-ctx.Done()
236 + _ = listener.Close()
237 + }()
238 +
239 + log.Info().Str("service", serviceName).Msg("")
240 + log.Info().Str("service", serviceName).Msg("Access via:")
241 + log.Info().Str("service", serviceName).Msgf("- Name: /peer/%s", serviceName)
242 + log.Info().Str("service", serviceName).Msgf("- Lease ID: /peer/%s", leaseID)
243 + log.Info().Str("service", serviceName).Msgf("- Example: http://%s/peer/%s", bootstrapServers[0], serviceName)
244 +
245 + log.Info().Str("service", serviceName).Msg("")
246 +
247 + connCount := 0
248 + var connWG sync.WaitGroup
249 + defer connWG.Wait()
250 + for {
251 + select {
252 + case <-ctx.Done():
253 + return nil
254 + default:
255 + }
256 +
257 + relayConn, err := listener.Accept()
258 + if err != nil {
259 + select {
260 + case <-ctx.Done():
261 + return nil
262 + default:
263 + log.Error().Str("service", serviceName).Err(err).Msg("Failed to accept connection")
264 + continue
265 + }
266 + }
267 +
268 + connCount++
269 + log.Info().Str("service", serviceName).Msgf("→ [#%d] New connection from %s", connCount, relayConn.RemoteAddr())
270 +
271 + connWG.Add(1)
272 + go func(relayConn net.Conn) {
273 + defer connWG.Done()
274 + if err := proxyConnection(ctx, localAddr, relayConn); err != nil {
275 + log.Error().Str("service", serviceName).Err(err).Msg("Proxy error")
276 + }
277 + log.Info().Str("service", serviceName).Msg("Connection closed")
278 + }(relayConn)
279 + }
280 +}
281 +
282 +// normalizeRelayURLs trims, de-duplicates, and filters empty relay URLs.
283 +func normalizeRelayURLs(urls []string) []string {
284 + seen := map[string]struct{}{}
285 + var out []string
286 + for _, u := range urls {
287 + u = strings.TrimSpace(u)
288 + if u == "" {
289 + continue
290 + }
291 + if _, ok := seen[u]; ok {
292 + continue
293 + }
294 + seen[u] = struct{}{}
295 + out = append(out, u)
296 + }
297 + return out
298 +}
go.mod
+4 -1
@@ -7,19 +7,22 @@ require (
7 github.com/hashicorp/yamux v0.1.2
8 github.com/planetscale/vtprotobuf v0.6.0
9 github.com/rs/zerolog v1.34.0
10 + github.com/spf13/cobra v1.10.2
11 github.com/stretchr/testify v1.11.1
12 github.com/valyala/bytebufferpool v1.0.0
13 golang.org/x/crypto v0.44.0
14 golang.org/x/net v0.47.0
15 google.golang.org/protobuf v1.36.10
16 + gopkg.in/yaml.v3 v3.0.1
17 )
18
19 require (
20 github.com/davecgh/go-spew v1.1.1 // indirect
21 + github.com/inconshreveable/mousetrap v1.1.0 // indirect
22 github.com/mattn/go-colorable v0.1.14 // indirect
23 github.com/mattn/go-isatty v0.0.20 // indirect
24 github.com/pmezard/go-difflib v1.0.0 // indirect
25 + github.com/spf13/pflag v1.0.9 // indirect
26 golang.org/x/sys v0.38.0 // indirect
27 golang.org/x/text v0.31.0 // indirect
24 - gopkg.in/yaml.v3 v3.0.1 // indirect
28 )
go.sum
+9
@@ -1,4 +1,5 @@
1 github.com/coreos/go-systemd/v22 v22.5.0/go.mod h1:Y58oyj3AT4RCenI/lSvhwexgC+NSVTIJ3seZv2GcEnc=
2 +github.com/cpuguy83/go-md2man/v2 v2.0.6/go.mod h1:oOW0eioCTA6cOiMLiUPZOpcVxMig6NIQQ7OS05n1F4g=
3 github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c=
4 github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
5 github.com/godbus/dbus/v5 v5.0.4/go.mod h1:xhWf0FNVPg57R7Z0UbKHbJfkEywrmjJnf7w5xrFpKfA=
@@ -8,6 +9,8 @@ github.com/gorilla/websocket v1.5.3 h1:saDtZ6Pbx/0u+bgYQ3q96pZgCzfhKXGPqt7kZ72aN
9 github.com/gorilla/websocket v1.5.3/go.mod h1:YR8l580nyteQvAITg2hZ9XVh4b55+EU/adAjf1fMHhE=
10 github.com/hashicorp/yamux v0.1.2 h1:XtB8kyFOyHXYVFnwT5C3+Bdo8gArse7j2AQ0DA0Uey8=
11 github.com/hashicorp/yamux v0.1.2/go.mod h1:C+zze2n6e/7wshOZep2A70/aQU6QBRWJO/G6FT1wIns=
12 +github.com/inconshreveable/mousetrap v1.1.0 h1:wN+x4NVGpMsO7ErUn/mUI3vEoE6Jt13X2s0bqwp9tc8=
13 +github.com/inconshreveable/mousetrap v1.1.0/go.mod h1:vpF70FUmC8bwa3OWnCshd2FqLfsEA9PFc4w1p2J65bw=
14 github.com/mattn/go-colorable v0.1.13/go.mod h1:7S9/ev0klgBDR4GtXTXX8a3vIGJpMovkB8vQcUbaXHg=
15 github.com/mattn/go-colorable v0.1.14 h1:9A9LHSqF/7dyVVX6g0U9cwm9pG3kP9gSzcuIPHPsaIE=
16 github.com/mattn/go-colorable v0.1.14/go.mod h1:6LmQG8QLFO4G5z1gPvYEzlUgJ2wF+stgPZH1UqBm1s8=
@@ -23,10 +26,16 @@ github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZN
26 github.com/rs/xid v1.6.0/go.mod h1:7XoLgs4eV+QndskICGsho+ADou8ySMSjJKDIan90Nz0=
27 github.com/rs/zerolog v1.34.0 h1:k43nTLIwcTVQAncfCw4KZ2VY6ukYoZaBPNOE8txlOeY=
28 github.com/rs/zerolog v1.34.0/go.mod h1:bJsvje4Z08ROH4Nhs5iH600c3IkWhwp44iRc54W6wYQ=
29 +github.com/russross/blackfriday/v2 v2.1.0/go.mod h1:+Rmxgy9KzJVeS9/2gXHxylqXiyQDYRxCVz55jmeOWTM=
30 +github.com/spf13/cobra v1.10.2 h1:DMTTonx5m65Ic0GOoRY2c16WCbHxOOw6xxezuLaBpcU=
31 +github.com/spf13/cobra v1.10.2/go.mod h1:7C1pvHqHw5A4vrJfjNwvOdzYu0Gml16OCs2GRiTUUS4=
32 +github.com/spf13/pflag v1.0.9 h1:9exaQaMOCwffKiiiYk6/BndUBv+iRViNW+4lEMi0PvY=
33 +github.com/spf13/pflag v1.0.9/go.mod h1:McXfInJRrz4CZXVZOBLb0bTZqETkiAhM9Iw0y3An2Bg=
34 github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U=
35 github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U=
36 github.com/valyala/bytebufferpool v1.0.0 h1:GqA5TC/0021Y/b9FG4Oi9Mr3q7XYx6KllzawFIhcdPw=
37 github.com/valyala/bytebufferpool v1.0.0/go.mod h1:6bBcMArwyJ5K/AmCkWv1jt77kVWyCJ6HpOuEn7z0Csc=
38 +go.yaml.in/yaml/v3 v3.0.4/go.mod h1:DhzuOOF2ATzADvBadXxruRBLzYTpT36CKvDb3+aBEFg=
39 golang.org/x/crypto v0.44.0 h1:A97SsFvM3AIwEEmTBiaxPPTYpDC47w720rdiiUvgoAU=
40 golang.org/x/crypto v0.44.0/go.mod h1:013i+Nw79BMiQiMsOPcVCB5ZIJbYkerPrGnOa00tvmc=
41 golang.org/x/net v0.47.0 h1:Mx+4dIFzqraBXUugkia1OOvlD6LemFo1ALMHjrXDOhY=