-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathhttp_scanner.go
More file actions
143 lines (129 loc) · 3.34 KB
/
Copy pathhttp_scanner.go
File metadata and controls
143 lines (129 loc) · 3.34 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
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
package main
import (
"bufio"
"context"
"errors"
"fmt"
"io"
"sync"
)
type httpProbeFunc func(context.Context, string) (exists, selected bool, err error)
type httpScanner struct {
context context.Context
cancel context.CancelFunc
probe httpProbeFunc
queue chan string
jobs sync.WaitGroup
workers sync.WaitGroup
stats counters
outputMu sync.Mutex
output *bufio.Writer
outputErr error
errorsMu sync.Mutex
errorSamples []string
}
func newHTTPScanner(ctx context.Context, probe httpProbeFunc, workers int, output io.Writer) *httpScanner {
scannerCtx, cancel := context.WithCancel(ctx)
scanner := &httpScanner{
context: scannerCtx,
cancel: cancel,
probe: probe,
queue: make(chan string, workers*4),
output: bufio.NewWriterSize(output, 64*1024),
}
scanner.workers.Add(workers)
for i := 0; i < workers; i++ {
go scanner.worker()
}
return scanner
}
func (scanner *httpScanner) submit(ctx context.Context, candidate string) error {
scanner.jobs.Add(1)
select {
case scanner.queue <- candidate:
scanner.stats.checked.Add(1)
return nil
case <-scanner.context.Done():
scanner.jobs.Done()
return scanner.context.Err()
case <-ctx.Done():
scanner.jobs.Done()
return ctx.Err()
}
}
func (scanner *httpScanner) worker() {
defer scanner.workers.Done()
for candidate := range scanner.queue {
select {
case <-scanner.context.Done():
scanner.stats.canceled.Add(1)
scanner.jobs.Done()
continue
default:
}
exists, selected, err := scanner.probe(scanner.context, candidate)
if exists {
scanner.stats.existing.Add(1)
}
if errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded) && scanner.context.Err() != nil {
scanner.stats.canceled.Add(1)
} else if err != nil {
scanner.stats.errors.Add(1)
scanner.recordError(candidate, err)
} else if selected && scanner.emit(candidate) {
scanner.stats.found.Add(1)
}
scanner.jobs.Done()
}
}
func (scanner *httpScanner) finish() error {
defer scanner.cancel()
scanner.jobs.Wait()
close(scanner.queue)
scanner.workers.Wait()
scanner.outputMu.Lock()
defer scanner.outputMu.Unlock()
if err := scanner.output.Flush(); scanner.outputErr == nil {
scanner.outputErr = err
}
return scanner.outputErr
}
func (scanner *httpScanner) emit(candidate string) bool {
scanner.outputMu.Lock()
defer scanner.outputMu.Unlock()
if scanner.outputErr != nil {
return false
}
if _, err := fmt.Fprintln(scanner.output, candidate); err != nil {
scanner.outputErr = err
scanner.cancel()
return false
}
if err := scanner.output.Flush(); err != nil {
scanner.outputErr = err
scanner.cancel()
return false
}
return true
}
func (scanner *httpScanner) snapshot() stats {
return stats{
Checked: scanner.stats.checked.Load(),
Existing: scanner.stats.existing.Load(),
Found: scanner.stats.found.Load(),
Errors: scanner.stats.errors.Load(),
Canceled: scanner.stats.canceled.Load(),
}
}
func (scanner *httpScanner) recordError(candidate string, err error) {
scanner.errorsMu.Lock()
defer scanner.errorsMu.Unlock()
if len(scanner.errorSamples) < maxErrorSamples {
scanner.errorSamples = append(scanner.errorSamples, fmt.Sprintf("%q: %v", candidate, err))
}
}
func (scanner *httpScanner) diagnostics() []string {
scanner.errorsMu.Lock()
defer scanner.errorsMu.Unlock()
return append([]string(nil), scanner.errorSamples...)
}