mirror of
https://github.com/buger/goreplay.git
synced 2024-04-21 12:32:02 +00:00
Merge branch 'master' of https://github.com/buger/gor
This commit is contained in:
+45
-14
@@ -45,12 +45,18 @@ func CopyMulty(src io.Reader, writers ...io.Writer) (err error) {
|
||||
buf := make([]byte, 5*1024*1024)
|
||||
wIndex := 0
|
||||
modifier := NewHTTPModifier(&Settings.modifierConfig)
|
||||
filteredRequests := make(map[string]time.Time)
|
||||
filteredRequestsLastCleanTime := time.Now()
|
||||
|
||||
i := 0
|
||||
|
||||
for {
|
||||
nr, er := src.Read(buf)
|
||||
|
||||
if nr > 0 && len(buf) > nr {
|
||||
payload := buf[:nr]
|
||||
meta := payloadMeta(payload)
|
||||
requestID := string(meta[1])
|
||||
|
||||
_maxN := nr
|
||||
if nr > 500 {
|
||||
@@ -61,23 +67,31 @@ func CopyMulty(src io.Reader, writers ...io.Writer) (err error) {
|
||||
Debug("[EMITTER] input:", string(payload[0:_maxN]), nr, "from:", src)
|
||||
}
|
||||
|
||||
if modifier != nil && isRequestPayload(payload) {
|
||||
headSize := bytes.IndexByte(payload, '\n') + 1
|
||||
body := payload[headSize:]
|
||||
originalBodyLen := len(body)
|
||||
body = modifier.Rewrite(body)
|
||||
if modifier != nil {
|
||||
if isRequestPayload(payload) {
|
||||
headSize := bytes.IndexByte(payload, '\n') + 1
|
||||
body := payload[headSize:]
|
||||
originalBodyLen := len(body)
|
||||
body = modifier.Rewrite(body)
|
||||
|
||||
// If modifier tells to skip request
|
||||
if len(body) == 0 {
|
||||
continue
|
||||
}
|
||||
// If modifier tells to skip request
|
||||
if len(body) == 0 {
|
||||
filteredRequests[requestID] = time.Now()
|
||||
continue
|
||||
}
|
||||
|
||||
if originalBodyLen != len(body) {
|
||||
payload = append(payload[:headSize], body...)
|
||||
}
|
||||
if originalBodyLen != len(body) {
|
||||
payload = append(payload[:headSize], body...)
|
||||
}
|
||||
|
||||
if Settings.debug {
|
||||
Debug("[EMITTER] Rewrittern input:", len(payload), "First 500 bytes:", string(payload[0:_maxN]))
|
||||
if Settings.debug {
|
||||
Debug("[EMITTER] Rewritten input:", len(payload), "First 500 bytes:", string(payload[0:_maxN]))
|
||||
}
|
||||
} else {
|
||||
if _, ok := filteredRequests[requestID]; ok {
|
||||
delete(filteredRequests, requestID);
|
||||
continue
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -114,6 +128,23 @@ func CopyMulty(src io.Reader, writers ...io.Writer) (err error) {
|
||||
err = er
|
||||
break
|
||||
}
|
||||
|
||||
// Run GC on each 1000 request
|
||||
if i % 1000 == 0 {
|
||||
// Clean up filtered requests for which we didn't get a response to filter
|
||||
now := time.Now()
|
||||
if now.Sub(filteredRequestsLastCleanTime) > 60 * time.Second {
|
||||
for k, v := range filteredRequests {
|
||||
if now.Sub(v) > 60 * time.Second {
|
||||
delete(filteredRequests, k)
|
||||
}
|
||||
}
|
||||
filteredRequestsLastCleanTime = time.Now()
|
||||
}
|
||||
}
|
||||
|
||||
i++
|
||||
}
|
||||
|
||||
return err
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user