diff options
Diffstat (limited to 'internal/server')
| -rw-r--r-- | internal/server/server.go | 39 | ||||
| -rw-r--r-- | internal/server/tcpserver.go | 68 |
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)) + } + } +} |
