(ctx context.Context)
| 104 | func (r *reader) SetFilter(f filter.Filter) { r.filter = f } |
| 105 | |
| 106 | func (r *reader) Read(ctx context.Context) (chan *event.Event, chan error) { |
| 107 | errsc := make(chan error, 100) |
| 108 | eventsc := make(chan *event.Event, 2000) |
| 109 | go func() { |
| 110 | r.mu.Lock() |
| 111 | defer r.mu.Unlock() |
| 112 | for { |
| 113 | select { |
| 114 | case <-ctx.Done(): |
| 115 | return |
| 116 | default: |
| 117 | } |
| 118 | |
| 119 | var sec section.Section |
| 120 | if _, err := io.ReadFull(r.zr, sec[:]); err != nil { |
| 121 | if err != io.EOF { |
| 122 | errsc <- err |
| 123 | continue |
| 124 | } |
| 125 | break |
| 126 | } |
| 127 | |
| 128 | l := sec.Size() |
| 129 | buf := make([]byte, l) |
| 130 | if _, err := io.ReadFull(r.zr, buf); err != nil { |
| 131 | if err != io.EOF { |
| 132 | errsc <- err |
| 133 | continue |
| 134 | } |
| 135 | break |
| 136 | } |
| 137 | evt, err := event.NewFromCapture(buf, sec.Version()) |
| 138 | if err != nil { |
| 139 | errsc <- fmt.Errorf("fail to unmarshal event: %v", err) |
| 140 | capEventUnmarshalErrors.Add(1) |
| 141 | continue |
| 142 | } |
| 143 | capReadBytes.Add(int64(len(buf))) |
| 144 | // update the state of the ps/handle snapshotters |
| 145 | if err := r.updateSnapshotters(evt); err != nil { |
| 146 | log.Warn(err) |
| 147 | } |
| 148 | // push the event to the chanel |
| 149 | r.read(evt, eventsc) |
| 150 | } |
| 151 | }() |
| 152 | |
| 153 | return eventsc, errsc |
| 154 | } |
| 155 | |
| 156 | func (r *reader) Close() error { |
| 157 | r.mu.Lock() |
nothing calls this directly
no test coverage detected