StreamLast 辅助函数 - 获取最后一个事件
(reader *stream.Reader[*session.Event])
| 203 | |
| 204 | // StreamLast 辅助函数 - 获取最后一个事件 |
| 205 | func StreamLast(reader *stream.Reader[*session.Event]) (*session.Event, error) { |
| 206 | var lastEvent *session.Event |
| 207 | var lastErr error |
| 208 | |
| 209 | for { |
| 210 | event, err := reader.Recv() |
| 211 | if err != nil { |
| 212 | if errors.Is(err, io.EOF) { |
| 213 | break |
| 214 | } |
| 215 | lastErr = err |
| 216 | continue |
| 217 | } |
| 218 | if event != nil { |
| 219 | lastEvent = event |
| 220 | } |
| 221 | } |
| 222 | |
| 223 | if lastErr != nil { |
| 224 | return lastEvent, lastErr |
| 225 | } |
| 226 | if lastEvent == nil { |
| 227 | return nil, errors.New("no events in stream") |
| 228 | } |
| 229 | return lastEvent, nil |
| 230 | } |
| 231 | |
| 232 | // StreamFilter 辅助函数 - 过滤事件 |
| 233 | func StreamFilter(reader *stream.Reader[*session.Event], predicate func(*session.Event) bool) *stream.Reader[*session.Event] { |
no test coverage detected