From f5528bdb14a0560cb8b7c47d500065dd027e8b4e Mon Sep 17 00:00:00 2001 From: Will Andrews Date: Thu, 14 Dec 2023 18:26:40 +0000 Subject: [PATCH] Merge pull request #4 from willdot/ci Add Ci and fix all the lint --- .github/workflows/workflow.yaml | 26 ++++++++++++++++++++++++++ example/main.go | 17 +++++++++++++---- example/server/main.go | 5 ++++- pubsub/subscriber_test.go | 2 +- server/peer/peer.go | 8 ++++---- server/server.go | 44 ++++++++++++++++++++++++++++++-------------- server/server_test.go | 10 ++++++++-- server/topic.go | 12 ++---------- 8 file(s) changed, 88 insertion(s)(+), 36 deletion(s)(-) diff --git a/.github/workflows/workflow.yaml b/.github/workflows/workflow.yaml new file mode 100644 --- /dev/null +++ b/.github/workflows/workflow.yaml @@ -0,0 +1,26 @@ +name: Go package + +on: [push] + +jobs: + build: + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v3 + + - name: Set up Go + uses: actions/setup-go@v4 + with: + go-version: '1.21' + + - name: golangci-lint + run: curl -sSfL https://raw.githubusercontent.com/golangci/golangci-lint/master/install.sh | sh -s -- -b $(go env GOPATH)/bin + + - name: Build + run: go build -v ./... + + - name: lint + run: golangci-lint run + + - name: Test + run: go test ./... -p 1 -count=1 -v diff --git a/example/main.go b/example/main.go --- a/example/main.go +++ b/example/main.go @@ -16,7 +16,7 @@ func main() { consumeOnly = flag.Bool("consume-only", false, "just consumes (doesn't start server and doesn't publish)") flag.Parse() - if *consumeOnly == false { + if !*consumeOnly { go sendMessages() } @@ -24,9 +24,15 @@ sub, err := pubsub.NewSubscriber(":3000") if err != nil { panic(err) } - defer sub.Close() + + defer func() { + _ = sub.Close() + }() - sub.SubscribeToTopics([]string{"topic a"}) + err = sub.SubscribeToTopics([]string{"topic a"}) + if err != nil { + panic(err) + } consumer := sub.Consume(context.Background()) if consumer.Err != nil { @@ -44,7 +50,10 @@ publisher, err := pubsub.NewPublisher("localhost:3000") if err != nil { panic(err) } - defer publisher.Close() + + defer func() { + _ = publisher.Close() + }() // send some messages i := 0 diff --git a/example/server/main.go b/example/server/main.go --- a/example/server/main.go +++ b/example/server/main.go @@ -14,7 +14,10 @@ srv, err := server.New(":3000") if err != nil { log.Fatal(err) } - defer srv.Shutdown() + + defer func() { + _ = srv.Shutdown() + }() signals := make(chan os.Signal, 1) signal.Notify(signals, syscall.SIGTERM, syscall.SIGINT) diff --git a/pubsub/subscriber_test.go b/pubsub/subscriber_test.go --- a/pubsub/subscriber_test.go +++ b/pubsub/subscriber_test.go @@ -23,7 +23,7 @@ server, err := server.New(serverAddr) require.NoError(t, err) t.Cleanup(func() { - server.Shutdown() + _ = server.Shutdown() }) } diff --git a/server/peer/peer.go b/server/peer/peer.go --- a/server/peer/peer.go +++ b/server/peer/peer.go @@ -5,25 +5,25 @@ "net" "sync" ) -// Peer represents a remote connection to the server such as a publisher or subscriber +// Peer represents a remote connection to the server such as a publisher or subscriber. type Peer struct { conn net.Conn connMu sync.Mutex } -// New returns a new peer +// New returns a new peer. func New(conn net.Conn) *Peer { return &Peer{ conn: conn, } } -// Addr returns the peers connections address +// Addr returns the peers connections address. func (p *Peer) Addr() net.Addr { return p.conn.RemoteAddr() } -// ConnOpp represents a set of actions on a connection that can be used synchrnously +// ConnOpp represents a set of actions on a connection that can be used synchrnously. type ConnOpp func(conn net.Conn) error // RunConnOperation will run the provided operation. It ensures that it is the only operation that is being diff --git a/server/server.go b/server/server.go --- a/server/server.go +++ b/server/server.go @@ -5,10 +5,12 @@ "encoding/binary" "encoding/json" "errors" "fmt" + "io" "log/slog" "net" "strings" "sync" + "syscall" "time" "github.com/willdot/messagebroker/server/peer" @@ -51,7 +53,7 @@ Addr string lis net.Listener mu sync.Mutex - topics map[string]topic + topics map[string]*topic } // New creates and starts a new server @@ -63,7 +65,7 @@ } srv := &Server{ lis: lis, - topics: map[string]topic{}, + topics: map[string]*topic{}, } go srv.start() @@ -97,7 +99,9 @@ peer := peer.New(conn) action, err := readAction(peer, 0) if err != nil { - slog.Error("failed to read action from peer", "error", err, "peer", peer.Addr()) + if !errors.Is(err, io.EOF) { + slog.Error("failed to read action from peer", "error", err, "peer", peer.Addr()) + } return } @@ -123,15 +127,19 @@ // once the peers connection ends, it will be unsubscribed from all topics and returned for { action, err := readAction(peer, time.Millisecond*100) if err != nil { + // if the error is a timeout, it means the peer hasn't sent an action indicating it wishes to do something so sleep + // for a little bit to allow for other actions to happen on the connection var neterr net.Error if errors.As(err, &neterr) && neterr.Timeout() { - time.Sleep(time.Second) + time.Sleep(time.Millisecond * 500) continue } - // TODO: see if there's a way to check if the peers connection has been ended etc - slog.Error("failed to read action from subscriber", "error", err, "peer", peer.Addr()) + + if !errors.Is(err, io.EOF) { + slog.Error("failed to read action from subscriber", "error", err, "peer", peer.Addr()) + } - s.unsubscribePeerFromAllTopics(*peer) + s.unsubscribePeerFromAllTopics(peer) return } @@ -218,7 +226,7 @@ writeStatus(Error, "invalid topic data provided", conn) return nil } - s.unsubscribeToTopics(*peer, topics) + s.unsubscribeToTopics(peer, topics) writeStatus(Unsubscribed, "", conn) return nil @@ -239,6 +247,9 @@ op := func(conn net.Conn) error { dataLen, err := dataLength(conn) if err != nil { + if errors.Is(err, io.EOF) { + return nil + } slog.Error("failed to read data length", "error", err, "peer", peer.Addr()) writeStatus(Error, "invalid data length of data provided", conn) return nil @@ -325,13 +336,13 @@ s.topics[topicName] = t } -func (s *Server) unsubscribeToTopics(peer peer.Peer, topics []string) { +func (s *Server) unsubscribeToTopics(peer *peer.Peer, topics []string) { for _, topic := range topics { s.removeSubsciberFromTopic(topic, peer) } } -func (s *Server) removeSubsciberFromTopic(topicName string, peer peer.Peer) { +func (s *Server) removeSubsciberFromTopic(topicName string, peer *peer.Peer) { s.mu.Lock() defer s.mu.Unlock() @@ -343,7 +354,7 @@ delete(t.subscriptions, peer.Addr()) } -func (s *Server) unsubscribePeerFromAllTopics(peer peer.Peer) { +func (s *Server) unsubscribePeerFromAllTopics(peer *peer.Peer) { s.mu.Lock() defer s.mu.Unlock() @@ -357,7 +368,7 @@ s.mu.Lock() defer s.mu.Unlock() if topic, ok := s.topics[topicName]; ok { - return &topic + return topic } return nil @@ -367,7 +378,10 @@ func readAction(peer *peer.Peer, timeout time.Duration) (Action, error) { var action Action op := func(conn net.Conn) error { if timeout > 0 { - conn.SetReadDeadline(time.Now().Add(timeout)) + err := conn.SetReadDeadline(time.Now().Add(timeout)) + if err != nil { + slog.Error("failed to set connection read deadline", "error", err, "peer", peer.Addr()) + } } err := binary.Read(conn, binary.BigEndian, &action) @@ -406,7 +420,9 @@ func writeStatus(status Status, message string, conn net.Conn) { err := binary.Write(conn, binary.BigEndian, status) if err != nil { - slog.Error("failed to write status to peers connection", "error", err, "peer", conn.RemoteAddr()) + if !errors.Is(err, syscall.EPIPE) { + slog.Error("failed to write status to peers connection", "error", err, "peer", conn.RemoteAddr()) + } return } diff --git a/server/server_test.go b/server/server_test.go --- a/server/server_test.go +++ b/server/server_test.go @@ -25,7 +25,7 @@ srv, err := New(serverAddr) require.NoError(t, err) t.Cleanup(func() { - srv.Shutdown() + _ = srv.Shutdown() }) return srv @@ -33,7 +33,7 @@ } func createServerWithExistingTopic(t *testing.T, topicName string) *Server { srv := createServer(t) - srv.topics[topicName] = topic{ + srv.topics[topicName] = &topic{ name: topicName, subscriptions: make(map[net.Addr]subscriber), } @@ -61,6 +61,7 @@ expectedRes := Subscribed var resp Status err = binary.Read(conn, binary.BigEndian, &resp) + require.NoError(t, err) assert.Equal(t, expectedRes, int(resp)) @@ -106,6 +107,7 @@ expectedRes := Unsubscribed var resp Status err = binary.Read(conn, binary.BigEndian, &resp) + require.NoError(t, err) assert.Equal(t, expectedRes, int(resp)) @@ -160,6 +162,7 @@ expectedRes := Error var resp Status err = binary.Read(conn, binary.BigEndian, &resp) + require.NoError(t, err) assert.Equal(t, expectedRes, int(resp)) @@ -167,6 +170,7 @@ expectedMessage := "unknown action" var dataLen uint32 err = binary.Read(conn, binary.BigEndian, &dataLen) + require.NoError(t, err) assert.Equal(t, len(expectedMessage), int(dataLen)) buf := make([]byte, dataLen) @@ -196,6 +200,7 @@ expectedRes := Error var resp Status err = binary.Read(publisherConn, binary.BigEndian, &resp) + require.NoError(t, err) assert.Equal(t, expectedRes, int(resp)) @@ -203,6 +208,7 @@ expectedMessage := "topic data does not contain 'topic:' prefix" var dataLen uint32 err = binary.Read(publisherConn, binary.BigEndian, &dataLen) + require.NoError(t, err) assert.Equal(t, len(expectedMessage), int(dataLen)) buf := make([]byte, dataLen) diff --git a/server/topic.go b/server/topic.go --- a/server/topic.go +++ b/server/topic.go @@ -21,19 +21,11 @@ peer *peer.Peer currentOffset int } -func newTopic(name string) topic { - return topic{ +func newTopic(name string) *topic { + return &topic{ name: name, subscriptions: make(map[net.Addr]subscriber), } -} - -func (t *topic) removeSubscriber(addr net.Addr) { - t.mu.Lock() - defer t.mu.Unlock() - - slog.Info("removing subscriber", "peer", addr) - delete(t.subscriptions, addr) } func (t *topic) sendMessageToSubscribers(msgData []byte) { -- tangled.sh