Skip to content

Commit e44e055

Browse files
authored
Merge pull request #368 from urnetwork/proxy-api
Proxy api
2 parents cece70c + 9440fe1 commit e44e055

5 files changed

Lines changed: 246 additions & 28 deletions

File tree

connect/resident_proxy.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -124,7 +124,7 @@ func NewResidentProxyDevice(
124124
deviceLocal.SetConnectLocation(initialDeviceState.Location)
125125
}
126126

127-
tnet, err := proxy.CreateNetTUN(
127+
tnet, err := proxy.CreateNetTun(
128128
cancelCtx,
129129
settings.Mtu,
130130
)

http.go

Lines changed: 26 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@ package server
33
import (
44
"bytes"
55
"context"
6+
"crypto/tls"
67
"encoding/json"
78
// "errors"
89
"fmt"
@@ -390,6 +391,7 @@ func HttpListenAndServeWithReusePort(ctx context.Context, addr string, handler h
390391
if err != nil {
391392
return err
392393
}
394+
defer listener.Close()
393395

394396
server := &http.Server{
395397
Addr: addr,
@@ -401,3 +403,27 @@ func HttpListenAndServeWithReusePort(ctx context.Context, addr string, handler h
401403

402404
return server.Serve(listener)
403405
}
406+
407+
func HttpListenAndServeTlsWithReusePort(ctx context.Context, addr string, handler http.Handler, reusePort bool, httpServerOptions HttpServerOptions, tlsConfig *tls.Config) error {
408+
listenConfig := net.ListenConfig{}
409+
if reusePort {
410+
listenConfig.Control = SoReusePort
411+
}
412+
413+
listener, err := listenConfig.Listen(ctx, "tcp", addr)
414+
if err != nil {
415+
return err
416+
}
417+
defer listener.Close()
418+
419+
server := &http.Server{
420+
Addr: addr,
421+
Handler: handler,
422+
TLSConfig: tlsConfig,
423+
ReadTimeout: httpServerOptions.ReadTimeout,
424+
WriteTimeout: httpServerOptions.WriteTimeout,
425+
IdleTimeout: httpServerOptions.IdleTimeout,
426+
}
427+
428+
return server.ServeTLS(listener, "", "")
429+
}

model/network_client_model.go

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -126,12 +126,14 @@ type ProxyConfigResult struct {
126126
SocksProxyUrl string `json:"socks_proxy_url"`
127127
HttpProxyUrl string `json:"http_proxy_url"`
128128
HttpsProxyUrl string `json:"https_proxy_url"`
129+
ApiBaseUrl string `json:"api_base_url"`
129130
AuthToken string `json:"auth_token"`
130131
InstanceId server.Id `json:"instance_id"`
131132
ProxyHost string `json:"proxy_host"`
132133
HttpProxyPort int `json:"http_proxy_port"`
133134
HttpsProxyPort int `json:"https_proxy_port"`
134135
SocksProxyPort int `json:"socks_proxy_port"`
136+
ApiPort int `json:"api_port"`
135137
}
136138

137139
type ProxyAuthResult struct {
@@ -286,6 +288,7 @@ func AuthNetworkClient(
286288
socksProxyPort := 8080
287289
httpProxyPort := 8081
288290
httpsProxyPort := 8082
291+
apiPort := 8083
289292

290293
proxyHost := fmt.Sprintf("%s.%s", "cosmic", server.RequireDomain())
291294

@@ -314,16 +317,24 @@ func AuthNetworkClient(
314317
)
315318
}
316319

320+
apiBaseUrl := fmt.Sprintf(
321+
"https://api.%s:%d",
322+
proxyHost,
323+
apiPort,
324+
)
325+
317326
authClientResult.ProxyConfigResult = &ProxyConfigResult{
318327
SocksProxyUrl: socksProxyUrl,
319328
HttpProxyUrl: httpProxyUrl,
320329
HttpsProxyUrl: httpsProxyUrl,
330+
ApiBaseUrl: apiBaseUrl,
321331
AuthToken: strings.ToLower(signedProxyId),
322332
InstanceId: proxyDeviceConfig.InstanceId,
323333
ProxyHost: proxyHost,
324334
SocksProxyPort: socksProxyPort,
325335
HttpProxyPort: httpProxyPort,
326336
HttpsProxyPort: httpsProxyPort,
337+
ApiPort: apiPort,
327338
}
328339
} else {
329340
authClientResult.Error = &AuthNetworkClientError{

proxy/main.go

Lines changed: 155 additions & 27 deletions
Original file line numberDiff line numberDiff line change
@@ -2,13 +2,16 @@ package main
22

33
import (
44
"context"
5-
// "crypto/tls"
5+
"crypto/tls"
66
"encoding/base64"
7+
"encoding/json"
78
"fmt"
9+
"io"
810
"net"
911
"net/http"
1012
"net/netip"
1113
"os"
14+
"strconv"
1215
"strings"
1316
"syscall"
1417
"time"
@@ -21,6 +24,7 @@ import (
2124
"github.com/urnetwork/proxy"
2225
"github.com/urnetwork/server"
2326
"github.com/urnetwork/server/model"
27+
"github.com/urnetwork/server/router"
2428
)
2529

2630
// FIXME this is meant to be deployed with no lb and no containers
@@ -35,9 +39,11 @@ import (
3539
// -ldflags "-X main.Version=$WARP_VERSION-$WARP_VERSION_CODE"
3640
var Version string
3741

42+
// FIXME warp port blocks
3843
const ListenSocksPort = 8080
3944
const ListenHttpPort = 8081
4045
const ListenHttpsPort = 8082
46+
const ListenApiPort = 8083
4147

4248
func DefaultProxySettings() *ProxySettings {
4349
return &ProxySettings{
@@ -89,7 +95,15 @@ func main() {
8995
panic(err)
9096
}
9197

92-
glog.Infof("Listen socks5 (:%d), http (:%d), https (:%d)", ListenSocksPort, ListenHttpPort, ListenHttpsPort)
98+
glog.Infof("Listen api (:%d), socks5 (:%d), http (:%d), https (:%d)", ListenApiPort, ListenSocksPort, ListenHttpPort, ListenHttpsPort)
99+
100+
newApiServer(
101+
ctx,
102+
cancel,
103+
proxyDeviceManager,
104+
transportTls,
105+
settings,
106+
)
93107

94108
newSocks5Server(
95109
ctx,
@@ -242,31 +256,8 @@ func (self *httpServer) run() {
242256
}
243257
}
244258

245-
headerAuth := r.Header.Get("Proxy-Authorization")
246-
247-
bearerPrefix := "bearer "
248-
basicPrefix := "basic "
249-
250-
if len(bearerPrefix) < len(headerAuth) && strings.ToLower(headerAuth[:len(bearerPrefix)]) == bearerPrefix {
251-
signedProxyId := headerAuth[len(bearerPrefix):]
252-
proxyId, err := model.ParseSignedProxyId(signedProxyId)
253-
if err == nil {
254-
return proxyId, nil
255-
}
256-
} else if len(basicPrefix) < len(headerAuth) && strings.ToLower(headerAuth[:len(basicPrefix)]) == basicPrefix {
257-
// user:pass
258-
combinedSignedProxyId, err := base64.StdEncoding.DecodeString(headerAuth[len(basicPrefix):])
259-
if err != nil {
260-
return server.Id{}, err
261-
}
262-
signedProxyId := strings.SplitN(string(combinedSignedProxyId), ":", 2)[0]
263-
proxyId, err := model.ParseSignedProxyId(signedProxyId)
264-
if err == nil {
265-
return proxyId, nil
266-
}
267-
}
268-
269-
return server.Id{}, fmt.Errorf("Not authorized")
259+
authHeader := r.Header.Get("Proxy-Authorization")
260+
return authHeaderProxyId(authHeader)
270261
}
271262

272263
connectDial := func(r *http.Request, network string, addr string) (net.Conn, error) {
@@ -333,3 +324,140 @@ func (self *httpServer) run() {
333324
case <-self.ctx.Done():
334325
}
335326
}
327+
328+
type apiServer struct {
329+
ctx context.Context
330+
cancel context.CancelFunc
331+
proxyDeviceManager *ProxyDeviceManager
332+
transportTls *server.TransportTls
333+
settings *ProxySettings
334+
}
335+
336+
func newApiServer(
337+
ctx context.Context,
338+
cancel context.CancelFunc,
339+
proxyDeviceManager *ProxyDeviceManager,
340+
transportTls *server.TransportTls,
341+
settings *ProxySettings,
342+
) *apiServer {
343+
s := &apiServer{
344+
ctx: ctx,
345+
cancel: cancel,
346+
proxyDeviceManager: proxyDeviceManager,
347+
transportTls: transportTls,
348+
settings: settings,
349+
}
350+
351+
go server.HandleError(s.run, cancel)
352+
353+
return s
354+
}
355+
356+
func (self *apiServer) run() {
357+
defer self.cancel()
358+
359+
routes := []*router.Route{
360+
router.NewRoute("POST", "/warmup", self.HandleWarmup),
361+
}
362+
363+
reusePort := false
364+
365+
httpServerOptions := server.HttpServerOptions{
366+
ReadTimeout: 15 * time.Second,
367+
WriteTimeout: 30 * time.Second,
368+
IdleTimeout: 5 * time.Minute,
369+
}
370+
371+
tlsConfig := &tls.Config{
372+
GetConfigForClient: self.transportTls.GetTlsConfigForClient,
373+
}
374+
375+
err := server.HttpListenAndServeTlsWithReusePort(
376+
self.ctx,
377+
net.JoinHostPort("", strconv.Itoa(ListenApiPort)),
378+
router.NewRouter(self.ctx, routes),
379+
reusePort,
380+
httpServerOptions,
381+
tlsConfig,
382+
)
383+
if err != nil {
384+
panic(err)
385+
}
386+
}
387+
388+
type WarmupRequest struct {
389+
TimeoutSeconds int `json:"timeout_seconds,omitempty"`
390+
}
391+
392+
type WarmupResponse struct {
393+
Ready bool `json:"ready"`
394+
}
395+
396+
func (self *apiServer) HandleWarmup(w http.ResponseWriter, r *http.Request) {
397+
authHeader := r.Header.Get("Authorization")
398+
proxyId, err := authHeaderProxyId(authHeader)
399+
if err != nil {
400+
http.Error(w, err.Error(), http.StatusUnauthorized)
401+
return
402+
}
403+
404+
var warmupRequest WarmupRequest
405+
406+
defer r.Body.Close()
407+
bodyBytes, err := io.ReadAll(r.Body)
408+
409+
if 0 < len(bodyBytes) {
410+
err = json.Unmarshal(bodyBytes, &warmupRequest)
411+
if err != nil {
412+
http.Error(w, err.Error(), http.StatusInternalServerError)
413+
return
414+
}
415+
}
416+
// else use the default object
417+
418+
proxyDevice, err := self.proxyDeviceManager.OpenProxyDevice(proxyId)
419+
if err != nil {
420+
http.Error(w, err.Error(), http.StatusInternalServerError)
421+
return
422+
}
423+
424+
timeout := time.Duration(warmupRequest.TimeoutSeconds) * time.Second
425+
ready := proxyDevice.WaitForReady(r.Context(), timeout)
426+
427+
warmupResponse := &WarmupResponse{
428+
Ready: ready,
429+
}
430+
431+
out, err := json.Marshal(warmupResponse)
432+
if err != nil {
433+
http.Error(w, err.Error(), http.StatusInternalServerError)
434+
return
435+
}
436+
w.Write(out)
437+
}
438+
439+
func authHeaderProxyId(authHeader string) (server.Id, error) {
440+
bearerPrefix := "bearer "
441+
basicPrefix := "basic "
442+
443+
if len(bearerPrefix) < len(authHeader) && strings.ToLower(authHeader[:len(bearerPrefix)]) == bearerPrefix {
444+
signedProxyId := authHeader[len(bearerPrefix):]
445+
proxyId, err := model.ParseSignedProxyId(signedProxyId)
446+
if err == nil {
447+
return proxyId, nil
448+
}
449+
} else if len(basicPrefix) < len(authHeader) && strings.ToLower(authHeader[:len(basicPrefix)]) == basicPrefix {
450+
// user:pass
451+
combinedSignedProxyId, err := base64.StdEncoding.DecodeString(authHeader[len(basicPrefix):])
452+
if err != nil {
453+
return server.Id{}, err
454+
}
455+
signedProxyId := strings.SplitN(string(combinedSignedProxyId), ":", 2)[0]
456+
proxyId, err := model.ParseSignedProxyId(signedProxyId)
457+
if err == nil {
458+
return proxyId, nil
459+
}
460+
}
461+
462+
return server.Id{}, fmt.Errorf("Not authorized")
463+
}

proxy/proxy_device.go

Lines changed: 53 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -276,6 +276,59 @@ func (self *ProxyDevice) Tun() *proxy.Net {
276276
return self.tnet
277277
}
278278

279+
func (self *ProxyDevice) WaitForReady(ctx context.Context, timeout time.Duration) bool {
280+
readyCtx, readyCancel := context.WithCancel(self.ctx)
281+
defer readyCancel()
282+
go server.HandleError(func() {
283+
defer readyCancel()
284+
select {
285+
case <-self.ctx.Done():
286+
case <-ctx.Done():
287+
}
288+
})
289+
290+
windowStatus := self.deviceLocal.GetWindowStatus()
291+
if windowStatus.MinSatisfied {
292+
return true
293+
}
294+
295+
if timeout == 0 {
296+
return false
297+
}
298+
299+
sub := self.deviceLocal.AddWindowStatusChangeListener(&windowStatusChangeListener{
300+
callback: func(windowStatus *sdk.WindowStatus) {
301+
if windowStatus.MinSatisfied {
302+
readyCancel()
303+
}
304+
},
305+
})
306+
defer sub.Close()
307+
308+
if 0 < timeout {
309+
select {
310+
case <-readyCtx.Done():
311+
return true
312+
case <-time.After(timeout):
313+
return false
314+
}
315+
} else {
316+
select {
317+
case <-readyCtx.Done():
318+
return true
319+
}
320+
}
321+
}
322+
323+
// conforms to `sdk.WindowStatusChangeListener`
324+
type windowStatusChangeListener struct {
325+
callback func(*sdk.WindowStatus)
326+
}
327+
328+
func (self *windowStatusChangeListener) WindowStatusChanged(windowStatus *sdk.WindowStatus) {
329+
self.callback(windowStatus)
330+
}
331+
279332
func (self *ProxyDevice) UpdateActivity() bool {
280333
self.stateLock.Lock()
281334
defer self.stateLock.Unlock()

0 commit comments

Comments
 (0)