Start.go 2.3 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394
  1. package client
  2. import (
  3. "log"
  4. "os"
  5. "os/signal"
  6. "strconv"
  7. "time"
  8. "github.com/gorilla/websocket"
  9. "github.com/inhies/go-bytesize"
  10. "github.com/mackerelio/go-osstat/memory"
  11. "github.com/shirou/gopsutil/cpu"
  12. )
  13. func (client *Client) Start() {
  14. // Setup interrupt handling
  15. interrupt := make(chan os.Signal, 1)
  16. signal.Notify(interrupt, os.Interrupt)
  17. // Connect to the server
  18. err := client.connectToServer()
  19. if err != nil {
  20. log.Fatal("ERROR: Failed connecting to server - ", err)
  21. }
  22. defer client.conn.Close()
  23. // Agent => Manager:
  24. // - register: register an agent with the manager
  25. // - update: agents sends CPU and Mem metrics (does this every second)
  26. // - solution: solution of a task given by the manager
  27. // - deregister: deregister an agent from the manager
  28. // Register
  29. // One-time at the start
  30. // Fetch the hostname
  31. hostname, err := os.Hostname()
  32. if err != nil {
  33. log.Fatal("ERROR: Failed to register - cannot get hostname - ", err)
  34. }
  35. // Fetch the number of logical cores
  36. cpustats, err := cpu.Counts(true)
  37. if err != nil {
  38. log.Fatal("ERROR: Failed to register - unable to fetch number of cores - ", err)
  39. }
  40. // Fetch memory stats in bytes
  41. mem, err := memory.Get()
  42. if err != nil {
  43. log.Println("ERROR: Failed to register - fetching memory info - ", err)
  44. }
  45. // Use the ByteSize package to allow for memory calculations
  46. b := bytesize.New(float64(mem.Total))
  47. // Register this agent with the scheduler
  48. client.writeToServer("register;" + hostname + ";" + strconv.Itoa(cpustats) + ";" + b.String())
  49. // Setup updater logic
  50. // Runs continuesly in the background
  51. go client.statusUpdater()
  52. // Start handler for incoming messages
  53. done := make(chan struct{})
  54. go client.handleIncoming(done)
  55. for {
  56. select {
  57. case <-done:
  58. return
  59. case <-interrupt:
  60. log.Println("Interrupt received -- Stopping")
  61. // Dereg the agent before dying off
  62. client.writeToServer("deregister;" + hostname)
  63. // Cleanly close the connection by sending a close message and then
  64. // waiting (with timeout) for the server to close the connection.
  65. err := client.conn.WriteMessage(websocket.CloseMessage, websocket.FormatCloseMessage(websocket.CloseNormalClosure, ""))
  66. if err != nil {
  67. log.Println("ERROR: Close connection - ", err)
  68. return
  69. }
  70. select {
  71. case <-done:
  72. case <-time.After(time.Second):
  73. }
  74. return
  75. }
  76. }
  77. }