-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathmonitor.go
More file actions
88 lines (83 loc) · 2.05 KB
/
Copy pathmonitor.go
File metadata and controls
88 lines (83 loc) · 2.05 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
package main
import (
"encoding/gob"
. "fmt"
"net"
"os"
"path/filepath"
"sync"
"time"
. "github.com/pkorotkov/diskloner/internal"
)
var progressLineUpdatePool = sync.Pool{
New: func() interface{} {
return &ProgressLineUpdate{}
},
}
func listenForMessage(ul *net.UnixListener, messages chan Message) {
// Register all types of messages.
gob.Register(&CloningMessage{})
gob.Register(&InquiringMessage{})
gob.Register(&CompletedMessage{})
gob.Register(&AbortedMessage{})
conn, err := ul.AcceptUnix()
if err != nil {
log.Error("got accept error: %s", err)
messages <- nil
return
}
defer conn.Close()
var mp Message
err = gob.NewDecoder(conn).Decode(&mp)
if err != nil {
log.Error("failed to read progress message: %s", err)
messages <- nil
return
}
messages <- mp
return
}
func MonitorStatus(quit chan os.Signal) {
// Create a UNIX domain socket server to wait connections from working cloners/inquirers.
file := filepath.Join(AppPath.ProgressDirectory, Sprintf("%d", time.Now().UnixNano()))
ul, err := net.ListenUnix("unix", &net.UnixAddr{file, "unix"})
if err != nil {
log.Error("failed to establish monitoring connection: %s", err)
return
}
defer os.Remove(file)
messages := make(chan Message)
sessions := make(map[string]int)
progress := NewProgress()
defer progress.Close()
for {
go listenForMessage(ul, messages)
select {
case message := <-messages:
if message != nil {
_, ok := sessions[message.UUID()]
if !ok {
sid := Sprintf("%s...%s", message.UUID()[0:6], message.UUID()[32:36])
sessions[message.UUID()] = progress.AddLine(NewProgressLine(sid))
}
l := progressLineUpdatePool.Get().(*ProgressLineUpdate)
l.Id = sessions[message.UUID()]
switch m := message.(type) {
case *CloningMessage:
l.State = "Cloning"
l.Current = m.CopiedBytes
l.Total = m.TotalBytes
case *CompletedMessage:
l.State = "Completed"
case *AbortedMessage:
l.State = "Aborted"
}
progress.UpdateLine(l)
progressLineUpdatePool.Put(l)
}
case <-quit:
return
}
}
return
}