readMessages.go 1.4 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152
  1. package server
  2. import (
  3. "fmt"
  4. "log"
  5. "strconv"
  6. "strings"
  7. )
  8. // Rewrite to become a general input handler
  9. func (server *Server) readMessages(agent Agent) {
  10. for {
  11. _, message, err := agent.conn.ReadMessage()
  12. if err != nil {
  13. log.Println("Error reading message:", err)
  14. break
  15. }
  16. // log.Printf("Received(%d): %s\n", messageType, message)
  17. msgs := server.parseMessage(message)
  18. if msgs[0] == "update" {
  19. cpu_usg, err := strconv.Atoi(msgs[1])
  20. if err != nil {
  21. fmt.Println("ERROR: converting string to int", err)
  22. }
  23. server.Agents[agent.Name].Cpu = append(server.Agents[agent.Name].Cpu, cpu_usg)
  24. server.Agents[agent.Name].Cpu = server.pruneIntSlice(server.Agents[agent.Name].Cpu)
  25. mem_usg_string := strings.ReplaceAll(msgs[2], "GB", "")
  26. mem_usg_float, err := strconv.ParseFloat(mem_usg_string, 64)
  27. if err != nil {
  28. log.Fatal("ERROR: failed strconv.ParseFloat():", err)
  29. }
  30. server.Agents[agent.Name].Mem_usg = append(server.Agents[agent.Name].Mem_usg, mem_usg_float)
  31. server.Agents[agent.Name].Mem_usg = server.pruneFloat64Slice(server.Agents[agent.Name].Mem_usg)
  32. taskId, err := strconv.Atoi(msgs[3])
  33. if err != nil {
  34. log.Println("ERROR: cannot convert taskId:", err)
  35. }
  36. server.Agents[agent.Name].TaskId = taskId
  37. }
  38. if msgs[0] == "deregister" {
  39. server.closeConnection(agent.Name)
  40. }
  41. // This sets the time-out for any incoming messages.
  42. // conn.SetReadDeadline(time.Now().Add(10 * time.Second))
  43. }
  44. }