-
Notifications
You must be signed in to change notification settings - Fork 10
/
Copy pathapp.go
101 lines (91 loc) · 1.7 KB
/
app.go
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
package main
import (
"log"
"sync"
"github.com/andviro/grayproxy/pkg/gelf"
)
type listener interface {
Listen(dest chan<- gelf.Chunk) (err error)
}
type sender interface {
Send(data []byte) (err error)
}
type queue interface {
Put(data []byte) error
ReadChan() <-chan []byte
Close() error
}
type app struct {
inputURLs urlList
outputURLs urlList
verbose bool
sendTimeout int
dataDir string
ins []listener
outs []sender
sendErrors []error
q queue
}
func (app *app) enqueue(msgs <-chan gelf.Chunk) {
for msg := range msgs {
if err := app.q.Put(msg); err != nil {
panic(err)
}
}
}
func (app *app) dequeue() {
for msg := range app.q.ReadChan() {
var sent bool
if app.verbose {
log.Println(string(msg))
}
for i, out := range app.outs {
err := out.Send(msg)
if err != nil {
if app.sendErrors[i] == nil {
log.Printf("out %d: %v", i, err)
app.sendErrors[i] = err
}
continue
}
if app.sendErrors[i] != nil {
log.Printf("out %d is now alive", i)
app.sendErrors[i] = nil
}
sent = true
break
}
if !sent {
if app.dataDir == "" {
continue
}
if err := app.q.Put(msg); err != nil {
panic(err)
}
}
}
}
func (app *app) run() (err error) {
if err = app.configure(); err != nil {
return
}
defer app.q.Close()
msgs := make(chan gelf.Chunk, len(app.ins)*1000000)
defer close(msgs)
var wg sync.WaitGroup
for i := range app.ins {
wg.Add(1)
go func(i int) {
defer wg.Done()
err := app.ins[i].Listen(msgs)
if err != nil {
log.Printf("Input %d exited with error: %+v", i, err)
}
}(i)
}
go app.enqueue(msgs)
go app.dequeue()
log.Println("starting grayproxy")
wg.Wait()
return
}