diff --git a/network.go b/network.go index 7836f17ae..082047207 100644 --- a/network.go +++ b/network.go @@ -5,6 +5,7 @@ import ( "errors" "fmt" "github.com/dotcloud/docker/iptables" + "github.com/dotcloud/docker/proxy" "github.com/dotcloud/docker/utils" "log" "net" @@ -205,9 +206,9 @@ func getIfaceAddr(name string) (net.Addr, error) { // It keeps track of all mappings and is able to unmap at will type PortMapper struct { tcpMapping map[int]*net.TCPAddr - tcpProxies map[int]Proxy + tcpProxies map[int]proxy.Proxy udpMapping map[int]*net.UDPAddr - udpProxies map[int]Proxy + udpProxies map[int]proxy.Proxy iptables *iptables.Chain } @@ -222,7 +223,7 @@ func (mapper *PortMapper) Map(port int, backendAddr net.Addr) error { } } mapper.tcpMapping[port] = backendAddr.(*net.TCPAddr) - proxy, err := NewProxy(&net.TCPAddr{IP: net.IPv4(0, 0, 0, 0), Port: port}, backendAddr) + proxy, err := proxy.NewProxy(&net.TCPAddr{IP: net.IPv4(0, 0, 0, 0), Port: port}, backendAddr) if err != nil { mapper.Unmap(port, "tcp") return err @@ -238,7 +239,7 @@ func (mapper *PortMapper) Map(port int, backendAddr net.Addr) error { } } mapper.udpMapping[port] = backendAddr.(*net.UDPAddr) - proxy, err := NewProxy(&net.UDPAddr{IP: net.IPv4(0, 0, 0, 0), Port: port}, backendAddr) + proxy, err := proxy.NewProxy(&net.UDPAddr{IP: net.IPv4(0, 0, 0, 0), Port: port}, backendAddr) if err != nil { mapper.Unmap(port, "udp") return err @@ -300,9 +301,9 @@ func newPortMapper(config *DaemonConfig) (*PortMapper, error) { mapper := &PortMapper{ tcpMapping: make(map[int]*net.TCPAddr), - tcpProxies: make(map[int]Proxy), + tcpProxies: make(map[int]proxy.Proxy), udpMapping: make(map[int]*net.UDPAddr), - udpProxies: make(map[int]Proxy), + udpProxies: make(map[int]proxy.Proxy), iptables: chain, } return mapper, nil diff --git a/network_proxy_test.go b/proxy/network_proxy_test.go similarity index 99% rename from network_proxy_test.go rename to proxy/network_proxy_test.go index c27393eb5..b57c23c1e 100644 --- a/network_proxy_test.go +++ b/proxy/network_proxy_test.go @@ -1,4 +1,4 @@ -package docker +package proxy import ( "bytes" diff --git a/proxy/proxy.go b/proxy/proxy.go new file mode 100644 index 000000000..7a711f657 --- /dev/null +++ b/proxy/proxy.go @@ -0,0 +1,29 @@ +package proxy + +import ( + "fmt" + "net" +) + +type Proxy interface { + // Start forwarding traffic back and forth the front and back-end + // addresses. + Run() + // Stop forwarding traffic and close both ends of the Proxy. + Close() + // Return the address on which the proxy is listening. + FrontendAddr() net.Addr + // Return the proxied address. + BackendAddr() net.Addr +} + +func NewProxy(frontendAddr, backendAddr net.Addr) (Proxy, error) { + switch frontendAddr.(type) { + case *net.UDPAddr: + return NewUDPProxy(frontendAddr.(*net.UDPAddr), backendAddr.(*net.UDPAddr)) + case *net.TCPAddr: + return NewTCPProxy(frontendAddr.(*net.TCPAddr), backendAddr.(*net.TCPAddr)) + default: + panic(fmt.Errorf("Unsupported protocol")) + } +} diff --git a/proxy/tcp_proxy.go b/proxy/tcp_proxy.go new file mode 100644 index 000000000..18d68355e --- /dev/null +++ b/proxy/tcp_proxy.go @@ -0,0 +1,94 @@ +package proxy + +import ( + "github.com/dotcloud/docker/utils" + "io" + "log" + "net" + "syscall" +) + +type TCPProxy struct { + listener *net.TCPListener + frontendAddr *net.TCPAddr + backendAddr *net.TCPAddr +} + +func NewTCPProxy(frontendAddr, backendAddr *net.TCPAddr) (*TCPProxy, error) { + listener, err := net.ListenTCP("tcp", frontendAddr) + if err != nil { + return nil, err + } + // If the port in frontendAddr was 0 then ListenTCP will have a picked + // a port to listen on, hence the call to Addr to get that actual port: + return &TCPProxy{ + listener: listener, + frontendAddr: listener.Addr().(*net.TCPAddr), + backendAddr: backendAddr, + }, nil +} + +func (proxy *TCPProxy) clientLoop(client *net.TCPConn, quit chan bool) { + backend, err := net.DialTCP("tcp", nil, proxy.backendAddr) + if err != nil { + log.Printf("Can't forward traffic to backend tcp/%v: %v\n", proxy.backendAddr, err.Error()) + client.Close() + return + } + + event := make(chan int64) + var broker = func(to, from *net.TCPConn) { + written, err := io.Copy(to, from) + if err != nil { + err, ok := err.(*net.OpError) + // If the socket we are writing to is shutdown with + // SHUT_WR, forward it to the other end of the pipe: + if ok && err.Err == syscall.EPIPE { + from.CloseWrite() + } + } + to.CloseRead() + event <- written + } + utils.Debugf("Forwarding traffic between tcp/%v and tcp/%v", client.RemoteAddr(), backend.RemoteAddr()) + go broker(client, backend) + go broker(backend, client) + + var transferred int64 = 0 + for i := 0; i < 2; i++ { + select { + case written := <-event: + transferred += written + case <-quit: + // Interrupt the two brokers and "join" them. + client.Close() + backend.Close() + for ; i < 2; i++ { + transferred += <-event + } + goto done + } + } + client.Close() + backend.Close() +done: + utils.Debugf("%v bytes transferred between tcp/%v and tcp/%v", transferred, client.RemoteAddr(), backend.RemoteAddr()) +} + +func (proxy *TCPProxy) Run() { + quit := make(chan bool) + defer close(quit) + utils.Debugf("Starting proxy on tcp/%v for tcp/%v", proxy.frontendAddr, proxy.backendAddr) + for { + client, err := proxy.listener.Accept() + if err != nil { + utils.Debugf("Stopping proxy on tcp/%v for tcp/%v (%v)", proxy.frontendAddr, proxy.backendAddr, err.Error()) + return + } + go proxy.clientLoop(client.(*net.TCPConn), quit) + } +} + +func (proxy *TCPProxy) Close() { proxy.listener.Close() } +func (proxy *TCPProxy) FrontendAddr() net.Addr { return proxy.frontendAddr } +func (proxy *TCPProxy) BackendAddr() net.Addr { return proxy.backendAddr } diff --git a/network_proxy.go b/proxy/udp_proxy.go similarity index 55% rename from network_proxy.go rename to proxy/udp_proxy.go index fb91cc1b3..6fcb6afd8 100644 --- a/network_proxy.go +++ b/proxy/udp_proxy.go @@ -1,10 +1,8 @@ -package docker +package proxy import ( "encoding/binary" - "fmt" "github.com/dotcloud/docker/utils" - "io" "log" "net" "sync" @@ -17,103 +15,6 @@ const ( UDPBufSize = 2048 ) -type Proxy interface { - // Start forwarding traffic back and forth the front and back-end - // addresses. - Run() - // Stop forwarding traffic and close both ends of the Proxy. - Close() - // Return the address on which the proxy is listening. - FrontendAddr() net.Addr - // Return the proxied address. - BackendAddr() net.Addr -} - -type TCPProxy struct { - listener *net.TCPListener - frontendAddr *net.TCPAddr - backendAddr *net.TCPAddr -} - -func NewTCPProxy(frontendAddr, backendAddr *net.TCPAddr) (*TCPProxy, error) { - listener, err := net.ListenTCP("tcp", frontendAddr) - if err != nil { - return nil, err - } - // If the port in frontendAddr was 0 then ListenTCP will have a picked - // a port to listen on, hence the call to Addr to get that actual port: - return &TCPProxy{ - listener: listener, - frontendAddr: listener.Addr().(*net.TCPAddr), - backendAddr: backendAddr, - }, nil -} - -func (proxy *TCPProxy) clientLoop(client *net.TCPConn, quit chan bool) { - backend, err := net.DialTCP("tcp", nil, proxy.backendAddr) - if err != nil { - log.Printf("Can't forward traffic to backend tcp/%v: %v\n", proxy.backendAddr, err.Error()) - client.Close() - return - } - - event := make(chan int64) - var broker = func(to, from *net.TCPConn) { - written, err := io.Copy(to, from) - if err != nil { - err, ok := err.(*net.OpError) - // If the socket we are writing to is shutdown with - // SHUT_WR, forward it to the other end of the pipe: - if ok && err.Err == syscall.EPIPE { - from.CloseWrite() - } - } - to.CloseRead() - event <- written - } - utils.Debugf("Forwarding traffic between tcp/%v and tcp/%v", client.RemoteAddr(), backend.RemoteAddr()) - go broker(client, backend) - go broker(backend, client) - - var transferred int64 = 0 - for i := 0; i < 2; i++ { - select { - case written := <-event: - transferred += written - case <-quit: - // Interrupt the two brokers and "join" them. - client.Close() - backend.Close() - for ; i < 2; i++ { - transferred += <-event - } - goto done - } - } - client.Close() - backend.Close() -done: - utils.Debugf("%v bytes transferred between tcp/%v and tcp/%v", transferred, client.RemoteAddr(), backend.RemoteAddr()) -} - -func (proxy *TCPProxy) Run() { - quit := make(chan bool) - defer close(quit) - utils.Debugf("Starting proxy on tcp/%v for tcp/%v", proxy.frontendAddr, proxy.backendAddr) - for { - client, err := proxy.listener.Accept() - if err != nil { - utils.Debugf("Stopping proxy on tcp/%v for tcp/%v (%v)", proxy.frontendAddr, proxy.backendAddr, err.Error()) - return - } - go proxy.clientLoop(client.(*net.TCPConn), quit) - } -} - -func (proxy *TCPProxy) Close() { proxy.listener.Close() } -func (proxy *TCPProxy) FrontendAddr() net.Addr { return proxy.frontendAddr } -func (proxy *TCPProxy) BackendAddr() net.Addr { return proxy.backendAddr } - // A net.Addr where the IP is split into two fields so you can use it as a key // in a map: type connTrackKey struct { @@ -245,14 +146,3 @@ func (proxy *UDPProxy) Close() { func (proxy *UDPProxy) FrontendAddr() net.Addr { return proxy.frontendAddr } func (proxy *UDPProxy) BackendAddr() net.Addr { return proxy.backendAddr } - -func NewProxy(frontendAddr, backendAddr net.Addr) (Proxy, error) { - switch frontendAddr.(type) { - case *net.UDPAddr: - return NewUDPProxy(frontendAddr.(*net.UDPAddr), backendAddr.(*net.UDPAddr)) - case *net.TCPAddr: - return NewTCPProxy(frontendAddr.(*net.TCPAddr), backendAddr.(*net.TCPAddr)) - default: - panic(fmt.Errorf("Unsupported protocol")) - } -}