-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathwatcherd.go
More file actions
114 lines (92 loc) · 2.2 KB
/
Copy pathwatcherd.go
File metadata and controls
114 lines (92 loc) · 2.2 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
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
package main
import (
"io/ioutil"
"os"
"regexp"
"github.com/fsnotify/fsnotify"
)
type watcherd struct {
logPrefix string
slapdAccesslogDir string
queues *queues
}
func NewWatcherd(qs *queues) *watcherd {
wd := new(watcherd)
wd.logPrefix = "[watcherd] "
wd.queues = qs
// default parameters
wd.slapdAccesslogDir = "/var/log/slapd/cn=accesslog/"
return wd
}
func (wd *watcherd) run() {
fs := wd.getFileLists(wd.slapdAccesslogDir)
for i := range fs {
targetFilePath := wd.slapdAccesslogDir + fs[i].Name()
mes, err := wd.makeMessage(targetFilePath)
if err != nil {
continue
}
*wd.queues.parser <- mes
}
go wd.daemonize()
}
var reFile = regexp.MustCompile(`\.ldif$`)
func (wd *watcherd) daemonize() {
watcher, err := fsnotify.NewWatcher()
if err != nil {
myLoggerInfo(wd.logPrefix + err.Error())
}
defer watcher.Close()
err = watcher.Add(wd.slapdAccesslogDir)
if err != nil {
myLoggerInfo(wd.logPrefix + err.Error())
}
for {
select {
case event := <-watcher.Events:
if event.Op.String() != "CREATE" || !reFile.MatchString(event.Name) {
continue
}
mes, err := wd.makeMessage(event.Name)
if err != nil {
continue
}
*wd.queues.parser <- mes
case err := <-watcher.Errors:
myLoggerInfo(wd.logPrefix + "Error: " + err.Error())
}
}
}
func (wd *watcherd) makeMessage(filepath string) (message, error) {
myLoggerInfo(wd.logPrefix + "New file found: " + filepath)
file, err := os.Open(filepath)
defer file.Close()
mes := new(message)
if err != nil {
myLoggerInfo(wd.logPrefix + "Failed to open file: " + err.Error())
return *mes, err
}
content, err := ioutil.ReadAll(file)
if err != nil {
myLoggerInfo(wd.logPrefix + "Failed to read file: " + err.Error())
return *mes, err
}
mes.filepath = filepath
mes.content = string(content)
return *mes, nil
}
func (wd *watcherd) getFileLists(path string) []os.FileInfo {
fileLists, err := ioutil.ReadDir(path)
if err != nil {
myLoggerInfo(wd.logPrefix + "Directory cannot read: " + err.Error())
return []os.FileInfo{}
}
nf := []os.FileInfo{}
for _, f := range fileLists {
if f.IsDir() || !reFile.MatchString(f.Name()) {
continue
}
nf = append(nf, f)
}
return nf
}