handleConnections.go 1.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566
  1. package server
  2. import (
  3. "log"
  4. "net/http"
  5. "strconv"
  6. "strings"
  7. "time"
  8. )
  9. // Handle incoming websocket connections
  10. func (server *Server) handleConnections(w http.ResponseWriter, r *http.Request) {
  11. conn, err := upgrader.Upgrade(w, r, nil)
  12. if err != nil {
  13. log.Println(err)
  14. return
  15. }
  16. // defer conn.Close()
  17. log.Println("Incoming connection")
  18. // Check for register
  19. var registered bool
  20. var msgs []string
  21. for !registered {
  22. _, message, err := conn.ReadMessage()
  23. if err != nil {
  24. log.Println("Error reading message:", err)
  25. break
  26. }
  27. msgs = server.parseMessage(message)
  28. if string(msgs[0]) == "register" {
  29. registered = true
  30. }
  31. }
  32. log.Println("Connected to agent")
  33. // Create var of type Agent
  34. cores_int, err := strconv.Atoi(msgs[2])
  35. if err != nil {
  36. log.Fatal("ERROR: failed strconv.Atoi(cpu cores):", err)
  37. }
  38. mem_max_string := strings.ReplaceAll(msgs[3], "GB", "")
  39. mem_max_float, err := strconv.ParseFloat(mem_max_string, 64)
  40. if err != nil {
  41. log.Fatal("ERROR: failed strconv.ParseFloat():", err)
  42. }
  43. agent := Agent{
  44. Name: msgs[1],
  45. Reg: time.Now(),
  46. Cores: cores_int,
  47. Mem_max: mem_max_float,
  48. TaskId: 0,
  49. conn: conn,
  50. }
  51. // Dump it into server.Agents
  52. server.Agents[msgs[1]] = &agent
  53. defer server.closeConnection(msgs[1])
  54. go server.readMessages(agent)
  55. server.writeMessages(conn)
  56. }