From d415a06126e5e0785002cdd2bc76cf1c2528b952 Mon Sep 17 00:00:00 2001 From: Will Date: Wed, 13 Dec 2023 20:53:41 +0000 Subject: [PATCH] fixed the peer read action pointer --- server/peer.go | 4 ++-- server/server.go | 18 ++++++++++-------- 2 files changed, 12 insertions(+), 10 deletions(-) diff --git a/server/peer.go b/server/peer.go index 18298f8..49f9f0a 100644 --- a/server/peer.go +++ b/server/peer.go @@ -32,8 +32,8 @@ type peer struct { connMu sync.Mutex } -func newPeer(conn net.Conn) peer { - return peer{ +func newPeer(conn net.Conn) *peer { + return &peer{ conn: conn, } } diff --git a/server/server.go b/server/server.go index b93fb66..ca2523b 100644 --- a/server/server.go +++ b/server/server.go @@ -70,7 +70,8 @@ func (s *Server) start() { func (s *Server) handleConn(conn net.Conn) { peer := newPeer(conn) - action, err := readAction(peer) + + action, err := readAction(peer, 0) if err != nil { slog.Error("failed to read action from peer", "error", err, "peer", peer.addr()) return @@ -78,11 +79,11 @@ func (s *Server) handleConn(conn net.Conn) { switch action { case Subscribe: - s.handleSubscribe(&peer) + s.handleSubscribe(peer) case Unsubscribe: - s.handleUnsubscribe(&peer) + s.handleUnsubscribe(peer) case Publish: - s.handlePublish(&peer) + s.handlePublish(peer) default: slog.Error("unknown action", "action", action, "peer", peer.addr()) writeStatus(Error, "unknown action", peer.conn) @@ -96,7 +97,7 @@ func (s *Server) handleSubscribe(peer *peer) { // keep handling the peers connection, getting the action from the peer when it wishes to do something else. // once the peers connection ends, it will be unsubscribed from all topics and returned for { - action, err := readAction(*peer) + action, err := readAction(peer, time.Millisecond*100) if err != nil { var neterr net.Error if errors.As(err, &neterr) && neterr.Timeout() { @@ -338,11 +339,12 @@ func (s *Server) getTopic(topicName string) *topic { return nil } -// TODO: work out why this can't take a pointer to the peer -func readAction(peer peer) (Action, error) { +func readAction(peer *peer, timeout time.Duration) (Action, error) { var action Action op := func(conn net.Conn) error { - conn.SetReadDeadline(time.Now().Add(time.Second)) + if timeout > 0 { + conn.SetReadDeadline(time.Now().Add(timeout)) + } err := binary.Read(conn, binary.BigEndian, &action) if err != nil { -- 2.51.2