From bb79f9f50a89507fb7db6d1d4c62630e431952f9 Mon Sep 17 00:00:00 2001 From: Paul Buetow Date: Mon, 22 May 2023 00:09:51 +0300 Subject: move server and client into their own packages --- internal/client.go | 10 ------- internal/client/client.go | 10 +++++++ internal/run.go | 6 ++-- internal/server.go | 39 ------------------------- internal/server/server.go | 39 +++++++++++++++++++++++++ internal/server/tcpserver.go | 68 ++++++++++++++++++++++++++++++++++++++++++++ internal/tcpserver.go | 68 -------------------------------------------- 7 files changed, 121 insertions(+), 119 deletions(-) delete mode 100644 internal/client.go create mode 100644 internal/client/client.go delete mode 100644 internal/server.go create mode 100644 internal/server/server.go create mode 100644 internal/server/tcpserver.go delete mode 100644 internal/tcpserver.go (limited to 'internal') diff --git a/internal/client.go b/internal/client.go deleted file mode 100644 index 6775ab2..0000000 --- a/internal/client.go +++ /dev/null @@ -1,10 +0,0 @@ -package internal - -import ( - "context" - - "codeberg.org/snonux/gorum/internal/config" -) - -func runClient(ctx context.Context, conf config.Config) { -} diff --git a/internal/client/client.go b/internal/client/client.go new file mode 100644 index 0000000..299623d --- /dev/null +++ b/internal/client/client.go @@ -0,0 +1,10 @@ +package client + +import ( + "context" + + "codeberg.org/snonux/gorum/internal/config" +) + +func Start(ctx context.Context, conf config.Config) { +} diff --git a/internal/run.go b/internal/run.go index d575ded..e1139d7 100644 --- a/internal/run.go +++ b/internal/run.go @@ -3,7 +3,9 @@ package internal import ( "context" + "codeberg.org/snonux/gorum/internal/client" "codeberg.org/snonux/gorum/internal/config" + "codeberg.org/snonux/gorum/internal/server" ) func Run(ctx context.Context, configFile string) { @@ -12,6 +14,6 @@ func Run(ctx context.Context, configFile string) { panic(err) } - go runClient(ctx, conf) - runServer(ctx, conf) + go client.Start(ctx, conf) + server.Start(ctx, conf) } diff --git a/internal/server.go b/internal/server.go deleted file mode 100644 index 3c37d6a..0000000 --- a/internal/server.go +++ /dev/null @@ -1,39 +0,0 @@ -package internal - -import ( - "context" - "log" - "time" - - "codeberg.org/snonux/gorum/internal/config" - "codeberg.org/snonux/gorum/internal/quorum" - "codeberg.org/snonux/gorum/internal/vote" -) - -func runServer(ctx context.Context, conf config.Config) { - ch := make(chan vote.Vote) - quo := make(quorum.Quorum) - - go func() { - for { - select { - case vote := <-ch: - quo.Vote(vote) - winner, err := quo.Winner(conf) - if err != nil { - log.Println(err.Error()) - continue - } - log.Printf("The current leader node is %s", winner) - case <-time.After(vote.Expiry): - quo.CleanExpired() - case <-ctx.Done(): - return - } - } - }() - - if err := startTcpServer(ctx, conf, ch); err != nil { - panic(err) - } -} diff --git a/internal/server/server.go b/internal/server/server.go new file mode 100644 index 0000000..55d7c9a --- /dev/null +++ b/internal/server/server.go @@ -0,0 +1,39 @@ +package server + +import ( + "context" + "log" + "time" + + "codeberg.org/snonux/gorum/internal/config" + "codeberg.org/snonux/gorum/internal/quorum" + "codeberg.org/snonux/gorum/internal/vote" +) + +func Start(ctx context.Context, conf config.Config) { + ch := make(chan vote.Vote) + quo := make(quorum.Quorum) + + go func() { + for { + select { + case vote := <-ch: + quo.Vote(vote) + winner, err := quo.Winner(conf) + if err != nil { + log.Println(err.Error()) + continue + } + log.Printf("The current leader node is %s", winner) + case <-time.After(vote.Expiry): + quo.CleanExpired() + case <-ctx.Done(): + return + } + } + }() + + if err := tcpServerStart(ctx, conf, ch); err != nil { + panic(err) + } +} diff --git a/internal/server/tcpserver.go b/internal/server/tcpserver.go new file mode 100644 index 0000000..4157c29 --- /dev/null +++ b/internal/server/tcpserver.go @@ -0,0 +1,68 @@ +package server + +import ( + "bufio" + "context" + "fmt" + "log" + "net" + + "codeberg.org/snonux/gorum/internal/config" + "codeberg.org/snonux/gorum/internal/vote" +) + +func tcpServerStart(ctx context.Context, conf config.Config, + ch chan<- vote.Vote) error { + + listener, err := net.Listen("tcp", conf.Address) + if err != nil { + return fmt.Errorf("Error starting TCP server: %s", err.Error()) + } + defer listener.Close() + + log.Printf("TCP server started on %s\n", conf.Address) + + for { + conn, err := listener.Accept() + if err != nil { + log.Printf("Error accepting connection: %s\n", err.Error()) + continue + } + + if !conf.IsParticipantWithLookup(conn.RemoteAddr().String(), net.LookupIP) { + log.Printf("Denying connection, peer not a participant: %v\n", conn.RemoteAddr().String()) + conn.Close() + continue + } + + log.Printf("Client connected: %s\n", conn.RemoteAddr().String()) + go handleConnection(ctx, conf, conn, ch) + } +} + +func handleConnection(ctx context.Context, conf config.Config, + conn net.Conn, ch chan<- vote.Vote) { + + defer conn.Close() + remoteAddr := conn.RemoteAddr().String() + + reader := bufio.NewReader(conn) + for { + select { + case <-ctx.Done(): + log.Printf("Server context done, disconnecting client %s\n", remoteAddr) + return + default: + message, err := reader.ReadString('\n') + if err != nil { + log.Printf("Client %s disconnected: %s\n", remoteAddr, err.Error()) + return + } + + log.Printf("Received message from %s: %s", remoteAddr, message) + ch <- vote.New(conf, remoteAddr, message) + + conn.Write([]byte(message)) + } + } +} diff --git a/internal/tcpserver.go b/internal/tcpserver.go deleted file mode 100644 index b1fd11e..0000000 --- a/internal/tcpserver.go +++ /dev/null @@ -1,68 +0,0 @@ -package internal - -import ( - "bufio" - "context" - "fmt" - "log" - "net" - - "codeberg.org/snonux/gorum/internal/config" - "codeberg.org/snonux/gorum/internal/vote" -) - -func startTcpServer(ctx context.Context, conf config.Config, - ch chan<- vote.Vote) error { - - listener, err := net.Listen("tcp", conf.Address) - if err != nil { - return fmt.Errorf("Error starting TCP server: %s", err.Error()) - } - defer listener.Close() - - log.Printf("TCP server started on %s\n", conf.Address) - - for { - conn, err := listener.Accept() - if err != nil { - log.Printf("Error accepting connection: %s\n", err.Error()) - continue - } - - if !conf.IsParticipantWithLookup(conn.RemoteAddr().String(), net.LookupIP) { - log.Printf("Denying connection, peer not a participant: %v\n", conn.RemoteAddr().String()) - conn.Close() - continue - } - - log.Printf("Client connected: %s\n", conn.RemoteAddr().String()) - go handleConnection(ctx, conf, conn, ch) - } -} - -func handleConnection(ctx context.Context, conf config.Config, - conn net.Conn, ch chan<- vote.Vote) { - - defer conn.Close() - remoteAddr := conn.RemoteAddr().String() - - reader := bufio.NewReader(conn) - for { - select { - case <-ctx.Done(): - log.Printf("Server context done, disconnecting client %s\n", remoteAddr) - return - default: - message, err := reader.ReadString('\n') - if err != nil { - log.Printf("Client %s disconnected: %s\n", remoteAddr, err.Error()) - return - } - - log.Printf("Received message from %s: %s", remoteAddr, message) - ch <- vote.New(conf, remoteAddr, message) - - conn.Write([]byte(message)) - } - } -} -- cgit v1.2.3