From cc4cd1c010ddc77f33a56298fd8a9e57ba795654 Mon Sep 17 00:00:00 2001 From: Alec Scott Date: Sun, 30 Jan 2022 19:23:35 -0700 Subject: [PATCH 1/6] Initial commit of stream reuse --- cli/up.go | 55 ++++++++++++++++++++++++++++++++++++++++++++----------- go.mod | 1 + go.sum | 2 ++ 3 files changed, 47 insertions(+), 11 deletions(-) diff --git a/cli/up.go b/cli/up.go index b8a27ea..61d5af5 100644 --- a/cli/up.go +++ b/cli/up.go @@ -33,6 +33,8 @@ var ( // RevLookup allow quick lookups of an incoming stream // for security before accepting or responding to any data. RevLookup map[string]bool + // activeStreams is a map of active streams to a peer + activeStreams map[string]network.Stream ) // Up creates and brings up a Hyprspace Interface. @@ -180,23 +182,36 @@ func UpRun(r *cmd.Root, c *cmd.Sub) { fmt.Println("[+] Network Setup Complete...Waiting on Node Discovery") // Listen For New Packets on TUN Interface - packet := make([]byte, 1420) + activeStreams = make(map[string]network.Stream) + var packet = make([]byte, 1420) var stream network.Stream - var header *ipv4.Header + var ok bool var plen int + var dst string for { - plen, err = tunDev.Iface.Read(packet) - checkErr(err) - header, _ = ipv4.ParseHeader(packet) - _, ok := cfg.Peers[header.Dst.String()] + plen, err = tunDev.Iface.Read([]byte(packet)) + if err != nil { + log.Println(err) + continue + } + dst = net.IPv4(packet[16], packet[17], packet[18], packet[19]).String() + stream, ok = activeStreams[dst] if ok { - stream, err = host.NewStream(ctx, peerTable[header.Dst.String()], p2p.Protocol) + _, err = stream.Write([]byte(packet[:plen])) + if err == nil { + continue + } + stream.Close() + ok = false + } + if !ok { + stream, err = host.NewStream(ctx, peerTable[dst]) if err != nil { log.Println(err) continue } - stream.Write(packet[:plen]) - stream.Close() + stream.Write([]byte(packet[:plen])) + activeStreams[dst] = stream } } } @@ -236,9 +251,27 @@ func streamHandler(stream network.Stream) { // If the remote node ID isn't in the list of known nodes don't respond. if _, ok := RevLookup[stream.Conn().RemotePeer().Pretty()]; !ok { stream.Reset() + return + } + headers := io.LimitReader(stream, 20) + var err error + var header *ipv4.Header + var packetHeader = make([]byte, 20) + + for { + _, err = headers.Read([]byte(packetHeader)) + if err != nil { + stream.Close() + return + } + header, err = ipv4.ParseHeader(packetHeader) + if err != nil { + log.Println(err) + continue + } + tunDev.Iface.Write(packetHeader) + io.CopyN(tunDev.Iface.ReadWriteCloser, stream, int64(header.TotalLen)) } - io.Copy(tunDev.Iface.ReadWriteCloser, stream) - stream.Close() } func prettyDiscovery(ctx context.Context, node host.Host, peerTable map[string]peer.ID) { diff --git a/go.mod b/go.mod index fa589ed..70ddeeb 100644 --- a/go.mod +++ b/go.mod @@ -14,6 +14,7 @@ require ( github.com/libp2p/go-libp2p-quic-transport v0.12.0 github.com/libp2p/go-tcp-transport v0.2.8 github.com/multiformats/go-multiaddr v0.3.3 + github.com/songgao/packets v0.0.0-20160404182456-549a10cd4091 github.com/songgao/water v0.0.0-20200317203138-2b4b6d7c09d8 github.com/tcnksm/go-latest v0.0.0-20170313132115-e3007ae9052e github.com/vishvananda/netlink v1.1.0 diff --git a/go.sum b/go.sum index 793b409..988b6af 100644 --- a/go.sum +++ b/go.sum @@ -838,6 +838,8 @@ github.com/smartystreets/assertions v0.0.0-20180927180507-b2de0cb4f26d/go.mod h1 github.com/smartystreets/goconvey v1.6.4/go.mod h1:syvi0/a8iFYH4r/RixwvyeAJjdLS9QV7WQ/tjFTllLA= github.com/smola/gocompat v0.2.0/go.mod h1:1B0MlxbmoZNo3h8guHp8HztB3BSYR5itql9qtVc0ypY= github.com/soheilhy/cmux v0.1.4/go.mod h1:IM3LyeVVIOuxMH7sFAkER9+bJ4dT7Ms6E4xg4kGIyLM= +github.com/songgao/packets v0.0.0-20160404182456-549a10cd4091 h1:1zN6ImoqhSJhN8hGXFaJlSC8msLmIbX8bFqOfWLKw0w= +github.com/songgao/packets v0.0.0-20160404182456-549a10cd4091/go.mod h1:N20Z5Y8oye9a7HmytmZ+tr8Q2vlP0tAHP13kTHzwvQY= github.com/songgao/water v0.0.0-20200317203138-2b4b6d7c09d8 h1:TG/diQgUe0pntT/2D9tmUCz4VNwm9MfrtPr0SU2qSX8= github.com/songgao/water v0.0.0-20200317203138-2b4b6d7c09d8/go.mod h1:P5HUIBuIWKbyjl083/loAegFkfbFNx5i2qEP4CNbm7E= github.com/sony/gobreaker v0.4.1/go.mod h1:ZKptC7FHNvhBz7dN2LGjPVBz2sZJmc0/PkyDJOjmxWY= From c750bc1d92e93c211e9a054269b76e5b6588c836 Mon Sep 17 00:00:00 2001 From: Alec Date: Sun, 30 Jan 2022 21:18:07 -0700 Subject: [PATCH 2/6] Fixes --- cli/up.go | 27 +++++++++------------------ go.mod | 1 - go.sum | 2 -- 3 files changed, 9 insertions(+), 21 deletions(-) diff --git a/cli/up.go b/cli/up.go index 61d5af5..d0801fa 100644 --- a/cli/up.go +++ b/cli/up.go @@ -5,7 +5,6 @@ import ( "context" "errors" "fmt" - "io" "log" "net" "os" @@ -23,7 +22,6 @@ import ( "github.com/libp2p/go-libp2p-core/host" "github.com/libp2p/go-libp2p-core/network" "github.com/libp2p/go-libp2p-core/peer" - "golang.org/x/net/ipv4" ) var ( @@ -189,7 +187,7 @@ func UpRun(r *cmd.Root, c *cmd.Sub) { var plen int var dst string for { - plen, err = tunDev.Iface.Read([]byte(packet)) + plen, err = tunDev.Iface.Read(packet) if err != nil { log.Println(err) continue @@ -197,7 +195,7 @@ func UpRun(r *cmd.Root, c *cmd.Sub) { dst = net.IPv4(packet[16], packet[17], packet[18], packet[19]).String() stream, ok = activeStreams[dst] if ok { - _, err = stream.Write([]byte(packet[:plen])) + _, err = stream.Write(packet[:plen]) if err == nil { continue } @@ -205,12 +203,12 @@ func UpRun(r *cmd.Root, c *cmd.Sub) { ok = false } if !ok { - stream, err = host.NewStream(ctx, peerTable[dst]) + stream, err = host.NewStream(ctx, peerTable[dst], p2p.Protocol) if err != nil { log.Println(err) continue } - stream.Write([]byte(packet[:plen])) + stream.Write(packet[:plen]) activeStreams[dst] = stream } } @@ -253,24 +251,17 @@ func streamHandler(stream network.Stream) { stream.Reset() return } - headers := io.LimitReader(stream, 20) var err error - var header *ipv4.Header - var packetHeader = make([]byte, 20) - + var packet = make([]byte, 1420) + var plen int for { - _, err = headers.Read([]byte(packetHeader)) + plen, err = stream.Read(packet) if err != nil { + log.Println(err) stream.Close() return } - header, err = ipv4.ParseHeader(packetHeader) - if err != nil { - log.Println(err) - continue - } - tunDev.Iface.Write(packetHeader) - io.CopyN(tunDev.Iface.ReadWriteCloser, stream, int64(header.TotalLen)) + tunDev.Iface.Write(packet[:plen]) } } diff --git a/go.mod b/go.mod index 70ddeeb..fa589ed 100644 --- a/go.mod +++ b/go.mod @@ -14,7 +14,6 @@ require ( github.com/libp2p/go-libp2p-quic-transport v0.12.0 github.com/libp2p/go-tcp-transport v0.2.8 github.com/multiformats/go-multiaddr v0.3.3 - github.com/songgao/packets v0.0.0-20160404182456-549a10cd4091 github.com/songgao/water v0.0.0-20200317203138-2b4b6d7c09d8 github.com/tcnksm/go-latest v0.0.0-20170313132115-e3007ae9052e github.com/vishvananda/netlink v1.1.0 diff --git a/go.sum b/go.sum index 988b6af..793b409 100644 --- a/go.sum +++ b/go.sum @@ -838,8 +838,6 @@ github.com/smartystreets/assertions v0.0.0-20180927180507-b2de0cb4f26d/go.mod h1 github.com/smartystreets/goconvey v1.6.4/go.mod h1:syvi0/a8iFYH4r/RixwvyeAJjdLS9QV7WQ/tjFTllLA= github.com/smola/gocompat v0.2.0/go.mod h1:1B0MlxbmoZNo3h8guHp8HztB3BSYR5itql9qtVc0ypY= github.com/soheilhy/cmux v0.1.4/go.mod h1:IM3LyeVVIOuxMH7sFAkER9+bJ4dT7Ms6E4xg4kGIyLM= -github.com/songgao/packets v0.0.0-20160404182456-549a10cd4091 h1:1zN6ImoqhSJhN8hGXFaJlSC8msLmIbX8bFqOfWLKw0w= -github.com/songgao/packets v0.0.0-20160404182456-549a10cd4091/go.mod h1:N20Z5Y8oye9a7HmytmZ+tr8Q2vlP0tAHP13kTHzwvQY= github.com/songgao/water v0.0.0-20200317203138-2b4b6d7c09d8 h1:TG/diQgUe0pntT/2D9tmUCz4VNwm9MfrtPr0SU2qSX8= github.com/songgao/water v0.0.0-20200317203138-2b4b6d7c09d8/go.mod h1:P5HUIBuIWKbyjl083/loAegFkfbFNx5i2qEP4CNbm7E= github.com/sony/gobreaker v0.4.1/go.mod h1:ZKptC7FHNvhBz7dN2LGjPVBz2sZJmc0/PkyDJOjmxWY= From 8da27edd39edab3aae25c6e028a75e80f3fc06c2 Mon Sep 17 00:00:00 2001 From: Alec Date: Sun, 30 Jan 2022 21:46:21 -0700 Subject: [PATCH 3/6] Add stream reuse fixes and performance improvements --- cli/up.go | 13 +++++++------ 1 file changed, 7 insertions(+), 6 deletions(-) diff --git a/cli/up.go b/cli/up.go index d0801fa..bf7a8bc 100644 --- a/cli/up.go +++ b/cli/up.go @@ -30,7 +30,7 @@ var ( tunDev *tun.TUN // RevLookup allow quick lookups of an incoming stream // for security before accepting or responding to any data. - RevLookup map[string]bool + RevLookup map[string]string // activeStreams is a map of active streams to a peer activeStreams map[string]network.Stream ) @@ -91,9 +91,9 @@ func UpRun(r *cmd.Root, c *cmd.Sub) { } // Setup reverse lookup hash map for authentication. - RevLookup = make(map[string]bool, len(cfg.Peers)) - for _, id := range cfg.Peers { - RevLookup[id.ID] = true + RevLookup = make(map[string]string, len(cfg.Peers)) + for ip, id := range cfg.Peers { + RevLookup[id.ID] = ip } fmt.Println("[+] Creating TUN Device") @@ -200,9 +200,10 @@ func UpRun(r *cmd.Root, c *cmd.Sub) { continue } stream.Close() + delete(activeStreams, dst) ok = false } - if !ok { + if _, ok := peerTable[dst]; ok { stream, err = host.NewStream(ctx, peerTable[dst], p2p.Protocol) if err != nil { log.Println(err) @@ -257,8 +258,8 @@ func streamHandler(stream network.Stream) { for { plen, err = stream.Read(packet) if err != nil { - log.Println(err) stream.Close() + delete(activeStreams, RevLookup[stream.Conn().RemotePeer().Pretty()]) return } tunDev.Iface.Write(packet[:plen]) From 9cc7c03eb340c2d020ea9eb215f7d63f9c27a8e6 Mon Sep 17 00:00:00 2001 From: Alec Date: Sun, 30 Jan 2022 21:49:34 -0700 Subject: [PATCH 4/6] Update language in delete cmd --- cli/down.go | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/cli/down.go b/cli/down.go index 289614c..c2414ab 100644 --- a/cli/down.go +++ b/cli/down.go @@ -26,7 +26,8 @@ func DownRun(r *cmd.Root, c *cmd.Sub) { // Parse Command Args args := c.Args.(*DownArgs) - fmt.Println("[+] ip link delete dev " + args.InterfaceName) err := tun.Delete(args.InterfaceName) checkErr(err) + + fmt.Println("[+] deleted hyprspace " + args.InterfaceName + " daemon") } From 39b23f4e91bf8367342aa3f19eefc859354af776 Mon Sep 17 00:00:00 2001 From: Alec Date: Sun, 30 Jan 2022 21:50:10 -0700 Subject: [PATCH 5/6] Update go.mod dependencies --- go.mod | 1 - 1 file changed, 1 deletion(-) diff --git a/go.mod b/go.mod index fa589ed..3f51957 100644 --- a/go.mod +++ b/go.mod @@ -17,7 +17,6 @@ require ( github.com/songgao/water v0.0.0-20200317203138-2b4b6d7c09d8 github.com/tcnksm/go-latest v0.0.0-20170313132115-e3007ae9052e github.com/vishvananda/netlink v1.1.0 - golang.org/x/net v0.0.0-20210428140749-89ef3d95e781 golang.org/x/sys v0.0.0-20210603125802-9665404d3644 // indirect gopkg.in/yaml.v2 v2.4.0 honnef.co/go/tools v0.0.1-2020.1.4 // indirect From 275eb18ca8e42547947f2ac237f00cead2cbf922 Mon Sep 17 00:00:00 2001 From: Alec Scott Date: Mon, 31 Jan 2022 16:44:06 -0700 Subject: [PATCH 6/6] Add support for daemon locks for better process removal --- cli/down.go | 31 ++++++++++++++++++++++++++++++- cli/up.go | 36 +++++++++++++++++++++++------------- config/config.go | 17 +++++++++++++---- 3 files changed, 66 insertions(+), 18 deletions(-) diff --git a/cli/down.go b/cli/down.go index c2414ab..3fc34fc 100644 --- a/cli/down.go +++ b/cli/down.go @@ -2,6 +2,9 @@ package cli import ( "fmt" + "os" + "path/filepath" + "strconv" "github.com/DataDrake/cli-ng/v2/cmd" "github.com/hyprspace/hyprspace/tun" @@ -26,8 +29,34 @@ func DownRun(r *cmd.Root, c *cmd.Sub) { // Parse Command Args args := c.Args.(*DownArgs) - err := tun.Delete(args.InterfaceName) + // Parse Global Config Flag for Custom Config Path + configPath := r.Flags.(*GlobalFlags).Config + if configPath == "" { + configPath = "/etc/hyprspace/" + args.InterfaceName + ".yaml" + } + + // Read lock from file system to stop process. + lockPath := filepath.Join(filepath.Dir(configPath), args.InterfaceName+".lock") + out, err := os.ReadFile(lockPath) checkErr(err) + pid, err := strconv.Atoi(string(out)) + checkErr(err) + + process, err := os.FindProcess(pid) + checkErr(err) + + err0 := process.Signal(os.Interrupt) + + err1 := tun.Delete(args.InterfaceName) + + // Different types of systems may need the tun devices destroyed first or + // the process to exit first don't worry as long as one of these two has + // suceeded. + if err0 != nil && err1 != nil { + checkErr(err0) + checkErr(err1) + } + fmt.Println("[+] deleted hyprspace " + args.InterfaceName + " daemon") } diff --git a/cli/up.go b/cli/up.go index bf7a8bc..6ca0fe4 100644 --- a/cli/up.go +++ b/cli/up.go @@ -9,6 +9,7 @@ import ( "net" "os" "os/signal" + "path/filepath" "runtime" "strconv" "strings" @@ -158,6 +159,9 @@ func UpRun(r *cmd.Root, c *cmd.Sub) { go p2p.Discover(ctx, host, dht, peerTable) go prettyDiscovery(ctx, host, peerTable) + // Configure path for lock + lockPath := filepath.Join(filepath.Dir(cfg.Path), cfg.Interface.Name+".lock") + go func() { // Wait for a SIGINT or SIGTERM signal ch := make(chan os.Signal, 1) @@ -169,9 +173,19 @@ func UpRun(r *cmd.Root, c *cmd.Sub) { if err := host.Close(); err != nil { panic(err) } + + // Remove daemon lock from file system. + err = os.Remove(lockPath) + checkErr(err) + + // Exit the application. os.Exit(0) }() + // Write lock to filesystem to indicate an existing running daemon. + err = os.WriteFile(lockPath, []byte(fmt.Sprint(os.Getpid())), os.ModePerm) + checkErr(err) + // Bring Up TUN Device err = tunDev.Up() if err != nil { @@ -182,18 +196,14 @@ func UpRun(r *cmd.Root, c *cmd.Sub) { // Listen For New Packets on TUN Interface activeStreams = make(map[string]network.Stream) var packet = make([]byte, 1420) - var stream network.Stream - var ok bool - var plen int - var dst string for { - plen, err = tunDev.Iface.Read(packet) + plen, err := tunDev.Iface.Read(packet) if err != nil { log.Println(err) continue } - dst = net.IPv4(packet[16], packet[17], packet[18], packet[19]).String() - stream, ok = activeStreams[dst] + dst := net.IPv4(packet[16], packet[17], packet[18], packet[19]).String() + stream, ok := activeStreams[dst] if ok { _, err = stream.Write(packet[:plen]) if err == nil { @@ -203,8 +213,8 @@ func UpRun(r *cmd.Root, c *cmd.Sub) { delete(activeStreams, dst) ok = false } - if _, ok := peerTable[dst]; ok { - stream, err = host.NewStream(ctx, peerTable[dst], p2p.Protocol) + if peer, ok := peerTable[dst]; ok { + stream, err = host.NewStream(ctx, peer, p2p.Protocol) if err != nil { log.Println(err) continue @@ -215,7 +225,7 @@ func UpRun(r *cmd.Root, c *cmd.Sub) { } } -func createDaemon(cfg config.Config, out chan<- error) { +func createDaemon(cfg *config.Config, out chan<- error) { path, err := os.Executable() checkErr(err) // Create Pipe to monitor for daemon output. @@ -238,6 +248,8 @@ func createDaemon(cfg config.Config, out chan<- error) { count++ } } + + // Release the created daemon err = process.Release() checkErr(err) if count < len(cfg.Peers) { @@ -252,11 +264,9 @@ func streamHandler(stream network.Stream) { stream.Reset() return } - var err error var packet = make([]byte, 1420) - var plen int for { - plen, err = stream.Read(packet) + plen, err := stream.Read(packet) if err != nil { stream.Close() delete(activeStreams, RevLookup[stream.Conn().RemotePeer().Pretty()]) diff --git a/config/config.go b/config/config.go index 61b0946..f281f6e 100644 --- a/config/config.go +++ b/config/config.go @@ -8,6 +8,7 @@ import ( // Config is the main Configuration Struct for Hyprspace. type Config struct { + Path string `yaml:"path,omitempty"` Interface Interface `yaml:"interface"` Peers map[string]Peer `yaml:"peers"` } @@ -27,12 +28,12 @@ type Peer struct { } // Read initializes a config from a file. -func Read(path string) (result Config, err error) { +func Read(path string) (*Config, error) { in, err := os.ReadFile(path) if err != nil { - return + return nil, err } - result = Config{ + result := Config{ Interface: Interface{ Name: "hs0", ListenPort: 8001, @@ -41,6 +42,14 @@ func Read(path string) (result Config, err error) { PrivateKey: "", }, } + + // Read in config settings from file. err = yaml.Unmarshal(in, &result) - return + if err != nil { + return nil, err + } + + // Overwrite path of config to input. + result.Path = path + return &result, nil }