summaryrefslogtreecommitdiff
path: root/internal/server
diff options
context:
space:
mode:
Diffstat (limited to 'internal/server')
-rw-r--r--internal/server/server.go39
-rw-r--r--internal/server/tcpserver.go68
2 files changed, 107 insertions, 0 deletions
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))
+ }
+ }
+}