| 50 | } |
| 51 | |
| 52 | func (sh *streamHandler) streamWithMutex(streamWriter io.Writer, streamReader io.Reader, mutex *sync.Mutex) { |
| 53 | if streamWriter == nil || streamReader == nil { |
| 54 | sh.log.Debug("nil-stream", lager.Data{ |
| 55 | "streamWriter-nil": streamWriter == nil, |
| 56 | "streamReader-nil": streamReader == nil, |
| 57 | }) |
| 58 | return |
| 59 | } |
| 60 | |
| 61 | sh.wg.Add(1) |
| 62 | go func() { |
| 63 | mutex.Lock() |
| 64 | defer mutex.Unlock() |
| 65 | defer sh.wg.Done() |
| 66 | |
| 67 | _, err := io.Copy(streamWriter, streamReader) |
| 68 | if err != nil { |
| 69 | sh.log.Debug("failed-to-copy-stream-data", lager.Data{"error": err}) |
| 70 | } |
| 71 | }() |
| 72 | } |
| 73 | |
| 74 | func (sh *streamHandler) wait(decoder *json.Decoder) (int, error) { |
| 75 | for { |