Skip to content

Commit 459cfca

Browse files
committed
proxy fixes
1 parent 539a49a commit 459cfca

17 files changed

Lines changed: 1092 additions & 479 deletions

connect/connect_proxy_test.go

Lines changed: 402 additions & 0 deletions
Large diffs are not rendered by default.

connect/connect_router.go

Lines changed: 91 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,91 @@
1+
package main
2+
3+
import (
4+
"context"
5+
6+
"fmt"
7+
"net/http"
8+
"strings"
9+
"time"
10+
11+
"github.com/urnetwork/server"
12+
"github.com/urnetwork/server/model"
13+
)
14+
15+
type ConnectRouter struct {
16+
ctx context.Context
17+
cancel context.CancelFunc
18+
exchange *Exchange
19+
service string
20+
envService string
21+
connectHandler *ConnectHandler
22+
proxyConnectHandler *ProxyConnectHandler
23+
}
24+
25+
func NewConnectRouterWithDefaults(
26+
ctx context.Context,
27+
cancel context.CancelFunc,
28+
exchange *Exchange,
29+
) *ConnectRouter {
30+
return NewConnectRouter(
31+
ctx,
32+
cancel,
33+
exchange,
34+
DefaultConnectHandlerSettings(),
35+
DefaultProxyConnectHandlerSettings(),
36+
)
37+
}
38+
39+
func NewConnectRouter(
40+
ctx context.Context,
41+
cancel context.CancelFunc,
42+
exchange *Exchange,
43+
connectHandlerSettings *ConnectHandlerSettings,
44+
proxyConnectHandlerSettings *ProxyConnectHandlerSettings,
45+
) *ConnectRouter {
46+
handlerId := model.CreateNetworkClientHandler(ctx)
47+
48+
// update the heartbeat
49+
go server.HandleError(func() {
50+
defer cancel()
51+
for {
52+
select {
53+
case <-ctx.Done():
54+
return
55+
case <-time.After(model.NetworkClientHandlerHeartbeatTimeout):
56+
}
57+
// try again after unhandled errors. these signal a transient issue such as db load
58+
server.HandleError(func() {
59+
err := model.HeartbeatNetworkClientHandler(ctx, handlerId)
60+
if err != nil {
61+
// shut down
62+
cancel()
63+
}
64+
})
65+
}
66+
})
67+
68+
service := strings.ToLower(server.RequireService())
69+
envService := strings.ToLower(fmt.Sprintf("%s-%s", server.RequireEnv(), server.RequireService()))
70+
71+
connectHandler := NewConnectHandler(ctx, handlerId, exchange, connectHandlerSettings)
72+
proxyConnectHandler := NewProxyConnectHandler(ctx, handlerId, exchange, proxyConnectHandlerSettings)
73+
74+
return &ConnectRouter{
75+
ctx: ctx,
76+
cancel: cancel,
77+
exchange: exchange,
78+
service: service,
79+
envService: envService,
80+
connectHandler: connectHandler,
81+
proxyConnectHandler: proxyConnectHandler,
82+
}
83+
}
84+
85+
func (self *ConnectRouter) Connect(w http.ResponseWriter, r *http.Request) {
86+
self.connectHandler.Connect(w, r)
87+
}
88+
89+
func (self *ConnectRouter) ProxyConnect(w http.ResponseWriter, r *http.Request) {
90+
self.proxyConnectHandler.Connect(w, r)
91+
}

connect/connect_test.go

Lines changed: 3 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -460,7 +460,7 @@ func testConnect(
460460
ctx, cancel := context.WithCancel(context.Background())
461461
defer cancel()
462462

463-
service := "testConnect"
463+
service := "connect"
464464
block := "test"
465465

466466
clientIdA := server.NewId()
@@ -508,10 +508,6 @@ func testConnect(
508508
return server
509509
}
510510

511-
select {
512-
case <-time.After(1 * time.Second):
513-
}
514-
515511
hostPorts := map[string]int{}
516512
exchanges := map[string]*Exchange{}
517513
servers := map[string]*http.Server{}
@@ -549,7 +545,7 @@ func testConnect(
549545
}
550546

551547
select {
552-
case <-time.After(1 * time.Second):
548+
case <-time.After(2 * time.Second):
553549
}
554550

555551
randServer := func() (string, int) {
@@ -1535,7 +1531,7 @@ func Testing_NewControllerOutOfBandControl(ctx context.Context, clientId server.
15351531
}
15361532

15371533
func (self *controllerOutOfBandControl) SendControl(frames []*protocol.Frame, callback func(resultFrames []*protocol.Frame, err error)) {
1538-
server.HandleError(func() {
1534+
go server.HandleError(func() {
15391535
resultFrames, err := controller.ConnectControlFrames(self.ctx, self.clientId, frames)
15401536
callback(resultFrames, err)
15411537
})

connect/main.go

Lines changed: 8 additions & 50 deletions
Original file line numberDiff line numberDiff line change
@@ -2,14 +2,14 @@ package main
22

33
import (
44
"context"
5-
"fmt"
5+
// "fmt"
66
"net"
7-
"net/http"
7+
// "net/http"
88
"os"
99
"strconv"
10-
"strings"
10+
// "strings"
1111
"syscall"
12-
"time"
12+
// "time"
1313

1414
"github.com/docopt/docopt-go"
1515
// "github.com/prometheus/client_golang/prometheus"
@@ -19,7 +19,7 @@ import (
1919

2020
"github.com/urnetwork/connect"
2121
"github.com/urnetwork/server"
22-
"github.com/urnetwork/server/model"
22+
// "github.com/urnetwork/server/model"
2323
"github.com/urnetwork/server/router"
2424
)
2525

@@ -58,30 +58,6 @@ Options:
5858
exchange := NewExchangeFromEnvWithDefaults(ctx)
5959
defer exchange.Close()
6060

61-
handlerId := model.CreateNetworkClientHandler(ctx)
62-
63-
connectHandler := NewConnectHandlerWithDefaults(ctx, handlerId, exchange)
64-
proxyConnectHandler := NewProxyConnectHandlerWithDefaults(ctx, handlerId, exchange)
65-
// update the heartbeat
66-
go server.HandleError(func() {
67-
defer cancel()
68-
for {
69-
select {
70-
case <-ctx.Done():
71-
return
72-
case <-time.After(model.NetworkClientHandlerHeartbeatTimeout):
73-
}
74-
// try again after unhandled errors. these signal a transient issue such as db load
75-
server.HandleError(func() {
76-
err := model.HeartbeatNetworkClientHandler(ctx, handlerId)
77-
if err != nil {
78-
// shut down
79-
cancel()
80-
}
81-
})
82-
}
83-
})
84-
8561
// drain on sigterm
8662
go server.HandleError(func() {
8763
defer cancel()
@@ -92,30 +68,12 @@ Options:
9268
}
9369
})
9470

95-
service := strings.ToLower(server.RequireService())
96-
envService := strings.ToLower(fmt.Sprintf("%s-%s", server.RequireEnv(), server.RequireService()))
97-
98-
connectRouter := func(w http.ResponseWriter, r *http.Request) {
99-
host := r.Header.Get("X-Forwarded-Host")
100-
if host == "" {
101-
host = r.Header.Get("Host")
102-
}
103-
104-
sub := strings.ToLower(strings.SplitN(host, ".", 2)[0])
105-
switch sub {
106-
case service, envService:
107-
// the host is connect.<domain> or <env>-connect.<domain>
108-
connectHandler.Connect(w, r)
109-
default:
110-
// the host is <auth>.connect.<domain>
111-
proxyConnectHandler.Connect(w, r)
112-
}
113-
}
71+
connectRouter := NewConnectRouterWithDefaults(ctx, cancel, exchange)
11472

115-
// FIXME multiplex connectHandler.Connect and proxyConnectHandler.Connect
11673
routes := []*router.Route{
11774
router.NewRoute("GET", "/status", router.WarpStatus),
118-
router.NewRoute("*", "/", connectRouter),
75+
router.NewRoute("GET", "/", connectRouter.Connect),
76+
router.NewRoute("CONNECT", "", connectRouter.ProxyConnect),
11977
}
12078

12179
port, _ := opts.Int("--port")

connect/resident.go

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -823,7 +823,7 @@ func (self *ExchangeBuffer) WriteMessage(conn net.Conn, transferFrameBytes []byt
823823
}
824824

825825
func (self *ExchangeBuffer) ReadMessage(conn net.Conn) ([]byte, error) {
826-
conn.SetWriteDeadline(time.Now().Add(self.settings.ExchangeReadTimeout))
826+
conn.SetReadDeadline(time.Now().Add(self.settings.ExchangeReadTimeout))
827827
return self.framer.Read(conn)
828828
}
829829

@@ -1033,7 +1033,8 @@ func (self *ExchangeConnection) Run() {
10331033
if !ok {
10341034
return
10351035
}
1036-
if err := self.sendBuffer.WriteMessage(self.conn, message); err != nil {
1036+
err := self.sendBuffer.WriteMessage(self.conn, message)
1037+
if err != nil {
10371038
return
10381039
}
10391040
glog.V(2).Infof("[ecs] %s/%s@%s:%d\n", self.clientId, self.residentId, self.host, self.port)
@@ -1480,7 +1481,6 @@ type Resident struct {
14801481
clientForwardUnsub func()
14811482
}
14821483

1483-
// FIXME check for proxy config on (clientId, instanceId) and configure the mode
14841484
func NewResident(
14851485
ctx context.Context,
14861486
exchange *Exchange,

connect/resident_proxy.go

Lines changed: 13 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@ package main
44

55
import (
66
"context"
7+
"fmt"
78
"time"
89

910
"github.com/urnetwork/connect"
@@ -45,7 +46,7 @@ func NewResidentProxyDevice(
4546

4647
cancelCtx, cancel := context.WithCancel(ctx)
4748

48-
networkSpace := newExchangeNetworkSpace(exchange)
49+
networkSpace := sdk.NewPlatformNetworkSpace(ctx, server.RequireEnv(), server.RequireHost())
4950

5051
generatorFunc := func(specs []*connect.ProviderSpec) connect.MultiClientGenerator {
5152
return newExchangeGenerator(
@@ -99,13 +100,16 @@ func (self *ResidentProxyDevice) AddTun() (
99100

100101
tunCtx, tunCancel := context.WithCancel(self.ctx)
101102

102-
server.HandleError(func() {
103+
go server.HandleError(func() {
103104
defer tunCancel()
104105
for {
105106
select {
106107
case <-tunCtx.Done():
107108
return
108-
case packet := <-receive:
109+
case packet, ok := <-receive:
110+
if !ok {
111+
return
112+
}
109113
self.deviceLocal.SendPacketNoCopy(packet, int32(len(packet)))
110114
case <-time.After(self.exchange.settings.WriteTimeout):
111115
// drop
@@ -129,6 +133,7 @@ func (self *ResidentProxyDevice) AddTun() (
129133

130134
// note `send` is not closed. This channel is left open.
131135
}
136+
132137
return
133138
}
134139

@@ -138,7 +143,8 @@ func (self *ResidentProxyDevice) Close() {
138143
self.deviceLocal.Close()
139144
}
140145

141-
func newExchangeNetworkSpace(exchange *Exchange) *sdk.NetworkSpace {
142-
// FIXME
143-
return nil
144-
}
146+
// func newExchangeNetworkSpace(exchange *Exchange) *sdk.NetworkSpace {
147+
// // FIXME
148+
// return nil
149+
150+
// }

0 commit comments

Comments
 (0)