StreamFilter 辅助函数 - 过滤事件
(reader *stream.Reader[*session.Event], predicate func(*session.Event) bool)
| 231 | |
| 232 | // StreamFilter 辅助函数 - 过滤事件 |
| 233 | func StreamFilter(reader *stream.Reader[*session.Event], predicate func(*session.Event) bool) *stream.Reader[*session.Event] { |
| 234 | outReader, outWriter := stream.Pipe[*session.Event](10) |
| 235 | |
| 236 | go func() { |
| 237 | defer outWriter.Close() |
| 238 | for { |
| 239 | event, err := reader.Recv() |
| 240 | if err != nil { |
| 241 | if errors.Is(err, io.EOF) { |
| 242 | break |
| 243 | } |
| 244 | outWriter.Send(nil, err) |
| 245 | return |
| 246 | } |
| 247 | if event != nil && predicate(event) { |
| 248 | if outWriter.Send(event, nil) { |
| 249 | return |
| 250 | } |
| 251 | } |
| 252 | } |
| 253 | }() |
| 254 | |
| 255 | return outReader |
| 256 | } |
| 257 | |
| 258 | // streamConfig 流式执行配置 |
| 259 | type streamConfig struct { |
no test coverage detected